Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
36 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
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
13 changes: 13 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -21,6 +21,19 @@ to include examples, links to docs, or any other relevant information.
### Added

- 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. The record on the wire
is `temporal.api.stream.v1.StreamRecord` on every provider.
`temporalio.streams.providers.memory.MemoryStreams` is the in-memory
reference provider the conformance tests run against.

### Changed

### Deprecated
Expand Down
101 changes: 101 additions & 0 deletions scripts/gen_stream_api_protos.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,101 @@
"""Regenerate the stream additions to the vendored public API protos.

The stream protos live on a branch of the api repo that the Core submodule
does not pin yet. Regenerating everything from that branch would also pull in
whatever else moved on the api main line since the pin, so this stages the
API protos the submodule pins, applies the branch's own diff on top, and
regenerates only the files that diff touches. Every other module under
``temporalio/api`` stays byte-identical.

uv run --python 3.10 --no-project --with "grpcio-tools==1.48.2" \\
--with "mypy-protobuf==3.3.0" --with "protobuf<4" \\
scripts/gen_stream_api_protos.py /path/to/api origin/main..origin/<branch>

The toolchain pins match ``scripts/_proto/Dockerfile``. Generated code refuses
to load on a protobuf runtime older than the one it was built against, and
these modules end up in applications that pin protobuf themselves. Run
``uv run poe format`` afterwards, as the ``gen-protos`` task does.
"""

from __future__ import annotations

import shutil
import subprocess
import sys
import tempfile
from pathlib import Path

BASE = Path(__file__).parent.parent
sys.path.insert(0, str(BASE / "scripts"))

import gen_protos # noqa: E402

API_OUT = BASE / "temporalio" / "api"


def main() -> None:
if len(sys.argv) != 3:
sys.exit("usage: gen_stream_api_protos.py /path/to/api-checkout <base>..<ref>")
api = Path(sys.argv[1]).resolve()
revisions = sys.argv[2]
diff = subprocess.run(
["git", "-C", str(api), "diff", revisions, "--", "temporal/"],
check=True,
capture_output=True,
text=True,
).stdout
touched = subprocess.run(
["git", "-C", str(api), "diff", "--name-only", revisions, "--", "temporal/"],
check=True,
capture_output=True,
text=True,
).stdout.split()
if not touched:
sys.exit(f"{revisions} touches no proto under temporal/")

with tempfile.TemporaryDirectory() as tmp:
stage = Path(tmp) / "stage"
shutil.copytree(gen_protos.api_proto_dir / "temporal", stage / "temporal")
subprocess.run(
["patch", "-p1", "--silent"],
check=True,
cwd=stage,
input=diff,
text=True,
)
out = Path(tmp) / "out"
out.mkdir()
subprocess.check_call(
[
sys.executable,
"-mgrpc_tools.protoc",
f"--proto_path={stage}",
f"--python_out={out}",
f"--mypy_out={out}",
*touched,
]
)
packages: set[Path] = set()
for proto in touched:
relative = Path(proto).relative_to("temporal/api").with_suffix("")
for suffix in ("_pb2.py", "_pb2.pyi"):
generated = (
out
/ "temporal"
/ "api"
/ relative.with_name(relative.name + suffix)
)
target = API_OUT / relative.with_name(relative.name + suffix)
target.parent.mkdir(parents=True, exist_ok=True)
shutil.copyfile(generated, target)
print(f"wrote {target.relative_to(BASE)}")
packages.add(target.parent)

for package in sorted(packages):
(package.parent / "__init__.py").touch()
gen_protos.fix_generated_output(package)
print(f"rewrote {(package / '__init__.py').relative_to(BASE)}")


if __name__ == "__main__":
main()
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()
Loading
Loading