A streamed answer is delivered as it is written. That is the whole appeal of streaming, and it is also the problem for moderation: once text has rendered, the user has read it, and no later verdict takes it back. Batch moderation checks a finished answer; streaming moderation has to make decisions about an answer that does not exist yet.

The basic remedy, holding text back until a sentence passes a check, is described in Output Filtering, in depth and, from the prompt side, in Prompt Design for Streaming UX. This article covers the engineering underneath that sketch: what each check should see, how often to score, how a classifier keeps pace with thousands of streams, a release watermark that decouples generation from release, the retraction protocol, stopping generation early, tool-call arguments, and the two numbers, lag and leak, that tell you whether the design is working.

Why naive rescoring does not scale

Start with the naive design: after every new token, score the whole answer so far. It is correct, because a check always sees full context, but its cost grows with the square of the answer length. A 600-token answer triggers 600 classifier calls that read a combined 180,300 tokens, about 300 times the answer itself. Multiply by concurrent streams and the moderation fleet is larger than the model fleet.

Every practical design reduces two things: how often you score (cadence) and how much each score reads (window). Scoring a fixed 256-token window every 32 tokens costs about 19 calls for that answer and reads under 5,000 tokens. The question is what you lose, and the answer is context and reaction time. A window that does not include the start of the answer may miss harm that only reads as harmful in light of what came before; a stride of 32 tokens means a problem is seen up to 32 tokens after it was written.

Architecture: three loops and a watermark

The design that resolves this treats generation, scoring and release as three independent loops joined by a buffer and a number.

Streaming moderation: generation, scoring and release run at different speedsModeltokens at decode speedHold bufferall tokens, numberedWindow scorerasync, batchedWatermarkhighest safe tokenwindowsverdictsRelaySSE to clientrelease up to marktokensAbortcancel upstreamblockstop generatingClientrenders, retractsdelta / hold / retractRelease lag = tokens generated minus tokens released. It is the price of safety, paid in latency.Leak = tokens released before a later window blocked. It is the price of speed, paid in harm.
The watermark is the index of the highest token that every covering window has passed. The relay never sends a token above it. Generation does not wait for scoring unless the buffer is full.
  • Generation appends numbered tokens to a hold buffer as fast as the model produces them.
  • Scoring takes windows from the buffer on its own schedule and runs them, batched across streams, through one or more classifiers.
  • Release advances a watermark: the highest token index that has been covered by a passing window and is not inside any pending window. The relay sends everything up to the watermark.

Because the loops are decoupled, a slow classifier increases release lag but never blocks the model, and a fast classifier lets the watermark track generation closely. The buffer has a ceiling; if lag exceeds it, generation pauses, which is backpressure, not failure.

What each check should see

Each window should carry three things: the user turn it answers, a short prefix of the answer for context, and the tokens being judged. The user turn matters because many categories are contextual; dosage information is benign as an answer to a pharmacist and not to a stated intent to self-harm. The prefix matters because a harmful instruction can be set up in one paragraph and delivered in the next.

Windows must overlap. If a harmful phrase can be as long as L tokens, consecutive windows need at least L tokens of overlap, or the phrase can straddle a boundary and score low in both halves. In practice, overlap a full stride: each token is judged twice, once near the end of a window and once near the start.

Keep a final full-text check. When generation ends, score the whole answer once with the best classifier you have, before releasing the tail. Windowed checks catch local harm early; the final check catches harm that only appears as a whole, such as a set of individually harmless steps that together form a procedure you do not allow. Moderation Bypass, in depth lists the evasions, such as splitting and encoding, that defeat windows tuned only for local phrases.

The engine in code

The sketch below runs the three loops for one stream with asyncio. The scorer is any async function that returns a verdict; in production it is a client for a batched classifier service. The final full-text check is folded into the last window for brevity.

import asyncio
from dataclasses import dataclass, field

@dataclass
class Stream:
    user_turn: str
    tokens: list[str] = field(default_factory=list)
    watermark: int = 0            # tokens [0, watermark) may be released
    sent: int = 0                 # tokens actually sent to the client
    done: bool = False
    blocked: bool = False

STRIDE, WINDOW, PREFIX, MAX_LAG = 32, 256, 64, 160

async def score_loop(s: Stream, scorer, emit, cancel_generation):
    """Advance the watermark one stride at a time; block on the first failing window."""
    judged = 0                                     # tokens covered by passing windows
    while not s.blocked:
        ready = len(s.tokens) - judged
        if ready < STRIDE and not s.done:
            await asyncio.sleep(0.01)
            continue
        end = len(s.tokens) if s.done else judged + STRIDE
        start = 0 if s.done else max(0, end - WINDOW)   # final pass reads the whole answer
        window = "".join(s.tokens[start:end])
        prefix = "".join(s.tokens[:min(PREFIX, start)])
        verdict = await scorer(s.user_turn, prefix, window)   # batched across streams by the scorer service
        if verdict.block:
            s.blocked = True
            cancel_generation()                    # early abort: stop paying for tokens nobody will see
            await emit({"type": "final", "verdict": "blocked", "category": verdict.category})
            return
        judged = end
        s.watermark = end - STRIDE if not s.done else end    # keep one stride unreleased for overlap
        if s.done and judged == len(s.tokens):
            return

async def release_loop(s: Stream, emit):
    while not s.blocked:
        if s.watermark > s.sent:
            await emit({"type": "delta", "seq": s.sent, "text": "".join(s.tokens[s.sent:s.watermark])})
            s.sent = s.watermark
        elif s.done and s.sent == len(s.tokens):
            await emit({"type": "final", "verdict": "pass"})
            return
        await asyncio.sleep(0.02)

async def generate(s: Stream, model_stream):
    async for tok in model_stream:
        while len(s.tokens) - s.sent > MAX_LAG and not s.blocked:   # backpressure
            await asyncio.sleep(0.01)
        if s.blocked:
            return
        s.tokens.append(tok)
    s.done = True

Two lines carry the design. The watermark trails the judged position by one stride, so every released token has been seen by two overlapping windows. And cancel_generation stops the upstream request when a window blocks, which matters for cost: a blocked 2,000-token answer that would otherwise run to completion is cut off at the point it went wrong.

Cadence can adapt. Score every stride while the risk score is low, and drop to a smaller stride, or hold release entirely until the end, once any window scores above a watch threshold. Most answers never approach the threshold and stream with minimal lag; the few that do pay more latency, which is the right place to spend it.

Keeping pace with decode

Scoring load is streams times decode speed divided by stride. With 10,000 concurrent streams at 50 tokens per second and a 32-token stride, the scorer handles about 15,600 windows per second, each up to 320 tokens with its prefix and the user turn on top. These are illustrative numbers; plug in your own.

That load is only affordable batched: a service that collects windows from many streams for a few milliseconds and scores them together on a GPU. A small encoder classifier handles this cheaply; a large guard model or LLM judge belongs only on flagged windows and the final check. Size the service so that scoring latency, p95, is well under the time it takes to generate one stride: at 50 tokens per second a stride of 32 takes 640 ms, so a scorer that answers in 60 ms adds little lag.

Measure release lag directly, in tokens and milliseconds, as generated minus released. Its floor is about two strides, one to fill a window and one held back for overlap, plus scoring latency. When it climbs, the scorer is saturated, and the right response is to add scorer capacity or widen the stride, never to release unchecked text.

Retraction as a protocol event

Even with a watermark, sometimes text must be withdrawn: a final full-text check fails after most of the answer was released, or a policy update arrives mid-stream. Make retraction a first-class event in the wire protocol, not an error.

event: delta
data: {"seq": 0, "text": "Here is how to configure the proxy..."}

event: delta
data: {"seq": 96, "text": " Next, set the timeout..."}

: optional hold event, shown as a subtle indicator
event: hold
data: {"reason": "checking"}

event: retract
data: {"upto": 160, "notice": "Part of this answer was removed by a safety check."}

event: final
data: {"verdict": "blocked", "category": "self_harm", "request_id": "r-91c2"}

The client keeps deltas keyed by sequence number. A retract event replaces everything up to the given position with a visible notice; silently deleting text confuses users and hides the event from them. The final event carries the verdict and a request id that support can look up. Log both, because a retraction means a user saw something your policy forbids, and that is an incident to review, not a routine metric.

Provider-side filters make the same trade. Azure OpenAI, for example, offers an Asynchronous Filter mode in which content streams without buffering and filter annotations arrive later in the stream; Microsoft's documentation says plainly that harmful sections may be displayed before the signal arrives, and asks applications to consume the annotations and redact. If you use such a mode, you own the retraction path.

Tool calls, structured output and code

Not every stream is prose for a person. Three cases need different rules.

  • Tool-call arguments. Never execute a tool call from a partial stream. Buffer the arguments until the call is complete, parse them, and moderate the decoded values, not the raw JSON. An email body or a SQL statement assembled token by token is only safe to judge once whole.
  • Structured output. If the client renders JSON fields as they arrive, moderate each string value as a small stream of its own; JSON punctuation and keys only add noise to a classifier.
  • Code. Text classifiers trained on conversation score code poorly in both directions. Route code blocks to a code-specific check, or hold them until the block closes and check once.

Hidden reasoning that is never shown needs no release gate, but scoring it as a signal can catch intent that should block the visible answer.

Worked example: a tutor answer that turns

A chemistry tutor streams a three-paragraph answer at 50 tokens per second with a 32-token stride. Paragraphs one and two explain reaction safety and release about two strides behind generation, roughly 1.3 seconds, which users barely notice once text is flowing.

At token 210 the third paragraph starts listing quantities for a hazardous preparation. The window ending at token 224 scores above the watch threshold but below block, so cadence drops to an 8-token stride and the watermark stops advancing. The window ending at token 232 blocks. The user has seen 192 tokens of safe text; generation is cancelled at token 240 instead of running to an estimated 700; the client shows paragraphs one and two plus a notice; the event goes to the review queue described in LLM moderation architecture.

With no buffer, the quantities would have rendered before any check ran.

Failure modes

SymptomCauseFix
Lag climbs at peakScorer saturatedBatch harder, add capacity, widen stride; never bypass
Harm split across windows passesOverlap shorter than the phraseOverlap a full stride; keep the final full-text check
Retractions after releaseFinal check stricter than window checksAlign thresholds; hold the tail until the final check
Tool ran with an unsafe argumentArguments executed while streamingExecute only complete, parsed, moderated calls
Tokens billed after a blockUpstream request not cancelledPropagate cancellation to the model call
False blocks on code answersProse classifier applied to codeRoute code blocks to a code-aware check

Testing with replayed streams

Record harmful and harmless answers as token sequences and replay them through the engine at real speed. Measure, per configuration: leak tokens (released tokens later inside a blocking window or retracted), release lag p50 and p95, stall rate (share of streams that hit backpressure), block recall on the harmful set and false-block rate on the harmless set. Then change stride, window and thresholds and compare. Leak should be zero on every case in the harmful set where the harm sits after the first stride; anything else is a bug, not a trade-off.

Trade-offs

Every knob trades latency against safety. A smaller stride reacts faster but costs more scorer capacity. A larger window sees more context but scores more tokens per call. Holding whole answers is the safest option and removes streaming's benefit. The risk-adaptive approach puts the cost where the risk is, but it depends on a well-calibrated watch threshold; calibrate it on your own traffic, not a vendor's.

Choose per category: where any exposure is the harm, hold to the final check; where brief exposure is recoverable, such as mild profanity, stream with a stride and retract if needed.

What to do next

  1. Measure today's leak: replay a harmful set through your current stream and count tokens shown before a block.
  2. Split generation, scoring and release into separate loops joined by a numbered buffer and a watermark.
  3. Pick a window and stride, overlap by a full stride, and keep a final full-text check.
  4. Run the classifier as a batched service and size it so p95 latency is a fraction of one stride's generation time.
  5. Add retract and final events to the wire protocol and render retractions as visible notices.
  6. Cancel generation upstream on block, and never execute partially streamed tool calls.
  7. Track release lag, stall rate and leak tokens on a dashboard, and re-run the replay suite on every model or threshold change.
Key takeaway: Streaming moderation is a scheduling problem. Run generation, scoring and release as separate loops, release only up to a watermark that every overlapping window has passed, and keep a final full-text check. Batch the classifier so it keeps pace with decode, adapt cadence to risk, cancel generation on a block, never execute partially streamed tool calls, and treat retraction as a visible protocol event. Measure release lag and leak tokens with replayed streams, and tune stride and thresholds against both.