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 @@ -46,6 +46,26 @@ to include examples, links to docs, or any other relevant information.
`workflow.linked_channel(name)` and ends a subscription with `unsubscribe()`. Notifications
arrive with the workflow's tasks. `WorkflowExecutionDescription.channel_subscriptions` lists
the channels a run listens on.
- **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
2 changes: 2 additions & 0 deletions temporalio/client/__init__.py
Original file line number Diff line number Diff line change
Expand Up @@ -70,6 +70,7 @@
ChannelKind,
ChannelListener,
ChannelSubscriptionInfo,
stream_channel,
)
from ._client import (
Client,
Expand Down Expand Up @@ -373,6 +374,7 @@
"ChannelListener",
"ChannelAddress",
"ChannelSubscriptionInfo",
"stream_channel",
"_ClientImpl",
"_apply_headers",
"_decode_user_metadata",
Expand Down
44 changes: 44 additions & 0 deletions temporalio/client/_channel.py
Original file line number Diff line number Diff line change
Expand Up @@ -10,6 +10,7 @@
import temporalio.api.notification.v1
import temporalio.api.workflow.v1
import temporalio.common
from temporalio.streams._ref import StreamRef
from temporalio.workflow import Notification

from ._callback import Callback
Expand All @@ -20,8 +21,12 @@
"ChannelKind",
"ChannelListener",
"ChannelSubscriptionInfo",
"stream_channel",
]

STREAM_CHANNEL_PREFIX = "stream/"
"""The first segment of the channel a native stream notifies."""


class ChannelKind(IntEnum):
"""Where a channel lives, which decides how a call addresses it.
Expand Down Expand Up @@ -226,3 +231,42 @@ def workflow_id(self) -> str | None:
):
return self.execution.business_id
return None


def stream_channel(ref: StreamRef) -> ChannelAddress:
"""The channel a native stream notifies on every append and on its close.

The server derives the name from the stream's identity, and this helper
derives the same one, so a client polls or registers a callback without
asking. A stream a workflow owns notifies ``stream/<topic>`` linked to the
owning workflow. A stream an activity owns notifies
``stream/<activity id>/<topic>`` linked to the workflow that scheduled
the activity, or ``stream/<topic>`` linked to the activity execution
itself when the activity is a standalone one. A standalone stream notifies
the independent channel ``stream/<stream id>``, whatever the topic, since
its topics share one stream on the server. The address names the owner
without a run, so it reaches the owner's current run.

Each change arrives as one notification: the stream's change sequence as
the counter, the head after the change as the position, and ``closed``
set in the metadata on the close.
"""
if ref.kind == "workflow":
assert ref.workflow_id is not None
return ChannelAddress(
STREAM_CHANNEL_PREFIX + ref.topic,
temporalio.common.Execution.workflow(ref.workflow_id),
)
if ref.kind == "activity":
assert ref.activity_id is not None
if not ref.workflow_id:
return ChannelAddress(
STREAM_CHANNEL_PREFIX + ref.topic,
temporalio.common.Execution.activity(ref.activity_id),
)
return ChannelAddress(
f"{STREAM_CHANNEL_PREFIX}{ref.activity_id}/{ref.topic}",
temporalio.common.Execution.workflow(ref.workflow_id),
)
assert ref.stream_id is not None
return ChannelAddress(STREAM_CHANNEL_PREFIX + ref.stream_id, None)
Loading
Loading