Repository navigation
Conversation
The stream is reachable over gRPC next to the workflow service: ids, names, offsets and page sizes are bounded on the frontend, every method declares its authorization scope and rate-limit priority, and a call for a namespace active elsewhere is forwarded like a workflow RPC.
This was referenced Sep 25, 2026
Registering the service turned the whole surface on everywhere, so it is off by default now and the calls that destroy other workflows' data need an operator. A workflow's owned streams share one byte budget, because their per-stream budgets multiplied out well past the mutable state limit. The notify task no longer strands its coalescing flag or drops a consumer's pin on a transport error.
2 of 3 tasks
…3-stream-service # Conflicts: # service/frontend/service.go
Main now builds under Go 1.27, where go fix rewrites this shape, so make fmt was left dirty on these two files.
A standalone activity gets its own stream map, and an activity a workflow scheduled keeps its streams in the workflow's map under reserved keys until it is a component of its own. One owner reference on the owned-stream calls reaches all three owners, so no second set of calls per owner kind is needed.
Covers the default and named streams of both owners, outside producers and readers, retry inheritance, and a standalone activity's streams ending at its terminal status.
2 of 3 tasks
Earliest, tail and last N are resolved against the frontier of the transaction that registers a consumer or serves a first poll, so none of them races with truncation. The resolved absolute offset is what gets pinned and recorded, which leaves replay and reset untouched.
The Core and Python chains now send a start position, so nothing depends on the sentinel. Refusing it with the field to use tells a caller still sending it, where reinterpreting it would not.
The public stream protos take the `streampb` alias the lint config requires, so the generated ones move to `streamlib` as in the tests. The two switches got their default cases, and the owner test says why it waits.
The component carries the tokens; the functional tests now read them the way an SDK does, by status code and message prefix over a raw gRPC client, so a change to the prose cannot silently break the mapping.
A reader that attaches before the producer's first write got NotFound and had to retry on its own schedule, while a reader of an owned stream could wait on the owner. The poll now rechecks on stream.createWaitRecheckInterval for its wait budget and reads the first batch when the stream appears; NotFound still comes back when the wait expires or the reference pins a run.
The poll RPCs sat on the frontend's long-poll quota and dynamic config bounded sizes, but nothing bounded how fast one namespace appends or polls and nothing counted what it moved. Three namespace rates now answer ResourceExhausted with the RPS cause, and five namespace-tagged counters give metering a unit to read.
An activity's streams take no more records once it is terminal, which is the same refusal a sealed standalone stream gives, so it carries the same token. The functional test reads the token off a sealed stream too.
The engine's already-started error is not a service error and reached the caller as Unknown. A create that repeats the existing lifecycle now answers AlreadyExists, and one that differs is refused with STREAM_POLICY_MISMATCH, so an SDK tells an idempotent retry from a change without a describe.
… to end. A negative cap is refused before the request is routed, like a negative record cap. The functional case appends until the cap refuses, reads the cap and the held bytes back from describe, and gets the room back by truncating.
The component arms the check and decides when to re-arm it; the task only supplies the cadence from stream.retentionRecheckInterval. The functional case watches a record age out of an open stream while a workflow's pin holds an equally old one on another.
Owner
Author
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
StreamService, the gRPC surface for streams, and serves it from the frontend.What changed?
chasm/lib/stream/proto/v1/{service,request_response}.protodefine the RPCs: standalone streams, owned streams, and the two internal calls history makes while it resolves a subscription.chasm/lib/stream/service/holds the history handler, the frontend forwarder, the library registration and the notify-consumers andstreamAgetasks.common/api/metadata.go.TruncateStream,DeleteStreamand the internal calls are admin.gen/streampb/v1/namespace.goexposes each request's namespace. A descriptor-driven test fails any new RPC that skips rate limits, validation, authorization or redirection.service/frontend/configs/quotas.gosets priorities, andservice/frontend/fx.gomakes the methods redirectable.stream.enabled, default off, gates the whole surface. The frontend bounds ids, names, producer and workflow ids, offsets, page sizes and the lifecycle caps. Per-namespace ratesstream.appendRecordsPerSecond,stream.appendBytesPerSecondandstream.pollsPerSecondanswerResourceExhausted. Five counters, fromstream_records_appendedtostream_polls, meter the traffic, and a deduplicated retry adds nothing.stream.maxConsumeItemsPerTaskis clamped to 1000, the largest page one read returns.stream.OwnedStreamsinowned.goholds the create path and the shared byte budget,stream.ownedStreamsMaxBytesPerWorkflow. A workflow (StreamsandStreamCursorsonchasm/lib/workflow/workflow.go) and a standalone activity (chasm/lib/activity/streams.go) both use it. Calls take oneowner { kind, id, run_id, activity_id }routed onowner.id. A request naming onlyworkflow_idfolds into aWORKFLOWowner, so old callers keep working.activity/<escaped id>/<name>keys. A standalone activity's streams end at its terminal status, and a later append getsSTREAM_CLOSED. A retry keeps writing to the same stream.CreateStreamsettles the lifecycle. Unset retention takes the namespace's, and a non-positive or larger one is refused. A repeat with the same lifecycle answersAlreadyExists. A different one gets the newSTREAM_POLICY_MISMATCHtoken, so an SDK tells a retry from a policy change.start.goresolves start positions where the reader is registered:earliestto the floor,tailto the head,last_ntomax(floor, head - N). A first poll resolves in the read that serves it. A blockingPollMessageson a standalone id that names nothing yet parks and rechecks everystream.createWaitRecheckInterval, then reads the first batch.RegisterStreamConsumeranswersstream_absentwhen no execution holds the stream, so the subscribe fallback fires on that alone. A workflow can't ownxand consume a standalonex, since both share one cursor keyspace. The notify task always validates, lowers its flag on discard and fans out through a bounded pool. A partial failure retries the whole set, which is safe because the pushes are idempotent. A truncate refused by a pin probes those runs.The commands that use the workflow's maps come in #6.
Part of AI-198 (epic AI-37).
Why?
The service is the whole external surface. Splitting it from the component keeps two reviews apart: what a caller can do and who's allowed to, versus what the component guarantees. Activity owners land here because the workflow map and the shared budget first exist on this layer.
How did you test it?
go build ./...,make protoandmake fmtare clean, and the lint findings are fixed. Unit tests ran for./chasm/lib/stream/...,./chasm/lib/workflow/...,./common/api/...,./common/rpc/interceptor/...and./service/frontend/configs/.... Thetests/stream_*_test.gofunctional tests ran with-tags test_dep. They cover the RPC surface, lifecycle and policy, metering, owners, start positions on all four owner kinds and ageing out of an open stream. They read the reason tokens the way an SDK does, by status code and message prefix over raw gRPC.