diff --git a/.claude/CLAUDE.md b/.claude/CLAUDE.md index 732f94a837..7a15cb83cc 100644 --- a/.claude/CLAUDE.md +++ b/.claude/CLAUDE.md @@ -17,7 +17,9 @@ relay-server code; smaller modules are `benchmark` (Android macrobenchmarks), `relayBench` (head-to-head relay benchmark — boots geode, strfry and other relay binaries, replays a shared deterministic corpus, measures ingest/query/ NIP-77 sync; `./relayBench/run.sh`, see `relayBench/README.md`) and -`quic-interop` (QUIC interop runner, lives at `quic/interop`). `nestsClient` runs +`quic-interop` (QUIC interop runner, lives at `quic/interop`). `marmotQuic` is the Marmot raw-QUIC transport +binding for agent text stream previews (`transports/quic.md`) on top of +`:quic` — its own ALPNs and framing, not WebTransport. `nestsClient` runs the audio-room protocol on top of `:quic` for the NIP-53 audio-rooms feature. It implements both IETF `draft-ietf-moq-transport-17` (under `moq/`) and **moq-lite Lite-03** (kixelated's variant, under `moq/lite/`); the production listener AND speaker paths both run on moq-lite to interop with the @@ -78,6 +80,11 @@ amethyst/ KMP project that needs MoQ. Has no Android-framework dependencies. - `nestsClient/` = MoQ + audio-rooms client; takes `:quic` as transport, Quartz for crypto, `MediaCodec` / `AudioRecord` / `AudioTrack` for audio. +- `marmotQuic/` = Marmot's raw-QUIC binding for agent text stream previews. + Takes `:quic` for the connection and `:quartz` for the record/envelope + codecs. Not WebTransport — the binding has its own ALPNs and writes frames + straight onto QUIC streams, so it deliberately does not reuse + `nestsClient`'s `WebTransportSession`. - `amethyst/` & `desktopApp/` = Platform-native layouts and navigation - `cli/` = Thin assembly layer over `quartz/` + `commons/` (no new logic allowed). May also depend on `:geode` (for `amy serve`, which embeds the diff --git a/marmotQuic/README.md b/marmotQuic/README.md new file mode 100644 index 0000000000..333ba2e29a --- /dev/null +++ b/marmotQuic/README.md @@ -0,0 +1,72 @@ +# marmotQuic + +Marmot's raw QUIC transport binding for agent text stream previews +(`transports/quic.md`), on top of the repo's own pure-Kotlin `:quic` stack. + +## Why this is not `nestsClient`'s WebTransport + +Both features move bytes over `:quic`, but they enter it at different layers. + +`nestsClient` speaks **WebTransport**: HTTP/3, an Extended CONNECT handshake, a +`:protocol` pseudo-header, QPACK, SETTINGS negotiation. Its +`WebTransportSession` abstraction starts *above* all of that. + +Marmot's binding is **raw QUIC**. It negotiates its own ALPN — +`marmot.quic_broker.v1` for the broker path, `marmot.quic_stream.v1` for the +direct one — and writes frames straight onto QUIC streams. There is no HTTP/3 +in it at all, so `WebTransportSession` is the wrong shape. + +What both share is everything below that line, which is the hard part and is +already built: the QUIC connection, TLS 1.3, ALPN negotiation, stream +multiplexing, loss recovery and the UDP socket. + +## Shape + +- A **publisher** opens a client-initiated *unidirectional* stream, writes a + `publish` control envelope, then record frames. +- A **subscriber** opens a client-initiated *bidirectional* stream, writes a + `subscribe` control envelope, and reads the fan-out on the return direction. + +A broker rejects the wrong pairing. Both roles frame everything the same way: +`uint32 frame_len || bytes`, the control envelope first and then each +`AgentTextStreamRecordV1`. + +The codecs — control envelope, frame reader/writer with both caps, `quic://` +candidate parsing — live in `quartz` next to the rest of agent-text-stream, +because they are pure bytes and belong with the feature. This module is only +the connection. + +## The broker sees nothing + +Records are encrypted under a key derived from the group's MLS exporter. A +broker holds no key and learns only the routing pair +`(stream_id, start_event_id)` plus ciphertext. It is an untrusted forwarder, +and a candidate that points somewhere hostile still cannot forge a record. + +## Interop tests + +`MarmotQuicBrokerInteropTest` drives our publisher and subscriber through +MDK's own reference broker. Start it from an MDK checkout: + +```bash +cargo build --release --bin marmot-quic-broker +./target/release/marmot-quic-broker --bind 127.0.0.1:4450 --json +``` + +then: + +```bash +./gradlew :marmotQuic:jvmTest -DmarmotQuicBroker=127.0.0.1:4450 +``` + +Without the property the cases skip visibly, so an ordinary `./gradlew test` +never needs a broker on the machine. + +## 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 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/build.gradle.kts b/marmotQuic/build.gradle.kts new file mode 100644 index 0000000000..8039d51a9d --- /dev/null +++ b/marmotQuic/build.gradle.kts @@ -0,0 +1,86 @@ +import org.jetbrains.kotlin.gradle.dsl.JvmTarget + +plugins { + alias(libs.plugins.kotlinMultiplatform) + alias(libs.plugins.androidKotlinMultiplatformLibrary) +} + +kotlin { + jvm { + compilerOptions { + jvmTarget.set(JvmTarget.JVM_21) + } + } + + android { + namespace = "com.vitorpamplona.marmotquic" + compileSdk = + libs.versions.android.compileSdk + .get() + .toInt() + minSdk = + libs.versions.android.minSdk + .get() + .toInt() + + compilerOptions { + jvmTarget.set(JvmTarget.JVM_21) + } + + withHostTest {} + } + + sourceSets { + commonMain { + dependencies { + implementation(libs.kotlin.stdlib) + implementation(libs.kotlinx.coroutines.core) + api(project(":quartz")) + implementation(project(":quic")) + } + } + + commonTest { + dependencies { + implementation(libs.kotlin.test) + implementation(libs.kotlinx.coroutines.test) + } + } + + val jvmAndroid = + create("jvmAndroid") { + dependsOn(commonMain.get()) + } + + jvmMain { + dependsOn(jvmAndroid) + } + + androidMain { + dependsOn(jvmAndroid) + } + + jvmTest { + dependencies { + implementation(libs.kotlin.test) + implementation(libs.kotlinx.coroutines.test) + implementation(libs.secp256k1.kmp.jni.jvm) + } + } + + getByName("androidHostTest") { + dependencies { + implementation(libs.kotlin.test) + implementation(libs.kotlinx.coroutines.test) + implementation(libs.secp256k1.kmp.jni.jvm) + } + } + } +} + +// Forward the broker opt-in from the Gradle JVM to the test workers. Without +// this, `-DmarmotQuicBroker=...` never reaches the test and every interop +// case silently skips. Mirrors the same forwarding in `:nestsClient`. +tasks.withType().configureEach { + System.getProperty("marmotQuicBroker")?.let { systemProperty("marmotQuicBroker", it) } +} diff --git a/marmotQuic/src/commonMain/kotlin/com/vitorpamplona/marmotquic/MarmotQuicStreamTransport.kt b/marmotQuic/src/commonMain/kotlin/com/vitorpamplona/marmotquic/MarmotQuicStreamTransport.kt new file mode 100644 index 0000000000..2bd759570a --- /dev/null +++ b/marmotQuic/src/commonMain/kotlin/com/vitorpamplona/marmotquic/MarmotQuicStreamTransport.kt @@ -0,0 +1,113 @@ +/* + * 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.marmotquic + +import com.vitorpamplona.quartz.marmot.appComponents.agentTextStream.AgentTextStreamRecordV1 +import kotlinx.coroutines.flow.Flow + +/** + * One delivery stream of an agent text stream preview, as + * `transports/quic.md` defines it: a QUIC stream carrying + * `uint32 frame_len || AgentTextStreamRecordV1` frames, preceded on the broker + * path by one control envelope framed the same way. + * + * The transport is deliberately narrow. It moves opaque ciphertext records and + * knows nothing about what they say: the record key comes from the group's MLS + * exporter, so a broker — and this layer — sees only the routing pair + * `(stream_id, start_event_id)` and bytes it cannot read. + */ +interface MarmotQuicStream { + /** Append one record to the stream. */ + suspend fun send(record: AgentTextStreamRecordV1) + + /** + * Records as they arrive, already de-framed and with the binding's + * stream-id pinning applied. Completes when the peer finishes the stream. + * + * Ordering, replay and gap handling belong to the caller — they need the + * transcript to decide, and this layer has no key to fold one with. + */ + fun incoming(): Flow + + /** Finish our write side cleanly; the stream ends when both sides have. */ + suspend fun finish() + + /** Tear the whole thing down, including the QUIC connection under it. */ + suspend fun close() +} + +/** + * Opens preview delivery streams against a `quic://` candidate. + * + * A candidate is advisory: one that fails to connect, fails TLS, or serves a + * different `(stream_id, start_event_id)` is unusable and the caller moves to + * the next. Implementations therefore surface a failure as + * [MarmotQuicException] rather than pretending a stream exists. + */ +interface MarmotQuicTransport { + /** + * Claim a broker room and stream records into it. + * + * The publisher path is a client-opened UNIDIRECTIONAL stream: it writes + * a `publish` control envelope and then the record frames. A broker + * rejects a publish envelope that arrives on a bidirectional stream. + */ + suspend fun publish( + candidate: String, + streamId: ByteArray, + startEventId: ByteArray, + ): MarmotQuicStream + + /** + * Join a broker room and read the fan-out. + * + * The subscriber path is a client-opened BIDIRECTIONAL stream: it writes a + * `subscribe` control envelope and reads record frames on the return + * direction. A broker rejects a subscribe envelope on a unidirectional + * stream, because it would have nowhere to answer. + */ + suspend fun subscribe( + candidate: String, + streamId: ByteArray, + startEventId: ByteArray, + ): MarmotQuicStream +} + +/** Why a candidate turned out to be unusable. */ +class MarmotQuicException( + val kind: Kind, + message: String, + cause: Throwable? = null, +) : RuntimeException(message, cause) { + enum class Kind { + /** The `quic://` candidate does not parse, or is over the 512-byte bound. */ + BadCandidate, + + /** UDP, QUIC or TLS never got as far as a connection. */ + HandshakeFailed, + + /** The endpoint does not speak our ALPN, so it is not a Marmot endpoint. */ + AlpnRejected, + + /** The peer closed the stream or the connection under us. */ + PeerClosed, + } +} diff --git a/marmotQuic/src/jvmAndroid/kotlin/com/vitorpamplona/marmotquic/QuicAgentTextStreamTransport.kt b/marmotQuic/src/jvmAndroid/kotlin/com/vitorpamplona/marmotquic/QuicAgentTextStreamTransport.kt new file mode 100644 index 0000000000..82caaad301 --- /dev/null +++ b/marmotQuic/src/jvmAndroid/kotlin/com/vitorpamplona/marmotquic/QuicAgentTextStreamTransport.kt @@ -0,0 +1,201 @@ +/* + * 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.marmotquic + +import com.vitorpamplona.quartz.marmot.appComponents.agentTextStream.AgentTextStreamRecordV1 +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.QuicBrokerControlEnvelopeV1 +import com.vitorpamplona.quartz.marmot.appComponents.agentTextStream.transport.QuicEndpointCandidate +import com.vitorpamplona.quic.connection.QuicConnection +import com.vitorpamplona.quic.connection.QuicConnectionConfig +import com.vitorpamplona.quic.connection.QuicConnectionDriver +import com.vitorpamplona.quic.stream.QuicStream +import com.vitorpamplona.quic.tls.CertificateValidator +import com.vitorpamplona.quic.transport.UdpSocket +import kotlinx.coroutines.CoroutineScope +import kotlinx.coroutines.Dispatchers +import kotlinx.coroutines.SupervisorJob +import kotlinx.coroutines.flow.Flow +import kotlinx.coroutines.flow.flow +import kotlinx.coroutines.withTimeoutOrNull + +/** + * `transports/quic.md` on top of the repo's own pure-Kotlin `:quic` stack. + * + * Marmot's binding is RAW QUIC, not WebTransport: it negotiates its own ALPN + * (`marmot.quic_broker.v1` / `marmot.quic_stream.v1`) and writes frames + * straight onto QUIC streams. So this deliberately does not reuse + * `nestsClient`'s `WebTransportSession` — that abstraction begins above HTTP/3 + * Extended CONNECT, which this binding has no part of. What it does reuse is + * everything under that: the QUIC connection, TLS 1.3, ALPN negotiation, + * stream multiplexing and the UDP socket. + * + * One stream per delivery: a publisher's uni stream or a subscriber's bidi + * stream owns its connection and closes it on [MarmotQuicStream.close]. That + * is the shape the binding describes — a room is a stream — and it keeps a + * failed candidate from leaving a connection behind. + */ +class QuicAgentTextStreamTransport( + private val parentScope: CoroutineScope = CoroutineScope(SupervisorJob() + Dispatchers.IO), + /** + * Preview endpoints and brokers are commonly self-signed, and the binding + * says so: a client MAY pin by DER or SHA-256 fingerprint through local + * configuration instead of the system trust store. That choice is the + * caller's, so the validator is required rather than defaulted — the type + * system should not let "forgot to decide" compile. + */ + private val certificateValidator: CertificateValidator, + private val handshakeTimeoutMillis: Long = 10_000L, + /** + * The group's `max_plaintext_frame_len`, when the caller knows it. A + * receiver that knows the policy must reject a frame above it; without one + * the broker's blind cap applies. + */ + private val maxPlaintextFrameLen: Long? = null, +) : MarmotQuicTransport { + override suspend fun publish( + candidate: String, + streamId: ByteArray, + startEventId: ByteArray, + ): MarmotQuicStream = open(candidate, streamId, startEventId, BrokerControlType.PUBLISH) + + override suspend fun subscribe( + candidate: String, + streamId: ByteArray, + startEventId: ByteArray, + ): MarmotQuicStream = open(candidate, streamId, startEventId, BrokerControlType.SUBSCRIBE) + + private suspend fun open( + candidate: String, + streamId: ByteArray, + startEventId: ByteArray, + role: BrokerControlType, + ): MarmotQuicStream { + val endpoint = + QuicEndpointCandidate.parse(candidate) + ?: throw MarmotQuicException(MarmotQuicException.Kind.BadCandidate, "unusable quic:// candidate") + + val socket = + try { + UdpSocket.connect(endpoint.host, endpoint.port) + } catch (t: Throwable) { + throw MarmotQuicException(MarmotQuicException.Kind.HandshakeFailed, "cannot reach the candidate", t) + } + + val connection = + QuicConnection( + // An IP literal is matched against an iPAddress SAN and never + // sent as SNI; `serverName` is only meaningful for a DNS name. + serverName = endpoint.serverNameIndication ?: endpoint.host, + config = QuicConnectionConfig(), + tlsCertificateValidator = certificateValidator, + alpnList = listOf(MarmotQuicAlpn.BROKER), + ) + val driver = QuicConnectionDriver(connection, socket, parentScope) + driver.start() + + try { + val completed = + withTimeoutOrNull(handshakeTimeoutMillis) { + connection.awaitHandshake() + true + } + if (completed == null || connection.status != QuicConnection.Status.CONNECTED) { + throw MarmotQuicException( + MarmotQuicException.Kind.HandshakeFailed, + "QUIC handshake did not complete (status=${connection.status})", + ) + } + // An endpoint that did not take our ALPN is not a Marmot endpoint, + // whatever else it may be. Fail here so the caller moves to the + // next candidate rather than waiting on records that never come. + val alpn = connection.tls.negotiatedAlpn + if (alpn == null || !alpn.contentEquals(MarmotQuicAlpn.BROKER)) { + throw MarmotQuicException( + MarmotQuicException.Kind.AlpnRejected, + "endpoint negotiated ${alpn?.decodeToString()} instead of ${MarmotQuicAlpn.BROKER.decodeToString()}", + ) + } + + // Stream direction IS the role: a publisher claims the room on a + // uni stream, a subscriber needs the return direction of a bidi + // one. A broker rejects the wrong pairing. + val stream = + when (role) { + BrokerControlType.PUBLISH -> connection.openUniStream() + BrokerControlType.SUBSCRIBE -> connection.openBidiStream() + } + + // The control envelope is the first frame, framed exactly like a + // record frame — length-prefixed the same way, so a broker reads + // both with one framer. + stream.send.enqueue(frameEnvelope(QuicBrokerControlEnvelopeV1(role, streamId, startEventId))) + driver.wakeup() + + return QuicStreamDelivery(stream, driver, maxPlaintextFrameLen) + } catch (t: Throwable) { + driver.close() + throw if (t is MarmotQuicException) t else MarmotQuicException(MarmotQuicException.Kind.PeerClosed, "${t.message}", t) + } + } + + private fun frameEnvelope(envelope: QuicBrokerControlEnvelopeV1): ByteArray { + val encoded = envelope.encode() + return byteArrayOf( + ((encoded.size shr 24) and 0xff).toByte(), + ((encoded.size shr 16) and 0xff).toByte(), + ((encoded.size shr 8) and 0xff).toByte(), + (encoded.size and 0xff).toByte(), + ) + encoded + } +} + +/** One QUIC stream carrying framed records, plus the connection it rides on. */ +private class QuicStreamDelivery( + private val stream: QuicStream, + private val driver: QuicConnectionDriver, + maxPlaintextFrameLen: Long?, +) : MarmotQuicStream { + private val reader = AgentTextStreamFraming.Reader(maxPlaintextFrameLen) + + override suspend fun send(record: AgentTextStreamRecordV1) { + stream.send.enqueue(AgentTextStreamFraming.frame(record)) + driver.wakeup() + } + + override fun incoming(): Flow = + flow { + stream.incoming.collect { chunk -> + for (record in reader.push(chunk)) emit(record) + } + } + + override suspend fun finish() { + stream.send.finish() + driver.wakeup() + } + + override suspend fun close() { + driver.close() + } +} diff --git a/marmotQuic/src/jvmTest/kotlin/com/vitorpamplona/marmotquic/MarmotQuicBrokerInteropTest.kt b/marmotQuic/src/jvmTest/kotlin/com/vitorpamplona/marmotquic/MarmotQuicBrokerInteropTest.kt new file mode 100644 index 0000000000..52daaf4d04 --- /dev/null +++ b/marmotQuic/src/jvmTest/kotlin/com/vitorpamplona/marmotquic/MarmotQuicBrokerInteropTest.kt @@ -0,0 +1,219 @@ +/* + * 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.marmotquic + +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.AgentTextStreamTranscriptV1 +import com.vitorpamplona.quartz.marmot.appComponents.agentTextStream.InMemoryAgentTextStreamSequenceStore +import com.vitorpamplona.quic.tls.PermissiveCertificateValidator +import kotlinx.coroutines.CoroutineScope +import kotlinx.coroutines.Dispatchers +import kotlinx.coroutines.SupervisorJob +import kotlinx.coroutines.async +import kotlinx.coroutines.cancel +import kotlinx.coroutines.delay +import kotlinx.coroutines.flow.take +import kotlinx.coroutines.flow.toList +import kotlinx.coroutines.runBlocking +import kotlinx.coroutines.withTimeout +import org.junit.Assume +import kotlin.random.Random +import kotlin.test.AfterTest +import kotlin.test.Test +import kotlin.test.assertContentEquals +import kotlin.test.assertEquals +import kotlin.test.assertTrue + +/** + * Drives our `transports/quic.md` client against MDK's own + * `marmot-quic-broker`, the reference implementation of the other side. + * + * This is the only way to know the binding is right. Everything it exercises + * is a place where two implementations have to agree byte for byte and where + * our own tests would happily agree with themselves: the ALPN string, the + * control envelope's layout and its literal protocol name, which stream + * direction each role uses, and the 4-byte frame prefix. The records + * themselves stay opaque to the broker — it fans out ciphertext and never + * holds a key. + * + * Opt in with `-DmarmotQuicBroker=127.0.0.1:4450` after starting: + * + * ``` + * cargo build --release --bin marmot-quic-broker # in the MDK checkout + * ./target/release/marmot-quic-broker --bind 127.0.0.1:4450 --json + * ``` + * + * Without the property the test skips, so a normal `./gradlew test` never + * needs a broker on the machine. + */ +class MarmotQuicBrokerInteropTest { + private val scope = CoroutineScope(SupervisorJob() + Dispatchers.IO) + + private val brokerAuthority: String? = System.getProperty("marmotQuicBroker") + + /** + * Report "no broker configured" as a JUnit skip rather than a silent pass, + * so a run that was meant to exercise the broker cannot look green because + * the property never reached the worker. + */ + private fun requireBroker(): String { + Assume.assumeTrue( + "set -DmarmotQuicBroker=host:port and start MDK's marmot-quic-broker to run the interop cases", + brokerAuthority != null, + ) + return brokerAuthority!! + } + + private val candidate get() = "quic://$brokerAuthority" + + @AfterTest + fun tearDown() { + scope.cancel() + } + + private fun transport() = + QuicAgentTextStreamTransport( + parentScope = scope, + // The broker generates a self-signed certificate on startup. The + // binding expects exactly that ("preview endpoints and brokers may + // be self-signed") and says a client MAY pin it locally; a test + // against a throwaway broker accepts it outright. + certificateValidator = PermissiveCertificateValidator(), + ) + + private fun keyContext( + streamId: ByteArray, + startEventId: ByteArray, + ) = AgentTextStreamKeyContextV1( + groupId = ByteArray(32) { 0x01 }, + streamId = streamId, + mlsEpoch = 3, + senderId = ByteArray(32) { 0x02 }, + startEventId = startEventId, + ) + + @Test + fun ourPublisherAndSubscriberMeetInsideTheReferenceBroker() { + requireBroker() + runBlocking { + val streamId = Random.nextBytes(32) + val startEventId = Random.nextBytes(32) + val secret = ByteArray(32) { 0x77 } + val crypto = AgentTextStreamCrypto(secret, keyContext(streamId, startEventId)) + + val transport = transport() + // Subscribe first: the broker's replay window is 0 by default, so + // a subscriber that arrives after the records were pushed sees + // nothing — which is the binding working as specified, not a bug. + val subscriber = transport.subscribe(candidate, streamId, startEventId) + val received = async { withTimeout(30_000) { subscriber.incoming().take(3).toList() } } + delay(500) + + val publisher = transport.publish(candidate, streamId, startEventId) + val sender = AgentTextStreamPublisher.open(crypto, InMemoryAgentTextStreamSequenceStore()) + val plaintexts = listOf("the ", "quick ", "brown fox") + for (text in plaintexts) { + publisher.send(sender.publish(AgentTextStreamRecordV1.TYPE_TEXT_DELTA, text.encodeToByteArray())) + } + publisher.finish() + + val records = received.await() + + assertEquals(listOf(1L, 2L, 3L), records.map { it.seq }) + assertEquals( + plaintexts.joinToString(""), + records.joinToString("") { crypto.open(it).frame.decodeToString() }, + "the broker relays ciphertext and cannot read a byte of it, so what comes back must open " + + "under the same group-derived key", + ) + + // A receiver that folded every record must agree with the + // publisher's transcript, which is what the final kind:9 carries. + val fold = AgentTextStreamTranscriptV1.start(streamId, startEventId) + records.forEach { fold.append(crypto.open(it)) } + assertContentEquals(sender.transcript.hash, fold.hash) + assertEquals(sender.transcript.chunkCount, fold.chunkCount) + + publisher.close() + subscriber.close() + } + } + + @Test + fun theBrokerKeepsRoomsApart() { + requireBroker() + runBlocking { + val startEventId = Random.nextBytes(32) + val mine = Random.nextBytes(32) + val theirs = Random.nextBytes(32) + val transport = transport() + + val subscriber = transport.subscribe(candidate, mine, startEventId) + val received = async { withTimeout(15_000) { subscriber.incoming().take(1).toList() } } + delay(500) + + // Same start event, different stream id: a different room. "A + // broker MUST NOT merge or cross-deliver records between different + // rooms." + val wrongRoom = transport.publish(candidate, theirs, startEventId) + val strayCrypto = AgentTextStreamCrypto(ByteArray(32) { 0x66 }, keyContext(theirs, startEventId)) + val stray = AgentTextStreamPublisher.open(strayCrypto, InMemoryAgentTextStreamSequenceStore()) + wrongRoom.send(stray.publish(AgentTextStreamRecordV1.TYPE_TEXT_DELTA, "not for you".encodeToByteArray())) + wrongRoom.finish() + + // Now the right room, so the test finishes on a positive signal + // rather than a timeout we cannot distinguish from a hang. + val rightRoom = transport.publish(candidate, mine, startEventId) + val mineCrypto = AgentTextStreamCrypto(ByteArray(32) { 0x77 }, keyContext(mine, startEventId)) + val ours = AgentTextStreamPublisher.open(mineCrypto, InMemoryAgentTextStreamSequenceStore()) + rightRoom.send(ours.publish(AgentTextStreamRecordV1.TYPE_TEXT_DELTA, "for you".encodeToByteArray())) + rightRoom.finish() + + val records = received.await() + assertEquals(1, records.size) + assertContentEquals(mine, records.single().streamId) + assertEquals("for you", mineCrypto.open(records.single()).frame.decodeToString()) + + wrongRoom.close() + rightRoom.close() + subscriber.close() + } + } + + @Test + fun theBrokerRefusesAnEndpointThatDoesNotSpeakOurAlpn() { + assertTrue(requireBroker().isNotEmpty()) + // Sanity: a candidate that parses but points nowhere must fail as a + // handshake, not hang or throw something unclassified. + runBlocking { + val failure = + runCatching { + withTimeout(30_000) { + transport().subscribe("quic://127.0.0.1:1", Random.nextBytes(32), Random.nextBytes(32)) + } + }.exceptionOrNull() + assertTrue(failure is MarmotQuicException, "expected a MarmotQuicException, got $failure") + } + } +} diff --git a/quartz/plans/2026-09-08-marmot-spec-resync.md b/quartz/plans/2026-09-08-marmot-spec-resync.md index 1dcbdf8f4f..fad088e64c 100644 --- a/quartz/plans/2026-09-08-marmot-spec-resync.md +++ b/quartz/plans/2026-09-08-marmot-spec-resync.md @@ -695,11 +695,31 @@ 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`), because - there is no data plane behind them yet and a role we cannot serve is worse for - the group than a role we do not claim. A group whose policy requires `send` is - refused at join rather than joined into a state every peer would reject us - from. -- **The QUIC transport for agent text streams is not wired.** The record layer - and the kind-1200 anchor are implemented; nothing yet opens a WebTransport - session to a broker and feeds it records. + 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. + +- **The QUIC transport binding is implemented and verified against MDK's + broker.** `transports/quic.md` is a RAW QUIC binding — its own ALPNs + (`marmot.quic_broker.v1` / `marmot.quic_stream.v1`), frames written straight + onto QUIC streams — so it does not go through `nestsClient`'s + `WebTransportSession`, which begins above HTTP/3 Extended CONNECT. It sits + directly on `:quic`, which already had everything under that line: the + connection, TLS 1.3, ALPN negotiation, stream multiplexing, the UDP socket. + + The pure protocol layer (control envelope, `uint32` frame codec with both + caps, `quic://` candidate parsing, stream-id pinning) is in `quartz`; the + connection layer is the new `:marmotQuic` module, which mirrors how + `:nestsClient` sits on `:quic`. `MarmotQuicBrokerInteropTest` drives our + publisher and subscriber through MDK's own `marmot-quic-broker` and checks + that the records come back, open under the group-derived key, and fold to the + 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. + diff --git a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/marmot/appComponents/agentTextStream/transport/MarmotQuicBinding.kt b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/marmot/appComponents/agentTextStream/transport/MarmotQuicBinding.kt new file mode 100644 index 0000000000..6fd5103b69 --- /dev/null +++ b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/marmot/appComponents/agentTextStream/transport/MarmotQuicBinding.kt @@ -0,0 +1,318 @@ +/* + * 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.transport + +import com.vitorpamplona.quartz.marmot.appComponents.agentTextStream.AgentTextStreamQuicPolicyV1 +import com.vitorpamplona.quartz.marmot.appComponents.agentTextStream.AgentTextStreamRecordV1 +import com.vitorpamplona.quartz.marmot.appComponents.agentTextStream.QuicVarInt +import com.vitorpamplona.quartz.marmot.appComponents.agentTextStream.addLengthPrefixed + +/** + * ALPN identifiers for `transports/quic.md`. + * + * QUIC requires ALPN on every connection, and the binding gives each delivery + * mode its own so the two can diverge without a version negotiation of their + * own: a peer that only speaks one simply fails the handshake. + * + * This is a raw QUIC binding, not WebTransport — there is no HTTP/3 layer, no + * Extended CONNECT and no `:protocol` pseudo-header. The records ride directly + * on QUIC streams. + */ +object MarmotQuicAlpn { + /** Broker-relayed delivery — the v1 discovery mechanism. */ + val BROKER = "marmot.quic_broker.v1".encodeToByteArray() + + /** Direct point-to-point delivery, where the sender already knows the receiver's endpoint. */ + val DIRECT = "marmot.quic_stream.v1".encodeToByteArray() +} + +/** Which side of a broker room a control envelope claims. */ +enum class BrokerControlType( + val code: Int, +) { + /** Claims the room and streams records into it. Client-opened unidirectional stream. */ + PUBLISH(1), + + /** Joins the room and reads the fan-out. Client-opened bidirectional stream. */ + SUBSCRIBE(2), + ; + + companion object { + fun fromCode(code: Int): BrokerControlType? = entries.firstOrNull { it.code == code } + } +} + +/** + * The first frame on a broker stream, framed exactly like a record frame. + * + * ``` + * struct { + * opaque marmot_broker<1..255>; // ASCII "marmot.quic_broker.v1" + * BrokerControlType control_type; // uint8 + * opaque stream_id<1..64>; + * opaque start_event_id<1..64>; + * } QuicBrokerControlEnvelopeV1; + * ``` + * + * `marmot_broker` is length-prefixed rather than a fixed array on purpose: a + * future protocol string of a different length still decodes for an old + * reader, which can then reject it cleanly instead of misparsing the fields + * behind it. + * + * The broker is an untrusted forwarder. Everything it learns is in here — the + * routing pair — plus ciphertext; it cannot read, author or alter preview + * plaintext. + */ +class QuicBrokerControlEnvelopeV1( + val controlType: BrokerControlType, + val streamId: ByteArray, + val startEventId: ByteArray, +) { + init { + require(streamId.size in 1..MAX_ID_LEN) { "broker control stream_id must be 1..$MAX_ID_LEN bytes" } + require(startEventId.size in 1..MAX_ID_LEN) { "broker control start_event_id must be 1..$MAX_ID_LEN bytes" } + } + + fun encode(): ByteArray { + val out = ArrayList() + out.addLengthPrefixed(PROTOCOL) + out.add(controlType.code.toByte()) + out.addLengthPrefixed(streamId) + out.addLengthPrefixed(startEventId) + return out.toByteArray() + } + + companion object { + /** The exact 21 ASCII bytes a broker matches on. */ + val PROTOCOL = "marmot.quic_broker.v1".encodeToByteArray() + + const val MAX_ID_LEN = 64 + + fun decode(bytes: ByteArray): QuicBrokerControlEnvelopeV1 { + var at = 0 + + val protocolLen = QuicVarInt.decode(bytes, at) + at += protocolLen.length + require(at + protocolLen.value <= bytes.size) { "broker control envelope is truncated while reading marmot_broker" } + val protocol = bytes.copyOfRange(at, at + protocolLen.value.toInt()) + at += protocolLen.value.toInt() + require(protocol.contentEquals(PROTOCOL)) { + "broker control envelope names another protocol: ${protocol.decodeToString()}" + } + + require(at < bytes.size) { "broker control envelope is truncated while reading control_type" } + val controlType = + BrokerControlType.fromCode(bytes[at].toInt() and 0xff) + ?: throw IllegalArgumentException("unknown broker control_type: ${bytes[at].toInt() and 0xff}") + at += 1 + + val streamIdLen = QuicVarInt.decode(bytes, at) + at += streamIdLen.length + require(at + streamIdLen.value <= bytes.size) { "broker control envelope is truncated while reading stream_id" } + val streamId = bytes.copyOfRange(at, at + streamIdLen.value.toInt()) + at += streamIdLen.value.toInt() + + val startEventIdLen = QuicVarInt.decode(bytes, at) + at += startEventIdLen.length + require(at + startEventIdLen.value <= bytes.size) { + "broker control envelope is truncated while reading start_event_id" + } + val startEventId = bytes.copyOfRange(at, at + startEventIdLen.value.toInt()) + at += startEventIdLen.value.toInt() + + // "A broker MUST reject an envelope whose frame carries trailing + // bytes after the envelope" — an envelope with a record glued on + // is a different message than the one we would be acting on. + require(at == bytes.size) { "broker control envelope has ${bytes.size - at} trailing byte(s)" } + + return QuicBrokerControlEnvelopeV1(controlType, streamId, startEventId) + } + } +} + +/** + * Length-delimited record frames on a delivery stream: + * `uint32 frame_len || AgentTextStreamRecordV1[frame_len]`. + * + * The 4-byte big-endian prefix is the transport's only framing; QUIC gives an + * ordered byte stream, not messages. + */ +object AgentTextStreamFraming { + /** + * Header + AEAD-tag allowance the binding adds on top of a group's + * `max_plaintext_frame_len` when sizing a frame. + */ + const val FRAME_OVERHEAD_ALLOWANCE = 1024 + + /** + * The cap a broker enforces. It cannot read group state, so it uses the + * v1 component's largest legal `max_plaintext_frame_len` plus the same + * allowance. A client that knows the group's actual policy uses that + * instead, which is always smaller. + */ + const val BROKER_MAX_FRAME_LEN = AgentTextStreamQuicPolicyV1.MAX_PLAINTEXT_FRAME_LEN.toInt() + FRAME_OVERHEAD_ALLOWANCE + + fun frame(record: AgentTextStreamRecordV1): ByteArray { + val encoded = record.encode() + require(encoded.size <= BROKER_MAX_FRAME_LEN) { "agent text stream frame is over the transport cap" } + val out = ByteArray(4 + encoded.size) + out[0] = ((encoded.size shr 24) and 0xff).toByte() + out[1] = ((encoded.size shr 16) and 0xff).toByte() + out[2] = ((encoded.size shr 8) and 0xff).toByte() + out[3] = (encoded.size and 0xff).toByte() + encoded.copyInto(out, 4) + return out + } + + /** + * Incremental frame reader over a QUIC stream's bytes. + * + * A QUIC stream hands over whatever arrived, so a record can straddle any + * number of reads; the reader buffers until a whole frame is present and + * returns the records it completed. + * + * @param maxPlaintextFrameLen the group's policy value when the caller + * knows it. A receiver that knows the policy must reject anything above + * it; without one the broker's blind cap applies. + */ + class Reader( + maxPlaintextFrameLen: Long? = null, + ) { + private val maxFrameLen = + maxPlaintextFrameLen?.let { (it + FRAME_OVERHEAD_ALLOWANCE).coerceAtMost(BROKER_MAX_FRAME_LEN.toLong()).toInt() } + ?: BROKER_MAX_FRAME_LEN + + private var buffer = ByteArray(0) + + /** The stream id of the first record, which every later record must repeat. */ + var pinnedStreamId: ByteArray? = null + private set + + /** Buffer the bytes and return whatever records they completed, in order. */ + fun push(chunk: ByteArray): List { + buffer += chunk + val out = mutableListOf() + while (true) { + if (buffer.size < 4) return out + val frameLen = + ((buffer[0].toLong() and 0xff) shl 24) or + ((buffer[1].toLong() and 0xff) shl 16) or + ((buffer[2].toLong() and 0xff) shl 8) or + (buffer[3].toLong() and 0xff) + require(frameLen <= maxFrameLen) { "agent text stream frame_len $frameLen is over the cap $maxFrameLen" } + if (buffer.size < 4 + frameLen) return out + + val record = AgentTextStreamRecordV1.decode(buffer.copyOfRange(4, (4 + frameLen).toInt())) + buffer = buffer.copyOfRange((4 + frameLen).toInt(), buffer.size) + + val pinned = pinnedStreamId + if (pinned == null) { + pinnedStreamId = record.streamId + } else { + // "A reader MUST reject records whose stream_id differs + // from the first record's stream_id on the same stream." + require(record.streamId.contentEquals(pinned)) { + "agent text stream record carries a different stream_id than the stream it arrived on" + } + } + out.add(record) + } + } + } +} + +/** + * One `["broker", "quic://"]` endpoint candidate. + * + * Candidates are advisory routing hints, not authenticated stream content: a + * candidate that points somewhere hostile still cannot forge a record, because + * every record is authenticated under the group-derived record key. So a + * candidate that does not parse, does not connect, or serves a different + * stream is skipped rather than treated as an error — hence [parse] returning + * null instead of throwing. + */ +class QuicEndpointCandidate( + val host: String, + val port: Int, + /** + * True for a DNS hostname, false for an IPv4/IPv6 literal. Decides the + * trust model: a name gets normal DNS-name/SNI validation, a literal is + * matched against an `iPAddress` subjectAltName and gets no SNI at all. + */ + val isDnsName: Boolean, +) { + /** The SNI to send, or null for an IP literal (which must not carry one). */ + val serverNameIndication: String? get() = host.takeIf { isDnsName } + + companion object { + const val SCHEME = "quic://" + + /** `quic://` + at most 505 authority bytes. */ + const val MAX_CANDIDATE_BYTES = 512 + + fun parse(candidate: String): QuicEndpointCandidate? { + val bytes = candidate.encodeToByteArray() + if (bytes.size > MAX_CANDIDATE_BYTES) return null + // Round-tripping catches a candidate that was not valid UTF-8 to + // begin with: the decoder substitutes U+FFFD and the bytes differ. + if (!bytes.decodeToString().encodeToByteArray().contentEquals(bytes)) return null + if (!candidate.startsWith(SCHEME)) return null + + // "The authority ends at the first /, ? or #; everything after + // that character is ignored." + val rest = candidate.substring(SCHEME.length) + val authority = rest.takeWhile { it != '/' && it != '?' && it != '#' } + if (authority.isEmpty()) return null + + val host: String + val portText: String + if (authority.startsWith("[")) { + val close = authority.indexOf(']') + if (close < 0) return null + host = authority.substring(1, close) + if (authority.getOrNull(close + 1) != ':') return null + portText = authority.substring(close + 2) + if (!host.contains(':')) return null + } else { + val colon = authority.lastIndexOf(':') + if (colon <= 0) return null + host = authority.substring(0, colon) + portText = authority.substring(colon + 1) + if (host.contains(':')) return null + } + + if (host.isEmpty()) return null + val port = portText.toIntOrNull() ?: return null + if (port !in 1..65535) return null + + return QuicEndpointCandidate(host, port, isDnsName = !looksLikeIpLiteral(host)) + } + + private fun looksLikeIpLiteral(host: String): Boolean { + if (host.contains(':')) return true + val parts = host.split('.') + if (parts.size != 4) return false + return parts.all { part -> + part.isNotEmpty() && part.length <= 3 && part.all { it.isDigit() } && (part.toIntOrNull() ?: 256) <= 255 + } + } + } +} diff --git a/quartz/src/jvmAndroidTest/kotlin/com/vitorpamplona/quartz/marmot/appComponents/AgentTextStreamQuicTransportTest.kt b/quartz/src/jvmAndroidTest/kotlin/com/vitorpamplona/quartz/marmot/appComponents/AgentTextStreamQuicTransportTest.kt new file mode 100644 index 0000000000..23993693f9 --- /dev/null +++ b/quartz/src/jvmAndroidTest/kotlin/com/vitorpamplona/quartz/marmot/appComponents/AgentTextStreamQuicTransportTest.kt @@ -0,0 +1,249 @@ +/* + * 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.AgentTextStreamRecordV1 +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.QuicBrokerControlEnvelopeV1 +import com.vitorpamplona.quartz.marmot.appComponents.agentTextStream.transport.QuicEndpointCandidate +import org.junit.Assert.assertArrayEquals +import org.junit.Assert.assertEquals +import org.junit.Assert.assertNull +import org.junit.Assert.assertThrows +import org.junit.Assert.assertTrue +import org.junit.Test + +/** + * `transports/quic.md` — the raw QUIC binding for agent text stream previews. + * + * This is the wire between us and a broker written by somebody else, so every + * rule here is one an independent implementation will hold us to: the exact + * ALPN strings, the control envelope's field layout and its literal protocol + * string, the 4-byte frame prefix and its caps, and what a candidate URL does + * and does not mean. + */ +class AgentTextStreamQuicTransportTest { + private val streamId = ByteArray(32) { 0x11 } + private val startEventId = ByteArray(32) { 0x22 } + + // --- ALPN ------------------------------------------------------------- + + @Test + fun theTwoAlpnsAreTheExactStringsTheSpecNames() { + assertEquals("marmot.quic_broker.v1", MarmotQuicAlpn.BROKER.decodeToString()) + assertEquals("marmot.quic_stream.v1", MarmotQuicAlpn.DIRECT.decodeToString()) + } + + // --- Broker control envelope ----------------------------------------- + + @Test + fun aControlEnvelopeRoundTrips() { + for (type in BrokerControlType.entries) { + val envelope = QuicBrokerControlEnvelopeV1(type, streamId, startEventId) + val decoded = QuicBrokerControlEnvelopeV1.decode(envelope.encode()) + assertEquals(type, decoded.controlType) + assertArrayEquals(streamId, decoded.streamId) + assertArrayEquals(startEventId, decoded.startEventId) + } + } + + @Test + fun theEnvelopeStartsWithTheLengthPrefixedProtocolString() { + val encoded = QuicBrokerControlEnvelopeV1(BrokerControlType.PUBLISH, streamId, startEventId).encode() + // varint(21) fits in one byte, so the protocol string starts at index 1. + assertEquals(21, encoded[0].toInt()) + assertEquals("marmot.quic_broker.v1", encoded.copyOfRange(1, 22).decodeToString()) + assertEquals(BrokerControlType.PUBLISH.code, encoded[22].toInt()) + } + + @Test + fun anEnvelopeNamingAnotherProtocolIsRejected() { + val good = QuicBrokerControlEnvelopeV1(BrokerControlType.SUBSCRIBE, streamId, startEventId).encode() + // Same length, different string: a broker must reject on the bytes, not the length. + val tampered = good.copyOf() + tampered[1] = 'M'.code.toByte() + assertThrows(IllegalArgumentException::class.java) { QuicBrokerControlEnvelopeV1.decode(tampered) } + } + + @Test + fun anUnknownControlTypeIsRejectedRatherThanIgnored() { + val encoded = QuicBrokerControlEnvelopeV1(BrokerControlType.PUBLISH, streamId, startEventId).encode() + encoded[22] = 0x7f + assertThrows(IllegalArgumentException::class.java) { QuicBrokerControlEnvelopeV1.decode(encoded) } + } + + @Test + fun trailingBytesAfterTheEnvelopeAreRejected() { + val encoded = QuicBrokerControlEnvelopeV1(BrokerControlType.PUBLISH, streamId, startEventId).encode() + assertThrows(IllegalArgumentException::class.java) { + QuicBrokerControlEnvelopeV1.decode(encoded + byteArrayOf(0x00)) + } + } + + @Test + fun anIdentityFieldOutsideItsBoundIsRejected() { + assertThrows(IllegalArgumentException::class.java) { + QuicBrokerControlEnvelopeV1(BrokerControlType.PUBLISH, ByteArray(0), startEventId).encode() + } + assertThrows(IllegalArgumentException::class.java) { + QuicBrokerControlEnvelopeV1(BrokerControlType.PUBLISH, ByteArray(65), startEventId).encode() + } + } + + // --- Record framing --------------------------------------------------- + + @Test + fun aFrameIsAFourByteBigEndianLengthFollowedByTheRecord() { + val record = AgentTextStreamRecordV1(streamId, seq = 1, recordType = 1, frame = ByteArray(7) { 0x5a }) + val encoded = record.encode() + val framed = AgentTextStreamFraming.frame(record) + + assertEquals(4 + encoded.size, framed.size) + assertEquals(0, framed[0].toInt()) + assertEquals(0, framed[1].toInt()) + assertEquals((encoded.size shr 8) and 0xff, framed[2].toInt() and 0xff) + assertEquals(encoded.size and 0xff, framed[3].toInt() and 0xff) + assertArrayEquals(encoded, framed.copyOfRange(4, framed.size)) + } + + @Test + fun theReaderDeliversRecordsAsTheirBytesArriveInAnySplit() { + val records = + (1..4).map { + AgentTextStreamRecordV1(streamId, seq = it.toLong(), recordType = 1, frame = "chunk $it".encodeToByteArray()) + } + val wire = records.fold(ByteArray(0)) { acc, r -> acc + AgentTextStreamFraming.frame(r) } + + // One byte at a time is the worst case a QUIC stream can hand us. + val reader = AgentTextStreamFraming.Reader() + val delivered = mutableListOf() + for (b in wire) delivered.addAll(reader.push(byteArrayOf(b))) + + assertEquals(records.size, delivered.size) + delivered.forEachIndexed { i, r -> + assertEquals(records[i].seq, r.seq) + assertArrayEquals(records[i].frame, r.frame) + } + } + + @Test + fun theReaderRefusesAFrameOverTheBrokerCap() { + val reader = AgentTextStreamFraming.Reader() + val tooBig = AgentTextStreamFraming.BROKER_MAX_FRAME_LEN + 1 + val header = + byteArrayOf( + ((tooBig shr 24) and 0xff).toByte(), + ((tooBig shr 16) and 0xff).toByte(), + ((tooBig shr 8) and 0xff).toByte(), + (tooBig and 0xff).toByte(), + ) + assertThrows(IllegalArgumentException::class.java) { reader.push(header) } + } + + @Test + fun aReaderThatKnowsTheGroupPolicyRefusesAnythingOverIt() { + // max_plaintext_frame_len + the spec's 1024-byte header/tag allowance. + val reader = AgentTextStreamFraming.Reader(maxPlaintextFrameLen = 16) + val record = AgentTextStreamRecordV1(streamId, seq = 1, recordType = 1, frame = ByteArray(2000)) + assertThrows(IllegalArgumentException::class.java) { reader.push(AgentTextStreamFraming.frame(record)) } + } + + @Test + fun aReaderPinsTheStreamIdOfItsFirstRecord() { + val reader = AgentTextStreamFraming.Reader() + reader.push(AgentTextStreamFraming.frame(AgentTextStreamRecordV1(streamId, 1, 1, frame = ByteArray(1)))) + val otherStream = ByteArray(32) { 0x33 } + assertThrows(IllegalArgumentException::class.java) { + reader.push(AgentTextStreamFraming.frame(AgentTextStreamRecordV1(otherStream, 2, 1, frame = ByteArray(1)))) + } + } + + // --- Endpoint candidates --------------------------------------------- + + @Test + fun aCandidateIsAnAuthorityAndNothingAfterIt() { + val parsed = QuicEndpointCandidate.parse("quic://broker.example:4433")!! + assertEquals("broker.example", parsed.host) + assertEquals(4433, parsed.port) + assertTrue(parsed.isDnsName) + + // A path, query or fragment is ignored, not rejected. + for (suffix in listOf("/room/1", "?x=1", "#frag")) { + val withSuffix = QuicEndpointCandidate.parse("quic://broker.example:4433$suffix")!! + assertEquals("broker.example", withSuffix.host) + assertEquals(4433, withSuffix.port) + } + } + + @Test + fun anIpv6LiteralKeepsItsBracketsOutOfTheHost() { + val parsed = QuicEndpointCandidate.parse("quic://[2001:db8::1]:443")!! + assertEquals("2001:db8::1", parsed.host) + assertEquals(443, parsed.port) + assertTrue( + "an IP literal is matched against an iPAddress SAN and is never sent as SNI", + !parsed.isDnsName, + ) + assertNull(parsed.serverNameIndication) + } + + @Test + fun anIpv4LiteralIsAlsoNotASniName() { + val parsed = QuicEndpointCandidate.parse("quic://192.0.2.7:4433")!! + assertEquals("192.0.2.7", parsed.host) + assertTrue(!parsed.isDnsName) + assertNull(parsed.serverNameIndication) + } + + @Test + fun aDnsCandidateCarriesItsOwnSni() { + assertEquals("broker.example", QuicEndpointCandidate.parse("quic://broker.example:4433")!!.serverNameIndication) + } + + @Test + fun anUnusableCandidateIsSkippedNotFatal() { + // The spec says a receiver moves to the next candidate; null is that. + val bad = + listOf( + "https://broker.example:4433", + "quic://broker.example", + "quic://broker.example:0", + "quic://broker.example:65536", + "quic://:4433", + "quic://[2001:db8::1:443", + "quic://" + "h".repeat(506) + ":443", + "", + ) + for (candidate in bad) { + assertNull("must not accept $candidate", QuicEndpointCandidate.parse(candidate)) + } + } + + @Test + fun aCandidateOverTheByteBoundIsRejected() { + val authority = "h".repeat(505 - 4) + ":443" + assertEquals(512, ("quic://$authority").length) + assertTrue(QuicEndpointCandidate.parse("quic://$authority") != null) + assertNull(QuicEndpointCandidate.parse("quic://x$authority")) + } +} diff --git a/settings.gradle.kts b/settings.gradle.kts index 129aa4fe3f..4e63b53707 100644 --- a/settings.gradle.kts +++ b/settings.gradle.kts @@ -39,6 +39,7 @@ include(":geode") include(":commons") include(":quic") include(":nestsClient") +include(":marmotQuic") include(":desktopApp") include(":cli") include(":relayBench")