Skip to content

Commit d8f85b9

Browse files
committed
Added the worked example that runs one agent on every provider.
1 parent 1771757 commit d8f85b9

10 files changed

Lines changed: 1164 additions & 5 deletions

File tree

‎examples/__init__.py‎

Whitespace-only changes.

‎examples/streams/__init__.py‎

Whitespace-only changes.

‎examples/streams/agent.py‎

Lines changed: 90 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,90 @@
1+
"""The workflow and activities. Identical on every provider.
2+
3+
Nothing here names a store, a transport, or an option. The loop reads its
4+
inbound stream, decides, publishes the decision, and runs an ordinary
5+
activity in the same workflow task, which is the shape the design doc calls
6+
Paths A, B and C together.
7+
"""
8+
9+
from __future__ import annotations
10+
11+
from datetime import timedelta
12+
13+
from temporalio import activity, streams, workflow
14+
15+
16+
@activity.defn
17+
async def generate(workflow_id: str, count: int) -> None:
18+
"""Stream model output into the workflow's inbound stream.
19+
20+
The producer carries this activity's own id and attempt, so a retry
21+
deduplicates and a new attempt is reported to readers as a supersession.
22+
"""
23+
producer = await streams.producer(
24+
activity.client(), workflow_id=workflow_id, stream="inputs"
25+
)
26+
for n in range(count):
27+
await producer.append({"n": n})
28+
await producer.finish()
29+
30+
31+
@activity.defn
32+
async def record_decision(decision: dict) -> str:
33+
"""An ordinary activity, run from the same task that read and published."""
34+
return f"recorded {decision['echo']}"
35+
36+
37+
@workflow.defn
38+
class Agent:
39+
def __init__(self) -> None:
40+
# Lets the provider install what it needs before the first task
41+
# completes. A no-op except on the transport that serves outside
42+
# readers through handlers on this workflow.
43+
streams.prepare()
44+
self._done = False
45+
46+
@workflow.signal
47+
def release(self) -> None:
48+
"""Lets the run end once a reader has finished following it.
49+
50+
Only needed because an Option 0 stream dies with its workflow, so a
51+
demo that ends immediately would leave nothing to read.
52+
"""
53+
self._done = True
54+
55+
@workflow.run
56+
async def run(self, count: int) -> int:
57+
inputs = streams.reader("inputs", type=dict, idle_timeout=timedelta(seconds=1))
58+
decisions = streams.writer("decisions")
59+
60+
await workflow.start_activity(
61+
generate,
62+
args=[workflow.info().workflow_id, count],
63+
start_to_close_timeout=timedelta(minutes=1),
64+
)
65+
66+
seen = 0
67+
async for record in inputs:
68+
if record.kind is streams.RecordKind.FINISH:
69+
break
70+
if record.kind is streams.RecordKind.SUPERSEDED:
71+
await decisions.publish(
72+
{"retracting_attempt": record.value.previous_attempt}
73+
)
74+
continue
75+
seen += 1
76+
decision = {"echo": record.value["n"]}
77+
await decisions.publish(decision)
78+
await workflow.execute_activity(
79+
record_decision,
80+
decision,
81+
start_to_close_timeout=timedelta(minutes=1),
82+
)
83+
84+
await decisions.finish()
85+
await workflow.wait_condition(lambda: self._done)
86+
# Lets go of anything the provider parked against this run. A no-op
87+
# except on the transport that parks an outside reader here.
88+
streams.drain()
89+
await workflow.wait_condition(workflow.all_handlers_finished)
90+
return seen

‎examples/streams/run.py‎

Lines changed: 172 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,172 @@
1+
"""Run the same agent on whichever provider is configured.
2+
3+
python -m examples.streams.run workflow_streams
4+
python -m examples.streams.run redis --redis redis://127.0.0.1:6399
5+
python -m examples.streams.run native --address 127.0.0.1:7233
6+
python -m examples.streams.run nexus --endpoint <endpoint-id>
7+
8+
The only provider-specific code in this file is :func:`configure_provider`,
9+
which turns a name into one ``streams.configure`` call. Everything below it,
10+
and all of ``agent.py``, is the same on every option.
11+
"""
12+
13+
from __future__ import annotations
14+
15+
import argparse
16+
import asyncio
17+
import uuid
18+
19+
from temporalio import streams
20+
from temporalio.client import Client
21+
from temporalio.worker import Worker
22+
23+
from examples.streams.agent import Agent, generate, record_decision
24+
25+
26+
def configure_provider(args: argparse.Namespace) -> None:
27+
"""The whole difference between the options."""
28+
if args.provider == "workflow_streams":
29+
streams.configure(provider="workflow_streams")
30+
elif args.provider == "redis":
31+
streams.configure(provider="redis", url=args.redis)
32+
elif args.provider == "native":
33+
streams.configure(provider="native")
34+
elif args.provider == "nexus":
35+
# The worker still needs a store; the endpoint is how everything
36+
# outside the worker reaches it without naming one.
37+
streams.configure(provider=args.behind)
38+
else:
39+
raise SystemExit(f"unknown provider {args.provider}")
40+
41+
42+
async def ensure_stream(args: argparse.Namespace, client: Client, workflow_id: str) -> None:
43+
"""Create the inbound stream when the provider needs it to pre-exist.
44+
45+
The two storage providers want opposite orders, which is worth knowing
46+
before writing an application against either. The native provider
47+
resolves a subscription against a stream the server already has, and a
48+
workflow cannot create one from workflow code because that would be I/O,
49+
so whoever starts the workflow creates it first. The client-side provider
50+
is the mirror image: it names streams under the run's chain key, so its
51+
producer needs the workflow to exist already. Option 0 mints on first
52+
publish and does not care.
53+
"""
54+
if args.provider != "native":
55+
return
56+
from temporalio.client_stream import StreamClient
57+
from temporalio.streams.providers.native import inbound_stream_id
58+
59+
streams_client = StreamClient.connect(args.address, client.namespace)
60+
await streams_client.create(inbound_stream_id(workflow_id, "inputs"))
61+
62+
63+
def outside_surface(args: argparse.Namespace):
64+
"""Where an outside producer or consumer connects.
65+
66+
The same handles either way: through the Nexus endpoint when there is
67+
one, otherwise straight at the configured provider.
68+
"""
69+
if args.provider == "nexus":
70+
from temporalio.streams._provider import instance
71+
72+
return instance("nexus", endpoint=args.endpoint, http_address=args.http)
73+
return streams
74+
75+
76+
async def main() -> None:
77+
parser = argparse.ArgumentParser()
78+
parser.add_argument(
79+
"provider", choices=["workflow_streams", "redis", "native", "nexus"]
80+
)
81+
parser.add_argument("--address", default="localhost:7233")
82+
parser.add_argument("--redis", default="redis://127.0.0.1:6399")
83+
parser.add_argument("--endpoint", default="", help="nexus endpoint id")
84+
parser.add_argument("--http", default="http://127.0.0.1:7243")
85+
parser.add_argument(
86+
"--behind",
87+
default="workflow_streams",
88+
help="the store a nexus handler serves",
89+
)
90+
parser.add_argument("--records", type=int, default=3)
91+
parser.add_argument("--cache", type=int, default=100)
92+
parser.add_argument(
93+
"--handler-queue",
94+
default="streams-handlers-e2e",
95+
help="task queue the nexus endpoint routes to",
96+
)
97+
args = parser.parse_args()
98+
99+
configure_provider(args)
100+
client = await Client.connect(args.address)
101+
workflow_id = f"streams-example-{uuid.uuid4().hex[:8]}"
102+
task_queue = f"tq-{workflow_id}"
103+
104+
workers = [
105+
Worker(
106+
client,
107+
task_queue=task_queue,
108+
workflows=[Agent],
109+
activities=[generate, record_decision],
110+
# Warm, because two of these transports park work against the
111+
# running workflow. The native provider also runs at zero, which
112+
# is its own result rather than something this example shows.
113+
max_cached_workflows=args.cache,
114+
**streams.worker_options(),
115+
)
116+
]
117+
if args.provider == "nexus":
118+
from temporalio.streams.providers.nexus import TemporalStreamsHandler
119+
120+
workers.append(
121+
Worker(
122+
client,
123+
task_queue=args.handler_queue,
124+
nexus_service_handlers=[
125+
TemporalStreamsHandler(client, provider=args.behind)
126+
],
127+
)
128+
)
129+
130+
async with workers[0]:
131+
async with workers[-1] if len(workers) > 1 else _null():
132+
await ensure_stream(args, client, workflow_id)
133+
handle = await client.start_workflow(
134+
Agent.run, args.records, id=workflow_id, task_queue=task_queue
135+
)
136+
print(f"provider={args.provider} workflow={workflow_id}")
137+
138+
reader = await outside_surface(args).consumer(
139+
client if args.provider != "nexus" else None,
140+
workflow_id=workflow_id,
141+
)
142+
seen = 0
143+
records = reader.read(type=dict, topic="decisions")
144+
try:
145+
async for record in records:
146+
print(
147+
f" {record.kind.name:11} {record.value} at {record.cursor.token}"
148+
)
149+
if record.kind is streams.RecordKind.FINISH:
150+
break
151+
seen += 1
152+
finally:
153+
# Closed before the workflow is released, because a reader
154+
# still polling this run would have its in-flight poll
155+
# cancelled when the workflow lets its readers go.
156+
await records.aclose()
157+
158+
await handle.signal(Agent.release)
159+
decided = await handle.result()
160+
print(f"workflow decided {decided}, reader saw {seen}")
161+
162+
163+
class _null:
164+
async def __aenter__(self):
165+
return self
166+
167+
async def __aexit__(self, *exc):
168+
return False
169+
170+
171+
if __name__ == "__main__":
172+
asyncio.run(main())

‎streams_demo/README.md‎

Lines changed: 7 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -62,3 +62,10 @@ handler worker.
6262
The conformance suite runs the same expectations on every provider:
6363
`pytest tests/streams/` (in-memory, no server), plus
6464
`STREAMS_LIVE=workflow_streams|nexus` for the live suites.
65+
66+
## The examples
67+
68+
`examples/streams/` is the same thing written as a worked example rather than
69+
a measurement run: `agent.py` holds the workflow and activities, identical on
70+
every provider, and `run.py` picks one. See its `configure_provider`, which
71+
is the whole difference between the options.

0 commit comments

Comments
 (0)