A guardrail is a single check. A content moderation pipeline is the system around those checks: a written policy with versions, the points in a turn where content is inspected, a sequence of increasingly expensive detectors, the action taken for each outcome, a queue where people review hard cases and appeals, and the measurements that tell you whether any of it is right. Without that system, guardrails drift: thresholds get tweaked in code, nobody knows the false positive rate, and a user who was wrongly blocked has nowhere to go.

This article builds such a pipeline for an ADK Java agent. How ADK orders plugins and callbacks is covered in ADK Java guardrails, and the model provider's own filters in Gemini safety settings; here they are building blocks. The focus is the policy, the tiers, the decision record, human review and measurement.

The policy comes first

Moderation starts with a policy that a reviewer, an engineer and a classifier can all apply the same way. Write it as data, not prose scattered through code: categories, a definition and examples for each, and an action per surface and severity.

CategoryInput actionOutput actionNotes
Harassment of a personblock above 0.9, judge 0.5-0.9block and regenerateContext matters: quoting abuse to report it is allowed
Self-harm intentdo not block; route to support flowreplace with support messageBlocking a person in crisis is the worst outcome
Sexual content involving minorsblock at any score, escalateblock, escalateLegal reporting duties may apply; involve counsel
Spam and scamsblock on rule matchn/aMostly deterministic patterns and rate signals
Credentials and secretsredactredactOverlaps with PII handling

Two properties make the policy usable. It is versioned: every decision records the policy version, so when the policy changes you can tell which decisions were made under which rules. And it separates detection from action: a classifier says how likely a category is, and the policy decides what that means for this surface. The same self-harm score leads to a block nowhere and to a support response everywhere.

enum Action { ALLOW, REDACT, SUPPORT, JUDGE, BLOCK, BLOCK_AND_ESCALATE }
enum Surface { USER_INPUT, TOOL_RESULT, MODEL_OUTPUT }

record Rule(String category, Surface surface, double judgeAbove, double blockAbove, Action onBlock) {}

record Policy(String version, List<Rule> rules) {
  Optional<Rule> rule(String category, Surface s) {
    return rules.stream().filter(r -> r.category().equals(category) && r.surface() == s).findFirst();
  }
}

record Verdict(Action action, String category, double score, int tier,
               String policyVersion, String reason, Duration took) {
  static Verdict allow(String v, int tier, Duration took) {
    return new Verdict(Action.ALLOW, "", 0, tier, v, "", took);
  }
}

Three checkpoints in a turn

Moderation checkpoints and tiers around one ADK Java turnUser messageC1 input checkonUserMessageCallbackAgent + modelbeforeModel short-circuitC3 output checkafterModelCallbackC2 tool resultsafterToolCallbacktool callTier 1: rulesregex, lists, ~1 msTier 2: classifierscores per categoryTier 3: LLM judgegray zone onlypass0.5-0.9Decision log + human review queue + appealspolicy version, scores, actionevery checkpointEach checkpoint runs the same tiered moderator with its own policy section and budget.
Three checkpoints inspect user input, tool results and model output; each runs the same tiered moderator and writes to one decision log.

Content enters and leaves an agent turn at three places, and each maps to an ADK Java plugin hook:

  • C1, user input, in onUserMessageCallback(InvocationContext, Content). The runner builds the user event from what this returns, so a replacement is what is stored and what the model sees on later turns.
  • C2, tool results, in afterToolCallback(...). Retrieved web pages, documents and other users' content arrive here. This is indirect input and the easiest to forget.
  • C3, model output, in afterModelCallback(CallbackContext, LlmResponse). Returning a response replaces the model's before it is stored and shown.

Input checks must be synchronous because the turn should not proceed on blocked content. Output checks are synchronous for non-streaming responses. Everything else, including the expensive judge on borderline allowed content, sampling for quality and review, can run asynchronously from the decision log. A useful rule: anything that decides whether the user sees something is on the hot path; anything that decides whether the policy is right is off it.

Tiered detectors and a latency budget

Detectors differ by orders of magnitude in cost and latency, so run them in sequence and stop as soon as one is confident. Typical figures, which you should measure on your own stack: rules and lists in about a millisecond, a dedicated classifier in tens of milliseconds, an LLM judge in hundreds of milliseconds to seconds.

interface Detector { Map<String, Double> score(String text, Surface surface); }

final class TieredModerator {
  private final RuleDetector rules;            // tier 1: deterministic, explainable
  private final Detector classifier;           // tier 2: your hosted or self-hosted model
  private final Detector judge;                // tier 3: an LLM prompted with the policy text
  private final Supplier<Policy> policy;       // hot-reloadable, versioned
  private final Duration budget;               // e.g. 150 ms for input

  Verdict check(String text, Surface surface) {
    Instant start = Instant.now();
    Policy pol = policy.get();
    Optional<Verdict> hit = rules.match(text, surface, pol);           // tier 1
    if (hit.isPresent()) return hit.get();

    Map<String, Double> scores;
    try {
      scores = CompletableFuture.supplyAsync(() -> classifier.score(text, surface))
          .get(budget.toMillis(), TimeUnit.MILLISECONDS);              // tier 2
    } catch (Exception e) {
      return failMode(surface, pol, start);                            // timeout or error
    }
    Verdict worst = decide(scores, surface, pol, 2, start);
    if (worst.action() != Action.JUDGE) return worst;

    Duration left = budget.minus(Duration.between(start, Instant.now()));
    if (left.isNegative()) return failMode(surface, pol, start);
    Map<String, Double> j = judge.score(text, surface);                // tier 3, gray zone only
    return decide(j, surface, pol, 3, start);
  }
}

decide looks up each category's rule for the surface and returns the most severe action: BLOCK above blockAbove, JUDGE between the two thresholds, ALLOW below. The judge's own scores go through the same function with the judge threshold disabled, so it can only allow or block.

The failMode decision is a policy choice, not an implementation detail. For input on a general assistant, failing open on a classifier timeout (allow and queue for async review) keeps the product usable; for categories with legal exposure, or for output on a children's product, fail closed. Record the choice per surface in the policy and count every fail-mode decision, because a classifier outage otherwise looks like a quiet day.

Wiring it into ADK Java

One plugin applies the moderator at all three checkpoints. Because plugins run before agent callbacks and cover every agent, nobody can forget to attach it to a new sub-agent. A blocked input is replaced with a placeholder that records the decision, and a flag in state makes beforeModelCallback answer with a fixed refusal instead of calling the model.

final class ModerationPlugin extends BasePlugin {
  private final TieredModerator mod;
  private final DecisionLog log;

  ModerationPlugin(TieredModerator mod, DecisionLog log) {
    super("moderation"); this.mod = mod; this.log = log;
  }

  @Override
  public Maybe<Content> onUserMessageCallback(InvocationContext ic, Content msg) {
    Verdict v = mod.check(Texts.of(msg), Surface.USER_INPUT);
    log.record(ic.invocationId(), Surface.USER_INPUT, v);
    return switch (v.action()) {
      case BLOCK, BLOCK_AND_ESCALATE -> {
        ic.session().state().put("mod:blocked", v.category());
        yield Maybe.just(Texts.content("user",
            "[message withheld: " + v.category() + ", policy " + v.policyVersion() + "]"));
      }
      case REDACT -> Maybe.just(Texts.redact(msg, v));
      default -> Maybe.empty();
    };
  }

  @Override
  public Maybe<LlmResponse> beforeModelCallback(CallbackContext ctx, LlmRequest.Builder req) {
    Object cat = ctx.state().get("mod:blocked");
    if (cat == null) return Maybe.empty();
    ctx.state().remove("mod:blocked");
    return Maybe.just(Texts.response(Refusals.forCategory(cat.toString())));
  }

  @Override
  public Maybe<LlmResponse> afterModelCallback(CallbackContext ctx, LlmResponse r) {
    if (r.partial().orElse(false)) return Maybe.empty();     // see streaming below
    Verdict v = mod.check(Texts.of(r), Surface.MODEL_OUTPUT);
    log.record(ctx.invocationId(), Surface.MODEL_OUTPUT, v);
    return v.action() == Action.ALLOW ? Maybe.empty()
        : Maybe.just(Texts.response(Refusals.forCategory(v.category())));
  }
}

Treat this as a shape to adapt rather than a drop-in: check how your ADK version persists state written from onUserMessageCallback (writing through ic.session().state() is a direct mutation, not a recorded delta), or pass the flag another way, for example a request-scoped map keyed by invocation ID. The afterToolCallback for C2 follows the same pattern, replacing a blocked result map with {"status": "withheld", "reason": category} so the model knows the content existed but cannot use it. The general callback model is in ADK Java callbacks.

Streaming output

Streaming breaks the clean picture: if tokens reach the user as they are generated, an output check on the final response is too late. There are three workable designs. Buffer the whole answer before display, which is simplest and costs time to first token. Hold back a window, releasing text a sentence or a few hundred characters behind generation and checking each window, which keeps most of the latency benefit. Retract, streaming freely and replacing the displayed message if the final check fails, which only works if your client supports it and your policy accepts that a user may have seen the text briefly.

Pick per surface. A held-back window is the usual compromise for general assistants; anything with legal or child-safety exposure should buffer. Whatever you choose, the final, complete response must also be checked, because a violation can span two windows.

Decision log, review queue and appeals

The decision log is the pipeline's memory. Each row holds the invocation ID, surface, category, scores per tier, action, policy version, latency and fail-mode flag, plus a content reference to a short-lived, access-controlled store rather than the text itself.

Three streams feed a human review queue from that log: every BLOCK_AND_ESCALATE; a random sample of blocks, to measure precision; and a random sample of allows, to measure misses. User appeals form a fourth stream with a response-time target. Reviewers label each item against the policy version that was in force, and their labels become the test set for the next policy change. Reviewing harmful content is draining work, so cap per-reviewer exposure, blur images by default and rotate assignments.

A policy change runs in shadow mode first: the new version scores live traffic alongside the current one, writes its verdicts to the log without acting, and you compare the two on the labelled set before switching.

Worked example: measuring a day of moderation

A community help agent receives 50,000 user messages a day. Tier 1 rules block 60, mostly known scam links. The classifier scores the rest: 250 score above 0.9 in some category and are blocked, and 1,250 fall in the 0.5 to 0.9 gray zone and go to the judge, which blocks 300 of them. Total blocks: 610, or 1.2% of messages.

Reviewers label a random 200 of the 610 blocks and agree with 168, so precision is about 0.84. They also label a random 2,000 of the 49,390 allowed messages and find 6 violations, a miss rate of 0.3%, which projects to about 148 missed per day. True violations caught are about 610 times 0.84, or 512, so estimated recall is 512 divided by 660, roughly 0.78.

Breaking the misses down by category shows most are harassment phrased politely. Lowering the harassment judge threshold to 0.4 in shadow mode catches 4 of the 6 missed examples and adds 90 blocks a day, of which reviewers judge 60 correct. That is a precision of about 0.67 on the added blocks: acceptable for harassment under this policy, so the change ships with a new version number and the judge's extra cost of 90 additional calls a day is noted. These numbers come from small samples; put confidence intervals on them before treating a one-point change as real.

Failure modes

  • Moderating only the user. Tool results carry other people's content and injected instructions. C2 is not optional for agents that browse or retrieve.
  • Classifier outage looks like safety. Fail-open timeouts silently allow everything. Alert on the rate of fail-mode verdicts.
  • Context-free judging. A message quoting abuse to report it scores like abuse. Give the judge the previous turn and the policy's examples of allowed quotation.
  • Language gaps. Classifiers trained mostly on English under-detect other languages. Measure recall per language you serve, not overall.
  • Threshold edits in code. Changes made outside the policy file are not versioned and cannot be audited. Load thresholds only from the policy.
  • Logging the content you blocked. The decision log becomes the most sensitive table you own. Store references with short retention and tight access instead. See PII redaction for the same problem in session storage.

Trade-offs

Every threshold trades false blocks against misses, and the right point differs by category and surface: blocking a self-harm disclosure is a failure, while missing child sexual abuse material is unacceptable. Express that as per-category actions rather than one global sensitivity. Tiering trades latency and cost against accuracy; the judge is your most accurate and least predictable detector, so keep it on the gray zone and measure its agreement with reviewers the way continuous evaluation measures any judge. Buffering output trades responsiveness for certainty. Provider safety filters are free and fast but use their own taxonomy and cannot see your policy, so treat them as one more input, not the pipeline.

What to do next

  1. Write the policy as versioned data: categories, definitions, examples and an action per surface and severity.
  2. Put your classifier behind a Detector interface and measure its latency on real traffic.
  3. Implement the tiered moderator with a budget and an explicit fail mode per surface.
  4. Attach one plugin covering user input, tool results and model output; add a test per checkpoint.
  5. Choose buffering, a held-back window or retraction for each streaming surface.
  6. Create the decision log without raw content, and a review queue fed by escalations, samples and appeals.
  7. Estimate precision and recall per category from samples, and run every policy change in shadow mode first.
Key takeaway: A moderation pipeline is a versioned policy, three checkpoints, tiered detectors and a measurement loop, not a single filter. In ADK Java, one plugin can cover user input, tool results and model output; a tiered moderator keeps cost and latency bounded with an explicit fail mode; and a decision log without raw content feeds human review, appeals and per-category precision and recall. Change the policy only through new versions tested in shadow mode.