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
4 changes: 4 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -42,6 +42,10 @@ to include examples, links to docs, or any other relevant information.
`unregister_channel_listener` reach a named channel on the server. The channel is either
independent or linked to an execution, which `execution=` (`temporalio.common.Execution`) or
the `workflow_id=` shorthand names.
- A workflow subscribes with `workflow.subscribe_channel(name)`, reads a channel linked to it with
`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.

### Changed

Expand Down
2 changes: 2 additions & 0 deletions temporalio/client/__init__.py
Original file line number Diff line number Diff line change
Expand Up @@ -69,6 +69,7 @@
ChannelDescription,
ChannelKind,
ChannelListener,
ChannelSubscriptionInfo,
)
from ._client import (
Client,
Expand Down Expand Up @@ -371,6 +372,7 @@
"ChannelKind",
"ChannelListener",
"ChannelAddress",
"ChannelSubscriptionInfo",
"_ClientImpl",
"_apply_headers",
"_decode_user_metadata",
Expand Down
67 changes: 67 additions & 0 deletions temporalio/client/_channel.py
Original file line number Diff line number Diff line change
Expand Up @@ -8,6 +8,7 @@
from enum import IntEnum

import temporalio.api.notification.v1
import temporalio.api.workflow.v1
import temporalio.common
from temporalio.workflow import Notification

Expand All @@ -18,6 +19,7 @@
"ChannelDescription",
"ChannelKind",
"ChannelListener",
"ChannelSubscriptionInfo",
]


Expand Down Expand Up @@ -128,6 +130,71 @@ class ChannelDescription:
"""


@dataclass(frozen=True)
class ChannelSubscriptionInfo:
"""A workflow's standing on one channel, as its description reports it.

An independent channel is listed from the subscribe event until the run
unsubscribes or closes. A linked channel is listed once it holds state. A
closed run keeps listing what it stood on, and a continue-as-new
successor starts with nothing.

.. warning::
This API is experimental and unstable.
"""

channel: str
"""Channel name."""

kind: ChannelKind
""":attr:`ChannelKind.INDEPENDENT` for a subscription the workflow made by
command, :attr:`ChannelKind.LINKED` for a channel linked to it."""

subscribed_event_id: int
"""Id of the event that recorded the subscription. Zero for the linked kind."""

last_counter: int
"""Highest counter the workflow has accepted from the channel. Zero when
none has arrived."""

pending_notification: Notification | None
"""The notification held for the workflow's next Workflow Task, when one
is pending."""

scheduled_counter: int
"""Counter carried by the scheduled event of a Workflow Task that has not
started yet. Zero otherwise."""

listener_count: int
"""Linked kind: callback listeners registered on the channel."""

retained_count: int
"""Linked kind: notifications retained for pollers."""

accepted_count: int
"""Linked kind: notifications the channel has accepted over its life."""

@staticmethod
def _from_proto(
proto: temporalio.api.workflow.v1.ChannelSubscriptionInfo,
) -> ChannelSubscriptionInfo:
return ChannelSubscriptionInfo(
channel=proto.channel,
kind=ChannelKind(proto.kind),
subscribed_event_id=proto.subscribed_event_id,
last_counter=proto.last_counter,
pending_notification=(
Notification._from_proto(proto.pending_notification)
if proto.HasField("pending_notification")
else None
),
scheduled_counter=proto.scheduled_counter,
listener_count=proto.listener_count,
retained_count=proto.retained_count,
accepted_count=proto.accepted_count,
)


@dataclass(frozen=True)
class ChannelAddress:
"""Where a channel call reaches a channel: its name and, when linked, its owner.
Expand Down
17 changes: 17 additions & 0 deletions temporalio/client/_workflow.py
Original file line number Diff line number Diff line change
Expand Up @@ -59,6 +59,7 @@
ReturnType,
SelfType,
)
from ._channel import ChannelSubscriptionInfo
from ._exceptions import (
WorkflowContinuedAsNewError,
WorkflowFailureError,
Expand Down Expand Up @@ -1418,6 +1419,18 @@ class WorkflowExecutionDescription(WorkflowExecution):
raw_description: temporalio.api.workflowservice.v1.DescribeWorkflowExecutionResponse
"""Underlying protobuf description."""

channel_subscriptions: Sequence[ChannelSubscriptionInfo] = ()
"""The notification channels this run stands on.

The independent channels it subscribed to and the channels linked to it
that hold any state, sorted by name with the independent kind first.
Empty when there are none. See
:py:class:`temporalio.client.ChannelSubscriptionInfo`.

.. warning::
This API is experimental and unstable.
"""

_static_summary: str | None = None
_static_details: str | None = None
_metadata_decoded: bool = False
Expand Down Expand Up @@ -1456,6 +1469,10 @@ async def _from_raw_description(
namespace=namespace,
converter=converter,
raw_description=description,
channel_subscriptions=tuple(
ChannelSubscriptionInfo._from_proto(info)
for info in description.channel_subscriptions
),
)


Expand Down
63 changes: 63 additions & 0 deletions temporalio/worker/_workflow_instance.py
Original file line number Diff line number Diff line change
Expand Up @@ -303,6 +303,16 @@ def __init__(self, det: WorkflowInstanceDetails) -> None:
det.worker_level_failure_exception_types
)
self._patch_activation_callback = det.patch_activation_callback
# Keyed by channel name; one subscription per channel per run
self._channel_subscriptions: dict[
str, temporalio.workflow.ChannelSubscription
] = {}
# The channels linked to this workflow, keyed by name as well. No
# command: the owner is the listener by construction, so the map only
# routes a notification carrying ``linked_to`` to its handle.
self._linked_channel_subscriptions: dict[
str, temporalio.workflow.ChannelSubscription
] = {}
self._default_workflow_logic_flags = det.default_workflow_logic_flags
self._primary_task: asyncio.Task[None] | None = None
self._cancel_primary_task_pending = False
Expand Down Expand Up @@ -624,6 +634,8 @@ def _apply(
self._apply_query_workflow(job.query_workflow)
elif job.HasField("notify_has_patch"):
self._apply_notify_has_patch(job.notify_has_patch)
elif job.HasField("notifications_received"):
self._apply_notifications_received(job.notifications_received)
elif job.HasField("remove_from_cache"):
self._apply_remove_from_cache(job.remove_from_cache)
elif job.HasField("resolve_activity"):
Expand Down Expand Up @@ -865,6 +877,28 @@ async def run_query() -> None:
# Schedule it
self.create_task(run_query(), name=f"query: {job.query_type}")

def _apply_notifications_received(
self, job: temporalio.bridge.proto.workflow_activation.NotificationsReceived
) -> None:
for proto in job.notifications:
# A name may be open as both kinds; the kind the server stamped
# on the notification picks the handle.
if proto.HasField("linked_to"):
subscription = self._linked_channel_subscriptions.get(proto.channel)
else:
subscription = self._channel_subscriptions.get(proto.channel)
if subscription is None:
# The server fans out to whatever listened at the time, so a
# channel this run never asked for is not the workflow's
# concern.
logger.debug(
"Dropping a notification on channel %r, which this run does "
"not listen on",
proto.channel,
)
continue
subscription._deliver(temporalio.workflow.Notification._from_proto(proto))

def _apply_notify_has_patch(
self, job: temporalio.bridge.proto.workflow_activation.NotifyHasPatch
) -> None:
Expand Down Expand Up @@ -1810,6 +1844,35 @@ async def workflow_start_nexus_operation(
)
)

def workflow_subscribe_channel(
self, channel: str
) -> temporalio.workflow.ChannelSubscription:
existing = self._channel_subscriptions.get(channel)
if existing is not None:
return existing
command = self._add_command()
command.subscribe_notification_channel.channel = channel
subscription = temporalio.workflow.ChannelSubscription(channel)
self._channel_subscriptions[channel] = subscription
return subscription

def workflow_unsubscribe_channel(self, channel: str) -> None:
command = self._add_command()
command.unsubscribe_notification_channel.channel = channel
# Out of the map before the next activation: the server may still hand
# this run a notification it folded onto a task ahead of the command.
self._channel_subscriptions.pop(channel, None)

def workflow_linked_channel(
self, channel: str
) -> temporalio.workflow.ChannelSubscription:
existing = self._linked_channel_subscriptions.get(channel)
if existing is not None:
return existing
subscription = temporalio.workflow.ChannelSubscription(channel, linked=True)
self._linked_channel_subscriptions[channel] = subscription
return subscription

def workflow_time_ns(self) -> int:
return self._time_ns

Expand Down
10 changes: 9 additions & 1 deletion temporalio/workflow/__init__.py
Original file line number Diff line number Diff line change
Expand Up @@ -57,7 +57,12 @@
as_completed,
wait,
)
from ._channels import Notification
from ._channels import (
ChannelSubscription,
Notification,
linked_channel,
subscribe_channel,
)
from ._context import (
Info,
ParentInfo,
Expand Down Expand Up @@ -259,7 +264,10 @@
"SandboxImportNotificationPolicy",
"logger",
"unsafe",
"ChannelSubscription",
"Notification",
"linked_channel",
"subscribe_channel",
"ChildWorkflowCancellationType",
"ChildWorkflowConfig",
"ChildWorkflowHandle",
Expand Down
Loading
Loading