A common assumption about the Agent Development Kit for Java is that, now that virtual threads make blocking cheap, its model interface is a plain blocking method, and streaming is a java.util.stream.Stream. That is not what the code does. Every seam in ADK Java, the model, tools, callbacks and the runner, is expressed in RxJava 3 reactive types, and synchronous code is supported only through specific adapters. Virtual threads change where your blocking code may run cheaply; they do not change the contract you must implement.
This article reads those contracts precisely, from the google/adk-java sources as of 2026-10-01: what each type promises, what empty and error mean, how to put a blocking HTTP client behind a reactive contract without the most common bug, and how to consume the reactive output from both blocking and streaming callers. It deliberately stays at the contract level. Mapping requests and responses for a new model is covered in implementing a custom LLM, and which threads run what is covered in the runtime thread model.
The contract table
| Seam | Signature returns | Meaning of the stream |
|---|---|---|
BaseLlm.generateContent(LlmRequest, boolean stream) | Flowable<LlmResponse> | one response when not streaming; possibly many when streaming |
BaseLlm.connect(LlmRequest) | BaseLlmConnection | live session: sends return Completable, receive() returns a Flowable |
BaseTool.runAsync(Map, ToolContext) | Single<Map<String, Object>> | exactly one result or an error |
BaseTool.processLlmRequest(...) | Completable | done or failed, no value |
| Model, agent and tool callbacks | Maybe<...> | empty continues; a value replaces the step |
| The *Sync callback twins | Optional<...> | same meaning, adapted with Maybe.fromOptional |
Runner.runAsync(...) | Flowable<Event> | every event of the invocation, in order |
The pattern is consistent: the cardinality of the reactive type encodes the protocol. A tool must produce exactly one result, so it is a Single. A callback may or may not intervene, so it is a Maybe. A model call may produce one or many chunks, so it is a Flowable. Once you read the types this way, most questions about how to implement a seam answer themselves.
What the stream flag promises
The Javadoc for generateContent says a non-streaming call yields one response, and a streaming call may yield more, with all responses treated as merged content parts. Two things follow. When stream is false, emit exactly one LlmResponse and complete; emitting several would be read as fragments of one answer. When it is true, emit chunks as they arrive, marking fragments with partial() and the end of the model's turn with turnComplete(); both accessors return Optional<Boolean>, so absence and false are distinct and you should set them explicitly.
Which mode is requested comes from RunConfig: its StreamingMode is NONE, SSE or BIDI, and the default is NONE. BIDI goes through connect instead of generateContent. A model with no live API can reasonably throw UnsupportedOperationException from connect, which fails loudly rather than pretending. The same builder sets setMaxLlmCalls, default 500, a budget that also bounds how badly a looping agent can go wrong.
Blocking code behind a reactive contract
Most model SDKs you will wrap are blocking. That is fine, provided the blocking happens at subscription time on a thread that may block. The single most common bug is eager evaluation: calling the client while building the reactive pipeline, then wrapping the already-computed result.
// WRONG: the HTTP call runs while the pipeline is being ASSEMBLED, on the caller's
// thread, before anyone subscribes, and again never on retry.
public Flowable<LlmResponse> generateContent(LlmRequest req, boolean stream) {
return Flowable.just(toResponse(client.complete(toWire(req))));
}
// RIGHT: nothing happens until subscription; retry re-executes the call.
public Flowable<LlmResponse> generateContent(LlmRequest req, boolean stream) {
if (!stream) {
return Flowable.fromCallable(() -> toResponse(client.complete(toWire(req))))
.subscribeOn(Schedulers.io()); // or a virtual-thread scheduler
}
return Flowable.<LlmResponse>create(emitter -> {
var sse = client.stream(toWire(req)); // blocking iterator of chunks
emitter.setCancellable(sse::close); // dispose() closes the socket
for (var chunk : sse) {
if (emitter.isCancelled()) return;
emitter.onNext(toPartialResponse(chunk)); // partial() = true
}
emitter.onNext(toFinalResponse(sse.summary())); // turnComplete() = true
emitter.onComplete();
}, BackpressureStrategy.BUFFER).subscribeOn(Schedulers.io());
}In the wrong version the network call runs on whatever thread called generateContent, before the framework subscribes, so timeouts and cancellation attached downstream cannot reach it, and a retry operator resubscribes to a constant and never calls the model again. In the right version fromCallable and create defer the work to subscription, setCancellable ties disposal to closing the connection, and subscribeOn moves the blocking to a scheduler meant for it. On JDK 21 and later you can build that scheduler from a virtual-thread executor with Schedulers.from(Executors.newVirtualThreadPerTaskExecutor()), which makes the blocking wait cheap without changing the contract.
toResponse and friends are your mapping functions; building LlmResponse objects correctly is the subject of the custom LLM article. The backpressure strategy matters less than it seems for text chunks, which are small; BUFFER is a reasonable default, but it does mean a stalled consumer holds the whole response in memory.
Sync callbacks and their async twins
Callbacks are where ADK Java offers genuine synchronous contracts. For every callback, before and after model, agent and tool, plus the on-error callbacks for model and tool, there is a Maybe-returning interface and a Sync twin that returns Optional. The agent builder adapts the twin with Maybe.fromOptional, so their semantics are identical: empty means carry on as normal, and a value replaces the step. A before-model callback returning a response skips the model call entirely, which is how response caches and policy refusals are built.
// Sync twin: return Optional.empty() to let the model call proceed.
LlmAgent agent = LlmAgent.builder()
.name("support")
.model("gemini-2.5-flash")
.beforeModelCallbackSync((ctx, request) ->
cache.lookup(request) // Optional<LlmResponse>
)
.build();
// Async form: needed when the check itself does I/O you want to keep non-blocking.
.beforeModelCallback((ctx, request) ->
Maybe.fromCallable(() -> remoteCache.get(key(request))) // null -> empty
.timeout(50, TimeUnit.MILLISECONDS)
.onErrorComplete() // cache down: proceed
)Use the sync twin when the callback is pure computation or a fast in-memory lookup. Use the async form when the check does I/O and you want timeouts and fallbacks expressed on the type, as above, where a slow or failing remote cache degrades to proceeding with the model call instead of failing the turn. The callback lattice itself is covered in ADK Java callbacks.
Tools: plain returns are adapted for you
BaseTool.runAsync returns Single<Map<String, Object>>, but if you write tools as annotated methods wrapped by FunctionTool, you choose. A plain object or map is converted to a map, falling back to a single result entry for scalars. A Single or Maybe is subscribed. A null or an empty Optional becomes an empty result. A Flowable is accepted only on the live, bidirectional path.
public class OrderTools {
// Plain return: FunctionTool converts the result to a map for you.
public static Map<String, Object> getOrder(@Schema(name = "orderId") String orderId) {
return Map.of("status", orders.status(orderId));
}
// Reactive return: Single (or Maybe) is subscribed by the framework.
public static Single<Map<String, Object>> refund(@Schema(name = "orderId") String orderId) {
return Single.fromCallable(() -> Map.<String, Object>of("refundId", payments.refund(orderId)));
}
}A plain return is the simpler contract, and with a blocking tool that is fine as long as the runtime's tool execution mode and thread choice account for it. Return a reactive type when the tool is already asynchronous or needs its own timeout, retry or cancellation on the type. See writing a FunctionTool for schemas and naming.
Consuming reactive output from sync and streaming callers
Runner.runAsync returns Flowable<Event>, and nothing happens until you subscribe. A blocking caller can collect it; a streaming caller should subscribe and forward. Both are legitimate, and the choice belongs to the edge of your system, not to the agent.
// 1. Blocking caller that wants the final answer (batch job, test, simple servlet).
List<Event> events = runner.runAsync(userId, sessionId, message)
.toList()
.blockingGet(); // fine on a virtual thread
// 2. Streaming caller (Spring SSE): subscribe, never block.
Disposable d = runner.runAsync(userId, sessionId, message, RunConfig.builder()
.setStreamingMode(RunConfig.StreamingMode.SSE).build())
.subscribe(ev -> sse.send(render(ev)), sse::completeWithError, sse::complete);
sse.onCompletion(d::dispose); // client left: cancel the run
// 3. Retry only before the first chunk: a mid-stream retry duplicates text the user saw.
Flowable<LlmResponse> guarded = Flowable.defer(() -> {
AtomicBoolean emitted = new AtomicBoolean(); // per subscription
AtomicInteger attempt = new AtomicInteger();
return Flowable.defer(() -> llm.generateContent(req, true))
.doOnNext(r -> emitted.set(true))
.timeout(30, TimeUnit.SECONDS)
.retryWhen(errs -> errs.flatMap(e ->
emitted.get() || attempt.incrementAndGet() > 3
? Flowable.<Long>error(e) // surface the real error
: Flowable.timer(attempt.get() * 500L, TimeUnit.MILLISECONDS)));
});Blocking with blockingGet is acceptable on a virtual thread or in a batch job, and harmful on an event-loop thread such as a Netty worker, where it stalls every other connection on that loop. For streaming, wire the client going away to dispose(), which propagates upstream to the model call and to the setCancellable hook, so an abandoned request stops consuming tokens. The third snippet enforces the retry rule for streams. The inner defer makes each retry a fresh call, but if the failure arrives after chunks were already forwarded, resubscribing would replay the answer from the start, so the handler retries only while nothing has been emitted, backs off between attempts, and after the third failure re-raises the real error instead of completing quietly with an empty answer. The outer defer gives each subscription its own flags. Buffer until complete instead when you need retries after output has started.
Worked example: one turn, two contracts
Consider a support agent with a before-model cache and one tool, called from a Spring controller in SSE mode. The runner subscribes to the agent flow. The sync cache callback returns empty on a miss, so generateContent is subscribed with stream true. Our adapter opens the HTTP stream on an I/O thread and emits three partial responses of text, then a final response whose content is a function call to getOrder. The runtime runs the tool; its plain map return is adapted into a Single, and the function response goes back to the model in a second generateContent call, which streams the final answer and marks the turn complete. The controller forwarded every event as it arrived.
If the user closes the tab during the second call, the SSE completion hook disposes the subscription, disposal reaches our create emitter, and setCancellable closes the socket. Had the adapter used Flowable.just with an eager call, the whole second response would already have been generated and billed before disposal could do anything. The contract, not the thread, is what made cancellation work.
Failure modes
- Eager work in assembly.
justaround a blocking call: runs too early, cannot be cancelled, and never retries. UsefromCallable,deferorcreate. - Blocking on the wrong thread. Blocking inside a reactive chain without
subscribeOnstalls whichever thread subscribed, possibly an event loop. - Several responses when not streaming. Downstream code treats them as fragments of one answer. Emit exactly one when
streamis false. - Errors as data, or data as errors. Transport failures should be
onError; a model that answered but refused or was cut off is a response carryingerrorCode()orerrorMessage(). Mixing them breaks retry and error callbacks. - Swallowed disposal. A streaming adapter without
setCancellablekeeps reading, and paying, after the client left. - Mid-stream retry. Resubscribing after partial output duplicates text for the user.
Choosing sync or async at each seam
A simple rule covers most cases. Implement the framework's types exactly at the model and tool seams. Inside them, write blocking code if it is simpler, but defer it and give it a thread that may block. Use sync callback twins for pure logic and async ones for I/O. At the edge, block only on threads built for it, such as virtual threads, and stream everything a user watches. The reactive types are the contract; sync code is an implementation detail you are allowed to hide inside them.
What to do next
- Grep your adapters and tools for
Flowable.just,Single.justandMaybe.justaround I/O, and replace them with deferred forms. - Check your model adapter emits exactly one response when
streamis false, and setspartialandturnCompleteexplicitly when true. - Add
setCancellableto every streaming adapter and test that disposing the runner's subscription closes the HTTP connection. - Move pure-logic callbacks to the
Synctwins; give I/O callbacks a timeout and a fallback. - Audit every
blockingGetandblockingSubscribeand confirm none runs on an event-loop thread. - Decide your retry policy per seam, and forbid retries after the first streamed chunk.
- Read streaming agent architecture next for backpressure and resumability.