Skip to content

Added the CHASM stream component. - #4

Closed
moetemp wants to merge 22 commits into
moe/AI-198-srv-1-api-pinfrom
moe/AI-198-srv-2-stream-component
Closed

moetemp wants to merge 22 commits into
moe/AI-198-srv-1-api-pinfrom
moe/AI-198-srv-2-stream-component

Conversation

@moetemp

@moetemp moetemp commented Sep 25, 2026 •

Copy link
Copy Markdown
Owner

This PR adds the CHASM stream component, with its retention, caps and producer dedup.

What changed?

  • chasm/lib/stream/stream.go is the component. It holds records in batches next to the owning workflow, with a frontier, producer cursors and consumer pins. Appends, closing, truncation, cap reclaim and the read window live here. Truncation never passes an active consumer's floor. A filtered read moves only past the messages it examined, so a short page skips nothing.
  • Producer dedup keeps one entry per producer id and matches on its sequence. A record can declare its identity under temporal.io/content-hash (fingerprint.go). That's a binary/plain payload with the 64 hex bytes of the SHA-256 over the converted body, taken before any codec (stream.ContentHashOf). The table compares that hash, so a codec with a per-call nonce can't turn a retry into a conflict. Only appends that name a producer read the key. A value that isn't hex falls back to the encoded fingerprint.
  • reason.go: refusals a caller must act on are FailedPrecondition with a documented message prefix. The tokens are STREAM_PRODUCER_CONFLICT, STREAM_PRODUCER_STALE_SEQUENCE, STREAM_CURSOR_BELOW_FLOOR and STREAM_CLOSED. An SDK maps them to typed errors until a typed detail carries the reason on the wire.
  • cursor.go: the consumer cursor lives on the reader, not the stream. A delivered range then commits with the event that records it.
  • config.go holds the caps as namespace dynamic config, plus a byte and item budget for a stream a workflow owns, well under the mutable state limit. The budget counts held records, head minus base, so a stream created above zero isn't over budget at birth. messages.go caps the bytes and count one read returns.
  • StreamLifecycle.max_bytes and StreamState.held_bytes give the creator a byte cap beside the record cap. Reclaim keeps the count exact. An append past the cap gets the storage-limit ResourceExhausted and reclaims nothing.
  • Retention is an age on an open stream. Each batch carries its append time. The first append arms a check that reclaims older batches and re-arms on stream.retentionRecheckInterval while records remain. It never moves the floor past an active consumer. Closed streams keep the post-close deletion.
  • proto/v1/{message,stream_state,tasks}.proto define the record, the batch, the state and the two side-effect tasks. A sequence of zero means the producer doesn't number. The close reason on Terminate is encoded with payload.EncodeString, so a data converter reads it.

This layer wires nothing up. The service, the commands and delivery come in #5 to #7.

Part of AI-198 (epic AI-37).

Why?

The payload lives in the component, not in History and not in a log table. Everything above this layer depends on that property, so it's worth reading before the wiring.

How did you test it?

go build ./..., make proto and make fmt are clean. The ./chasm/lib/stream/... unit tests pass, with the age, cursor, fingerprint and read-cap suites. Nothing here runs in a cluster yet. Limits: a registered workflow consumer makes max_items a lifetime quota until it deregisters, since replay re-reads every range. One entry per producer id allows one append in flight, so pipelining needs an id per lane. A query or reset of a closed run is refused once truncation takes the floor ForgetConsumer released. Every limit treats zero as the default.

  • Unit Tests
  • Staging
  • End to End Tests

Records live in the component next to the workflow that owns them, so a publish rides the same commit. Retention, truncation and the consumer floor are the component's own state, and the caps that bound resource use come from namespace-scoped dynamic config.
The item budget compared an absolute head offset against a count, so a stream
whose floor had moved, or one created above zero, was over budget the moment
it existed. The close reason now carries its encoding, and the record kind is
settled on a copy rather than through the caller's protos.
The stored record mirrors the public one field for field, and the public one no longer reserves -1 for an unnumbered record.
A divergent repeat, a stale sequence and a read below the floor were generic
failures the SDK could only tell apart by matching prose. Each is now a
FailedPrecondition whose message begins with a documented token, so the SDK
maps it to its typed error until a wire detail carries the reason.
The producer table fingerprinted the encoded batch, so a payload codec with a
fresh nonce per call turned every retry into a divergent repeat. A record may
carry the hex SHA-256 of its converted body under a metadata key the codec
never touches, and the table compares that instead; without the key nothing
changes.
The SDK maps a sealed stream to its own error, and the refusal was a bare
FailedPrecondition it could only tell apart by prose. It now carries the
fourth documented token next to the other three.
… for producers.

The value is a binary/plain payload holding the hex SHA-256 over the
deterministic serialization of the converted body. A workflow's own publish
carries no producer to repeat under and its codec may encode the value, so
the key is read only on appends that name a producer, and a value that is
not hex falls back to the encoded fingerprint instead of refusing the append.
The creator could cap records but not bytes, so a stream of large records had
no bound short of the namespace's. The stream now counts the bytes it holds,
reclaim keeps that count exact, and an append that would take it past
max_bytes is refused with the storage-limit refusal a budgeted stream gives.
Retention was only a time to deletion after close, while the interface
promises an age on an open stream too. Each batch now carries when it was
appended, an age check armed by the first append reclaims batches older than
the retention on the recheck interval, and it never moves the floor past an
active consumer's replay floor. The fingerprint is taken over the records
alone, so the stamp cannot turn a retry into a conflict.
The registration's start field changes shape on the layer above, and the
test only needs a consumer whose floor holds, so it writes the entry itself.
@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