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.

Advertisement

The contract table

ADK Java contracts: reactive at every seam, sync only where the framework adapts itYour callerservlet, CLI, SSE, testRunner.runAsyncFlowable of EventCallbacksMaybe, or Optional (Sync)BaseLlm.generateContentFlowable of LlmResponseBaseTool.runAsyncSingle of MapBaseLlm.connectlive: Completable sends,Flowable receive()subscribeone or manyempty = go ontool callBIDIEmpty means continue; one value short-circuits; onError means the call failed; dispose means stop.Blocking code is allowed inside these types only if it is deferred to subscription time.
Each seam of the runtime and the reactive type it uses. Sync code enters through adapters, never by replacing the type.
SeamSignature returnsMeaning of the stream
BaseLlm.generateContent(LlmRequest, boolean stream)Flowable<LlmResponse>one response when not streaming; possibly many when streaming
BaseLlm.connect(LlmRequest)BaseLlmConnectionlive session: sends return Completable, receive() returns a Flowable
BaseTool.runAsync(Map, ToolContext)Single<Map<String, Object>>exactly one result or an error
BaseTool.processLlmRequest(...)Completabledone or failed, no value
Model, agent and tool callbacksMaybe<...>empty continues; a value replaces the step
The *Sync callback twinsOptional<...>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.

Advertisement

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. just around a blocking call: runs too early, cannot be cancelled, and never retries. Use fromCallable, defer or create.
  • Blocking on the wrong thread. Blocking inside a reactive chain without subscribeOn stalls 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 stream is 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 carrying errorCode() or errorMessage(). Mixing them breaks retry and error callbacks.
  • Swallowed disposal. A streaming adapter without setCancellable keeps 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

  1. Grep your adapters and tools for Flowable.just, Single.just and Maybe.just around I/O, and replace them with deferred forms.
  2. Check your model adapter emits exactly one response when stream is false, and sets partial and turnComplete explicitly when true.
  3. Add setCancellable to every streaming adapter and test that disposing the runner's subscription closes the HTTP connection.
  4. Move pure-logic callbacks to the Sync twins; give I/O callbacks a timeout and a fallback.
  5. Audit every blockingGet and blockingSubscribe and confirm none runs on an event-loop thread.
  6. Decide your retry policy per seam, and forbid retries after the first streamed chunk.
  7. Read streaming agent architecture next for backpressure and resumability.
Key takeaway: ADK Java's contracts are reactive throughout: models return a Flowable of responses, tools a Single, callbacks a Maybe, the runner a Flowable of events, and live connections Completable sends plus a Flowable receive. Synchronous code is welcome, through Sync callback twins, plain FunctionTool returns and deferred blocking inside adapters, but it must never replace the type. Defer work to subscription, block only on threads meant for it, wire disposal to your sockets, and pick the blocking or streaming style at the edge of your system.