Skip to content

Added the notification channel: listeners woken on their scheduled task. - #9

Closed
moetemp wants to merge 18 commits into
moe/AI-198-srv-6-reset-failoverfrom
moe/AI-198-srv-7-notification-channel
Closed

moetemp wants to merge 18 commits into
moe/AI-198-srv-6-reset-failoverfrom
moe/AI-198-srv-7-notification-channel

Conversation

@moetemp

@moetemp moetemp commented Oct 2, 2026 •

Copy link
Copy Markdown
Owner

This PR adds the notification channel, a pub/sub point a writer notifies by name while the server wakes every listener.

What changed?

  • chasm/lib/channel holds a channel as an execution of its own, keyed by the channel name in the namespace. It keeps the listener table, a ring of recent notifications for pollers (channel.retainedNotifications) and the latest notification by counter. A notify writes the state and one retained notification, not a node per listener.
  • NotifyChannel to a channel with no listeners, or to one that does not exist yet, is accepted: it creates the channel, retains the notification and returns listener_count 0. It is never an error, so a poller that arrives later catches up through after_counter.
  • A listener that registers is handed the channel's latest notification once. A workflow takes it as pending, so it rides the next scheduled event; a callback is posted it. This closes the race where a reader checks an external store, finds nothing and subscribes while a writer appends and notifies before the subscription reaches the server.
  • One fan-out task is outstanding per channel, coalesced by a flag the way StreamNotifyConsumersTask is. It hands the latest notification to every listener: each workflow listener through DeliverChannelNotification on its own shard, each callback listener inside the channel's transition.
  • A workflow listener whose run on record has closed is followed to the current run. If that run subscribed, the listener is re-keyed to it; otherwise the channel drops it. A redelivered counter is a duplicate and changes nothing.
  • SubscribeNotificationChannel writes WorkflowNotificationChannelSubscribed, records the subscription on the Workflow component and registers the run on the channel's shard before the commit. A refusal fails the task with BAD_SUBSCRIBE_NOTIFICATION_CHANNEL_ATTRIBUTES. A subscribe to a channel the run already listens to records its own event, since every command needs one, and changes nothing else: no new listener and no delivery. The subscription lives for the run: a continue-as-new successor subscribes again, while a reset run keeps it from the copied event and the next notify re-keys the listener to it.
  • A run keeps the latest notification per channel that no scheduled event has carried, and a pending one schedules a Workflow Task through the transaction close. The scheduled event of a task no worker has seen takes them into WorkflowTaskScheduledEventAttributes.notifications and clears them, so History is the acknowledgment. A retry's scheduled event, written for an attempt a worker already ran, carries none, so a notification is never carried twice. A speculative task that started carries none either, and query tasks carry nothing.
  • Callback listeners take URL callbacks. A channel task on the outbound queue posts the notification as JSON with the channel in Temporal-Notification-Channel, through the callback library's HTTP caller, retry policy and per-destination circuit breaker. One post is in flight per listener; what arrives meanwhile folds into one pending notification, sent when the post completes.
  • PollChannel long-polls on the channel and returns retained notifications above after_counter. DescribeChannel returns the listeners, the latest notification and the retained count. RegisterChannelListener is idempotent by request id, and UnregisterChannelListener of a gone listener succeeds.
  • The frontend serves the five RPCs with Signal-style validation and replaces the stubs from Pinned api-go to the branch carrying the stream protos. #3. A channel with no listeners and no activity for channel.retention is closed and deleted by an idle task. Limits are channel.maxListeners, channel.maxMetadataBytes, channel.notifyPerSecond and channel.maxSubscriptionsPerWorkflow, and the counters are channel_notifications_accepted, channel_notifications_delivered (tagged by listener kind), channel_notifications_folded and channel_pollers.
  • Writes per notification: NotifyChannel reads the channel first. A counter above the latest writes the channel state and one retained notification, never a node per listener, and the coalesced fan-out task writes the channel once per batch. A counter at or below the latest is not retained and still reaches every listener, folded per listener: it writes the channel only when a callback listener does not hold it, in flight or pending, and reaches each workflow listener directly through DeliverChannelNotification. That call reads the run first. A run that holds the channel's notification pending at that counter or above, or whose scheduled task has not started and carries it, folds it with no write. A run whose task carrying it has started gets a new pending entry and a task, since that task may have read before the write, so a watcher that sends the same counter again after its reader's task ran still wakes it. A callback attempt writes its listener node once. PollChannel and DescribeChannel are reads, except that a poll records itself as activity at most once per half retention.
  • Pending notifications keep a table of their own on the Workflow component, beside the stream cursors, since the scheduled event that carries them is their acknowledgment.
  • A channel that does not exist is created in a transition of its own before it is updated. The update half of an update-with-start runs after the new execution's structure is synced, so a listener added there was dropped.

Part of AI-198 (epic AI-37).

Why?

Waking one workflow by name adds little over a Signal: the writer has to know each consumer. A channel lets a writer announce that a source moved without knowing who listens, and the server wakes every listener. Notifications a workflow sees are on its scheduled event, so they are in History, replay sees the same, and workflow code may act on them.

How did you test it?

go build ./... is clean, and golangci-lint and the errortype vet are clean on the touched packages. Unit tests ran for ./chasm/lib/channel/..., ./chasm/lib/workflow/..., ./service/history/..., ./service/frontend/... and the touched ./common/... packages, covering the ring, folding, the listener table, re-keying, the idle check and the scheduled-event snapshot. The functional tests ran with -tags test_dep: subscribe and notify, folding while a task is open, several listeners, continue-as-new, re-keying to a reset run, a callback listener that folds while busy, a waiting poll, retention with no listeners, validation and limits, the notify rate, a retry that does not carry a notification again, a late subscriber handed the latest notification, a subscribe on an empty channel, a repeated subscribe that records its event and adds no listener, and the same counter notified again after the reader's task ran, which schedules a new task carrying it, while a repeat that is still pending, or that reaches a scheduled task before it starts, writes neither the run nor the channel. The stream functional suite passes.

  • Unit Tests
  • Staging
  • End to End Tests

moetemp added 12 commits October 1, 2026 17:43
A channel is an execution keyed by its name. It keeps the listener table and a ring of recent notifications, and one coalesced fan-out task hands the latest to each workflow listener's shard and to each callback listener.
…low.

A run keeps the channels it subscribed to and the latest notification per channel that no scheduled event has carried, and a pending one schedules a Workflow Task through the transaction close.
The update half of an update-with-start runs after the new execution's structure is synced, so a listener added there was dropped. The poll wait is a duration, as in the public request.
A reader that found its source empty and subscribed would otherwise miss a write that landed in between, since the notify that followed found nobody to tell.
…ents.

The scheduled event of a task no worker has seen takes the pending notifications and clears them, so History is the acknowledgment and replay sees what the workflow saw.
Every command needs its event, or the SDK's command matching desyncs on replay. The repeat changes nothing else.
A notify whose counter is not above the channel's latest is read and answered, and a redelivery a run already has is folded without touching the run.
A watcher that finds a record after its reader's task ran sends the same counter again, and folding that against the channel's latest left the reader waiting. Only a listener that holds it pending folds it.
A task that has not read anything yet will see the write a repeat stands for, so the producer and its watcher notifying the same counter cost one task. Once the task starts, a repeat goes to the next one.
…notification-channel

# Conflicts:
#	chasm/lib/workflow/gen/workflowpb/v1/state.go-helpers.pb.go
#	chasm/lib/workflow/gen/workflowpb/v1/state.pb.go
#	chasm/lib/workflow/proto/v1/state.proto
#	chasm/lib/workflow/workflow.go
#	common/metrics/metric_defs.go
#	service/frontend/errors.go
#	service/history/workflow/mutable_state_impl.go
…notification-channel

# Conflicts:
#	service/frontend/workflow_handler.go
Timers with equal deadlines are grouped into one state machine timer, and the wall clock on macOS does not always tick between two nodes of one refresh walk, so the task refresher test counted two timer groups where it expected three.
@moetemp

moetemp commented Oct 3, 2026

Copy link
Copy Markdown
Owner Author

Replaced by #18, #19, #20, #21, #22, #23, #24, #25, #26, #27, #28, #29, #30, #31, #32, #33, #34, #35, #36 and #37.

Same content, split into 20 PRs in the v3 series: the notification channel first, then the streaming interface, then native streams and the rest. The branch stays as a pin.

@moetemp moetemp closed this Oct 3, 2026
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

1 participant