Repository navigation
Conversation
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.
moetemp
force-pushed
the
moe/AI-198-stream-protos
branch
from
September 16, 2026 22:30
d51ed68 to
2b05496
Compare
1 of 3 tasks
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.
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.
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 public API for server-side streams and the notification channel that wakes their listeners.
What changed?
temporal/api/stream/v1/message.protoholdsStreamRecord, the one wire form every store keeps. It carrieskind,producer_id,attemptandsequencenext tobody,metadataandtopic. An unsetkindreads as data. An emptyproducer_idmeans the owning Workflow wrote the record.command/v1/message.protoaddsAppendStreamRecordsandSubscribeStream. The append addressesstream_name. The subscribe takesstream_name_or_idand astart_position(StreamStartPosition:offset,last_n,earliestortail).history/v1/message.protoaddsWorkflowStreamRecordsAppended, bounded byfrom_offsetandto_offset.WorkflowTaskCompletedEventAttributes.consumed_stream_rangesrecords what a task read asStreamRangevalues.reserved 14 to 19keeps numbers for the main line.PollWorkflowTaskQueueResponse.stream_slicescarries the records a task reads, andreserved 21;on that response keeps a number the series used for one round and dropped.failed_cause.protoadds fiveWorkflowTaskFailedCausevalues: 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.notification/v1(Notification,ChannelListener,ChannelKind), theNotifyChannel,RegisterChannelListener,UnregisterChannelListener,PollChannelandDescribeChannelRPCs, theSubscribeNotificationChannelcommand and its subscribed event, andWorkflowTaskScheduledEventAttributes.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.executionon the five channel requests, acommon.v1.Executionwithtype,business_idandrun_id, so a workflow or a standalone activity can own one and listens without a subscribe command.Notification.linked_tonames the owner as the same type, andDescribeChannelResponseanswerskindandlinked_to. Themetadatacomment states that it is folded state, not a log, and theidentitycomment states that it is for audit and is not copied into the notification.UnsubscribeNotificationChannelends a run's subscription to an independent channel before the run closes: the command incommand/v1, itsWorkflowNotificationChannelUnsubscribedevent inhistory/v1carrying the channel and the id of the subscribe event it ended, and theCommandType,EventTypeandWorkflowTaskFailedCausevalues. A command for a channel the run is not subscribed to still records its event, so replay matches every command.DescribeWorkflowExecutionResponse.channel_subscriptionslists the channels a run stands on asChannelSubscriptionInfovalues inworkflow/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./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.StreamRecord.bodynames 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 breakingagainst upstreammain,api-linterand the codegen checks, and they pass at the head. The additions are additive, sobuf breakingstays 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.