Decide on a stream as it arrives
You evaluate a model's reply against a pack while it streams: the pack's remedy is applied before text reaches the user, and the stream stops the moment a fatal rule fires, even when the value arrives split across several tokens. With the built-in example pack (data protection), contact details are redacted and an SSN halts the stream.
About the example
This tutorial 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.
What you'll build#
- A push-based session using
create_streaming_session(),push()andflush(). - A fatal halt with
block_on_fatal=True. - The drop-in generators
stream_filter()andastream_filter(), including how they signal a block. - A production wrapper,
stream_and_persist(), that stores an authoritative remedied copy of each reply.
Prerequisites#
- Make your first decision.
- Read Streaming for the release rule and parameters.
Steps#
- Fake a model stream
The replies below split an email address, a phone number and an SSN across tokens, which is what makes streaming hard: no single token contains the whole value the pack must see.
"""Stand-in for a language model's token stream. Replace with your client.""" import asyncio def reply_with_contact(): yield from [ "Thanks for waiting. ", "Your account manager is ", "Dana. ", "Email her at dana", ".reyes", "@example", ".com", " or call (415) 555-", "0199", " before five. ", "Have a good day!", ] def reply_with_ssn(): yield from [ "Intake summary: ", "applicant verified. ", "SSN on file is ", "123", "-45", "-6789", ". Next step: ", "schedule review.", ] async def areply_with_contact(): for token in reply_with_contact(): await asyncio.sleep(0.01) yield token - Push tokens through a session
A session buffers the tail of the stream as raw text, re-evaluates the buffer on every push, and releases text up to a safe delimiter once it is at least
lookahead_charsbehind the newest character.flush()releases the rest when the stream ends. The second half of the script feeds a stream that contains an SSN.from helixor_runtime import HelixorEngine from fake_model import reply_with_contact, reply_with_ssn engine = HelixorEngine() print("== redact in flight ==") session = engine.create_streaming_session(lookahead_chars=28, block_on_fatal=True) for token in reply_with_contact(): for chunk in session.push(token): print(f"emit #{chunk.chunk_index}: {chunk.text!r}") for chunk in session.flush(): print(f"final #{chunk.chunk_index}: {chunk.text!r} is_final={chunk.is_final}") print("\n== halt on fatal ==") session = engine.create_streaming_session(block_on_fatal=True) for token in reply_with_ssn(): for chunk in session.push(token): if chunk.blocked: print(f"HALTED: {chunk.action}") print(f" reason: {chunk.reason}") print(f" rules : {[t.rule_id for t in chunk.triggers]}") else: print(f"emit #{chunk.chunk_index}: {chunk.text!r}") if session.is_blocked: break print("pushes after the halt return:", session.push("more text"))== redact in flight == emit #1: 'Thanks for ' emit #2: 'waiting. ' emit #3: 'Your account ' emit #4: 'manager is ' emit #5: 'Dana. ' emit #6: 'Email ' emit #7: 'her at ' emit #8: '[REDACTED_EMAIL] or call ' emit #9: '[REDACTED_PHONE] ' final #10: 'before five. Have a good day!' is_final=True == halt on fatal == emit #1: 'Intake ' emit #2: 'summary: ' emit #3: 'applicant ' HALTED: block_glba_ssn_leakage reason: Statutory violation: Social Security Number (GLBA / FCRA) detected rules : ['RULE-GLBA-SSN-BLOCK'] pushes after the halt return: []
Read the output closely:
- Chunks end at word boundaries, not at tokens. The session cuts only where no value straddles the cut, so a redaction marker such as
[REDACTED_EMAIL]always arrives whole. Still concatenate chunks before you show or store them. - On a fatal value the session returns one chunk with
blocked=True, emptytext, the action, the reason and the fatal triggers only. Every laterpush()returns[]. The SSN itself was never released. chunk.redacteddescribes the chunk. It isTrueexactly on the chunks whose text was changed, here chunks 8 and 9.
- Chunks end at word boundaries, not at tokens. The session cuts only where no value straddles the cut, so a redaction marker such as
- Use the generator filters
When you already have an iterator of tokens, wrap it.
stream_filter()yields safe text;astream_filter()does the same for an async iterator. Unlikepush(), both raiseStreamBlockedErroron a fatal block, so wrap them intry. It is exported fromhelixor_runtime, carriesactionandreason, and is a subclass ofRuntimeError, so existingexcept RuntimeErrorhandlers still catch it.import asyncio from helixor_runtime import HelixorEngine, StreamBlockedError from fake_model import areply_with_contact, reply_with_contact, reply_with_ssn engine = HelixorEngine() print("== stream_filter ==") print("".join(engine.stream_filter(reply_with_contact()))) print("\n== stream_filter on a fatal stream ==") shown = [] try: for text in engine.stream_filter(reply_with_ssn()): shown.append(text) except StreamBlockedError as exc: print("shown before halt:", repr("".join(shown))) print("blocked by:", exc.action) async def main(): print("\n== astream_filter ==") parts = [] async for text in engine.astream_filter(areply_with_contact()): parts.append(text) print("".join(parts)) asyncio.run(main())== stream_filter == Thanks for waiting. Your account manager is Dana. Email her at [REDACTED_EMAIL] or call [REDACTED_PHONE] before five. Have a good day! == stream_filter on a fatal stream == shown before halt: 'Intake summary: applicant ' blocked by: block_glba_ssn_leakage == astream_filter == Thanks for waiting. Your account manager is Dana. Email her at [REDACTED_EMAIL] or call [REDACTED_PHONE] before five. Have a good day!
Text shown before the halt stays shown. Decide in your UI what to do with it, for example replace the partial reply with a notice.
- Store a copy of record
What the user saw is not what you should store. Keep the raw tokens of one reply in memory, and when the stream ends, decide on the complete reply with one
evaluate()call. Store that result'sclean_text, and nothing when it blocks. The script below feeds the worst case, one character per token.from helixor_runtime import HelixorEngine, StreamBlockedError engine = HelixorEngine() REPLY = "Reach support at support@example.com or security@example.com." def one_char_tokens(text): yield from text # worst case: every character is its own token def stream_and_persist(tokens): raw_parts, shown_parts = [], [] def tee(source): for token in source: raw_parts.append(token) yield token try: for text in engine.stream_filter(tee(tokens)): shown_parts.append(text) # send to the user here except StreamBlockedError: return "".join(shown_parts), None # nothing to store # Before persisting, decide on the complete reply in one call. final = engine.evaluate("".join(raw_parts)) if final.action.startswith("block"): return "".join(shown_parts), None return "".join(shown_parts), final.remedy.clean_text shown, stored = stream_and_persist(one_char_tokens(REPLY)) print("shown :", repr(shown)) print("stored :", repr(stored)) print("same :", shown == stored)shown : 'Reach support at [REDACTED_EMAIL] or [REDACTED_EMAIL].' stored : 'Reach support at [REDACTED_EMAIL] or [REDACTED_EMAIL].' same : True
The streamed text and the stored copy match, even with one-character tokens: the session releases a prefix only where redacting it separately gives the same text as redacting the whole buffer. Keep the final
evaluate()anyway. It is the copy of record, and its receipt is the one to log.Changed in 0.2.1
In 0.2.0 a session redacted a value that was still arriving, which produced text such as
[REDACTED_EMAIL]omwith short tokens. 0.2.1 keeps the buffer raw, so that no longer happens and you no longer need to coalesce tokens into words before you push them.Keep the raw parts in memory only for the length of one reply, and never log them.
How it works#
Streaming is a sliding window over raw text. Each push() appends the token and evaluates the whole unreleased buffer with the same pack as evaluate(). It then looks for the latest safe delimiter (space, newline, tab, carriage return, .,;!?()[]{}) that is at least lookahead_chars characters behind the end and where redacting the text before it and the text after it separately gives the same result as redacting the whole buffer. It releases the remedied text up to that point and keeps the rest raw. If there is no delimiter and the buffer grows beyond twice the lookahead, it tries every cut that keeps the last lookahead_chars characters, with the same check. lookahead_chars defaults to 28 and cannot go below 16.
Because the unreleased buffer is always re-checked as a whole, a value split across tokens ("123", "-45", "-6789") is caught once it completes, before it is released. Because the buffer stays raw and a cut is never placed inside a value, a value that is still arriving is not redacted early.
| API | On a fatal value with block_on_fatal=True |
|---|---|
session.push() / flush() | Returns one chunk with blocked=True; later pushes return []; session.is_blocked is True. |
stream_filter() / astream_filter() | Raises StreamBlockedError("Stream blocked by statutory invariant: <reason> (<action>)"), with action and reason attributes. |
With block_on_fatal=False, fatal matches are remedied like warnings (in the example pack, redacted like contact details) and the stream continues. The parameters and the chunk fields are listed in the Python reference.