Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
79 changes: 35 additions & 44 deletions docs/dev/migrate-streaming-v0.8.md
Original file line number Diff line number Diff line change
Expand Up @@ -8,9 +8,10 @@ finish.

That two-task model is replaced by a single-task primitive: `stream()` returns a
`Streamer` you consume directly with `async for`, on your own task. There is no
background orchestrator and no `acomplete()`. Typed events now come from the
`streaming_event` plugin hook rather than an iterator on the result. This is a
**breaking change with no deprecation shim** — call sites must be updated.
background orchestrator and no `acomplete()`. Typed events now come from an
`EventStreamer`, returned by `stream(as_events=True)` and iterated with `async
for`. This is a **breaking change with no deprecation shim** — call sites must be
updated.

## API mapping

Expand All @@ -19,7 +20,7 @@ background orchestrator and no `acomplete()`. Typed events now come from the
| `stream_with_chunking(...)` | `stream(...)` | Same arguments, except `chunking` now defaults to `None` (raw deltas) instead of `"sentence"`; pass `chunking="sentence"` to preserve v0.7 chunk boundaries |
| returns `StreamChunkingResult` | returns `Streamer` | Consume with `async for`, ideally inside `async with` |
| `async for chunk in result.astream()` | `async for chunk in streamer` | Iterate the `Streamer` directly |
| `async for event in result.events()` | `@hook("streaming_event")` plugin | Events move to the hook (see below) |
| `async for event in result.events()` | `async for event in streamer` | Iterate the `EventStreamer` directly (`stream(as_events=True)`) |
| `await result.acomplete()` | *(removed)* | Consuming the stream drives it to completion |
| `STREAMING_ORCHESTRATION_START`/`_END` hooks | *(removed)* | No replacement — the whole run is on one task now, so there is nothing to reattach a span across |
| `result.completed` | `not streamer.failed_early` | |
Expand Down Expand Up @@ -55,10 +56,11 @@ subclass it or pass an instance). Passing a string alias — `chunking="sentence

## Before and after: observing events

If you consumed typed events with `result.events()`, the events move to a
`streaming_event` plugin. Both snippets below produce the same output — they are
the `main()` from `docs/examples/streaming/validated_streaming.py`, v0.7 then
v0.8, using the same requirement, prompt, and chunking.
If you consumed typed events with `result.events()`, you now iterate an
`EventStreamer` (`stream(as_events=True)`). Both snippets below produce the same
output — they are the `main()` from
`docs/examples/streaming/validated_streaming.py`, v0.7 then v0.8, using the same
requirement, prompt, and chunking.

### v0.7

Expand Down Expand Up @@ -97,52 +99,41 @@ print(f"Full text: {result.full_text!r}")
### v0.8

```python
@hook("streaming_event")
async def print_events(payload, ctx) -> None:
event = payload.event
match event:
case ChunkEvent():
print(f" CHUNK[{event.chunk_index}]: {event.text!r}")
case QuickCheckEvent(passed=False):
print(
f" QUICK_CHECK[{event.chunk_index}]: FAIL — "
f"{event.results[0].reason if event.results else 'unknown reason'}"
)
case QuickCheckEvent():
print(f" QUICK_CHECK[{event.chunk_index}]: pass")
case StreamingDoneEvent():
print(f" STREAMING_DONE: {len(event.full_text)} chars accumulated")
case FullValidationEvent():
print(f" FULL_VALIDATION: {'PASS' if event.passed else 'FAIL'}")
case CompletedEvent():
print(f" COMPLETED: success={event.success}")
case _:
pass


register(print_events)

print("Stream events as they arrive:")
async with await stream(
action, backend, ctx, requirements=[req], chunking="sentence"
action, backend, ctx, requirements=[req], chunking="sentence", as_events=True
) as streamer:
# Draining the stream fires the events; the hook does the printing.
async for _chunk in streamer:
pass
async for event in streamer:
match event:
case ChunkEvent():
print(f" CHUNK[{event.chunk_index}]: {event.text!r}")
case QuickCheckEvent(passed=False):
print(
f" QUICK_CHECK[{event.chunk_index}]: FAIL — "
f"{event.results[0].reason if event.results else 'unknown reason'}"
)
case QuickCheckEvent():
print(f" QUICK_CHECK[{event.chunk_index}]: pass")
case StreamingDoneEvent():
print(f" STREAMING_DONE: {len(event.full_text)} chars accumulated")
case FullValidationEvent():
print(f" FULL_VALIDATION: {'PASS' if event.passed else 'FAIL'}")
case CompletedEvent():
print(f" COMPLETED: success={event.success}")
case _:
pass

print(f"\nCompleted normally: {streamer.completed_normally}")
print(f"Full text: {streamer.full_text!r}")
```

The three transformations to make:

1. **`stream_with_chunking(...)` → `async with await stream(...) as streamer:`**,
and iterate `streamer` directly. The `async with` replaces `acomplete()` for
cleanup.
2. **The `result.events()` match loop → a `@hook("streaming_event")` plugin.**
The event vocabulary and payloads are unchanged; register the plugin (with
`register()` or a `plugin_scope`) before consuming. Draining the stream is
what fires the events.
1. **`stream_with_chunking(...)` → `async with await stream(..., as_events=True)
as streamer:`**, and iterate `streamer` directly. The `async with` replaces
`acomplete()` for cleanup.
2. **`async for event in result.events()` → `async for event in streamer`.** The
event vocabulary and payloads are unchanged.
3. **`result.<attr>` → `streamer.<attr>`**. `result.completed` maps directly to
`not streamer.failed_early`; prefer the new `streamer.completed_normally` if
you need a signal that also excludes an early `break` (used above).
Expand Down
41 changes: 38 additions & 3 deletions docs/docs/how-to/use-async-and-streaming.md
Original file line number Diff line number Diff line change
Expand Up @@ -310,9 +310,44 @@ async def main() -> None:
asyncio.run(main())
```

To observe the run through typed `StreamEvent` objects instead, register a
plugin on the `streaming_event` hook — see the
[streaming validation tutorial](../tutorials/06-streaming-validation.md).
### Consuming events with `EventStreamer`

The `Streamer` above yields validated chunks. To consume the run's typed events
instead, pass `as_events=True` to `stream()`: it returns an `EventStreamer` you
iterate the same way, each step yielding an event rather than a chunk:

```python
from mellea.stdlib.streaming import (
ChunkEvent,
CompletedEvent,
FullValidationEvent,
QuickCheckEvent,
StreamingDoneEvent,
)

async with await stream(
action, m.backend, m.ctx, requirements=[req], chunking="sentence", as_events=True
) as streamer:
async for event in streamer:
match event:
case ChunkEvent():
Comment thread
planetf1 marked this conversation as resolved.
print(f" chunk[{event.chunk_index}]: {event.text!r}")
case QuickCheckEvent(passed=False):
print(f" quick-check[{event.chunk_index}]: FAIL")
case StreamingDoneEvent():
print(f" done — {len(event.full_text)} chars")
case FullValidationEvent():
print(f" final validation: {'pass' if event.passed else 'fail'}")
case CompletedEvent():
print(f" completed — success={event.success}")
case _:
pass # ErrorEvent

print(f"Completed normally: {streamer.completed_normally}")
```

See the [Streaming Validation tutorial](../tutorials/06-streaming-validation.md)
for a full walkthrough.

### The `stream_validate` tri-state

Expand Down
51 changes: 21 additions & 30 deletions docs/docs/tutorials/06-streaming-validation.md
Original file line number Diff line number Diff line change
Expand Up @@ -18,7 +18,7 @@ By the end you will have covered:
- Early-exit cancellation and reading `streaming_failures`
- Choosing between `"word"`, `"sentence"`, and `"paragraph"` chunking
- Observing the typed event vocabulary (`ChunkEvent`, `QuickCheckEvent`, …)
through the `streaming_event` hook
by iterating `stream(as_events=True)`
- Subclassing `ChunkingStrategy` to define a custom split boundary

**Prerequisites:** [Tutorial 02](./streaming-and-async) (async and streaming),
Expand Down Expand Up @@ -322,11 +322,10 @@ common English word like `"and"` or `"the"`.

## Step 4: Observing the event lifecycle

The `async for` loop gives you validated chunks. To observe the full lifecycle —
per-chunk validation results, stream completion, final validation, errors —
subscribe to the `streaming_event` hook. `stream()` fires one typed `StreamEvent`
per lifecycle moment through this hook, so a plugin can watch a run without
touching the chunk iterator.
The `async for` loop above yields validated chunks. To observe the full lifecycle
instead — per-chunk validation results, stream completion, final validation,
errors — pass `as_events=True`: `stream()` then returns an `EventStreamer` that
yields one typed `StreamEvent` per lifecycle moment in place of chunks.

```python
# Requires: mellea
Expand All @@ -337,7 +336,6 @@ import re
from mellea.core.backend import Backend
from mellea.core.base import Context
from mellea.core.requirement import PartialValidationResult, Requirement, ValidationResult
from mellea.plugins import hook, register
from mellea.stdlib.components import Instruction
from mellea.stdlib.streaming import (
ChunkEvent,
Expand Down Expand Up @@ -376,40 +374,33 @@ class MaxSentencesReq(Requirement):
return ValidationResult(result=self._count <= self._limit)


@hook("streaming_event")
async def print_events(payload, ctx) -> None:
event = payload.event
match event:
case ChunkEvent():
print(f" chunk[{event.chunk_index}]: {event.text!r}")
case QuickCheckEvent(passed=False):
print(f" FAIL at chunk {event.chunk_index}: {event.results[0].reason}")
case StreamingDoneEvent():
print(f" stream done — {len(event.full_text)} chars")
case FullValidationEvent():
print(f" final validation: {'pass' if event.passed else 'fail'}")
case CompletedEvent():
print(f" completed — success={event.success}")
case _:
pass


async def main() -> None:
from mellea.stdlib.session import start_session

m = start_session()
register(print_events)

# Draining the stream drives generation; print_events fires per event.
async with await stream(
Instruction("Write a two-sentence summary of the water cycle."),
m.backend,
m.ctx,
requirements=[MaxSentencesReq(limit=3)],
chunking="sentence",
as_events=True,
) as streamer:
async for _chunk in streamer:
pass
async for event in streamer:
match event:
case ChunkEvent():
print(f" chunk[{event.chunk_index}]: {event.text!r}")
case QuickCheckEvent(passed=False):
print(f" FAIL at chunk {event.chunk_index}: {event.results[0].reason}")
case StreamingDoneEvent():
print(f" stream done — {len(event.full_text)} chars")
case FullValidationEvent():
print(f" final validation: {'pass' if event.passed else 'fail'}")
case CompletedEvent():
print(f" completed — success={event.success}")
case _:
pass


asyncio.run(main())
Expand Down Expand Up @@ -577,7 +568,7 @@ explicitly or subclass to override `flush()`.
| `stream()` + `requirements=` | Per-chunk validation with automatic early exit |
| `async for chunk in streamer` | Validated chunks as they arrive, inside `async with` for safe cleanup |
| `streamer.failed_early` / `streamer.streaming_failures` | Detect and inspect a mid-stream requirement failure |
| `streaming_event` hook | Typed event stream — observe every chunk, validation result, and lifecycle signal |
| `stream(as_events=True)` | Typed event stream (an `EventStreamer`) — observe every chunk, validation result, and lifecycle signal |
| `"word"` / `"sentence"` / `"paragraph"` | Built-in chunking strategies trading reaction speed for context |
| `ChunkingStrategy` subclass | Custom split boundaries for structured output (lists, code, CSV) |

Expand Down
4 changes: 2 additions & 2 deletions docs/examples/streaming/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -22,7 +22,7 @@ Demonstrates:
- Streaming token-by-token generation
- Sentence-level chunking via `stream()`
- Per-chunk validation with custom `stream_validate()` methods
- Accessing stream events via the STREAMING_EVENT hook
- Iterating typed stream events via `stream(as_events=True)`
- Early exit on validation failure

### Word-Level Chunking
Expand Down Expand Up @@ -76,7 +76,7 @@ Demonstrates:

**Stream Validation**: Apply requirements at chunk level for early exit—stop generation when a constraint is violated.

**Stream Events**: Process stream events through the `STREAMING_EVENT` hook to monitor generation progress:
**Stream Events**: Observe a run's lifecycle as typed events — iterate them for one stream with `stream(as_events=True)`, or subscribe the `STREAMING_EVENT` hook to watch many streams at once:
- `ChunkEvent` — A new chunk of text
- `QuickCheckEvent` — Initial validation result
- `FullValidationEvent` — Complete validation after full generation
Expand Down
52 changes: 21 additions & 31 deletions docs/examples/streaming/custom_chunking.py
Original file line number Diff line number Diff line change
Expand Up @@ -26,7 +26,6 @@

import asyncio
import re
from typing import Any

from mellea.core.backend import Backend
from mellea.core.base import Context
Expand All @@ -35,7 +34,6 @@
Requirement,
ValidationResult,
)
from mellea.plugins import hook, register
from mellea.stdlib.chunking import ChunkingStrategy
from mellea.stdlib.components import Instruction
from mellea.stdlib.streaming import (
Expand Down Expand Up @@ -143,37 +141,29 @@ async def main() -> None:
chunker = LineChunking()
req = NumberedLineReq()

@hook("streaming_event")
async def print_events(payload: Any, ctx: Any) -> None:
event = payload.event
match event:
case ChunkEvent():
print(f" LINE[{event.chunk_index}]: {event.text!r}")
case QuickCheckEvent(passed=False):
print(
f" QUICK_CHECK[line {event.chunk_index}]: FAIL — "
f"{event.results[0].reason if event.results else 'unknown'}"
)
case QuickCheckEvent():
print(f" QUICK_CHECK[line {event.chunk_index}]: pass")
case StreamingDoneEvent():
print(f" STREAMING_DONE: {len(event.full_text)} chars accumulated")
case FullValidationEvent():
print(f" FULL_VALIDATION: {'PASS' if event.passed else 'FAIL'}")
case CompletedEvent():
print(f" COMPLETED: success={event.success}")
case _:
pass

register(print_events)

print("Stream events as they arrive (one per line):")
print("Stream events as they arrive (one ChunkEvent per line):")
async with await stream(
action, backend, ctx, requirements=[req], chunking=chunker
action, backend, ctx, requirements=[req], chunking=chunker, as_events=True
) as streamer:
# Draining the stream fires the events; the hook does the printing.
async for _line in streamer:
pass
async for event in streamer:
match event:
case ChunkEvent():
print(f" LINE[{event.chunk_index}]: {event.text!r}")
case QuickCheckEvent(passed=False):
print(
f" QUICK_CHECK[line {event.chunk_index}]: FAIL — "
f"{event.results[0].reason if event.results else 'unknown'}"
)
case QuickCheckEvent():
print(f" QUICK_CHECK[line {event.chunk_index}]: pass")
case StreamingDoneEvent():
print(f" STREAMING_DONE: {len(event.full_text)} chars accumulated")
case FullValidationEvent():
print(f" FULL_VALIDATION: {'PASS' if event.passed else 'FAIL'}")
case CompletedEvent():
print(f" COMPLETED: success={event.success}")
case _:
pass

print(f"\nCompleted normally: {streamer.completed_normally}")
if streamer.streaming_failures:
Expand Down
3 changes: 2 additions & 1 deletion docs/examples/streaming/multi_stream_events.py
Original file line number Diff line number Diff line change
Expand Up @@ -5,7 +5,8 @@
`stream()` yields validated chunks through `async for`, but its typed lifecycle
events (`ChunkEvent`, `QuickCheckEvent`, `StreamingDoneEvent`,
`FullValidationEvent`, `CompletedEvent`, `ErrorEvent`) are surfaced through the
`STREAMING_EVENT` plugin hook rather than the iterator.
`STREAMING_EVENT` plugin hook — or, for a single stream, from
`stream(as_events=True)`.

Demonstrates:
- Registering a `@hook("streaming_event")` function to receive every stream's events
Expand Down
Loading
Loading