Skip to content

Added the stream service and its frontend wiring. - #5

Closed
moetemp wants to merge 34 commits into
moe/AI-198-srv-2-stream-componentfrom
moe/AI-198-srv-3-stream-service
Closed

moetemp wants to merge 34 commits into
moe/AI-198-srv-2-stream-componentfrom
moe/AI-198-srv-3-stream-service

Conversation

@moetemp

@moetemp moetemp commented Sep 25, 2026 •

Copy link
Copy Markdown
Owner

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}.proto define 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 and streamAge tasks.
  • Admission: each method declares scope and access in common/api/metadata.go. TruncateStream, DeleteStream and the internal calls are admin. gen/streampb/v1/namespace.go exposes each request's namespace. A descriptor-driven test fails any new RPC that skips rate limits, validation, authorization or redirection. service/frontend/configs/quotas.go sets priorities, and service/frontend/fx.go makes 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 rates stream.appendRecordsPerSecond, stream.appendBytesPerSecond and stream.pollsPerSecond answer ResourceExhausted. Five counters, from stream_records_appended to stream_polls, meter the traffic, and a deduplicated retry adds nothing. stream.maxConsumeItemsPerTask is clamped to 1000, the largest page one read returns.
  • Owners: stream.OwnedStreams in owned.go holds the create path and the shared byte budget, stream.ownedStreamsMaxBytesPerWorkflow. A workflow (Streams and StreamCursors on chasm/lib/workflow/workflow.go) and a standalone activity (chasm/lib/activity/streams.go) both use it. Calls take one owner { kind, id, run_id, activity_id } routed on owner.id. A request naming only workflow_id folds into a WORKFLOW owner, so old callers keep working.
  • A workflow's activities keep their streams in the workflow's map under reserved activity/<escaped id>/<name> keys. A standalone activity's streams end at its terminal status, and a later append gets STREAM_CLOSED. A retry keeps writing to the same stream.
  • CreateStream settles the lifecycle. Unset retention takes the namespace's, and a non-positive or larger one is refused. A repeat with the same lifecycle answers AlreadyExists. A different one gets the new STREAM_POLICY_MISMATCH token, so an SDK tells a retry from a policy change.
  • start.go resolves start positions where the reader is registered: earliest to the floor, tail to the head, last_n to max(floor, head - N). A first poll resolves in the read that serves it. A blocking PollMessages on a standalone id that names nothing yet parks and rechecks every stream.createWaitRecheckInterval, then reads the first batch.
  • RegisterStreamConsumer answers stream_absent when no execution holds the stream, so the subscribe fallback fires on that alone. A workflow can't own x and consume a standalone x, 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 proto and make fmt are 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/.... The tests/stream_*_test.go functional 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.

  • Unit Tests
  • Staging
  • End to End Tests

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.
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.
…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.
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.
@moetemp moetemp closed this Oct 1, 2026
@moetemp moetemp reopened this Oct 1, 2026
@moetemp

moetemp commented Oct 3, 2026

Copy link
Copy Markdown
Owner Author

Replaced by #18, #19, #20, #21, #22, #23, #24, #25, #26, #27, #28, #29, #30, #31, #32, #33, #34, #35, #36 and #37.

Same content, split into 20 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