Skip to content

Added the temporal stream and temporal channel command groups. - #2

Closed
moetemp wants to merge 17 commits into
mainfrom
moe/AI-198-stream-cli
Closed

moetemp wants to merge 17 commits into
mainfrom
moe/AI-198-stream-cli

Conversation

@moetemp

@moetemp moetemp commented Sep 30, 2026 •

Copy link
Copy Markdown
Owner

This PR adds a temporal stream command group, a temporal channel command group and a dev server that serves streams and notification channels.

What changed?

  • go.mod replaces go.temporal.io/server with github.com/moetemp/temporal at 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 replaces go.temporal.io/api with the moetemp/api-go head the server series is moving to, which carries ChannelSubscriptionInfo, the unsubscribe command, event and failed cause, and temporal.api.common.v1.Execution on the five channel requests and both linked_to fields, at the numbers workflow_execution and linked_to held. The Go directive is 1.27.0.
  • internal/devserver/server.go sets stream.enabled as a dynamic-config default, next to the CHASM and standalone-activity flags. It also sets callback.allowedAddresses to 127.0.0.1 and localhost on any port over plain HTTP, so a local receiver can be a channel's callback listener. An explicit --dynamic-config-value for either key still wins.
  • commands.stream.go adds create, list, describe, read, append, close, truncate and delete. A stream is --stream-id for a standalone one, or --workflow-id or --activity-id, with an optional --run-id and --name, for an owned one.
  • read starts at the beginning, --from-offset, --from-tail or --last N, filters with --topic, and --follow long-polls until close. create takes --retention, --max-items and --max-bytes. describe shows held and appended bytes next to the frontier and floor.
  • describe also names the notification channel the stream notifies on each append and close, derived from the reference without a call: stream/NAME linked to the owner for a workflow's or a standalone activity's stream, stream/ACTIVITY_ID/NAME linked to the workflow for an activity it scheduled, and the independent stream/STREAM_ID for a standalone one, with the owner line LinkedTo workflow ID or LinkedTo activity ID, ready for temporal channel poll.
  • A stream refusal that starts with one of the server's reason tokens prints in plain words, with the server's detail after it. The token comes from the server's own stream.ReasonOf.
  • The stream commands reach the stream service over the connection the SDK dialed, captured by a unary interceptor. TLS, API keys and headers match every other command. They apply the payload codec by hand on bodies, metadata and the close reason, because the codec interceptor walks only the public API's messages.
  • commands.channel.go adds notify, describe, poll, listener add and listener remove over 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-id or --activity-id and an optional --run-id, sent as the request's execution with the matching type. The two owner flags cannot be combined, and --run-id alone is refused.
  • notify takes --position, --counter and --metadata KEY=VALUE with JSON values, sent as JSON payloads the way --input sends a Signal's arguments, and prints how many listeners it reached. describe shows the kind, the owner of a linked channel as LinkedTo workflow ID or LinkedTo activity ID with 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 carries linkedTo as the server sends it. poll waits up to --wait for notifications above --after-counter, and --follow keeps polling from the highest counter seen. listener add registers a callback URL with --header values and prints the listener ID that listener remove takes.
  • The channel commands use the SDK's workflow service client, so the codec interceptor covers their payloads. A missing channel, a linked channel whose workflow has closed, and the server's limits (ResourceExhausted on the listener cap and on the notify rate) print in plain words, with the server's detail after them.
  • temporal workflow describe prints 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 carries channelSubscriptions as the server sends it.
  • docs/stream.md describes the stream commands, the refusals and the dev-server flag. The listener add help 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. create isn't on the gap list, but without it no standalone stream exists to close, truncate or delete. list covers standalone streams only, because visibility holds no owned streams, and its --query matches the stream ID on WorkflowId.

The channel group gives the notification channel the same reach from a terminal: a writer can be imitated with notify, a client listener with poll, and a callback listener with listener add, and describe shows 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 its callback.allowedAddresses dynamic-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=1 and the cmd/gen-commands diff pass. The TestStream tests cover the standalone lifecycle with dedup, each refusal, start positions, topic filter, paging, byte caps, truncation, --follow ending on close, an owned stream, list with a query and flag validation. A unit test against an in-process fake of the stream service checks the channel describe names 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 polls stream/STREAM_ID for the notification and the closed metadata. At an earlier server pin, June scenario s1 ran on the native provider against temporal server start-dev, and the CLI read back its scores stream. No test covers the stream codec path.

The TestChannel unit tests run the commands against an in-process fake of the workflow service: flag validation that sends nothing, the requests each command sends, the execution each 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, --follow advancing its cursor and stopping on interrupt, the plain-language refusals for either owner, and the workflow describe table with an independent and a linked entry, a blank pending column, the JSON pass-through and no section on an empty list. The TestChannel suite 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 leaves describe once removed, and an unknown channel is refused in plain words. On a running workflow, a notify to its linked channel reaches the owner, describe shows the linked kind, the owner and the retained notification, poll returns it with the owner set, a callback listener comes and goes in describe, 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 channel workflow describe lists the linked entry in text and JSON. On a standalone activity started with temporal activity start, a notify to its linked channel reaches nobody and is retained, describe shows the linked kind and the activity owner with its run, poll returns the notification with linkedTo of the activity type, a callback listener comes and goes, an append to the activity's scores stream lands on stream/scores linked 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 against temporal server start-dev built from this branch, with the JSON output carrying linkedTo.type as EXECUTION_TYPE_ACTIVITY on describe and on the polled notification, along with a localhost callback accepted by default, another host refused, and an empty callback.allowedAddresses overriding 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 and closed=true.

  • Unit Tests
  • Staging
  • End to End Tests

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.
@moetemp moetemp changed the title Added the temporal stream command group and a dev server that serves streams. Added temporal stream commands and a dev server that serves streams. Oct 1, 2026
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.
@moetemp moetemp changed the title Added temporal stream commands and a dev server that serves streams. Added the temporal stream and temporal channel command groups. Oct 2, 2026
moetemp added 10 commits October 1, 2026 23:35
--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.
@moetemp

moetemp commented Oct 3, 2026

Copy link
Copy Markdown
Owner Author

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

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