graphflow-stream
The astream_events LangGraph gives Python, for graph-flow in Rust.
Install
Usage
use ;
use ;
;
let = spawn_task;
while let Some = rx.recv.await
let result = handle.await??; // same TaskResult you'd get from task.run()
Just want the full streamed text back, no manual loop? One line:
let text = collect_text.await?;
emit_token/emit_started/emit_finished/emit_failed are ambient — call them from anywhere inside Task::run, no new trait to implement, no-op if nothing is listening. spawn_graph(flow_runner, session_id, buffer) does the same for a whole FlowRunner run; SubgraphTask wraps a nested Graph as one Task and streams through automatically.
Replay debugging: record(rx).await turns a run into a Recording (serializable, so it can be saved to disk), and recording.replay(buffer) plays it back on a fresh channel with the original timing — inspect a past run, or demo a UI without hitting an LLM again.
More examples (full_graph, sse_axum, websocket_axum) in examples/.
Benchmarks
cargo bench (criterion, benches/overhead.rs):
| Scenario | Time |
|---|---|
task.run() direct — no graphflow-stream involved |
~1.0 µs |
emit_token() with nobody listening (ambient no-op) |
~30 ns / call |
spawn_task() streaming 100 tokens to a draining receiver |
~2.0 µs / token |
record() capturing 100 streamed tokens |
~1.4 µs / token |
License
MIT