Repository navigation
Conversation
The server dependency now points at the stream branch head of moedash/temporal, with the api-go pin that head needs. The server ships stream.enabled off so a deployment opts in; the dev server is that opt-in and sets it as one of its own dynamic-config defaults.
create, list, describe, read, append, close, truncate and delete over the stream service, addressing a standalone stream by id or an owned one by its owner and name. The commands reach the service over the connection the SDK dialed and apply the remote payload codec by hand, since the codec interceptor walks only the public API's messages.
That head adds a byte cap to the stream lifecycle, counts held bytes, and prefixes refusals with a reason token.
create takes the lifecycle's byte cap, describe shows held next to appended bytes, and create, read and append read the server's reason token off a FailedPrecondition and say what to do instead of the token.
temporal stream command group and a dev server that serves streams.temporal stream commands and a dev server that serves streams.
notify, describe, poll and listener add/remove over the notification channel calls, pinned to the srv-7 server fork and its api-go head so the dev server serves them. The dev server allows localhost callbacks, since the server calls no callback address until one is allowed.
temporal stream commands and a dev server that serves streams.temporal stream and temporal channel command groups.
--workflow-id and --run-id on every channel command address the channel linked to that workflow; without them the commands use the independent channel. The server pin moves to srv-8, which serves the linked kind.
The head carries the channel subscriptions on the workflow description and the unsubscribe command, event and failed cause, which workflow show gets through the generated code.
The table shows the channel, kind, counters and the linked counts when the description carries any. The JSON output carries the field as the server sends it.
The name follows from the stream's identity, so the CLI derives it without a call and a user can pass it to channel poll. The default owned stream prints the name the server resolves it to.
The head serves the channel subscriptions on the workflow description, the stream-named channels and the unsubscribe command, so the live cases can cover them.
…ver. A notify to a workflow's linked channel lists it on describe and a fresh workflow lists nothing. A standalone stream's append and close reach the channel the CLI names for it.
The head adds the channel.linkedKindEnabled dynamic config, on by default, so nothing the CLI relies on changes.
--workflow-id or --activity-id on every channel command sends the owner as the request's execution with the matching type, and the owner line reads it back from linked_to. A standalone activity's stream notifies the linked stream/<name> on the activity, so the server pin moves to the layer that resolves the execution.
The execution fields take the numbers the workflow_execution fields held, since the earlier shape was never released. The server pin builds against it unchanged because the Go field names are the same.
The head mirrors the public channel field numbers into the server's internal protos, with no Go change and nothing new on the wire for the CLI.
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 a
temporal streamcommand group, atemporal channelcommand group and a dev server that serves streams and notification channels.What changed?
go.modreplacesgo.temporal.io/serverwithgithub.com/moetemp/temporalat the srv-14 head (moetemp/temporal PR Bump go.temporal.io/sdk from 1.17.0 to 1.18.1 temporalio/cli#17) that serves both kinds of notification channel, the channels native streams notify, the channel subscriptions on the workflow description, the unsubscribe command, and a linked channel addressed by execution. It replacesgo.temporal.io/apiwith themoetemp/api-gohead the server series is moving to, which carriesChannelSubscriptionInfo, the unsubscribe command, event and failed cause, andtemporal.api.common.v1.Executionon the five channel requests and bothlinked_tofields, at the numbersworkflow_executionandlinked_toheld. The Go directive is 1.27.0.internal/devserver/server.gosetsstream.enabledas a dynamic-config default, next to the CHASM and standalone-activity flags. It also setscallback.allowedAddressesto127.0.0.1andlocalhoston any port over plain HTTP, so a local receiver can be a channel's callback listener. An explicit--dynamic-config-valuefor either key still wins.commands.stream.goaddscreate,list,describe,read,append,close,truncateanddelete. A stream is--stream-idfor a standalone one, or--workflow-idor--activity-id, with an optional--run-idand--name, for an owned one.readstarts at the beginning,--from-offset,--from-tailor--last N, filters with--topic, and--followlong-polls until close.createtakes--retention,--max-itemsand--max-bytes.describeshows held and appended bytes next to the frontier and floor.describealso names the notification channel the stream notifies on each append and close, derived from the reference without a call:stream/NAMElinked to the owner for a workflow's or a standalone activity's stream,stream/ACTIVITY_ID/NAMElinked to the workflow for an activity it scheduled, and the independentstream/STREAM_IDfor a standalone one, with the owner lineLinkedTo workflow IDorLinkedTo activity ID, ready fortemporal channel poll.stream.ReasonOf.commands.channel.goaddsnotify,describe,poll,listener addandlistener removeover the channel calls of the workflow service. A channel is--channel(-c), a name the writer and its listeners agree on. An independent channel exists on its own and any number of listeners attach to it. A channel linked to a workflow lives with that Workflow Execution, which listens on it without subscribing, and a channel linked to a standalone activity lives with that activity the same way. Every channel command names the owner with--workflow-idor--activity-idand an optional--run-id, sent as the request'sexecutionwith the matching type. The two owner flags cannot be combined, and--run-idalone is refused.notifytakes--position,--counterand--metadata KEY=VALUEwith JSON values, sent as JSON payloads the way--inputsends a Signal's arguments, and prints how many listeners it reached.describeshows the kind, the owner of a linked channel asLinkedTo workflow IDorLinkedTo activity IDwith the run when the server names one, the retained count, the latest notification, and the workflow and callback listeners in separate tables. The JSON output carrieslinkedToas the server sends it.pollwaits up to--waitfor notifications above--after-counter, and--followkeeps polling from the highest counter seen.listener addregisters a callback URL with--headervalues and prints the listener ID thatlistener removetakes.ResourceExhaustedon the listener cap and on the notify rate) print in plain words, with the server's detail after them.temporal workflow describeprints a "Notification Channels" table when the description carries any: channel, kind, last counter, pending counter (blank without a pending notification), scheduled counter, listeners and retained. The JSON output carrieschannelSubscriptionsas the server sends it.docs/stream.mddescribes the stream commands, the refusals and the dev-server flag. Thelistener addhelp names the allowlist setting for servers other than the dev server.Part of AI-198 (epic AI-37).
Why?
Two prototype gaps sit outside the server series: a CLI view of streams, and a local server that carries the stream service without a build from the branch.
createisn't on the gap list, but without it no standalone stream exists to close, truncate or delete.listcovers standalone streams only, because visibility holds no owned streams, and its--querymatches the stream ID onWorkflowId.The channel group gives the notification channel the same reach from a terminal: a writer can be imitated with
notify, a client listener withpoll, and a callback listener withlistener add, anddescribeshows what the server holds. Both kinds are reachable, so a stream owned by a workflow can be tried on its linked channel and a standalone stream on an independent one. The server only calls callback addresses itscallback.allowedAddressesdynamic-config value allows, and upstream's default allows none, so the dev server allows local addresses the same way it turns streams on.How did you test it?
go build ./...,go test ./internal/temporalcli/ -count=1and thecmd/gen-commandsdiff pass. TheTestStreamtests cover the standalone lifecycle with dedup, each refusal, start positions, topic filter, paging, byte caps, truncation,--followending on close, an owned stream,listwith a query and flag validation. A unit test against an in-process fake of the stream service checks the channeldescribenames and owner lines for a standalone stream, a workflow's default and named streams, an activity the workflow scheduled and a standalone activity. A suite test creates a standalone stream, appends, closes, and pollsstream/STREAM_IDfor the notification and theclosedmetadata. At an earlier server pin, June scenario s1 ran on the native provider againsttemporal server start-dev, and the CLI read back itsscoresstream. No test covers the stream codec path.The
TestChannelunit tests run the commands against an in-process fake of the workflow service: flag validation that sends nothing, the requests each command sends, theexecutioneach owner flag sends with its type and the run only where given, text and JSON output with the owner line for a workflow and for a standalone activity,--followadvancing its cursor and stopping on interrupt, the plain-language refusals for either owner, and theworkflow describetable with an independent and a linked entry, a blank pending column, the JSON pass-through and no section on an empty list. TheTestChannelsuite tests run against the shared dev server at the srv-14 pin (moetemp/temporal PR temporalio#17): a notify with no listeners reaches none and is retained, an equal or lower counter is folded and not retained, a poll returns only higher counters and answers when a notification lands, a callback listener is called with its header and leavesdescribeonce removed, and an unknown channel is refused in plain words. On a running workflow, a notify to its linked channel reaches the owner,describeshows the linked kind, the owner and the retained notification,pollreturns it with the owner set, a callback listener comes and goes indescribe, an untouched name describes as linked with nothing held, and the independent channel of the same name stays separate. A fresh workflow describes with no channel section, and after a notify to its linked channelworkflow describelists the linked entry in text and JSON. On a standalone activity started withtemporal activity start, a notify to its linked channel reaches nobody and is retained,describeshows the linked kind and the activity owner with its run,pollreturns the notification withlinkedToof the activity type, a callback listener comes and goes, an append to the activity'sscoresstream lands onstream/scoreslinked to the activity with the run and head as its position, and an activity nobody started is refused in plain words. The same checks passed by hand againsttemporal server start-devbuilt from this branch, with the JSON output carryinglinkedTo.typeasEXECUTION_TYPE_ACTIVITYondescribeand on the polled notification, along with alocalhostcallback accepted by default, another host refused, and an emptycallback.allowedAddressesoverriding the default. A notify to the linked channel of a terminated workflow is refused in plain words. By hand, two appends and a close on a standalone stream arrived on its channel as one notification with the latest counter andclosed=true.