mirror of
https://github.com/vitorpamplona/amethyst.git
synced 2026-10-05 11:18:24 +00:00
feat(marmot): the direct QUIC path, and a pin instead of blind trust
Two halves of the same gap in `transports/quic.md`. **The direct path.** The binding has a second delivery mode we had not built: the sender dials the receiver, opens one unidirectional stream and writes records with no control envelope at all. It is deliberately smaller than the broker path — the dialed endpoint is already the one receiver, so there is no room to claim — and it negotiates its own ALPN so an incompatible change to either mode cannot reach the other. Note the inverted direction: here the RECEIVER listens and the SENDER dials, which is also why v1 gives it no start-payload discovery and it is only usable against an endpoint known out of band. Only the sending half is here. `:quic` is a client stack with no server role, so this module can dial a direct receiver but cannot be one; that is recorded in the README rather than half-built. **The pin.** Preview endpoints and brokers are commonly self-signed and the binding expects that, saying a client MAY pin by exact DER or SHA-256 fingerprint. What we had instead was `PermissiveCertificateValidator` on the CLI path, which is not a weaker trust model — it is none, and anyone on the path can be the broker. `PinnedCertificateValidator` replaces the chain and the hostname check and nothing else: the peer still has to sign the TLS transcript with the pinned certificate's private key, so copying a public certificate off the wire buys an attacker nothing. `amy marmot stream send|watch` takes `--pin-sha256`, and `--insecure` still exists for a throwaway local broker but now has to be asked for by name. Both are verified against the reference implementation, which is the only thing that can tell an ALPN string, a stream direction, an absent envelope and a frame prefix from an implementation agreeing with itself: our direct sender against `wn stream receive`, and the pin — accepted and refused — against a real handshake with `marmot-quic-broker`. One thing that only showed up under a real handshake: a certificate the validator refuses closes the connection before it is established, and the transport was reporting that as PeerClosed. A caller walking a candidate list reads that kind to decide what to do next, and "never connected" is not "the peer hung up on us", so it is classified on the connection's actual status now. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_016kCuA6tc4JQzHPCDd39GHq
This commit is contained in:
@@ -36,7 +36,10 @@ import com.vitorpamplona.quartz.marmot.mls.crypto.MlsCryptoProvider
|
||||
import com.vitorpamplona.quartz.nip01Core.core.Event
|
||||
import com.vitorpamplona.quartz.nip01Core.core.hexToByteArray
|
||||
import com.vitorpamplona.quartz.nip01Core.core.toHexKey
|
||||
import com.vitorpamplona.quic.tls.CertificateValidator
|
||||
import com.vitorpamplona.quic.tls.JdkCertificateValidator
|
||||
import com.vitorpamplona.quic.tls.PermissiveCertificateValidator
|
||||
import com.vitorpamplona.quic.tls.PinnedCertificateValidator
|
||||
import kotlinx.coroutines.withTimeoutOrNull
|
||||
|
||||
/**
|
||||
@@ -56,8 +59,10 @@ object StreamCommands {
|
||||
| marmot stream start GID [--stream-id HEX] [--broker quic://HOST:PORT[,…]]
|
||||
| publish the kind:1200 that anchors a stream; prints stream_id + start_event_id
|
||||
|
|
||||
| marmot stream send GID --stream-id HEX --start-event-id HEX --broker URI TEXT…
|
||||
| push TEXT as TextDelta records to the broker; prints the transcript to finish with
|
||||
| marmot stream send GID --stream-id HEX --start-event-id HEX
|
||||
| (--broker URI | --direct quic://HOST:PORT) TEXT…
|
||||
| push TEXT as TextDelta records; --direct dials the receiver point to point
|
||||
| (ALPN marmot.quic_stream.v1, no control envelope) instead of a broker
|
||||
|
|
||||
| marmot stream watch GID [--stream-id HEX] [--timeout SECS]
|
||||
| find the kind:1200 in the group, subscribe over QUIC, fold the preview
|
||||
@@ -67,6 +72,12 @@ object StreamCommands {
|
||||
|
|
||||
|Every record is encrypted under the group's own MLS exporter secret, so a
|
||||
|broker relays ciphertext and learns only which room it belongs to.
|
||||
|
|
||||
|TLS trust for the QUIC hop (send and watch):
|
||||
| --pin-sha256 HEX[,HEX…] trust exactly these leaf certificates (self-signed
|
||||
| endpoints; colons and whitespace are ignored)
|
||||
| --insecure accept any certificate — local testing only
|
||||
|Without either, the platform trust store decides.
|
||||
""".trimMargin()
|
||||
|
||||
suspend fun dispatch(
|
||||
@@ -141,8 +152,16 @@ object StreamCommands {
|
||||
val streamId = args.flag("stream-id")
|
||||
val startEventId = args.flag("start-event-id")
|
||||
val broker = args.flag("broker")
|
||||
if (positional.size < 2 || streamId == null || startEventId == null || broker == null) {
|
||||
return Output.error("bad_args", "stream send GID --stream-id HEX --start-event-id HEX --broker URI TEXT…")
|
||||
// The two delivery modes are alternatives, not a fallback chain: one
|
||||
// dials a broker room, the other dials the receiver itself, and they
|
||||
// negotiate different ALPNs. Picking silently when both are given
|
||||
// would hide which one actually carried the records.
|
||||
val direct = args.flag("direct")
|
||||
if (positional.size < 2 || streamId == null || startEventId == null || (broker == null) == (direct == null)) {
|
||||
return Output.error(
|
||||
"bad_args",
|
||||
"stream send GID --stream-id HEX --start-event-id HEX (--broker URI | --direct quic://HOST:PORT) TEXT…",
|
||||
)
|
||||
}
|
||||
|
||||
Context.open(dataDir).use { ctx ->
|
||||
@@ -165,13 +184,20 @@ object StreamCommands {
|
||||
epoch = anchorEpoch,
|
||||
)
|
||||
val publisher = AgentTextStreamPublisher.open(crypto, InMemoryAgentTextStreamSequenceStore())
|
||||
val transport = QuicAgentTextStreamTransport(certificateValidator = PermissiveCertificateValidator())
|
||||
val transport = QuicAgentTextStreamTransport(certificateValidator = certificateValidator(args))
|
||||
|
||||
val stream =
|
||||
try {
|
||||
transport.publish(broker, streamId.hexToByteArray(), startEventId.hexToByteArray())
|
||||
if (direct != null) {
|
||||
transport.sendDirect(direct, streamId.hexToByteArray(), startEventId.hexToByteArray())
|
||||
} else {
|
||||
transport.publish(broker!!, streamId.hexToByteArray(), startEventId.hexToByteArray())
|
||||
}
|
||||
} catch (e: Exception) {
|
||||
return Output.error("broker_unreachable", "${e.message}")
|
||||
return Output.error(
|
||||
if (direct != null) "receiver_unreachable" else "broker_unreachable",
|
||||
"${e.message}",
|
||||
)
|
||||
}
|
||||
try {
|
||||
for (text in positional.drop(1)) {
|
||||
@@ -187,6 +213,8 @@ object StreamCommands {
|
||||
"group_id" to gid,
|
||||
"stream_id" to streamId,
|
||||
"start_event_id" to startEventId,
|
||||
"mode" to if (direct != null) "direct" else "broker",
|
||||
"endpoint" to (direct ?: broker),
|
||||
"records" to positional.size - 1,
|
||||
"epoch" to crypto.context.mlsEpoch,
|
||||
// What `stream finish` has to publish so a receiver can
|
||||
@@ -199,6 +227,29 @@ object StreamCommands {
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* The TLS trust policy for the QUIC hop, from the flags.
|
||||
*
|
||||
* Pinning is the interesting one and the binding calls it out: preview
|
||||
* endpoints and brokers are commonly self-signed, so a client MAY pin the
|
||||
* endpoint certificate by SHA-256 fingerprint instead of chaining to a CA.
|
||||
* `--insecure` stays available because a local test broker mints a fresh
|
||||
* certificate on every boot, but it is not a weaker trust model — it is
|
||||
* none, so it has to be asked for by name.
|
||||
*/
|
||||
private fun certificateValidator(args: Args): CertificateValidator {
|
||||
val pins =
|
||||
args
|
||||
.flag("pin-sha256")
|
||||
?.split(',')
|
||||
?.map { it.trim() }
|
||||
?.filter { it.isNotEmpty() }
|
||||
.orEmpty()
|
||||
if (pins.isNotEmpty()) return PinnedCertificateValidator.ofSha256Hex(*pins.toTypedArray())
|
||||
if (args.bool("insecure")) return PermissiveCertificateValidator()
|
||||
return JdkCertificateValidator()
|
||||
}
|
||||
|
||||
private suspend fun watch(
|
||||
dataDir: DataDir,
|
||||
rest: Array<String>,
|
||||
@@ -250,7 +301,7 @@ object StreamCommands {
|
||||
epoch = anchorEpoch,
|
||||
)
|
||||
val subscriber = AgentTextStreamSubscriber(crypto)
|
||||
val transport = QuicAgentTextStreamTransport(certificateValidator = PermissiveCertificateValidator())
|
||||
val transport = QuicAgentTextStreamTransport(certificateValidator = certificateValidator(args))
|
||||
|
||||
// "A receiver tries advertised candidates in listed order"; the
|
||||
// first that yields the matching stream wins.
|
||||
|
||||
+6
@@ -69,6 +69,12 @@ class MarmotAgentStreamWatcherTest {
|
||||
startEventId: ByteArray,
|
||||
): MarmotQuicStream = error("the watcher never publishes")
|
||||
|
||||
override suspend fun sendDirect(
|
||||
candidate: String,
|
||||
streamId: ByteArray,
|
||||
startEventId: ByteArray,
|
||||
): MarmotQuicStream = error("the watcher never sends")
|
||||
|
||||
override suspend fun subscribe(
|
||||
candidate: String,
|
||||
streamId: ByteArray,
|
||||
|
||||
+46
-12
@@ -26,10 +26,19 @@ multiplexing, loss recovery and the UDP socket.
|
||||
`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 **direct sender** opens a client-initiated *unidirectional* stream and
|
||||
writes record frames with **no** control envelope — the dialed endpoint is
|
||||
already the one receiver, so there is no room to name.
|
||||
|
||||
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`.
|
||||
A broker rejects the wrong pairing. Every role frames the same way:
|
||||
`uint32 frame_len || bytes` — on the broker path the control envelope first and
|
||||
then each `AgentTextStreamRecordV1`, on the direct path records from the first
|
||||
byte.
|
||||
|
||||
Note the direct path's connection direction: the **receiver** listens and the
|
||||
**sender** dials, inverted from the broker path where both ends dial the
|
||||
broker. Only the sender half is here; `:quic` is a client stack with no server
|
||||
role, so this module cannot expose a direct-path endpoint of its own.
|
||||
|
||||
The codecs — control envelope, frame reader/writer with both caps, `quic://`
|
||||
candidate parsing — live in `quartz` next to the rest of agent-text-stream,
|
||||
@@ -43,24 +52,47 @@ 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.
|
||||
|
||||
## TLS trust
|
||||
|
||||
Preview endpoints and brokers are commonly self-signed, and the binding says so:
|
||||
a client MAY pin the endpoint certificate by exact DER or SHA-256 fingerprint
|
||||
instead of chaining to a CA. `PinnedCertificateValidator` (in `:quic`) is that
|
||||
pin. It replaces the chain and the hostname check and nothing else — the peer
|
||||
still has to sign the TLS transcript with the pinned certificate's private key,
|
||||
so copying a public certificate off the wire buys an attacker nothing.
|
||||
|
||||
`amy marmot stream send|watch` takes `--pin-sha256 HEX[,HEX…]`; the reference
|
||||
broker prints its own `server_cert_sha256_fingerprint` in its startup JSON.
|
||||
|
||||
## Interop tests
|
||||
|
||||
`MarmotQuicBrokerInteropTest` drives our publisher and subscriber through
|
||||
MDK's own reference broker. Start it from an MDK checkout:
|
||||
`MarmotQuicBrokerInteropTest` drives our publisher and subscriber through MDK's
|
||||
own reference broker, and `MarmotQuicDirectInteropTest` drives our direct
|
||||
sender against MDK's direct receiver (`wn stream receive`). Both are the only
|
||||
way to know the binding is right: an ALPN string, a stream direction, a missing
|
||||
control envelope and a frame prefix are all things an implementation will
|
||||
happily agree with itself about.
|
||||
|
||||
Start the broker from an MDK checkout:
|
||||
|
||||
```bash
|
||||
cargo build --release --bin marmot-quic-broker
|
||||
cargo build --release --bin marmot-quic-broker --bin wn
|
||||
./target/release/marmot-quic-broker --bind 127.0.0.1:4450 --json
|
||||
```
|
||||
|
||||
then:
|
||||
|
||||
```bash
|
||||
./gradlew :marmotQuic:jvmTest -DmarmotQuicBroker=127.0.0.1:4450
|
||||
./gradlew :marmotQuic:jvmTest \
|
||||
-DmarmotQuicBroker=127.0.0.1:4450 \
|
||||
-DmarmotQuicBrokerPin=<server_cert_sha256_fingerprint from that JSON> \
|
||||
-DmarmotWn=/path/to/mdk/target/release/wn
|
||||
```
|
||||
|
||||
Without the property the cases skip visibly, so an ordinary `./gradlew test`
|
||||
never needs a broker on the machine.
|
||||
Each property gates its own cases and they skip visibly without it, so an
|
||||
ordinary `./gradlew test` never needs the reference implementation on the
|
||||
machine. `-DmarmotWn` needs no running process: the test spawns
|
||||
`wn stream receive` itself on a free port.
|
||||
|
||||
## Using it
|
||||
|
||||
@@ -70,6 +102,7 @@ run it in both directions against MDK.
|
||||
```bash
|
||||
amy marmot stream start GID --broker quic://127.0.0.1:4450
|
||||
amy marmot stream send GID --stream-id … --start-event-id … --broker … "hello"
|
||||
amy marmot stream send GID --stream-id … --start-event-id … --direct quic://host:port "hello"
|
||||
amy marmot stream watch GID --stream-id …
|
||||
amy marmot stream finish GID --stream-id … --transcript-hash … --chunk-count N "hello"
|
||||
```
|
||||
@@ -80,6 +113,7 @@ amy marmot stream finish GID --stream-id … --transcript-hash … --chunk-count
|
||||
an agent's job, and no agent runs in the app yet. Only `amy` publishes one.
|
||||
- The desktop app has no Marmot chat screen at all, so there is nothing to
|
||||
render a preview into. The watcher it would use already lives in `commons`.
|
||||
- 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.
|
||||
- The direct path's **receiving** half. `:quic` has no server role, so this
|
||||
module can dial a direct receiver but cannot be one. v1 also defines no
|
||||
start-payload candidate format for the direct path, so a sender only reaches
|
||||
a receiver whose endpoint it already knows out of band.
|
||||
|
||||
@@ -78,9 +78,11 @@ kotlin {
|
||||
}
|
||||
}
|
||||
|
||||
// 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`.
|
||||
// Forward the interop opt-ins from the Gradle JVM to the test workers.
|
||||
// Without this, `-DmarmotQuicBroker=...` / `-DmarmotWn=...` never reach the
|
||||
// tests and every interop case silently skips. Mirrors `:nestsClient`.
|
||||
tasks.withType<Test>().configureEach {
|
||||
System.getProperty("marmotQuicBroker")?.let { systemProperty("marmotQuicBroker", it) }
|
||||
System.getProperty("marmotQuicBrokerPin")?.let { systemProperty("marmotQuicBrokerPin", it) }
|
||||
System.getProperty("marmotWn")?.let { systemProperty("marmotWn", it) }
|
||||
}
|
||||
|
||||
+68
-14
@@ -54,8 +54,9 @@ import kotlinx.coroutines.withTimeoutOrNull
|
||||
* 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
|
||||
* One stream per delivery: a publisher's uni stream, a direct sender'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.
|
||||
*/
|
||||
@@ -89,16 +90,53 @@ class QuicAgentTextStreamTransport(
|
||||
startEventId: ByteArray,
|
||||
): MarmotQuicStream = open(candidate, streamId, startEventId, BrokerControlType.SUBSCRIBE)
|
||||
|
||||
/**
|
||||
* The direct path: dial the receiver, open one uni stream, write records.
|
||||
*
|
||||
* Two things separate it from [publish] beyond the ALPN. There is no
|
||||
* control envelope — the dialed endpoint is already the one receiver, so
|
||||
* there is no room to name, and the first bytes on the stream are a record
|
||||
* frame. And [startEventId] never leaves this process: it is validated for
|
||||
* shape so a caller cannot pass a placeholder that would later disagree
|
||||
* with the record key and transcript hash it is bound into, but nothing is
|
||||
* written for it. A direct endpoint learns it only if the out-of-band
|
||||
* setup supplied it separately.
|
||||
*
|
||||
* Only the SENDER half lives here. The receiver half has to listen, and
|
||||
* `:quic` is a client stack with no server role — so a direct-path
|
||||
* receiver is not something this module can offer yet.
|
||||
*/
|
||||
override suspend fun sendDirect(
|
||||
candidate: String,
|
||||
streamId: ByteArray,
|
||||
startEventId: ByteArray,
|
||||
): MarmotQuicStream = open(candidate, streamId, startEventId, role = null)
|
||||
|
||||
private suspend fun open(
|
||||
candidate: String,
|
||||
streamId: ByteArray,
|
||||
startEventId: ByteArray,
|
||||
role: BrokerControlType,
|
||||
/** The broker role to claim, or null for the envelope-less direct path. */
|
||||
role: BrokerControlType?,
|
||||
): MarmotQuicStream {
|
||||
val endpoint =
|
||||
QuicEndpointCandidate.parse(candidate)
|
||||
?: throw MarmotQuicException(MarmotQuicException.Kind.BadCandidate, "unusable quic:// candidate")
|
||||
|
||||
val alpn = if (role == null) MarmotQuicAlpn.DIRECT else MarmotQuicAlpn.BROKER
|
||||
if (role == null) {
|
||||
// Same bounds the broker envelope enforces, applied even though
|
||||
// nothing is encoded: a stream id or start event id this layer
|
||||
// would refuse to route is one the record key and transcript hash
|
||||
// should not be built on either.
|
||||
require(streamId.size in 1..QuicBrokerControlEnvelopeV1.MAX_ID_LEN) {
|
||||
"direct stream_id must be 1..${QuicBrokerControlEnvelopeV1.MAX_ID_LEN} bytes"
|
||||
}
|
||||
require(startEventId.size in 1..QuicBrokerControlEnvelopeV1.MAX_ID_LEN) {
|
||||
"direct start_event_id must be 1..${QuicBrokerControlEnvelopeV1.MAX_ID_LEN} bytes"
|
||||
}
|
||||
}
|
||||
|
||||
val socket =
|
||||
try {
|
||||
UdpSocket.connect(endpoint.host, endpoint.port)
|
||||
@@ -113,7 +151,7 @@ class QuicAgentTextStreamTransport(
|
||||
serverName = endpoint.serverNameIndication ?: endpoint.host,
|
||||
config = QuicConnectionConfig(),
|
||||
tlsCertificateValidator = certificateValidator,
|
||||
alpnList = listOf(MarmotQuicAlpn.BROKER),
|
||||
alpnList = listOf(alpn),
|
||||
)
|
||||
val driver = QuicConnectionDriver(connection, socket, parentScope)
|
||||
driver.start()
|
||||
@@ -133,11 +171,11 @@ class QuicAgentTextStreamTransport(
|
||||
// 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)) {
|
||||
val negotiated = connection.tls.negotiatedAlpn
|
||||
if (negotiated == null || !negotiated.contentEquals(alpn)) {
|
||||
throw MarmotQuicException(
|
||||
MarmotQuicException.Kind.AlpnRejected,
|
||||
"endpoint negotiated ${alpn?.decodeToString()} instead of ${MarmotQuicAlpn.BROKER.decodeToString()}",
|
||||
"endpoint negotiated ${negotiated?.decodeToString()} instead of ${alpn.decodeToString()}",
|
||||
)
|
||||
}
|
||||
|
||||
@@ -146,20 +184,36 @@ class QuicAgentTextStreamTransport(
|
||||
// one. A broker rejects the wrong pairing.
|
||||
val stream =
|
||||
when (role) {
|
||||
BrokerControlType.PUBLISH -> connection.openUniStream()
|
||||
BrokerControlType.PUBLISH, null -> 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()
|
||||
if (role != null) {
|
||||
// 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. The direct path writes none: its
|
||||
// stream opens straight into records.
|
||||
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)
|
||||
if (t is MarmotQuicException) throw t
|
||||
// A connection that never reached CONNECTED did not fail as a
|
||||
// peer closing on us mid-stream — it failed to be established at
|
||||
// all, and that is a different decision for a caller walking its
|
||||
// candidate list. Certificate rejection lands here: our own
|
||||
// validator refuses, we send a TLS alert, and the connection
|
||||
// closes before the handshake ever completes.
|
||||
val kind =
|
||||
if (connection.status == QuicConnection.Status.CONNECTED) {
|
||||
MarmotQuicException.Kind.PeerClosed
|
||||
} else {
|
||||
MarmotQuicException.Kind.HandshakeFailed
|
||||
}
|
||||
throw MarmotQuicException(kind, "${t.message}", t)
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
+60
@@ -28,6 +28,7 @@ import com.vitorpamplona.quartz.marmot.appComponents.agentTextStream.AgentTextSt
|
||||
import com.vitorpamplona.quartz.marmot.appComponents.agentTextStream.InMemoryAgentTextStreamSequenceStore
|
||||
import com.vitorpamplona.quartz.marmot.appComponents.agentTextStream.transport.MarmotQuicException
|
||||
import com.vitorpamplona.quic.tls.PermissiveCertificateValidator
|
||||
import com.vitorpamplona.quic.tls.PinnedCertificateValidator
|
||||
import kotlinx.coroutines.CoroutineScope
|
||||
import kotlinx.coroutines.Dispatchers
|
||||
import kotlinx.coroutines.SupervisorJob
|
||||
@@ -44,6 +45,7 @@ import kotlin.test.AfterTest
|
||||
import kotlin.test.Test
|
||||
import kotlin.test.assertContentEquals
|
||||
import kotlin.test.assertEquals
|
||||
import kotlin.test.assertFailsWith
|
||||
import kotlin.test.assertTrue
|
||||
|
||||
/**
|
||||
@@ -73,6 +75,13 @@ class MarmotQuicBrokerInteropTest {
|
||||
|
||||
private val brokerAuthority: String? = System.getProperty("marmotQuicBroker")
|
||||
|
||||
/**
|
||||
* The broker's leaf-certificate SHA-256, which it prints as
|
||||
* `server_cert_sha256_fingerprint` in its startup JSON. Supplying it opts
|
||||
* into the pinning cases.
|
||||
*/
|
||||
private val brokerPin: String? = System.getProperty("marmotQuicBrokerPin")
|
||||
|
||||
/**
|
||||
* 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
|
||||
@@ -93,6 +102,57 @@ class MarmotQuicBrokerInteropTest {
|
||||
scope.cancel()
|
||||
}
|
||||
|
||||
/**
|
||||
* The binding says a client MAY pin a self-signed endpoint by SHA-256
|
||||
* fingerprint. A unit test can only prove the comparison; whether pinning
|
||||
* actually admits the right peer is a question about a real TLS 1.3
|
||||
* handshake, and only a real one answers it.
|
||||
*/
|
||||
@Test
|
||||
fun aPinnedFingerprintCompletesTheHandshake() {
|
||||
requireBroker()
|
||||
val pin = requirePin()
|
||||
runBlocking {
|
||||
val transport =
|
||||
QuicAgentTextStreamTransport(
|
||||
parentScope = scope,
|
||||
certificateValidator = PinnedCertificateValidator.ofSha256Hex(pin),
|
||||
)
|
||||
val stream = transport.publish(candidate, Random.nextBytes(32), Random.nextBytes(32))
|
||||
stream.finish()
|
||||
stream.close()
|
||||
}
|
||||
}
|
||||
|
||||
@Test
|
||||
fun aPinForAnotherCertificateIsRefused() {
|
||||
requireBroker()
|
||||
requirePin()
|
||||
runBlocking {
|
||||
// Same broker, wrong pin. If this connected, the pin would be
|
||||
// decoration — which is exactly the failure mode that makes a
|
||||
// misconfigured pin dangerous rather than merely broken.
|
||||
val transport =
|
||||
QuicAgentTextStreamTransport(
|
||||
parentScope = scope,
|
||||
certificateValidator = PinnedCertificateValidator.ofSha256Hex("00".repeat(32)),
|
||||
)
|
||||
val failure =
|
||||
assertFailsWith<MarmotQuicException> {
|
||||
transport.publish(candidate, Random.nextBytes(32), Random.nextBytes(32))
|
||||
}
|
||||
assertEquals(MarmotQuicException.Kind.HandshakeFailed, failure.kind, "${failure.message}")
|
||||
}
|
||||
}
|
||||
|
||||
private fun requirePin(): String {
|
||||
Assume.assumeTrue(
|
||||
"set -DmarmotQuicBrokerPin=<server_cert_sha256_fingerprint from the broker's startup JSON>",
|
||||
brokerPin != null,
|
||||
)
|
||||
return brokerPin!!
|
||||
}
|
||||
|
||||
private fun transport() =
|
||||
QuicAgentTextStreamTransport(
|
||||
parentScope = scope,
|
||||
|
||||
+214
@@ -0,0 +1,214 @@
|
||||
/*
|
||||
* 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.InMemoryAgentTextStreamSequenceStore
|
||||
import com.vitorpamplona.quartz.marmot.appComponents.agentTextStream.transport.MarmotQuicException
|
||||
import com.vitorpamplona.quic.tls.PermissiveCertificateValidator
|
||||
import com.vitorpamplona.quic.tls.PinnedCertificateValidator
|
||||
import kotlinx.coroutines.CoroutineScope
|
||||
import kotlinx.coroutines.Dispatchers
|
||||
import kotlinx.coroutines.SupervisorJob
|
||||
import kotlinx.coroutines.cancel
|
||||
import kotlinx.coroutines.delay
|
||||
import kotlinx.coroutines.runBlocking
|
||||
import org.junit.Assume
|
||||
import java.io.File
|
||||
import java.net.DatagramSocket
|
||||
import java.util.concurrent.TimeUnit
|
||||
import kotlin.random.Random
|
||||
import kotlin.test.AfterTest
|
||||
import kotlin.test.Test
|
||||
import kotlin.test.assertEquals
|
||||
import kotlin.test.assertFailsWith
|
||||
import kotlin.test.assertTrue
|
||||
|
||||
/**
|
||||
* Drives our direct-path sender against MDK's own direct-path receiver
|
||||
* (`wn stream receive`), the reference implementation of the listening half.
|
||||
*
|
||||
* The direct path inverts the broker path's connection direction — the
|
||||
* RECEIVER listens and the SENDER dials — and drops the control envelope
|
||||
* entirely, so the very first bytes on the stream are a record frame. Both of
|
||||
* those are exactly the kind of thing an implementation happily agrees with
|
||||
* itself about: a self-test would pass with an envelope still on the wire, or
|
||||
* with the wrong ALPN, as long as both ends made the same mistake. Only the
|
||||
* reference receiver can say otherwise.
|
||||
*
|
||||
* `:quic` is a client stack with no server role, so we can only drive the
|
||||
* sender half here. That is also the half the spec makes usable in v1: there
|
||||
* is no start-payload candidate by which a direct receiver advertises its own
|
||||
* endpoint, so the sender always has the address from somewhere else.
|
||||
*
|
||||
* Opt in with `-DmarmotWn=/path/to/wn` (MDK's CLI, `cargo build --release
|
||||
* --bin wn`). Without it the cases skip, so an ordinary `./gradlew test`
|
||||
* never needs the reference implementation on the machine.
|
||||
*/
|
||||
class MarmotQuicDirectInteropTest {
|
||||
private val scope = CoroutineScope(SupervisorJob() + Dispatchers.IO)
|
||||
|
||||
private val wnPath: String? = System.getProperty("marmotWn")
|
||||
|
||||
@AfterTest
|
||||
fun tearDown() {
|
||||
scope.cancel()
|
||||
}
|
||||
|
||||
private fun requireWn(): File {
|
||||
val file = wnPath?.let { File(it) }
|
||||
Assume.assumeTrue(
|
||||
"set -DmarmotWn=/path/to/wn (MDK's CLI) to run the direct-path interop cases",
|
||||
file != null && file.canExecute(),
|
||||
)
|
||||
return file!!
|
||||
}
|
||||
|
||||
/** A UDP port nothing is listening on right now. */
|
||||
private fun freeUdpPort(): Int = DatagramSocket(0).use { it.localPort }
|
||||
|
||||
/**
|
||||
* One run of `wn stream receive`: it binds, waits for a single direct
|
||||
* stream, and prints its JSON result when the stream finishes.
|
||||
*/
|
||||
private class ReferenceReceiver(
|
||||
val process: Process,
|
||||
val port: Int,
|
||||
) {
|
||||
fun awaitResult(timeoutSeconds: Long): String {
|
||||
val finished = process.waitFor(timeoutSeconds, TimeUnit.SECONDS)
|
||||
val out = process.inputStream.readBytes().decodeToString()
|
||||
val err = process.errorStream.readBytes().decodeToString()
|
||||
if (!finished) {
|
||||
process.destroyForcibly()
|
||||
throw AssertionError("the reference receiver never finished. stdout=$out stderr=$err")
|
||||
}
|
||||
return out.ifBlank { throw AssertionError("the reference receiver printed nothing. stderr=$err") }
|
||||
}
|
||||
}
|
||||
|
||||
private fun startReceiver(
|
||||
wn: File,
|
||||
startEventId: ByteArray,
|
||||
): ReferenceReceiver {
|
||||
val port = freeUdpPort()
|
||||
val process =
|
||||
ProcessBuilder(
|
||||
wn.absolutePath,
|
||||
"--json",
|
||||
"stream",
|
||||
"receive",
|
||||
"--bind",
|
||||
"127.0.0.1:$port",
|
||||
"--start-event-id",
|
||||
startEventId.toHex(),
|
||||
).start()
|
||||
return ReferenceReceiver(process, port)
|
||||
}
|
||||
|
||||
private fun crypto(
|
||||
streamId: ByteArray,
|
||||
startEventId: ByteArray,
|
||||
) = AgentTextStreamCrypto(
|
||||
ByteArray(32) { 0x77 },
|
||||
AgentTextStreamKeyContextV1(
|
||||
groupId = ByteArray(32) { 0x01 },
|
||||
streamId = streamId,
|
||||
mlsEpoch = 3,
|
||||
senderId = ByteArray(32) { 0x02 },
|
||||
startEventId = startEventId,
|
||||
),
|
||||
)
|
||||
|
||||
@Test
|
||||
fun ourDirectSenderReachesTheReferenceReceiver() {
|
||||
val wn = requireWn()
|
||||
runBlocking {
|
||||
val streamId = Random.nextBytes(32)
|
||||
val startEventId = Random.nextBytes(32)
|
||||
val receiver = startReceiver(wn, startEventId)
|
||||
// The receiver binds before it accepts; give it a moment so the
|
||||
// dial does not race the bind and report the port as unreachable.
|
||||
delay(1_000)
|
||||
|
||||
val transport =
|
||||
QuicAgentTextStreamTransport(
|
||||
parentScope = scope,
|
||||
// `wn stream receive` mints a throwaway self-signed
|
||||
// certificate per run and only prints it in its final
|
||||
// JSON, so there is nothing to pin ahead of the dial. The
|
||||
// pin is enforced in its own case below.
|
||||
certificateValidator = PermissiveCertificateValidator(),
|
||||
)
|
||||
val stream = transport.sendDirect("quic://127.0.0.1:${receiver.port}", streamId, startEventId)
|
||||
val sender = AgentTextStreamPublisher.open(crypto(streamId, startEventId), InMemoryAgentTextStreamSequenceStore())
|
||||
val chunks = listOf("direct ", "path ", "records")
|
||||
for (text in chunks) {
|
||||
stream.send(sender.publish(AgentTextStreamRecordV1.TYPE_TEXT_DELTA, text.encodeToByteArray()))
|
||||
}
|
||||
stream.finish()
|
||||
stream.close()
|
||||
|
||||
val json = receiver.awaitResult(30)
|
||||
// The reference receiver read our frames, so the ALPN, the
|
||||
// stream direction, the absence of a control envelope and the
|
||||
// 4-byte frame prefix all matched. It was given no key, so the
|
||||
// payloads stay ciphertext to it — what it can confirm is the
|
||||
// routing identity and the sequence.
|
||||
assertTrue(json.contains("\"stream_id\":\"${streamId.toHex()}\""), json)
|
||||
assertTrue(json.contains("\"chunk_count\":${chunks.size}"), json)
|
||||
assertTrue(json.contains("\"seq\":1"), json)
|
||||
assertTrue(json.contains("\"seq\":${chunks.size}"), json)
|
||||
}
|
||||
}
|
||||
|
||||
@Test
|
||||
fun aWrongPinIsRefusedAgainstARealHandshake() {
|
||||
val wn = requireWn()
|
||||
runBlocking {
|
||||
val streamId = Random.nextBytes(32)
|
||||
val startEventId = Random.nextBytes(32)
|
||||
val receiver = startReceiver(wn, startEventId)
|
||||
delay(1_000)
|
||||
|
||||
// A pin over a certificate the receiver is not holding. The
|
||||
// handshake has to fail here rather than at the first record:
|
||||
// once bytes are flowing, "pinned" would have meant nothing.
|
||||
val transport =
|
||||
QuicAgentTextStreamTransport(
|
||||
parentScope = scope,
|
||||
certificateValidator = PinnedCertificateValidator.ofSha256Hex("00".repeat(32)),
|
||||
)
|
||||
val failure =
|
||||
assertFailsWith<MarmotQuicException> {
|
||||
transport.sendDirect("quic://127.0.0.1:${receiver.port}", streamId, startEventId)
|
||||
}
|
||||
assertEquals(MarmotQuicException.Kind.HandshakeFailed, failure.kind, "${failure.message}")
|
||||
|
||||
receiver.process.destroyForcibly()
|
||||
}
|
||||
}
|
||||
|
||||
private fun ByteArray.toHex(): String = joinToString("") { "%02x".format(it) }
|
||||
}
|
||||
+29
@@ -89,6 +89,35 @@ interface MarmotQuicTransport {
|
||||
streamId: ByteArray,
|
||||
startEventId: ByteArray,
|
||||
): MarmotQuicStream
|
||||
|
||||
/**
|
||||
* Dial a receiver directly and stream records to it, point to point.
|
||||
*
|
||||
* The direct path is the binding's other delivery mode, and it is
|
||||
* deliberately smaller than the broker one: the dialed endpoint already
|
||||
* corresponds to one receiver, so there is no room to claim and NO control
|
||||
* envelope — the very first bytes on the stream are a record frame. It
|
||||
* negotiates its own ALPN (`marmot.quic_stream.v1`) so an incompatible
|
||||
* change to either mode cannot reach the other.
|
||||
*
|
||||
* Note the connection direction: the RECEIVER listens and the SENDER
|
||||
* dials. That is inverted from the broker path, where both ends dial the
|
||||
* broker, and it is why v1 has no start-payload discovery for this mode —
|
||||
* a start payload advertises broker candidates only, and there is no
|
||||
* candidate shape by which a direct receiver publishes its own endpoint.
|
||||
* So this is usable only when the sender already knows where to dial:
|
||||
* out-of-band configuration, a dev/test peer, a preconfigured pair.
|
||||
*
|
||||
* [startEventId] never crosses the wire here. It stays in the signature
|
||||
* because the caller still binds it into the record key and transcript
|
||||
* hash, and because a direct endpoint that was told it out of band should
|
||||
* be checking the same pair we are.
|
||||
*/
|
||||
suspend fun sendDirect(
|
||||
candidate: String,
|
||||
streamId: ByteArray,
|
||||
startEventId: ByteArray,
|
||||
): MarmotQuicStream
|
||||
}
|
||||
|
||||
/** Why a candidate turned out to be unusable. */
|
||||
|
||||
@@ -0,0 +1,119 @@
|
||||
/*
|
||||
* 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.quic.tls
|
||||
|
||||
import com.vitorpamplona.quic.QuicCodecException
|
||||
import java.security.NoSuchAlgorithmException
|
||||
import java.security.PublicKey
|
||||
import java.security.Signature
|
||||
import java.security.spec.MGF1ParameterSpec
|
||||
import java.security.spec.PSSParameterSpec
|
||||
|
||||
/**
|
||||
* TLS 1.3 `CertificateVerify` verification, shared by every JDK/Android-backed
|
||||
* [CertificateValidator].
|
||||
*
|
||||
* Chain policy and CertificateVerify are separate questions, and the split
|
||||
* matters: how a validator decides it likes a certificate (a trust store, a
|
||||
* pinned fingerprint) is policy, but proving the peer holds the matching
|
||||
* private key is not optional under any policy. A validator that pinned a
|
||||
* fingerprint and skipped this would accept anyone who could copy a public
|
||||
* certificate off the wire.
|
||||
*/
|
||||
internal object CertificateVerifySignature {
|
||||
fun verify(
|
||||
publicKey: PublicKey,
|
||||
signatureAlgorithm: Int,
|
||||
signature: ByteArray,
|
||||
transcriptHash: ByteArray,
|
||||
) {
|
||||
// RFC 8446 §4.4.3 — the signed content is:
|
||||
// 64 spaces || "TLS 1.3, server CertificateVerify" || 0x00 || transcript_hash
|
||||
val context = "TLS 1.3, server CertificateVerify".encodeToByteArray()
|
||||
val signedData = ByteArray(64 + context.size + 1 + transcriptHash.size)
|
||||
for (i in 0 until 64) signedData[i] = 0x20
|
||||
context.copyInto(signedData, 64)
|
||||
signedData[64 + context.size] = 0x00
|
||||
transcriptHash.copyInto(signedData, 64 + context.size + 1)
|
||||
|
||||
val sig = jcaSignatureFor(signatureAlgorithm)
|
||||
sig.initVerify(publicKey)
|
||||
sig.update(signedData)
|
||||
if (!sig.verify(signature)) {
|
||||
throw QuicCodecException("CertificateVerify signature did not verify")
|
||||
}
|
||||
}
|
||||
|
||||
private fun jcaSignatureFor(algorithm: Int): Signature =
|
||||
when (algorithm) {
|
||||
TlsConstants.SIG_ECDSA_SECP256R1_SHA256 -> {
|
||||
Signature.getInstance("SHA256withECDSA")
|
||||
}
|
||||
|
||||
TlsConstants.SIG_ECDSA_SECP384R1_SHA384 -> {
|
||||
Signature.getInstance("SHA384withECDSA")
|
||||
}
|
||||
|
||||
TlsConstants.SIG_RSA_PSS_RSAE_SHA256 -> {
|
||||
rsaPss("SHA-256", 32)
|
||||
}
|
||||
|
||||
TlsConstants.SIG_RSA_PSS_RSAE_SHA384 -> {
|
||||
rsaPss("SHA-384", 48)
|
||||
}
|
||||
|
||||
TlsConstants.SIG_RSA_PSS_RSAE_SHA512 -> {
|
||||
rsaPss("SHA-512", 64)
|
||||
}
|
||||
|
||||
TlsConstants.SIG_ED25519 -> {
|
||||
try {
|
||||
// JCA "Ed25519" was added to Android Conscrypt in API 33.
|
||||
// On API 26–32 (our minSdk floor) this throws — surface
|
||||
// it as a clean QuicCodecException so the read loop maps
|
||||
// to CONNECTION_CLOSE rather than crashing the parser.
|
||||
Signature.getInstance("Ed25519")
|
||||
} catch (_: NoSuchAlgorithmException) {
|
||||
throw QuicCodecException(
|
||||
"Ed25519 not supported on this platform " +
|
||||
"(requires Android API 33+ or a JDK with the EdDSA provider)",
|
||||
)
|
||||
}
|
||||
}
|
||||
|
||||
// Audit-4 #2: rsa_pkcs1_* schemes are forbidden in CertificateVerify
|
||||
// by RFC 8446 §4.2.3 (only allowed in CertificateRequest for
|
||||
// legacy compat). Accepting them allowed a server to sign with
|
||||
// weaker PKCS#1 v1.5 instead of RSA-PSS.
|
||||
else -> {
|
||||
throw QuicCodecException("unsupported signature algorithm 0x${algorithm.toString(16)}")
|
||||
}
|
||||
}
|
||||
|
||||
private fun rsaPss(
|
||||
digest: String,
|
||||
saltLen: Int,
|
||||
): Signature {
|
||||
val sig = Signature.getInstance("RSASSA-PSS")
|
||||
sig.setParameter(PSSParameterSpec(digest, "MGF1", MGF1ParameterSpec(digest), saltLen, 1))
|
||||
return sig
|
||||
}
|
||||
}
|
||||
@@ -26,12 +26,8 @@ import java.lang.reflect.InvocationTargetException
|
||||
import java.net.IDN
|
||||
import java.net.InetAddress
|
||||
import java.security.KeyStore
|
||||
import java.security.NoSuchAlgorithmException
|
||||
import java.security.Signature
|
||||
import java.security.cert.CertificateFactory
|
||||
import java.security.cert.X509Certificate
|
||||
import java.security.spec.MGF1ParameterSpec
|
||||
import java.security.spec.PSSParameterSpec
|
||||
import javax.net.ssl.TrustManagerFactory
|
||||
import javax.net.ssl.X509TrustManager
|
||||
|
||||
@@ -117,77 +113,7 @@ class JdkCertificateValidator(
|
||||
transcriptHash: ByteArray,
|
||||
) {
|
||||
val cert = leafCert ?: throw QuicCodecException("CertificateVerify before Certificate")
|
||||
|
||||
// RFC 8446 §4.4.3 — the signed content is:
|
||||
// 64 spaces || "TLS 1.3, server CertificateVerify" || 0x00 || transcript_hash
|
||||
val context = "TLS 1.3, server CertificateVerify".encodeToByteArray()
|
||||
val signedData = ByteArray(64 + context.size + 1 + transcriptHash.size)
|
||||
for (i in 0 until 64) signedData[i] = 0x20
|
||||
context.copyInto(signedData, 64)
|
||||
signedData[64 + context.size] = 0x00
|
||||
transcriptHash.copyInto(signedData, 64 + context.size + 1)
|
||||
|
||||
val sig = jcaSignatureFor(signatureAlgorithm)
|
||||
sig.initVerify(cert.publicKey)
|
||||
sig.update(signedData)
|
||||
if (!sig.verify(signature)) {
|
||||
throw QuicCodecException("CertificateVerify signature did not verify")
|
||||
}
|
||||
}
|
||||
|
||||
private fun jcaSignatureFor(algorithm: Int): Signature =
|
||||
when (algorithm) {
|
||||
TlsConstants.SIG_ECDSA_SECP256R1_SHA256 -> {
|
||||
Signature.getInstance("SHA256withECDSA")
|
||||
}
|
||||
|
||||
TlsConstants.SIG_ECDSA_SECP384R1_SHA384 -> {
|
||||
Signature.getInstance("SHA384withECDSA")
|
||||
}
|
||||
|
||||
TlsConstants.SIG_RSA_PSS_RSAE_SHA256 -> {
|
||||
rsaPss("SHA-256", 32)
|
||||
}
|
||||
|
||||
TlsConstants.SIG_RSA_PSS_RSAE_SHA384 -> {
|
||||
rsaPss("SHA-384", 48)
|
||||
}
|
||||
|
||||
TlsConstants.SIG_RSA_PSS_RSAE_SHA512 -> {
|
||||
rsaPss("SHA-512", 64)
|
||||
}
|
||||
|
||||
TlsConstants.SIG_ED25519 -> {
|
||||
try {
|
||||
// JCA "Ed25519" was added to Android Conscrypt in API 33.
|
||||
// On API 26–32 (our minSdk floor) this throws — surface
|
||||
// it as a clean QuicCodecException so the read loop maps
|
||||
// to CONNECTION_CLOSE rather than crashing the parser.
|
||||
Signature.getInstance("Ed25519")
|
||||
} catch (_: NoSuchAlgorithmException) {
|
||||
throw QuicCodecException(
|
||||
"Ed25519 not supported on this platform " +
|
||||
"(requires Android API 33+ or a JDK with the EdDSA provider)",
|
||||
)
|
||||
}
|
||||
}
|
||||
|
||||
// Audit-4 #2: rsa_pkcs1_* schemes are forbidden in CertificateVerify
|
||||
// by RFC 8446 §4.2.3 (only allowed in CertificateRequest for
|
||||
// legacy compat). Accepting them allowed a server to sign with
|
||||
// weaker PKCS#1 v1.5 instead of RSA-PSS.
|
||||
else -> {
|
||||
throw QuicCodecException("unsupported signature algorithm 0x${algorithm.toString(16)}")
|
||||
}
|
||||
}
|
||||
|
||||
private fun rsaPss(
|
||||
digest: String,
|
||||
saltLen: Int,
|
||||
): Signature {
|
||||
val sig = Signature.getInstance("RSASSA-PSS")
|
||||
sig.setParameter(PSSParameterSpec(digest, "MGF1", MGF1ParameterSpec(digest), saltLen, 1))
|
||||
return sig
|
||||
CertificateVerifySignature.verify(cert.publicKey, signatureAlgorithm, signature, transcriptHash)
|
||||
}
|
||||
|
||||
private fun hostnameMatches(
|
||||
|
||||
@@ -0,0 +1,147 @@
|
||||
/*
|
||||
* 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.quic.tls
|
||||
|
||||
import com.vitorpamplona.quic.QuicCodecException
|
||||
import java.io.ByteArrayInputStream
|
||||
import java.security.MessageDigest
|
||||
import java.security.cert.CertificateFactory
|
||||
import java.security.cert.X509Certificate
|
||||
|
||||
/**
|
||||
* A certificate validator that trusts exactly the endpoints whose leaf
|
||||
* certificate matches a configured SHA-256 fingerprint.
|
||||
*
|
||||
* This is the "self-signed endpoint" case, and it is a real one: Marmot's raw
|
||||
* QUIC binding says a preview endpoint or broker may be self-signed and that a
|
||||
* client MAY pin it by exact DER or by SHA-256 fingerprint through local
|
||||
* configuration. The alternative in use until now was accepting every
|
||||
* certificate, which is not a weaker trust model — it is no trust model, and
|
||||
* anyone on the path can be the broker.
|
||||
*
|
||||
* Pinning replaces the chain and the hostname check, and only those. It does
|
||||
* NOT replace proof of possession: the peer still has to sign the TLS
|
||||
* transcript with the pinned certificate's private key ([verifySignature]),
|
||||
* so copying a public certificate off the wire buys an attacker nothing.
|
||||
*
|
||||
* The pin is over the leaf's DER bytes exactly as the peer sent them, which is
|
||||
* what `openssl x509 -outform der | sha256sum` prints and what a broker
|
||||
* operator can therefore publish alongside its address. An expired or
|
||||
* not-yet-valid pinned certificate is still refused: pinning says WHICH
|
||||
* certificate, not that any certificate will do forever.
|
||||
*/
|
||||
class PinnedCertificateValidator(
|
||||
pins: Collection<ByteArray>,
|
||||
) : CertificateValidator {
|
||||
private val pins: List<ByteArray> =
|
||||
pins.map {
|
||||
require(it.size == SHA256_LEN) { "a certificate pin is a $SHA256_LEN-byte SHA-256 digest, got ${it.size}" }
|
||||
it.copyOf()
|
||||
}
|
||||
|
||||
private var leafCert: X509Certificate? = null
|
||||
|
||||
init {
|
||||
require(this.pins.isNotEmpty()) { "a pinned validator needs at least one pin" }
|
||||
}
|
||||
|
||||
override fun validateChain(
|
||||
chain: List<ByteArray>,
|
||||
expectedHost: String,
|
||||
) {
|
||||
if (chain.isEmpty()) throw QuicCodecException("server sent empty certificate chain")
|
||||
|
||||
// Only the leaf is pinned. The rest of the chain is not consulted at
|
||||
// all — with a pin there is no path to build and no issuer to trust,
|
||||
// and a self-signed endpoint has no chain to speak of.
|
||||
val leafDer = chain[0]
|
||||
val fingerprint = MessageDigest.getInstance("SHA-256").digest(leafDer)
|
||||
if (pins.none { it.contentEqualsConstantTime(fingerprint) }) {
|
||||
throw QuicCodecException("certificate does not match any pinned SHA-256 fingerprint")
|
||||
}
|
||||
|
||||
val parsed =
|
||||
try {
|
||||
CertificateFactory
|
||||
.getInstance("X.509")
|
||||
.generateCertificate(ByteArrayInputStream(leafDer)) as X509Certificate
|
||||
} catch (t: Throwable) {
|
||||
throw QuicCodecException("pinned certificate parse failed: ${t.message}", t)
|
||||
}
|
||||
try {
|
||||
parsed.checkValidity()
|
||||
} catch (t: Throwable) {
|
||||
throw QuicCodecException("pinned certificate is not currently valid: ${t.message}", t)
|
||||
}
|
||||
|
||||
// No hostname verification: the pin already names one certificate, and
|
||||
// a self-signed preview endpoint reached by IP literal typically has no
|
||||
// name to check against. `expectedHost` stays in the signature because
|
||||
// the interface is shared with trust-store validation.
|
||||
leafCert = parsed
|
||||
}
|
||||
|
||||
override fun verifySignature(
|
||||
signatureAlgorithm: Int,
|
||||
signature: ByteArray,
|
||||
transcriptHash: ByteArray,
|
||||
) {
|
||||
val cert = leafCert ?: throw QuicCodecException("CertificateVerify before Certificate")
|
||||
CertificateVerifySignature.verify(cert.publicKey, signatureAlgorithm, signature, transcriptHash)
|
||||
}
|
||||
|
||||
companion object {
|
||||
const val SHA256_LEN = 32
|
||||
|
||||
/**
|
||||
* Pin by SHA-256 fingerprint, written as hex.
|
||||
*
|
||||
* Colons and whitespace are accepted and ignored so the output of
|
||||
* `openssl x509 -fingerprint -sha256` can be pasted in as-is.
|
||||
*/
|
||||
fun ofSha256Hex(vararg fingerprints: String): PinnedCertificateValidator = PinnedCertificateValidator(fingerprints.map { parseHexDigest(it) })
|
||||
|
||||
/**
|
||||
* Pin by the certificate's exact DER bytes.
|
||||
*
|
||||
* The DER is reduced to its own SHA-256 immediately: "exact DER" and
|
||||
* "its fingerprint" are the same pin, and keeping one representation
|
||||
* means one comparison path to get right.
|
||||
*/
|
||||
fun ofDer(vararg certificates: ByteArray): PinnedCertificateValidator =
|
||||
PinnedCertificateValidator(
|
||||
certificates.map { MessageDigest.getInstance("SHA-256").digest(it) },
|
||||
)
|
||||
|
||||
private fun parseHexDigest(raw: String): ByteArray {
|
||||
val cleaned = raw.filterNot { it == ':' || it.isWhitespace() }
|
||||
require(cleaned.length == SHA256_LEN * 2) {
|
||||
"a SHA-256 fingerprint is ${SHA256_LEN * 2} hex characters, got ${cleaned.length}"
|
||||
}
|
||||
return ByteArray(SHA256_LEN) { i ->
|
||||
val hi = Character.digit(cleaned[i * 2], 16)
|
||||
val lo = Character.digit(cleaned[i * 2 + 1], 16)
|
||||
require(hi >= 0 && lo >= 0) { "a SHA-256 fingerprint must be hex" }
|
||||
((hi shl 4) or lo).toByte()
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,163 @@
|
||||
/*
|
||||
* 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.quic.tls
|
||||
|
||||
import com.vitorpamplona.quic.QuicCodecException
|
||||
import org.junit.Test
|
||||
import java.security.MessageDigest
|
||||
import kotlin.test.assertEquals
|
||||
import kotlin.test.assertFailsWith
|
||||
import kotlin.test.assertTrue
|
||||
|
||||
/**
|
||||
* Pinning is the trust model for a self-signed Marmot preview endpoint, so the
|
||||
* interesting cases are all the ways it must REFUSE. A pin that quietly
|
||||
* accepts the wrong certificate is worse than no pin: the operator believes
|
||||
* they configured something.
|
||||
*/
|
||||
class PinnedCertificateValidatorTest {
|
||||
/**
|
||||
* A self-signed P-256 certificate for `marmot-preview.test` with an
|
||||
* `IP:127.0.0.1` SAN, valid for a century so this test does not become a
|
||||
* time bomb. Nothing signs with it — only its DER bytes matter here.
|
||||
*/
|
||||
private val leafDer =
|
||||
(
|
||||
"308201a43082014aa00302010202146633303a60bb8854f6c7128247d681e5e72afff2300a06082a8648ce3d040302301e31" +
|
||||
"1c301a06035504030c136d61726d6f742d707265766965772e746573743020170d3236303930393135323432325a180f3231" +
|
||||
"3236303831363135323432325a301e311c301a06035504030c136d61726d6f742d707265766965772e746573743059301306" +
|
||||
"072a8648ce3d020106082a8648ce3d030107034200042498e9233c2eb1e6302fb98d0761205c0cd38e9eb72ea89651acb8f1" +
|
||||
"a97c32f2f658dfef6a2c1e118f2874f2ae6607c3499814c00d58cf6ebd894b2445742d40a3643062301d0603551d0e041604" +
|
||||
"14b9b33ffe3a961e40004413767b7d84acc1069cc5301f0603551d23041830168014b9b33ffe3a961e40004413767b7d84ac" +
|
||||
"c1069cc5300f0603551d130101ff040530030101ff300f0603551d110408300687047f000001300a06082a8648ce3d040302" +
|
||||
"0348003045022100fea992edecb7f0b3e92d798fac6eca4f728784d879a2d96999e6ec458fa7813002207b921e876f9f38f2" +
|
||||
"ede262f55fa8b8968a4373c5bc0ec8326cdd1f8c8c759743"
|
||||
).hexToBytes()
|
||||
|
||||
private val fingerprintHex = "c0f7a502abf9b8f0657ec5c39eaaf4bcff21bb530c8afb81ecbf0b38b76f676c"
|
||||
|
||||
@Test
|
||||
fun `the pin is the SHA-256 of the leaf DER exactly as sent`() {
|
||||
// The digest a broker operator publishes comes from
|
||||
// `openssl x509 -outform der | sha256sum`, so this is the value the
|
||||
// whole design hangs on. If it were over anything else — the PEM, the
|
||||
// public key, a re-encoded cert — a correctly configured pin would
|
||||
// reject a correct endpoint.
|
||||
val digest = MessageDigest.getInstance("SHA-256").digest(leafDer)
|
||||
assertEquals(fingerprintHex, digest.joinToString("") { "%02x".format(it) })
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `a matching fingerprint validates`() {
|
||||
PinnedCertificateValidator.ofSha256Hex(fingerprintHex).validateChain(listOf(leafDer), "marmot-preview.test")
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `the host is not checked because the pin already named the certificate`() {
|
||||
// A self-signed preview endpoint reached by IP literal usually has no
|
||||
// name worth checking, and the pin is a stronger statement than any
|
||||
// name would be. This asserts the deliberate difference from
|
||||
// JdkCertificateValidator rather than an accident.
|
||||
PinnedCertificateValidator.ofSha256Hex(fingerprintHex).validateChain(listOf(leafDer), "not-the-cert-name.example")
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `a different fingerprint is refused`() {
|
||||
val other = "00".repeat(32)
|
||||
val e =
|
||||
assertFailsWith<QuicCodecException> {
|
||||
PinnedCertificateValidator.ofSha256Hex(other).validateChain(listOf(leafDer), "marmot-preview.test")
|
||||
}
|
||||
assertTrue(e.message!!.contains("pinned SHA-256"), e.message)
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `one matching pin among several is enough`() {
|
||||
PinnedCertificateValidator
|
||||
.ofSha256Hex("11".repeat(32), fingerprintHex, "22".repeat(32))
|
||||
.validateChain(listOf(leafDer), "marmot-preview.test")
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `pinning by DER is the same pin as pinning by its fingerprint`() {
|
||||
PinnedCertificateValidator.ofDer(leafDer).validateChain(listOf(leafDer), "marmot-preview.test")
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `an openssl-formatted fingerprint is accepted verbatim`() {
|
||||
// `openssl x509 -fingerprint -sha256` prints colon-separated upper
|
||||
// case. Making the operator strip that by hand is how a pin ends up
|
||||
// mistyped.
|
||||
val colonised = fingerprintHex.chunked(2).joinToString(":").uppercase()
|
||||
PinnedCertificateValidator.ofSha256Hex(colonised).validateChain(listOf(leafDer), "marmot-preview.test")
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `an empty chain is refused`() {
|
||||
assertFailsWith<QuicCodecException> {
|
||||
PinnedCertificateValidator.ofSha256Hex(fingerprintHex).validateChain(emptyList(), "marmot-preview.test")
|
||||
}
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `only the leaf is pinned, so a matching cert deeper in the chain does not count`() {
|
||||
// Pinning the leaf and then honouring a match anywhere in the chain
|
||||
// would let a peer present any certificate it likes and append the
|
||||
// pinned one behind it.
|
||||
assertFailsWith<QuicCodecException> {
|
||||
PinnedCertificateValidator
|
||||
.ofSha256Hex(fingerprintHex)
|
||||
.validateChain(listOf(byteArrayOf(1, 2, 3), leafDer), "marmot-preview.test")
|
||||
}
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `a pin that is not a SHA-256 digest is rejected at construction`() {
|
||||
assertFailsWith<IllegalArgumentException> { PinnedCertificateValidator.ofSha256Hex("abcd") }
|
||||
assertFailsWith<IllegalArgumentException> { PinnedCertificateValidator.ofSha256Hex("zz".repeat(32)) }
|
||||
assertFailsWith<IllegalArgumentException> { PinnedCertificateValidator(emptyList()) }
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `CertificateVerify before Certificate is refused`() {
|
||||
// Order matters: without a leaf there is no key to check the signature
|
||||
// against, and silently passing would make the pin decorative.
|
||||
assertFailsWith<QuicCodecException> {
|
||||
PinnedCertificateValidator
|
||||
.ofSha256Hex(fingerprintHex)
|
||||
.verifySignature(TlsConstants.SIG_ECDSA_SECP256R1_SHA256, ByteArray(64), ByteArray(32))
|
||||
}
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `a garbage signature does not verify against the pinned key`() {
|
||||
val validator = PinnedCertificateValidator.ofSha256Hex(fingerprintHex)
|
||||
validator.validateChain(listOf(leafDer), "marmot-preview.test")
|
||||
assertFailsWith<Exception> {
|
||||
validator.verifySignature(TlsConstants.SIG_ECDSA_SECP256R1_SHA256, ByteArray(70), ByteArray(32))
|
||||
}
|
||||
}
|
||||
|
||||
private fun String.hexToBytes(): ByteArray =
|
||||
ByteArray(length / 2) { i ->
|
||||
((Character.digit(this[i * 2], 16) shl 4) or Character.digit(this[i * 2 + 1], 16)).toByte()
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user