Streaming
Text generated by a language model arrives in fragments, and a fact a pack must see can be split across several of them. A streaming session evaluates the pack as text arrives and holds back a short look-ahead window so it can catch an entity that spans chunks, while still releasing text that has passed quickly.
About the example
This page uses the built-in data-protection pack so you can run everything without writing a pack first. The same calls work for any decision pack: eligibility, limits, routing and so on. See Core concepts.
The chunk boundary problem#
Checking each fragment on its own misses entities that straddle fragments. With the example pack, an SSN split across four fragments looks like this:
"Identified SSN is " → clean "123" → clean "-45" → clean "-6789" → clean ← the SSN leaked, one piece at a time
A session keeps the most recent text in a raw buffer and re-evaluates the buffer on every push. It only releases text that sits before the look-ahead window, ends at a safe delimiter (space, newline, tab, carriage return, or one of . , ; ! ? ( ) [ ] { }), and cuts no value in two: redacting the released part and the held part separately must give the same text as redacting the whole buffer. The released part is remedied; the held part stays raw until it can be released.
Three ways to stream#
| API | Use when | On a fatal rule |
|---|---|---|
create_streaming_session() | You need per-chunk control: flags, indexes, custom handling. | Returns one chunk with blocked=True; later pushes return nothing. |
stream_filter(iterable) | You want a drop-in wrapper around a synchronous generator. | Raises StreamBlockedError. |
astream_filter(async_iterable) | Same, for async generators. | Raises StreamBlockedError. |
Sessions#
from helixor_runtime import HelixorEngine
engine = HelixorEngine()
session = engine.create_streaming_session(lookahead_chars=28, block_on_fatal=True)
for fragment in model_stream(): # your source of text fragments
for chunk in session.push(fragment):
if chunk.blocked:
abort_response(chunk.action, chunk.reason)
break
send_to_client(chunk.text)
if session.is_blocked:
break
else:
for chunk in session.flush(): # release what is left in the buffer
send_to_client(chunk.text)
Each emitted StreamChunkResult has chunk_index, text, redacted, blocked, is_final, action, reason and triggers. Always call flush() at the end of a stream that was not blocked, or the last window of text is never released.
Choosing lookahead_chars#
The window must be longer than the longest entity you need to catch across a boundary. For the example pack, the default of 28 covers formatted card numbers (19 characters) and most email addresses; values below 16 are raised to 16. A larger window catches longer entities but delays output by that many characters.
block_on_fatal#
With True (the default), a fatal match stops the stream: nothing more is released, including the text before the match that is still in the buffer. With False, fatal matches are remedied like warnings (in the example pack, redacted) and the stream continues. Keep the default unless a human reviews the output before it is used.
Generator wrappers#
from helixor_runtime import StreamBlockedError
try:
for text in engine.stream_filter(model_stream(), lookahead_chars=28):
send_to_client(text)
except StreamBlockedError as blocked:
abort_response(blocked.action, blocked.reason)
try:
async for text in engine.astream_filter(model_stream_async()):
await send_to_client(text)
except StreamBlockedError as blocked:
await abort_response(blocked.action, blocked.reason)
StreamBlockedError is exported from helixor_runtime, derives from HelixorRuntimeError and RuntimeError, and carries action and reason.
What to expect from the output#
- Streamed text equals the block result. Concatenated chunks match
evaluate(full_text).remedy.clean_text, whatever the fragment size, down to one character per fragment. Still store the result of oneevaluate()call on the assembled reply as the copy of record; its receipt is the one to log. - Chunks follow the text, not the fragments. A chunk ends at a safe delimiter, and a redaction token such as
[REDACTED_EMAIL]is never split across two chunks. chunk.redacteddescribes that chunk. It isTrueexactly when the chunk's text differs from the raw text it replaces.- Cost per push grows with the window, not the stream. Every push re-checks the unreleased buffer, and finding a safe cut can take a few more evaluations of its parts.
Changed in 0.2.1
In 0.2.0 a partial match could be redacted while the value was still arriving, which produced text such as [REDACTED_EMAIL]m, and redacted described the buffer rather than the chunk. Both are fixed in 0.2.1; you no longer need to join fragments into whole words before you push them. stream_filter() and astream_filter() now raise StreamBlockedError, a subclass of the RuntimeError they raised before.
Over the network#
The decision service exposes the same session over WebSocket (stream_start, stream_chunk, stream_flush) and a simulated stream over SSE. See Serving decisions and the streaming tutorial.