Repository navigation
Conversation
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.
…notification-channel
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.
This was referenced Oct 2, 2026
…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.
2 of 3 tasks
…notification-channel
…notification-channel
…notification-channel
Owner
Author
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Sign up for free
to join this conversation on GitHub.
Already have an account?
Sign in to comment
Add this suggestion to a batch that can be applied as a single commit.This suggestion is invalid because no changes were made to the code.Suggestions cannot be applied while the pull request is closed.Suggestions cannot be applied while viewing a subset of changes.Only one suggestion per line can be applied in a batch.Add this suggestion to a batch that can be applied as a single commit.Applying suggestions on deleted lines is not supported.You must change the existing code in this line in order to create a valid suggestion.Outdated suggestions cannot be applied.This suggestion has been applied or marked resolved.Suggestions cannot be applied from pending reviews.Suggestions cannot be applied on multi-line comments.Suggestions cannot be applied while the pull request is queued to merge.Suggestion cannot be applied right now. Please check back later.
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/channelholds 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.NotifyChannelto a channel with no listeners, or to one that does not exist yet, is accepted: it creates the channel, retains the notification and returnslistener_count0. It is never an error, so a poller that arrives later catches up throughafter_counter.StreamNotifyConsumersTaskis. It hands the latest notification to every listener: each workflow listener throughDeliverChannelNotificationon its own shard, each callback listener inside the channel's transition.SubscribeNotificationChannelwritesWorkflowNotificationChannelSubscribed, records the subscription on the Workflow component and registers the run on the channel's shard before the commit. A refusal fails the task withBAD_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.WorkflowTaskScheduledEventAttributes.notificationsand 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.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.PollChannellong-polls on the channel and returns retained notifications aboveafter_counter.DescribeChannelreturns the listeners, the latest notification and the retained count.RegisterChannelListeneris idempotent by request id, andUnregisterChannelListenerof a gone listener succeeds.channel.retentionis closed and deleted by an idle task. Limits arechannel.maxListeners,channel.maxMetadataBytes,channel.notifyPerSecondandchannel.maxSubscriptionsPerWorkflow, and the counters arechannel_notifications_accepted,channel_notifications_delivered(tagged by listener kind),channel_notifications_foldedandchannel_pollers.NotifyChannelreads 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 throughDeliverChannelNotification. 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.PollChannelandDescribeChannelare reads, except that a poll records itself as activity at most once per half retention.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 theerrortypevet 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.