helixordevelopers

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.

Tier 030 minutesIntermediatePython 3.9+

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() and flush().
  • A fatal halt with block_on_fatal=True.
  • The drop-in generators stream_filter() and astream_filter(), including how they signal a block.
  • A production wrapper, stream_and_persist(), that stores an authoritative remedied copy of each reply.

Prerequisites#

Steps#

  1. 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
    
  2. 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_chars behind 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, empty text, the action, the reason and the fatal triggers only. Every later push() returns []. The SSN itself was never released.
    • chunk.redacted describes the chunk. It is True exactly on the chunks whose text was changed, here chunks 8 and 9.
  3. 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. Unlike push(), both raise StreamBlockedError on a fatal block, so wrap them in try. It is exported from helixor_runtime, carries action and reason, and is a subclass of RuntimeError, so existing except RuntimeError handlers 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.

  4. 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's clean_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]om with 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.

APIOn 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.

Next steps#