feat(marmot): wire agent text streams end to end, both directions

The transport was there and the codecs were there; nothing joined them to
a group. Now `amy marmot stream start|send|watch|finish` does: a hidden
kind:1200 anchors the stream over MLS, records ride raw QUIC through a
broker, and a kind:9 closes it carrying the transcript a receiver checks
its own fold against.

`AgentTextStreamSubscriber` is the receive discipline the binding spells
out, and it matters because a preview that quietly diverges is worse than
no preview: `seq` accepted at most once and never folded out of order, a
replayed record (which a broker WILL send from the start of its replay
window on reconnect) discarded silently and never stream-fatal, a gap
that cannot be backfilled marking the preview unverifiable because the
transcript hash can no longer complete. Only TextDelta and Checkpoint
reach the answer text — progress and status are chrome the spec forbids
from ever reaching notifications, indexes or automation input.

The start payload also grew the tags it was missing: `stream-type`,
`final-kind` and the optional `parent`, plus the rule that a final
payload whose kind disagrees with `final-kind` is ignored.

Verified in both directions against MDK in harness tests 18 and 19: `wn
stream verify` confirms our transcript from our own kind:1200 + kind:9,
and our subscriber folds MDK's stream to a transcript hash identical to
the one `wn stream send` computed. That equality is the key schedule, key
context, AEAD, framing and transcript construction all agreeing with an
implementation that is not ours. 19 of 19 harness tests pass, twice.

Two defects only that exercise could have found:

  - The epoch belongs to the stream, not to the clock. The record key
    context binds mls_epoch, and both sides were resolving it as "the
    group's current epoch" at each command, so a commit landing between
    the start and the send put them on different keys and produced an
    empty preview. The epoch that DELIVERED the kind:1200 is the
    stream's; it is persisted with the message now and read back by
    publisher and receiver alike.

  - close() dropped the tail of a stream. enqueue only fills the send
    buffer, so tearing the connection down before the driver flushed it
    lost records silently — the publisher had already counted them. QUIC
    ACKs a FIN only once everything ahead of it arrived, so finish() now
    waits for finAcked. This is exactly why the test passed alone and
    failed inside a full run.

The `send` (0xF2D2) and `fanout` (0xF2D4) role capabilities stay
unadvertised: a role is a promise to the whole group, and only the CLI
originates a stream so far.

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 12:43:44 +00:00
parent 6d43d1941a
commit 7014d88b2e
19 changed files with 1404 additions and 25 deletions
@@ -111,16 +111,64 @@ class AndroidMarmotMessageStore(
override suspend fun delete(nostrGroupId: String) {
withContext(Dispatchers.IO) {
writeMutex.withLock {
val file = messagesFile(nostrGroupId)
if (file.exists() && !file.delete()) {
Log.w(TAG) { "delete($nostrGroupId): failed to remove ${file.absolutePath}" }
for (file in listOf(messagesFile(nostrGroupId), epochsFile(nostrGroupId))) {
if (file.exists() && !file.delete()) {
Log.w(TAG) { "delete($nostrGroupId): failed to remove ${file.absolutePath}" }
}
}
}
}
}
private fun readAll(nostrGroupId: String): List<String> {
val file = messagesFile(nostrGroupId)
private fun epochsFile(nostrGroupId: String): File = File(groupDir(nostrGroupId), "epochs")
/**
* Which MLS epoch delivered an inner event. Agent text streams bind the
* epoch into their record key context, so a receiver needs the epoch that
* carried the stream's kind:1200 anchor rather than the group's current
* one — a commit landing in between would otherwise derive a different key
* and render nothing.
*
* Stored through the same encrypted codec as the messages: the ids are as
* sensitive as the payloads they point at.
*/
override suspend fun recordEpoch(
nostrGroupId: String,
innerEventId: String,
epoch: Long,
) = withContext(Dispatchers.IO) {
writeMutex.withLock {
try {
val line = "$innerEventId $epoch"
val existing = readAllFrom(epochsFile(nostrGroupId)).toMutableList()
if (line in existing) return@withLock
existing.add(line)
writeAllTo(epochsFile(nostrGroupId), existing)
} catch (e: Exception) {
Log.e(TAG, "recordEpoch($nostrGroupId) FAILED: ${e.message}", e)
}
}
}
override suspend fun loadEpochs(nostrGroupId: String): Map<String, Long> =
withContext(Dispatchers.IO) {
try {
readAllFrom(epochsFile(nostrGroupId))
.mapNotNull { line ->
val parts = line.trim().split(' ')
if (parts.size != 2) return@mapNotNull null
val epoch = parts[1].toLongOrNull() ?: return@mapNotNull null
parts[0] to epoch
}.toMap()
} catch (e: Exception) {
Log.e(TAG, "loadEpochs($nostrGroupId) FAILED: ${e.message}", e)
emptyMap()
}
}
private fun readAll(nostrGroupId: String): List<String> = readAllFrom(messagesFile(nostrGroupId))
private fun readAllFrom(file: File): List<String> {
if (!file.exists()) return emptyList()
val encrypted = file.readBytes()
val plain = encryption.decrypt(encrypted) ?: return emptyList()
@@ -151,8 +199,12 @@ class AndroidMarmotMessageStore(
private fun writeAll(
nostrGroupId: String,
messages: List<String>,
) = writeAllTo(messagesFile(nostrGroupId), messages)
private fun writeAllTo(
file: File,
messages: List<String>,
) {
val file = messagesFile(nostrGroupId)
file.parentFile?.mkdirs()
val encodedEntries = messages.map { it.encodeToByteArray() }
+4
View File
@@ -43,6 +43,10 @@ tasks.named<Test>("test") {
dependencies {
implementation(project(":quartz"))
implementation(project(":commons"))
// Agent text stream previews: the raw-QUIC binding, and the QUIC
// stack under it for the certificate validator the transport requires.
implementation(project(":marmotQuic"))
implementation(project(":quic"))
// `amy serve` embeds geode (the standalone Ktor relay built on quartz's
// relay-server code). geode depends only on :quartz, never on :amethyst.
implementation(project(":geode"))
@@ -71,6 +71,7 @@ import com.vitorpamplona.amethyst.cli.commands.SearchCommand
import com.vitorpamplona.amethyst.cli.commands.ServeCommand
import com.vitorpamplona.amethyst.cli.commands.StatusCommand
import com.vitorpamplona.amethyst.cli.commands.StoreCommands
import com.vitorpamplona.amethyst.cli.commands.StreamCommands
import com.vitorpamplona.amethyst.cli.commands.SubscribeCommand
import com.vitorpamplona.amethyst.cli.commands.SyncCommand
import com.vitorpamplona.amethyst.cli.commands.UseCommand
@@ -370,12 +371,13 @@ private suspend fun marmotDispatch(
route(
name = "marmot",
tail = tail,
usage = "marmot <key-package|group|message|await|reset>",
usage = "marmot <key-package|group|message|stream|await|reset>",
routes =
mapOf(
"key-package" to { rest -> KeyPackageCommands.dispatch(dataDir, rest) },
"group" to { rest -> GroupCommands.dispatch(dataDir, rest) },
"message" to { rest -> MessageCommands.dispatch(dataDir, rest) },
"stream" to { rest -> StreamCommands.dispatch(dataDir, rest) },
"await" to { rest -> AwaitCommands.dispatch(dataDir, rest) },
"reset" to { rest -> MarmotResetCommand.run(dataDir, rest) },
),
@@ -0,0 +1,373 @@
/*
* 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.amethyst.cli.commands
import com.vitorpamplona.amethyst.cli.Args
import com.vitorpamplona.amethyst.cli.Context
import com.vitorpamplona.amethyst.cli.DataDir
import com.vitorpamplona.amethyst.cli.Output
import com.vitorpamplona.marmotquic.QuicAgentTextStreamTransport
import com.vitorpamplona.quartz.marmot.appComponents.agentTextStream.AgentTextStreamPublisher
import com.vitorpamplona.quartz.marmot.appComponents.agentTextStream.AgentTextStreamRecordV1
import com.vitorpamplona.quartz.marmot.appComponents.agentTextStream.AgentTextStreamStart
import com.vitorpamplona.quartz.marmot.appComponents.agentTextStream.AgentTextStreamSubscriber
import com.vitorpamplona.quartz.marmot.appComponents.agentTextStream.InMemoryAgentTextStreamSequenceStore
import com.vitorpamplona.quartz.marmot.appComponents.agentTextStream.PreviewStatus
import com.vitorpamplona.quartz.marmot.appComponents.agentTextStream.RecordOutcome
import com.vitorpamplona.quartz.marmot.mls.crypto.MlsCryptoProvider
import com.vitorpamplona.quartz.nip01Core.core.Event
import com.vitorpamplona.quartz.nip01Core.core.hexToByteArray
import com.vitorpamplona.quartz.nip01Core.core.toHexKey
import com.vitorpamplona.quic.tls.PermissiveCertificateValidator
import kotlinx.coroutines.withTimeoutOrNull
/**
* `amy marmot stream` — agent text stream previews (`0x8006`).
*
* The durable half is ordinary Marmot messaging: a hidden kind:1200 anchors
* the stream and a kind:9 closes it, both over MLS. The live half is raw QUIC
* to a broker, and it is strictly a progressive enhancement — a member that
* never opens a QUIC connection still reads the whole answer from the final
* kind:9.
*/
object StreamCommands {
val USAGE: String =
"""
|amy marmot stream — agent text stream previews over QUIC
|
| marmot stream start GID [--stream-id HEX] [--broker quic://HOST:PORT[,…]]
| publish the kind:1200 that anchors a stream; prints stream_id + start_event_id
|
| marmot stream send GID --stream-id HEX --start-event-id HEX --broker URI TEXT…
| push TEXT as TextDelta records to the broker; prints the transcript to finish with
|
| marmot stream watch GID [--stream-id HEX] [--timeout SECS]
| find the kind:1200 in the group, subscribe over QUIC, fold the preview
|
| marmot stream finish GID --stream-id HEX --transcript-hash HEX --chunk-count N TEXT…
| publish the authoritative kind:9 carrying the transcript a receiver checks against
|
|Every record is encrypted under the group's own MLS exporter secret, so a
|broker relays ciphertext and learns only which room it belongs to.
""".trimMargin()
suspend fun dispatch(
dataDir: DataDir,
tail: Array<String>,
): Int =
route(
"stream",
tail,
"stream <start|send|watch|finish> …",
mapOf(
"start" to { rest -> start(dataDir, rest) },
"send" to { rest -> send(dataDir, rest) },
"watch" to { rest -> watch(dataDir, rest) },
"finish" to { rest -> finish(dataDir, rest) },
),
help = USAGE,
)
private suspend fun start(
dataDir: DataDir,
rest: Array<String>,
): Int {
val args = Args(rest)
val positional = args.positional
if (positional.isEmpty()) return Output.error("bad_args", "stream start GID [--stream-id HEX] [--broker URI]…")
val streamId = args.flag("stream-id") ?: MlsCryptoProvider.randomBytes(32).toHexKey()
if (streamId.length != 64) return Output.error("bad_args", "--stream-id must be 32 bytes of hex")
// Repeatable in the spec, comma-separated here: `Args` keeps one
// value per flag and a receiver tries them in the order given.
val brokers =
args
.flag("broker")
?.split(',')
?.map { it.trim() }
?.filter { it.isNotEmpty() } ?: emptyList()
Context.open(dataDir).use { ctx ->
ctx.prepare()
val gid = ctx.resolveGroupId(positional[0])
ctx.syncIncoming()
if (!ctx.marmot.isMember(gid)) return Output.error("not_member", "not a member of group $gid")
val bundle = ctx.marmot.buildAgentStreamStart(gid, streamId, brokers, parentEventId = args.flag("parent"))
val targets = ctx.marmotGroupRelays(gid).ifEmpty { ctx.outboxRelays() }
val ack = ctx.publish(bundle.outbound.signedEvent, targets)
RawEventSupport.publishGuard(ack, bundle.outbound.signedEvent.id)?.let { return it }
Output.emit(
mapOf(
"group_id" to gid,
"stream_id" to streamId,
// The start payload's OWN id is the stream anchor that
// goes into the key context — not the kind:445 that
// carried it, and not the MLS message id.
"start_event_id" to bundle.innerEvent.id,
"epoch" to ctx.marmot.currentEpoch(gid),
"brokers" to brokers,
) + RawEventSupport.ackFields(ack),
)
return 0
}
}
private suspend fun send(
dataDir: DataDir,
rest: Array<String>,
): Int {
val args = Args(rest)
val positional = args.positional
val streamId = args.flag("stream-id")
val startEventId = args.flag("start-event-id")
val broker = args.flag("broker")
if (positional.size < 2 || streamId == null || startEventId == null || broker == null) {
return Output.error("bad_args", "stream send GID --stream-id HEX --start-event-id HEX --broker URI TEXT…")
}
Context.open(dataDir).use { ctx ->
ctx.prepare()
val gid = ctx.resolveGroupId(positional[0])
if (!ctx.marmot.isMember(gid)) return Output.error("not_member", "not a member of group $gid")
// The stream's epoch is the one that carried its kind:1200, not
// whatever the group has reached by now — a commit between the
// start and the first record would otherwise put the publisher on
// a key no receiver derives.
val anchorEpoch =
args.flag("epoch")?.toLongOrNull()
?: ctx.marmot.storedEpochs(gid)[startEventId]
val crypto =
ctx.marmot.agentTextStreamCrypto(
nostrGroupId = gid,
streamId = streamId.hexToByteArray(),
startEventId = startEventId.hexToByteArray(),
epoch = anchorEpoch,
)
val publisher = AgentTextStreamPublisher.open(crypto, InMemoryAgentTextStreamSequenceStore())
val transport = QuicAgentTextStreamTransport(certificateValidator = PermissiveCertificateValidator())
val stream =
try {
transport.publish(broker, streamId.hexToByteArray(), startEventId.hexToByteArray())
} catch (e: Exception) {
return Output.error("broker_unreachable", "${e.message}")
}
try {
for (text in positional.drop(1)) {
stream.send(publisher.publish(AgentTextStreamRecordV1.TYPE_TEXT_DELTA, text.encodeToByteArray()))
}
stream.finish()
} finally {
stream.close()
}
Output.emit(
mapOf(
"group_id" to gid,
"stream_id" to streamId,
"start_event_id" to startEventId,
"records" to positional.size - 1,
"epoch" to crypto.context.mlsEpoch,
// What `stream finish` has to publish so a receiver can
// prove it saw this exact stream.
"transcript_hash" to publisher.transcript.hash.toHexKey(),
"chunk_count" to publisher.transcript.chunkCount,
),
)
return 0
}
}
private suspend fun watch(
dataDir: DataDir,
rest: Array<String>,
): Int {
val args = Args(rest)
val positional = args.positional
if (positional.isEmpty()) return Output.error("bad_args", "stream watch GID [--stream-id HEX] [--timeout SECS]")
val timeoutMs = (args.flag("timeout")?.toLongOrNull() ?: 30L) * 1000
Context.open(dataDir).use { ctx ->
ctx.prepare()
val gid = ctx.resolveGroupId(positional[0])
ctx.syncIncoming()
if (!ctx.marmot.isMember(gid)) return Output.error("not_member", "not a member of group $gid")
val wanted = args.flag("stream-id")
val anchor =
findStart(ctx, gid, wanted)
?: return Output.error("no_stream", "no kind:1200 stream start in group $gid")
val (startEvent, start) = anchor
if (!start.isTextProfile) {
return Output.error("unsupported_stream", "stream-type=${start.streamType} final-kind=${start.finalKind}")
}
if (!start.isQuicRoute) {
return Output.error("unsupported_route", "route=${start.route} — only the raw QUIC binding is implemented")
}
if (start.brokerCandidates.isEmpty()) {
return Output.error("no_candidate", "the start payload advertises no broker; the final kind:9 is the answer")
}
// The key context is the PUBLISHER's, not ours: the sender id and
// the epoch are theirs, and every member of that epoch derives the
// same record key from the group exporter.
//
// The epoch is the one that DELIVERED the anchor, not the group's
// current one. A commit landing between the start and the watch
// moves the group on, and deriving under the newer epoch produces
// a different key and an empty preview.
val anchorEpoch =
args.flag("epoch")?.toLongOrNull()
?: ctx.marmot.storedEpochs(gid)[startEvent.id]
val crypto =
ctx.marmot.agentTextStreamCrypto(
nostrGroupId = gid,
streamId = start.streamId.hexToByteArray(),
startEventId = startEvent.id.hexToByteArray(),
senderPubKey = startEvent.pubKey,
epoch = anchorEpoch,
)
val subscriber = AgentTextStreamSubscriber(crypto)
val transport = QuicAgentTextStreamTransport(certificateValidator = PermissiveCertificateValidator())
// "A receiver tries advertised candidates in listed order"; the
// first that yields the matching stream wins.
var lastError: String? = null
for (candidate in start.brokerCandidates) {
val stream =
try {
transport.subscribe(candidate, start.streamId.hexToByteArray(), startEvent.id.hexToByteArray())
} catch (e: Exception) {
lastError = "${e.message}"
continue
}
val outcomes = mutableMapOf<String, Int>()
try {
withTimeoutOrNull(timeoutMs) {
stream.incoming().collect { record ->
val outcome = subscriber.accept(record)
outcomes[outcome.name] = (outcomes[outcome.name] ?: 0) + 1
if (outcome == RecordOutcome.Accepted &&
(subscriber.status == PreviewStatus.FINISHED || subscriber.status == PreviewStatus.ABORTED)
) {
throw StreamComplete()
}
}
}
} catch (_: StreamComplete) {
// The publisher said the final message is coming.
} finally {
stream.close()
}
Output.emit(
mapOf(
"group_id" to gid,
"stream_id" to start.streamId,
"start_event_id" to startEvent.id,
"author" to startEvent.pubKey,
"broker" to candidate,
"preview" to subscriber.previewText,
"status" to subscriber.status.name,
"records" to subscriber.highWaterMark,
"transcript_hash" to subscriber.transcript.hash.toHexKey(),
"chunk_count" to subscriber.transcript.chunkCount,
"epoch" to crypto.context.mlsEpoch,
"outcomes" to outcomes,
"latest_status" to subscriber.latestStatus,
"latest_progress" to subscriber.latestProgress,
),
)
return 0
}
return Output.error("no_candidate_worked", lastError ?: "every advertised broker candidate was unusable")
}
}
private suspend fun finish(
dataDir: DataDir,
rest: Array<String>,
): Int {
val args = Args(rest)
val positional = args.positional
val streamId = args.flag("stream-id")
val transcriptHash = args.flag("transcript-hash")
val chunkCount = args.flag("chunk-count")?.toLongOrNull()
if (positional.size < 2 || streamId == null || transcriptHash == null || chunkCount == null) {
return Output.error(
"bad_args",
"stream finish GID --stream-id HEX --transcript-hash HEX --chunk-count N TEXT…",
)
}
Context.open(dataDir).use { ctx ->
ctx.prepare()
val gid = ctx.resolveGroupId(positional[0])
ctx.syncIncoming()
if (!ctx.marmot.isMember(gid)) return Output.error("not_member", "not a member of group $gid")
val text = positional.drop(1).joinToString(" ")
val bundle = ctx.marmot.buildAgentStreamFinal(gid, streamId, transcriptHash, chunkCount, text)
val targets = ctx.marmotGroupRelays(gid).ifEmpty { ctx.outboxRelays() }
val ack = ctx.publish(bundle.outbound.signedEvent, targets)
RawEventSupport.publishGuard(ack, bundle.outbound.signedEvent.id)?.let { return it }
Output.emit(
mapOf(
"group_id" to gid,
"stream_id" to streamId,
"final_event_id" to bundle.innerEvent.id,
"transcript_hash" to transcriptHash,
"chunk_count" to chunkCount,
"content" to text,
) + RawEventSupport.ackFields(ack),
)
return 0
}
}
/**
* The newest kind:1200 in the group's decrypted log, optionally pinned to
* one stream id. Newest wins because a group can carry many streams over
* its life and a watcher almost always means the current one.
*/
private suspend fun findStart(
ctx: Context,
nostrGroupId: String,
streamId: String?,
): Pair<Event, AgentTextStreamStart>? {
var best: Pair<Event, AgentTextStreamStart>? = null
for (line in ctx.marmot.loadStoredMessages(nostrGroupId)) {
val parsed = Event.fromJsonOrNull(line) ?: continue
val start = AgentTextStreamStart.fromTags(parsed.kind, parsed.tags) ?: continue
if (streamId != null && !start.streamId.equals(streamId, ignoreCase = true)) continue
if (best == null || parsed.createdAt >= best.first.createdAt) best = parsed to start
}
return best
}
/** Unwinds the collect loop once the publisher signalled the end. */
private class StreamComplete : RuntimeException(null, null, false, false)
}
@@ -151,7 +151,32 @@ class FileMarmotMessageStore(
override suspend fun delete(nostrGroupId: String) {
file(nostrGroupId).deleteOrWarn("FileMarmotMessageStore", "group messages")
epochFile(nostrGroupId).deleteOrWarn("FileMarmotMessageStore", "group message epochs")
}
private fun epochFile(id: String) = File(dir, "$id.epochs")
override suspend fun recordEpoch(
nostrGroupId: String,
innerEventId: String,
epoch: Long,
) {
val line = "$innerEventId $epoch"
val target = epochFile(nostrGroupId)
if (target.exists() && target.readLines().any { it == line }) return
SecureFileIO.appendText(target, line + "\n")
}
override suspend fun loadEpochs(nostrGroupId: String): Map<String, Long> =
epochFile(nostrGroupId)
.takeIf { it.exists() }
?.readLines()
?.mapNotNull { line ->
val parts = line.trim().split(' ')
if (parts.size != 2) return@mapNotNull null
val epoch = parts[1].toLongOrNull() ?: return@mapNotNull null
parts[0] to epoch
}?.toMap() ?: emptyMap()
}
/**
@@ -53,6 +53,15 @@ RELAY_BIN="$RELAY_REPO/target/release/nostr-rs-relay"
RELAY_DATA="$STATE_DIR/relay"
RELAY_PORT="${RELAY_PORT:-8080}"
RELAY_URL="ws://$RELAY_HOST:$RELAY_PORT"
# MDK's reference QUIC broker, for the agent-text-stream tests. Loopback like
# everything else; the tests skip when the binary was never built.
BROKER_BIN="$WN_REPO/target/release/marmot-quic-broker"
BROKER_HOST="${BROKER_HOST:-127.0.0.1}"
BROKER_PORT="${BROKER_PORT:-4455}"
BROKER_URI="quic://$BROKER_HOST:$BROKER_PORT"
BROKER_PID=""
NO_BUILD=0
# Every run starts from empty stores. wnd already wipes B's and C's data dirs
# on each start, but A's amy home and the relay's SQLite file used to survive,
@@ -129,6 +138,7 @@ cleanup() {
local rc=$?
trap - EXIT INT TERM HUP
stop_daemons
stop_quic_broker
stop_local_relay
print_summary
exit "$rc"
@@ -141,6 +151,7 @@ trap 'exit 129' HUP
banner "Marmot headless interop harness ($RUN_TS)"
preflight
start_local_relay
start_quic_broker || true
start_daemon B "$B_DIR" "$B_SOCKET"
start_daemon C "$C_DIR" "$C_SOCKET"
ensure_identity_a
@@ -166,6 +177,8 @@ ALL_TESTS=(
test_14_wn_removes_a
test_15_wn_member_leaves
test_16_wn_keypackage_rotation
test_18_agent_stream_amy_publishes
test_19_agent_stream_wn_publishes
)
# --tests runs a subset in the order given. Most tests read state a previous
+42
View File
@@ -120,6 +120,48 @@ preflight() {
info "relay bin: $RELAY_BIN"
}
# --- local QUIC broker -------------------------------------------------------
# MDK's own `marmot-quic-broker`, the reference implementation of the other
# side of `transports/quic.md`. Agent text stream previews are the only tests
# that need it, and they are the only way to know our binding is right — the
# ALPN, the control envelope, the frame prefix and the record key schedule all
# have to agree with an implementation that is not ours.
#
# `--replay-ttl-secs` is what lets a subscriber that connects after the
# records were pushed still see them; with the default 0 a test would have to
# race the publisher.
start_quic_broker() {
if [[ ! -x "$BROKER_BIN" ]]; then
info "marmot-quic-broker not built — agent text stream tests will skip"
return 1
fi
step "starting QUIC broker on $BROKER_HOST:$BROKER_PORT"
mkdir -p "$STATE_DIR/broker"
nohup "$BROKER_BIN" --bind "$BROKER_HOST:$BROKER_PORT" --replay-ttl-secs 60 --json \
>"$STATE_DIR/broker/stdout.log" 2>"$STATE_DIR/broker/stderr.log" &
BROKER_PID=$!
local deadline=$(( $(date +%s) + 15 ))
while [[ $(date +%s) -lt $deadline ]]; do
if grep -q '"local_addr"' "$STATE_DIR/broker/stdout.log" 2>/dev/null; then
info "broker pid $BROKER_PID ready"
return 0
fi
if ! kill -0 "$BROKER_PID" 2>/dev/null; then break; fi
sleep 1
done
fail_msg "broker never came up (see $STATE_DIR/broker/stderr.log)"
tail -n 20 "$STATE_DIR/broker/stderr.log" 2>/dev/null | sed 's/^/ /' >&2 || true
BROKER_PID=""
return 1
}
stop_quic_broker() {
[[ -n "${BROKER_PID:-}" ]] || return 0
step "stopping broker pid $BROKER_PID"
kill "$BROKER_PID" 2>/dev/null || true
BROKER_PID=""
}
# --- local relay -------------------------------------------------------------
# Start nostr-rs-relay on $RELAY_PORT with a minimal config. Every test
# runs against this one loopback endpoint — no external network traffic.
+150
View File
@@ -409,3 +409,153 @@ test_16_wn_keypackage_rotation() {
record_result "$id" fail "amy kept seeing the pre-rotation KP"
fi
}
# --- Agent text streams (0x8006) --------------------------------------------
# The live-preview half of an agent turn: a hidden kind:1200 anchors the
# stream over MLS, encrypted records ride raw QUIC through a broker, and a
# kind:9 closes it with the transcript a receiver checks its own fold against.
#
# Both tests need MDK's `marmot-quic-broker` — the reference implementation of
# the other side. Without it there is no honest way to claim the binding is
# right, so they skip rather than pretending.
test_18_agent_stream_amy_publishes() {
banner "Test 18 — amy publishes an agent text stream; wn verifies it"
local id="18 agent stream amy->wn"
if [[ -z "${BROKER_PID:-}" ]]; then record_result "$id" skip "no QUIC broker"; return; fi
# Its own group: by this point in the run A has left GROUP_02 and been
# removed from others, and a stream needs both parties actually present.
local out gid mls_gid
out=$(amy_json marmot group create --name "Interop-18") || {
record_result "$id" fail "amy group create failed"; return
}
gid=$(printf '%s' "$out" | jq -r '.group_id')
mls_gid=$(printf '%s' "$out" | jq -r '.mls_group_id')
amy_json marmot group add "$gid" "$B_NPUB" >/dev/null || {
record_result "$id" fail "amy group add B failed"; return
}
local b_gid
if ! b_gid=$(wait_for_invite B 60); then
record_result "$id" fail "B never received the invite"; return
fi
wn_b groups accept "$b_gid" >/dev/null 2>&1 || true
save_state GROUP_STREAM "$gid"
save_state GROUP_STREAM_MLS "$mls_gid"
local start_json sid seid
start_json=$(amy_json marmot stream start "$gid" --broker "$BROKER_URI") || {
record_result "$id" fail "amy stream start failed"; return
}
sid=$(printf '%s' "$start_json" | jq -r '.stream_id // empty')
seid=$(printf '%s' "$start_json" | jq -r '.start_event_id // empty')
if [[ -z "$sid" || -z "$seid" ]]; then
record_result "$id" fail "stream start reported no ids"; return
fi
local send_json thash chunks
send_json=$(amy_json marmot stream send "$gid" --stream-id "$sid" --start-event-id "$seid" \
--broker "$BROKER_URI" "Hello " "from " "amethyst") || {
record_result "$id" fail "amy stream send failed"; return
}
thash=$(printf '%s' "$send_json" | jq -r '.transcript_hash // empty')
chunks=$(printf '%s' "$send_json" | jq -r '.chunk_count // empty')
printf 'stream18 start=%s\nstream18 send=%s\n' "$start_json" "$send_json" >>"$LOG_FILE"
# Our own subscriber must recover the stream from the broker's replay window
# and fold it to the same transcript the publisher computed.
local watch_json
watch_json=$(amy_json marmot stream watch "$gid" --stream-id "$sid" --timeout 15) || {
record_result "$id" fail "amy stream watch failed"; return
}
printf 'stream18 watch=%s\n' "$watch_json" >>"$LOG_FILE"
if [[ "$(printf '%s' "$watch_json" | jq -r '.transcript_hash')" != "$thash" ]]; then
record_result "$id" fail "amy's own fold disagrees with what it published"; return
fi
if [[ "$(printf '%s' "$watch_json" | jq -r '.preview')" != "Hello from amethyst" ]]; then
record_result "$id" fail "preview text did not survive the round trip"; return
fi
amy_json marmot stream finish "$gid" --stream-id "$sid" \
--transcript-hash "$thash" --chunk-count "$chunks" "Hello from amethyst" >/dev/null || {
record_result "$id" fail "amy stream finish failed"; return
}
# The real check: MDK reads our kind:1200 + kind:9 and confirms the
# transcript itself.
local deadline=$(( $(date +%s) + 60 )) verified="false"
while [[ $(date +%s) -lt $deadline ]]; do
verified=$(wn_b --json stream verify "$mls_gid" --stream-id "$sid" --transcript-hash "$thash" 2>/dev/null \
| jq -r '.result.verified // false')
[[ "$verified" == "true" ]] && break
sleep 3
done
if [[ "$verified" == "true" ]]; then
record_result "$id" pass
else
record_result "$id" fail "wn could not verify amy's transcript"
fi
}
test_19_agent_stream_wn_publishes() {
banner "Test 19 — wn publishes an agent text stream; amy watches it"
local id="19 agent stream wn->amy"
local gid mls_gid
gid=$(load_state GROUP_STREAM || true)
mls_gid=$(load_state GROUP_STREAM_MLS || true)
if [[ -z "${gid:-}" ]]; then record_result "$id" skip "no stream group (test 18 did not run)"; return; fi
if [[ -z "${BROKER_PID:-}" ]]; then record_result "$id" skip "no QUIC broker"; return; fi
local wn_start wsid wseid
wn_start=$(wn_b --json stream start "$mls_gid" --quic-candidate "$BROKER_URI" 2>>"$LOG_FILE") || {
record_result "$id" fail "wn stream start failed"; return
}
wsid=$(printf '%s' "$wn_start" | jq -r '.result.stream_id // empty')
# wn reports the kind:1200's own Marmot app event id as message_ids[0] —
# the same value amy resolves as start_event_id from the payload itself.
wseid=$(printf '%s' "$wn_start" | jq -r '.result.message_ids[0] // empty')
if [[ -z "$wsid" || -z "$wseid" ]]; then
record_result "$id" fail "wn stream start reported no ids"; return
fi
# amy has to have the kind:1200 before it can derive the stream's keys: the
# anchor's own event id is part of the key context. Sync until it lands.
local anchor_deadline=$(( $(date +%s) + 45 )) saw_anchor=0
while [[ $(date +%s) -lt $anchor_deadline ]]; do
if amy_a marmot message list "$gid" --limit 50 2>/dev/null \
| jq -e --arg id "$wseid" '.messages[]? | select(.event_id == $id)' >/dev/null 2>&1; then
saw_anchor=1; break
fi
sleep 3
done
if [[ "$saw_anchor" -ne 1 ]]; then
record_result "$id" fail "amy never received wn's kind:1200 anchor"; return
fi
local watch_out="$STATE_DIR/stream-19-watch.json"
( amy_a marmot stream watch "$gid" --stream-id "$wsid" --timeout 25 >"$watch_out" 2>>"$LOG_FILE" ) &
local watch_pid=$!
sleep 4
local send_json wthash
send_json=$(wn_b --json stream send --broker --connect "$BROKER_HOST:$BROKER_PORT" --insecure-local \
--stream-id "$wsid" --start-event-id "$wseid" "Hello from whitenoise" 2>>"$LOG_FILE")
wthash=$(printf '%s' "$send_json" | jq -r '.result.transcript_hash // empty')
wait "$watch_pid" || true
local preview athash
preview=$(tail -n 1 "$watch_out" 2>/dev/null | jq -r '.preview // empty')
athash=$(tail -n 1 "$watch_out" 2>/dev/null | jq -r '.transcript_hash // empty')
if [[ "$preview" != "Hello from whitenoise" ]]; then
record_result "$id" fail "amy rendered '$preview' instead of wn's text"; return
fi
# The decisive one: our record key schedule, key context, AEAD and transcript
# construction all have to match MDK's exactly for these to agree.
if [[ -n "$wthash" && "$athash" != "$wthash" ]]; then
record_result "$id" fail "transcript hash disagrees with wn's (${athash:0:12}… vs ${wthash:0:12}…)"; return
fi
record_result "$id" pass
}
@@ -162,7 +162,7 @@ private suspend fun MarmotManager.ingestGroupEvent(ge: GroupEvent): MarmotIngest
is GroupEventResult.ApplicationMessage -> {
// MLS ratchets once we decrypt; future reads of the same ciphertext
// would fail — persist the plaintext now so restarts/replays see it.
persistDecryptedMessage(result.groupId, result.innerEventJson)
persistDecryptedMessage(result.groupId, result.innerEventJson, result.epoch)
MarmotIngestResult.Message(result)
}
@@ -37,6 +37,10 @@ import com.vitorpamplona.quartz.marmot.appComponents.GroupBlossomImageV1
import com.vitorpamplona.quartz.marmot.appComponents.GroupProfileV1
import com.vitorpamplona.quartz.marmot.appComponents.MarmotGroupState
import com.vitorpamplona.quartz.marmot.appComponents.MessageRetentionV1
import com.vitorpamplona.quartz.marmot.appComponents.agentTextStream.AgentTextStreamCrypto
import com.vitorpamplona.quartz.marmot.appComponents.agentTextStream.AgentTextStreamFinal
import com.vitorpamplona.quartz.marmot.appComponents.agentTextStream.AgentTextStreamKeyContextV1
import com.vitorpamplona.quartz.marmot.appComponents.agentTextStream.AgentTextStreamStart
import com.vitorpamplona.quartz.marmot.mip00KeyPackages.KeyPackageBundleStore
import com.vitorpamplona.quartz.marmot.mip00KeyPackages.KeyPackageEvent
import com.vitorpamplona.quartz.marmot.mip00KeyPackages.KeyPackageRotationManager
@@ -458,6 +462,119 @@ class MarmotManager(
return TextMessageBundle(outbound = outbound, innerEvent = innerEvent)
}
/**
* Build the hidden kind:1200 payload that anchors one agent text stream.
*
* The start payload is what makes a live preview renderable at all: its
* own Marmot app event id goes into [AgentTextStreamKeyContextV1], so a
* receiver can only derive record keys for a stream it has already seen
* announced inside the group. That is also why the anchor is an ordinary
* in-group payload rather than something the broker hands out — the broker
* relays ciphertext and learns nothing.
*
* The returned bundle's `innerEvent.id` IS the `start_event_id`; a caller
* needs it before it can derive the stream's crypto, which is why this
* returns the built payload instead of publishing and forgetting it.
*
* @param brokerCandidates `quic://host:port` endpoints a receiver may try,
* in preference order. Zero is valid — the preview is then unavailable
* and every member still gets the final message.
* @param parentEventId the prompt this stream answers, when there is one.
*/
suspend fun buildAgentStreamStart(
nostrGroupId: HexKey,
streamId: HexKey,
brokerCandidates: List<String> = emptyList(),
parentEventId: HexKey? = null,
persistOwn: Boolean = true,
): TextMessageBundle {
val template =
com.vitorpamplona.quartz.nip01Core.signers
.eventTemplate<Event>(kind = AgentTextStreamStart.KIND, description = "") {
AgentTextStreamStart
.tags(streamId, brokerCandidates, parentEventId = parentEventId)
.forEach { addUnique(it) }
}
val innerEvent =
com.vitorpamplona.quartz.nip59Giftwrap.rumors.RumorAssembler
.assembleRumor<Event>(signer.pubKey, template)
// The epoch this went out at is the one the stream's key context binds,
// so remember it the same way an inbound anchor's epoch is remembered.
val epoch = currentEpoch(nostrGroupId)
val outbound = buildGroupMessage(nostrGroupId, innerEvent)
if (persistOwn) persistDecryptedMessage(nostrGroupId, innerEvent.toJson(), epoch)
return TextMessageBundle(outbound = outbound, innerEvent = innerEvent)
}
/**
* Build the durable kind:9 that closes an agent text stream out.
*
* This is the authoritative message. A receiver that rendered a preview
* compares its own fold against [transcriptHash] / [chunkCount]: agreement
* means it saw exactly the stream the publisher sent, disagreement means
* records were dropped, reordered or injected even though each one opened.
* A receiver that skipped the preview just reads this as normal chat.
*/
suspend fun buildAgentStreamFinal(
nostrGroupId: HexKey,
streamId: HexKey,
transcriptHash: HexKey,
chunkCount: Long,
text: String,
persistOwn: Boolean = true,
): TextMessageBundle {
val template =
com.vitorpamplona.quartz.nip01Core.signers
.eventTemplate<Event>(kind = 9, description = text) {
AgentTextStreamFinal
.tags(streamId, transcriptHash, chunkCount)
.forEach { addUnique(it) }
}
val innerEvent =
com.vitorpamplona.quartz.nip59Giftwrap.rumors.RumorAssembler
.assembleRumor<Event>(signer.pubKey, template)
val outbound = buildGroupMessage(nostrGroupId, innerEvent)
if (persistOwn) persistDecryptedMessage(nostrGroupId, innerEvent.toJson())
return TextMessageBundle(outbound = outbound, innerEvent = innerEvent)
}
/**
* The record AEAD for one stream in [nostrGroupId].
*
* The stream secret is the group's own
* `MLS-Exporter("marmot", "agent-text-stream-quic", 32)`, so every member
* of the epoch derives the same one and no key ever crosses the wire. All
* the per-stream separation comes from the key context: change the stream,
* the epoch, the sender or the anchoring kind:1200 event and the record
* key changes with it.
*
* [senderPubKey] is the stream's author, which is not necessarily us — a
* receiver derives the publisher's context, not its own.
*/
fun agentTextStreamCrypto(
nostrGroupId: HexKey,
streamId: ByteArray,
startEventId: ByteArray,
senderPubKey: HexKey = signer.pubKey,
epoch: Long? = null,
): AgentTextStreamCrypto {
val group = groupManager.getGroup(nostrGroupId) ?: error("not a member of group $nostrGroupId")
return AgentTextStreamCrypto(
streamSecret = group.agentTextStreamSecret(),
context =
AgentTextStreamKeyContextV1(
groupId = group.groupId,
streamId = streamId,
mlsEpoch = epoch ?: group.epoch,
senderId = senderPubKey.hexToByteArray(),
startEventId = startEventId,
),
)
}
/** The group's current MLS epoch, which the stream key context binds. */
fun currentEpoch(nostrGroupId: HexKey): Long? = groupManager.getGroup(nostrGroupId)?.epoch
/**
* Build a kind:5 deletion inner event targeting one or more prior inner
* events in the same group. Unsigned rumor (MIP-03); e-tag + k-tag for
@@ -907,14 +1024,33 @@ class MarmotManager(
suspend fun persistDecryptedMessage(
nostrGroupId: HexKey,
innerEventJson: String,
/**
* The MLS epoch that delivered this payload, when the caller knows it.
* Only agent text streams read it back — their record key context
* binds the epoch, so a receiver has to derive keys under the epoch
* that carried the stream's anchor rather than the group's current one.
*/
epoch: Long? = null,
) {
try {
messageStore?.appendMessage(nostrGroupId, innerEventJson)
if (epoch != null) {
Event.fromJsonOrNull(innerEventJson)?.let { messageStore?.recordEpoch(nostrGroupId, it.id, epoch) }
}
} catch (e: Exception) {
Log.w("MarmotManager", "Failed to persist Marmot message for $nostrGroupId", e)
}
}
/** Inner event id → delivering MLS epoch, for whatever the store kept. */
suspend fun storedEpochs(nostrGroupId: HexKey): Map<String, Long> =
try {
messageStore?.loadEpochs(nostrGroupId) ?: emptyMap()
} catch (e: Exception) {
Log.w("MarmotManager", "Failed to read Marmot message epochs for $nostrGroupId", e)
emptyMap()
}
/**
* Load all persisted inner event JSONs for a group, in append order.
* Returns an empty list if no message store is configured or none exist.
+15 -3
View File
@@ -62,11 +62,23 @@ then:
Without the property the cases skip visibly, so an ordinary `./gradlew test`
never needs a broker on the machine.
## Using it
`amy marmot stream` drives the whole feature; the harness's tests 18 and 19
run it in both directions against MDK.
```bash
amy marmot stream start GID --broker quic://127.0.0.1:4450
amy marmot stream send GID --stream-id … --start-event-id … --broker … "hello"
amy marmot stream watch GID --stream-id …
amy marmot stream finish GID --stream-id … --transcript-hash … --chunk-count N "hello"
```
## Not done
- Nothing in the app yet mints a kind-1200 start payload, chooses a broker
candidate, or renders a live preview — this is the transport, not the
feature wiring.
- The GUIs do not originate or render a stream yet, which is why the `send`
(`0xF2D2`) and `fanout` (`0xF2D4`) role capabilities stay unadvertised — a
role is a promise to the whole group.
- The direct path (`marmot.quic_stream.v1`) is unimplemented. v1 defines no
start-payload candidate format for it, so it is only reachable with an
endpoint known out of band.
@@ -24,6 +24,9 @@ import com.vitorpamplona.quartz.marmot.appComponents.agentTextStream.AgentTextSt
import com.vitorpamplona.quartz.marmot.appComponents.agentTextStream.transport.AgentTextStreamFraming
import com.vitorpamplona.quartz.marmot.appComponents.agentTextStream.transport.BrokerControlType
import com.vitorpamplona.quartz.marmot.appComponents.agentTextStream.transport.MarmotQuicAlpn
import com.vitorpamplona.quartz.marmot.appComponents.agentTextStream.transport.MarmotQuicException
import com.vitorpamplona.quartz.marmot.appComponents.agentTextStream.transport.MarmotQuicStream
import com.vitorpamplona.quartz.marmot.appComponents.agentTextStream.transport.MarmotQuicTransport
import com.vitorpamplona.quartz.marmot.appComponents.agentTextStream.transport.QuicBrokerControlEnvelopeV1
import com.vitorpamplona.quartz.marmot.appComponents.agentTextStream.transport.QuicEndpointCandidate
import com.vitorpamplona.quic.connection.QuicConnection
@@ -35,6 +38,7 @@ import com.vitorpamplona.quic.transport.UdpSocket
import kotlinx.coroutines.CoroutineScope
import kotlinx.coroutines.Dispatchers
import kotlinx.coroutines.SupervisorJob
import kotlinx.coroutines.delay
import kotlinx.coroutines.flow.Flow
import kotlinx.coroutines.flow.flow
import kotlinx.coroutines.withTimeoutOrNull
@@ -175,6 +179,7 @@ private class QuicStreamDelivery(
private val stream: QuicStream,
private val driver: QuicConnectionDriver,
maxPlaintextFrameLen: Long?,
private val flushTimeoutMillis: Long = DEFAULT_FLUSH_TIMEOUT_MILLIS,
) : MarmotQuicStream {
private val reader = AgentTextStreamFraming.Reader(maxPlaintextFrameLen)
@@ -190,12 +195,45 @@ private class QuicStreamDelivery(
}
}
/**
* FIN our write side and wait until the peer acknowledges it.
*
* The wait is the point. `enqueue` only puts bytes in the send buffer; the
* driver still has to put them on the wire and the peer still has to ACK
* them. Returning before that and letting the caller [close] tears the
* connection down with records still buffered, and they are simply lost —
* silently, because the publisher already counted them. QUIC only ACKs a
* FIN once everything ahead of it arrived, so `finAcked` is exactly the
* "the broker has all of it" signal.
*/
override suspend fun finish() {
stream.send.finish()
driver.wakeup()
withTimeoutOrNull(flushTimeoutMillis) {
while (!stream.send.finAcked) {
driver.wakeup()
delay(FLUSH_POLL_MILLIS)
}
}
}
/**
* Tear down the connection. A caller that wrote records is expected to
* [finish] first; this still gives an unacknowledged FIN a bounded moment
* rather than dropping the tail of a stream on the floor.
*/
override suspend fun close() {
if (stream.send.finSent && !stream.send.finAcked) {
withTimeoutOrNull(flushTimeoutMillis) {
while (!stream.send.finAcked) {
driver.wakeup()
delay(FLUSH_POLL_MILLIS)
}
}
}
driver.close()
}
}
private const val DEFAULT_FLUSH_TIMEOUT_MILLIS = 10_000L
private const val FLUSH_POLL_MILLIS = 20L
@@ -26,6 +26,7 @@ import com.vitorpamplona.quartz.marmot.appComponents.agentTextStream.AgentTextSt
import com.vitorpamplona.quartz.marmot.appComponents.agentTextStream.AgentTextStreamRecordV1
import com.vitorpamplona.quartz.marmot.appComponents.agentTextStream.AgentTextStreamTranscriptV1
import com.vitorpamplona.quartz.marmot.appComponents.agentTextStream.InMemoryAgentTextStreamSequenceStore
import com.vitorpamplona.quartz.marmot.appComponents.agentTextStream.transport.MarmotQuicException
import com.vitorpamplona.quic.tls.PermissiveCertificateValidator
import kotlinx.coroutines.CoroutineScope
import kotlinx.coroutines.Dispatchers
+37 -9
View File
@@ -695,11 +695,12 @@ test we have.
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`). Not for
want of a transport any more — see below — but because nothing in the app yet
originates a stream, and a role we do not 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.
We still do NOT advertise `send` (`0xF2D2`) or `fanout` (`0xF2D4`) in
KeyPackage capabilities. Everything behind them now works end to end, but the
roles are a promise to a whole group and the GUIs do not yet originate or
render a stream — only the CLI does. 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 binding is implemented and verified against MDK's
broker.** `transports/quic.md` is a RAW QUIC binding — its own ALPNs
@@ -718,8 +719,35 @@ test we have.
publisher's transcript hash — plus that the broker keeps rooms apart. Opt in
with `-DmarmotQuicBroker=host:port`; it skips visibly without one.
What is left is the application wiring: nothing yet mints a kind-1200 start,
picks a broker candidate, or renders a live preview. The direct path
(`marmot.quic_stream.v1`) is also unimplemented — v1 has no start-payload
candidate format for it, so it is only usable with an out-of-band endpoint.
The direct path (`marmot.quic_stream.v1`) is unimplemented — v1 has no
start-payload candidate format for it, so it is only usable with an
out-of-band endpoint.
- **The feature is wired end to end, both directions, against MDK.**
`amy marmot stream start|send|watch|finish` mints the kind-1200 anchor,
pushes records through a broker, folds a preview under the receive discipline
(`seq` high-water mark, silent replay discard, gap → unverifiable) and
publishes the authoritative kind-9 carrying the transcript. Harness tests 18
and 19 run both directions: MDK's `wn stream verify` confirms our transcript
from our own kind-1200 + kind-9, and our subscriber folds MDK's stream to a
transcript hash identical to the one `wn stream send` computed. That equality
is the whole key schedule, key context, AEAD, framing and transcript
construction agreeing with an implementation that is not ours.
Two defects only that exercise could have found:
- **The epoch belongs to the stream, not to the clock.** The record key
context binds `mls_epoch`, and we were resolving it as "the group's current
epoch" at each command. A commit landing between the start and the send (or
the watch) put the two sides on different keys and produced an empty
preview. The epoch that DELIVERED the kind-1200 is the stream's, so it is
persisted with the message now (`MarmotMessageStore.recordEpoch`, which
MDK's own storage has as `source_epoch`) and read back by both sides.
- **`close()` dropped the tail of a stream.** `enqueue` only fills the send
buffer; tearing the connection down before the driver flushed it lost
records silently, because the publisher had already counted them. QUIC ACKs
a FIN only after everything ahead of it arrived, so `finish()` now waits for
`finAcked` and `close()` gives an unacknowledged FIN a bounded moment. This
is why the test passed alone and failed inside a full run.
@@ -35,10 +35,18 @@ import com.vitorpamplona.quartz.nip01Core.core.HexKey
*
* ```
* ["stream", <32-byte stream id, hex>]
* ["stream-type", "text"]
* ["final-kind", "9"]
* ["route", "quic"] (optional; "quic" when absent)
* ["parent", <prompt event id>] (optional)
* ["broker", <candidate URL>] (repeatable)
* ```
*
* `stream`, `stream-type` and `final-kind` are owned by the feature;
* `route` and `broker` belong to the transport binding, and a client that
* implements `receive` but not raw QUIC may ignore them entirely and wait for
* the final message.
*
* The matching END of a stream is an ordinary kind-9 chat carrying
* [STREAM_TAG], [STREAM_HASH_TAG] and [STREAM_CHUNKS_TAG]; see
* [AgentTextStreamFinal].
@@ -47,33 +55,69 @@ class AgentTextStreamStart(
val streamId: HexKey,
val route: String,
val brokerCandidates: List<String>,
val streamType: String = TYPE_TEXT,
val finalKind: Int = FINAL_KIND_TEXT,
val parentEventId: HexKey? = null,
) {
val isQuicRoute: Boolean get() = route == ROUTE_QUIC
/**
* True when this start is one we know how to render live. A `stream-type`
* we don't implement is not an error — the final message still arrives —
* so callers skip the preview rather than rejecting the payload.
*/
val isTextProfile: Boolean get() = streamType == TYPE_TEXT && finalKind == FINAL_KIND_TEXT
/**
* "Receivers MUST ignore a final payload whose kind does not match the
* start payload's `final-kind`."
*/
fun acceptsFinalKind(kind: Int): Boolean = kind == finalKind
companion object {
const val KIND = 1200
const val STREAM_TAG = "stream"
const val STREAM_TYPE_TAG = "stream-type"
const val FINAL_KIND_TAG = "final-kind"
const val ROUTE_TAG = "route"
const val PARENT_TAG = "parent"
const val BROKER_TAG = "broker"
const val ROUTE_QUIC = "quic"
/** The first — and so far only — stream type. */
const val TYPE_TEXT = "text"
/** For `stream-type=text`, `final-kind` MUST be 9. */
const val FINAL_KIND_TEXT = 9
fun tags(
streamId: HexKey,
brokerCandidates: List<String>,
route: String = ROUTE_QUIC,
streamType: String = TYPE_TEXT,
finalKind: Int = FINAL_KIND_TEXT,
parentEventId: HexKey? = null,
): Array<Array<String>> =
buildList {
add(arrayOf(STREAM_TAG, streamId))
add(arrayOf(STREAM_TYPE_TAG, streamType))
add(arrayOf(FINAL_KIND_TAG, finalKind.toString()))
add(arrayOf(ROUTE_TAG, route))
if (parentEventId != null) add(arrayOf(PARENT_TAG, parentEventId))
brokerCandidates.forEach { add(arrayOf(BROKER_TAG, it)) }
}.toTypedArray()
/**
* Read the start view from a kind-1200 payload's tags, or null when
* this is not a stream start. A missing `route` means `quic`; a
* missing `stream` tag means the payload is not usable as an anchor at
* all, so it is rejected rather than defaulted.
* this is not a stream start.
*
* A missing `stream` tag means the payload cannot anchor anything, so
* it is rejected rather than defaulted. The rest default to the first
* text profile: an older or terser sender that omits `stream-type` /
* `final-kind` / `route` still describes exactly that profile, and
* refusing it would drop a stream we can render.
*/
fun fromTags(
kind: Int,
@@ -82,8 +126,13 @@ class AgentTextStreamStart(
if (kind != KIND) return null
val streamId = tags.firstOrNull { it.size >= 2 && it[0] == STREAM_TAG }?.get(1) ?: return null
val route = tags.firstOrNull { it.size >= 2 && it[0] == ROUTE_TAG }?.get(1) ?: ROUTE_QUIC
val streamType = tags.firstOrNull { it.size >= 2 && it[0] == STREAM_TYPE_TAG }?.get(1) ?: TYPE_TEXT
val finalKind =
tags.firstOrNull { it.size >= 2 && it[0] == FINAL_KIND_TAG }?.get(1)?.toIntOrNull()
?: FINAL_KIND_TEXT
val parent = tags.firstOrNull { it.size >= 2 && it[0] == PARENT_TAG }?.get(1)
val brokers = tags.filter { it.size >= 2 && it[0] == BROKER_TAG }.map { it[1] }
return AgentTextStreamStart(streamId, route, brokers)
return AgentTextStreamStart(streamId, route, brokers, streamType, finalKind, parent)
}
}
}
@@ -0,0 +1,189 @@
/*
* 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
/** What the subscriber did with one inbound record. */
enum class RecordOutcome {
/** Folded into the transcript and applied to the preview. */
Accepted,
/** At or below the high-water mark. Discarded silently; not stream-fatal. */
Replay,
/** Ahead of the high-water mark: records are missing and were not folded. */
Gap,
/** Failed its AEAD. Not ours, or altered in flight. */
Undecryptable,
/** Belongs to a different stream than the one being rendered. */
WrongStream,
}
/** How much a renderer may claim about the preview it is showing. */
enum class PreviewStatus {
/** Every record so far arrived in order and opened. */
LIVE,
/** A gap could not be backfilled, so the transcript hash can never complete. */
UNVERIFIABLE,
/** The publisher aborted; there is no durable text coming from this preview. */
ABORTED,
/** The publisher signalled that the final MLS message is on its way. */
FINISHED,
}
/**
* The receiving half of one agent text stream preview.
*
* Everything here is provisional. The authority is the final kind-9 MLS
* message, and [matchesFinal] is the only thing that turns "we rendered
* something" into "we rendered what the publisher sent" — every record can
* open individually and the stream still be wrong, if one was dropped,
* reordered or injected.
*
* A renderer must show this text as visibly distinct from confirmed content
* until that check passes.
*/
class AgentTextStreamSubscriber(
private val crypto: AgentTextStreamCrypto,
) {
private val builder = StringBuilder()
/** The stream's rolling transcript over every accepted record. */
val transcript: AgentTextStreamTranscriptV1 =
AgentTextStreamTranscriptV1.start(crypto.context.streamId, crypto.context.startEventId)
/** Highest `seq` folded so far; the next accepted record is this plus one. */
var highWaterMark: Long = 0
private set
var status: PreviewStatus = PreviewStatus.LIVE
private set
/** Provisional answer text: `TextDelta` appended, `Checkpoint` replacing. */
val previewText: String get() = builder.toString()
/** Latest `Status` label, for local UI chrome only. Never part of the answer. */
var latestStatus: String? = null
private set
/** Latest `ProgressDelta`, for live non-chat progress chrome only. */
var latestProgress: String? = null
private set
/**
* True once a gap was seen and not yet backfilled. The transcript is
* incomplete, so it can never be compared against the final message.
*/
private var sawUnbackfilledGap = false
/**
* Fold one record in, judged against the high-water mark.
*
* Order is not negotiable: `seq` is XORed into the record nonce, so an
* out-of-order fold would also produce a transcript nobody else computes.
* Hence a record ahead of the mark is reported rather than applied — the
* caller backfills it from a replay source, or lives with an unverifiable
* preview and waits for the final message.
*/
fun accept(record: AgentTextStreamRecordV1): RecordOutcome {
if (!record.streamId.contentEquals(crypto.context.streamId)) return RecordOutcome.WrongStream
// "A record whose seq is at or below the high-water mark — for example,
// a record a broker replays on reconnect — MUST be discarded silently
// without affecting the stream."
if (record.seq <= highWaterMark) return RecordOutcome.Replay
if (record.seq > highWaterMark + 1) {
sawUnbackfilledGap = true
if (status == PreviewStatus.LIVE) status = PreviewStatus.UNVERIFIABLE
return RecordOutcome.Gap
}
// Opening before advancing anything: a record that fails its AEAD must
// leave the stream exactly as it was, or a single injected frame could
// burn the sequence value the real record needs.
val opened = crypto.openOrNull(record) ?: return RecordOutcome.Undecryptable
highWaterMark = record.seq
transcript.append(opened)
when (opened.recordType) {
AgentTextStreamRecordV1.TYPE_TEXT_DELTA -> builder.append(opened.frame.decodeToString())
// "Receivers that support checkpoints replace the provisional
// preview text with the checkpoint plaintext, then continue
// applying later TextDelta records."
AgentTextStreamRecordV1.TYPE_CHECKPOINT -> {
builder.setLength(0)
builder.append(opened.frame.decodeToString())
}
AgentTextStreamRecordV1.TYPE_STATUS -> latestStatus = opened.frame.decodeToString()
AgentTextStreamRecordV1.TYPE_PROGRESS_DELTA -> latestProgress = opened.frame.decodeToString()
AgentTextStreamRecordV1.TYPE_ABORT -> {
// "Receivers remove or mark the preview as cancelled and wait
// for later durable events." Nothing here becomes chat text.
builder.setLength(0)
status = PreviewStatus.ABORTED
}
AgentTextStreamRecordV1.TYPE_FINAL_NOTICE ->
if (status == PreviewStatus.LIVE) status = PreviewStatus.FINISHED
// An unknown record type still counts toward the transcript — the
// hash covers the stream the publisher sent, not the subset we
// happen to understand — but contributes nothing to the preview.
else -> Unit
}
// A backfill that closes the last gap makes the preview whole again.
if (sawUnbackfilledGap && status == PreviewStatus.UNVERIFIABLE) {
sawUnbackfilledGap = false
status = PreviewStatus.LIVE
}
return RecordOutcome.Accepted
}
/**
* Does what we folded match what the final kind-9 says the publisher sent?
*
* A disagreement means records were dropped, reordered or injected even
* though each one opened, so the preview must be discarded in favour of
* the final message. A preview already known to be incomplete answers
* false without comparing: it cannot have folded the stream, whatever its
* running hash happens to be.
*/
fun matchesFinal(
transcriptHash: ByteArray,
chunkCount: Long,
): Boolean {
if (sawUnbackfilledGap) return false
if (status == PreviewStatus.ABORTED) return false
return transcript.chunkCount == chunkCount && transcript.hash.contentEquals(transcriptHash)
}
}
@@ -18,7 +18,7 @@
* 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.marmotquic
package com.vitorpamplona.quartz.marmot.appComponents.agentTextStream.transport
import com.vitorpamplona.quartz.marmot.appComponents.agentTextStream.AgentTextStreamRecordV1
import kotlinx.coroutines.flow.Flow
@@ -68,4 +68,27 @@ interface MarmotMessageStore {
* @param nostrGroupId hex-encoded Nostr group ID
*/
suspend fun delete(nostrGroupId: String)
/**
* Remember which MLS epoch delivered [innerEventId].
*
* Almost nothing needs this — the inner event is the message. Agent text
* streams do: their record key context binds `mls_epoch`, so a receiver
* that derives keys for a stream must use the epoch that carried the
* stream's kind:1200 anchor, not whatever epoch the group has reached by
* the time someone watches. A commit landing in between would otherwise
* silently produce a different key and an empty preview.
*
* Optional: a store that does not keep it simply cannot render a live
* preview for a stream anchored in an older epoch, which is a degraded
* feature and not a broken group.
*/
suspend fun recordEpoch(
nostrGroupId: String,
innerEventId: String,
epoch: Long,
) = Unit
/** Inner event id → the MLS epoch that delivered it, for what was recorded. */
suspend fun loadEpochs(nostrGroupId: String): Map<String, Long> = emptyMap()
}
@@ -0,0 +1,242 @@
/*
* 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.AgentTextStreamSubscriber
import com.vitorpamplona.quartz.marmot.appComponents.agentTextStream.InMemoryAgentTextStreamSequenceStore
import com.vitorpamplona.quartz.marmot.appComponents.agentTextStream.PreviewStatus
import com.vitorpamplona.quartz.marmot.appComponents.agentTextStream.RecordOutcome
import kotlinx.coroutines.runBlocking
import org.junit.Assert.assertEquals
import org.junit.Assert.assertFalse
import org.junit.Assert.assertTrue
import org.junit.Test
/**
* The receive half of `transports/quic.md`, which is where a preview either
* stays honest or quietly stops being one.
*
* The rules it has to hold: `seq` is accepted at most once and never folded
* out of order; a replayed record — which a broker WILL send after a
* reconnect, from the start of its replay window — is discarded silently and
* is never stream-fatal; a gap that cannot be backfilled makes the preview
* unverifiable, because the transcript hash can no longer be completed and
* therefore can no longer be checked against the final MLS message.
*/
class AgentTextStreamSubscriberTest {
private val streamId = ByteArray(32) { 0x31 }
private val startEventId = ByteArray(32) { 0x32 }
private val secret = ByteArray(32) { 0x33 }
private val crypto =
AgentTextStreamCrypto(
secret,
AgentTextStreamKeyContextV1(
groupId = ByteArray(32) { 0x34 },
streamId = streamId,
mlsEpoch = 11,
senderId = ByteArray(32) { 0x35 },
startEventId = startEventId,
),
)
/** Sealed records 1..n of a stream, as the publisher would have sent them. */
private fun sealedRecords(vararg texts: String): List<AgentTextStreamRecordV1> =
runBlocking {
val publisher = AgentTextStreamPublisher.open(crypto, InMemoryAgentTextStreamSequenceStore())
texts.map { publisher.publish(AgentTextStreamRecordV1.TYPE_TEXT_DELTA, it.encodeToByteArray()) }
}
private fun subscriber() = AgentTextStreamSubscriber(crypto)
@Test
fun textDeltasConcatenateIntoTheProvisionalPreview() {
val sub = subscriber()
sealedRecords("the ", "quick ", "brown fox").forEach {
assertEquals(RecordOutcome.Accepted, sub.accept(it))
}
assertEquals("the quick brown fox", sub.previewText)
assertEquals(PreviewStatus.LIVE, sub.status)
assertEquals(3L, sub.transcript.chunkCount)
}
@Test
fun aReplayedRecordIsDiscardedSilentlyAndChangesNothing() {
val records = sealedRecords("a", "b")
val sub = subscriber()
records.forEach { sub.accept(it) }
val hashBefore = sub.transcript.hash.toList()
// A broker replays its backlog from the start of the window on
// reconnect. Every one of these is at or below the high-water mark.
records.forEach {
assertEquals(
"a replayed record is never stream-fatal",
RecordOutcome.Replay,
sub.accept(it),
)
}
assertEquals("ab", sub.previewText)
assertEquals(2L, sub.transcript.chunkCount)
assertEquals(hashBefore, sub.transcript.hash.toList())
assertEquals(PreviewStatus.LIVE, sub.status)
}
@Test
fun aGapMakesThePreviewUnverifiableRatherThanWrong() {
val records = sealedRecords("one", "two", "three")
val sub = subscriber()
sub.accept(records[0])
// seq 2 never arrives. Folding seq 3 anyway would produce a transcript
// hash that cannot match the publisher's, so the preview is marked
// instead — the final MLS message is still authoritative.
assertEquals(RecordOutcome.Gap, sub.accept(records[2]))
assertEquals(PreviewStatus.UNVERIFIABLE, sub.status)
assertEquals("one", sub.previewText)
assertEquals(1L, sub.transcript.chunkCount)
}
@Test
fun aGapCanBeBackfilledBeforeItIsFatal() {
val records = sealedRecords("one", "two", "three")
val sub = subscriber()
sub.accept(records[0])
sub.accept(records[2]) // gap
assertEquals(PreviewStatus.UNVERIFIABLE, sub.status)
// The binding says a receiver backfills the missing records from a
// replay source; a stream that completes is verifiable again.
assertEquals(RecordOutcome.Accepted, sub.accept(records[1]))
assertEquals(RecordOutcome.Accepted, sub.accept(records[2]))
assertEquals("onetwothree", sub.previewText)
assertEquals(PreviewStatus.LIVE, sub.status)
}
@Test
fun aRecordThatDoesNotOpenIsRejectedWithoutTouchingTheStream() {
val records = sealedRecords("one", "two")
val sub = subscriber()
sub.accept(records[0])
val tampered = records[1].frame.copyOf().also { it[0] = (it[0].toInt() xor 0xff).toByte() }
assertEquals(RecordOutcome.Undecryptable, sub.accept(records[1].copyWithFrame(tampered)))
assertEquals("one", sub.previewText)
assertEquals(1L, sub.transcript.chunkCount)
assertEquals(
"a record that fails its AEAD never advances the high-water mark",
1L,
sub.highWaterMark,
)
}
@Test
fun aRecordForAnotherStreamIsRefused() {
val sub = subscriber()
val alien = AgentTextStreamRecordV1(ByteArray(32) { 0x66 }, seq = 1, recordType = 1, frame = ByteArray(32))
assertEquals(RecordOutcome.WrongStream, sub.accept(alien))
assertEquals(0L, sub.highWaterMark)
}
@Test
fun onlyTextDeltasAppendToThePreview() {
val publisher = runBlocking { AgentTextStreamPublisher.open(crypto, InMemoryAgentTextStreamSequenceStore()) }
val sub = subscriber()
runBlocking {
sub.accept(publisher.publish(AgentTextStreamRecordV1.TYPE_TEXT_DELTA, "answer".encodeToByteArray()))
// Progress and status are agent chrome. The spec is explicit that
// they MUST NOT reach preview text, notifications, indexes or
// automation input.
sub.accept(publisher.publish(AgentTextStreamRecordV1.TYPE_PROGRESS_DELTA, "reading files".encodeToByteArray()))
sub.accept(publisher.publish(AgentTextStreamRecordV1.TYPE_STATUS, "thinking".encodeToByteArray()))
}
assertEquals("answer", sub.previewText)
assertEquals("thinking", sub.latestStatus)
assertEquals("reading files", sub.latestProgress)
// Every accepted record still folds into the transcript — the hash
// covers the stream, not just the text.
assertEquals(3L, sub.transcript.chunkCount)
}
@Test
fun aCheckpointReplacesThePreviewAndLaterDeltasAppendToIt() {
val publisher = runBlocking { AgentTextStreamPublisher.open(crypto, InMemoryAgentTextStreamSequenceStore()) }
val sub = subscriber()
runBlocking {
sub.accept(publisher.publish(AgentTextStreamRecordV1.TYPE_TEXT_DELTA, "draft".encodeToByteArray()))
sub.accept(publisher.publish(AgentTextStreamRecordV1.TYPE_CHECKPOINT, "the whole answer".encodeToByteArray()))
sub.accept(publisher.publish(AgentTextStreamRecordV1.TYPE_TEXT_DELTA, " so far".encodeToByteArray()))
}
assertEquals("the whole answer so far", sub.previewText)
}
@Test
fun anAbortCancelsThePreviewWithoutProducingText() {
val publisher = runBlocking { AgentTextStreamPublisher.open(crypto, InMemoryAgentTextStreamSequenceStore()) }
val sub = subscriber()
runBlocking {
sub.accept(publisher.publish(AgentTextStreamRecordV1.TYPE_TEXT_DELTA, "half an ans".encodeToByteArray()))
sub.accept(publisher.publish(AgentTextStreamRecordV1.TYPE_ABORT, ByteArray(0)))
}
assertEquals(PreviewStatus.ABORTED, sub.status)
assertTrue(sub.previewText.isEmpty())
}
@Test
fun theFinalMessageIsWhatDecidesWhetherWeSawTheRealStream() {
val records = sealedRecords("a", "b", "c")
val sub = subscriber()
records.forEach { sub.accept(it) }
val publisherSide = runBlocking { AgentTextStreamPublisher.open(crypto, InMemoryAgentTextStreamSequenceStore()) }
runBlocking { listOf("a", "b", "c").forEach { publisherSide.publish(AgentTextStreamRecordV1.TYPE_TEXT_DELTA, it.encodeToByteArray()) } }
assertTrue(
sub.matchesFinal(publisherSide.transcript.hash, publisherSide.transcript.chunkCount),
)
assertFalse(
"a hash that disagrees means we saw a different stream, even though every record opened",
sub.matchesFinal(ByteArray(32), 3),
)
assertFalse(
"the same hash with a different count is still a different stream",
sub.matchesFinal(publisherSide.transcript.hash, 2),
)
}
@Test
fun anUnverifiablePreviewNeverClaimsToMatchAFinal() {
val records = sealedRecords("a", "b", "c")
val sub = subscriber()
sub.accept(records[0])
sub.accept(records[2])
assertEquals(PreviewStatus.UNVERIFIABLE, sub.status)
// Even a hash that happens to agree cannot rehabilitate it: the
// receiver knows it is missing a record it never folded.
assertFalse(sub.matchesFinal(sub.transcript.hash, sub.transcript.chunkCount))
}
}