diff --git a/amethyst/src/main/java/com/vitorpamplona/amethyst/model/marmot/AndroidMarmotMessageStore.kt b/amethyst/src/main/java/com/vitorpamplona/amethyst/model/marmot/AndroidMarmotMessageStore.kt index 829126eb67..d313488a3e 100644 --- a/amethyst/src/main/java/com/vitorpamplona/amethyst/model/marmot/AndroidMarmotMessageStore.kt +++ b/amethyst/src/main/java/com/vitorpamplona/amethyst/model/marmot/AndroidMarmotMessageStore.kt @@ -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 { - 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 = + 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 = readAllFrom(messagesFile(nostrGroupId)) + + private fun readAllFrom(file: File): List { 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, + ) = writeAllTo(messagesFile(nostrGroupId), messages) + + private fun writeAllTo( + file: File, + messages: List, ) { - val file = messagesFile(nostrGroupId) file.parentFile?.mkdirs() val encodedEntries = messages.map { it.encodeToByteArray() } diff --git a/cli/build.gradle.kts b/cli/build.gradle.kts index 6aa88d026f..cd33bc9080 100644 --- a/cli/build.gradle.kts +++ b/cli/build.gradle.kts @@ -43,6 +43,10 @@ tasks.named("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")) diff --git a/cli/src/main/kotlin/com/vitorpamplona/amethyst/cli/Main.kt b/cli/src/main/kotlin/com/vitorpamplona/amethyst/cli/Main.kt index 1375852f28..5d5997dfa6 100644 --- a/cli/src/main/kotlin/com/vitorpamplona/amethyst/cli/Main.kt +++ b/cli/src/main/kotlin/com/vitorpamplona/amethyst/cli/Main.kt @@ -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 ", + usage = "marmot ", 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) }, ), diff --git a/cli/src/main/kotlin/com/vitorpamplona/amethyst/cli/commands/StreamCommands.kt b/cli/src/main/kotlin/com/vitorpamplona/amethyst/cli/commands/StreamCommands.kt new file mode 100644 index 0000000000..e7e1ef3b70 --- /dev/null +++ b/cli/src/main/kotlin/com/vitorpamplona/amethyst/cli/commands/StreamCommands.kt @@ -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, + ): Int = + route( + "stream", + tail, + "stream …", + 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, + ): 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, + ): 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, + ): 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() + 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, + ): 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? { + var best: Pair? = 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) +} diff --git a/cli/src/main/kotlin/com/vitorpamplona/amethyst/cli/stores/FileStores.kt b/cli/src/main/kotlin/com/vitorpamplona/amethyst/cli/stores/FileStores.kt index 2fe55c7dd4..9b6d367c41 100644 --- a/cli/src/main/kotlin/com/vitorpamplona/amethyst/cli/stores/FileStores.kt +++ b/cli/src/main/kotlin/com/vitorpamplona/amethyst/cli/stores/FileStores.kt @@ -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 = + 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() } /** diff --git a/cli/tests/marmot/marmot-interop-headless.sh b/cli/tests/marmot/marmot-interop-headless.sh index a1f058ca53..eadd5f00c9 100755 --- a/cli/tests/marmot/marmot-interop-headless.sh +++ b/cli/tests/marmot/marmot-interop-headless.sh @@ -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 diff --git a/cli/tests/marmot/setup.sh b/cli/tests/marmot/setup.sh index 5d64c2d654..82fe0e3c1b 100644 --- a/cli/tests/marmot/setup.sh +++ b/cli/tests/marmot/setup.sh @@ -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. diff --git a/cli/tests/marmot/tests-extras.sh b/cli/tests/marmot/tests-extras.sh index 54c962e20e..9df6cb7a34 100644 --- a/cli/tests/marmot/tests-extras.sh +++ b/cli/tests/marmot/tests-extras.sh @@ -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 +} diff --git a/commons/src/commonMain/kotlin/com/vitorpamplona/amethyst/commons/marmot/MarmotIngest.kt b/commons/src/commonMain/kotlin/com/vitorpamplona/amethyst/commons/marmot/MarmotIngest.kt index 54f51999d4..217df5d7c2 100644 --- a/commons/src/commonMain/kotlin/com/vitorpamplona/amethyst/commons/marmot/MarmotIngest.kt +++ b/commons/src/commonMain/kotlin/com/vitorpamplona/amethyst/commons/marmot/MarmotIngest.kt @@ -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) } diff --git a/commons/src/commonMain/kotlin/com/vitorpamplona/amethyst/commons/marmot/MarmotManager.kt b/commons/src/commonMain/kotlin/com/vitorpamplona/amethyst/commons/marmot/MarmotManager.kt index b1187d3abb..cef8174ea6 100644 --- a/commons/src/commonMain/kotlin/com/vitorpamplona/amethyst/commons/marmot/MarmotManager.kt +++ b/commons/src/commonMain/kotlin/com/vitorpamplona/amethyst/commons/marmot/MarmotManager.kt @@ -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 = emptyList(), + parentEventId: HexKey? = null, + persistOwn: Boolean = true, + ): TextMessageBundle { + val template = + com.vitorpamplona.quartz.nip01Core.signers + .eventTemplate(kind = AgentTextStreamStart.KIND, description = "") { + AgentTextStreamStart + .tags(streamId, brokerCandidates, parentEventId = parentEventId) + .forEach { addUnique(it) } + } + val innerEvent = + com.vitorpamplona.quartz.nip59Giftwrap.rumors.RumorAssembler + .assembleRumor(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(kind = 9, description = text) { + AgentTextStreamFinal + .tags(streamId, transcriptHash, chunkCount) + .forEach { addUnique(it) } + } + val innerEvent = + com.vitorpamplona.quartz.nip59Giftwrap.rumors.RumorAssembler + .assembleRumor(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 = + 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. diff --git a/marmotQuic/README.md b/marmotQuic/README.md index 333ba2e29a..3f9c1e1837 100644 --- a/marmotQuic/README.md +++ b/marmotQuic/README.md @@ -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. diff --git a/marmotQuic/src/jvmAndroid/kotlin/com/vitorpamplona/marmotquic/QuicAgentTextStreamTransport.kt b/marmotQuic/src/jvmAndroid/kotlin/com/vitorpamplona/marmotquic/QuicAgentTextStreamTransport.kt index 82caaad301..5cd46aa81e 100644 --- a/marmotQuic/src/jvmAndroid/kotlin/com/vitorpamplona/marmotquic/QuicAgentTextStreamTransport.kt +++ b/marmotQuic/src/jvmAndroid/kotlin/com/vitorpamplona/marmotquic/QuicAgentTextStreamTransport.kt @@ -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 diff --git a/marmotQuic/src/jvmTest/kotlin/com/vitorpamplona/marmotquic/MarmotQuicBrokerInteropTest.kt b/marmotQuic/src/jvmTest/kotlin/com/vitorpamplona/marmotquic/MarmotQuicBrokerInteropTest.kt index 52daaf4d04..7178e50ca4 100644 --- a/marmotQuic/src/jvmTest/kotlin/com/vitorpamplona/marmotquic/MarmotQuicBrokerInteropTest.kt +++ b/marmotQuic/src/jvmTest/kotlin/com/vitorpamplona/marmotquic/MarmotQuicBrokerInteropTest.kt @@ -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 diff --git a/quartz/plans/2026-09-08-marmot-spec-resync.md b/quartz/plans/2026-09-08-marmot-spec-resync.md index fad088e64c..c6865b72cf 100644 --- a/quartz/plans/2026-09-08-marmot-spec-resync.md +++ b/quartz/plans/2026-09-08-marmot-spec-resync.md @@ -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. diff --git a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/marmot/appComponents/agentTextStream/AgentTextStreamStart.kt b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/marmot/appComponents/agentTextStream/AgentTextStreamStart.kt index 0b0c09ebe3..cb2e5623bf 100644 --- a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/marmot/appComponents/agentTextStream/AgentTextStreamStart.kt +++ b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/marmot/appComponents/agentTextStream/AgentTextStreamStart.kt @@ -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", ] (optional) * ["broker", ] (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, + 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, route: String = ROUTE_QUIC, + streamType: String = TYPE_TEXT, + finalKind: Int = FINAL_KIND_TEXT, + parentEventId: HexKey? = null, ): Array> = 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) } } } diff --git a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/marmot/appComponents/agentTextStream/AgentTextStreamSubscriber.kt b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/marmot/appComponents/agentTextStream/AgentTextStreamSubscriber.kt new file mode 100644 index 0000000000..1b819d917f --- /dev/null +++ b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/marmot/appComponents/agentTextStream/AgentTextStreamSubscriber.kt @@ -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) + } +} diff --git a/marmotQuic/src/commonMain/kotlin/com/vitorpamplona/marmotquic/MarmotQuicStreamTransport.kt b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/marmot/appComponents/agentTextStream/transport/MarmotQuicStreamTransport.kt similarity index 98% rename from marmotQuic/src/commonMain/kotlin/com/vitorpamplona/marmotquic/MarmotQuicStreamTransport.kt rename to quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/marmot/appComponents/agentTextStream/transport/MarmotQuicStreamTransport.kt index 2bd759570a..7a9d627eab 100644 --- a/marmotQuic/src/commonMain/kotlin/com/vitorpamplona/marmotquic/MarmotQuicStreamTransport.kt +++ b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/marmot/appComponents/agentTextStream/transport/MarmotQuicStreamTransport.kt @@ -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 diff --git a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/marmot/mls/group/MarmotMessageStore.kt b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/marmot/mls/group/MarmotMessageStore.kt index 05440a35b7..e2073dc765 100644 --- a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/marmot/mls/group/MarmotMessageStore.kt +++ b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/marmot/mls/group/MarmotMessageStore.kt @@ -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 = emptyMap() } diff --git a/quartz/src/jvmAndroidTest/kotlin/com/vitorpamplona/quartz/marmot/appComponents/AgentTextStreamSubscriberTest.kt b/quartz/src/jvmAndroidTest/kotlin/com/vitorpamplona/quartz/marmot/appComponents/AgentTextStreamSubscriberTest.kt new file mode 100644 index 0000000000..723094c391 --- /dev/null +++ b/quartz/src/jvmAndroidTest/kotlin/com/vitorpamplona/quartz/marmot/appComponents/AgentTextStreamSubscriberTest.kt @@ -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 = + 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)) + } +}