Stream progress
Stream progress
Workflow streams expose ordered lifecycle, step, result-preview, and diagnostic events. They are for observation. Retrieve the stored result after termination.
Start and consume a stream
from makra import EventTypes, Makra
run_id = None
with Makra() as client:
for event in client.extract_stream(
["https://shop.example/products/atlas-lamp"],
{"price": "The current selling price"},
):
run_id = event.run_id or run_id
print(event.sequence, event.type, event.detail_type)
if event.type == EventTypes.DIAGNOSTIC:
print(event.payload)
if not run_id:
raise RuntimeError("The stream did not provide a run ID")
result = client.get_run_result(run_id)
WorkflowEvent.payload preserves the API event object. Convenience properties expose detail_type, status, reason, and terminal success when present.
Terminal events
The stream ends after completed, failed, cancelled, or budget-exhausted events. event.is_terminal identifies all terminal event types. A terminal completed event can still report domain-level failure inside its payload, so inspect event.success and fetch the result envelope.
Resume an existing run
If your application stores the last processed event sequence, it can attach to the durable run stream later.
from makra import Makra
run_id = "your-run-id"
last_sequence = 42
with Makra() as client:
for event in client.stream_run_events(
run_id,
last_event_id=last_sequence,
):
last_sequence = event.sequence
print(event.type)
Persist the sequence after processing the event. Event handling should be idempotent because reconnect boundaries can overlap with application persistence boundaries.
Async iteration
from makra import AsyncMakra
async def watch() -> None:
async with AsyncMakra() as client:
events = client.extract_stream(
["https://shop.example/products/atlas-lamp"],
{"price": "The current selling price"},
)
async for event in events:
print(event.type)
A heartbeat keeps a healthy idle stream alive. stream_idle_timeout bounds the gap between received bytes, not total workflow duration.
Next, manage deferred runs.