feat(marmot): sequence discipline for agent text stream publishing

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 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_016kCuA6tc4JQzHPCDd39GHq
This commit is contained in:
Claude
2026-09-09 09:11:40 +00:00
parent fa5e14e605
commit 7cfa0758d9
3 changed files with 508 additions and 7 deletions
+16 -7
View File
@@ -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.
@@ -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<String, AgentTextStreamSequenceState>()
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,
)
}
}
}
@@ -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)
}
}