mirror of
https://github.com/vitorpamplona/amethyst.git
synced 2026-10-05 19:28:25 +00:00
feat(marmot): agent text stream previews over the repo's own QUIC stack
The transport binding was the last piece missing from agent text streams, and it did not need a new QUIC implementation — `:quic` already had the whole hard part. What it needed was entering that stack at the right layer. `nestsClient` speaks WebTransport: HTTP/3, Extended CONNECT, QPACK, SETTINGS. Its `WebTransportSession` abstraction begins above all of that. Marmot's binding is raw QUIC — it negotiates its own ALPN (`marmot.quic_broker.v1` / `marmot.quic_stream.v1`) and writes frames straight onto QUIC streams, with no HTTP/3 anywhere in it. So this reuses everything below that line — connection, TLS 1.3, ALPN negotiation, stream multiplexing, loss recovery, the UDP socket — and none of the WebTransport wrapper. The codecs are in quartz next to the rest of agent-text-stream, because they are pure bytes and that is where the conformance risk lives: the control envelope with its literal 21-byte protocol string and its trailing-byte rejection, the uint32 frame codec with both the broker's blind cap and a policy-aware one, `quic://` candidate parsing down to ignoring everything after the authority and never sending an IP literal as SNI, and the first record's stream id pinning the rest. `:marmotQuic` is the connection layer, mirroring how `:nestsClient` sits on `:quic`. A publisher claims a room on a uni stream, a subscriber reads the fan-out on a bidi one, and an endpoint that does not take our ALPN is reported as unusable so the caller moves to the next candidate rather than waiting on records that never come. Verified against MDK's own `marmot-quic-broker`, which is the only way to know a wire format is right: our publisher and subscriber meet inside the reference broker, the records come back, open under the group-derived key and fold to the publisher's transcript hash, and the broker keeps rooms apart. Opt in with -DmarmotQuicBroker=host:port; the cases skip visibly without one, so an ordinary test run needs no broker. Still not wired at the app layer: nothing yet mints a kind-1200 start, picks a candidate, or renders a live preview, so `send` (0xF2D2) and `fanout` (0xF2D4) stay unadvertised. The direct path has no start-payload candidate format in v1 and is unimplemented. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_016kCuA6tc4JQzHPCDd39GHq
This commit is contained in:
+8
-1
@@ -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
|
||||
|
||||
@@ -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.
|
||||
@@ -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<Test>().configureEach {
|
||||
System.getProperty("marmotQuicBroker")?.let { systemProperty("marmotQuicBroker", it) }
|
||||
}
|
||||
+113
@@ -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<AgentTextStreamRecordV1>
|
||||
|
||||
/** 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,
|
||||
}
|
||||
}
|
||||
+201
@@ -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<AgentTextStreamRecordV1> =
|
||||
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()
|
||||
}
|
||||
}
|
||||
+219
@@ -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")
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -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.
|
||||
|
||||
|
||||
+318
@@ -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<Byte>()
|
||||
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<AgentTextStreamRecordV1> {
|
||||
buffer += chunk
|
||||
val out = mutableListOf<AgentTextStreamRecordV1>()
|
||||
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://<authority>"]` 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
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
+249
@@ -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<AgentTextStreamRecordV1>()
|
||||
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"))
|
||||
}
|
||||
}
|
||||
@@ -39,6 +39,7 @@ include(":geode")
|
||||
include(":commons")
|
||||
include(":quic")
|
||||
include(":nestsClient")
|
||||
include(":marmotQuic")
|
||||
include(":desktopApp")
|
||||
include(":cli")
|
||||
include(":relayBench")
|
||||
|
||||
Reference in New Issue
Block a user