Skip to content

Output filters see each streamed chunk alone, so a value split across two SSE events passes output_redactor #1036

Description

@ninadphalak

Summary

Output filters on a model listener receive each raw upstream chunk on its own. A value that the upstream splits across two SSE events is never seen whole by any filter call, so it reaches the client unfiltered. The repo's own output_redactor demo shows it: SECRET_TOKEN inside one event comes back as [REDACTED]; the same text split as SECRET_ + TOKEN. comes back intact.

Reproduced on katanemo/plano:0.4.36 (tag 0.4.36, 003c36a) with demos/filter_chains/model_listener_filter/output_filter.py unmodified. Three runs, same result each time:

marker in one SSE event : The token is [REDACTED].
marker split across two : The token is SECRET_TOKEN.

Reproduction

Run from demos/filter_chains/model_listener_filter/. Needs Docker and Python with fastapi and uvicorn; no API keys. It starts a fake OpenAI-compatible upstream on 8765, the demo's output_filter.py on 10502, and the release image with only output_filters: [output_redactor] configured.

repro.sh
#!/usr/bin/env bash
# Output filter vs a value split across two SSE events.
# Run from demos/filter_chains/model_listener_filter/ in a checkout of katanemo/plano.
# Needs: docker, python3 with fastapi + uvicorn. No API keys.
set -u
IMAGE="${PLANO_IMAGE:-katanemo/plano:0.4.36}"
WORK="$(mktemp -d)"
cp output_filter.py "$WORK/"

cat > "$WORK/split_upstream.py" <<'PY'
import asyncio, json, time
from fastapi import FastAPI, Request
from fastapi.responses import StreamingResponse

app = FastAPI()
HEAD, TAIL = "The token is SECRET_", "TOKEN."

def ev(content):
    return "data: " + json.dumps({"id": "x", "object": "chat.completion.chunk", "created": int(time.time()),
        "model": "gpt-4o-mini", "choices": [{"index": 0, "delta": {"content": content}, "finish_reason": None}]}) + "\n\n"

@app.post("/v1/chat/completions")
async def completions(req: Request):
    body = await req.json()
    split = body["messages"][-1]["content"] == "split"
    async def gen():
        if split:
            yield ev(HEAD); await asyncio.sleep(0.5)
            yield ev(TAIL); await asyncio.sleep(0.5)
        else:
            yield ev(HEAD + TAIL); await asyncio.sleep(0.5)
        yield "data: [DONE]\n\n"
    return StreamingResponse(gen(), media_type="text/event-stream")
PY

cat > "$WORK/plano_config.yaml" <<'YAML'
version: v0.3.0
filters:
  - id: output_redactor
    url: http://host.docker.internal:10502
    type: http
model_providers:
  - model: openai/gpt-4o-mini
    access_key: local-demo-key
    base_url: http://host.docker.internal:8765/v1
    default: true
listeners:
  - type: model
    name: llm_gateway
    port: 12000
    output_filters:
      - output_redactor
YAML

(cd "$WORK" && python3 -m uvicorn split_upstream:app --host 0.0.0.0 --port 8765 >/dev/null 2>&1 & echo $! > "$WORK/up.pid")
(cd "$WORK" && python3 -m uvicorn output_filter:app --host 0.0.0.0 --port 10502 >/dev/null 2>&1 & echo $! > "$WORK/filter.pid")
docker run -d --name plano-split-repro -p 12000:12000 \
  -v "$WORK/plano_config.yaml:/app/plano_config.yaml:ro" \
  --add-host host.docker.internal:host-gateway "$IMAGE" >/dev/null
for _ in $(seq 1 60); do curl -s -o /dev/null localhost:12000/v1/models && break; sleep 1; done
sleep 2

ask() {
  curl -sN http://localhost:12000/v1/chat/completions -H 'Content-Type: application/json' \
    -d "{\"model\":\"gpt-4o-mini\",\"messages\":[{\"role\":\"user\",\"content\":\"$1\"}],\"stream\":true}" |
  python3 -c '
import sys, json
text = ""
for line in sys.stdin:
    line = line.strip()
    if line.startswith("data:") and "[DONE]" not in line:
        text += json.loads(line[5:])["choices"][0]["delta"].get("content") or ""
print(text)'
}

echo "marker in one SSE event : $(ask single)"
echo "marker split across two : $(ask split)"

docker rm -f plano-split-repro >/dev/null
kill "$(cat "$WORK/up.pid")" "$(cat "$WORK/filter.pid")" 2>/dev/null

Cause

crates/brightstaff/src/streaming.rs, create_streaming_response_with_output_filter (line 681 onward at 003c36a): the loop calls process_raw_filter_chain(&chunk, ...) for each chunk from the upstream byte stream and forwards the result. Nothing is carried from one call to the next, and there is no end-of-stream call, so a filter has no way to see text that spans a boundary. The demo filter matches correctly; it is never given the whole value.

The same function passes the original chunk through when a filter errors (lines 729-745), which produces a second case: with stream: false the demo filter also does not redact. The response reaches the filter gzip-compressed by Plano's Envoy and in two raw chunks (208 + 10 bytes in my run). gzip.decompress on the first chunk raises EOFError, the filter returns 500, and the gateway forwards the original bytes. test.sh only asserts redaction for stream: true, so this does not show up there.

Possible directions

The same split-boundary pattern has shipped in several streaming guardrails; a short write-up with a single test that catches it: https://doi.org/10.5281/zenodo.22909585

Disclosure: I maintain a streaming-privacy gateway and a benchmark in this area. I ran this against your published image only.

Activity

  1. ninadphalak commented on Oct 7, 2026

    @ninadphalak
    Author

    Opened #1040 for this: an opt-in output_filter_mode: buffered on the model listener. The filter chain runs once on the whole response, so a split value is redacted, and a filter error withholds the body instead of forwarding the original bytes. The default stays streaming, unchanged.

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Metadata

Metadata

Assignees

No one assigned

    Labels

    No labels
    No labels

    Type

    No type

    Projects

    No projects

      Milestone

      No milestone

      Relationships

      None yet

      Development

      No branches or pull requests

      Issue actions