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
20 changes: 20 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -37,6 +37,26 @@ to include examples, links to docs, or any other relevant information.
worker-side factories registered with `StrandsPlugin(sandboxes=...)`.

- Added the `temporalio.contrib.gcp.cloud_run.id` module with the `CloudRunIdPlugin` client plugin to set the worker identity on Cloud Run.
- **Experimental**: `temporalio.streams` defines one stream interface a workflow
can read, decide on, and write. A provider is registered once as a plugin,
`Client.connect(plugins=[provider])`, and workers built from that client
inherit it; each context then asks for its stream the same way:
`workflow.stream_reader()` and `workflow.stream_writer()` in workflow code,
`activity.stream_handle()` in an activity, and `client.get_stream_handle()`
anywhere a client is held. A topic is a typed definition,
`streams.topic("inputs", Token)`, shared by workflow, activity and client
code; a plain string names a topic decided at runtime, and a call that names
no topic addresses the default topic, `streams.DEFAULT_TOPIC` (`"output"`,
the server's default stream name). The record on the wire
is `temporal.api.stream.v1.StreamRecord` on every provider. A stream is
handed to another process as a `streams.StreamRef`, plain data naming the
owner and the topic, which `client.get_stream_handle(ref)` and
`activity.stream_handle(ref)` open; `client.create_stream(stream_id, ...)`
creates a standalone stream with a retention policy, and its handle's
`close()` seals it. A provider runs record bodies through the client's data
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.

### Changed

Expand Down
128 changes: 128 additions & 0 deletions temporalio/activity.py
Original file line number Diff line number Diff line change
Expand Up @@ -20,6 +20,7 @@
from typing import (
TYPE_CHECKING,
Any,
Literal,
NoReturn,
overload,
)
Expand All @@ -29,9 +30,11 @@
import temporalio.bridge.proto.activity_task
import temporalio.common
import temporalio.converter
import temporalio.streams
from temporalio.converter._payload_converter import (
_TemporalTransferTypePayloadConverter,
)
from temporalio.streams._ref import open_ref

from .types import CallableType

Expand Down Expand Up @@ -209,6 +212,10 @@ class _Context:
runtime_metric_meter: temporalio.common.MetricMeter | None
client: Client | None
cancellation_details: _ActivityCancellationDetailsHolder
stream_provider: temporalio.streams.StreamProvider | None = None
# A ``def`` activity is handed no client and no stream provider, so what
# it is missing cannot be read off the fields that are absent.
sync: bool = False
_logger_details: Mapping[str, Any] | None = None
_payload_converter: temporalio.converter.PayloadConverter | None = None
_metric_meter: temporalio.common.MetricMeter | None = None
Expand Down Expand Up @@ -294,6 +301,127 @@ def client() -> Client:
return client


def stream_handle(
workflow_id: str | temporalio.streams.StreamRef | None = None,
*,
run_id: str | None = None,
scope: Literal["workflow", "activity"] | None = None,
) -> temporalio.streams.StreamHandle:
"""Return a stream handle from the provider the worker was given.

A :py:class:`temporalio.streams.StreamRef` in place of ``workflow_id``,
such as one this activity received as an argument, opens the stream it
names, whatever owns it, and takes no other argument; the handle's calls
that name no topic then address the ref's topic.

Which stream a call with no ``workflow_id`` reaches is decided by where
the activity runs, never by what exists:

- In an activity a workflow scheduled, it is that workflow's stream,
pinned to the run the activity belongs to.
- In a standalone activity, it is the activity's own stream.
- ``scope="activity"`` gives an activity a workflow scheduled its own
streams instead, apart from the workflow's. ``scope="workflow"`` asks
for the workflow explicitly, and a standalone activity has none.

The rule is static because a stream is created by its first write, so a
rule that looked for one would send attempt 1 to the workflow and a retry
to the stream attempt 1 created. An activity's own streams are one per
activity execution, not per attempt: a retry writes to the same stream,
under a new attempt, and they end when the activity reaches a terminal
status.

Name a ``workflow_id`` to address another workflow; ``run_id`` then pins
the handle to one run and its absence follows the execution chain. A
``read``, ``latest`` or ``producer`` that names no topic addresses the
owner's default topic, :py:data:`temporalio.streams.DEFAULT_TOPIC`, the
one :py:func:`temporalio.workflow.stream_reader` and
:py:func:`temporalio.workflow.stream_writer` use without a topic. A
``read`` starts at :py:data:`temporalio.streams.BEGINNING`, at
:py:data:`temporalio.streams.END` or at the last ``N`` records with
``last=N``. See :py:mod:`temporalio.streams`.

Like :py:func:`client`, this is only available in ``async def``
activities.

Args:
workflow_id: Another workflow whose stream to address, or a
:py:class:`temporalio.streams.StreamRef` naming the stream.
run_id: The run of ``workflow_id`` to pin to.
scope: ``"activity"`` for this activity's own streams,
``"workflow"`` for its workflow's. Without it the rule above
decides.

Returns:
:py:class:`temporalio.streams.StreamHandle` for use in the current
activity.

Raises:
RuntimeError: When the client is not available, which is what a
``def`` activity gets, or when ``scope="workflow"`` is asked of
an activity that belongs to no workflow.
temporalio.streams.StreamUnsupportedError: The worker has no stream
provider, or its provider cannot hold a stream an activity owns.
Register one with ``Client.connect(plugins=[provider])`` or
``Worker(plugins=[provider])``.
ValueError: ``run_id`` was given without ``workflow_id``,
``scope="activity"`` with one, or a ref with either.
"""
context = _Context.current()
if context.sync:
# A sync activity is handed neither a client nor a provider. Saying
# the worker has none would send the reader to fix a registration
# that is not the problem.
raise RuntimeError(
"No stream handle available. Stream handles are only available in "
"`async def` activities; not in `def` activities, which are handed no "
"client to reach the store with."
)
provider = context.stream_provider
if provider is None:
raise temporalio.streams.StreamUnsupportedError(
"no stream provider is configured on this worker; register one with "
"Client.connect(plugins=[provider]) or Worker(plugins=[provider])"
)
if isinstance(workflow_id, temporalio.streams.StreamRef):
if run_id is not None or scope is not None:
raise ValueError(
"a StreamRef names the stream in full, so it takes no run_id or scope"
)
return open_ref(provider, client(), workflow_id)
if workflow_id is not None:
if scope == "activity":
raise ValueError(
"scope='activity' addresses this activity's own streams, so it takes "
"no workflow_id"
)
return provider.get_stream_handle(client(), workflow_id, run_id=run_id)
if run_id is not None:
raise ValueError("run_id needs a workflow_id")
info = context.info()
if scope is None:
scope = "workflow" if info.in_workflow else "activity"
if scope == "workflow":
if info.workflow_id is None:
raise RuntimeError(
"this activity belongs to no workflow, so name the workflow_id to "
"address, or leave scope unset for the activity's own streams"
)
return provider.get_stream_handle(
client(), info.workflow_id, run_id=info.workflow_run_id
)
if info.workflow_id is not None:
return provider.get_activity_stream_handle(
client(),
info.activity_id,
workflow_id=info.workflow_id,
run_id=info.workflow_run_id,
)
return provider.get_activity_stream_handle(
client(), info.activity_id, run_id=info.activity_run_id
)


def in_activity() -> bool:
"""Whether the current code is inside an activity.

Expand Down
142 changes: 142 additions & 0 deletions temporalio/client/_client.py
Original file line number Diff line number Diff line change
Expand Up @@ -40,6 +40,7 @@
ServiceClient,
TLSConfig,
)
from temporalio.streams._ref import open_ref

from ..common import HeaderCodecBehavior
from ..types import (
Expand Down Expand Up @@ -156,6 +157,7 @@ async def connect(
grpc_compression: GrpcCompression = GrpcCompression.GZIP,
payload_limits: PayloadLimitsConfig = PayloadLimitsConfig(),
header_codec_behavior: HeaderCodecBehavior = HeaderCodecBehavior.NO_CODEC,
stream_provider: temporalio.streams.StreamProvider | None = None,
) -> Self:
"""Connect to a Temporal server.

Expand Down Expand Up @@ -221,6 +223,11 @@ async def connect(
payload_limits: Warning thresholds for outbound payload/memo sizes. Over-threshold
fields are logged but still sent. Set a threshold to 0 to disable it.
header_codec_behavior: Encoding behavior for headers sent by the client.
stream_provider: Experimental. The stream provider
:py:meth:`get_stream_handle` opens handles from, see
:py:mod:`temporalio.streams`. A provider that is also a
:py:class:`Plugin` sets this itself when passed in ``plugins``,
and workers built from this client inherit it.
"""
connect_config = temporalio.service.ConnectConfig(
target_host=target_host,
Expand Down Expand Up @@ -257,6 +264,7 @@ def make_lambda(
default_workflow_query_reject_condition=default_workflow_query_reject_condition,
header_codec_behavior=header_codec_behavior,
plugins=plugins,
stream_provider=stream_provider,
)

def __init__(
Expand Down Expand Up @@ -889,6 +897,139 @@ def get_workflow_handle(
result_type=result_type,
)

def get_stream_handle(
self,
workflow_id: str | temporalio.streams.StreamRef | None = None,
*,
run_id: str | None = None,
activity_id: str | None = None,
stream_id: str | None = None,
) -> temporalio.streams.StreamHandle:
"""Get a handle on a stream from the provider registered on this client.

Mirrors :py:meth:`get_workflow_handle`: without ``run_id`` the handle
follows the workflow's execution chain across continue-as-new, with
one it is pinned to that run. With ``activity_id`` the handle is on
the streams that activity owns: a standalone activity's when
``workflow_id`` is left out, and ``run_id`` then pins the activity's
run, or an activity that ``workflow_id`` scheduled. With
``stream_id`` it is on a standalone stream, one with an id of its own
and no owner, which :py:meth:`create_stream` made; it takes no other
argument. A :py:class:`temporalio.streams.StreamRef` in place of
``workflow_id`` opens the stream the ref names, whatever owns it, and
takes no other argument either: the handle's calls that name no topic
then address the ref's topic. The provider is the one registered
with ``plugins=[provider]`` at :py:meth:`connect`, or passed as
``stream_provider``. The handle's ``read``, ``latest`` and
``producer`` take a topic, and without one address the owner's
default topic, :py:data:`temporalio.streams.DEFAULT_TOPIC`. A
``read`` starts at :py:data:`temporalio.streams.BEGINNING`, at
:py:data:`temporalio.streams.END` or at the last ``N`` records with
``last=N``, and resumes only after a cursor it was handed. See
:py:mod:`temporalio.streams`.

Args:
workflow_id: Workflow ID whose stream to get a handle to, the
workflow that scheduled ``activity_id``, or a
:py:class:`temporalio.streams.StreamRef` naming the stream.
run_id: Run ID to pin the handle to.
activity_id: Activity ID whose own streams to get a handle to.
stream_id: ID of the standalone stream to get a handle to.

Returns:
The stream handle.

Raises:
ValueError: No owner was named, or a ref or ``stream_id`` was
given together with another argument.
temporalio.streams.StreamUnsupportedError: No stream provider is
registered on this client, or it cannot hold a stream an
activity owns or a stream without an owner.
"""
provider = self._config.get("stream_provider")
if provider is None:
raise temporalio.streams.StreamUnsupportedError(
"no stream provider is registered on this client; connect with "
"plugins=[provider]"
)
if isinstance(workflow_id, temporalio.streams.StreamRef):
if run_id is not None or activity_id is not None or stream_id is not None:
raise ValueError(
"a StreamRef names the stream in full, so it takes no run_id, "
"activity_id or stream_id"
)
return open_ref(provider, self, workflow_id)
if stream_id is not None:
if workflow_id is not None or run_id is not None or activity_id is not None:
raise ValueError(
"stream_id names a standalone stream, which has no workflow_id, "
"run_id or activity_id"
)
return provider.get_standalone_stream_handle(self, stream_id)
if activity_id is not None:
return provider.get_activity_stream_handle(
self, activity_id, workflow_id=workflow_id, run_id=run_id
)
if workflow_id is None:
raise ValueError(
"name the workflow_id, the activity_id or the stream_id to address"
)
return provider.get_stream_handle(self, workflow_id, run_id=run_id)

async def create_stream(
self,
stream_id: str,
*,
retention: timedelta | None = None,
max_records: int | None = None,
max_bytes: int | None = None,
) -> temporalio.streams.StreamHandle:
"""Create a standalone stream and get a handle on it.

A standalone stream has an id of its own and no owner, so it is
created here on purpose rather than by its first write, and it is
sealed on purpose with the handle's ``close()``, after which appends
are refused and the retained records stay readable. The three policy
arguments bound what it retains: records older than ``retention``,
beyond the newest ``max_records`` or past ``max_bytes`` of stored
records are dropped, and ``None`` leaves a bound to the provider's
default. Creating a stream that exists with the same policy returns a
handle on it, so a retried create is harmless. Another process
reaches the stream with ``get_stream_handle(stream_id=...)`` or with
the handle's ``ref()``.

Args:
stream_id: ID of the stream to create.
retention: How long a record is kept.
max_records: How many of the newest records are kept.
max_bytes: How many bytes of records are kept.

Returns:
A handle on the new or existing stream.

Raises:
ValueError: ``stream_id`` is empty, a bound is not positive, or
the stream exists with a different policy.
temporalio.streams.StreamUnsupportedError: No stream provider is
registered on this client, or it cannot hold a stream without
an owner.
"""
provider = self._config.get("stream_provider")
if provider is None:
raise temporalio.streams.StreamUnsupportedError(
"no stream provider is registered on this client; connect with "
"plugins=[provider]"
)
if not stream_id:
raise ValueError("stream_id must not be empty")
return await provider.create_standalone_stream(
self,
stream_id,
retention=retention,
max_records=max_records,
max_bytes=max_bytes,
)

def get_workflow_handle_for(
self,
workflow: (
Expand Down Expand Up @@ -3021,6 +3162,7 @@ class ClientConnectConfig(TypedDict, total=False):
grpc_compression: GrpcCompression
payload_limits: PayloadLimitsConfig
header_codec_behavior: HeaderCodecBehavior
stream_provider: temporalio.streams.StreamProvider | None


class ClientConfig(TypedDict, total=False):
Expand Down
15 changes: 9 additions & 6 deletions temporalio/streams/__init__.py
Original file line number Diff line number Diff line change
Expand Up @@ -43,12 +43,15 @@
``Client.connect(plugins=[provider])``; workers built from that client inherit
it, and ``Worker(plugins=[provider])`` or ``Replayer(plugins=[provider])``
registers it on a worker alone. Each context then asks for its stream the
same way. Workflow code uses ``temporalio.workflow.stream_reader`` and
``temporalio.workflow.stream_writer``. An activity uses
``temporalio.activity.stream_handle``, which is its own workflow pinned
to its run unless told otherwise. Any process holding a client uses
``temporalio.client.Client.get_stream_handle``, which mirrors
``get_workflow_handle``. The explicit form,
same way. Workflow code uses :func:`temporalio.workflow.stream_reader` and
:func:`temporalio.workflow.stream_writer`. An activity uses
:func:`temporalio.activity.stream_handle`: an activity a workflow scheduled
reaches that workflow's stream pinned to its run, a standalone activity
reaches its own, and ``scope="activity"`` gives the first kind its own
streams too. Any process holding a client uses
:meth:`temporalio.client.Client.get_stream_handle`, which mirrors
``get_workflow_handle`` and takes an ``activity_id`` for an activity's
streams. The explicit form,
``provider.get_stream_handle(client, workflow_id)``, stays for a process that
talks to two stores. This module keeps the shared types, the errors and the
protocols a provider implements; nothing here that workflow code imports does
Expand Down
Loading
Loading