Every production agent service needs behaviour the framework does not ship: audit records, rate limits, tool policies, model fallbacks, trace IDs across threads. In the Agent Development Kit for Java you can add all of these without forking the runtime, because it exposes a small set of extension points. The hard part is choosing the right one, since the same guardrail can be written as a plugin, an agent callback, a custom agent or a model wrapper, and each choice has a different scope, ordering and failure behaviour.
This article maps those extension points, explains the plugin contract (the runtime's middleware layer), and builds a policy plugin, a fallback agent, a model decorator and a context-propagating executor. API names were checked against the adk-java source; the library still changes, so confirm signatures against your version. How one turn flows through the runtime is covered in the execution loop anatomy, and which thread does the work in the thread model article.
The extension points and their scope
| Extension point | Scope | Mechanism | Use it for |
|---|---|---|---|
| Plugin | Every agent, model and tool call in one runner | Implement Plugin or extend BasePlugin | Cross-cutting policy, audit, metrics, rate limits |
| Agent callbacks | One agent | Callbacks on the agent builder | Agent-specific guardrails and caching |
| Custom agent | A node in the agent tree | Subclass BaseAgent | New control flow: fallback, gating, voting |
| Model decorator | Every call through one model instance | Subclass BaseLlm and delegate | Deadlines, retries, routing, request shaping |
| Services | Persistence for the whole app | Implement the session, artifact or memory service | Durable sessions, object storage, recall |
| Executors | Threads that run tools and HTTP I/O | Pass an Executor to the agent or model builder | Bounded concurrency, context propagation |
A useful rule: put behaviour at the narrowest scope that covers every case it must cover. A redaction rule for all agents belongs in a plugin, because a callback is easy to forget on the next agent. A cache for one agent belongs in its callbacks, where the callback lattice article covers the patterns. Control flow belongs in a custom agent, and transport concerns belong in the model decorator or the executor.
The plugin contract
A plugin implements the com.google.adk.plugins.Plugin interface, whose only abstract method is getName(); every hook is a default method that does nothing, so you override only what you need. BasePlugin is a convenience base class that takes the name in its constructor. The hooks fall into three groups:
- Run level:
onUserMessageCallbackandbeforeRunCallback(both returnMaybe<Content>),onEventCallback(Maybe<Event>), andafterRunCallback,onRunErrorCallbackandclose(Completable). - Agent and model:
beforeAgentCallbackandafterAgentCallback(Maybe<Content>);beforeModelCallback, which receives a mutableLlmRequest.Builder;afterModelCallback; andonModelErrorCallback(Maybe<LlmResponse>). - Tool:
beforeToolCallback,afterToolCallbackandonToolErrorCallback, which see the tool, its arguments and theToolContextand returnMaybe<Map<String, Object>>.
The return type carries the semantics. An empty Maybe means carry on. A value means "use this instead": a before-tool value becomes the tool result and the tool never runs; a before-model value becomes the model response and no request is sent; an on-error value replaces the error with a response; a beforeRunCallback value ends the run without calling the agent. The PluginManager calls plugins in registration order and stops at the first non-empty result, so later plugins never see that hook for that call. Plugins are consulted before the agent's own callbacks, and in the code checked a plugin result skips those callbacks entirely. Registering two plugins with the same name throws IllegalArgumentException.
Two hooks behave differently. The value from onEventCallback replaces the event the caller receives; in the code checked that happens after the event is appended to the session, so treat it as a way to shape output, not to change stored history. And onRunErrorCallback is for observation only: the runner swallows any error it raises.
A policy plugin in code
The plugin below does two jobs that must apply to every agent: it blocks named tools, returning an explanatory result the model can read, and it turns a model failure into a polite reply instead of a stack trace. It does no work at assembly time and catches its own failures. Do not rely on how the plugin manager treats a throwing plugin; decide explicitly whether a failure should block (fail closed) or be ignored (fail open).
import com.google.adk.agents.CallbackContext;
import com.google.adk.models.LlmRequest;
import com.google.adk.models.LlmResponse;
import com.google.adk.plugins.BasePlugin;
import com.google.adk.tools.BaseTool;
import com.google.adk.tools.ToolContext;
import com.google.genai.types.Content;
import com.google.genai.types.Part;
import io.reactivex.rxjava3.core.Maybe;
import java.util.Map;
import java.util.Set;
/** Runner-wide policy: block risky tools, degrade gracefully when the model fails. */
public final class PolicyPlugin extends BasePlugin {
private final Set<String> blockedTools;
private final AuditSink audit; // your own interface; must never throw
public PolicyPlugin(Set<String> blockedTools, AuditSink audit) {
super("policy"); // names must be unique per runner
this.blockedTools = blockedTools;
this.audit = audit;
}
@Override
public Maybe<Map<String, Object>> beforeToolCallback(
BaseTool tool, Map<String, Object> args, ToolContext ctx) {
return Maybe.defer(() -> {
if (!blockedTools.contains(tool.name())) {
return Maybe.empty(); // empty: let the tool run
}
audit.record(ctx.invocationId(), ctx.agentName(), "tool_blocked", tool.name());
// Non-empty: becomes the tool's result; the tool body never runs.
return Maybe.just(Map.of("error", "Tool " + tool.name() + " is disabled by policy."));
});
}
@Override
public Maybe<LlmResponse> onModelErrorCallback(
CallbackContext ctx, LlmRequest.Builder request, Throwable error) {
return Maybe.fromCallable(() -> {
try {
audit.record(ctx.invocationId(), ctx.agentName(), "model_error", error.toString());
} catch (RuntimeException auditFailure) {
// A failing audit must not cost the user the fallback reply.
}
return LlmResponse.builder()
.content(Content.fromParts(Part.fromText(
"The assistant is temporarily unavailable. Please try again shortly.")))
.build();
});
}
}import com.google.adk.runner.Runner;
Runner runner = Runner.builder()
.appName("support")
.agent(rootAgent)
.sessionService(sessionService) // defaults to InMemorySessionService
.artifactService(artifactService) // defaults to InMemoryArtifactService
.plugins(
new MetricsPlugin(meterRegistry), // observes only: always returns empty
new RateLimitPlugin(limiter), // may short-circuit beforeRun
new PolicyPlugin(Set.of("delete_account", "refund"), audit))
.build();Registration order is part of the design. Metrics comes first and always returns empty, so it sees every call. Rate limiting comes next, so a throttled run is rejected before any policy work. Policy comes last among the before-hooks because a value it returns ends the chain. Put policy first and metrics would never count blocked calls.
Custom agents: new control flow
When you need behaviour the built-in sequential, parallel and loop agents do not provide, subclass BaseAgent. The constructor takes a name, a description, optional sub-agents and optional before and after agent callbacks. You implement two methods, runAsyncImpl and runLiveImpl, each taking an InvocationContext and returning a Flowable<Event>. The public runAsync wraps your implementation with the plugin and agent callbacks, so you get that behaviour for free.
import com.google.adk.agents.BaseAgent;
import com.google.adk.agents.InvocationContext;
import com.google.adk.events.Event;
import io.reactivex.rxjava3.core.Flowable;
import java.util.List;
/** Runs the primary sub-agent; if it fails before emitting anything, runs the fallback. */
public final class FallbackAgent extends BaseAgent {
private final BaseAgent primary, fallback;
public FallbackAgent(String name, BaseAgent primary, BaseAgent fallback) {
super(name, "Primary agent with a fallback on early failure.",
List.of(primary, fallback), null, null);
this.primary = primary;
this.fallback = fallback;
}
@Override
protected Flowable<Event> runAsyncImpl(InvocationContext ctx) {
return Flowable.defer(() -> {
boolean[] emitted = {false};
return primary.runAsync(ctx)
.doOnNext(e -> emitted[0] = true)
.onErrorResumeNext(err -> emitted[0]
? Flowable.error(err) // partial output already recorded: do not mix
: fallback.runAsync(ctx));
});
}
@Override
protected Flowable<Event> runLiveImpl(InvocationContext ctx) {
return Flowable.error(new UnsupportedOperationException("live mode not supported"));
}
}The fallback rule is deliberately narrow. Events from the primary are appended to the session as they are emitted, so if it fails halfway through, starting the fallback would leave two agents' partial work interleaved in the history. The agent therefore falls back only when nothing was emitted, and otherwise surfaces the error. Keep custom agents cold until subscription, as Flowable.defer does above. Declare unsupported live mode with an error rather than an empty stream, so a misconfiguration fails loudly.
Decorating the model
Transport concerns such as deadlines, retries, provider routing and request shaping sit most naturally around the model. BaseLlm has a constructor taking the model name and two abstract methods, generateContent(LlmRequest, boolean stream) returning Flowable<LlmResponse>, and connect(LlmRequest) for live connections. A decorator holds a delegate and forwards to it:
import com.google.adk.models.BaseLlm;
import com.google.adk.models.BaseLlmConnection;
import com.google.adk.models.LlmRequest;
import com.google.adk.models.LlmResponse;
import io.reactivex.rxjava3.core.Flowable;
import java.util.concurrent.TimeUnit;
/** Decorator: deadline plus bounded retry for non-streaming calls; everything else delegates. */
public final class ResilientLlm extends BaseLlm {
private final BaseLlm delegate;
private final long timeoutSeconds;
private final int retries;
public ResilientLlm(BaseLlm delegate, long timeoutSeconds, int retries) {
super(delegate.model());
this.delegate = delegate;
this.timeoutSeconds = timeoutSeconds;
this.retries = retries;
}
@Override
public Flowable<LlmResponse> generateContent(LlmRequest request, boolean stream) {
Flowable<LlmResponse> call = Flowable.defer(() -> delegate.generateContent(request, stream))
.timeout(timeoutSeconds, TimeUnit.SECONDS);
// Retrying a stream would replay chunks the caller has already seen, so only retry unary calls.
return stream ? call : call.retry(retries, ResilientLlm::isTransient);
}
@Override
public BaseLlmConnection connect(LlmRequest request) {
return delegate.connect(request);
}
private static boolean isTransient(Throwable t) {
// Not TimeoutException: the timed-out request is still in flight, so a retry duplicates it.
return t instanceof java.io.IOException; // extend with your client's 429/503 types
}
}The deadline matters more than it looks. The thread-model article explains that the built-in Gemini client waits on its future with a blocking call and configures no HTTP timeouts, so a downstream timeout signals an error but may not free the waiting thread. A timeout does not free that thread or cancel the request, so retrying after one duplicates it; the decorator does not replace HTTP client timeouts. Retry only unary calls, only on errors you have classified as transient, and only a small number of times; a retried stream replays text the user has already seen.
Services and executors
Sessions, artifacts and memory are interfaces. Runner.builder() defaults to InMemorySessionService and InMemoryArtifactService, which lose everything on restart and are not shared between replicas. For production, implement the session service over a durable store, keeping event appends atomic and ordered per session; the Postgres session schema article shows one layout. Implement against the interface in the version you build with, since service method signatures have changed between releases.
Executors are covered in depth in the thread-model article: LlmAgent.builder().executor(...) sets the workers for the parallel-subscribe tool mode and the live send loop, and the Gemini builder accepts an HTTP dispatcher executor. The extension to add here is wrapping. Thread-local context such as the SLF4J MDC does not follow work onto a worker thread, so logs from tools lose their request ID. Wrap the executor once and pass the wrapper everywhere:
import java.util.Map;
import java.util.concurrent.Executor;
import org.slf4j.MDC;
/** Copies the submitting thread's MDC (request id, tenant) onto whichever thread runs the task. */
public final class MdcExecutor implements Executor {
private final Executor delegate;
public MdcExecutor(Executor delegate) { this.delegate = delegate; }
@Override public void execute(Runnable task) {
Map<String, String> captured = MDC.getCopyOfContextMap();
delegate.execute(() -> {
Map<String, String> previous = MDC.getCopyOfContextMap();
if (captured != null) MDC.setContextMap(captured); else MDC.clear();
try { task.run(); }
finally { if (previous != null) MDC.setContextMap(previous); else MDC.clear(); }
});
}
}OpenTelemetry offers the same idea for trace context through its context-wrapping executors; see ADK Java observability for spans and metrics. Restore the previous context in a finally block, as above, because pooled threads are reused and leaked context attributes one request's logs to another.
Worked example: composing a support deployment
A support assistant has a triage agent, an orders agent with read tools, and a refunds agent with a refund tool. Requirements: every tool call audited; refunds disabled for one tenant; at most 30 runs per minute per user; a 20-second model deadline with two retries; logs correlated by request ID. The mapping is:
- A metrics-and-audit plugin records
beforeToolCallbackandafterToolCallbackfor all three agents, returning empty every time. - A rate-limit plugin checks the user in
beforeRunCallbackand returns a short "slow down" content value when the budget is spent, so no model call is made. - The policy plugin reads the tenant from session state through the
ToolContextand blocksrefundfor that tenant. - Each agent's model is wrapped once in
ResilientLlm(gemini, 20, 2). - One
MdcExecutoraround a bounded pool is passed to every agent builder.
Worst case per run: a transient failure on every attempt costs three model calls of up to 20 seconds, 60 seconds in all, so cap the run with a limit on model calls in the run configuration and an overall deadline in the host.
Failure modes and operations
- Silent short-circuits. A plugin that returns a value by mistake (for example
Maybe.justof an empty map) skips the tool and every later plugin. Log every non-empty return. - Blocking in hooks. Hooks run on the subscribing thread, so a synchronous database call in a hook adds latency to every call.
- Ordering regressions. Test the plugin order: blocked calls are still counted, throttled runs never reach the model.
- Mutable request edits. Request edits in one plugin are visible to later plugins and the model; keep them additive.
- Resource leaks. Implement
close()for plugins that hold clients. - Version drift. Pin the ADK version and cover each extension with an integration test against a fake model.
What to do next
- List every cross-cutting behaviour your agents need and assign each one to a plugin, callback, custom agent, model decorator, service or executor, narrowest scope first.
- Move any guardrail that is copied across agent callbacks into one plugin, and write a test that proves it fires for a newly added agent.
- Order your plugins deliberately, observe-only first, and assert the order in a unit test.
- Wrap each model in a decorator with a deadline and a narrow transient-error retry, and configure HTTP client timeouts as well.
- Replace the in-memory session and artifact services before running more than one replica.
- Wrap every executor you pass to ADK so that logging and trace context reach tool threads, and check a tool's log line for the request ID.