From 7cfa0758d929c48d6fd9e03c1250ca96ff7ddc8e Mon Sep 17 00:00:00 2001 From: Claude Date: Wed, 9 Sep 2026 09:11:40 +0000 Subject: [PATCH] feat(marmot): sequence discipline for agent text stream publishing MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Publishing preview records was blocked on one thing: the spec requires a publisher to never restart or reuse a `seq` for one key context — "including after reconnect, retry, process restart, or daemon resume" — and to stop publishing entirely when it cannot prove which value is next. That is a cryptographic requirement, not bookkeeping. `seq` is XORed into the ChaCha20-Poly1305 record nonce and the key is fixed for the stream, so a repeated `seq` repeats a (key, nonce) pair, which leaks the XOR of the two plaintexts and forfeits authentication for every record under that key. `AgentTextStreamPublisher` owns that discipline. Sequence values are reserved in the durable store before a record is handed out, in windows so a chatty stream is not a write per record — a crash then skips the unused tail of a window rather than replaying it, and a gap is something the transport binding already handles while a repeat is a nonce collision. `resume` returns null, rather than starting over at 1, both when nothing was retained and when the stream was finished or aborted; the caller falls back to the authoritative final kind:9 and a later preview needs a fresh stream id. A frame the group's `max_plaintext_frame_len` refuses is rejected before it claims a value, so a refused frame does not leave every receiver with a permanent gap where no record ever existed. Only an in-memory sequence store ships here. `send` (0xF2D2) and `fanout` (0xF2D4) stay unadvertised: there is no QUIC data plane behind them yet, and claiming a role we cannot serve is worse for a group than not claiming it. Co-Authored-By: Claude Opus 5 Claude-Session: https://claude.ai/code/session_016kCuA6tc4JQzHPCDd39GHq --- quartz/plans/2026-09-08-marmot-spec-resync.md | 23 +- .../AgentTextStreamPublisher.kt | 274 ++++++++++++++++++ .../AgentTextStreamPublisherTest.kt | 218 ++++++++++++++ 3 files changed, 508 insertions(+), 7 deletions(-) create mode 100644 quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/marmot/appComponents/agentTextStream/AgentTextStreamPublisher.kt create mode 100644 quartz/src/jvmAndroidTest/kotlin/com/vitorpamplona/quartz/marmot/appComponents/AgentTextStreamPublisherTest.kt diff --git a/quartz/plans/2026-09-08-marmot-spec-resync.md b/quartz/plans/2026-09-08-marmot-spec-resync.md index 648d3df42a..025cbbae47 100644 --- a/quartz/plans/2026-09-08-marmot-spec-resync.md +++ b/quartz/plans/2026-09-08-marmot-spec-resync.md @@ -680,13 +680,22 @@ test we have. harness-relay flake, not a protocol failure, and it costs whichever test is running at the time. Worth making the CLI's publish confirmation tolerate a reconnect rather than papering over it in the tests. -- **Agent-text-stream is receive-only.** We decode the `0x8006` policy, derive - per-stream record keys, open records and fold the transcript, and we advertise - the `0xF2D1` receive capability. We do NOT advertise `send` (`0xF2D2`) or - `fanout` (`0xF2D4`): publishing needs durable per-stream sequence state to - avoid reusing an AEAD nonce across a restart, and there is none. A group whose - policy requires `send` is refused at join rather than joined into a state - every peer would reject us from. +- **Agent-text-stream publishes records but has nowhere to send them.** We + decode the `0x8006` policy, derive per-stream record keys, open records and + fold the transcript, we advertise the `0xF2D1` receive capability, and + `AgentTextStreamPublisher` now seals records under the sequence discipline the + spec demands: values reserved durably ahead of use (in windows, so the hot + path is not a write per record), never restarted or replayed across a crash, + and a publisher that cannot prove which value is next refuses to publish at + all rather than colliding a ChaCha20-Poly1305 (key, nonce) pair. Only an + in-memory `AgentTextStreamSequenceStore` exists; a platform-backed one lands + with the transport that needs it. + + We still do NOT advertise `send` (`0xF2D2`) or `fanout` (`0xF2D4`), because + there is no data plane behind them yet and a role we cannot serve is worse for + the group than a role we do not claim. A group whose policy requires `send` is + refused at join rather than joined into a state every peer would reject us + from. - **The QUIC transport for agent text streams is not wired.** The record layer and the kind-1200 anchor are implemented; nothing yet opens a WebTransport session to a broker and feeds it records. diff --git a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/marmot/appComponents/agentTextStream/AgentTextStreamPublisher.kt b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/marmot/appComponents/agentTextStream/AgentTextStreamPublisher.kt new file mode 100644 index 0000000000..9c31d7652f --- /dev/null +++ b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/marmot/appComponents/agentTextStream/AgentTextStreamPublisher.kt @@ -0,0 +1,274 @@ +/* + * Copyright (c) 2025 Vitor Pamplona + * + * Permission is hereby granted, free of charge, to any person obtaining a copy of + * this software and associated documentation files (the "Software"), to deal in + * the Software without restriction, including without limitation the rights to use, + * copy, modify, merge, publish, distribute, sublicense, and/or sell copies of the + * Software, and to permit persons to whom the Software is furnished to do so, + * subject to the following conditions: + * + * The above copyright notice and this permission notice shall be included in all + * copies or substantial portions of the Software. + * + * THE SOFTWARE IS PROVIDED "AS IS", WITHOUT WARRANTY OF ANY KIND, EXPRESS OR + * IMPLIED, INCLUDING BUT NOT LIMITED TO THE WARRANTIES OF MERCHANTABILITY, FITNESS + * FOR A PARTICULAR PURPOSE AND NONINFRINGEMENT. IN NO EVENT SHALL THE AUTHORS OR + * COPYRIGHT HOLDERS BE LIABLE FOR ANY CLAIM, DAMAGES OR OTHER LIABILITY, WHETHER IN + * AN ACTION OF CONTRACT, TORT OR OTHERWISE, ARISING FROM, OUT OF OR IN CONNECTION + * WITH THE SOFTWARE OR THE USE OR OTHER DEALINGS IN THE SOFTWARE. + */ +package com.vitorpamplona.quartz.marmot.appComponents.agentTextStream + +import com.vitorpamplona.quartz.nip01Core.core.toHexKey +import kotlinx.coroutines.sync.Mutex +import kotlinx.coroutines.sync.withLock + +/** + * What a publisher must remember about one agent text stream to be allowed to + * keep publishing it. + * + * @param nextSeq the lowest sequence value that has NEVER been handed out. + * Reserved ahead of use, so it may sit above the last record actually sent. + * @param closed true once the stream was finished or aborted. A closed stream + * is closed for good — a later preview needs a fresh `stream_id`, which + * produces a new start payload and therefore a new key context. + */ +class AgentTextStreamSequenceState( + val nextSeq: Long, + val closed: Boolean, +) + +/** + * Durable sequence bookkeeping for agent text streams we publish. + * + * This is not caching. `features/agent-text-streams-quic.md` requires that a + * publisher never restart or reuse a `seq` for one + * [AgentTextStreamKeyContextV1] — "including after reconnect, retry, process + * restart, or daemon resume" — and that a publisher which cannot prove which + * value is next stop publishing. The reason is cryptographic rather than + * clerical: `seq` is XORed into the ChaCha20-Poly1305 record nonce, and the + * key is fixed for the stream, so a repeated `seq` repeats a (key, nonce) pair. + * That leaks the XOR of the two plaintexts and forfeits authentication for + * every record under that key. + * + * Keys are [AgentTextStreamKeyContextV1.encode] bytes, so a different stream + * id, epoch, sender or start event is a different entry with its own sequence. + * + * Implementations SHOULD encrypt at rest for the same reason the message store + * does: the key context contains the group id and the stream anchor. + */ +interface AgentTextStreamSequenceStore { + suspend fun load(keyContext: ByteArray): AgentTextStreamSequenceState? + + suspend fun save( + keyContext: ByteArray, + state: AgentTextStreamSequenceState, + ) +} + +/** Process-lifetime store. Suitable for tests and for a publisher that never restarts. */ +class InMemoryAgentTextStreamSequenceStore : AgentTextStreamSequenceStore { + private val states = mutableMapOf() + + override suspend fun load(keyContext: ByteArray): AgentTextStreamSequenceState? = states[keyContext.toHexKey()] + + override suspend fun save( + keyContext: ByteArray, + state: AgentTextStreamSequenceState, + ) { + states[keyContext.toHexKey()] = state + } +} + +/** + * Publishes the QUIC record side of one agent text stream. + * + * The publisher owns exactly one cryptographic session, the one the stream's + * kind-1200 start payload authorized, and it owns the sequence discipline that + * makes that session safe. Sequence values are reserved in the durable store + * BEFORE a record is handed out, in windows of [RESERVATION_WINDOW] so the hot + * path is not a write per record. A crash therefore skips the unused tail of a + * window rather than replaying it: a gap is something the transport binding + * already has to handle, and a repeat is a nonce collision. + * + * This produces records. It does not move them — the QUIC/WebTransport data + * plane that carries them to a broker is a separate layer, and until that + * exists a group's `send` role (`0xF2D2`) stays unadvertised, because + * advertising a role we cannot serve is worse for the group than not + * advertising it. + */ +class AgentTextStreamPublisher private constructor( + private val crypto: AgentTextStreamCrypto, + private val store: AgentTextStreamSequenceStore, + private val maxPlaintextFrameLen: Long, + val transcript: AgentTextStreamTranscriptV1, + private var nextSeq: Long, + private var reservedThrough: Long, +) { + private val mutex = Mutex() + private var closed = false + + private val keyContext = crypto.context.encode() + + /** The next sequence value this publisher would use. Exposed for diagnostics. */ + val peekNextSeq: Long get() = nextSeq + + val isClosed: Boolean get() = closed + + /** + * Seal [plaintextFrame] as the next record in the stream. + * + * The frame is validated against the group's `max_plaintext_frame_len` + * before a sequence value is claimed, so a frame the policy refuses does + * not burn one — a burnt value would show up at every receiver as a + * permanent gap in a stream that never had a record there. + */ + suspend fun publish( + recordType: Int, + plaintextFrame: ByteArray, + ): AgentTextStreamRecordV1 { + require(plaintextFrame.size <= maxPlaintextFrameLen) { + "agent text stream plaintext frame is larger than the group's limit" + } + return mutex.withLock { + check(!closed) { "agent text stream ${crypto.context.streamId.toHexKey()} is closed" } + val seq = claimNextSeqUnlocked() + val sealed = + crypto.seal( + AgentTextStreamRecordV1( + streamId = crypto.context.streamId, + seq = seq, + recordType = recordType, + frame = plaintextFrame, + ), + maxPlaintextFrameLen, + ) + transcript.append(seq, recordType, plaintextFrame) + sealed + } + } + + /** Publish a `FinalNotice` and close the stream. */ + suspend fun finish(notice: ByteArray = ByteArray(0)): AgentTextStreamRecordV1 { + val record = publish(AgentTextStreamRecordV1.TYPE_FINAL_NOTICE, notice) + close() + return record + } + + /** Publish an `Abort` and close the stream without producing durable text. */ + suspend fun abort(reason: ByteArray = ByteArray(0)): AgentTextStreamRecordV1 { + val record = publish(AgentTextStreamRecordV1.TYPE_ABORT, reason) + close() + return record + } + + /** + * Close the stream permanently. Idempotent, and safe to call without ever + * having published — a start payload whose preview never materialised is + * still a spent key context. + */ + suspend fun close() = + mutex.withLock { + if (closed) return@withLock + closed = true + store.save(keyContext, AgentTextStreamSequenceState(nextSeq = nextSeq, closed = true)) + } + + /** Caller holds [mutex]. Persists the window before returning a value from it. */ + private suspend fun claimNextSeqUnlocked(): Long { + if (nextSeq > reservedThrough) { + reservedThrough = nextSeq + RESERVATION_WINDOW - 1 + store.save(keyContext, AgentTextStreamSequenceState(nextSeq = reservedThrough + 1, closed = false)) + } + return nextSeq++ + } + + companion object { + /** + * How many sequence values one durable write reserves. Larger means + * fewer writes on a chatty stream and a longer gap after a crash; + * neither is a correctness question, only a cost one. + */ + const val RESERVATION_WINDOW = 64L + + /** The first sequence value of a stream, per the record spec. */ + const val FIRST_SEQ = 1L + + /** + * Start publishing a stream whose kind-1200 start payload was just + * emitted. + * + * Refuses when the store already knows this key context: that means + * either a live publisher or a spent one, and starting over would + * re-issue sequence values under a key that has already used them. + */ + suspend fun open( + crypto: AgentTextStreamCrypto, + store: AgentTextStreamSequenceStore, + maxPlaintextFrameLen: Long = AgentTextStreamQuicPolicyV1.MAX_PLAINTEXT_FRAME_LEN, + ): AgentTextStreamPublisher { + val keyContext = crypto.context.encode() + val existing = store.load(keyContext) + check(existing == null) { + "agent text stream ${crypto.context.streamId.toHexKey()} already has publisher state — " + + "a new preview needs a fresh stream id and start payload" + } + return AgentTextStreamPublisher( + crypto = crypto, + store = store, + maxPlaintextFrameLen = maxPlaintextFrameLen, + transcript = AgentTextStreamTranscriptV1.start(crypto.context.streamId, crypto.context.startEventId), + nextSeq = FIRST_SEQ, + reservedThrough = FIRST_SEQ - 1, + ) + } + + /** + * Resume publishing after a restart, or null when this publisher may + * not publish for this start payload any more. + * + * Null is the spec's required outcome in both cases it covers: nothing + * retained (we cannot prove which value is next) and a closed stream. + * The caller falls back to the authoritative final kind-9 message, and + * a later preview attempt starts a new stream. + * + * The transcript resumes empty because the records it would have + * covered were already sent; a resumed publisher's `stream-hash` is + * therefore only meaningful when [transcriptHash] and [chunkCount] are + * carried across the restart alongside the sequence state. + */ + suspend fun resume( + crypto: AgentTextStreamCrypto, + store: AgentTextStreamSequenceStore, + maxPlaintextFrameLen: Long = AgentTextStreamQuicPolicyV1.MAX_PLAINTEXT_FRAME_LEN, + transcriptHash: ByteArray? = null, + chunkCount: Long = 0, + ): AgentTextStreamPublisher? { + val keyContext = crypto.context.encode() + val state = store.load(keyContext) ?: return null + if (state.closed) return null + return AgentTextStreamPublisher( + crypto = crypto, + store = store, + maxPlaintextFrameLen = maxPlaintextFrameLen, + transcript = + if (transcriptHash != null) { + AgentTextStreamTranscriptV1.resume( + crypto.context.streamId, + crypto.context.startEventId, + transcriptHash, + chunkCount, + ) + } else { + AgentTextStreamTranscriptV1.start(crypto.context.streamId, crypto.context.startEventId) + }, + nextSeq = state.nextSeq, + // Everything the previous process reserved is spent as far as + // this one is concerned: it cannot tell which of those values + // actually reached the wire, so it uses none of them. + reservedThrough = state.nextSeq - 1, + ) + } + } +} diff --git a/quartz/src/jvmAndroidTest/kotlin/com/vitorpamplona/quartz/marmot/appComponents/AgentTextStreamPublisherTest.kt b/quartz/src/jvmAndroidTest/kotlin/com/vitorpamplona/quartz/marmot/appComponents/AgentTextStreamPublisherTest.kt new file mode 100644 index 0000000000..cf87553b7d --- /dev/null +++ b/quartz/src/jvmAndroidTest/kotlin/com/vitorpamplona/quartz/marmot/appComponents/AgentTextStreamPublisherTest.kt @@ -0,0 +1,218 @@ +/* + * Copyright (c) 2025 Vitor Pamplona + * + * Permission is hereby granted, free of charge, to any person obtaining a copy of + * this software and associated documentation files (the "Software"), to deal in + * the Software without restriction, including without limitation the rights to use, + * copy, modify, merge, publish, distribute, sublicense, and/or sell copies of the + * Software, and to permit persons to whom the Software is furnished to do so, + * subject to the following conditions: + * + * The above copyright notice and this permission notice shall be included in all + * copies or substantial portions of the Software. + * + * THE SOFTWARE IS PROVIDED "AS IS", WITHOUT WARRANTY OF ANY KIND, EXPRESS OR + * IMPLIED, INCLUDING BUT NOT LIMITED TO THE WARRANTIES OF MERCHANTABILITY, FITNESS + * FOR A PARTICULAR PURPOSE AND NONINFRINGEMENT. IN NO EVENT SHALL THE AUTHORS OR + * COPYRIGHT HOLDERS BE LIABLE FOR ANY CLAIM, DAMAGES OR OTHER LIABILITY, WHETHER IN + * AN ACTION OF CONTRACT, TORT OR OTHERWISE, ARISING FROM, OUT OF OR IN CONNECTION + * WITH THE SOFTWARE OR THE USE OR OTHER DEALINGS IN THE SOFTWARE. + */ +package com.vitorpamplona.quartz.marmot.appComponents + +import com.vitorpamplona.quartz.marmot.appComponents.agentTextStream.AgentTextStreamCrypto +import com.vitorpamplona.quartz.marmot.appComponents.agentTextStream.AgentTextStreamKeyContextV1 +import com.vitorpamplona.quartz.marmot.appComponents.agentTextStream.AgentTextStreamPublisher +import com.vitorpamplona.quartz.marmot.appComponents.agentTextStream.AgentTextStreamRecordV1 +import com.vitorpamplona.quartz.marmot.appComponents.agentTextStream.AgentTextStreamSequenceState +import com.vitorpamplona.quartz.marmot.appComponents.agentTextStream.AgentTextStreamSequenceStore +import com.vitorpamplona.quartz.marmot.appComponents.agentTextStream.InMemoryAgentTextStreamSequenceStore +import kotlinx.coroutines.runBlocking +import org.junit.Assert.assertEquals +import org.junit.Assert.assertFalse +import org.junit.Assert.assertNull +import org.junit.Assert.assertThrows +import org.junit.Assert.assertTrue +import org.junit.Test + +/** + * `features/agent-text-streams-quic.md`: "A publisher MUST NOT restart `seq` + * or reuse any prior `seq` value for the same `AgentTextStreamKeyContextV1`, + * including after reconnect, retry, process restart, or daemon resume. A + * publisher MAY resume only when it has retained the next unused sequence + * value. If it cannot prove which sequence value is next, it MUST stop + * publishing preview records for that start payload." + * + * That is a durability requirement, not a bookkeeping one: the sequence number + * is XORed into the record nonce, so re-using one under the same key context + * reuses a ChaCha20-Poly1305 (key, nonce) pair — which leaks the XOR of two + * plaintexts and forfeits authentication for the whole stream. The publisher + * therefore reserves sequence values durably ahead of use and refuses to + * publish at all when it cannot prove which value is next. + */ +class AgentTextStreamPublisherTest { + private val groupId = ByteArray(32) { 0x51 } + private val streamId = ByteArray(32) { 0x52 } + private val senderId = ByteArray(32) { 0x53 } + private val startEventId = ByteArray(32) { 0x54 } + private val secret = ByteArray(32) { 0x55 } + + private fun context(epoch: Long = 7) = + AgentTextStreamKeyContextV1( + groupId = groupId, + streamId = streamId, + mlsEpoch = epoch, + senderId = senderId, + startEventId = startEventId, + ) + + private fun crypto(epoch: Long = 7) = AgentTextStreamCrypto(secret, context(epoch)) + + @Test + fun theFirstRecordIsSeqOneAndEachRecordAdvancesByOne() = + runBlocking { + val publisher = AgentTextStreamPublisher.open(crypto(), InMemoryAgentTextStreamSequenceStore()) + val seqs = + (1..5).map { + publisher.publish(AgentTextStreamRecordV1.TYPE_TEXT_DELTA, "chunk $it".encodeToByteArray()).seq + } + assertEquals(listOf(1L, 2L, 3L, 4L, 5L), seqs) + } + + @Test + fun aRestartNeverReplaysASequenceValueItAlreadyHandedOut() = + runBlocking { + val store = InMemoryAgentTextStreamSequenceStore() + val before = AgentTextStreamPublisher.open(crypto(), store) + val used = (1..3).map { before.publish(AgentTextStreamRecordV1.TYPE_TEXT_DELTA, "a".encodeToByteArray()).seq } + + // Process dies here — no close, no flush. + val after = AgentTextStreamPublisher.resume(crypto(), store)!! + val next = after.publish(AgentTextStreamRecordV1.TYPE_TEXT_DELTA, "b".encodeToByteArray()).seq + + assertTrue( + "a resumed publisher must hand out a sequence value strictly above every value used before the restart", + next > used.max(), + ) + } + + @Test + fun reservationIsDurableBeforeTheRecordIsHandedOut() = + runBlocking { + val store = InMemoryAgentTextStreamSequenceStore() + val publisher = AgentTextStreamPublisher.open(crypto(), store) + val record = publisher.publish(AgentTextStreamRecordV1.TYPE_TEXT_DELTA, "a".encodeToByteArray()) + + val persisted = store.load(context().encode()) + assertTrue( + "the watermark must already cover the record we handed out — persisting after the fact would " + + "let a crash re-issue the same nonce", + persisted!!.nextSeq > record.seq, + ) + } + + @Test + fun aPublisherThatCannotProveTheNextSequenceRefusesToPublish() = + runBlocking { + // Nothing retained for this key context: the spec says stop, not + // start over from 1. + assertNull( + "resume must fail rather than restart the sequence", + AgentTextStreamPublisher.resume(crypto(), InMemoryAgentTextStreamSequenceStore()), + ) + } + + @Test + fun aStreamThatEndedCannotBeResumed() { + runBlocking { + val store = InMemoryAgentTextStreamSequenceStore() + val publisher = AgentTextStreamPublisher.open(crypto(), store) + publisher.publish(AgentTextStreamRecordV1.TYPE_TEXT_DELTA, "a".encodeToByteArray()) + publisher.finish() + + assertNull( + "a finished stream is closed for good — a later preview needs a fresh stream id and start payload", + AgentTextStreamPublisher.resume(crypto(), store), + ) + assertThrows(IllegalStateException::class.java) { + runBlocking { publisher.publish(AgentTextStreamRecordV1.TYPE_TEXT_DELTA, "b".encodeToByteArray()) } + } + } + } + + @Test + fun anAbortClosesTheStreamToo() = + runBlocking { + val store = InMemoryAgentTextStreamSequenceStore() + val publisher = AgentTextStreamPublisher.open(crypto(), store) + val abort = publisher.abort() + assertEquals(AgentTextStreamRecordV1.TYPE_ABORT, abort.recordType) + assertNull(AgentTextStreamPublisher.resume(crypto(), store)) + } + + @Test + fun aDifferentStartEventIsADifferentStreamWithItsOwnSequence() = + runBlocking { + val store = InMemoryAgentTextStreamSequenceStore() + AgentTextStreamPublisher.open(crypto(epoch = 7), store).also { + it.publish(AgentTextStreamRecordV1.TYPE_TEXT_DELTA, "a".encodeToByteArray()) + it.publish(AgentTextStreamRecordV1.TYPE_TEXT_DELTA, "b".encodeToByteArray()) + } + + // A new epoch is a different key context, so a fresh sequence is + // correct here — the (key, nonce) pair cannot collide across it. + val other = AgentTextStreamPublisher.open(crypto(epoch = 8), store) + assertEquals(1L, other.publish(AgentTextStreamRecordV1.TYPE_TEXT_DELTA, "a".encodeToByteArray()).seq) + } + + @Test + fun everyPublishedRecordIsSealedAndOpensBackToItsPlaintext() = + runBlocking { + val publisher = AgentTextStreamPublisher.open(crypto(), InMemoryAgentTextStreamSequenceStore()) + val plaintext = "hello from the agent".encodeToByteArray() + val record = publisher.publish(AgentTextStreamRecordV1.TYPE_TEXT_DELTA, plaintext) + + assertFalse("the wire record must carry ciphertext", record.frame.contentEquals(plaintext)) + assertEquals(plaintext.size + AgentTextStreamRecordV1.AEAD_TAG_LEN, record.frame.size) + assertTrue(crypto().open(record).frame.contentEquals(plaintext)) + } + + @Test + fun theTranscriptTracksWhatWasPublishedSoTheFinalMessageCanCarryIt() = + runBlocking { + val publisher = AgentTextStreamPublisher.open(crypto(), InMemoryAgentTextStreamSequenceStore()) + publisher.publish(AgentTextStreamRecordV1.TYPE_TEXT_DELTA, "one ".encodeToByteArray()) + publisher.publish(AgentTextStreamRecordV1.TYPE_TEXT_DELTA, "two".encodeToByteArray()) + + assertEquals(2L, publisher.transcript.chunkCount) + assertEquals(32, publisher.transcript.hash.size) + } + + @Test + fun aFrameOverTheGroupsLimitIsRefusedBeforeItBurnsASequenceValue() { + runBlocking { + val store = InMemoryAgentTextStreamSequenceStore() + val publisher = AgentTextStreamPublisher.open(crypto(), store, maxPlaintextFrameLen = 8) + assertThrows(IllegalArgumentException::class.java) { + runBlocking { publisher.publish(AgentTextStreamRecordV1.TYPE_TEXT_DELTA, ByteArray(9)) } + } + assertEquals( + "a refused frame must not consume a sequence value", + 1L, + publisher.publish(AgentTextStreamRecordV1.TYPE_TEXT_DELTA, ByteArray(8)).seq, + ) + } + } + + @Test + fun theStoreRoundTripsItsState() = + runBlocking { + val store: AgentTextStreamSequenceStore = InMemoryAgentTextStreamSequenceStore() + val key = context().encode() + store.save(key, AgentTextStreamSequenceState(nextSeq = 42, closed = false)) + assertEquals(42L, store.load(key)!!.nextSeq) + assertFalse(store.load(key)!!.closed) + store.save(key, AgentTextStreamSequenceState(nextSeq = 42, closed = true)) + assertTrue(store.load(key)!!.closed) + } +}