Repository navigation
Conversation
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.
This was referenced Sep 25, 2026
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.
2 of 3 tasks
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.
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 the CHASM stream component, with its retention, caps and producer dedup.
What changed?
chasm/lib/stream/stream.gois 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.temporal.io/content-hash(fingerprint.go). That's abinary/plainpayload 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 areFailedPreconditionwith a documented message prefix. The tokens areSTREAM_PRODUCER_CONFLICT,STREAM_PRODUCER_STALE_SEQUENCE,STREAM_CURSOR_BELOW_FLOORandSTREAM_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.goholds 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.gocaps the bytes and count one read returns.StreamLifecycle.max_bytesandStreamState.held_bytesgive the creator a byte cap beside the record cap. Reclaim keeps the count exact. An append past the cap gets the storage-limitResourceExhaustedand reclaims nothing.stream.retentionRecheckIntervalwhile records remain. It never moves the floor past an active consumer. Closed streams keep the post-close deletion.proto/v1/{message,stream_state,tasks}.protodefine the record, the batch, the state and the two side-effect tasks. Asequenceof zero means the producer doesn't number. The close reason onTerminateis encoded withpayload.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 protoandmake fmtare 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 makesmax_itemsa 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 floorForgetConsumerreleased. Every limit treats zero as the default.