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
5 changes: 5 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -37,6 +37,11 @@ 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**: notification channels. Requires a server that serves notification channels.
- `Client.notify_channel`, `poll_channel`, `describe_channel`, `register_channel_listener` and
`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.

### Changed

Expand Down
20 changes: 20 additions & 0 deletions temporalio/client/__init__.py
Original file line number Diff line number Diff line change
Expand Up @@ -64,6 +64,12 @@
from ._callback import (
Callback,
)
from ._channel import (
ChannelAddress,
ChannelDescription,
ChannelKind,
ChannelListener,
)
from ._client import (
Client,
ClientConfig,
Expand Down Expand Up @@ -106,6 +112,7 @@
CreateScheduleInput,
DeleteScheduleInput,
DescribeActivityInput,
DescribeChannelInput,
DescribeNexusOperationInput,
DescribeScheduleInput,
DescribeWorkflowInput,
Expand All @@ -120,10 +127,13 @@
ListNexusOperationsInput,
ListSchedulesInput,
ListWorkflowsInput,
NotifyChannelInput,
OutboundInterceptor,
PauseActivityInput,
PauseScheduleInput,
PollChannelInput,
QueryWorkflowInput,
RegisterChannelListenerInput,
ReportCancellationAsyncActivityInput,
SignalWorkflowInput,
StartActivityInput,
Expand All @@ -137,6 +147,7 @@
TriggerScheduleInput,
UnpauseActivityInput,
UnpauseScheduleInput,
UnregisterChannelListenerInput,
UpdateActivityOptionsInput,
UpdateScheduleInput,
UpdateWithStartStartWorkflowInput,
Expand Down Expand Up @@ -315,6 +326,11 @@
"TerminateNexusOperationInput",
"ListNexusOperationsInput",
"CountNexusOperationsInput",
"DescribeChannelInput",
"NotifyChannelInput",
"PollChannelInput",
"RegisterChannelListenerInput",
"UnregisterChannelListenerInput",
"StartWorkflowUpdateInput",
"UpdateWithStartUpdateWorkflowInput",
"UpdateWithStartStartWorkflowInput",
Expand Down Expand Up @@ -351,6 +367,10 @@
"CloudOperationsClient",
"Plugin",
"Callback",
"ChannelDescription",
"ChannelKind",
"ChannelListener",
"ChannelAddress",
"_ClientImpl",
"_apply_headers",
"_decode_user_metadata",
Expand Down
161 changes: 161 additions & 0 deletions temporalio/client/_channel.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,161 @@
"""Notification channel descriptions as the client reports them."""

from __future__ import annotations

from collections.abc import Sequence
from dataclasses import dataclass
from datetime import datetime, timezone
from enum import IntEnum

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

from ._callback import Callback

__all__ = [
"ChannelAddress",
"ChannelDescription",
"ChannelKind",
"ChannelListener",
]


class ChannelKind(IntEnum):
"""Where a channel lives, which decides how a call addresses it.

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

UNSPECIFIED = int(
temporalio.api.notification.v1.ChannelKind.CHANNEL_KIND_UNSPECIFIED
)
"""The server did not say; an older server answers this."""

INDEPENDENT = int(
temporalio.api.notification.v1.ChannelKind.CHANNEL_KIND_INDEPENDENT
)
"""Its own execution, keyed by namespace and channel name.

Any number of workflows subscribe to it and callbacks register on it.
"""

LINKED = int(temporalio.api.notification.v1.ChannelKind.CHANNEL_KIND_LINKED)
"""Kept in one execution's state, keyed by namespace, execution and name.

The owning execution, a workflow or a standalone activity, is its listener
by construction. A call reaches it with the ``execution`` argument, or
with ``workflow_id`` when the owner is a workflow.
"""


@dataclass(frozen=True)
class ChannelListener:
"""One listener of a channel: a workflow or a callback.

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

listener_id: str
"""Assigned by the server when the listener registered."""

workflow_id: str | None
"""The subscribed workflow, when the listener is one."""

run_id: str | None
"""The run that subscribed. Delivery follows the chain's current run."""

callback: Callback | None
"""The callback the server invokes, when the listener is one."""

registered_time: datetime | None
"""When the listener registered."""

@staticmethod
def _from_proto(
proto: temporalio.api.notification.v1.ChannelListener,
) -> ChannelListener:
callback: Callback | None = None
if proto.HasField("callback") and proto.callback.HasField("nexus"):
callback = Callback(
url=proto.callback.nexus.url, headers=dict(proto.callback.nexus.header)
)
workflow = proto.workflow if proto.HasField("workflow") else None
return ChannelListener(
listener_id=proto.listener_id,
workflow_id=workflow.workflow_id if workflow else None,
run_id=workflow.run_id if workflow else None,
callback=callback,
registered_time=(
proto.registered_time.ToDatetime(tzinfo=timezone.utc)
if proto.HasField("registered_time")
else None
),
)


@dataclass(frozen=True)
class ChannelDescription:
"""What the server knows about a channel.

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

listeners: Sequence[ChannelListener]
"""Who is listening, workflows and callbacks alike."""

latest: Notification | None
"""The notification with the highest counter the channel retains."""

retained_count: int
"""How many notifications the channel keeps for pollers."""

kind: ChannelKind = ChannelKind.UNSPECIFIED
"""Which kind of channel this is.

A linked channel of a running execution exists by construction, so a
describe with ``execution`` answers :attr:`ChannelKind.LINKED` with no
listeners and nothing retained for a name nobody has notified yet.
"""

linked_to: temporalio.common.Execution | None = None
"""The owner of a linked channel and the run that holds it.

``None`` for an independent channel.
"""


@dataclass(frozen=True)
class ChannelAddress:
"""Where a channel call reaches a channel: its name and, when linked, its owner.

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

channel: str
"""The channel name."""

execution: temporalio.common.Execution | None
"""The execution the channel is linked to, or ``None`` for an independent one.

Pass both to :py:meth:`temporalio.client.Client.poll_channel` and the
other channel calls as ``channel`` and ``execution``.
"""

@property
def workflow_id(self) -> str | None:
"""The owning workflow's id, when the owner is a workflow.

``None`` for an independent channel and for one a standalone activity
owns, which only ``execution`` reaches.
"""
if (
self.execution is not None
and self.execution.type == temporalio.common.ExecutionType.WORKFLOW
):
return self.execution.business_id
return None
Loading
Loading