Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
75 commits
Select commit Hold shift + click to select a range
7af879d
Added the temporalio.streams interface with a provider registry.
moedash Sep 15, 2026
c2f0251
Added the in-memory reference provider and the conformance tests.
moedash Sep 15, 2026
2815caf
Added the shared demo loop.
moedash Sep 15, 2026
ccae2b4
Added the prepare and drain hooks the Workflow Streams transport needs.
moedash Sep 16, 2026
8926913
Made cursors exclusive, added latest(), and allowed topic appends.
moedash Sep 16, 2026
8e36829
Sorted and formatted the streams package for ruff.
moedash Sep 16, 2026
a627022
Added the changelog entry for the stream interface.
moedash Sep 16, 2026
0fbac29
Made the streams package pass the type and doc linters.
moedash Sep 16, 2026
6b8b35e
Made the demo call the provider lifecycle hooks.
moedash Sep 16, 2026
696d355
Woke outside readers across threads and ignored idle_timeout in memory.
moedash Sep 18, 2026
ba75070
Named the last record in append's cursor and allowed None for it.
moedash Sep 18, 2026
9de9342
Resolved producer identity once in the package.
moedash Sep 18, 2026
fe4ec72
Required an explicit provider and exported instance.
moedash Sep 18, 2026
0af0389
Typed the handles and records honestly.
moedash Sep 18, 2026
8a2bf8d
Kept topics off inbound frames and skipped unreadable frames with a w…
moedash Sep 18, 2026
35b1212
Let a broken provider import surface instead of vanishing from the re…
moedash Sep 18, 2026
6e48b54
Added the workflow-side conformance tests.
moedash Sep 18, 2026
360d054
Replaced stale wording and renamed the demo types.
moedash Sep 18, 2026
936a725
Keyed stores by one helper that survives a colon in the workflow id.
moedash Sep 19, 2026
ad7c577
Added an async close hook next to prepare and drain.
moedash Sep 19, 2026
27fb594
Parametrised the conformance suite over providers.
moedash Sep 19, 2026
0972c50
Renamed the store key helper to inbound_stream_id for the provider br…
moedash Sep 19, 2026
144580c
Added the StreamError family for stream conditions.
moedash Sep 21, 2026
22d79ec
Vendored the stream additions to the public API protos.
moedash Sep 21, 2026
1effd48
Split the stream provider in two and made the proto the record.
moedash Sep 21, 2026
ebbdfa2
Carried the stream provider from the worker into the workflow runtime.
moedash Sep 21, 2026
ffae7a4
Moved the stream conformance suite onto the new surface.
moedash Sep 21, 2026
3f16fd3
Moved the stream demo and changelog entry onto the new surface.
moedash Sep 21, 2026
e5f21df
Let a conformance case start the workflow that owns a stream.
moedash Sep 21, 2026
2799b58
Kept the stream finish hook off evicted runs and collected coroutines.
moedash Sep 21, 2026
297d3c5
Typed the hook test's workflow input as awaitable.
moedash Sep 21, 2026
099c52c
Added activity.stream_handle and Client.get_stream_handle.
moedash Sep 21, 2026
22e3864
Registered the demo's provider on the client and reported from the ac…
moedash Sep 21, 2026
40bc91e
Added typed topic definitions.
moedash Sep 21, 2026
b068712
Listed stream_provider in ClientConnectConfig.
moedash Sep 22, 2026
1a8d19f
Repinned Core to a protos-only branch on the commit main pins.
moedash Sep 22, 2026
55200ec
Added activity.stream_handle.
moedash Sep 25, 2026
e94d11d
Added Client.get_stream_handle and the connect keyword.
moedash Sep 25, 2026
fe446e6
Moved the conformance cases onto the client accessor.
moedash Sep 25, 2026
717f58b
Added the streams demo and the changelog entry.
moedash Sep 25, 2026
2337651
Tied the series to the original branch head.
moedash Sep 25, 2026
c84eb66
Merge branch 'moe/AI-198-py-03-workflow-runtime' into moe/AI-198-py-0…
moedash Sep 25, 2026
fc92161
Told a sync activity what it is actually missing.
moedash Sep 25, 2026
e57bd02
Put the streams demo under the type checkers.
moedash Sep 25, 2026
a03d091
Merge branch 'moe/AI-198-py-03-workflow-runtime' into moe/AI-198-py-0…
moedash Sep 25, 2026
64530d0
Merge branch 'moe/AI-198-py-03-workflow-runtime' into moe/AI-198-py-0…
moedash Sep 25, 2026
9a41d49
Merge branch 'moe/AI-198-py-03-workflow-runtime' into moe/AI-198-py-0…
moedash Sep 26, 2026
150e87d
Merge branch 'moe/AI-198-py-03-workflow-runtime' into moe/AI-198-py-0…
moedash Sep 28, 2026
1466028
Documented the default topic on the stream accessors.
moedash Sep 28, 2026
92c4b8c
Let activities own streams and resolve stream_handle() by a static rule.
moedash Sep 28, 2026
be51186
Tested the activity stream accessor rule and retry supersession.
moedash Sep 28, 2026
703065a
Merge branch 'moe/AI-198-py-03-workflow-runtime' into moe/AI-198-py-0…
moedash Sep 28, 2026
06cfa01
Documented the read starts on both accessors and tested an activity s…
moedash Sep 28, 2026
d216812
Merge branch 'moe/AI-198-py-03-workflow-runtime' into moe/AI-198-py-0…
moedash Sep 30, 2026
61a58f7
Let the client and activity accessors open refs and standalone streams.
moedash Sep 30, 2026
a95c913
Hosted standalone streams in the memory provider.
moedash Sep 30, 2026
5287f6e
Tested refs, standalone streams and body encoding through the accessors.
moedash Sep 30, 2026
dc2f448
Merge branch 'moe/AI-198-py-03-workflow-runtime' into moe/AI-198-py-0…
moedash Sep 30, 2026
f5a7e1b
Held the retention case to the declared capabilities and used the cas…
moedash Sep 30, 2026
5627560
Merge branch 'moe/AI-198-py-03-workflow-runtime' into moe/AI-198-py-0…
moedash Oct 1, 2026
f05f6c3
Let the retention case accept a refused append and a coarse age sweep.
moedash Oct 1, 2026
bf31846
Merge branch 'moe/AI-198-py-03-workflow-runtime' into moe/AI-198-py-0…
moedash Oct 1, 2026
e69e3e0
Merge branch 'moe/AI-198-py-03-workflow-runtime' into moe/AI-198-py-0…
moedash Oct 1, 2026
1580fc9
Merge branch 'moe/AI-198-py-03-workflow-runtime' into moe/AI-198-py-0…
moedash Oct 2, 2026
09d8ab0
Merge branch 'moe/AI-198-py-03-workflow-runtime' into moe/AI-198-py-0…
moedash Oct 2, 2026
651c32e
Merge branch 'moe/AI-198-py-03-workflow-runtime' into moe/AI-198-py-0…
moedash Oct 2, 2026
5017569
Merge branch 'moe/AI-198-py-03-workflow-runtime' into moe/AI-198-py-0…
moedash Oct 2, 2026
99df1e6
Merge branch 'moe/AI-198-py-03-workflow-runtime' into moe/AI-198-py-0…
moedash Oct 2, 2026
035d371
Merge branch 'moe/AI-198-py-03-workflow-runtime' into moe/AI-198-py-0…
moedash Oct 2, 2026
270c760
Merge branch 'moe/AI-198-py-03-workflow-runtime' into moe/AI-198-py-0…
moedash Oct 2, 2026
86d4d50
Merge branch 'moe/AI-198-py-03-workflow-runtime' into moe/AI-198-py-0…
moedash Oct 2, 2026
4eb1047
Merge branch 'moe/AI-198-py-03-workflow-runtime' into moe/AI-198-py-0…
moedash Oct 2, 2026
10dd23f
Merge branch 'moe/AI-198-py-03-workflow-runtime' into moe/AI-198-py-0…
moedash Oct 3, 2026
4af0f93
Merge branch 'moe/AI-198-py-03-workflow-runtime' into moe/AI-198-py-0…
moedash Oct 3, 2026
44dfd4f
Merge branch 'moe/AI-198-py-03-workflow-runtime' 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
20 changes: 20 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -37,6 +37,26 @@ to include examples, links to docs, or any other relevant information.
worker-side factories registered with `StrandsPlugin(sandboxes=...)`.

- Added the `temporalio.contrib.gcp.cloud_run.id` module with the `CloudRunIdPlugin` client plugin to set the worker identity on Cloud Run.
- **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, and a call that names
no topic addresses the default topic, `streams.DEFAULT_TOPIC` (`"output"`,
the server's default stream name). 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 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.

### Changed

Expand Down
2 changes: 1 addition & 1 deletion pyproject.toml
Original file line number Diff line number Diff line change
Expand Up @@ -267,7 +267,7 @@ reportUnnecessaryIsInstance = "none"
reportUnnecessaryTypeIgnoreComment = "none"
reportUnusedCallResult = "none"
reportUnknownLambdaType = "none"
include = ["temporalio", "tests"]
include = ["temporalio", "tests", "streams_demo"]
exclude = [
# Exclude auto generated files
"temporalio/api",
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()
174 changes: 174 additions & 0 deletions streams_demo/run_demo.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,174 @@
"""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

# Run as a script rather than imported as a module, so the two siblings are
# reached by name off this directory rather than through a package path.
sys.path.insert(0, str(Path(__file__).resolve().parent))
import provider_setup # noqa: E402 # pyright: ignore[reportImplicitRelativeImport]
from agent_loop import ( # noqa: E402 # pyright: ignore[reportImplicitRelativeImport]
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