Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
145 changes: 145 additions & 0 deletions docs/stream.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,145 @@
# `temporal stream`

Streams are append-only logs the Temporal Service carries next to Workflow
Executions. This fork adds a `temporal stream` command group on top of the
stream service, and a dev server that serves it.

## Dev server

`temporal server start-dev` registers the stream service and turns it on for
every namespace it creates. The server ships with `stream.enabled` off, so the
dev server sets it as one of its own dynamic-config defaults, the same way it
sets the CHASM and standalone-activity flags. An explicit value still wins:

```sh
temporal server start-dev --dynamic-config-value 'stream.enabled=false'
```

The server behind the dev server is pinned in `go.mod` to the stream branch of
`moedash/temporal`, with the matching `moedash/api-go` pin the server needs.

## Addressing a stream

Two kinds of stream exist. A standalone stream has an ID of its own. An owned
stream lives inside a Workflow Execution or an Activity and is named by its
owner and a stream name.

- `--stream-id ID`: the standalone stream `ID`.
- `--workflow-id W [--run-id R] [--name N]`: the stream `N` the workflow owns.
Without `--name` it is the workflow's default stream.
- `--workflow-id W --activity-id A [--name N]`: the stream of an activity the
workflow scheduled.
- `--activity-id A [--run-id R] [--name N]`: the stream of a standalone
activity.

`describe`, `read` and `append` take either form. `create`, `close`,
`truncate` and `delete` act on a standalone stream only, because an owned
stream's lifecycle is its owner's.

## Commands

`stream create --stream-id ID [--retention D] [--max-items N] [--max-bytes N]`
: Creates a standalone stream. Retention defaults to the namespace's; records
older than it are reclaimed. `--max-items` is a rolling window on records,
`--max-bytes` a ceiling on held bytes that refuses appends until records
are reclaimed. Repeating a create with the same lifecycle answers "already
exists"; a different lifecycle is refused.

`stream list [--query Q] [--limit N] [--page-size N]`
: Lists standalone streams from visibility. `WorkflowId` in the query is the
stream ID. Owned streams are not listed; reach them through their owner.

`stream describe <ref>`
: Frontier, floor, readable record count, held and appended bytes, close
state and reason, retention, caps, budget, producers and consumers. Also
the notification channel the stream notifies on each append and close,
derived from the reference alone: `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 stream. The card names the owner as
`LinkedTo workflow ID` or `LinkedTo activity ID`, with the run when one was
given. `temporal channel poll --channel C` with the same `--workflow-id` or
`--activity-id` follows it.

`stream read <ref> [--from-offset N | --from-tail | --last N] [--follow]
[--topic T]... [--limit N]`
: Prints records with offset, kind, topic, producer, attempt, sequence and
body. Stops once caught up, or with `--follow` when the stream closes.

`stream append <ref> --input V... [--topic T] [--producer-id P] [--attempt A]
[--sequence S] [--expected-offset N] [--finish]`
: Appends one record per `--input` or `--input-file`, or per line of stdin.
`--producer-id` with `--sequence` deduplicates a retry. `--finish` adds a
`FINISH` record after the values.

`stream close --stream-id ID [--reason R]`
: Seals the stream. Records stay readable until retention passes.

`stream truncate --stream-id ID --to N`
: Moves the floor to `N`. Refused below an active consumer's floor.

`stream delete --stream-id ID [--force]`
: Deletes the stream and its records. Refused while a workflow consumes it,
unless forced.

`--output json` prints one JSON object per record or list entry, and the
stream state proto for `describe`.

Authorization follows the server's declaration: `truncate` and `delete` need
the admin role on the namespace, `list`, `describe` and `read` the read role,
and the rest the write role.

## Notification channels

`temporal channel notify|describe|poll|listener add|listener remove` reach the
channel a stream notifies, or any other channel, by `--channel` (`-c`). Without
an owner the commands use the independent channel of that name. With
`--workflow-id W` or `--activity-id A`, and an optional `--run-id`, they use
the channel linked to that Workflow Execution or standalone Activity, sent as
the request's `execution` with the matching type. The two owner flags cannot
be combined, and `--run-id` alone is refused.

`channel describe` prints the kind and, for a linked channel, the owner line
`LinkedTo workflow W (run R)` or `LinkedTo activity A (run R)`, the run only
when the Service names one. The JSON output carries `linkedTo` as the Service
sends it.

## Refusals

A refusal the caller has to act on comes back from the server as a
`FailedPrecondition` whose message starts with a reason token. `create`,
`read` and `append` read the token and print what to do instead of the token,
with the server's own detail after it:

- `STREAM_PRODUCER_CONFLICT`: the producer already appended different content
at this sequence; nothing was written.
- `STREAM_PRODUCER_STALE_SEQUENCE`: the sequence is below the producer's
latest; nothing was written.
- `STREAM_CURSOR_BELOW_FLOOR`: the read starts below the floor; those records
are gone.
- `STREAM_CLOSED`: the stream is closed; its records stay readable.
- `STREAM_POLICY_MISMATCH`: a stream with this id exists with a different
lifecycle.

## Payload codec

The records' bodies and metadata go through the same remote codec as every
other payload the CLI shows or sends, configured with `--codec-endpoint`,
`--codec-auth` and `--codec-header`. The gRPC interceptor that applies the
codec walks only the public API's messages, so the stream commands apply it by
hand on the records and on a close reason.

## Examples

```sh
temporal server start-dev --port 7533

temporal stream create --address localhost:7533 --stream-id scores
temporal stream append --address localhost:7533 --stream-id scores \
--producer-id game --sequence 1 --input '{"home": 1}' --input '{"home": 2}'
temporal stream read --address localhost:7533 --stream-id scores --follow
temporal stream describe --address localhost:7533 --stream-id scores

temporal stream read --address localhost:7533 \
--workflow-id june-s1-abcdef12 --name scores
```
12 changes: 8 additions & 4 deletions go.mod
Original file line number Diff line number Diff line change
@@ -1,6 +1,6 @@
module github.com/temporalio/cli

go 1.26.8
go 1.27.0

require (
github.com/BurntSushi/toml v1.4.0
Expand All @@ -10,15 +10,15 @@ require (
github.com/fatih/color v1.19.0
github.com/google/uuid v1.6.0
github.com/mattn/go-isatty v0.0.23
github.com/nexus-rpc/sdk-go v0.6.0
github.com/nexus-rpc/sdk-go v0.7.0
github.com/olekukonko/tablewriter v0.0.5
github.com/spf13/cobra v1.10.2
github.com/spf13/pflag v1.0.10
github.com/stretchr/testify v1.12.1
github.com/temporalio/cli/cliext v0.0.0
github.com/temporalio/ui-server/v2 v2.54.1
go.temporal.io/api v1.63.5
go.temporal.io/sdk v1.47.0
go.temporal.io/api v1.63.6-0.20260909222256-20151aa90480
go.temporal.io/sdk v1.48.0
go.temporal.io/sdk/contrib/envconfig v1.0.2
go.temporal.io/server v1.32.0
golang.org/x/exp v0.0.0-20260611194520-c48552f49976
Expand Down Expand Up @@ -232,3 +232,7 @@ require (
sigs.k8s.io/structured-merge-diff/v6 v6.4.0 // indirect
sigs.k8s.io/yaml v1.6.0 // indirect
)

replace go.temporal.io/server => github.com/moedash/temporal v0.0.0-20261003023939-91a1c125ec02

replace go.temporal.io/api => github.com/moedash/api-go v1.63.6-0.20261003021656-71792676a9fd
16 changes: 8 additions & 8 deletions go.sum
Original file line number Diff line number Diff line change
Expand Up @@ -334,14 +334,18 @@ github.com/modern-go/concurrent v0.0.0-20180306012644-bacd9c7ef1dd/go.mod h1:6dJ
github.com/modern-go/reflect2 v1.0.2/go.mod h1:yWuevngMOJpCy52FWWMvUC8ws7m/LJsjYzDa0/r8luk=
github.com/modern-go/reflect2 v1.0.3-0.20250322232337-35a7c28c31ee h1:W5t00kpgFdJifH4BDsTlE89Zl93FEloxaWZfGcifgq8=
github.com/modern-go/reflect2 v1.0.3-0.20250322232337-35a7c28c31ee/go.mod h1:yWuevngMOJpCy52FWWMvUC8ws7m/LJsjYzDa0/r8luk=
github.com/moedash/api-go v1.63.6-0.20261003021656-71792676a9fd h1:40ANlQLRwkA2/uY3o7U7v8kElCSclu55fA/RtRTyfQk=
github.com/moedash/api-go v1.63.6-0.20261003021656-71792676a9fd/go.mod h1:acM0I9WPuYg8W3Pd9jOZvEgi7mRUttUQ4+e7fowKVnM=
github.com/moedash/temporal v0.0.0-20261003023939-91a1c125ec02 h1:HOOvpfHzduzr7F3TaWM7a09/waOqjYULOgoMZ2MzPM4=
github.com/moedash/temporal v0.0.0-20261003023939-91a1c125ec02/go.mod h1:tReLM/fxuPU/oPY6leZDno9TkXXy+aO8rRMC3ufDygE=
github.com/munnerz/goautoneg v0.0.0-20191010083416-a7dc8b61c822 h1:C3w9PqII01/Oq1c1nUAm88MOHcQC9l5mIlSMApZMrHA=
github.com/munnerz/goautoneg v0.0.0-20191010083416-a7dc8b61c822/go.mod h1:+n7T8mK8HuQTcFwEeznm/DIxMOiR9yIdICNftLE1DvQ=
github.com/ncruces/go-strftime v1.0.0 h1:HMFp8mLCTPp341M/ZnA4qaf7ZlsbTc+miZjCLOFAw7w=
github.com/ncruces/go-strftime v1.0.0/go.mod h1:Fwc5htZGVVkseilnfgOVb9mKy6w1naJmn9CehxcKcls=
github.com/nexus-rpc/nexus-proto-annotations v0.1.0 h1:2fELd+9sqUtNu6Fg//pw8YFsxOvp8vZ8hfP0nHhNI80=
github.com/nexus-rpc/nexus-proto-annotations v0.1.0/go.mod h1:n3UjF1bPCW8llR8tHvbxJ+27yPWrhpo8w/Yg1IOuY0Y=
github.com/nexus-rpc/sdk-go v0.6.0 h1:QRgnP2zTbxEbiyWG/aXH8uSC5LV/Mg1fqb19jb4DBlo=
github.com/nexus-rpc/sdk-go v0.6.0/go.mod h1:FHdPfVQwRuJFZFTF0Y2GOAxCrbIBNrcPna9slkGKPYk=
github.com/nexus-rpc/sdk-go v0.7.0 h1:38NrfY5rLnZAiMMs2ZfCKI/CSDzdfJG+27iAgfA8bUI=
github.com/nexus-rpc/sdk-go v0.7.0/go.mod h1:FHdPfVQwRuJFZFTF0Y2GOAxCrbIBNrcPna9slkGKPYk=
github.com/niemeyer/pretty v0.0.0-20200227124842-a10e7caefd8e/go.mod h1:zD1mROLANZcx1PVRCS0qkT7pwLkGfwJo4zjcN/Tysno=
github.com/olekukonko/tablewriter v0.0.5 h1:P2Ga83D34wi1o9J6Wh1mRuqd4mF/x/lgBS7N7AbDhec=
github.com/olekukonko/tablewriter v0.0.5/go.mod h1:hPp6KlRPjbx+hW8ykQs1w3UBbZlj6HuIJcUGPhkA7kY=
Expand Down Expand Up @@ -481,16 +485,12 @@ go.opentelemetry.io/otel/trace v1.44.0 h1:jxF5CsGYCe74MCRx2X4g7WsY/VBKRqqpNvXlX/
go.opentelemetry.io/otel/trace v1.44.0/go.mod h1:oLl1jrMQAVo6v3GAggN+1VH9VIz9iUSvW53sW1Q8PIE=
go.opentelemetry.io/proto/otlp v1.10.0 h1:IQRWgT5srOCYfiWnpqUYz9CVmbO8bFmKcwYxpuCSL2g=
go.opentelemetry.io/proto/otlp v1.10.0/go.mod h1:/CV4QoCR/S9yaPj8utp3lvQPoqMtxXdzn7ozvvozVqk=
go.temporal.io/api v1.63.5 h1:c11+kPYHkXXL3UiShPdbMD+xtvqGsbTibUA9ypmiCa4=
go.temporal.io/api v1.63.5/go.mod h1:SrlW2JMwVlDP4nRWSNznUFqnSHd+YeMDS1BkYo63HCQ=
go.temporal.io/auto-scaled-workers v0.2.0-1.32.0.158.0 h1:l+Rj0cIHMC2VB/DC+axrLLnZ2ISTgUpUw8TMiP4B9F8=
go.temporal.io/auto-scaled-workers v0.2.0-1.32.0.158.0/go.mod h1:ZGNY7kCU0EZpXx02D81/jGXHbUcFcVIzzguJ/f8pvUQ=
go.temporal.io/sdk v1.47.0 h1:lZ39w1+uWSjHTL0F3mSc0t4XUnKX8CCcWxqSiLbaHnc=
go.temporal.io/sdk v1.47.0/go.mod h1:ilKs0twgP4JpP8pfhIgZumnOEyBiYn6ZO/ta//NnKMU=
go.temporal.io/sdk v1.48.0 h1:WDctKDVuh0Z8Nf7euAyqs/EwcPg1JTIIq1Fut8Tq118=
go.temporal.io/sdk v1.48.0/go.mod h1:SHv3+fLzD0GGZAwf0xNSvu8UmO1nFgG9WBSYoowApIk=
go.temporal.io/sdk/contrib/envconfig v1.0.2 h1:MGHfsuPUtsf7X9M6WYn3zYJj/mWsuYHnA1uuiL0KEuE=
go.temporal.io/sdk/contrib/envconfig v1.0.2/go.mod h1:MuMiH7hksps2uXnmKuAWaP9P6WbkSDy62kl64t1VJVg=
go.temporal.io/server v1.32.0 h1:JQoqsVREaGc8vcuDNfuRlySiV3tVi81XUOVC9LeL/1Y=
go.temporal.io/server v1.32.0/go.mod h1:SuxEWp1bDjSB7kHtUjyaUJBh7Qjyv5wB8lCrS0VFSvw=
go.uber.org/atomic v1.5.0/go.mod h1:sABNBOSYdrvTF6hTgEIbc7YasKWGhgEQZyfxyTvoXHQ=
go.uber.org/atomic v1.7.0/go.mod h1:fEN4uk6kAWBTFdckzkM89CLk9XfWZrxpCo0nPH17wJc=
go.uber.org/atomic v1.11.0 h1:ZvwS0R+56ePWxUNi+Atn9dWONBPp/AUETXlHW0DxSjE=
Expand Down
14 changes: 14 additions & 0 deletions internal/devserver/server.go
Original file line number Diff line number Diff line change
Expand Up @@ -40,6 +40,8 @@ import (
uiserveroptions "github.com/temporalio/ui-server/v2/server/server_options"
"go.temporal.io/api/enums/v1"
"go.temporal.io/server/chasm/lib/activity"
chasmcallback "go.temporal.io/server/chasm/lib/callback"
"go.temporal.io/server/chasm/lib/stream"
"go.temporal.io/server/common/authorization"
"go.temporal.io/server/common/cluster"
"go.temporal.io/server/common/config"
Expand Down Expand Up @@ -250,6 +252,18 @@ func (s *StartOptions) buildServerOptions() ([]temporal.ServerOption, *slog.Leve
dynConf[activity.EnableStandaloneActivityOperatorCommands.Key()] = true
dynConf[dynamicconfig.FrontendEnableBatchOperationsForStandaloneActivities.Key()] = true

// The server ships streams off so a deployment opts in. The dev server is
// the opt-in: local development is what it exists for, and an explicit
// value below still wins.
dynConf[stream.EnabledSetting.Key()] = true
// The server calls no callback address until one is allowed. A local
// receiver is the usual target of a notification channel's callback on
// a dev server, and it rarely serves TLS.
dynConf[chasmcallback.AllowedAddresses.Key()] = []any{
map[string]any{"Pattern": "127.0.0.1:*", "AllowInsecure": true},
map[string]any{"Pattern": "localhost:*", "AllowInsecure": true},
}

// Dynamic config if set
for k, v := range s.DynamicConfigValues {
dynConf[dynamicconfig.MakeKey(k)] = v
Expand Down
52 changes: 47 additions & 5 deletions internal/temporalcli/client.go
Original file line number Diff line number Diff line change
Expand Up @@ -7,6 +7,7 @@ import (
"io/fs"
"os"
"os/user"
"sync/atomic"

"github.com/temporalio/cli/cliext"
"go.temporal.io/api/common/v1"
Expand All @@ -32,8 +33,30 @@ func dialClient(cctx *CommandContext, c *cliext.ClientOptions) (client.Client, e
// used by the gRPC interceptor; callers can use it to decode payloads nested inside
// opaque proto bytes (e.g. the request/response of a system Nexus operation).
func dialClientWithCodec(cctx *CommandContext, c *cliext.ClientOptions) (client.Client, converter.PayloadCodec, error) {
cl, err := dialClientWithConn(cctx, c)
if err != nil {
return nil, nil, err
}
return cl.Client, cl.PayloadCodec, nil
}

// dialedClient is an SDK client together with what the SDK dialed for it.
type dialedClient struct {
client.Client
// The connection the SDK dialed, carrying its TLS, credentials, headers and
// interceptors, for services the SDK client does not wrap.
Conn grpc.ClientConnInterface
// As returned by [dialClientWithCodec]. The interceptor applying it walks
// the public API's messages only, so payloads inside any other message are
// encoded and decoded by hand.
PayloadCodec converter.PayloadCodec
}

// dialClientWithConn is like [dialClient] but also returns the connection the
// SDK dialed and the remote payload codec.
func dialClientWithConn(cctx *CommandContext, c *cliext.ClientOptions) (*dialedClient, error) {
if cctx.RootCommand == nil {
return nil, nil, fmt.Errorf("root command unexpectedly missing when dialing client")
return nil, fmt.Errorf("root command unexpectedly missing when dialing client")
}

// Set default identity if not provided
Expand Down Expand Up @@ -61,12 +84,12 @@ func dialClientWithCodec(cctx *CommandContext, c *cliext.ClientOptions) (client.
// original setup error instead of attaching a guessed address or profile.
var pathErr *fs.PathError
if errors.As(err, &pathErr) {
return nil, nil, newConnectError(&connectDiagnosis{
return nil, newConnectError(&connectDiagnosis{
Cause: causeCertFileUnreadable,
Detail: pathErr.Path,
}, connectMeta{}, err)
}
return nil, nil, err
return nil, err
}

// We do not put codec on data converter here, it is applied via
Expand All @@ -82,6 +105,20 @@ func dialClientWithCodec(cctx *CommandContext, c *cliext.ClientOptions) (client.
clientOpts.ConnectionOptions.DialOptions = append(
clientOpts.ConnectionOptions.DialOptions, grpc.WithChainUnaryInterceptor(fixedHeaderOverrideInterceptor))

// The SDK does not hand out the connection it dials, but every interceptor
// on it is given the connection, and the dial itself makes a call.
var conn atomic.Pointer[grpc.ClientConn]
clientOpts.ConnectionOptions.DialOptions = append(
clientOpts.ConnectionOptions.DialOptions, grpc.WithChainUnaryInterceptor(
func(
ctx context.Context,
method string, req, reply any,
cc *grpc.ClientConn, invoker grpc.UnaryInvoker, opts ...grpc.CallOption,
) error {
conn.CompareAndSwap(nil, cc)
return invoker(ctx, method, req, reply, cc, opts...)
}))

// Additional gRPC options
clientOpts.ConnectionOptions.DialOptions = append(
clientOpts.ConnectionOptions.DialOptions, cctx.Options.AdditionalClientGRPCDialOptions...)
Expand All @@ -97,14 +134,19 @@ func dialClientWithCodec(cctx *CommandContext, c *cliext.ClientOptions) (client.

cl, err := client.DialContext(dialCtx, clientOpts)
if err != nil {
return nil, nil, dialConnectError(cctx, dialCtx, clientOpts, err)
return nil, dialConnectError(cctx, dialCtx, clientOpts, err)
}

// Since this namespace value is used by many commands after this call,
// we are mutating it to be the derived one
c.Namespace = clientOpts.Namespace

return cl, builder.PayloadCodec, nil
cc := conn.Load()
if cc == nil {
cl.Close()
return nil, fmt.Errorf("dialing made no call, so the connection was not captured")
}
return &dialedClient{Client: cl, Conn: cc, PayloadCodec: builder.PayloadCodec}, nil
}

// dialConnectError enriches a client.DialContext failure with a staged
Expand Down
Loading
Loading