Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
39 commits
Select commit Hold shift + click to select a range
c27bfcf
Added the temporalio.streams interface on the external base.
moedash Sep 25, 2026
dad0743
Answered a legacy query without the stream snapshot beside it.
moedash Sep 25, 2026
c2a63e1
Merged the external repairs with the review-round fixes.
moedash Sep 25, 2026
588b41c
Refused a stream publish from a query handler.
moedash Sep 25, 2026
d69847f
Said so when a producer's attempt goes backwards.
moedash Sep 25, 2026
9b60101
Merge branch 'moe/AI-198-py-08-external-repairs' into moe/AI-198-py-0…
moedash Sep 26, 2026
375a352
Let activities own streams and resolve stream_handle() by a static rule.
moedash Sep 29, 2026
3504572
Tested the activity stream accessor rule and retry supersession.
moedash Sep 29, 2026
60e8e03
Mapped BEGINNING to the oldest record held and added END and last=N.
moedash Sep 28, 2026
9cc17dc
Added conformance cases for each read start, on a truncated topic too.
moedash Sep 28, 2026
b81e932
Let stream_reader() start at END or at the last N records.
moedash Sep 28, 2026
24452e3
Documented the read starts on both accessors and tested an activity s…
moedash Sep 28, 2026
4191d4d
Ran stream bodies through the data converter with a plaintext hash.
moedash Sep 30, 2026
f52bb82
Declared standalone streams and a handle close on the provider surface.
moedash Sep 30, 2026
1897c91
Added StreamRef, a serializable name for one stream.
moedash Sep 30, 2026
3019a58
Let the client and activity accessors open refs and standalone streams.
moedash Sep 30, 2026
4e20489
Hosted standalone streams in the memory provider.
moedash Sep 30, 2026
333e3f7
Tested refs, standalone streams and body encoding through the accessors.
moedash Sep 30, 2026
f4e314d
Fitted the ported interface commits to this chain.
moedash Sep 30, 2026
e1f115f
Declared the byte bound and age trim as standalone capabilities.
moedash Sep 30, 2026
1aceee7
Held the retention case to the declared capabilities and used the cas…
moedash Sep 30, 2026
485fb08
Merge remote-tracking branch 'origin/moe/AI-198-py-08-external-repair…
moedash Oct 1, 2026
48874dd
Merge remote-tracking branch 'origin/moe/AI-198-py-08-external-repair…
moedash Oct 1, 2026
94a2d5d
Merge remote-tracking branch 'origin/moe/AI-198-py-08-external-repair…
moedash Oct 1, 2026
41ddb61
Added a conformance case for a publish from the constructor.
moedash Oct 1, 2026
bdb7d94
Merge branch 'moe/AI-198-py-08-external-repairs' into moe/AI-198-py-0…
moedash Oct 2, 2026
fa7ae83
Merge branch 'moe/AI-198-py-08-external-repairs' into moe/AI-198-py-0…
moedash Oct 2, 2026
c0a0d18
Added a conformance case for a reader woken through the channel.
moedash Oct 2, 2026
58c87ba
Marked the channel conformance case for a channel server.
moedash Oct 2, 2026
55212a7
Merge branch 'moe/AI-198-py-08-external-repairs' into moe/AI-198-py-0…
moedash Oct 2, 2026
fb67a8f
Added a conformance case for a reader woken through its linked channel.
moedash Oct 2, 2026
c88c9d8
Merge branch 'moe/AI-198-py-08-external-repairs' into moe/AI-198-py-0…
moedash Oct 2, 2026
278a377
Merge branch 'moe/AI-198-py-08-external-repairs' into moe/AI-198-py-0…
moedash Oct 2, 2026
494d87c
Merge branch 'moe/AI-198-py-08-external-repairs' into moe/AI-198-py-0…
moedash Oct 2, 2026
790f833
Added conformance cases for a reader leaving its channel.
moedash Oct 2, 2026
662eba9
Checked that the subscription waits for the completion that ends the …
moedash Oct 2, 2026
fbbaf0c
Merge branch 'moe/AI-198-py-08-external-repairs' into moe/AI-198-py-0…
moedash Oct 3, 2026
58baa2a
Read the linked owner as an execution in the conformance case.
moedash Oct 3, 2026
e9ae44d
Merge branch 'moe/AI-198-py-08-external-repairs' into moe/AI-198-py-0…
moedash Oct 3, 2026
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
3 changes: 3 additions & 0 deletions .gitignore
Original file line number Diff line number Diff line change
Expand Up @@ -15,3 +15,6 @@ temporalio/bridge/temporal_sdk_bridge*
tags
/.claude
tmpclaude-*

# Demo run output, written per provider; the numbers live in the doc.
streams_demo/results-*/
19 changes: 19 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -35,6 +35,25 @@ to include examples, links to docs, or any other relevant information.

### Added

- **Experimental**: `temporalio.streams` defines one stream interface a workflow
can read, decide on, and write. A provider is registered once as a plugin,
`Client.connect(plugins=[provider])`, and workers built from that client
inherit it; each context then asks for its stream the same way:
`workflow.stream_reader()` and `workflow.stream_writer()` in workflow code,
`activity.stream_handle()` in an activity, and `client.get_stream_handle()`
anywhere a client is held. A topic is a typed definition,
`streams.topic("inputs", Token)`, shared by workflow, activity and client
code; a plain string names a topic decided at runtime. The record on the wire
is `temporal.api.stream.v1.StreamRecord` on every provider. A stream is
handed to another process as a `streams.StreamRef`, plain data naming the
owner and, when it has one, the topic, which `client.get_stream_handle(ref)`
and `activity.stream_handle(ref)` open; `client.create_stream(stream_id, ...)`
creates a standalone stream with a retention policy, and its handle's
`close()` seals it. A provider runs record bodies through the client's data
converter, so a payload codec and external storage apply to them.
`temporalio.streams.providers.memory.MemoryStreams` is the in-memory
reference provider the conformance tests run against.

- Added experimental External Workflow Streams in
`temporalio.contrib.external_workflow_streams`. Workflow stream payloads are
stored in a configured external backend instead of Temporal History, with a
Expand Down
110 changes: 110 additions & 0 deletions streams_demo/agent_loop.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,110 @@
"""One agent loop, written once, run on every stream provider.

Reads input, decides on it, publishes the decision, and runs an ordinary
Activity in the same workflow task; the Activity reports on its own
workflow's stream in turn. Also handles the two control records the contract
defines, so a retried producer and a finished topic are exercised rather than
described.

This file is byte-identical in the server-side tree and the client-side tree.
Nothing in it names a provider: the workflow asks its runtime, the Activity
asks its context, and the process that runs them registered the provider once.
"""

from __future__ import annotations

from datetime import timedelta
from typing import Any

from temporalio import activity, streams, workflow
from temporalio.streams import RecordKind

# Defined once and shared by the workflow, the Activity and the demo's
# reader, so the type each topic carries is stated in one place.
DECISIONS = streams.topic("decisions", dict[str, Any])
INPUTS = streams.topic("inputs", dict[str, Any])
RECEIPTS = streams.topic("receipts", dict[str, Any])


@activity.defn(name="RecordDecision")
async def record_decision(decision: dict[str, Any]) -> str:
"""An ordinary command in the same task as the publish.

Appends its receipt onto its own workflow's stream too, under the
Activity's own identity, so a reader outside sees the decision and the
record of it side by side.
"""
receipt = f"recorded:{decision['source']}:{decision['branch']}"
await activity.stream_handle().producer(topic=RECEIPTS).append({"receipt": receipt})
return receipt


def decide(token: dict[str, Any]) -> dict[str, Any]:
"""The decision the workflow is here to make."""
if token["value"] % 2 == 0:
return {
"source": token["id"],
"branch": "even",
"computed": token["value"] * 10,
}
return {"source": token["id"], "branch": "odd", "computed": token["value"] + 100}


@workflow.defn(name="StreamContractDemo", sandboxed=False)
class AgentLoop:
"""Read, decide, write, until the producer says it has finished."""

@workflow.run
async def run(self, limit: int) -> list[dict[str, Any]]:
"""Decide on at most ``limit`` inputs, then return the trace."""
inputs = workflow.stream_reader(INPUTS)
decisions = workflow.stream_writer(DECISIONS)
trace: list[dict[str, Any]] = []
accepted = 0
try:
async for record in inputs:
if record.kind is RecordKind.SUPERSEDED:
# A newer attempt of the same producer started writing. The
# decisions already published stand, so the workflow says so
# rather than pretending they can be withdrawn.
assert record.supersession is not None
trace.append(
{
"kind": "superseded",
"producer": record.producer_id,
"replaced": record.supersession.previous_attempt,
"attempt": record.supersession.attempt,
}
)
decisions.publish(
{"retracting_attempt": record.supersession.previous_attempt}
)
continue
if record.kind is RecordKind.FINISH:
# The producer says it is done, which is what ends the loop.
# Counting decisions instead would leave the terminal record
# unread and let the workflow finish while its producer is
# still writing.
trace.append({"kind": "finish", "producer": record.producer_id})
break
assert record.value is not None
decision = decide(record.value)
decisions.publish(decision)
receipt = await workflow.execute_activity(
record_decision,
decision,
activity_id=f"decision-{decision['source']}",
start_to_close_timeout=timedelta(seconds=10),
)
trace.append(
{"kind": "decision", "value": decision, "receipt": receipt}
)
accepted += 1
if accepted >= limit:
# A bound so a stuck producer cannot run this forever. The
# terminal record above is the ordinary way out.
break
finally:
inputs.close()
decisions.finish()
return trace
36 changes: 36 additions & 0 deletions streams_demo/provider_setup.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,36 @@
"""Pick the provider for a demo run from the environment.

``STREAMS_PROVIDER`` names a provider; this base tree carries only ``memory``,
and each provider branch adds its own name here. The demo needs a Temporal
server to run the workflow either way; ``TEMPORAL_ADDRESS`` points at it.
"""

from __future__ import annotations

import os

from temporalio.streams.providers import ProviderPlugin
from temporalio.streams.providers.memory import MemoryStreams

NAME = os.environ.get("STREAMS_PROVIDER", "memory")

# The memory provider is not replay-safe, so its demo keeps the cache warm.
# Storage providers run with the smallest cache they support instead.
WORKFLOW_CACHE = int(os.environ.get("STREAMS_WORKFLOW_CACHE", "512"))


async def open() -> tuple[str, ProviderPlugin]:
"""The server to connect to and the provider the worker and the client share."""
if NAME != "memory":
raise SystemExit(f"this tree carries no stream provider named {NAME!r}")
return os.environ.get("TEMPORAL_ADDRESS", "localhost:7233"), MemoryStreams()


async def close(provider: ProviderPlugin) -> None:
"""Let go of whatever :func:`open` acquired.

The memory provider holds no connection, so this is its ``close()`` and
nothing more. A provider branch that opens one closes it the same way, so
the demo's teardown reads the same on every provider.
"""
await provider.close()
172 changes: 172 additions & 0 deletions streams_demo/run_demo.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,172 @@
"""Run the shared agent loop against whichever provider is configured.

Three cases, the same on every provider:

- read, decide, publish and an ordinary Activity in the same workflow task,
with the smallest workflow cache the provider supports, so that as much of
the run as it allows is rebuilt rather than remembered;
- a producer whose second attempt supersedes its first, which the reader has
to report and the workflow has to act on;
- an outside reader following what the workflow published, and the receipts
the Activity appended on the workflow's stream from inside its own context.

The provider is registered once, on the client; the worker inherits it and
every context asks for its stream without naming it. Byte-identical in every
tree. ``provider_setup`` is what differs, and it is the only import here that
names a provider.
"""

from __future__ import annotations

import asyncio
import json
import sys
import time
import uuid
from pathlib import Path
from typing import Any

from temporalio.api.enums.v1 import EventType
from temporalio.client import Client
from temporalio.streams import RecordKind, StreamHandle
from temporalio.worker import Worker

sys.path.insert(0, str(Path(__file__).resolve().parent))
import provider_setup # noqa: E402
from agent_loop import ( # noqa: E402
DECISIONS,
INPUTS,
RECEIPTS,
AgentLoop,
record_decision,
)

DECISION_LIMIT = 8
EXPECTED_OUTPUT = 6


async def collect_output(stream: StreamHandle, want: int) -> list[dict]:
"""Read ``want`` decisions off the workflow's stream from outside it."""
seen: list[dict] = []
async for record in stream.read(topic=DECISIONS):
seen.append({"kind": record.kind.name, "value": record.value})
if len(seen) >= want:
break
return seen


async def main() -> int:
"""Run the demo once and write what happened next to this file."""
out = Path(__file__).resolve().parent / f"results-{provider_setup.NAME}"
out.mkdir(exist_ok=True)
target, provider = await provider_setup.open()
client = await Client.connect(target, namespace="default", plugins=[provider])

uid = f"ai198-contract-{provider_setup.NAME}-" + uuid.uuid4().hex
record: dict[str, Any] = {
"provider": provider_setup.NAME,
"workflow_id": uid,
"target": target,
"max_cached_workflows": provider_setup.WORKFLOW_CACHE,
}

async with Worker(
client,
task_queue=uid,
workflows=[AgentLoop],
activities=[record_decision],
max_cached_workflows=provider_setup.WORKFLOW_CACHE,
):
handle = await client.start_workflow(
AgentLoop.run, DECISION_LIMIT, id=uid, task_queue=uid
)
stream = client.get_stream_handle(uid)
output = asyncio.create_task(collect_output(stream, EXPECTED_OUTPUT))

# The first attempt writes two records and then stops, as a failed
# activity would. The second writes different inputs under the same
# logical producer, which is what the reader has to report.
first = stream.producer(topic=INPUTS, producer_id="model", attempt=1)
await first.append({"id": "r1", "value": 1}, {"id": "r2", "value": 2})
await asyncio.sleep(0.5)
second = stream.producer(topic=INPUTS, producer_id="model", attempt=2)
await second.append({"id": "r3", "value": 3}, {"id": "r4", "value": 4})
await second.finish()

# A failed Workflow Task is not an outcome. The server rejects a
# completion that raced newly buffered events, and the retry usually
# gets through, so reading the first failure as the result reports a
# working run as a broken one. Wait for a terminal event, then report
# the retries separately so they are neither the headline nor hidden.
deadline = time.monotonic() + 120
# Taken from the enum rather than written out, because guessing these
# numbers is how a run that completed gets reported as terminated.
terminal = {
EventType.EVENT_TYPE_WORKFLOW_EXECUTION_COMPLETED: "completed",
EventType.EVENT_TYPE_WORKFLOW_EXECUTION_FAILED: "workflow_failed",
EventType.EVENT_TYPE_WORKFLOW_EXECUTION_TIMED_OUT: "execution_timed_out",
EventType.EVENT_TYPE_WORKFLOW_EXECUTION_CANCELED: "canceled",
EventType.EVENT_TYPE_WORKFLOW_EXECUTION_TERMINATED: "terminated",
EventType.EVENT_TYPE_WORKFLOW_EXECUTION_CONTINUED_AS_NEW: "continued_as_new",
}
while time.monotonic() < deadline:
history = await handle.fetch_history()
reached = [
terminal[e.event_type]
for e in history.events
if e.event_type in terminal
]
if reached:
record["outcome"] = reached[-1]
if record["outcome"] == "completed":
record["trace"] = await handle.result()
break
await asyncio.sleep(0.1)
else:
record["outcome"] = "timed_out_waiting"
history = await handle.fetch_history()

retried = [
e
for e in history.events
if e.event_type == EventType.EVENT_TYPE_WORKFLOW_TASK_FAILED
]
record["task_failures"] = [
{
"event_id": e.event_id,
"cause": int(e.workflow_task_failed_event_attributes.cause),
"message": e.workflow_task_failed_event_attributes.failure.message,
}
for e in retried
]

try:
record["observed_output"] = await asyncio.wait_for(output, timeout=20)
except asyncio.TimeoutError:
output.cancel()
record["observed_output"] = "timed_out"

async def receipts() -> list[Any]:
# Ends by itself once the workflow is closed and the tail served.
return [
r.value
async for r in stream.read(topic=RECEIPTS)
if r.kind is RecordKind.DATA
]

try:
record["receipts"] = await asyncio.wait_for(receipts(), timeout=20)
except asyncio.TimeoutError:
record["receipts"] = "timed_out"

history = await handle.fetch_history()
record["history_events"] = len(history.events)
(out / "history.json").write_text(history.to_json())
await provider_setup.close(provider)
(out / "results.json").write_text(json.dumps(record, indent=2) + "\n")
print(json.dumps(record, indent=2))
return 0 if record["outcome"] == "completed" else 1


if __name__ == "__main__":
sys.exit(asyncio.run(main()))
Loading
Loading