Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
53 commits
Select commit Hold shift + click to select a range
09dccf4
Added the Workflow Streams provider over the shipped wire format.
moedash Sep 15, 2026
f088496
Let a workflow drain parked pollers before returning.
moedash Sep 15, 2026
3393901
Added the prepare and drain hooks the Workflow Streams transport needs.
moedash Sep 16, 2026
ee1f64d
Made cursors exclusive, added latest(), and allowed topic appends.
moedash Sep 16, 2026
91d11f6
Served a closed workflow's stream tail by Query on the shipped transp…
moedash Sep 16, 2026
87ae2c8
Made the Workflow Streams provider pass the linters.
moedash Sep 16, 2026
305a0d2
Added the changelog entry for the Workflow Streams provider.
moedash Sep 16, 2026
f33de01
Added public accessors to Workflow Streams for code that rides its wire.
moedash Sep 18, 2026
643bcdd
Keyed the Workflow Streams runtime by workflow instance, not run id.
moedash Sep 18, 2026
d015a6f
Committed producer sequences only after the publish signal was accepted.
moedash Sep 18, 2026
0e8236e
Adapted the Workflow Streams provider to the interface changes.
moedash Sep 18, 2026
dff4341
Narrowed the tail query's error handling and filtered owner reads by …
moedash Sep 18, 2026
50b69e4
Ran the Workflow Streams provider tests on the test environment's ser…
moedash Sep 18, 2026
06cbd78
Rewrote the Workflow Streams provider on the plugin surface.
moedash Sep 21, 2026
0de91c2
Ran the conformance suite over Workflow Streams and covered chain reads.
moedash Sep 21, 2026
c44ba22
Detached the stream this instance captured, not the thread's current …
moedash Sep 21, 2026
309dcec
Drove the poll Update directly so a cancelled read stays cancelled.
moedash Sep 21, 2026
1842f6e
Registered the Workflow Streams conformance provider on the client.
moedash Sep 21, 2026
1802e8c
Took topic definitions on the Workflow Streams handle.
moedash Sep 21, 2026
c29e057
Merged the interface branch's connect config fix and Core repin.
moedash Sep 22, 2026
f171df2
Keyed the publish dedupe on where a producer's records end.
moedash Sep 25, 2026
ffd14d1
Made the read path page its tail and name its failures.
moedash Sep 25, 2026
1508cb2
Merge branch 'moe/AI-198-py-04-accessors' into moe/AI-198-streams-pro…
moedash Sep 26, 2026
bae3a83
Numbered the Workflow Streams producer's records from one.
moedash Sep 26, 2026
cc6dbbd
Typed the empty history iterator in the provider test double.
moedash Sep 26, 2026
3dc0d66
Retried a poll that reaches a run before its first task registers.
moedash Sep 29, 2026
c0dcbf7
Merge remote-tracking branch 'origin/moe/AI-198-py-04-accessors' into…
moedash Sep 29, 2026
2332342
Absorbed the current stream contract on Workflow Streams.
moedash Sep 29, 2026
0e567a7
Hosted a workflow activity's own streams as reserved topics in the log.
moedash Sep 30, 2026
93a4e45
Followed a reset run from the position the read had reached.
moedash Sep 30, 2026
b532125
Carried outside publishes by an Update that answers and refuses.
moedash Sep 30, 2026
a910ebf
Merge remote-tracking branch 'origin/moe/AI-198-py-04-accessors' into…
moedash Sep 30, 2026
3db4601
Absorbed refs, the body hash and the standalone refusal on Workflow S…
moedash Sep 30, 2026
a213089
Merge remote-tracking branch 'origin/moe/AI-198-py-04-accessors' into…
moedash Sep 30, 2026
e07c3c9
Merge remote-tracking branch 'origin/moe/AI-198-py-04-accessors' into…
moedash Sep 30, 2026
b9ee545
Declared the standalone byte bound and age trim absent on Workflow St…
moedash Sep 30, 2026
cb6c773
Registered the provider's handlers before the first task's Updates run.
moedash Oct 1, 2026
303ed81
Merge remote-tracking branch 'origin/moe/AI-198-py-04-accessors' into…
moedash Oct 1, 2026
8ed2523
Declared the byte-cap refusal absent on Workflow Streams.
moedash Oct 1, 2026
b473638
Merge branch 'moe/AI-198-py-04-accessors' into moe/AI-198-streams-pro…
moedash Oct 1, 2026
bb7188f
Merge branch 'moe/AI-198-py-04-accessors' into moe/AI-198-streams-pro…
moedash Oct 1, 2026
d7e7492
Merge branch 'moe/AI-198-py-04-accessors' into moe/AI-198-streams-pro…
moedash Oct 2, 2026
d91ec39
Merge branch 'moe/AI-198-py-04-accessors' into moe/AI-198-streams-pro…
moedash Oct 2, 2026
6a77ad9
Merge branch 'moe/AI-198-py-04-accessors' into moe/AI-198-streams-pro…
moedash Oct 2, 2026
b680133
Merge branch 'moe/AI-198-py-04-accessors' into moe/AI-198-streams-pro…
moedash Oct 2, 2026
26384f2
Merge branch 'moe/AI-198-py-04-accessors' into moe/AI-198-streams-pro…
moedash Oct 2, 2026
fce9be2
Merge branch 'moe/AI-198-py-04-accessors' into moe/AI-198-streams-pro…
moedash Oct 2, 2026
5f338e8
Merge branch 'moe/AI-198-py-04-accessors' into moe/AI-198-streams-pro…
moedash Oct 2, 2026
f121ad6
Merge branch 'moe/AI-198-py-04-accessors' into moe/AI-198-streams-pro…
moedash Oct 2, 2026
8dfcc08
Merge branch 'moe/AI-198-py-04-accessors' into moe/AI-198-streams-pro…
moedash Oct 2, 2026
6c73610
Merge branch 'moe/AI-198-py-04-accessors' into moe/AI-198-streams-pro…
moedash Oct 3, 2026
0dce72e
Merge branch 'moe/AI-198-py-04-accessors' into moe/AI-198-streams-pro…
moedash Oct 3, 2026
7ef8805
Merge branch 'moe/AI-198-py-04-accessors' into moe/AI-198-streams-pro…
moedash Oct 3, 2026
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
10 changes: 10 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -57,6 +57,16 @@ to include examples, links to docs, or any other relevant information.
converter, so a payload codec and external storage apply to them.
`temporalio.streams.providers.memory.MemoryStreams` is the in-memory
reference provider the conformance tests run against.
- **Experimental**: `temporalio.streams.providers.workflow_streams.WorkflowStreamsProvider`
serves the stream interface over the shipped Workflow Streams transport as a
worker plugin, so a workflow reads and publishes through
`temporalio.contrib.workflow_streams` without naming it. Records are the
`StreamRecord` proto inside the shipped item payload, and a handle without a
run id follows continue-as-new run by run and a reset into the run reset to.
An outside publish is an Update that answers with the batch's position and
refuses a conflicting repeat, falling back to the shipped Signal on a
workflow whose worker predates it. A workflow's activity keeps its own
streams in the workflow's log under `activity/<id>/<name>`.

### Changed

Expand Down
12 changes: 11 additions & 1 deletion temporalio/contrib/workflow_streams/__init__.py
Original file line number Diff line number Diff line change
Expand Up @@ -13,12 +13,18 @@
"""

from temporalio.contrib.workflow_streams._client import WorkflowStreamClient
from temporalio.contrib.workflow_streams._stream import WorkflowStream
from temporalio.contrib.workflow_streams._stream import (
POLL_UPDATE_NAME,
PUBLISH_SIGNAL_NAME,
WorkflowStream,
)
from temporalio.contrib.workflow_streams._topic_handle import (
TopicHandle,
WorkflowTopicHandle,
)
from temporalio.contrib.workflow_streams._types import (
STREAM_DRAINING_ERROR_TYPE,
TRUNCATED_OFFSET_ERROR_TYPE,
PollInput,
PollResult,
PublishEntry,
Expand All @@ -29,6 +35,10 @@
)

__all__ = [
"POLL_UPDATE_NAME",
"PUBLISH_SIGNAL_NAME",
"STREAM_DRAINING_ERROR_TYPE",
"TRUNCATED_OFFSET_ERROR_TYPE",
"PollInput",
"PollResult",
"PublishEntry",
Expand Down
14 changes: 14 additions & 0 deletions temporalio/contrib/workflow_streams/_client.py
Original file line number Diff line number Diff line change
Expand Up @@ -325,6 +325,20 @@ def topic(
self._topic_types[name] = bound
return TopicHandle(self, name, bound)

@property
def handle(self) -> WorkflowHandle[Any, Any]:
"""The workflow handle this client publishes to and polls.

Re-targeted when :py:meth:`subscribe` follows a continue-as-new, so
read it when needed rather than caching it.
"""
return self._handle

@property
def payload_converter(self) -> PayloadConverter:
"""The sync payload converter used for per-item encode and decode."""
return self._payload_converter()

async def flush(self) -> None:
"""Flush buffered (and pending) items and wait for server confirmation.

Expand Down
37 changes: 35 additions & 2 deletions temporalio/contrib/workflow_streams/_stream.py
Original file line number Diff line number Diff line change
Expand Up @@ -50,8 +50,22 @@
_WorkflowStreamWireItem,
)

_PUBLISH_SIGNAL = "__temporal_workflow_stream_publish"
_POLL_UPDATE = "__temporal_workflow_stream_poll"
PUBLISH_SIGNAL_NAME = "__temporal_workflow_stream_publish"
"""The signal :class:`WorkflowStream` registers for external publishes.

Public so code that sends the signal itself, with its own publisher identity,
does not have to copy the name.
"""

POLL_UPDATE_NAME = "__temporal_workflow_stream_poll"
"""The update :class:`WorkflowStream` registers for long polls.

Public so code that drives the poll itself, with its own retry and
cancellation rules, does not have to copy the name.
"""

_PUBLISH_SIGNAL = PUBLISH_SIGNAL_NAME
_POLL_UPDATE = POLL_UPDATE_NAME
_OFFSET_QUERY = "__temporal_workflow_stream_offset"

_MAX_POLL_RESPONSE_BYTES = 1_000_000
Expand Down Expand Up @@ -234,6 +248,25 @@ def topic(
self._topic_types[name] = bound
return WorkflowTopicHandle(self, name, bound)

@property
def next_offset(self) -> int:
"""The global offset the next published item will receive."""
return self._base_offset + len(self._log)

def items_from(self, offset: int) -> list[tuple[int, str, Payload]]:
"""Return ``(offset, topic, payload)`` for every item at or past ``offset``.

Reads the log in place, so it is safe to call from a
:func:`temporalio.workflow.wait_condition` predicate. An ``offset``
below the truncation base starts at the base instead; the offsets in
the result say where the items actually sit.
"""
start = max(offset, self._base_offset) - self._base_offset
return [
(self._base_offset + index, item.topic, item.data)
for index, item in enumerate(self._log[start:], start)
]

def get_state(
self, *, publisher_ttl: timedelta = timedelta(seconds=900)
) -> WorkflowStreamState:
Expand Down
9 changes: 6 additions & 3 deletions temporalio/streams/_provider.py
Original file line number Diff line number Diff line change
Expand Up @@ -322,10 +322,13 @@ def open_writer(self, topic: str) -> WriteSink:
...

def on_workflow_start(self) -> None:
"""Called before the workflow function runs.
"""Called before the workflow function runs, and before the first task's handlers.

A provider that serves outside readers through handlers on the
workflow registers them here, before the first task completes.
After the workflow's own ``__init__`` and before any Signal or Update
of the first task is handled, which the SDK does ahead of the
workflow function. A provider that serves outside readers through
handlers on the workflow registers them here, so an Update that
arrives with the first task finds them.
"""
...

Expand Down
Loading
Loading