Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
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 @@ -66,6 +66,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
Loading
Loading