diff --git a/CHANGELOG.md b/CHANGELOG.md index 45360d889..f70089514 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -130,6 +130,8 @@ to include examples, links to docs, or any other relevant information. through `ExternalOutputStreamClient`. Workflow output is staged outside History and becomes readable only after its compact Workflow Task marker is committed. +- Added `examples/streams`, one agent loop that runs unchanged on every + stream provider and on the Nexus front. ### Changed diff --git a/examples/__init__.py b/examples/__init__.py new file mode 100644 index 000000000..e282df462 --- /dev/null +++ b/examples/__init__.py @@ -0,0 +1 @@ +"""Worked examples that run against a Temporal server.""" diff --git a/examples/streams/README.md b/examples/streams/README.md new file mode 100644 index 000000000..9513974c4 --- /dev/null +++ b/examples/streams/README.md @@ -0,0 +1,57 @@ +# Streams, by path + +One provider, registered once on the client. Workers built from that client +inherit it, and every context asks for its stream the same way. A topic is +defined once, with the type its records carry, and every context refers to +that definition, so no call names a type again. + +```python +client = await Client.connect("localhost:7233", plugins=[provider]) + +INPUTS = streams.topic("inputs", Token) +DECISIONS = streams.topic("decisions", Decision) +``` + +| Path | Who | Call | Example file | +|---|---|---|---| +| A: the workflow publishes | workflow code | `workflow.stream_writer(PROGRESS).publish(Progress(...))`, then `.finish()` | `path_a_publish.py` | +| A: a backend follows | any process with a client | `stream = client.get_stream_handle(workflow_id)`, then `stream.read(topic=PROGRESS, after=await stream.latest(topic=PROGRESS))` | `path_a_publish.py` | +| B: an Activity produces | activity code | `await activity.stream_handle().producer(topic=INPUTS).append(Token(...))` | `path_b_produce.py` | +| B: a backend produces | any process with a client | `client.get_stream_handle(workflow_id).producer(topic=NOTES, producer_id=..., attempt=...)` | `path_b_produce.py` | +| B: a backend consumes | any process with a client | `client.get_stream_handle(workflow_id).read(topic=INPUTS)` | `path_b_produce.py` | +| C: the workflow consumes | workflow code | `async for record in workflow.stream_reader(COMMANDS)` | `path_c_consume.py` | + +A plain string names a topic decided at runtime, with `result_type=` on the +call; the examples never need one. + +`agent.py` and `run.py` compose all three paths in one agent, on every provider +and behind the Nexus front. `_setup.py` is the one place a store is named. +`june_scenarios/` maps every scenario in Roey's June design notes onto this +surface, one file per scenario family, with a status on each. + +## Running + +Each example takes the provider's name and runs against a dev server: + +```sh +python -m examples.streams.path_a_publish workflow_streams +python -m examples.streams.path_a_publish native --address 127.0.0.1:7333 +python -m examples.streams.path_a_publish redis --redis redis://127.0.0.1:6379 +``` + +Swap `path_a_publish` for `path_b_produce`, `path_c_consume` or `run`. The +`native` provider needs a server built from the stream-carrying branch; the +`redis` provider needs a Redis to point at. `run.py` also takes `nexus`, with +an endpoint routed to the handler worker's task queue. + +## Why the workflow's verbs differ + +An Activity and a backend hold the same `StreamHandle`, with the same verbs, +because both act on the store at once: a producer's records are visible as +soon as the store accepts them, and a read follows the store live. Workflow +code gets two verbs of its own because its semantics differ. `publish` is +buffered and commits with the Workflow Task, so no reader can see a record +from a task that failed, and a `stream_reader` is an observation the SDK +records, so replay re-supplies the same records in the same order. That is +why `publish` is a plain call and the reader is an async iterator, and why +neither takes a client. diff --git a/examples/streams/__init__.py b/examples/streams/__init__.py new file mode 100644 index 000000000..f4ee8e3c9 --- /dev/null +++ b/examples/streams/__init__.py @@ -0,0 +1 @@ +"""One agent loop run on every stream provider.""" diff --git a/examples/streams/_setup.py b/examples/streams/_setup.py new file mode 100644 index 000000000..9765674e5 --- /dev/null +++ b/examples/streams/_setup.py @@ -0,0 +1,60 @@ +"""Provider selection for the examples: the one place a store is named. + +Every example takes the provider's name on the command line, builds it here, +and registers it once on the client. Nothing else in the examples names a +store: workers built from the client inherit the provider, and each context +asks for its stream through ``workflow.stream_reader`` or +``workflow.stream_writer``, ``activity.stream_handle()`` and +``client.get_stream_handle()``. +""" + +from __future__ import annotations + +import argparse + +from temporalio.client import Client +from temporalio.streams.providers import ProviderPlugin + +PROVIDERS = ("workflow_streams", "native", "redis") + + +def parser( + description: str, providers: tuple[str, ...] = PROVIDERS +) -> argparse.ArgumentParser: + """The flags every example shares.""" + parser = argparse.ArgumentParser(description=description) + parser.add_argument("provider", choices=providers) + parser.add_argument("--address", default="localhost:7233") + parser.add_argument("--redis", default="redis://127.0.0.1:6379") + return parser + + +def make_provider(name: str, args: argparse.Namespace) -> ProviderPlugin: + """The whole difference between the stores: one constructor call.""" + if name == "workflow_streams": + from temporalio.streams.providers.workflow_streams import ( + WorkflowStreamsProvider, + ) + + return WorkflowStreamsProvider() + if name == "redis": + from temporalio.streams.providers.redis import RedisStreams + + return RedisStreams(url=args.redis) + if name == "native": + from temporalio.streams.providers.native import NativeStreams + + return NativeStreams() + if name == "memory": + # Only for examples that keep a warm cache: this provider is not + # replay-safe, which is why PROVIDERS leaves it out. + from temporalio.streams.providers.memory import MemoryStreams + + return MemoryStreams() + raise SystemExit(f"unknown provider {name}") + + +async def connect(args: argparse.Namespace) -> tuple[Client, ProviderPlugin]: + """A client with the provider registered on it, and the provider to close later.""" + provider = make_provider(args.provider, args) + return await Client.connect(args.address, plugins=[provider]), provider diff --git a/examples/streams/agent.py b/examples/streams/agent.py new file mode 100644 index 000000000..b815c98ee --- /dev/null +++ b/examples/streams/agent.py @@ -0,0 +1,125 @@ +"""The workflow and activities. Identical on every provider. + +Nothing here names a store, a transport, or an option. The two topics are +defined once, with the types their records carry, and the workflow, the +Activity and the backend in ``run.py`` all refer to them. The loop reads its +``inputs`` topic, decides, publishes the decision, and runs an ordinary +activity in the same workflow task, which is the shape the design doc calls +Paths A, B and C together. The Activity that streams model output asks its +context for its own workflow's stream, the way workflow code asks its +runtime, so the file is the same whichever provider the process registered. +""" + +from __future__ import annotations + +import asyncio +from dataclasses import dataclass +from datetime import timedelta + +from temporalio import activity, streams, workflow +from temporalio.common import RetryPolicy +from temporalio.streams import RecordKind + + +@dataclass +class Token: + """One piece of model output.""" + + n: int + + +@dataclass +class Decision: + """What the workflow decided about a token, or which attempt it retracted.""" + + echo: int | None = None + retracting_attempt: int | None = None + + +INPUTS = streams.topic("inputs", Token) +DECISIONS = streams.topic("decisions", Decision) + + +@activity.defn +async def generate(count: int) -> None: + """Stream model output onto this workflow's ``inputs`` topic. + + No workflow id and no run id: the handle is this Activity's own + workflow, pinned to its run. The producer carries the Activity's own id + and attempt, so a retry deduplicates and a new attempt is reported to + readers as a supersession. + """ + model = activity.stream_handle().producer(topic=INPUTS) + for n in range(count): + await model.append(Token(n)) + await model.finish() + + +@activity.defn +async def record_decision(decision: Decision) -> str: + """An ordinary activity, run from the same task that read and published.""" + return f"recorded {decision.echo}" + + +@workflow.defn +class Agent: + """Reads ``inputs``, publishes a decision each time, ends on FINISH.""" + + @workflow.run + async def run(self, count: int) -> int: + """Decide on at most ``count`` inputs, then return how many landed.""" + decisions = workflow.stream_writer(DECISIONS) + + generating = workflow.start_activity( + generate, + count, + start_to_close_timeout=timedelta(minutes=1), + # Bounded, so a generator that cannot finish gives up instead of + # retrying forever while every attempt streams from the start. + retry_policy=RetryPolicy(maximum_attempts=3), + ) + consuming = asyncio.create_task(self._consume(count, decisions)) + + # Raced rather than awaited in turn: an attempt that fails writes no + # FINISH, so a generator that exhausts its attempts leaves the reader + # waiting forever. Its failure ends the run with its cause instead. + done, _ = await workflow.wait( + [consuming, generating], return_when=asyncio.FIRST_COMPLETED + ) + if generating in done and consuming not in done: + try: + await generating + except BaseException: + consuming.cancel() + raise + seen = await consuming + await generating + + decisions.finish() + return seen + + async def _consume( + self, count: int, decisions: workflow.StreamWriter[Decision] + ) -> int: + seen = 0 + async for record in workflow.stream_reader(INPUTS): + if record.kind is RecordKind.FINISH: + break + if record.kind is RecordKind.SUPERSEDED: + assert record.supersession is not None + decisions.publish( + Decision(retracting_attempt=record.supersession.previous_attempt) + ) + continue + assert record.value is not None + seen += 1 + decision = Decision(echo=record.value.n) + decisions.publish(decision) + await workflow.execute_activity( + record_decision, + decision, + start_to_close_timeout=timedelta(minutes=1), + ) + if seen >= count: + break + return seen diff --git a/examples/streams/june_scenarios/README.md b/examples/streams/june_scenarios/README.md new file mode 100644 index 000000000..e5e6388af --- /dev/null +++ b/examples/streams/june_scenarios/README.md @@ -0,0 +1,53 @@ +# Roey's June scenarios, on the shipped surface + +Every scenario in Roey's Notion page "Streaming Design Discussion Prep Notes" +(June 2, under "Streaming Links") mapped onto `temporalio.streams` as it +ships on this branch. Each file opens with his scenario heading, a status, +and one sentence why. Where his sketch uses a call shape we do not have, the +docstring shows his shape in one line and the code uses ours. Nothing here +reaches into private SDK code or adds a feature. + +| Roey's scenario | File | Status | Note | +|---|---|---|---| +| Client starts and consume stream: primary and named | `s1_client_consumes.py` | implemented | Default topic with no name, a typed topic, `last=N`, `after=END`, `BEGINNING` on a moved floor (memory only) | +| Client starts and consume stream: standalone alt 1, 2, 3 | `s2_standalone_streams.py` | alts 1, 2 and 3 implemented | Alt 1 as a read that parks until `create_stream` and an append land (native; memory and Redis answer `StreamNotFoundError`); alt 2 as `client.create_stream`, a policy floor, `close()` and `StreamClosedError`; alt 3 as a stream created first and passed into the workflow start as a `StreamRef`, opened in the activity with `activity.stream_handle(ref)`. A start that commits the stream with the workflow remains the design question. Workflow Streams declines | +| Workflow as Producer: as named handle | `s3_workflow_producer.py` | implemented | His turn loop with continue-as-new; the client follows the chain live | +| Workflow as Producer: as return type | `s4_workflow_as_generator.py` | emulated | Default-topic publishes plus `FINISH`, result from the workflow; the generator signature is sugar not built | +| Activity as Producer: as named handle | `s5_activity_producers.py` | implemented | Workflow topic (Path B), `scope="activity"`, standalone activity; all three on native, memory and Redis, the last two skipped on Workflow Streams | +| Activity as Producer: as return type | `s6_activity_as_generator.py` | emulated | Appends plus a heartbeat checkpoint; the retry resumes and readers see `SUPERSEDED` | +| Workflow as Consumer | `s7_workflow_consumer.py` | implemented; foreign stream unsupported | Own inbound topic across continue-as-new, handing over per batch with one producer per batch and carrying a checkpoint; runs on native, Workflow Streams and Redis, memory skips by design; reading a foreign stream from a workflow is rule 5 | +| Client as Consumer over Standalone Nexus | `s8_nexus_consumers.py` (a) | implemented | Activity reads through the `NexusStreams` front and resumes from a heartbeat cursor | +| Nexus operation handler | `s8_nexus_consumers.py` (b) | implemented | The operation returns `temporalio.streams.StreamRef`, taken from the producing workflow's handle, and the client opens it with `get_stream_handle(ref)` on a client whose provider is the front; a stream type of its own in the operation IDL is the nexgen follow-on | +| Workflow as Consumer over Nexus | `s8_nexus_consumers.py` docstring | unsupported by design | A workflow's reads ride its Workflow Task and never cross Nexus | + +## Running + +Each file runs on its own and takes the provider's name, the same way the +examples one directory up do. `run.py` runs them all in order: + +```sh +python -m examples.streams.june_scenarios.run native --address 127.0.0.1:7433 --http http://127.0.0.1:7343 +python -m examples.streams.june_scenarios.run workflow_streams --address 127.0.0.1:7433 --http http://127.0.0.1:7343 +python -m examples.streams.june_scenarios.run memory --address 127.0.0.1:7433 --http http://127.0.0.1:7343 +python -m examples.streams.june_scenarios.run redis --address 127.0.0.1:7433 --http http://127.0.0.1:7343 --redis redis://127.0.0.1:6379 +python -m examples.streams.june_scenarios.s2_standalone_streams native --address 127.0.0.1:7433 +``` + +`native` needs a server built from the stream-carrying branch, and `s2` +alt 1's read that parks until the stream is created needs one built from +its current head. `s5` (b) and (c) need a server with standalone activities +and activity-owned streams, and `s8` needs the server's Nexus HTTP ingress +(`--http`, default `http://127.0.0.1:7243`, `7343` on the server above); +the stream-carrying server has all of them, so the commands above point +every provider at it. `s8` creates and deletes its own Nexus endpoint. +`memory` is offered here, not in the parent examples, because it is not +replay-safe; these scenarios keep a warm cache. `memory`, `native` and +`redis` hold standalone streams, so `s2` runs on all three; `redis` runs +with `--redis` naming a local Redis. Every scenario ran green on all four +providers against that server, apart from the refusals below. + +A scenario a provider cannot serve says so in its output and moves on: +`s2` on `workflow_streams`, which keeps a stream inside a workflow's log, +`s5` (b) and (c) on `workflow_streams`, `s1` (d) on anything but `memory`, +and `s7` on `memory`, which keeps one topic across a chain rather than one +per run. diff --git a/examples/streams/june_scenarios/__init__.py b/examples/streams/june_scenarios/__init__.py new file mode 100644 index 000000000..646be7dae --- /dev/null +++ b/examples/streams/june_scenarios/__init__.py @@ -0,0 +1 @@ +"""Roey's June streaming scenarios, each mapped onto the shipped streams surface.""" diff --git a/examples/streams/june_scenarios/_common.py b/examples/streams/june_scenarios/_common.py new file mode 100644 index 000000000..837097444 --- /dev/null +++ b/examples/streams/june_scenarios/_common.py @@ -0,0 +1,38 @@ +"""What every scenario shares: its flags, its ids and one way to print. + +The store is still named in one place, ``examples.streams._setup``. The +scenarios add the memory provider to the choices because two of them need an +activity-owned stream or truncation, and memory is the one in-process store +that has both. +""" + +from __future__ import annotations + +import argparse +import uuid + +from examples.streams import _setup + +PROVIDERS = (*_setup.PROVIDERS, "memory") + + +def parser(description: str) -> argparse.ArgumentParser: + """The example flags, plus the memory provider and the Nexus ingress.""" + parser = _setup.parser(description, PROVIDERS) + parser.add_argument( + "--http", + default="http://127.0.0.1:7243", + help="the server's Nexus HTTP ingress, for s8", + ) + return parser + + +def ids(prefix: str) -> tuple[str, str]: + """A fresh workflow id and its own task queue, so reruns never collide.""" + workflow_id = f"{prefix}-{uuid.uuid4().hex[:8]}" + return workflow_id, f"tq-{workflow_id}" + + +def banner(title: str, provider: str) -> None: + """Head each scenario's output, so a run of all of them reads in sections.""" + print(f"\n== {title} [{provider}]") diff --git a/examples/streams/june_scenarios/run.py b/examples/streams/june_scenarios/run.py new file mode 100644 index 000000000..e6d08dec5 --- /dev/null +++ b/examples/streams/june_scenarios/run.py @@ -0,0 +1,58 @@ +r"""Run every June scenario, in order, on one provider. + + python -m examples.streams.june_scenarios.run native --address 127.0.0.1:7333 + python -m examples.streams.june_scenarios.run workflow_streams --address 127.0.0.1:7333 + python -m examples.streams.june_scenarios.run memory --address 127.0.0.1:7333 + python -m examples.streams.june_scenarios.run native --address 127.0.0.1:7333 --only s2,s5 + +Each scenario connects, runs and closes on its own, exactly as it does when +run as its own module, so one that a provider cannot serve prints why and +the next one starts clean. ``s8`` needs the server's Nexus HTTP ingress, +which ``--http`` names. +""" + +from __future__ import annotations + +import asyncio + +from examples.streams.june_scenarios import ( + _common, + s1_client_consumes, + s2_standalone_streams, + s3_workflow_producer, + s4_workflow_as_generator, + s5_activity_producers, + s6_activity_as_generator, + s7_workflow_consumer, + s8_nexus_consumers, +) + +SCENARIOS = { + "s1": s1_client_consumes.run, + "s2": s2_standalone_streams.run, + "s3": s3_workflow_producer.run, + "s4": s4_workflow_as_generator.run, + "s5": s5_activity_producers.run, + "s6": s6_activity_as_generator.run, + "s7": s7_workflow_consumer.run, + "s8": s8_nexus_consumers.run, +} + + +async def main() -> None: + """Run the scenarios named by ``--only``, or all of them.""" + parser = _common.parser(__doc__ or "") + parser.add_argument( + "--only", default="", help="comma-separated scenarios, for example s1,s5" + ) + args = parser.parse_args() + chosen = [name for name in args.only.split(",") if name] or list(SCENARIOS) + unknown = [name for name in chosen if name not in SCENARIOS] + if unknown: + parser.error(f"unknown scenarios {unknown}; choose from {list(SCENARIOS)}") + for name in chosen: + await SCENARIOS[name](args) + + +if __name__ == "__main__": + asyncio.run(main()) diff --git a/examples/streams/june_scenarios/s1_client_consumes.py b/examples/streams/june_scenarios/s1_client_consumes.py new file mode 100644 index 000000000..f99c7a7a1 --- /dev/null +++ b/examples/streams/june_scenarios/s1_client_consumes.py @@ -0,0 +1,166 @@ +"""Scenario "Client starts and consume stream": primary and named execution stream. + +Status: implemented. + +Why: a workflow's default topic is the primary stream he reads without a +name, and a typed topic is his named stream, so both sketches are one call. + + python -m examples.streams.june_scenarios.s1_client_consumes workflow_streams + python -m examples.streams.june_scenarios.s1_client_consumes native --address 127.0.0.1:7333 + python -m examples.streams.june_scenarios.s1_client_consumes memory --address 127.0.0.1:7333 + +His shape is ``handle.stream(ScoreUpdate)`` and +``handle.stream(ScoreUpdate, name="scores")``. Ours keeps the address on the +handle and the name and type on the topic: +``client.get_stream_handle(wid).read(result_type=ScoreUpdate)`` for the +default topic and ``.read(topic=SCORES)`` for the named one. + +The game plays two halves and waits for a signal between them, so the reads +below land at known places. The client reads the first half from the start +of the default topic, then joins the named topic at ``END`` and sees only +the second half, then asks for the ``last=3`` records, which count the +``FINISH`` record too. Last, it reads a truncated topic from ``BEGINNING``, +which starts at the oldest record left rather than at offset zero. Only the +memory provider can truncate from here; an owned topic on the native server +is budgeted and never truncated, and ``s2`` shows a moved floor on a +standalone stream there. +""" + +from __future__ import annotations + +import argparse +import asyncio +import contextlib +from dataclasses import dataclass +from datetime import timedelta + +from examples.streams import _setup +from examples.streams.june_scenarios import _common +from temporalio import streams, workflow +from temporalio.streams import END, RecordKind, StreamRecord +from temporalio.streams.providers.memory import MemoryStreams +from temporalio.worker import Worker + + +@dataclass +class ScoreUpdate: + """The score at one moment of the game.""" + + home_score: int + away_score: int + clock: str + + +SCORES = streams.topic("scores", ScoreUpdate) + + +@workflow.defn +class Game: + """Publishes every score change on the default topic and on ``scores``.""" + + def __init__(self) -> None: + """Start in the first half.""" + self._second_half = False + + @workflow.signal + def second_half(self) -> None: + """Let the second half start.""" + self._second_half = True + + @workflow.run + async def run(self, per_half: int) -> str: + """Play two halves of ``per_half`` updates each and return the final score.""" + primary = workflow.stream_writer() + named = workflow.stream_writer(SCORES) + home = away = 0 + for half in (1, 2): + if half == 2: + await workflow.wait_condition(lambda: self._second_half) + for minute in range(per_half): + if minute % 2 == 0: + home += 1 + else: + away += 1 + update = ScoreUpdate(home, away, clock=f"H{half} {minute:02}'") + primary.publish(update) + named.publish(update) + await workflow.sleep(timedelta(milliseconds=200)) + primary.finish() + named.finish() + return f"{home}-{away}" + + +def show(record: StreamRecord[ScoreUpdate]) -> None: + """One line per record.""" + print(f" {record.kind.name:6} {record.value}") + + +async def run(args: argparse.Namespace) -> None: + """Play the game and read it four ways.""" + _common.banner("s1 client consumes", args.provider) + client, provider = await _setup.connect(args) + workflow_id, task_queue = _common.ids("june-s1") + per_half = 3 + try: + async with Worker(client, task_queue=task_queue, workflows=[Game]): + handle = await client.start_workflow( + Game.run, per_half, id=workflow_id, task_queue=task_queue + ) + stream = client.get_stream_handle(workflow_id) + + print(" (a) default topic, no name, from the start: the first half") + seen = 0 + async with contextlib.aclosing( + stream.read(result_type=ScoreUpdate) + ) as primary: + async for record in primary: + show(record) + seen += record.kind is RecordKind.DATA + if seen == per_half: + break + + print(" (b) named topic from END: only what lands after joining") + + async def follow_from_end() -> None: + joined = stream.read(topic=SCORES, after=END) + async with contextlib.aclosing(joined) as records: + async for record in records: + show(record) + if record.kind is RecordKind.FINISH: + break + + joined = asyncio.create_task(follow_from_end()) + # END resolves on the reader's first poll, not at this call, so the + # half-time whistle waits until that poll has had time to land. + await asyncio.sleep(1.0) + await handle.signal(Game.second_half) + await joined + print(f" final score {await handle.result()}") + + print(" (c) named topic, last=3: FINISH counts as one of the three") + async for record in stream.read(topic=SCORES, last=3): + show(record) + + print(" (d) BEGINNING after the floor moved") + if isinstance(provider, MemoryStreams): + # Truncation is a store's retention, not part of the provider + # contract, so only the in-process store offers it by hand. + provider.truncate(workflow_id, SCORES.name, keep=2) + async for record in stream.read(topic=SCORES): + show(record) + else: + print( + f" {args.provider} cannot truncate a workflow's topic from " + "outside; see s2 for BEGINNING on a moved floor" + ) + finally: + await provider.close() + + +async def main() -> None: + """Parse the flags and run the scenario.""" + await run(_common.parser(__doc__ or "").parse_args()) + + +if __name__ == "__main__": + asyncio.run(main()) diff --git a/examples/streams/june_scenarios/s2_standalone_streams.py b/examples/streams/june_scenarios/s2_standalone_streams.py new file mode 100644 index 000000000..80eb51b4c --- /dev/null +++ b/examples/streams/june_scenarios/s2_standalone_streams.py @@ -0,0 +1,201 @@ +"""Scenario "Client starts and consume stream", the three standalone alternatives. + +Status: implemented for alts 1, 2 and 3; alt 3 in the ref-argument shape. + +Why: a standalone stream has an id of its own and no owner, the client +creates and seals it on purpose, and a ``StreamRef`` carries it as data into +a workflow start, an activity or a Nexus operation. + + python -m examples.streams.june_scenarios.s2_standalone_streams native --address 127.0.0.1:7433 + python -m examples.streams.june_scenarios.s2_standalone_streams memory --address 127.0.0.1:7433 + python -m examples.streams.june_scenarios.s2_standalone_streams redis --address 127.0.0.1:7433 + +His alt 1 is a read that blocks until the stream exists. On ``native`` a +read on an id nobody has created parks on the server and delivers the first +record once a ``create_stream`` and an append land; ``memory`` and ``redis`` +answer such a read with ``StreamNotFoundError`` instead, so the code shows +whichever the provider does. + +His alt 2 is ``stream = await client.create_stream(ProgressUpdate, +stream_id=...)``. Ours is ``await client.create_stream(stream_id, +max_records=3)``: the policy is the stream's, the type is the topic's. An +outside producer appends, a second handle opened by id reads from +``BEGINNING``, which the policy has moved past the oldest record, and +``close()`` seals the stream: a late append raises ``StreamClosedError`` and +a read opened afterwards ends by itself once the retained tail is delivered. + +His alt 3 is a client-side stream made durable by the start call. Ours +creates the stream first and passes ``handle.ref()`` as the workflow's +argument; the workflow hands the ref to its activity, which opens it with +``activity.stream_handle(ref)`` and appends, and the client follows the same +ref. A start that commits the stream with the workflow, so a crash between +the two calls cannot leave an orphan, remains the design question. + +``workflow_streams`` keeps a stream inside a workflow's own log, so it has +nowhere to put a stream with no owner and declines. +""" + +from __future__ import annotations + +import argparse +import asyncio +import contextlib +from dataclasses import dataclass +from datetime import timedelta +from typing import Any + +from examples.streams import _setup +from examples.streams.june_scenarios import _common +from temporalio import activity, streams, workflow +from temporalio.streams import ( + RecordKind, + StreamClosedError, + StreamHandle, + StreamNotFoundError, + StreamRecord, + StreamRef, +) +from temporalio.worker import Worker + + +@dataclass +class ProgressUpdate: + """One line of progress on the session.""" + + message: str + + +PROGRESS = streams.topic("progress", ProgressUpdate) + + +def show(record: StreamRecord[Any]) -> None: + """One line per record.""" + print(f" {record.kind.name:6} {record.producer_id:8} {record.value}") + + +async def first(reader: StreamHandle) -> StreamRecord[Any]: + """The first record of a read on ``PROGRESS``, closing the read after it. + + The read is opened in here, so a provider that refuses a missing stream at + the call refuses it inside the task that waits for the record. + """ + async with contextlib.aclosing(reader.read(topic=PROGRESS)) as reading: + async for record in reading: + return record + raise RuntimeError("the read ended before a record arrived") + + +@activity.defn +async def report_progress(session: StreamRef) -> int: + """Append three steps to the stream the ref names. + + Inside an activity the producer's identity is the activity's own, so a + retry of this activity lands each step once. + """ + producer = activity.stream_handle(session).producer() + for step in range(1, 4): + await producer.append(ProgressUpdate(f"step {step} done")) + await asyncio.sleep(0.2) + await producer.finish() + return 3 + + +@workflow.defn +class Session: + """Receives the stream as a ref and hands it to its activity.""" + + @workflow.run + async def run(self, session: StreamRef) -> int: + """Run the activity that writes to the stream; the workflow never touches it.""" + return await workflow.execute_activity( + report_progress, session, start_to_close_timeout=timedelta(minutes=1) + ) + + +async def run(args: argparse.Namespace) -> None: + """Wait for a stream, run one through its life, then pass one into a workflow.""" + _common.banner("s2 standalone streams", args.provider) + if args.provider == "workflow_streams": + print( + " workflow_streams keeps a stream inside a workflow's log, so it has " + "nowhere to put a stream with no owner; skipped" + ) + return + client, provider = await _setup.connect(args) + workflow_id, task_queue = _common.ids("june-s2") + stream_id = f"session-{workflow_id.rsplit('-', 1)[1]}" + try: + print(" alt 1: a reader asks for a stream nobody has created yet") + reader = client.get_stream_handle(stream_id=stream_id) + waiting = asyncio.create_task(first(reader)) + await asyncio.sleep(0.5) + parked: asyncio.Task[StreamRecord[Any]] | None = waiting + if waiting.done(): + # memory and redis do not wait for a stream to be created. + try: + waiting.result() + except StreamNotFoundError as error: + print(f" StreamNotFoundError: {error}") + parked = None + else: + print(" the read is parked on the server until the stream exists") + created = await client.create_stream(stream_id, max_records=3) + await created.producer( + topic=PROGRESS, producer_id="operator", attempt=1 + ).append(ProgressUpdate("created")) + if parked is not None: + record = await parked + print( + " and it delivered the first record once the create and an append landed:" + ) + show(record) + + print(" alt 2: produce from outside, read by id, seal it") + operator = created.producer(topic=PROGRESS, producer_id="operator", attempt=2) + for note in ("watching", "looks good", "wrapping up"): + await operator.append(ProgressUpdate(note)) + await operator.finish() + await created.close() + print(" max_records=3 moved the floor: BEGINNING is the oldest record kept") + by_id = client.get_stream_handle(stream_id=stream_id) + async for record in by_id.read(topic=PROGRESS): + show(record) + try: + await created.producer( + topic=PROGRESS, producer_id="late", attempt=1 + ).append(ProgressUpdate("too late")) + except StreamClosedError as error: + print(f" a late append is refused: StreamClosedError: {error}") + + print(" alt 3: create, then start the workflow with the stream's ref") + session = await client.create_stream(f"{stream_id}-run") + ref = session.ref(topic=PROGRESS) + print(f" passing {ref}") + async with Worker( + client, + task_queue=task_queue, + workflows=[Session], + activities=[report_progress], + ): + handle = await client.start_workflow( + Session.run, ref, id=workflow_id, task_queue=task_queue + ) + # The ref names the topic, not its type; the typed topic says + # what to decode as. + async for record in client.get_stream_handle(ref).read(topic=PROGRESS): + show(record) + if record.kind is RecordKind.FINISH: + break + print(f" the workflow returned {await handle.result()}") + await session.close() + finally: + await provider.close() + + +async def main() -> None: + """Parse the flags and run the scenario.""" + await run(_common.parser(__doc__ or "").parse_args()) + + +if __name__ == "__main__": + asyncio.run(main()) diff --git a/examples/streams/june_scenarios/s3_workflow_producer.py b/examples/streams/june_scenarios/s3_workflow_producer.py new file mode 100644 index 000000000..9bf60e29e --- /dev/null +++ b/examples/streams/june_scenarios/s3_workflow_producer.py @@ -0,0 +1,157 @@ +"""Scenario "Workflow as Producer", as a named handle. + +Status: implemented. + +Why: a workflow publishes to a typed topic of its own stream, and a client +handle without a run id follows the execution chain across continue-as-new. + + python -m examples.streams.june_scenarios.s3_workflow_producer workflow_streams + python -m examples.streams.june_scenarios.s3_workflow_producer native --address 127.0.0.1:7333 + +His shape is ``workflow.StreamHandle[ProgressUpdate](name="progress")`` +and ``.send(...)``. Ours is ``workflow.stream_writer(PROGRESS)`` and +``.publish(...)``, which commits with the Workflow Task. His standalone +variant, a workflow writing to a stream it does not own, is not offered: +a workflow publishes only to its own topics, and an activity writes to a +standalone stream instead, as in ``s2``. + +The loop is his, turn by turn: publish, call the model, publish, call a +tool, publish. Continue-as-new is taken when the server suggests it, and +also every two turns here so a short demo crosses runs. Each run's topic +belongs to that run, and the client's read walks from one run to the next +without being told. +""" + +from __future__ import annotations + +import argparse +import asyncio +from dataclasses import dataclass +from datetime import timedelta + +from examples.streams import _setup +from examples.streams.june_scenarios import _common +from temporalio import activity, streams, workflow +from temporalio.worker import Worker + + +@dataclass +class ProgressUpdate: + """What the agent is doing now, and which run of the chain said it.""" + + message: str + run: int + + +@dataclass +class LlmResult: + """The model's answer: done, or a tool to call.""" + + done: bool + tool: str = "" + + +@dataclass +class AgentInput: + """The prompt, and where a continued run picks up.""" + + prompt: str + turn: int = 0 + run: int = 1 + + +PROGRESS = streams.topic("progress", ProgressUpdate) +TURNS_PER_RUN = 2 +TOTAL_TURNS = 5 + + +@activity.defn +async def call_llm(input: AgentInput) -> LlmResult: + """A stand-in model: asks for a tool until the last turn.""" + if input.turn + 1 >= TOTAL_TURNS: + return LlmResult(done=True) + return LlmResult(done=False, tool=f"search #{input.turn + 1}") + + +@activity.defn +async def call_tool(result: LlmResult) -> str: + """A stand-in tool: its answer becomes the next prompt.""" + return f"results of {result.tool}" + + +@workflow.defn +class AgentWorkflow: + """His turn loop, publishing progress on the ``progress`` topic.""" + + @workflow.run + async def run(self, input: AgentInput) -> str: + """Run turns until the model is done, continuing as new along the way.""" + progress = workflow.stream_writer(PROGRESS) + prompt, turn, run = input.prompt, input.turn, input.run + timeout = timedelta(seconds=30) + + def send(message: str) -> None: + progress.publish(ProgressUpdate(message, run)) + + while True: + if ( + workflow.info().is_continue_as_new_suggested() + or turn - input.turn >= TURNS_PER_RUN + ): + workflow.continue_as_new(AgentInput(prompt, turn, run + 1)) + send(f"turn {turn} started") + result = await workflow.execute_activity( + call_llm, AgentInput(prompt, turn), start_to_close_timeout=timeout + ) + send(f"turn {turn} ended") + if result.done: + break + send(f"tool call {result.tool} started") + prompt = await workflow.execute_activity( + call_tool, result, start_to_close_timeout=timeout + ) + send(f"tool call {result.tool} ended") + turn += 1 + progress.finish() + return f"done after {turn + 1} turns" + + +async def run(args: argparse.Namespace) -> None: + """Start the agent and follow its progress across every run of the chain.""" + _common.banner("s3 workflow producer", args.provider) + client, provider = await _setup.connect(args) + workflow_id, task_queue = _common.ids("june-s3") + try: + async with Worker( + client, + task_queue=task_queue, + workflows=[AgentWorkflow], + activities=[call_llm, call_tool], + ): + handle = await client.start_workflow( + AgentWorkflow.run, + AgentInput("plan a trip"), + id=workflow_id, + task_queue=task_queue, + ) + # No run id on the handle, so the read follows the chain and ends + # when its last run is closed and the tail has been delivered. + async for record in client.get_stream_handle(workflow_id).read( + topic=PROGRESS + ): + if record.value is None: + print(f" {record.kind.name}") + continue + print(f" run {record.value.run} {record.value.message}") + print(f" {await handle.result()}") + finally: + await provider.close() + + +async def main() -> None: + """Parse the flags and run the scenario.""" + await run(_common.parser(__doc__ or "").parse_args()) + + +if __name__ == "__main__": + asyncio.run(main()) diff --git a/examples/streams/june_scenarios/s4_workflow_as_generator.py b/examples/streams/june_scenarios/s4_workflow_as_generator.py new file mode 100644 index 000000000..0c067c2e6 --- /dev/null +++ b/examples/streams/june_scenarios/s4_workflow_as_generator.py @@ -0,0 +1,156 @@ +r"""Scenario "Workflow as Producer", as return type (async generator). + +Status: emulated. + +Why: the async-generator run signature is sugar that is not built, and +publishing each value to the default topic, then ``FINISH``, then returning +the result is what that sugar would lower to. + + python -m examples.streams.june_scenarios.s4_workflow_as_generator workflow_streams + python -m examples.streams.june_scenarios.s4_workflow_as_generator \ + native --address 127.0.0.1:7333 + +His shape is ``async def run(...) -> workflow.Stream[ScoreUpdate, +GameFinalResult]`` with ``yield`` per update and ``return`` for the result. +Ours replaces each ``yield`` with ``workflow.stream_writer().publish(...)`` +on the default topic, so a reader needs no topic name, and the ``return`` +stays an ordinary workflow result. The client reads the updates until +``FINISH`` and then fetches the result from the workflow handle, which is +the two halves his generator type would have handed it at once. +""" + +from __future__ import annotations + +import argparse +import asyncio +import contextlib +from dataclasses import dataclass +from datetime import timedelta + +from examples.streams import _setup +from examples.streams.june_scenarios import _common +from temporalio import activity, workflow +from temporalio.streams import RecordKind +from temporalio.worker import Worker + + +@dataclass +class ScoreUpdate: + """What each ``yield`` would have produced.""" + + home_score: int + away_score: int + clock: str + + +@dataclass +class GameFinalResult: + """What the ``return`` produces.""" + + winner: str + final_score: str + + +@dataclass +class GameState: + """One poll of the game feed.""" + + home_score: int + away_score: int + clock: str + is_finished: bool + + +@dataclass +class GamePoll: + """Which game to poll, and which poll this is.""" + + game_id: str + tick: int + + +@activity.defn +async def get_game_state(poll: GamePoll) -> GameState: + """A stand-in feed that reports the game over on its fourth poll.""" + tick = poll.tick + return GameState( + home_score=(tick + 1) // 2 + tick // 3, + away_score=tick // 2, + clock=f"{tick * 20}'", + is_finished=tick >= 3, + ) + + +@workflow.defn +class LiveScoresWorkflow: + """His generator, with ``yield`` spelled as a publish to the default topic.""" + + @workflow.run + async def run(self, game_id: str) -> GameFinalResult: + """Publish each score, then finish the topic and return the result.""" + scores = workflow.stream_writer() + tick = 0 + while True: + game = await workflow.execute_activity( + get_game_state, + GamePoll(game_id, tick), + start_to_close_timeout=timedelta(seconds=30), + ) + # yield ScoreUpdate(...) + scores.publish(ScoreUpdate(game.home_score, game.away_score, game.clock)) + if game.is_finished: + break + tick += 1 + await workflow.sleep(timedelta(milliseconds=200)) + # The end of the generator, so a reader knows no more values follow. + scores.finish() + home, away = game.home_score, game.away_score + return GameFinalResult( + winner="home" if home > away else "away" if away > home else "draw", + final_score=f"{home}-{away}", + ) + + +async def run(args: argparse.Namespace) -> None: + """Read the yielded values, then the returned one.""" + _common.banner("s4 workflow as generator", args.provider) + client, provider = await _setup.connect(args) + workflow_id, task_queue = _common.ids("june-s4") + try: + async with Worker( + client, + task_queue=task_queue, + workflows=[LiveScoresWorkflow], + activities=[get_game_state], + ): + handle = await client.start_workflow( + LiveScoresWorkflow.run, "game-7", id=workflow_id, task_queue=task_queue + ) + # FINISH is the end of the generator, so the read stops there + # rather than waiting for the run to close. + records = client.get_stream_handle(workflow_id).read( + result_type=ScoreUpdate + ) + async with contextlib.aclosing(records) as updates: + async for record in updates: + if record.kind is RecordKind.FINISH: + print(" FINISH: the generator is exhausted") + break + assert record.value is not None + update = record.value + print( + f" Score: {update.home_score}-{update.away_score} " + f"at {update.clock}" + ) + print(f" returned {await handle.result()}") + finally: + await provider.close() + + +async def main() -> None: + """Parse the flags and run the scenario.""" + await run(_common.parser(__doc__ or "").parse_args()) + + +if __name__ == "__main__": + asyncio.run(main()) diff --git a/examples/streams/june_scenarios/s5_activity_producers.py b/examples/streams/june_scenarios/s5_activity_producers.py new file mode 100644 index 000000000..ba5c10983 --- /dev/null +++ b/examples/streams/june_scenarios/s5_activity_producers.py @@ -0,0 +1,194 @@ +"""Scenario "Activity as Producer", as named handle, separate from return type. + +Status: implemented. + +Why: ``activity.stream_handle()`` reaches the workflow's topics, the +activity's own streams with ``scope="activity"``, or a standalone +activity's own streams, by a static rule and with the same producer verbs. + + python -m examples.streams.june_scenarios.s5_activity_producers native --address 127.0.0.1:7333 + python -m examples.streams.june_scenarios.s5_activity_producers memory --address 127.0.0.1:7333 + python -m examples.streams.june_scenarios.s5_activity_producers workflow_streams + +His three shapes and ours, one line each: + +- ``activity.StreamHandle[ProgressUpdate](workflow_id=..., name="progress")`` + is ``activity.stream_handle().producer(topic=PROGRESS)``: the workflow + that scheduled the activity, pinned to its run, which is Path B. +- ``activity.StreamHandle[ProgressUpdate](name="progress")`` in an activity + that owns its streams is ``activity.stream_handle(scope="activity")`` + inside a workflow, and plain ``activity.stream_handle()`` in a standalone + activity. The rule is static because a stream is created by its first + write, so a rule that probed for one would send a retry somewhere else. +- ``client.StreamHandle[...](stream_id=...)``, a standalone stream, is the + stream client in ``s2``. + +A topic on an activity's own streams and the same topic on its workflow are +two streams, which part (b) shows by writing ``progress`` to both. Parts (b) +and (c) need a store that can hold an activity-owned stream: native and +memory can, and on any other provider the client's first handle raises +``StreamUnsupportedError``, which is printed in their place. +""" + +from __future__ import annotations + +import argparse +import asyncio +import contextlib +from dataclasses import dataclass +from datetime import timedelta + +from examples.streams import _setup +from examples.streams.june_scenarios import _common +from temporalio import activity, streams, workflow +from temporalio.client import Client +from temporalio.streams import RecordKind, StreamHandle, StreamUnsupportedError +from temporalio.worker import Worker + + +@dataclass +class ProgressUpdate: + """One line of progress from a tool call.""" + + message: str + + +@dataclass +class ToolCallInput: + """The call, and which stream its progress goes to. + + ``target`` is ``"workflow"`` for the scheduling workflow's topic, + ``"own"`` for the activity's own streams inside a workflow, and + ``"standalone"`` for a standalone activity's own default topic. + """ + + tool: str + target: str + + +PROGRESS = streams.topic("progress", ProgressUpdate) + + +@activity.defn +async def tool_call_activity(input: ToolCallInput) -> str: + """His tool call: progress before and after, and the answer as the result.""" + if input.target == "workflow": + progress = activity.stream_handle().producer(topic=PROGRESS) + elif input.target == "own": + progress = activity.stream_handle(scope="activity").producer(topic=PROGRESS) + else: + # A standalone activity owns its streams without a scope, and naming + # no topic writes its default one. + progress = activity.stream_handle().producer() + await progress.append(ProgressUpdate(f"{input.tool} started")) + await asyncio.sleep(0.3) + await progress.append(ProgressUpdate(f"{input.tool} ended")) + await progress.finish() + return f"{input.tool} answered" + + +@workflow.defn +class ToolWorkflow: + """Schedules the tool call once per stream target it is asked for.""" + + @workflow.run + async def run(self, own_streams: bool) -> list[str]: + """Run (a), and (b) when the store can hold an activity's own streams.""" + timeout = timedelta(seconds=30) + answers = [ + await workflow.execute_activity( + tool_call_activity, + ToolCallInput("search", "workflow"), + activity_id="tool-a", + start_to_close_timeout=timeout, + ) + ] + if own_streams: + answers.append( + await workflow.execute_activity( + tool_call_activity, + ToolCallInput("fetch", "own"), + activity_id="tool-b", + start_to_close_timeout=timeout, + ) + ) + return answers + + +async def show(stream: StreamHandle, topic: streams.StreamTopic | None) -> None: + """Print one stream until its producer finishes.""" + records = ( + stream.read(topic=topic) + if topic is not None + else stream.read(result_type=ProgressUpdate) + ) + async with contextlib.aclosing(records) as reading: + async for record in reading: + value = record.value.message if record.value else "" + print(f" {record.kind.name:6} {record.producer_id:6} {value}") + if record.kind is RecordKind.FINISH: + break + + +def supports_activity_owners(client: Client) -> bool: + """Ask the provider for an activity handle; one that cannot hold it says so.""" + try: + client.get_stream_handle(activity_id="probe") + except StreamUnsupportedError as error: + print(f" StreamUnsupportedError: {error}") + return False + return True + + +async def run(args: argparse.Namespace) -> None: + """Produce from an activity three ways and read each stream from outside.""" + _common.banner("s5 activity producers", args.provider) + client, provider = await _setup.connect(args) + workflow_id, task_queue = _common.ids("june-s5") + try: + async with Worker( + client, + task_queue=task_queue, + workflows=[ToolWorkflow], + activities=[tool_call_activity], + ): + print(" can this provider hold a stream an activity owns?") + own = supports_activity_owners(client) + handle = await client.start_workflow( + ToolWorkflow.run, own, id=workflow_id, task_queue=task_queue + ) + print(f" workflow returned {await handle.result()}") + + print(" (a) activity in a workflow, onto the workflow's progress topic") + await show(client.get_stream_handle(workflow_id), PROGRESS) + if not own: + print(" (b) and (c) skipped: they need an activity-owned stream") + return + + print(" (b) same activity with scope='activity', onto its own progress") + await show( + client.get_stream_handle(workflow_id, activity_id="tool-b"), PROGRESS + ) + + print(" (c) standalone activity, onto its own default topic") + activity_id = workflow_id.replace("june-s5", "saa") + started = await client.start_activity( + tool_call_activity, + ToolCallInput("summarize", "standalone"), + id=activity_id, + task_queue=task_queue, + start_to_close_timeout=timedelta(seconds=30), + ) + await show(client.get_stream_handle(activity_id=activity_id), None) + print(f" result {await started.result()}") + finally: + await provider.close() + + +async def main() -> None: + """Parse the flags and run the scenario.""" + await run(_common.parser(__doc__ or "").parse_args()) + + +if __name__ == "__main__": + asyncio.run(main()) diff --git a/examples/streams/june_scenarios/s6_activity_as_generator.py b/examples/streams/june_scenarios/s6_activity_as_generator.py new file mode 100644 index 000000000..170a86915 --- /dev/null +++ b/examples/streams/june_scenarios/s6_activity_as_generator.py @@ -0,0 +1,153 @@ +r"""Scenario "Activity as Producer", as return type (async generator). + +Status: emulated. + +Why: the yield-based activity signature is not built; each ``yield`` becomes +an append on the workflow's default topic, and his "load last heartbeat and +resume from snapshot" is a heartbeat checkpoint the retry reads back. + + python -m examples.streams.june_scenarios.s6_activity_as_generator workflow_streams + python -m examples.streams.june_scenarios.s6_activity_as_generator \ + native --address 127.0.0.1:7333 + +His shape is ``yield ProgressUpdate(...)`` inside the activity. Ours is +``await progress.append(ProgressUpdate(...))`` on +``activity.stream_handle().producer()``, followed by +``activity.heartbeat(next_step)`` once the step is safely written. + +The first attempt dies after writing step 3 and before checkpointing it. The +retry reads the checkpoint, so it starts again at step 3 rather than at the +top, and writes under a new attempt. The reader sees that change as a +``SUPERSEDED`` record and keys what it keeps by step, so step 3 arriving +twice, once per attempt, leaves one copy. The heartbeat is a hint, not the +truth: it can lag the stream, which is why the retry repeats a step rather +than skipping one. +""" + +from __future__ import annotations + +import argparse +import asyncio +import contextlib +from dataclasses import dataclass +from datetime import timedelta + +from examples.streams import _setup +from examples.streams.june_scenarios import _common +from temporalio import activity, workflow +from temporalio.common import RetryPolicy +from temporalio.streams import RecordKind +from temporalio.worker import Worker + + +@dataclass +class ProgressUpdate: + """One step of the agent loop.""" + + step: int + message: str + + +STEPS = ( + "Turn started", + "Turn ended", + "Tool call started", + "Tool call ended", + "Turn started", + "Turn ended", +) +CRASH_AFTER_STEP = 3 + + +@activity.defn +async def agent_activity(prompt: str) -> str: + """His generator: each ``yield`` is an append, each safe point a heartbeat.""" + info = activity.info() + # load last heartbeat and resume from snapshot + start = int(info.heartbeat_details[0]) if info.heartbeat_details else 0 + progress = activity.stream_handle().producer() + for step in range(start, len(STEPS)): + await progress.append(ProgressUpdate(step, STEPS[step])) + if step == CRASH_AFTER_STEP and info.attempt == 1: + raise RuntimeError("the worker died mid-stream") + # The checkpoint is the next step to write. A throttled heartbeat + # still reaches the server, because the worker sends the last one + # with the failure. + activity.heartbeat(step + 1) + await asyncio.sleep(0.1) + await progress.finish() + return f"answered {prompt!r} in attempt {info.attempt} from step {start}" + + +@workflow.defn +class AgentSession: + """Runs the generator activity and returns what it answered.""" + + @workflow.run + async def run(self, prompt: str) -> str: + """Run the activity with a short retry, so the crash retries at once.""" + return await workflow.execute_activity( + agent_activity, + prompt, + start_to_close_timeout=timedelta(minutes=1), + heartbeat_timeout=timedelta(seconds=10), + retry_policy=RetryPolicy( + initial_interval=timedelta(milliseconds=200), maximum_attempts=3 + ), + ) + + +async def run(args: argparse.Namespace) -> None: + """Read the generator's values across the crash and the resumed retry.""" + _common.banner("s6 activity as generator", args.provider) + client, provider = await _setup.connect(args) + workflow_id, task_queue = _common.ids("june-s6") + try: + async with Worker( + client, + task_queue=task_queue, + workflows=[AgentSession], + activities=[agent_activity], + ): + handle = await client.start_workflow( + AgentSession.run, "plan a trip", id=workflow_id, task_queue=task_queue + ) + kept: dict[int, str] = {} + # A failed attempt writes no FINISH, so the first one seen is the + # attempt that completed, and the read can stop there. + records = client.get_stream_handle(workflow_id).read( + result_type=ProgressUpdate + ) + async with contextlib.aclosing(records) as reading: + async for record in reading: + if record.kind is RecordKind.SUPERSEDED: + assert record.supersession is not None + change = record.supersession + print( + f" SUPERSEDED attempt {change.previous_attempt} " + f"by attempt {change.attempt}" + ) + elif record.kind is RecordKind.FINISH: + print(f" FINISH attempt {record.attempt}") + break + else: + assert record.value is not None + update = record.value + kept[update.step] = update.message + print( + f" DATA attempt {record.attempt} " + f"step {update.step} {update.message}" + ) + print(f" kept one copy of steps {sorted(kept)}") + print(f" {await handle.result()}") + finally: + await provider.close() + + +async def main() -> None: + """Parse the flags and run the scenario.""" + await run(_common.parser(__doc__ or "").parse_args()) + + +if __name__ == "__main__": + asyncio.run(main()) diff --git a/examples/streams/june_scenarios/s7_workflow_consumer.py b/examples/streams/june_scenarios/s7_workflow_consumer.py new file mode 100644 index 000000000..d16c0e409 --- /dev/null +++ b/examples/streams/june_scenarios/s7_workflow_consumer.py @@ -0,0 +1,157 @@ +"""Scenario "Workflow as Consumer". + +Status: implemented for the workflow's own inbound topic; consuming a foreign +standalone stream from workflow code is unsupported by design. + +Why: a workflow reads its own topics as recorded observations that replay +re-supplies, and rule 5 keeps reading somebody else's stream out of this +release. + + python -m examples.streams.june_scenarios.s7_workflow_consumer workflow_streams + python -m examples.streams.june_scenarios.s7_workflow_consumer native --address 127.0.0.1:7333 + python -m examples.streams.june_scenarios.s7_workflow_consumer redis --address 127.0.0.1:7333 --redis redis://127.0.0.1:6379 + +His shape is ``workflow.StreamHandle[ProgressUpdate](stream_id="...", +offset=offset)`` with ``continue_as_new(update.offset)``. Ours is +``workflow.stream_reader(COMMANDS)`` on the workflow's own ``commands`` +topic, fed from outside by a producer. + +Two deliberate differences. First, the stream: his reads a standalone +stream by id, and ours reads only a topic the workflow owns. Reading a +foreign stream from workflow code is not exposed; the server half exists +(external subscriptions), and the memo's foreign-read sketch covers the SDK +half. Second, what crosses continue-as-new: his offset works because his +stream outlives the run. A workflow's topic belongs to one run on the native +and ``workflow_streams`` providers, a successor's starts empty, and a cursor +from the previous run is refused there. So the run hands over at a batch +boundary, marked by the sender's ``FINISH``, and what it carries is its own +checkpoint: what it has applied so far. The sender waits for the successor +before writing the next batch, which is the one piece of coordination the +per-run topic asks for. On ``redis`` the stream spans the chain and a +successor resumes where its predecessor committed, so the same handover +works there; each batch is its own producer, because a producer's identity +must not repeat with different content on one stream. + +The memory provider keys a topic by workflow rather than by run, so a +successor would read its predecessor's batches again; it is skipped there. +""" + +from __future__ import annotations + +import argparse +import asyncio +from dataclasses import dataclass, field + +from examples.streams import _setup +from examples.streams.june_scenarios import _common +from temporalio import streams, workflow +from temporalio.client import Client +from temporalio.streams import RecordKind +from temporalio.worker import Worker + + +@dataclass +class Command: + """One instruction from outside.""" + + op: str + + +@dataclass +class Checkpoint: + """What a run carries into its successor instead of a cursor.""" + + run: int = 1 + applied: list[str] = field(default_factory=list) + + +COMMANDS = streams.topic("commands", Command) +STOP = "stop" + + +@workflow.defn +class ConsumerWorkflow: + """Reads its commands topic batch by batch, one batch per run.""" + + @workflow.run + async def run(self, checkpoint: Checkpoint) -> list[str]: + """Apply this run's batch, then continue as new or stop.""" + stop = False + async for record in workflow.stream_reader(COMMANDS): + if record.kind is RecordKind.FINISH: + break + if record.kind is not RecordKind.DATA: + continue + assert record.value is not None + workflow.logger.info("got command %s", record.value.op) + if record.value.op == STOP: + stop = True + continue + checkpoint.applied.append(f"run {checkpoint.run}: {record.value.op}") + if stop: + return checkpoint.applied + # Taken at every batch boundary here, so the demo crosses runs; a + # real consumer would also wait for is_continue_as_new_suggested(). + workflow.continue_as_new( + Checkpoint(run=checkpoint.run + 1, applied=checkpoint.applied) + ) + + +async def current_run(client: Client, workflow_id: str, previous: str | None) -> str: + """Wait until the chain's newest run is a new one and still open.""" + while True: + description = await client.get_workflow_handle(workflow_id).describe() + assert description.run_id is not None + if description.run_id != previous and description.close_time is None: + return description.run_id + await asyncio.sleep(0.2) + + +async def run(args: argparse.Namespace) -> None: + """Send three batches, one per run, and print what the chain applied.""" + _common.banner("s7 workflow consumer", args.provider) + if args.provider == "memory": + print( + " the memory provider keeps one topic across the chain, not one per " + "run; skipped" + ) + return + client, provider = await _setup.connect(args) + workflow_id, task_queue = _common.ids("june-s7") + batches = [["open", "resize"], ["rotate"], ["close", STOP]] + try: + async with Worker(client, task_queue=task_queue, workflows=[ConsumerWorkflow]): + handle = await client.start_workflow( + ConsumerWorkflow.run, + Checkpoint(), + id=workflow_id, + task_queue=task_queue, + ) + run_id: str | None = None + for number, batch in enumerate(batches, start=1): + run_id = await current_run(client, workflow_id, run_id) + # A producer made now writes to the run that is current now, + # and each batch is its own producer: a producer's (id, attempt, + # sequence) must not repeat with different content on the same + # stream, and on a store whose stream spans the chain every + # batch lands on one stream. + sender = client.get_stream_handle(workflow_id).producer( + topic=COMMANDS, producer_id=f"console-{number}", attempt=1 + ) + for op in batch: + await sender.append(Command(op)) + await sender.finish() + print(f" batch {number} sent to run ...{run_id[-6:]}: {batch}") + for line in await handle.result(): + print(f" applied {line}") + finally: + await provider.close() + + +async def main() -> None: + """Parse the flags and run the scenario.""" + await run(_common.parser(__doc__ or "").parse_args()) + + +if __name__ == "__main__": + asyncio.run(main()) diff --git a/examples/streams/june_scenarios/s8_nexus_consumers.py b/examples/streams/june_scenarios/s8_nexus_consumers.py new file mode 100644 index 000000000..3405316d8 --- /dev/null +++ b/examples/streams/june_scenarios/s8_nexus_consumers.py @@ -0,0 +1,281 @@ +r"""Scenarios over Nexus: client consumer, operation handler, workflow consumer. + +Status: "Client as Consumer over Standalone Nexus" and "Nexus operation +handler" returning a stream are implemented; "Workflow as Consumer over +Nexus" is unsupported by design. + +Why: the ``NexusStreams`` front serves outside reads through one endpoint, +so a client consumes over Nexus today, while a workflow's reads ride its +Workflow Task and never cross Nexus. + + python -m examples.streams.june_scenarios.s8_nexus_consumers native \ + --address 127.0.0.1:7433 --http http://127.0.0.1:7343 + python -m examples.streams.june_scenarios.s8_nexus_consumers workflow_streams \ + --address 127.0.0.1:7433 --http http://127.0.0.1:7343 + +The run creates its own Nexus endpoint through the operator service, routed +to its worker's task queue, and deletes it at the end. ``--http`` is the +server's Nexus HTTP ingress, which the front posts to. + +(a) His consumer activity starts a stream operation, or on a retry gets the +operation handle by id, and reads from the offset in its heartbeat. Ours +reads through the front with ``front.get_stream_handle(client, +workflow_id).read(topic=SCORES, after=cursor)`` and heartbeats each record's +cursor. No operation handle is needed to resume: every read is a short sync +operation, so the cursor alone is the resume point. The first attempt dies +after two records and the retry picks up after the second. + +(b) His handler returns ``Stream[ProgressUpdate]``. Ours returns a +``temporalio.streams.StreamRef``, the SDK's name for one stream, taken from +the producing workflow's handle with ``ref(topic=SCORES)``; the calling +workflow passes it on as its result and the client opens it with +``get_stream_handle(ref)`` on a client whose provider is the front, so the +read goes over Nexus. The ref is plain data to the operation's IDL; a stream +type of its own there is the nexgen follow-on. + +Not shown, by design: a workflow consuming over Nexus. The workflow half of +a provider records its reads on the Workflow Task, which cannot cross an +RPC, so the front has no workflow half. A workflow reads its own topics, as +in ``s7``. +""" + +from __future__ import annotations + +import argparse +import asyncio +import contextlib +from dataclasses import dataclass +from datetime import timedelta + +import nexusrpc +import nexusrpc.handler + +from examples.streams import _setup +from examples.streams.june_scenarios import _common +from temporalio import activity, nexus, streams, workflow +from temporalio.api.nexus.v1 import EndpointSpec, EndpointTarget +from temporalio.api.operatorservice.v1 import ( + CreateNexusEndpointRequest, + DeleteNexusEndpointRequest, +) +from temporalio.client import Client +from temporalio.common import RetryPolicy, WorkflowIDConflictPolicy +from temporalio.streams import BEGINNING, Cursor, RecordKind, StreamRef +from temporalio.streams.providers.nexus import NexusStreams, TemporalStreamsHandler +from temporalio.worker import Worker + + +@dataclass +class ScoreUpdate: + """One score change.""" + + home_score: int + away_score: int + + +@dataclass +class Consumed: + """What the consumer activity's last attempt read, and where it resumed.""" + + attempt: int + resumed_after: str + scores: list[str] + + +@dataclass +class GameRequest: + """Which game the operation should start.""" + + game_id: str + + +@dataclass +class CallerInput: + """The endpoint to call and the game to ask for.""" + + endpoint: str + game_id: str + + +SCORES = streams.topic("scores", ScoreUpdate) + + +@workflow.defn +class ScoresProducer: + """Publishes a short game on ``scores``.""" + + @workflow.run + async def run(self, updates: int) -> None: + """Publish ``updates`` scores, spaced out, then finish the topic.""" + scores = workflow.stream_writer(SCORES) + for n in range(updates): + scores.publish(ScoreUpdate(home_score=(n + 1) // 2, away_score=n // 2)) + await workflow.sleep(timedelta(milliseconds=300)) + scores.finish() + + +class Consumer: + """The consumer activity, holding the front it reads through.""" + + def __init__(self, front: NexusStreams) -> None: + """Read through ``front``.""" + self._front = front + + @activity.defn + async def consume_scores(self, workflow_id: str) -> Consumed: + """His standalone-Nexus consumer: resume from the cursor in the heartbeat.""" + info = activity.info() + token = str(info.heartbeat_details[0]) if info.heartbeat_details else "" + after = Cursor(token) if token else BEGINNING + stream = self._front.get_stream_handle(activity.client(), workflow_id) + scores: list[str] = [] + async with contextlib.aclosing( + stream.read(topic=SCORES, after=after) + ) as reading: + async for record in reading: + if record.kind is RecordKind.FINISH: + break + assert record.value is not None + scores.append(f"{record.value.home_score}-{record.value.away_score}") + activity.heartbeat(record.cursor.token) + if info.attempt == 1 and len(scores) == 2: + raise RuntimeError("the consumer died after two records") + return Consumed(info.attempt, token or "BEGINNING", scores) + + +@nexusrpc.service +class ScoresService: + """A service whose operation hands back a stream reference.""" + + start_game: nexusrpc.Operation[GameRequest, StreamRef] + + +@nexusrpc.handler.service_handler(service=ScoresService) +class ScoresHandler: + """Starts the producing workflow and returns where to read it.""" + + def __init__(self, task_queue: str) -> None: + """Start producers on ``task_queue``.""" + self._task_queue = task_queue + + @nexusrpc.handler.sync_operation + async def start_game( + self, _ctx: nexusrpc.handler.StartOperationContext, input: GameRequest + ) -> StreamRef: + """His "pre-create and start" handler, returning a ref to the stream.""" + # Reusing a running workflow makes a retried start of this sync + # operation hand back the same ref. + client = nexus.client() + await client.start_workflow( + ScoresProducer.run, + 3, + id=input.game_id, + task_queue=self._task_queue, + id_conflict_policy=WorkflowIDConflictPolicy.USE_EXISTING, + ) + return client.get_stream_handle(input.game_id).ref(topic=SCORES) + + +@workflow.defn +class Caller: + """Calls the operation and passes the reference on; it never reads the stream.""" + + @workflow.run + async def run(self, input: CallerInput) -> StreamRef: + """Return the stream reference the operation handed back.""" + client = workflow.create_nexus_client( + service=ScoresService, endpoint=input.endpoint + ) + return await client.execute_operation( + ScoresService.start_game, GameRequest(input.game_id) + ) + + +async def run(args: argparse.Namespace) -> None: + """Consume through the front from an activity, then read a returned reference.""" + _common.banner("s8 nexus consumers", args.provider) + client, provider = await _setup.connect(args) + workflow_id, task_queue = _common.ids("june-s8") + endpoint_name = workflow_id + created = await client.operator_service.create_nexus_endpoint( + CreateNexusEndpointRequest( + spec=EndpointSpec( + name=endpoint_name, + target=EndpointTarget( + worker=EndpointTarget.Worker( + namespace=client.namespace, task_queue=task_queue + ) + ), + ) + ) + ) + front = NexusStreams(endpoint=endpoint_name, http_address=args.http) + consumer_activities = Consumer(front) + try: + async with Worker( + client, + task_queue=task_queue, + workflows=[ScoresProducer, Caller], + activities=[consumer_activities.consume_scores], + # The front's handler serves the store this worker writes to, and + # the scores service sits beside it behind the same endpoint. + nexus_service_handlers=[ + TemporalStreamsHandler(provider, client), + ScoresHandler(task_queue), + ], + ): + print(" (a) activity consumer through the front, resuming from heartbeat") + producer = await client.start_workflow( + ScoresProducer.run, 5, id=f"{workflow_id}-a", task_queue=task_queue + ) + consumer = await client.start_activity( + consumer_activities.consume_scores, + producer.id, + id=f"{workflow_id}-consumer", + task_queue=task_queue, + start_to_close_timeout=timedelta(minutes=1), + retry_policy=RetryPolicy(initial_interval=timedelta(milliseconds=200)), + ) + consumed = await consumer.result() + print( + f" attempt {consumed.attempt} resumed after {consumed.resumed_after}" + ) + print(f" and read {consumed.scores}") + await producer.result() + + print( + " (b) an operation returns a StreamRef; the client opens it on the front" + ) + ref = await client.execute_workflow( + Caller.run, + CallerInput(endpoint_name, f"{workflow_id}-b"), + id=f"{workflow_id}-caller", + task_queue=task_queue, + ) + print(f" operation returned {ref}") + # The ref is opened on whatever provider the client carries; this + # one carries the front, so the read goes over Nexus. + config = client.config() + config["plugins"] = [front] + fronted = Client(**config) + async for record in fronted.get_stream_handle(ref).read( + result_type=ScoreUpdate + ): + print(f" {record.kind.name:6} {record.value}") + finally: + await front.close() + await client.operator_service.delete_nexus_endpoint( + DeleteNexusEndpointRequest( + id=created.endpoint.id, version=created.endpoint.version + ) + ) + await provider.close() + + +async def main() -> None: + """Parse the flags and run the scenario.""" + await run(_common.parser(__doc__ or "").parse_args()) + + +if __name__ == "__main__": + asyncio.run(main()) diff --git a/examples/streams/path_a_publish.py b/examples/streams/path_a_publish.py new file mode 100644 index 000000000..1e1eaf77d --- /dev/null +++ b/examples/streams/path_a_publish.py @@ -0,0 +1,80 @@ +"""Path A: the workflow publishes, and a backend follows from now. + + python -m examples.streams.path_a_publish workflow_streams + python -m examples.streams.path_a_publish native --address 127.0.0.1:7333 + python -m examples.streams.path_a_publish redis --redis redis://127.0.0.1:6379 + +The topic is defined once, with the type its records carry, and both sides +refer to that definition. The workflow reports its progress on it with +``workflow.stream_writer``; ``publish`` is a plain call, the record is +buffered and commits with the Workflow Task, so a reader never sees a step +the workflow did not commit. The backend positions itself with ``latest()`` +and reads what comes after, the way a UI attaches to a job that is already +running; the read ends by itself once the workflow is closed and the tail +has been delivered. +""" + +from __future__ import annotations + +import asyncio +import uuid +from dataclasses import dataclass +from datetime import timedelta + +from examples.streams import _setup +from temporalio import streams, workflow +from temporalio.worker import Worker + + +@dataclass +class Progress: + """One step of the job, as published.""" + + step: int + of: int + + +PROGRESS = streams.topic("progress", Progress) + + +@workflow.defn +class Job: + """Works through ``steps`` and publishes each one on ``progress``.""" + + @workflow.run + async def run(self, steps: int) -> int: + """Publish one record per step, then finish the topic.""" + progress = workflow.stream_writer(PROGRESS) + for step in range(steps): + progress.publish(Progress(step=step, of=steps)) + # A timer between steps, so each record commits with its own + # Workflow Task and a follower sees them arrive one at a time. + await workflow.sleep(timedelta(milliseconds=200)) + progress.finish() + return steps + + +async def main() -> None: + """Run the job and follow its progress from outside.""" + args = _setup.parser(__doc__ or "").parse_args() + client, provider = await _setup.connect(args) + workflow_id = f"path-a-{uuid.uuid4().hex[:8]}" + task_queue = f"tq-{workflow_id}" + try: + async with Worker(client, task_queue=task_queue, workflows=[Job]): + handle = await client.start_workflow( + Job.run, 5, id=workflow_id, task_queue=task_queue + ) + stream = client.get_stream_handle(workflow_id) + # Follow from now: whatever landed before this point is not + # re-read, which is how a client attaches to a running job. + since = await stream.latest(topic=PROGRESS) + async for record in stream.read(topic=PROGRESS, after=since): + print(f" {record.kind.name:7} {record.value} at {record.cursor.token}") + print(f"workflow {workflow_id} took {await handle.result()} steps") + finally: + await provider.close() + + +if __name__ == "__main__": + asyncio.run(main()) diff --git a/examples/streams/path_b_produce.py b/examples/streams/path_b_produce.py new file mode 100644 index 000000000..773d0a972 --- /dev/null +++ b/examples/streams/path_b_produce.py @@ -0,0 +1,156 @@ +"""Path B: an Activity and a backend produce, and a backend consumes. + + python -m examples.streams.path_b_produce workflow_streams + python -m examples.streams.path_b_produce native --address 127.0.0.1:7333 + python -m examples.streams.path_b_produce redis --redis redis://127.0.0.1:6379 + +Outside workflow code a stream is reached through a handle: an Activity asks +its context with ``activity.stream_handle()``, which is its own workflow +pinned to its run, and a backend asks its client with +``client.get_stream_handle(workflow_id)``. Both hand out the same handle, +with the same verbs, and both refer to the topics defined once below, so the +producer's ``append`` and the consumer's records are typed the same way. A +producer writes on its own account, visible as soon as the store accepts the +record, under an identity that lets readers tell a retry from a new attempt: +the Activity's producer takes the Activity's own id and attempt, a backend +names its own. + +The model Activity here fails halfway through its first attempt. Its retry +starts over, and the consumer sees that as a ``SUPERSEDED`` record before the +new attempt's first token, so it can drop what the earlier attempt produced. +""" + +from __future__ import annotations + +import asyncio +import contextlib +import uuid +from dataclasses import dataclass +from datetime import timedelta + +from examples.streams import _setup +from temporalio import activity, streams, workflow +from temporalio.common import RetryPolicy +from temporalio.streams import RecordKind +from temporalio.worker import Worker + + +@dataclass +class Token: + """One piece of model output.""" + + n: int + + +@dataclass +class Note: + """A remark a backend attaches to the session.""" + + text: str + + +INPUTS = streams.topic("inputs", Token) +NOTES = streams.topic("notes", Note) + + +@activity.defn +async def generate(count: int) -> None: + """Stream ``count`` tokens onto this workflow's ``inputs`` topic. + + No workflow id and no run id: the handle is this Activity's own + workflow, pinned to its run, and the producer's identity is the + Activity's, so a retry deduplicates and a new attempt is reported. + """ + model = activity.stream_handle().producer(topic=INPUTS) + for n in range(count): + await model.append(Token(n)) + if n == 1 and activity.info().attempt == 1: + raise RuntimeError("the model connection dropped") + await model.finish() + + +@workflow.defn +class Session: + """Runs the model, then stays open until the backend has said its piece.""" + + def __init__(self) -> None: + """Start open.""" + self._closed = False + + @workflow.signal + def close(self) -> None: + """Let the run end; a backend producer needs the run open to append.""" + self._closed = True + + @workflow.run + async def run(self, count: int) -> None: + """Generate ``count`` tokens, then wait to be closed.""" + await workflow.execute_activity( + generate, + count, + start_to_close_timeout=timedelta(minutes=1), + retry_policy=RetryPolicy( + initial_interval=timedelta(milliseconds=100), maximum_attempts=3 + ), + ) + await workflow.wait_condition(lambda: self._closed) + + +async def main() -> None: + """Run the session; produce from the backend; consume both topics.""" + args = _setup.parser(__doc__ or "").parse_args() + client, provider = await _setup.connect(args) + workflow_id = f"path-b-{uuid.uuid4().hex[:8]}" + task_queue = f"tq-{workflow_id}" + try: + async with Worker( + client, task_queue=task_queue, workflows=[Session], activities=[generate] + ): + handle = await client.start_workflow( + Session.run, 3, id=workflow_id, task_queue=task_queue + ) + stream = client.get_stream_handle(workflow_id) + + # A backend producer names its own identity. + notes = stream.producer(topic=NOTES, producer_id="operator", attempt=1) + await notes.append(Note("reviewing this session")) + await notes.finish() + + # A backend consumer keeps the tokens per attempt and drops an + # attempt the moment a newer one starts writing. + # Both reads stop at FINISH while the run is still open, so each one + # is closed on the way out rather than left for the collector. + tokens: dict[int, list[int]] = {} + async with contextlib.aclosing(stream.read(topic=INPUTS)) as records: + async for record in records: + if record.kind is RecordKind.SUPERSEDED: + assert record.supersession is not None + attempt = record.supersession.previous_attempt + dropped = tokens.pop(attempt, []) + print(f" attempt {attempt} superseded") + print(f" dropped {dropped}") + continue + if record.kind is RecordKind.FINISH: + print( + f" {record.producer_id} attempt {record.attempt} finished" + ) + break + assert record.value is not None + tokens.setdefault(record.attempt, []).append(record.value.n) + print(f"kept {tokens}") + + async with contextlib.aclosing(stream.read(topic=NOTES)) as notes_read: + async for note in notes_read: + if note.kind is RecordKind.FINISH: + break + assert note.value is not None + print(f" note from {note.producer_id}: {note.value.text}") + + await handle.signal(Session.close) + await handle.result() + finally: + await provider.close() + + +if __name__ == "__main__": + asyncio.run(main()) diff --git a/examples/streams/path_c_consume.py b/examples/streams/path_c_consume.py new file mode 100644 index 000000000..4537e2318 --- /dev/null +++ b/examples/streams/path_c_consume.py @@ -0,0 +1,101 @@ +"""Path C: the workflow consumes a topic fed from outside, with a cold cache. + + python -m examples.streams.path_c_consume workflow_streams + python -m examples.streams.path_c_consume native --address 127.0.0.1:7333 + python -m examples.streams.path_c_consume redis --redis redis://127.0.0.1:6379 + +The workflow reads its ``commands`` topic with ``workflow.stream_reader`` and +runs an Activity for each record; the topic is defined once, so the reader's +records and the Activity's argument share one type. A read is an observation +the SDK records: what the reader handed to workflow code commits with the +Workflow Task, so replay re-supplies the same records in the same order and +each Activity result is matched to the command that caused it. The worker +runs with the workflow cache off, so every Workflow Task rebuilds the +workflow from History and the loop completing at all is the proof. +""" + +from __future__ import annotations + +import asyncio +import uuid +from dataclasses import dataclass +from datetime import timedelta + +from examples.streams import _setup +from temporalio import activity, streams, workflow +from temporalio.streams import RecordKind +from temporalio.worker import Worker + + +@dataclass +class Command: + """One instruction from the console.""" + + op: str + + +COMMANDS = streams.topic("commands", Command) + + +@activity.defn +async def apply(command: Command) -> str: + """Carry out one command.""" + return f"applied {command.op}" + + +@workflow.defn +class Controller: + """Acts on each command as it arrives, until the sender finishes.""" + + @workflow.run + async def run(self) -> list[str]: + """Return what was applied, in the order the commands arrived.""" + commands = workflow.stream_reader(COMMANDS) + applied: list[str] = [] + async for record in commands: + if record.kind is RecordKind.FINISH: + break + if record.kind is not RecordKind.DATA: + continue + assert record.value is not None + applied.append( + await workflow.execute_activity( + apply, record.value, start_to_close_timeout=timedelta(minutes=1) + ) + ) + return applied + + +async def main() -> None: + """Feed commands from the backend, one at a time, and print what was applied.""" + args = _setup.parser(__doc__ or "").parse_args() + client, provider = await _setup.connect(args) + workflow_id = f"path-c-{uuid.uuid4().hex[:8]}" + task_queue = f"tq-{workflow_id}" + try: + async with Worker( + client, + task_queue=task_queue, + workflows=[Controller], + activities=[apply], + # Off, so every task replays the recorded reads from History. + max_cached_workflows=0, + ): + handle = await client.start_workflow( + Controller.run, id=workflow_id, task_queue=task_queue + ) + console = client.get_stream_handle(workflow_id).producer( + topic=COMMANDS, producer_id="console", attempt=1 + ) + for op in ("open", "resize", "close"): + await console.append(Command(op)) + # Spaced out, so the commands arrive across several tasks. + await asyncio.sleep(0.3) + await console.finish() + print(f"workflow {workflow_id} applied {await handle.result()}") + finally: + await provider.close() + + +if __name__ == "__main__": + asyncio.run(main()) diff --git a/examples/streams/run.py b/examples/streams/run.py new file mode 100644 index 000000000..ab21af954 --- /dev/null +++ b/examples/streams/run.py @@ -0,0 +1,147 @@ +r"""Run the same agent on whichever provider is configured. + + python -m examples.streams.run workflow_streams + python -m examples.streams.run redis --redis redis://127.0.0.1:6379 + python -m examples.streams.run native --address 127.0.0.1:7333 + python -m examples.streams.run nexus --endpoint streams-e2e \ + --http http://127.0.0.1:7243 + +The Nexus mode needs an endpoint that routes to the handler worker's task +queue. The flag takes the endpoint's name, which the front resolves to its id +through the client, and ``--http`` is the server's HTTP address as a full URL:: + + temporal operator nexus endpoint create --name streams-e2e \ + --target-task-queue streams-handlers-e2e + +The provider is registered once, on the client, and that is the only +provider-specific line here. The workers inherit it, the Activity reaches its +workflow's stream through ``activity.stream_handle()``, and the backend below +reads through ``client.get_stream_handle()``, or through the Nexus front +standing in for the store when there is one. +""" + +from __future__ import annotations + +import argparse +import asyncio +import contextlib +import uuid + +from examples.streams import _setup +from examples.streams.agent import DECISIONS, Agent, generate, record_decision +from temporalio.client import Client +from temporalio.streams import RecordKind, StreamHandle +from temporalio.worker import Worker + + +async def main() -> None: + """Run the loop once on the provider named on the command line.""" + parser = argparse.ArgumentParser() + parser.add_argument("provider", choices=[*_setup.PROVIDERS, "nexus"]) + parser.add_argument("--address", default="localhost:7233") + parser.add_argument("--redis", default="redis://127.0.0.1:6379") + parser.add_argument( + "--endpoint", + default="", + help="nexus endpoint name; see streams_demo/README.md", + ) + parser.add_argument("--http", default="http://127.0.0.1:7243") + parser.add_argument( + "--behind", + default="workflow_streams", + help="the store a nexus handler serves", + ) + parser.add_argument("--records", type=int, default=3) + parser.add_argument("--cache", type=int, default=100) + parser.add_argument( + "--handler-queue", + default="streams-handlers-e2e", + help="task queue the nexus endpoint routes to", + ) + args = parser.parse_args() + # Checked before any provider exists, so a missing flag is a usage error + # rather than a traceback with a provider left open. + if args.provider == "nexus" and not args.endpoint: + parser.error("nexus needs --endpoint ") + + store = args.behind if args.provider == "nexus" else args.provider + provider = _setup.make_provider(store, args) + client = await Client.connect(args.address, plugins=[provider]) + workflow_id = f"streams-example-{uuid.uuid4().hex[:8]}" + task_queue = f"tq-{workflow_id}" + + workers = [ + Worker( + client, + task_queue=task_queue, + workflows=[Agent], + activities=[generate, record_decision], + # Warm, because two of these transports park work against the + # running workflow. The native provider also runs at zero, which + # is its own result rather than something this example shows. + max_cached_workflows=args.cache, + ) + ] + front = None + if args.provider == "nexus": + from temporalio.streams.providers.nexus import ( + NexusStreams, + TemporalStreamsHandler, + ) + + workers.append( + Worker( + client, + task_queue=args.handler_queue, + nexus_service_handlers=[TemporalStreamsHandler(provider, client)], + ) + ) + front = NexusStreams(endpoint=args.endpoint, http_address=args.http) + + try: + async with contextlib.AsyncExitStack() as running: + for worker in workers: + await running.enter_async_context(worker) + handle = await client.start_workflow( + Agent.run, args.records, id=workflow_id, task_queue=task_queue + ) + print(f"provider={args.provider} workflow={workflow_id}") + + # The same handle either way: from the Nexus front when there is + # one, otherwise from the provider registered on the client. + stream: StreamHandle = ( + front.get_stream_handle(client, workflow_id) + if front is not None + else client.get_stream_handle(workflow_id) + ) + # Counted apart: a retried generator makes the workflow retract the + # earlier attempt, and those records are correct output rather than + # echoes that the workflow's own count would have to agree with. + echoes = retractions = 0 + # The read ends by itself once the workflow is closed and the tail + # has been delivered, on every provider. + async for record in stream.read(topic=DECISIONS): + print( + f" {record.kind.name:11} {record.value} at {record.cursor.token}" + ) + if record.kind is not RecordKind.DATA: + continue + assert record.value is not None + if record.value.echo is not None: + echoes += 1 + else: + retractions += 1 + + decided = await handle.result() + print( + f"workflow decided {decided}; reader saw {echoes} echoes and " + f"{retractions} retractions" + ) + finally: + if front is not None: + await front.close() + await provider.close() + + +if __name__ == "__main__": + asyncio.run(main()) diff --git a/pyproject.toml b/pyproject.toml index 3e50534cd..943fca250 100644 --- a/pyproject.toml +++ b/pyproject.toml @@ -278,7 +278,7 @@ reportUnnecessaryIsInstance = "none" reportUnnecessaryTypeIgnoreComment = "none" reportUnusedCallResult = "none" reportUnknownLambdaType = "none" -include = ["temporalio", "tests", "streams_demo"] +include = ["temporalio", "tests", "examples", "streams_demo"] exclude = [ # Exclude auto generated files "temporalio/api", diff --git a/streams_demo/README.md b/streams_demo/README.md index a52e56ae8..58faec26f 100644 --- a/streams_demo/README.md +++ b/streams_demo/README.md @@ -51,14 +51,30 @@ provider it is configured with. Workers still configure a storage provider, because workflow reads and writes ride the Workflow Task. ```python -front = NexusStreams(endpoint=endpoint_id) +front = NexusStreams(endpoint="streams-e2e") stream = front.get_stream_handle(client, workflow_id) producer = stream.producer(topic="inputs", producer_id="model", attempt=1) ``` -See `tests/streams/test_nexus_provider.py` for the endpoint setup and the -handler worker. +The endpoint has to exist and route to the handler worker's task queue. The +provider takes its name and resolves it to the id through the client: + +```sh +temporal operator nexus endpoint create --name streams-e2e \ + --target-task-queue streams-handlers-e2e +``` + +See `tests/streams/test_nexus_provider.py` for the handler worker, and +`examples/streams/run.py nexus --endpoint streams-e2e --http http://:` +for the whole loop behind it. The conformance suite runs the same expectations on every provider: `pytest tests/streams/` (memory and Workflow Streams on the test server), plus `STREAMS_LIVE=native|redis|nexus` for the suites that need a store. + +## The examples + +`examples/streams/` is the same thing written as a worked example rather than +a measurement run: `agent.py` holds the workflow and activities, identical on +every provider, and `run.py` picks one. See its `make_provider`, which +is the whole difference between the options: one constructor call.