Skip to content

Added the stream records, commands, events, and task carriage. - #1

Closed
moetemp wants to merge 30 commits into
mainfrom
moe/AI-198-stream-protos
Closed

moetemp wants to merge 30 commits into
mainfrom
moe/AI-198-stream-protos

Conversation

@moetemp

@moetemp moetemp commented Sep 14, 2026 •

Copy link
Copy Markdown
Owner

This PR adds the public API for server-side streams and the notification channel that wakes their listeners.

What changed?

  • temporal/api/stream/v1/message.proto holds StreamRecord, the one wire form every store keeps. It carries kind, producer_id, attempt and sequence next to body, metadata and topic. An unset kind reads as data. An empty producer_id means the owning Workflow wrote the record.
  • command/v1/message.proto adds AppendStreamRecords and SubscribeStream. The append addresses stream_name. The subscribe takes stream_name_or_id and a start_position (StreamStartPosition: offset, last_n, earliest or tail).
  • history/v1/message.proto adds WorkflowStreamRecordsAppended, bounded by from_offset and to_offset. WorkflowTaskCompletedEventAttributes.consumed_stream_ranges records what a task read as StreamRange values. reserved 14 to 19 keeps numbers for the main line. PollWorkflowTaskQueueResponse.stream_slices carries the records a task reads, and reserved 21; on that response keeps a number the series used for one round and dropped.
  • failed_cause.proto adds five WorkflowTaskFailedCause values: a refused append, a refused stream subscribe, a task whose recorded range the server can't serve any more, a refused channel subscription and a refused channel unsubscription.
  • The notification channel adds notification/v1 (Notification, ChannelListener, ChannelKind), the NotifyChannel, RegisterChannelListener, UnregisterChannelListener, PollChannel and DescribeChannel RPCs, the SubscribeNotificationChannel command and its subscribed event, and WorkflowTaskScheduledEventAttributes.notifications. A writer names a channel, not a Workflow, and the notifications a woken task carries are recorded on its scheduled event so replay sees them.
  • A channel comes in two kinds. An independent channel is addressed by name alone. A channel linked to an execution is addressed through the optional execution on the five channel requests, a common.v1.Execution with type, business_id and run_id, so a workflow or a standalone activity can own one and listens without a subscribe command. Notification.linked_to names the owner as the same type, and DescribeChannelResponse answers kind and linked_to. The metadata comment states that it is folded state, not a log, and the identity comment states that it is for audit and is not copied into the notification.
  • UnsubscribeNotificationChannel ends a run's subscription to an independent channel before the run closes: the command in command/v1, its WorkflowNotificationChannelUnsubscribed event in history/v1 carrying the channel and the id of the subscribe event it ended, and the CommandType, EventType and WorkflowTaskFailedCause values. A command for a channel the run is not subscribed to still records its event, so replay matches every command.
  • DescribeWorkflowExecutionResponse.channel_subscriptions lists the channels a run stands on as ChannelSubscriptionInfo values in workflow/v1: the independent channels it subscribed to and the channels linked to it that hold state, each with its kind, the subscribe event id, the last accepted counter, the held and scheduled notification counters, and the linked channel's listener, retained and accepted counts.
  • The HTTP bindings put the independent routes under /namespaces/{namespace}/channels/{channel} and the linked routes under /namespaces/{namespace}/workflows/{execution.business_id}/channels/{channel} and /namespaces/{namespace}/activities/{execution.business_id}/channels/{channel}. A binding can't set the execution type, so the path segment implies it. Commands aren't on the HTTP surface. openapi/ is regenerated.
  • The payload codec claim on StreamRecord.body names the paths this API owns. The stream service messages aren't part of this API, since streams have no public service API yet.

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

Why?

A workflow's publish has to commit with its Workflow Task, and its reads have to replay. That needs first-class commands and events, not markers. The channel covers stores outside Temporal, such as Redis, where each write would cost a Signal: one History event and one Action per record. A notification runs a task, writes nothing to History beyond the entry on the task's scheduled event, and folds a burst of writes into one. The writer names the channel and never learns who listens, and a channel linked to the listening workflow costs no subscription at all.

How did you test it?

This repo has no unit tests for the protos. CI runs buf lint, buf breaking against upstream main, api-linter and the codegen checks, and they pass at the head. The additions are additive, so buf breaking stays green. The server, Core and SDK PRs in the series exercise these shapes end to end, but that isn't this PR's own coverage.

  • Unit Tests
  • Staging
  • End to End Tests

@moetemp moetemp closed this Sep 16, 2026
@moetemp moetemp reopened this Sep 16, 2026
A half-open range reads as `from_offset`/`to_offset` in the server's own
protos, and AIP-140 rejects both names.
`make http-api-docs` picks up the stream slices on the poll response and
the two new event types, so the committed specs were stale.
The server fails the task instead of the RPC when a stream command is rejected or a recorded range cannot be served, and needs causes the SDK can show.
The publish command comment claimed no event. The run id, message fields and name precedence had no comments.
The record is the wire format every store keeps and every language reads, so it carries its kind, producer, attempt and sequence itself instead of a private envelope inside the body. The body is the user's payload, so a codec applies to it like to any other payload.
The SDKs call an entry a record and a cursor a single position, and `temporal.api` already uses "message" for protocol envelopes. One noun on the wire keeps History readable next to user code.
@moetemp moetemp changed the title Added the stream messages, commands, events, and task carriage. Added the stream records, commands, events, and task carriage. Sep 21, 2026
The append command only ever names a stream the workflow owns, and the subscribe command takes a name or an id and resolves them in that order. Calling both of them stream_id told a reader the wrong thing.
A comment saying 14 through 19 are left free stops nothing. Reserving them makes a rebase that lands an upstream field there fail loudly instead of colliding at 20 later.
A batch was named by a first offset and a count while StreamRange and StreamSlice bound the same idea with from_offset and to_offset. One vocabulary for one concept. The path-wide linter exemption gives way to the per-field disables the rest of this API uses.
An unset int64 reads as zero, so a record whose producer does not number its records claimed position zero while the comment reserved -1 for that case. Zero is now what it looks like, and the stream's offsets are what order a read.
A codec reaches records on the append command and on the Workflow Task response, both of which this API's payload visitor walks. It does not reach the stream service, whose messages live outside this module.
Proto3 cannot tell an unset start offset from zero, and there was no way to ask for the
earliest record of a truncated stream or the last N records. The position is a new message
at a new field number, so the change is additive.
A store Temporal does not host needs a way to make a consumer run a task without writing an event per append. Wakes fold by source, so a burst of writes costs one task.
A wake a task already received cannot be folded against, since only the receiver knows what it read; folding against it stalls a record the watcher read after the task ran.
A writer notifies a named channel and the server wakes its listeners, so writers never address a consumer. Notifications a woken Workflow Task carries are recorded on its scheduled event, so replay sees them. The wake is marked as superseded.
The notification channel covers what the wake was added for, so the PoC does not need a second path. Field 21 on the poll response stays reserved.
A linked channel lives in one workflow's state and is addressed by workflow id, so the five channel calls take an optional `workflow_execution` and get routes under the workflow's path. `Notification.linked_to` lets a listener that holds both kinds route a notification, and `DescribeChannel` reports the kind and the owner.
The surfaces that show a workflow (the CLI's describe, the UI's workflow page, the SDK description) can list the channels it stands on without one call per channel. `ChannelSubscriptionInfo` covers both kinds: the independent subscriptions a command recorded and the linked channels that hold state.
Until now a subscription ended only with the run. The command removes the run's subscription and its listener registration, and the event records it so replay matches the command. A command for a channel the run is not subscribed to still records its event, since every command must produce one.
The channel and the subscribe event it ended are what a reader of the event looks for, so they lead and the task-completed id follows, as the agreed wire change lists them.
Every event in this file leads with the task-completed id, the subscribe event included, and a public API review reads the new event next to them. The channel and the subscribe event id follow it.
A standalone activity can own a channel too, so the five channel calls and both `linked_to` fields take `common.v1.Execution` in place of `WorkflowExecution`, with the old numbers reserved. The linked routes bind `execution.business_id` under the workflow and the activity paths, since a binding cannot set the type. The comments on `metadata` and `identity` spell out that one is folded state and the other is for audit only.
The earlier shape was never released, so the numbers are free and nothing needs reserving.
@moetemp

moetemp commented Oct 3, 2026

Copy link
Copy Markdown
Owner Author

Replaced by #2, #3, #4 and #5.

Same content, split into 4 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