Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
34 commits
Select commit Hold shift + click to select a range
e5feee8
Added the stream service RPCs and the frontend wiring.
moetemp Sep 25, 2026
1764c71
Merge branch 'moe/AI-198-srv-2-stream-component' into moe/AI-198-srv-…
moetemp Sep 25, 2026
7abbfa9
Tightened stream service admission and consumer notification.
moetemp Sep 25, 2026
db9530e
Merge branch 'moe/AI-198-srv-2-stream-component' into moe/AI-198-srv-…
moetemp Sep 25, 2026
9b49920
Merge branch 'moe/AI-198-srv-2-stream-component' into moe/AI-198-srv-…
moetemp Sep 25, 2026
11a2666
Merge branch 'moe/AI-198-srv-2-stream-component' into moe/AI-198-srv-…
moetemp Sep 25, 2026
acffaab
Took the stream handler to errors.AsType.
moetemp Sep 25, 2026
1043385
Merge branch 'moe/AI-198-srv-2-stream-component' into moe/AI-198-srv-…
moetemp Sep 26, 2026
3fcbf9f
Let activities own streams, addressed by an owner reference.
moetemp Sep 28, 2026
72dbd53
Tested streams owned by standalone and workflow activities.
moetemp Sep 28, 2026
cfb3233
Merge branch 'moe/AI-198-srv-2-stream-component' into moe/AI-198-srv-…
moetemp Sep 28, 2026
5e7624a
Resolved stream start positions where the reader is registered.
moetemp Sep 28, 2026
67cf5d9
Tested first polls from each start position, on activity streams too.
moetemp Sep 28, 2026
3b15cce
Tested that a subscription records where its start position resolved.
moetemp Sep 28, 2026
c2fa4de
Refused a negative start offset instead of reading it as the head.
moetemp Sep 28, 2026
7afd8b8
Merge branch 'moe/AI-198-srv-2-stream-component' into moe/AI-198-srv-…
moetemp Sep 29, 2026
6c2c61c
Fixed the lint findings in the stream service and its owner test.
moetemp Sep 29, 2026
79f4475
Merge branch 'moe/AI-198-srv-2-stream-component' into moe/AI-198-srv-…
moetemp Sep 30, 2026
8a9a458
Asserted the reason tokens on the stream service's refusals.
moetemp Sep 30, 2026
49c5856
Parked a blocking poll on a standalone stream until it is created.
moetemp Sep 30, 2026
e0a6296
Rate-limited and metered the stream service per namespace.
moetemp Sep 30, 2026
23f5cea
Merge branch 'moe/AI-198-srv-2-stream-component' into moe/AI-198-srv-…
moetemp Sep 30, 2026
1bc33b7
Carried STREAM_CLOSED on an append to an ended activity's stream.
moetemp Sep 30, 2026
ced0f5b
Merge branch 'moe/AI-198-srv-2-stream-component' into moe/AI-198-srv-…
moetemp Sep 30, 2026
b0b82fb
Told a repeated CreateStream from a policy change on an existing id.
moetemp Sep 30, 2026
f8117a3
Validated the lifecycle's byte cap on the frontend and covered it end…
moetemp Sep 30, 2026
74b5cd4
Ran the stream age check as a pure task on the recheck interval.
moetemp Sep 30, 2026
4abd955
Merge branch 'moe/AI-198-srv-2-stream-component' into moe/AI-198-srv-…
moetemp Sep 30, 2026
f159016
Merge branch 'moe/AI-198-srv-2-stream-component' into moe/AI-198-srv-…
moetemp Oct 1, 2026
286e177
Merge branch 'moe/AI-198-srv-2-stream-component' into moe/AI-198-srv-…
moetemp Oct 2, 2026
658948d
Merge branch 'moe/AI-198-srv-2-stream-component' into moe/AI-198-srv-…
moetemp Oct 2, 2026
4e5353a
Merge branch 'moe/AI-198-srv-2-stream-component' into moe/AI-198-srv-…
moetemp Oct 2, 2026
06d3538
Merge branch 'moe/AI-198-srv-2-stream-component' into moe/AI-198-srv-…
moetemp Oct 3, 2026
00a612a
Merge branch 'moe/AI-198-srv-2-stream-component' into moe/AI-198-srv-…
moetemp Oct 3, 2026
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
6 changes: 6 additions & 0 deletions chasm/lib/activity/activity.go
Original file line number Diff line number Diff line change
Expand Up @@ -17,6 +17,7 @@ import (
"go.temporal.io/server/chasm"
"go.temporal.io/server/chasm/lib/activity/gen/activitypb/v1"
"go.temporal.io/server/chasm/lib/callback"
"go.temporal.io/server/chasm/lib/stream"
"go.temporal.io/server/common"
"go.temporal.io/server/common/contextutil"
"go.temporal.io/server/common/metrics"
Expand Down Expand Up @@ -74,6 +75,11 @@ type Activity struct {
// Callbacks holds completion callbacks to be invoked when this standalone activity reaches a terminal state. Nil
// for workflow-embedded activities as the workflow handles its own callbacks.
Callbacks chasm.Map[string, *callback.Callback]

// Streams this activity owns, keyed by stream name. One map per execution
// rather than per attempt, so a retry keeps writing to the stream its
// earlier attempts wrote to.
Streams chasm.Map[string, *stream.Stream]
}

// WithToken wraps a request with its deserialized task token.
Expand Down
44 changes: 44 additions & 0 deletions chasm/lib/activity/streams.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,44 @@
package activity

import (
"go.temporal.io/server/chasm"
"go.temporal.io/server/chasm/lib/stream"
)

var _ stream.Owner = (*Activity)(nil)

func (a *Activity) ownedStreams() stream.OwnedStreams {
return stream.OwnedStreams{Streams: &a.Streams, Kind: "activity"}
}

// OwnedStream returns the stream under key, or nil when nothing has created it.
func (a *Activity) OwnedStream(ctx chasm.Context, key string) *stream.Stream {
return a.ownedStreams().Get(ctx, key)
}

// OwnedStreamEnded reports that the activity reached a terminal status. A
// failed attempt that will be retried is not the end: the next attempt writes
// to the same stream.
func (a *Activity) OwnedStreamEnded(_ chasm.Context, _ string) bool {
return a.isTerminal()
}

// AppendToOwnedStream appends on behalf of a writer outside the execution,
// which for an activity is every writer: an activity publishes over the
// stream service, never through commands.
//
// Refused once the activity is terminal. A worker still running an attempt
// the server already timed out would otherwise keep adding to a stream its
// readers were told had ended.
func (a *Activity) AppendToOwnedStream(
mctx chasm.MutableContext,
key string,
req stream.AddMessagesRequest,
) (stream.AddMessagesResult, error) {
if a.isTerminal() {
return stream.AddMessagesResult{}, stream.Refusal(stream.ReasonStreamClosed,
"activity execution closed with status %v, so its streams take no more records",
InternalStatusToAPIStatus(a.GetStatus()))
}
return a.ownedStreams().Append(mctx, key, req)
}
104 changes: 104 additions & 0 deletions chasm/lib/activity/streams_test.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,104 @@
package activity

import (
"testing"

"github.com/stretchr/testify/require"
commonpb "go.temporal.io/api/common/v1"
"go.temporal.io/api/serviceerror"
streampb "go.temporal.io/api/stream/v1"
"go.temporal.io/server/chasm"
"go.temporal.io/server/chasm/lib/activity/gen/activitypb/v1"
"go.temporal.io/server/chasm/lib/stream"
streamlib "go.temporal.io/server/chasm/lib/stream/gen/streampb/v1"
)

func streamRecords(bodies ...string) []*streamlib.StreamRecord {
out := make([]*streamlib.StreamRecord, len(bodies))
for i, b := range bodies {
out[i] = &streamlib.StreamRecord{
Body: &commonpb.Payload{Data: []byte(b)},
Kind: streampb.STREAM_RECORD_KIND_DATA,
}
}
return out
}

func activityInStatus(status activitypb.ActivityExecutionStatus) *Activity {
return &Activity{ActivityState: &activitypb.ActivityState{Status: status}}
}

// A standalone activity owns a map of named streams, created on first write,
// and a retry scheduled after a failed attempt keeps writing to the same one.
func TestActivityOwnsStreamsAcrossAttempts(t *testing.T) {
ctx := &chasm.MockMutableContext{}
a := activityInStatus(activitypb.ACTIVITY_EXECUTION_STATUS_STARTED)

first, err := a.AppendToOwnedStream(ctx, stream.DefaultStreamName, stream.AddMessagesRequest{
Records: streamRecords("attempt one"),
})
require.NoError(t, err)
require.Equal(t, int64(0), first.FirstOffset)
require.NotNil(t, a.OwnedStream(ctx, stream.DefaultStreamName))
require.Nil(t, a.OwnedStream(ctx, "progress"), "a name nothing wrote to has no stream")

// The attempt failed and a retry is scheduled. That is not the end of the
// activity, so the stream is not ended either.
a.Status = activitypb.ACTIVITY_EXECUTION_STATUS_SCHEDULED
require.False(t, a.OwnedStreamEnded(ctx, stream.DefaultStreamName))

a.Status = activitypb.ACTIVITY_EXECUTION_STATUS_STARTED
second, err := a.AppendToOwnedStream(ctx, stream.DefaultStreamName, stream.AddMessagesRequest{
Records: streamRecords("attempt two"),
})
require.NoError(t, err)
require.Equal(t, int64(1), second.FirstOffset, "the retry continues the same log")
}

// Once the activity is terminal its streams are ended for readers, and an
// append from a worker still running an attempt the server gave up on is
// refused rather than added after readers were told the stream was over.
func TestActivityStreamsEndAtTerminalStatus(t *testing.T) {
for _, status := range []activitypb.ActivityExecutionStatus{
activitypb.ACTIVITY_EXECUTION_STATUS_COMPLETED,
activitypb.ACTIVITY_EXECUTION_STATUS_FAILED,
activitypb.ACTIVITY_EXECUTION_STATUS_CANCELED,
activitypb.ACTIVITY_EXECUTION_STATUS_TERMINATED,
activitypb.ACTIVITY_EXECUTION_STATUS_TIMED_OUT,
} {
t.Run(status.String(), func(t *testing.T) {
ctx := &chasm.MockMutableContext{}
a := activityInStatus(status)
require.True(t, a.OwnedStreamEnded(ctx, stream.DefaultStreamName))

_, err := a.AppendToOwnedStream(ctx, stream.DefaultStreamName, stream.AddMessagesRequest{
Records: streamRecords("too late"),
})
var precondition *serviceerror.FailedPrecondition
require.ErrorAs(t, err, &precondition)
require.Equal(t, stream.ReasonStreamClosed, stream.ReasonOf(err.Error()),
"an SDK maps it to its closed-stream error by the token")
require.Nil(t, a.OwnedStream(ctx, stream.DefaultStreamName),
"a refused append must not leave a stream behind")
})
}
}

// The count bound and the shared byte budget are the same code a workflow
// uses, so an activity cannot grow its state past what a workflow may.
func TestActivityStreamsAreBoundedLikeAWorkflowsAre(t *testing.T) {
ctx := &chasm.MockMutableContext{}
a := activityInStatus(activitypb.ACTIVITY_EXECUTION_STATUS_STARTED)
limits := stream.Limits{MaxOwnedStreamsPerWorkflow: 2}

for _, name := range []string{"a", "b"} {
_, err := a.AppendToOwnedStream(ctx, name, stream.AddMessagesRequest{
Records: streamRecords("x"), Limits: limits,
})
require.NoError(t, err)
}
_, err := a.AppendToOwnedStream(ctx, "c", stream.AddMessagesRequest{
Records: streamRecords("x"), Limits: limits,
})
require.ErrorContains(t, err, "activity already owns 2 streams")
}
101 changes: 100 additions & 1 deletion chasm/lib/stream/config.go
Original file line number Diff line number Diff line change
Expand Up @@ -95,13 +95,32 @@ const (
// the same reason: each batch is a node, and many small ones cost state
// that the byte budget alone does not see.
OwnedStreamMaxItems = 10_000

// OwnedStreamsMaxBytesPerWorkflow bounds every stream one execution owns
// taken together. The per-stream budget multiplied by the stream count
// comes to far more than limit.mutableStateSize.error, so without this an
// outside writer can name enough streams to terminate the execution while
// every single stream stays inside its own budget. Half the error limit,
// which leaves the rest of mutable state its own room.
OwnedStreamsMaxBytesPerWorkflow = 4 << 20
)

var (
// EnabledSetting gates the whole feature. Off by default: registering
// StreamService on the frontend otherwise turns a large new surface on in
// every deployment the moment it ships, and an operator needs a way to take
// it back without a rollback.
EnabledSetting = dynamicconfig.NewNamespaceBoolSetting(
"stream.enabled",
false,
`Whether the stream service and the workflow stream commands are available to a
namespace. Off by default.`,
)
MaxConsumeItemsPerTaskSetting = dynamicconfig.NewNamespaceIntSetting(
"stream.maxConsumeItemsPerTask",
MaxConsumeItemsPerTask,
`Most stream records one workflow task carries per subscription.`,
`Most stream records one workflow task carries per subscription. Clamped to 1000,
which is the largest page a single stream read returns.`,
)
MaxConsumeBytesPerTaskSetting = dynamicconfig.NewNamespaceIntSetting(
"stream.maxConsumeBytesPerTask",
Expand Down Expand Up @@ -144,20 +163,66 @@ under limit.mutableStateSize.error, which would otherwise terminate the workflow
OwnedStreamMaxItems,
`Message budget of a stream a workflow owns. Appends past it are refused.`,
)
OwnedStreamsMaxBytesPerWorkflowSetting = dynamicconfig.NewNamespaceIntSetting(
"stream.ownedStreamsMaxBytesPerWorkflow",
OwnedStreamsMaxBytesPerWorkflow,
`Byte budget of every stream one workflow execution owns, taken together. Appends
past it are refused. Keep it under limit.mutableStateSize.error, which would otherwise
terminate the workflow.`,
)
RetentionRecheckIntervalSetting = dynamicconfig.NewGlobalDurationSetting(
"stream.retentionRecheckInterval",
time.Minute,
`How long a closed stream past its retention waits before asking again whether the
consumers holding it are still running.`,
)
CreateWaitRecheckIntervalSetting = dynamicconfig.NewGlobalDurationSetting(
"stream.createWaitRecheckInterval",
250*time.Millisecond,
`How often a blocking poll on a standalone stream id that names nothing yet asks
again whether the stream has been created, until the poll's wait expires.`,
)

// The rates below are enforced by the stream service per namespace and
// answered with ResourceExhausted when exceeded. Zero means unlimited: a
// rate is something an operator may want off, and the enablement setting
// already covers turning the feature off.
AppendRecordsPerSecondSetting = dynamicconfig.NewNamespaceIntSetting(
"stream.appendRecordsPerSecond",
AppendRecordsPerSecond,
`Most stream records a namespace may append per second through the stream service.
A burst of one full batch is always admitted. Zero means unlimited.`,
)
AppendBytesPerSecondSetting = dynamicconfig.NewNamespaceIntSetting(
"stream.appendBytesPerSecond",
AppendBytesPerSecond,
`Most stream record bytes a namespace may append per second through the stream
service. A burst of one full batch is always admitted. Zero means unlimited.`,
)
PollsPerSecondSetting = dynamicconfig.NewNamespaceIntSetting(
"stream.pollsPerSecond",
PollsPerSecond,
`Most stream polls a namespace may make per second, blocking or not. Zero means
unlimited.`,
)
)

// Default rates. Generous: they exist so a runaway producer or reader can be
// contained by config, not to shape ordinary traffic.
const (
AppendRecordsPerSecond = 100_000
AppendBytesPerSecond = 256 << 20
PollsPerSecond = 10_000
)

// Config holds the settings as live property functions.
type Config struct {
Enabled dynamicconfig.BoolPropertyFnWithNamespaceFilter
// The id length limit shared with workflow ids. A stream id becomes an
// execution's business id, and a stream name a key in mutable state.
MaxIDLength dynamicconfig.IntPropertyFn
RetentionRecheckInterval dynamicconfig.DurationPropertyFn
CreateWaitRecheckInterval dynamicconfig.DurationPropertyFn
MaxConsumeItemsPerTask dynamicconfig.IntPropertyFnWithNamespaceFilter
MaxConsumeBytesPerTask dynamicconfig.IntPropertyFnWithNamespaceFilter
MaxProducersPerStream dynamicconfig.IntPropertyFnWithNamespaceFilter
Expand All @@ -167,12 +232,21 @@ type Config struct {
MaxOwnedStreamsPerWorkflow dynamicconfig.IntPropertyFnWithNamespaceFilter
OwnedStreamMaxBytes dynamicconfig.IntPropertyFnWithNamespaceFilter
OwnedStreamMaxItems dynamicconfig.IntPropertyFnWithNamespaceFilter
// Bounds every stream one execution owns taken together, which the
// per-stream budget cannot do.
OwnedStreamsMaxBytesPerWorkflow dynamicconfig.IntPropertyFnWithNamespaceFilter

AppendRecordsPerSecond dynamicconfig.IntPropertyFnWithNamespaceFilter
AppendBytesPerSecond dynamicconfig.IntPropertyFnWithNamespaceFilter
PollsPerSecond dynamicconfig.IntPropertyFnWithNamespaceFilter
}

func NewConfig(dc *dynamicconfig.Collection) *Config {
return &Config{
Enabled: EnabledSetting.Get(dc),
MaxIDLength: dynamicconfig.MaxIDLengthLimit.Get(dc),
RetentionRecheckInterval: RetentionRecheckIntervalSetting.Get(dc),
CreateWaitRecheckInterval: CreateWaitRecheckIntervalSetting.Get(dc),
MaxConsumeItemsPerTask: MaxConsumeItemsPerTaskSetting.Get(dc),
MaxConsumeBytesPerTask: MaxConsumeBytesPerTaskSetting.Get(dc),
MaxProducersPerStream: MaxProducersPerStreamSetting.Get(dc),
Expand All @@ -182,6 +256,12 @@ func NewConfig(dc *dynamicconfig.Collection) *Config {
MaxOwnedStreamsPerWorkflow: MaxOwnedStreamsPerWorkflowSetting.Get(dc),
OwnedStreamMaxBytes: OwnedStreamMaxBytesSetting.Get(dc),
OwnedStreamMaxItems: OwnedStreamMaxItemsSetting.Get(dc),

OwnedStreamsMaxBytesPerWorkflow: OwnedStreamsMaxBytesPerWorkflowSetting.Get(dc),

AppendRecordsPerSecond: AppendRecordsPerSecondSetting.Get(dc),
AppendBytesPerSecond: AppendBytesPerSecondSetting.Get(dc),
PollsPerSecond: PollsPerSecondSetting.Get(dc),
}
}

Expand All @@ -197,6 +277,8 @@ type Limits struct {
MaxOwnedStreamsPerWorkflow int
OwnedStreamMaxBytes int
OwnedStreamMaxItems int

OwnedStreamsMaxBytesPerWorkflow int
}

// LimitsFor resolves the limits for a namespace. A nil Config, which is what
Expand All @@ -215,9 +297,21 @@ func (c *Config) LimitsFor(namespaceName string) Limits {
MaxOwnedStreamsPerWorkflow: c.MaxOwnedStreamsPerWorkflow(namespaceName),
OwnedStreamMaxBytes: c.OwnedStreamMaxBytes(namespaceName),
OwnedStreamMaxItems: c.OwnedStreamMaxItems(namespaceName),

OwnedStreamsMaxBytesPerWorkflow: c.OwnedStreamsMaxBytesPerWorkflow(namespaceName),
}.withDefaults()
}

// EnabledFor reports whether a namespace may use streams. A nil Config, which
// is what component code driven without a service gets, reads as enabled: the
// gate is a deployment switch, and a unit test is not a deployment.
func (c *Config) EnabledFor(namespaceName string) bool {
if c == nil || c.Enabled == nil {
return true
}
return c.Enabled(namespaceName)
}

// DefaultLimits is the constant set above.
func DefaultLimits() Limits {
return Limits{}.withDefaults()
Expand All @@ -236,6 +330,10 @@ func (l Limits) withDefaults() Limits {
}
}
fill(&l.MaxConsumeItemsPerTask, MaxConsumeItemsPerTask)
// A slice is built from one stream read, which serves at most a page, so a
// larger setting than that cannot take effect. Clamped here rather than
// left to disagree with what delivery does.
l.MaxConsumeItemsPerTask = min(l.MaxConsumeItemsPerTask, DefaultMaxMessagesPerPoll)
fill(&l.MaxConsumeBytesPerTask, MaxConsumeBytesPerTask)
fill(&l.MaxProducersPerStream, MaxProducersPerStream)
fill(&l.MaxConsumersPerStream, MaxConsumersPerStream)
Expand All @@ -244,5 +342,6 @@ func (l Limits) withDefaults() Limits {
fill(&l.MaxOwnedStreamsPerWorkflow, MaxOwnedStreamsPerWorkflow)
fill(&l.OwnedStreamMaxBytes, OwnedStreamMaxBytes)
fill(&l.OwnedStreamMaxItems, OwnedStreamMaxItems)
fill(&l.OwnedStreamsMaxBytesPerWorkflow, OwnedStreamsMaxBytesPerWorkflow)
return l
}
3 changes: 3 additions & 0 deletions chasm/lib/stream/doc.go
Original file line number Diff line number Diff line change
Expand Up @@ -22,6 +22,9 @@
// starts at.
// - STREAM_CLOSED: an append on a stream that has been sealed. Its records
// stay readable; nothing more goes in.
// - STREAM_POLICY_MISMATCH: a create naming a stream that already exists
// with a different lifecycle. A create that repeats the existing
// lifecycle is an idempotent retry and answers AlreadyExists instead.
//
// [Refusal] builds one and [ReasonOf] reads one back.
//
Expand Down
Loading
Loading