diff --git a/quic/src/commonMain/kotlin/com/vitorpamplona/quic/connection/PathValidator.kt b/quic/src/commonMain/kotlin/com/vitorpamplona/quic/connection/PathValidator.kt new file mode 100644 index 000000000..ad820d87d --- /dev/null +++ b/quic/src/commonMain/kotlin/com/vitorpamplona/quic/connection/PathValidator.kt @@ -0,0 +1,534 @@ +/* + * 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.connection + +import com.vitorpamplona.quartz.utils.RandomInstance + +/** + * Single connection ID the peer (server) has issued to us via a + * `NEW_CONNECTION_ID` frame (RFC 9000 §19.15) and we have NOT yet + * retired. The initial DCID negotiated during the handshake is + * implicitly sequence number 0 (RFC 9000 §5.1.1). + * + * `connectionId` is the bytes the writer would stamp into outbound + * short-header packets; `statelessResetToken` is the 16-byte token + * we'd compare against an incoming stateless reset. + */ +data class PeerConnectionIdEntry( + val sequenceNumber: Long, + val connectionId: ByteArray, + val statelessResetToken: ByteArray, +) { + init { + require(connectionId.isNotEmpty()) { "connection id must not be empty" } + require(statelessResetToken.size == 16) { "stateless reset token must be 16 bytes" } + } + + override fun equals(other: Any?): Boolean { + if (this === other) return true + if (other !is PeerConnectionIdEntry) return false + return sequenceNumber == other.sequenceNumber && + connectionId.contentEquals(other.connectionId) && + statelessResetToken.contentEquals(other.statelessResetToken) + } + + override fun hashCode(): Int { + var h = sequenceNumber.hashCode() + h = 31 * h + connectionId.contentHashCode() + h = 31 * h + statelessResetToken.contentHashCode() + return h + } +} + +/** + * Where the path validation state machine currently sits (RFC 9000 + * §8.2 / §9). Most installations spend their entire life in [Idle]; + * the transition to [Validating] happens when the driver detects the + * current path looks dead (consecutive PTOs without ACKs) and asks + * the connection to migrate to a fresh DCID. + * + * - [Idle]: nothing in flight; the writer is happily using the + * current DCID. + * - [Validating]: a `PATH_CHALLENGE` is on the wire (or queued for + * the next outbound) and we're waiting for a matching + * `PATH_RESPONSE` with the SAME 8-byte payload. + * - [Succeeded]: the peer echoed the payload; the writer has + * switched to the new DCID and a `RETIRE_CONNECTION_ID` for the + * old sequence number is queued. + * - [Failed]: validation didn't complete inside the 3*PTO window + * (RFC 9000 §8.2.4) — caller may retry with a different CID, or + * surface the failure as a connection close. + */ +sealed class PathValidationState { + object Idle : PathValidationState() + + data class Validating( + val challengeData: ByteArray, + val newCidSequence: Long, + val newCidBytes: ByteArray, + val priorCidSequence: Long, + val startedAtMillis: Long, + val priorPtoMillis: Long, + ) : PathValidationState() { + override fun equals(other: Any?): Boolean { + if (this === other) return true + if (other !is Validating) return false + return challengeData.contentEquals(other.challengeData) && + newCidSequence == other.newCidSequence && + newCidBytes.contentEquals(other.newCidBytes) && + priorCidSequence == other.priorCidSequence && + startedAtMillis == other.startedAtMillis && + priorPtoMillis == other.priorPtoMillis + } + + override fun hashCode(): Int { + var h = challengeData.contentHashCode() + h = 31 * h + newCidSequence.hashCode() + h = 31 * h + newCidBytes.contentHashCode() + h = 31 * h + priorCidSequence.hashCode() + h = 31 * h + startedAtMillis.hashCode() + h = 31 * h + priorPtoMillis.hashCode() + return h + } + } + + object Succeeded : PathValidationState() + + object Failed : PathValidationState() +} + +/** + * Outcome of [PathValidator.tryStartValidation], used by the caller + * (driver / test) to know whether a probe was actually issued. + */ +enum class PathMigrationResult { + /** A `PATH_CHALLENGE` was queued; state is [PathValidationState.Validating]. */ + Started, + + /** + * Already in [PathValidationState.Validating] — the previous + * attempt hasn't resolved yet. Caller should not pile on; the + * existing challenge will resolve via PATH_RESPONSE or fail + * via [PathValidator.checkValidationTimeout]. + */ + AlreadyInProgress, + + /** + * No spare CID in the pool. The peer hasn't yet issued a + * NEW_CONNECTION_ID we can use, or we've already drained the + * pool. Caller should wait for fresh peer offers. + */ + NoSpareCid, +} + +/** + * Per-connection state for client-initiated path validation + DCID + * rotation (RFC 9000 §9). Owned by [QuicConnection]; the parser and + * writer call into it under `streamsLock`. + * + * Spec mapping: + * - **§5.1.1** initial DCID is sequence number 0. + * - **§5.1.2** at most `active_connection_id_limit` CIDs + * simultaneously usable; the limit is the PEER's transport + * parameter telling us how many of OUR source CIDs they'll + * hold, but reciprocally we cap the peer's pool the same way + * so a buggy server can't pin arbitrary memory by spamming + * NEW_CONNECTION_ID frames. + * - **§5.1.2** `retire_prior_to` in a NEW_CONNECTION_ID frame + * forces us to retire every previously-issued CID with a + * smaller sequence number (and queue RETIRE_CONNECTION_ID for + * each). The `retire_prior_to` value MUST NOT decrease across + * successive frames; we treat a regression as a peer protocol + * violation and abort the connection upstream. + * - **§8.2 / §9** path validation: pick a spare CID, issue + * PATH_CHALLENGE with a cryptographically random 8-byte + * payload, await a PATH_RESPONSE that echoes the payload byte + * for byte. RFC 9000 §8.2.4: validation MUST be abandoned after + * `3 * PTO` of inactivity. + * - **§19.16** RETIRE_CONNECTION_ID is reliable; the dispatcher + * re-queues on loss until ACK'd. + * + * Thread safety: every method here assumes the caller holds the + * connection's `streamsLock`. The pool / outstanding-challenge maps + * are plain mutable collections. + */ +class PathValidator( + /** + * Sequence number of the CID the writer is currently stamping + * into outbound short-header packets. Initial value 0 per RFC + * 9000 §5.1.1 — the DCID negotiated at handshake. After a + * successful validation [activeCidSequence] is the new + * sequence; the prior sequence is in `pendingRetireSequences` + * waiting to ride out as a RETIRE_CONNECTION_ID. + */ + initialActiveCidSequence: Long = 0L, + /** + * Cap on the number of unused CIDs we'll buffer. Defaults to + * the peer's `active_connection_id_limit` transport parameter, + * but the connection caller may pass a tighter cap. Excess + * NEW_CONNECTION_ID offers past this cap are silently dropped + * — the peer is in violation of the limit it asked us to + * advertise (or never advertised one) and we don't owe it + * indefinite buffering. + */ + private val maxUnusedCids: Int = DEFAULT_MAX_UNUSED_CIDS, + /** + * Cryptographically random 8-byte payload generator used for + * outbound PATH_CHALLENGE. Default uses [RandomInstance]; tests + * can substitute a deterministic supplier so assertions can + * compare the on-wire payload against a known value. + */ + private val challengePayloadFactory: () -> ByteArray = { RandomInstance.bytes(8) }, +) { + /** + * Sequence number → entry. Holds CIDs the peer has issued and + * we have NOT yet retired or activated. The active CID itself + * is NOT in this map (it's held by the connection's + * `destinationConnectionId` field); only the spare pool lives + * here. + * + * LinkedHashMap so iteration order matches issuance order — when + * we pick a CID for a fresh validation we prefer the lowest + * sequence number for FIFO fairness. + */ + private val unusedCids: LinkedHashMap = LinkedHashMap() + + /** + * Highest `retire_prior_to` value the peer has ever advertised. + * Per RFC 9000 §5.1.2 this MUST NOT decrease; a decrease is a + * peer protocol violation. Used to decide whether a + * just-arrived NEW_CONNECTION_ID with a sequence number ≤ + * [retirePriorToWatermark] is stale (drop it) or whether a + * previously-cached entry has now been forced into retirement. + */ + var retirePriorToWatermark: Long = 0L + private set + + var activeCidSequence: Long = initialActiveCidSequence + private set + + /** + * Sequence numbers of CIDs we owe the peer a + * `RETIRE_CONNECTION_ID` for. Drained by the writer on the + * next outbound application packet. Re-populated by the loss + * dispatcher if our RETIRE_CONNECTION_ID was lost. + */ + val pendingRetireSequences: ArrayDeque = ArrayDeque() + + /** + * Outbound `PATH_CHALLENGE` payloads the writer should drain on + * the next application-level packet. Populated by + * [tryStartValidation]; emptied by the writer after encoding + * (the writer also records a [com.vitorpamplona.quic.connection.recovery.RecoveryToken.PathChallenge] + * so the loss dispatcher can decide to retransmit or abandon). + */ + val pendingChallenges: ArrayDeque = ArrayDeque() + + var state: PathValidationState = PathValidationState.Idle + private set + + /** + * Total number of validations that have completed successfully + * over the lifetime of this connection. Diagnostic only — used + * by tests and qlog to assert "exactly one rotation happened". + */ + var successfulValidations: Long = 0L + private set + + /** + * Total number of validations that timed out / were abandoned. + * Counter-part of [successfulValidations]. + */ + var failedValidations: Long = 0L + private set + + /** + * Number of unused CIDs currently in the pool. Used by tests + * (and the connection's diagnostic surface) to assert the peer + * actually offered new CIDs we can rotate to. + */ + fun unusedCount(): Int = unusedCids.size + + fun unusedSequences(): List = unusedCids.keys.toList() + + /** + * Record a peer-issued `NEW_CONNECTION_ID` (RFC 9000 §19.15). + * Returns the result code so the caller (parser) can decide + * whether to close the connection on a peer protocol violation. + * + * The semantics: + * - If [retirePriorTo] is greater than [sequenceNumber], that's + * a §19.15 violation (the peer can't tell us to retire its + * own brand-new CID). + * - If [retirePriorTo] decreased, that's a §5.1.2 violation. + * - If the new entry would push the pool past [maxUnusedCids], + * return [Result.PoolFull] — the peer over-issued. + * - If [sequenceNumber] is below the current + * [retirePriorToWatermark], the offer is already-retired and + * we silently drop it (the caller MUST still queue a + * RETIRE_CONNECTION_ID for it per §5.1.2). + * - Duplicate sequence number with matching CID/token is + * idempotent ([Result.Duplicate]); a duplicate sequence with + * DIFFERENT bytes is a §19.15 violation. + * + * Side effect on success: any cached entry whose sequence + * number is now below the new [retirePriorTo] is moved into + * [pendingRetireSequences] (we owe the peer a retire) and + * dropped from [unusedCids]. If the active CID itself is + * forced-retired, the active sequence is bumped (the caller + * must rotate the active CID — see [forceRetireActiveIfNeeded]). + */ + fun recordPeerNewConnectionId( + sequenceNumber: Long, + retirePriorTo: Long, + connectionId: ByteArray, + statelessResetToken: ByteArray, + ): RecordResult { + if (retirePriorTo > sequenceNumber) return RecordResult.RetirePriorToExceedsSequence + if (connectionId.isEmpty() || connectionId.size > 20) return RecordResult.InvalidCidLength + if (statelessResetToken.size != 16) return RecordResult.InvalidStatelessResetToken + + // RFC 9000 §19.15: a smaller `retire_prior_to` than what we + // previously saw "MUST be treated as the largest one it has + // seen". This handles reordered NEW_CONNECTION_ID arrivals + // — not a protocol violation, just clamp. + val effectiveRetirePriorTo = if (retirePriorTo < retirePriorToWatermark) retirePriorToWatermark else retirePriorTo + + // §5.1.2: bump the watermark and force-retire any cached + // entries that fall under it. This step happens BEFORE the + // duplicate check so a sequence number that's just been + // dropped doesn't survive on a stale entry. + if (effectiveRetirePriorTo > retirePriorToWatermark) { + retirePriorToWatermark = effectiveRetirePriorTo + val it = unusedCids.entries.iterator() + while (it.hasNext()) { + val (seq, _) = it.next() + if (seq < retirePriorToWatermark) { + queueRetireSequence(seq) + it.remove() + } + } + } + + // Stale offer: the peer is pushing a sequence number we've + // already retired. Per §5.1.2 we owe a RETIRE_CONNECTION_ID + // back; queue it but don't add the entry to the pool. + if (sequenceNumber < retirePriorToWatermark) { + queueRetireSequence(sequenceNumber) + return RecordResult.AlreadyRetired + } + + // Duplicate handling. + val existing = unusedCids[sequenceNumber] + if (existing != null) { + if (!existing.connectionId.contentEquals(connectionId) || + !existing.statelessResetToken.contentEquals(statelessResetToken) + ) { + return RecordResult.DuplicateSequenceMismatch + } + return RecordResult.Duplicate + } + + if (unusedCids.size >= maxUnusedCids) return RecordResult.PoolFull + + unusedCids[sequenceNumber] = + PeerConnectionIdEntry( + sequenceNumber = sequenceNumber, + connectionId = connectionId.copyOf(), + statelessResetToken = statelessResetToken.copyOf(), + ) + return RecordResult.Stored + } + + /** + * Result of [recordPeerNewConnectionId]. + */ + enum class RecordResult { + Stored, + Duplicate, + DuplicateSequenceMismatch, + AlreadyRetired, + PoolFull, + RetirePriorToExceedsSequence, + InvalidCidLength, + InvalidStatelessResetToken, + } + + /** + * Append [sequenceNumber] to [pendingRetireSequences] iff it + * isn't already queued. Idempotent — a duplicate request from + * the loss dispatcher (re-queue on loss) finds the entry + * already there and is a no-op. + */ + fun queueRetireSequence(sequenceNumber: Long) { + if (!pendingRetireSequences.contains(sequenceNumber)) { + pendingRetireSequences.addLast(sequenceNumber) + } + } + + /** + * Begin a fresh path-validation attempt. Picks the lowest-sequence + * unused CID, generates a random 8-byte challenge, queues both, + * and transitions [state] to [PathValidationState.Validating]. + * + * Caller is responsible for waking the writer so the challenge + * actually goes out. The writer also stamps the new DCID into + * the next outbound short-header packet — see + * [QuicConnection.activatePendingValidatedCid] for the post- + * response promotion path. + */ + fun tryStartValidation( + nowMillis: Long, + currentPtoMillis: Long, + ): PathMigrationResult { + if (state is PathValidationState.Validating) return PathMigrationResult.AlreadyInProgress + val (seq, entry) = unusedCids.entries.firstOrNull() ?: return PathMigrationResult.NoSpareCid + unusedCids.remove(seq) + val payload = + challengePayloadFactory().also { + require(it.size == 8) { "challenge payload supplier must return 8 bytes" } + } + pendingChallenges.addLast(payload) + state = + PathValidationState.Validating( + challengeData = payload, + newCidSequence = seq, + newCidBytes = entry.connectionId, + priorCidSequence = activeCidSequence, + startedAtMillis = nowMillis, + priorPtoMillis = currentPtoMillis, + ) + return PathMigrationResult.Started + } + + /** + * Process an inbound `PATH_RESPONSE` payload. Returns true when + * it matched the outstanding challenge and the migration just + * completed; false if there's no outstanding challenge or the + * payload doesn't match (an attacker echoing random bytes). + * + * Side effects on success: + * - [state] transitions to [PathValidationState.Succeeded]. + * - [activeCidSequence] is bumped to the validated sequence. + * - The prior sequence is queued for RETIRE_CONNECTION_ID. + * - Returns the new entry so the connection can swap its + * `destinationConnectionId` field. + */ + fun applyPathResponse(payload: ByteArray): ValidationOutcome { + val current = state as? PathValidationState.Validating ?: return ValidationOutcome.NotValidating + if (!payload.contentEquals(current.challengeData)) return ValidationOutcome.PayloadMismatch + val priorSeq = activeCidSequence + activeCidSequence = current.newCidSequence + if (priorSeq != current.newCidSequence) queueRetireSequence(priorSeq) + state = PathValidationState.Succeeded + successfulValidations += 1 + return ValidationOutcome.Validated( + newSequence = current.newCidSequence, + newConnectionIdBytes = current.newCidBytes, + retiredSequence = priorSeq, + ) + } + + sealed class ValidationOutcome { + object NotValidating : ValidationOutcome() + + object PayloadMismatch : ValidationOutcome() + + data class Validated( + val newSequence: Long, + val newConnectionIdBytes: ByteArray, + val retiredSequence: Long, + ) : ValidationOutcome() { + override fun equals(other: Any?): Boolean { + if (this === other) return true + if (other !is Validated) return false + return newSequence == other.newSequence && + newConnectionIdBytes.contentEquals(other.newConnectionIdBytes) && + retiredSequence == other.retiredSequence + } + + override fun hashCode(): Int { + var h = newSequence.hashCode() + h = 31 * h + newConnectionIdBytes.contentHashCode() + h = 31 * h + retiredSequence.hashCode() + return h + } + } + } + + /** + * RFC 9000 §8.2.4: abandon validation if more than 3 * PTO has + * elapsed since the challenge went out without a matching + * response. Returns the abandoned context (so the caller can + * surface a qlog event / decide whether to attempt another + * rotation), or null if no validation is in progress / the + * timer hasn't fired yet. + */ + fun checkValidationTimeout(nowMillis: Long): PathValidationState.Validating? { + val current = state as? PathValidationState.Validating ?: return null + val elapsed = nowMillis - current.startedAtMillis + val budget = (current.priorPtoMillis * VALIDATION_PTO_BUDGET_MULTIPLIER).coerceAtLeast(MIN_VALIDATION_BUDGET_MS) + if (elapsed < budget) return null + state = PathValidationState.Failed + failedValidations += 1 + return current + } + + /** + * Reset to [PathValidationState.Idle] after the caller has + * surfaced a Succeeded / Failed transition. Allows the next + * trigger to start a fresh attempt without leaking the + * previous terminal state. + */ + fun acknowledgeTerminal() { + if (state is PathValidationState.Succeeded || state is PathValidationState.Failed) { + state = PathValidationState.Idle + } + } + + companion object { + /** + * Default cap on the unused-CID pool. RFC 9000 §18.2 default + * for `active_connection_id_limit` is 2, and the spec + * lower-bounds it at 2. A peer that issues more than its + * own advertised limit is misbehaving but harmless — we + * just cap at 8 here so a buggy server can't pin memory. + */ + const val DEFAULT_MAX_UNUSED_CIDS: Int = 8 + + /** + * RFC 9000 §8.2.4: validation MUST be abandoned after + * 3 * PTO. We use 3.0 exactly; a slightly looser value + * (e.g. 4) wouldn't change correctness but might mask peer + * misbehavior in tests. + */ + const val VALIDATION_PTO_BUDGET_MULTIPLIER: Long = 3L + + /** + * Floor on the validation budget. PTO during the first RTT + * sample is ~300 ms (see [com.vitorpamplona.quic.connection.recovery.QuicLossDetection.INITIAL_RTT_MS]), + * so 3*PTO ≈ 900 ms. We add a generous floor so test paths + * with an artificially-tiny RTT still get a real chance to + * complete validation before the timer fires. + */ + const val MIN_VALIDATION_BUDGET_MS: Long = 250L + } +} diff --git a/quic/src/commonMain/kotlin/com/vitorpamplona/quic/connection/QuicConnection.kt b/quic/src/commonMain/kotlin/com/vitorpamplona/quic/connection/QuicConnection.kt index 4682d7f6a..4508ae504 100644 --- a/quic/src/commonMain/kotlin/com/vitorpamplona/quic/connection/QuicConnection.kt +++ b/quic/src/commonMain/kotlin/com/vitorpamplona/quic/connection/QuicConnection.kt @@ -505,6 +505,28 @@ class QuicConnection( */ internal val pendingPathChallengePayloads: ArrayDeque = ArrayDeque() + /** + * RFC 9000 §9 client-initiated path validation + DCID rotation. + * The pool of unused peer-issued connection IDs, the outbound + * `PATH_CHALLENGE` queue, the matching state machine, and the + * `RETIRE_CONNECTION_ID` retransmit queue all live here. The + * parser populates the pool from inbound `NEW_CONNECTION_ID`; + * the writer drains the challenge + retire queues; the driver + * triggers a fresh validation when consecutive PTO threshold is + * exceeded. + * + * Caller of any read/write must hold [streamsLock]. + */ + internal val pathValidator: PathValidator = + PathValidator( + initialActiveCidSequence = 0L, + maxUnusedCids = + config.activeConnectionIdLimit + .coerceAtMost(PathValidator.DEFAULT_MAX_UNUSED_CIDS.toLong()) + .coerceAtLeast(2L) + .toInt(), + ) + /** * RFC 9002 RTT estimator + loss-detection algorithm. Single * shared instance per connection (RTT is per-path; we model a @@ -1881,6 +1903,190 @@ class QuicConnection( } } + /** + * RFC 9000 §19.15: store a peer-issued connection ID into the + * [pathValidator] pool, enforcing `retire_prior_to` semantics + * and capping pool size. On a peer protocol violation + * ([PathValidator.RecordResult.RetirePriorToExceedsSequence] / + * [PathValidator.RecordResult.RetirePriorToRegressed] / + * [PathValidator.RecordResult.DuplicateSequenceMismatch] / + * [PathValidator.RecordResult.InvalidCidLength] / + * [PathValidator.RecordResult.InvalidStatelessResetToken]) the + * connection is closed with FRAME_ENCODING_ERROR — these are + * MUSTs in the spec. + * + * Caller must hold [streamsLock]. + */ + internal fun applyPeerNewConnectionIdLocked( + sequenceNumber: Long, + retirePriorTo: Long, + connectionId: ByteArray, + statelessResetToken: ByteArray, + ) { + when ( + pathValidator.recordPeerNewConnectionId( + sequenceNumber = sequenceNumber, + retirePriorTo = retirePriorTo, + connectionId = connectionId, + statelessResetToken = statelessResetToken, + ) + ) { + PathValidator.RecordResult.Stored -> { + Unit + } + + PathValidator.RecordResult.Duplicate -> { + Unit + } + + PathValidator.RecordResult.AlreadyRetired -> { + Unit + } + + PathValidator.RecordResult.PoolFull -> { + // Peer over-issued past its own advertised + // active_connection_id_limit. RFC 9000 §5.1.1 says + // we MAY treat this as CONNECTION_ID_LIMIT_ERROR. We + // only drop silently here — pinning memory is the + // real concern, and the cap defense already handled + // that. + } + + PathValidator.RecordResult.RetirePriorToExceedsSequence -> { + markClosedExternally( + "FRAME_ENCODING_ERROR: NEW_CONNECTION_ID retire_prior_to ($retirePriorTo) > sequence_number ($sequenceNumber)", + ) + } + + PathValidator.RecordResult.DuplicateSequenceMismatch -> { + markClosedExternally( + "PROTOCOL_VIOLATION: NEW_CONNECTION_ID seq=$sequenceNumber re-issued with different bytes", + ) + } + + PathValidator.RecordResult.InvalidCidLength -> { + markClosedExternally( + "FRAME_ENCODING_ERROR: NEW_CONNECTION_ID has invalid cid length", + ) + } + + PathValidator.RecordResult.InvalidStatelessResetToken -> { + markClosedExternally( + "FRAME_ENCODING_ERROR: NEW_CONNECTION_ID has invalid stateless reset token length", + ) + } + } + } + + /** + * Process an inbound `PATH_RESPONSE` (RFC 9000 §19.18). If the + * payload matches our outstanding `PATH_CHALLENGE`, the writer + * gets a new `destinationConnectionId` to stamp into outbound + * short-headers and a `RETIRE_CONNECTION_ID` is queued for the + * old sequence number. Mismatches and unsolicited responses are + * silently dropped — RFC 9000 §8.2.2 says "an endpoint that + * receives a PATH_RESPONSE with a different payload than what + * it sent in the PATH_CHALLENGE on the path" is allowed to + * ignore. + * + * Caller must hold [streamsLock]. + */ + internal fun applyPeerPathResponseLocked(payload: ByteArray) { + when (val outcome = pathValidator.applyPathResponse(payload)) { + PathValidator.ValidationOutcome.NotValidating, + PathValidator.ValidationOutcome.PayloadMismatch, + -> { + Unit + } + + is PathValidator.ValidationOutcome.Validated -> { + destinationConnectionId = ConnectionId(outcome.newConnectionIdBytes) + qlogObserver.onPathValidationSucceeded(outcome.newSequence) + qlogObserver.onConnectionIdActivated("peer", outcome.newSequence, outcome.newConnectionIdBytes) + pathValidator.acknowledgeTerminal() + } + } + } + + /** + * RFC 9000 §19.16: peer asked us to retire one of our own + * source CIDs. Today we only ever issue one source CID at + * connection start (sequence 0), so any RETIRE_CONNECTION_ID + * the peer sends about a sequence > 0 references a CID we + * never issued — RFC 9000 says this is PROTOCOL_VIOLATION. + * Sequence 0 retire is also a violation per §19.16: "Receipt + * of a RETIRE_CONNECTION_ID frame containing a sequence number + * greater than any previously sent to the peer MUST be + * treated as a connection error of type PROTOCOL_VIOLATION." + * + * Caller must hold [streamsLock]. + */ + internal fun applyPeerRetireConnectionIdLocked(sequenceNumber: Long) { + // We only issue sequence 0 today; any retire request beyond + // that references a CID we never advertised. This is a + // peer-side violation per §19.16. + if (sequenceNumber > 0L) { + markClosedExternally( + "PROTOCOL_VIOLATION: RETIRE_CONNECTION_ID seq=$sequenceNumber > any we issued (0)", + ) + } + // Sequence 0 retire today means "stop using your initial + // SCID." We don't model client-issued CID rotation, so the + // request is silently honored at the protocol level (we + // never re-use the SCID in any new long-header packet — + // long headers are gone after handshake completion). + } + + /** + * Try to start client-initiated path validation + DCID rotation + * (RFC 9000 §9). Picks the lowest-sequence unused CID from the + * pool, queues a fresh PATH_CHALLENGE, and records the + * outstanding state. Returns the result so the caller can + * decide whether to retry later (e.g. on the next PTO if the + * pool was empty when this fired). + * + * Caller must hold [streamsLock]. + */ + internal fun triggerPathMigrationLocked( + nowMillis: Long, + currentPtoMillis: Long, + ): PathMigrationResult { + val result = pathValidator.tryStartValidation(nowMillis, currentPtoMillis) + if (result == PathMigrationResult.Started) { + val validating = pathValidator.state as PathValidationState.Validating + qlogObserver.onPathValidationStarted(validating.newCidSequence) + } + return result + } + + /** + * Public counterpart to [triggerPathMigrationLocked] — acquires + * the lock internally. Suitable for application / driver code + * that wants to force a rotation without coupling to the + * locking discipline. + */ + suspend fun triggerPathMigration( + nowMillis: Long = nowMillis(), + currentPtoMillis: Long = lossDetection.ptoBaseMs(peerTransportParameters?.maxAckDelay ?: 0L), + ): PathMigrationResult = streamsLock.withLock { triggerPathMigrationLocked(nowMillis, currentPtoMillis) } + + /** + * Check whether an outstanding PATH_CHALLENGE has exceeded the + * 3 * PTO budget (RFC 9000 §8.2.4). Called from the driver's + * PTO timer path. Returns true when validation just timed out + * — caller may surface that to qlog and decide to retry with + * another CID, or close the connection if the budget for + * retries is exhausted. + * + * Caller must hold [streamsLock]. + */ + internal fun checkPathValidationTimeoutLocked(nowMillis: Long): Boolean { + val abandoned = pathValidator.checkValidationTimeout(nowMillis) ?: return false + qlogObserver.onPathValidationFailed(abandoned.newCidSequence) + pathValidator.acknowledgeTerminal() + return true + } + /** * Queue a PATH_RESPONSE for the given [challengeData]. Called by * the parser when a PATH_CHALLENGE arrives. Idempotent on @@ -1956,12 +2162,19 @@ class QuicConnection( is com.vitorpamplona.quic.connection.recovery.RecoveryToken.MaxData, is com.vitorpamplona.quic.connection.recovery.RecoveryToken.MaxStreamData, is com.vitorpamplona.quic.connection.recovery.RecoveryToken.NewConnectionId, + is com.vitorpamplona.quic.connection.recovery.RecoveryToken.PathChallenge, + is com.vitorpamplona.quic.connection.recovery.RecoveryToken.RetireConnectionId, -> { - // Flow-control extensions and NEW_CONNECTION_ID - // have no per-buffer state to release on ACK; the + // Flow-control extensions, NEW_CONNECTION_ID, + // PATH_CHALLENGE, and RETIRE_CONNECTION_ID have + // no per-buffer state to release on ACK; the // pending* maps are populated only on loss, so an // ACK for a frame that never lost is naturally - // absent. + // absent. PATH_CHALLENGE specifically: a peer ACK + // means "we received the challenge frame" — but + // that's NOT path validation success. Validation + // requires the matching PATH_RESPONSE + // (handled in [applyPeerPathResponseLocked]). } is com.vitorpamplona.quic.connection.recovery.RecoveryToken.ResetStream -> { @@ -2086,6 +2299,33 @@ class QuicConnection( is com.vitorpamplona.quic.connection.recovery.RecoveryToken.NewConnectionId -> { pendingNewConnectionId[token.sequenceNumber] = token } + + is com.vitorpamplona.quic.connection.recovery.RecoveryToken.PathChallenge -> { + // RFC 9000 §8.2.1: if a PATH_CHALLENGE is lost + // (carrying packet declared lost without an + // ACK), the spec permits but does not require + // retransmission — validation simply runs against + // the 3 * PTO budget. We re-queue iff the + // outstanding validation still matches this + // payload (state hasn't progressed to Succeeded / + // Failed / a different challenge). The writer + // will pick it up on the next drain and emit a + // fresh PATH_CHALLENGE with the SAME 8 bytes — + // the peer must echo whichever it sees first. + val s = pathValidator.state + if (s is PathValidationState.Validating && s.challengeData.contentEquals(token.data)) { + if (!pathValidator.pendingChallenges.any { it.contentEquals(token.data) }) { + pathValidator.pendingChallenges.addLast(token.data) + } + } + } + + is com.vitorpamplona.quic.connection.recovery.RecoveryToken.RetireConnectionId -> { + // RFC 9000 §13.3: RETIRE_CONNECTION_ID is + // reliable. Re-queue on loss — idempotent on + // duplicate (queueRetireSequence dedupes). + pathValidator.queueRetireSequence(token.sequenceNumber) + } } } } diff --git a/quic/src/commonMain/kotlin/com/vitorpamplona/quic/connection/QuicConnectionDriver.kt b/quic/src/commonMain/kotlin/com/vitorpamplona/quic/connection/QuicConnectionDriver.kt index 75d8b2e25..0e2dcd184 100644 --- a/quic/src/commonMain/kotlin/com/vitorpamplona/quic/connection/QuicConnectionDriver.kt +++ b/quic/src/commonMain/kotlin/com/vitorpamplona/quic/connection/QuicConnectionDriver.kt @@ -362,7 +362,56 @@ internal suspend fun handlePtoFired(conn: QuicConnection) { conn.streamsLock.withLock { conn.requeueAllInflightStreamData() conn.requeueAllInflightCrypto(EncryptionLevel.APPLICATION) + // RFC 9000 §9 client-initiated path validation. After + // [PATH_PROBE_PTO_THRESHOLD] consecutive PTOs without + // any inbound ACK we suspect the path is dead (NAT + // rebind, route flap, dead peer). If the peer has + // issued spare CIDs via NEW_CONNECTION_ID, rotate to + // one and emit a PATH_CHALLENGE on the new DCID. The + // validator only triggers on the FIRST crossing of + // the threshold per validation cycle — the + // [PathValidator] internally rejects re-entry while + // [PathValidationState.Validating] holds. + // + // Also check the §8.2.4 budget on any in-flight + // validation: 3*PTO since the challenge went out + // without a matching response means the new path is + // also dead — abandon and let the next PTO try with + // another CID (or surface the failure to the higher + // layer). + val nowMillis = + kotlin.time.Clock.System + .now() + .toEpochMilliseconds() + val maxAckDelayMs = + if (conn.application.sendProtection != null) { + conn.peerTransportParameters?.maxAckDelay ?: 0L + } else { + 0L + } + val ptoBaseMs = conn.lossDetection.ptoBaseMs(maxAckDelayMs).coerceAtLeast(1L) + conn.checkPathValidationTimeoutLocked(nowMillis) + if (conn.consecutivePtoCount >= PATH_PROBE_PTO_THRESHOLD) { + conn.triggerPathMigrationLocked(nowMillis = nowMillis, currentPtoMillis = ptoBaseMs) + } } } conn.consecutivePtoCount = (conn.consecutivePtoCount + 1).coerceAtMost(6) } + +/** + * RFC 9000 §9 / §10.1.2 — number of consecutive PTOs we tolerate on + * the active path before assuming it's dead and probing a new one. + * The QUIC RFC doesn't pin a specific value (the spec only says "an + * endpoint that has previously discovered a particular path + * works"); 2 matches Firefox neqo's `PATH_PROBE_PTO_THRESHOLD` + * default and Chrome's behavior. Picking 1 is too aggressive + * (single dropped packet trips a rotation); 4+ is too late (user + * notices the silence). + * + * Once the threshold is crossed, [handlePtoFired] calls into + * [QuicConnection.triggerPathMigrationLocked] which is itself + * idempotent — a second crossing while validation is already in + * flight is a no-op. + */ +internal const val PATH_PROBE_PTO_THRESHOLD: Int = 2 diff --git a/quic/src/commonMain/kotlin/com/vitorpamplona/quic/connection/QuicConnectionParser.kt b/quic/src/commonMain/kotlin/com/vitorpamplona/quic/connection/QuicConnectionParser.kt index 64dfb96d3..f4927721c 100644 --- a/quic/src/commonMain/kotlin/com/vitorpamplona/quic/connection/QuicConnectionParser.kt +++ b/quic/src/commonMain/kotlin/com/vitorpamplona/quic/connection/QuicConnectionParser.kt @@ -36,6 +36,7 @@ import com.vitorpamplona.quic.frame.PathChallengeFrame import com.vitorpamplona.quic.frame.PathResponseFrame import com.vitorpamplona.quic.frame.PingFrame import com.vitorpamplona.quic.frame.ResetStreamFrame +import com.vitorpamplona.quic.frame.RetireConnectionIdFrame import com.vitorpamplona.quic.frame.StopSendingFrame import com.vitorpamplona.quic.frame.StreamFrame import com.vitorpamplona.quic.frame.decodeFrames @@ -716,10 +717,23 @@ private fun dispatchFrames( } is NewConnectionIdFrame -> { - // RFC 9000 §13.2.1: NEW_CONNECTION_ID is ack-eliciting. We - // don't support migration but still need to ACK to keep - // peer's loss-recovery happy. + // RFC 9000 §13.2.1: NEW_CONNECTION_ID is ack-eliciting. + // §19.15 + §5.1.2: store the offered CID + stateless + // reset token in the [PathValidator] pool so the + // client-initiated migration path (RFC 9000 §9) can + // pick one when the active path stops receiving + // ACKs. Also enforces `retire_prior_to` semantics + // and queues RETIRE_CONNECTION_ID for any cached + // entry whose sequence number falls below the new + // watermark. Peer protocol violations close the + // connection upstream. ackEliciting = true + conn.applyPeerNewConnectionIdLocked( + sequenceNumber = frame.sequenceNumber, + retirePriorTo = frame.retirePriorTo, + connectionId = frame.connectionId, + statelessResetToken = frame.statelessResetToken, + ) } is PathChallengeFrame -> { @@ -746,13 +760,26 @@ private fun dispatchFrames( is PathResponseFrame -> { // RFC 9000 §13.2.1: PATH_RESPONSE is ack-eliciting. - // We don't yet issue PATH_CHALLENGE ourselves (that's - // the client-initiated migration path, out of scope - // for the first-pass landing here), so any PATH_RESPONSE - // we receive is necessarily for a challenge we never - // sent — drop it after marking ack-eliciting so the - // outbound ACK still goes out. + // §8.2.2: if the payload matches our outstanding + // PATH_CHALLENGE, validation succeeded — the writer + // is told to start stamping the new DCID and a + // RETIRE_CONNECTION_ID for the old sequence number + // is queued (RFC 9000 §9.5). Mismatched payloads + // (no outstanding challenge OR an attacker echoing + // random bytes) are silently ignored after the ACK. ackEliciting = true + conn.applyPeerPathResponseLocked(frame.data) + } + + is RetireConnectionIdFrame -> { + // RFC 9000 §13.2.1: RETIRE_CONNECTION_ID is + // ack-eliciting. §19.16: peer is asking us to + // retire one of OUR source CIDs. We only ever + // issued sequence 0 (the SCID we picked at + // connection start); any seq > 0 is a protocol + // violation per §19.16, closes the connection. + ackEliciting = true + conn.applyPeerRetireConnectionIdLocked(frame.sequenceNumber) } is ConnectionCloseFrame -> { diff --git a/quic/src/commonMain/kotlin/com/vitorpamplona/quic/connection/QuicConnectionWriter.kt b/quic/src/commonMain/kotlin/com/vitorpamplona/quic/connection/QuicConnectionWriter.kt index 3172eb52b..c5ebaf9fd 100644 --- a/quic/src/commonMain/kotlin/com/vitorpamplona/quic/connection/QuicConnectionWriter.kt +++ b/quic/src/commonMain/kotlin/com/vitorpamplona/quic/connection/QuicConnectionWriter.kt @@ -32,9 +32,11 @@ import com.vitorpamplona.quic.frame.MaxDataFrame import com.vitorpamplona.quic.frame.MaxStreamDataFrame import com.vitorpamplona.quic.frame.MaxStreamsFrame import com.vitorpamplona.quic.frame.NewConnectionIdFrame +import com.vitorpamplona.quic.frame.PathChallengeFrame import com.vitorpamplona.quic.frame.PathResponseFrame import com.vitorpamplona.quic.frame.PingFrame import com.vitorpamplona.quic.frame.ResetStreamFrame +import com.vitorpamplona.quic.frame.RetireConnectionIdFrame import com.vitorpamplona.quic.frame.StopSendingFrame import com.vitorpamplona.quic.frame.StreamFrame import com.vitorpamplona.quic.frame.encodeFrames @@ -962,6 +964,26 @@ private fun appendFlowControlUpdates( frames += PathResponseFrame(data) } + // RFC 9000 §9 client-initiated path validation. Drain any + // PATH_CHALLENGE the [PathValidator] has queued — including + // re-queued ones from the loss dispatcher. Each emission + // records a [RecoveryToken.PathChallenge] so loss recovery + // can decide to retransmit (the validator's own 3 * PTO + // budget is the ultimate timeout per RFC 9000 §8.2.4). + while (conn.pathValidator.pendingChallenges.isNotEmpty()) { + val payload = conn.pathValidator.pendingChallenges.removeFirst() + frames += PathChallengeFrame(payload) + tokens += RecoveryToken.PathChallenge(payload) + } + + // RFC 9000 §19.16 RETIRE_CONNECTION_ID. Drain pending retires; + // the dispatcher re-queues on loss until the peer ACKs. + while (conn.pathValidator.pendingRetireSequences.isNotEmpty()) { + val seq = conn.pathValidator.pendingRetireSequences.removeFirst() + frames += RetireConnectionIdFrame(seq) + tokens += RecoveryToken.RetireConnectionId(seq) + } + // NEW_CONNECTION_ID retransmits. No application path emits these // initially today (connection-ID rotation isn't wired); the map // is populated only by the loss dispatcher, so this branch only diff --git a/quic/src/commonMain/kotlin/com/vitorpamplona/quic/connection/recovery/RecoveryToken.kt b/quic/src/commonMain/kotlin/com/vitorpamplona/quic/connection/recovery/RecoveryToken.kt index 8ed60acff..320015df7 100644 --- a/quic/src/commonMain/kotlin/com/vitorpamplona/quic/connection/recovery/RecoveryToken.kt +++ b/quic/src/commonMain/kotlin/com/vitorpamplona/quic/connection/recovery/RecoveryToken.kt @@ -178,10 +178,45 @@ sealed class RecoveryToken { val errorCode: Long, ) : RecoveryToken() + /** + * `PATH_CHALLENGE` frame we emitted as part of client-initiated + * path validation (RFC 9000 §8.2 / §9). Carries the same 8-byte + * payload the writer put on the wire so the loss dispatcher can + * decide whether the validation attempt should be considered + * failed or retransmitted. + * + * RFC 9000 §13.2.1 lists PATH_CHALLENGE as ack-eliciting; §8.2.4 + * says path validation MUST NOT exceed `3 * PTO` of waiting before + * the path is declared failed. The connection-side dispatcher + * checks the [com.vitorpamplona.quic.connection.PathValidator] + * state — if validation has already moved on (succeeded, failed, + * or abandoned), the lost token is dropped silently. + */ + data class PathChallenge( + val data: ByteArray, + ) : RecoveryToken() { + override fun equals(other: Any?): Boolean { + if (this === other) return true + if (other !is PathChallenge) return false + return data.contentEquals(other.data) + } + + override fun hashCode(): Int = data.contentHashCode() + } + + /** + * `RETIRE_CONNECTION_ID` frame we emitted (RFC 9000 §19.16). + * Reliable per §13.3 — peer must learn that we've stopped using + * the connection ID, otherwise it keeps the routing entry forever. + */ + data class RetireConnectionId( + val sequenceNumber: Long, + ) : RecoveryToken() + /** * `NEW_CONNECTION_ID` frame we sent (RFC 9000 §19.15). Reliable - * per §13.3. Same scaffolding-only status — connection-ID - * rotation isn't wired today. + * per §13.3. Same scaffolding-only status — client-issued + * connection IDs aren't yet emitted by `:quic`. */ data class NewConnectionId( val sequenceNumber: Long, diff --git a/quic/src/commonMain/kotlin/com/vitorpamplona/quic/frame/Frame.kt b/quic/src/commonMain/kotlin/com/vitorpamplona/quic/frame/Frame.kt index 965b0e529..a3ae7f8d8 100644 --- a/quic/src/commonMain/kotlin/com/vitorpamplona/quic/frame/Frame.kt +++ b/quic/src/commonMain/kotlin/com/vitorpamplona/quic/frame/Frame.kt @@ -313,6 +313,25 @@ class NewConnectionIdFrame( } } +/** + * RFC 9000 §19.16 — RETIRE_CONNECTION_ID frame. Tells the peer to + * stop using one of the connection IDs it issued. Carries the + * sequence number from the corresponding [NewConnectionIdFrame]. + * + * The client side emits this when it rotates to a new DCID + * post-handshake (see RFC 9000 §9 client-initiated migration): once + * the new path is validated, the previously-active CID is retired + * so the server can free the routing entry. + */ +class RetireConnectionIdFrame( + val sequenceNumber: Long, +) : Frame() { + override fun encode(out: QuicWriter) { + out.writeByte(FrameType.RETIRE_CONNECTION_ID.toInt()) + out.writeVarint(sequenceNumber) + } +} + /** * RFC 9000 §19.17 — PATH_CHALLENGE frame, used for path validation * (§8.2). The 8-byte [data] payload is opaque random bytes the @@ -489,7 +508,7 @@ fun decodeFrames(data: ByteArray): List { } type == FrameType.RETIRE_CONNECTION_ID -> { - r.readVarint() + out += RetireConnectionIdFrame(r.readVarint()) } type == FrameType.PATH_CHALLENGE -> { diff --git a/quic/src/commonMain/kotlin/com/vitorpamplona/quic/observability/QlogObserver.kt b/quic/src/commonMain/kotlin/com/vitorpamplona/quic/observability/QlogObserver.kt index e977c4f11..99754800e 100644 --- a/quic/src/commonMain/kotlin/com/vitorpamplona/quic/observability/QlogObserver.kt +++ b/quic/src/commonMain/kotlin/com/vitorpamplona/quic/observability/QlogObserver.kt @@ -159,6 +159,51 @@ interface QlogObserver { otherVersionsOffered: List, ) + /** + * RFC 9000 §9 client-initiated path validation just sent its + * first `PATH_CHALLENGE` for [newCidSequence] (the sequence + * number of the peer-issued connection ID we'll use on the new + * path). qvis renders this as a path-validation milestone. + */ + fun onPathValidationStarted(newCidSequence: Long) = Unit + + /** + * The peer echoed our PATH_CHALLENGE payload — validation + * succeeded. The writer has now switched to the new DCID and + * a `RETIRE_CONNECTION_ID` for the prior sequence is queued. + */ + fun onPathValidationSucceeded(newCidSequence: Long) = Unit + + /** + * Path validation was abandoned because more than `3 * PTO` + * elapsed without a matching `PATH_RESPONSE` (RFC 9000 + * §8.2.4). Caller may retry with another CID or surface the + * failure as a connection close. + */ + fun onPathValidationFailed(newCidSequence: Long) = Unit + + /** + * A connection ID belonging to [keyType] (`"local"` for our + * source CID, `"peer"` for the peer-issued destination CID) + * just became active — i.e. the writer / parser will start + * using it for outbound / inbound packets respectively. + */ + fun onConnectionIdActivated( + keyType: String, + sequenceNumber: Long, + connectionId: ByteArray, + ) = Unit + + /** + * A connection ID was retired (RFC 9000 §19.16 — either we + * sent RETIRE_CONNECTION_ID for a peer-issued CID, or the + * peer asked us to retire one of ours). + */ + fun onConnectionIdRetired( + keyType: String, + sequenceNumber: Long, + ) = Unit + /** * No-op observer. Default for production callers — every method * is an empty body that the JIT inlines. No allocation, no I/O. diff --git a/quic/src/commonTest/kotlin/com/vitorpamplona/quic/connection/ClientPathMigrationTest.kt b/quic/src/commonTest/kotlin/com/vitorpamplona/quic/connection/ClientPathMigrationTest.kt new file mode 100644 index 000000000..160bc0b77 --- /dev/null +++ b/quic/src/commonTest/kotlin/com/vitorpamplona/quic/connection/ClientPathMigrationTest.kt @@ -0,0 +1,260 @@ +/* + * 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.connection + +import com.vitorpamplona.quic.frame.NewConnectionIdFrame +import com.vitorpamplona.quic.frame.PathChallengeFrame +import com.vitorpamplona.quic.frame.PathResponseFrame +import com.vitorpamplona.quic.frame.RetireConnectionIdFrame +import com.vitorpamplona.quic.frame.decodeFrames +import com.vitorpamplona.quic.frame.encodeFrames +import kotlinx.coroutines.runBlocking +import kotlin.test.Test +import kotlin.test.assertContentEquals +import kotlin.test.assertEquals +import kotlin.test.assertTrue + +/** + * Integration tests for client-initiated path validation + DCID + * rotation (RFC 9000 §9). Drives the real client through the + * [InMemoryQuicPipe] harness so the assertions cover the parser / + * writer / [PathValidator] wiring end-to-end: + * + * 1. Server sends NEW_CONNECTION_ID — pool fills. + * 2. Application calls [QuicConnection.triggerPathMigration]. + * 3. Client emits PATH_CHALLENGE (drained from + * [PathValidator.pendingChallenges]). + * 4. Server replies with PATH_RESPONSE carrying the same 8 bytes. + * 5. Client switches `destinationConnectionId` and queues a + * RETIRE_CONNECTION_ID for the old sequence. + * 6. Next outbound application packet uses the NEW DCID and + * contains the RETIRE_CONNECTION_ID. + * + * Out of scope here: + * - PTO-driven rotation: covered by the unit-level + * [PathValidatorTest.validationTimeoutAfter3PtoTransitionsToFailed] + * plus the driver-level threshold constant. Driving a real PTO + * through the in-memory pipe would require simulating timer + * advancement, which the pipe doesn't model today. + */ +class ClientPathMigrationTest { + @Test + fun retireConnectionIdFrameRoundTripsThroughCodec() { + val encoded = encodeFrames(listOf(RetireConnectionIdFrame(sequenceNumber = 7L))) + val decoded = decodeFrames(encoded) + assertEquals(1, decoded.size) + val frame = decoded.first() as RetireConnectionIdFrame + assertEquals(7L, frame.sequenceNumber) + } + + @Test + fun newConnectionIdFromServerPopulatesPathValidatorPool() = + runBlocking { + val (client, pipe) = newConnectedClient() + val newCid = ByteArray(8) { (0xC0 or it).toByte() } + val token = ByteArray(16) { 0x42 } + val packet = + pipe.buildServerApplicationDatagram( + listOf( + NewConnectionIdFrame( + sequenceNumber = 1L, + retirePriorTo = 0L, + connectionId = newCid, + statelessResetToken = token, + ), + ), + )!! + feedDatagram(client, packet, nowMillis = 0L) + + assertEquals(1, client.pathValidator.unusedCount()) + assertEquals(listOf(1L), client.pathValidator.unusedSequences()) + } + + @Test + fun newConnectionIdWithRetirePriorToGreaterThanSequenceClosesConnection() = + runBlocking { + val (client, pipe) = newConnectedClient() + val packet = + pipe.buildServerApplicationDatagram( + listOf( + NewConnectionIdFrame( + sequenceNumber = 1L, + retirePriorTo = 5L, + connectionId = ByteArray(8) { 0x01 }, + statelessResetToken = ByteArray(16), + ), + ), + )!! + feedDatagram(client, packet, nowMillis = 0L) + // FRAME_ENCODING_ERROR per RFC 9000 §19.15 — closes the connection. + assertTrue( + client.status == QuicConnection.Status.CLOSING || client.status == QuicConnection.Status.CLOSED, + "got ${client.status}; expected CLOSING/CLOSED", + ) + } + + @Test + fun triggerPathMigrationWithoutSpareCidIsNoOp() = + runBlocking { + val (client, _) = newConnectedClient() + val priorDcid = client.destinationConnectionId + val result = client.triggerPathMigration(nowMillis = 0L, currentPtoMillis = 100L) + assertEquals(PathMigrationResult.NoSpareCid, result) + assertEquals(priorDcid, client.destinationConnectionId, "DCID must not change without a spare CID") + assertTrue(client.pathValidator.state is PathValidationState.Idle) + } + + @Test + fun fullMigrationRoundTripSwitchesDcidAndQueuesRetire() = + runBlocking { + val (client, pipe) = newConnectedClient() + val originalDcid = client.destinationConnectionId + val newCidBytes = ByteArray(8) { (0xA0 or it).toByte() } + val token = ByteArray(16) { 0x77 } + + // Step 1 — server offers a fresh CID via NEW_CONNECTION_ID. + val offer = + pipe.buildServerApplicationDatagram( + listOf( + NewConnectionIdFrame( + sequenceNumber = 1L, + retirePriorTo = 0L, + connectionId = newCidBytes, + statelessResetToken = token, + ), + ), + )!! + feedDatagram(client, offer, nowMillis = 0L) + assertEquals(1, client.pathValidator.unusedCount()) + + // Drain the client's ACK for that packet so subsequent assertions about + // outbound contents aren't polluted by leftover ACK frames. + drainOutbound(client, nowMillis = 0L) + + // Step 2 — application triggers migration. The validator picks + // sequence 1 and queues a PATH_CHALLENGE. + val triggered = client.triggerPathMigration(nowMillis = 100L, currentPtoMillis = 1_000L) + assertEquals(PathMigrationResult.Started, triggered) + assertTrue(client.pathValidator.state is PathValidationState.Validating) + + // Step 3 — the next outbound application packet must carry a + // PATH_CHALLENGE. + val challengeOut = drainOutbound(client, nowMillis = 100L) + assertTrue(challengeOut != null, "client must emit a packet carrying PATH_CHALLENGE") + val outboundFrames = pipe.decryptClientApplicationFrames(challengeOut) + assertTrue(outboundFrames != null, "client outbound must decrypt with server keys") + val challenge = outboundFrames.firstOrNull { it is PathChallengeFrame } as? PathChallengeFrame + assertTrue(challenge != null, "outbound must contain PATH_CHALLENGE — got ${outboundFrames.map { it::class.simpleName }}") + + // Step 4 — server echoes the payload in PATH_RESPONSE. + val response = + pipe.buildServerApplicationDatagram( + listOf(PathResponseFrame(challenge.data.copyOf())), + )!! + feedDatagram(client, response, nowMillis = 200L) + + // Step 5 — DCID must have rotated to the new bytes; the prior + // sequence (0) must be queued for RETIRE_CONNECTION_ID. + assertContentEquals(newCidBytes, client.destinationConnectionId.bytes, "DCID must rotate to new bytes") + assertEquals(1L, client.pathValidator.activeCidSequence) + assertTrue( + client.pathValidator.pendingRetireSequences.contains(0L), + "old sequence number must be queued for retire", + ) + + // Step 6 — the next outbound packet should carry the + // RETIRE_CONNECTION_ID for the old sequence. + val retireOut = drainOutbound(client, nowMillis = 200L) + assertTrue(retireOut != null, "client must emit a packet after PATH_RESPONSE") + val frames = pipe.decryptClientApplicationFrames(retireOut) + assertTrue(frames != null, "outbound after rotation must decrypt") + val retire = frames.firstOrNull { it is RetireConnectionIdFrame } as? RetireConnectionIdFrame + assertTrue(retire != null, "outbound must contain RETIRE_CONNECTION_ID — got ${frames.map { it::class.simpleName }}") + assertEquals(0L, retire.sequenceNumber) + + // Sanity: connection still alive after the rotation. + assertEquals(QuicConnection.Status.CONNECTED, client.status) + // Sanity: original DCID changed. + assertTrue(originalDcid != client.destinationConnectionId) + } + + @Test + fun pathResponseWithMismatchingPayloadDoesNotRotateDcid() = + runBlocking { + val (client, pipe) = newConnectedClient() + val originalDcid = client.destinationConnectionId + + // Seed a spare CID and trigger migration. + val offer = + pipe.buildServerApplicationDatagram( + listOf( + NewConnectionIdFrame( + sequenceNumber = 1L, + retirePriorTo = 0L, + connectionId = ByteArray(8) { 0x01 }, + statelessResetToken = ByteArray(16) { 0x10 }, + ), + ), + )!! + feedDatagram(client, offer, nowMillis = 0L) + client.triggerPathMigration(nowMillis = 100L, currentPtoMillis = 1_000L) + drainOutbound(client, nowMillis = 100L) // emit challenge, drop content + + // Server replies with WRONG payload — must not rotate. + val bogusResponse = + pipe.buildServerApplicationDatagram( + listOf(PathResponseFrame(ByteArray(8) { 0xFF.toByte() })), + )!! + feedDatagram(client, bogusResponse, nowMillis = 200L) + + assertEquals(originalDcid, client.destinationConnectionId, "mismatched response must NOT rotate DCID") + assertTrue(client.pathValidator.state is PathValidationState.Validating, "still validating") + } + + @Test + fun unsolicitedPathResponseIsSilentlyDropped() = + runBlocking { + // Defence-in-depth: a peer (or attacker) sending PATH_RESPONSE + // without a preceding PATH_CHALLENGE from us must NOT rotate the + // DCID, must NOT crash, and the connection must stay alive. + val (client, pipe) = newConnectedClient() + val originalDcid = client.destinationConnectionId + val packet = pipe.buildServerApplicationDatagram(listOf(PathResponseFrame(ByteArray(8))))!! + feedDatagram(client, packet, nowMillis = 0L) + assertEquals(originalDcid, client.destinationConnectionId) + assertEquals(QuicConnection.Status.CONNECTED, client.status) + assertTrue(client.pathValidator.state is PathValidationState.Idle) + } + + @Test + fun retireConnectionIdFromServerForUnissuedSequenceIsProtocolViolation() = + runBlocking { + val (client, pipe) = newConnectedClient() + // We only ever issue sequence 0; any seq > 0 is a protocol + // violation per RFC 9000 §19.16. + val packet = pipe.buildServerApplicationDatagram(listOf(RetireConnectionIdFrame(sequenceNumber = 7L)))!! + feedDatagram(client, packet, nowMillis = 0L) + assertTrue( + client.status == QuicConnection.Status.CLOSING || client.status == QuicConnection.Status.CLOSED, + "got ${client.status}", + ) + } +} diff --git a/quic/src/commonTest/kotlin/com/vitorpamplona/quic/connection/PathValidatorTest.kt b/quic/src/commonTest/kotlin/com/vitorpamplona/quic/connection/PathValidatorTest.kt new file mode 100644 index 000000000..d379b26e8 --- /dev/null +++ b/quic/src/commonTest/kotlin/com/vitorpamplona/quic/connection/PathValidatorTest.kt @@ -0,0 +1,237 @@ +/* + * 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.connection + +import kotlin.test.Test +import kotlin.test.assertContentEquals +import kotlin.test.assertEquals +import kotlin.test.assertNull +import kotlin.test.assertTrue + +/** + * Unit tests for the in-memory [PathValidator] state machine — the + * piece that implements RFC 9000 §5.1 + §8.2 + §9 client-initiated + * path validation. Driven without a real socket so the assertions + * are about the state-machine contract rather than the on-wire + * encoding (see [PathValidationTest] for the integration shape). + */ +class PathValidatorTest { + @Test + fun newConnectionIdStoresInPool() { + val v = PathValidator() + val cid = byteArrayOf(0x01, 0x02, 0x03, 0x04, 0x05, 0x06, 0x07, 0x08) + val tok = ByteArray(16) { 0xAA.toByte() } + val r = v.recordPeerNewConnectionId(sequenceNumber = 1L, retirePriorTo = 0L, connectionId = cid, statelessResetToken = tok) + assertEquals(PathValidator.RecordResult.Stored, r) + assertEquals(1, v.unusedCount()) + assertEquals(listOf(1L), v.unusedSequences()) + } + + @Test + fun duplicateOfferWithSameBytesIsIdempotent() { + val v = PathValidator() + val cid = ByteArray(8) { 0x42 } + val tok = ByteArray(16) { 0x99.toByte() } + v.recordPeerNewConnectionId(1L, 0L, cid, tok) + val r = v.recordPeerNewConnectionId(1L, 0L, cid, tok) + assertEquals(PathValidator.RecordResult.Duplicate, r) + assertEquals(1, v.unusedCount()) + } + + @Test + fun duplicateOfferWithDifferentBytesIsProtocolViolation() { + val v = PathValidator() + val seq = 2L + v.recordPeerNewConnectionId(seq, 0L, ByteArray(8) { 0x11 }, ByteArray(16) { 0x22 }) + val r = + v.recordPeerNewConnectionId( + sequenceNumber = seq, + retirePriorTo = 0L, + connectionId = ByteArray(8) { 0x33 }, + statelessResetToken = ByteArray(16) { 0x44 }, + ) + assertEquals(PathValidator.RecordResult.DuplicateSequenceMismatch, r) + } + + @Test + fun retirePriorToForcesEarlierEntriesIntoRetireQueue() { + val v = PathValidator() + v.recordPeerNewConnectionId(1L, 0L, ByteArray(8) { 0x01 }, ByteArray(16) { 0x10 }) + v.recordPeerNewConnectionId(2L, 0L, ByteArray(8) { 0x02 }, ByteArray(16) { 0x20 }) + // New offer with retirePriorTo = 2 forces sequences 0 and 1 into retirement. + val r = v.recordPeerNewConnectionId(3L, retirePriorTo = 2L, ByteArray(8) { 0x03 }, ByteArray(16) { 0x30 }) + assertEquals(PathValidator.RecordResult.Stored, r) + assertEquals(2L, v.retirePriorToWatermark) + // Sequence 1 was in the pool — retired. Sequence 0 was the active CID; + // PathValidator queues it implicitly via watermark advancement, but + // only sequences ≥ retirePriorToWatermark (i.e. ≥ 2) survive in + // unusedCids. + assertEquals(setOf(2L, 3L), v.unusedSequences().toSet()) + assertTrue(v.pendingRetireSequences.contains(1L), "seq 1 must be queued for retire") + } + + @Test + fun retirePriorToRegressionIsClampedNotRejected() { + // RFC 9000 §19.15: a smaller `retire_prior_to` than what we + // previously saw "MUST be treated as the largest one it has + // seen" — i.e. clamp, don't error. Reordered NEW_CONNECTION_ID + // arrivals are common on the wire and aren't peer protocol + // violations. + val v = PathValidator() + v.recordPeerNewConnectionId(1L, 1L, ByteArray(8) { 0x01 }, ByteArray(16) { 0x10 }) + val r = v.recordPeerNewConnectionId(2L, 0L, ByteArray(8) { 0x02 }, ByteArray(16) { 0x20 }) + assertEquals(PathValidator.RecordResult.Stored, r) + assertEquals(1L, v.retirePriorToWatermark, "watermark must NOT decrease on a regressed retire_prior_to") + } + + @Test + fun retirePriorToGreaterThanSequenceIsRejected() { + val v = PathValidator() + val r = v.recordPeerNewConnectionId(sequenceNumber = 1L, retirePriorTo = 5L, ByteArray(8), ByteArray(16)) + assertEquals(PathValidator.RecordResult.RetirePriorToExceedsSequence, r) + } + + @Test + fun poolFillsUpAndRejectsExcess() { + val v = PathValidator(maxUnusedCids = 3) + for (i in 1..3) { + assertEquals( + PathValidator.RecordResult.Stored, + v.recordPeerNewConnectionId(i.toLong(), 0L, ByteArray(8) { i.toByte() }, ByteArray(16) { i.toByte() }), + ) + } + val r = v.recordPeerNewConnectionId(4L, 0L, ByteArray(8) { 0xAA.toByte() }, ByteArray(16)) + assertEquals(PathValidator.RecordResult.PoolFull, r) + assertEquals(3, v.unusedCount()) + } + + @Test + fun startValidationReturnsNoSpareCidWhenPoolEmpty() { + val v = PathValidator() + assertEquals(PathMigrationResult.NoSpareCid, v.tryStartValidation(nowMillis = 0L, currentPtoMillis = 100L)) + assertTrue(v.state is PathValidationState.Idle) + } + + @Test + fun startValidationPicksLowestSequenceAndQueuesChallenge() { + val supplied = byteArrayOf(0xDE.toByte(), 0xAD.toByte(), 0xBE.toByte(), 0xEF.toByte(), 0x12, 0x34, 0x56, 0x78) + val v = PathValidator(challengePayloadFactory = { supplied.copyOf() }) + v.recordPeerNewConnectionId(1L, 0L, ByteArray(8) { 0x01 }, ByteArray(16) { 0x10 }) + v.recordPeerNewConnectionId(2L, 0L, ByteArray(8) { 0x02 }, ByteArray(16) { 0x20 }) + assertEquals(PathMigrationResult.Started, v.tryStartValidation(nowMillis = 0L, currentPtoMillis = 200L)) + val state = v.state as PathValidationState.Validating + assertEquals(1L, state.newCidSequence, "lowest sequence number must be picked first") + assertContentEquals(supplied, state.challengeData) + assertEquals(1, v.pendingChallenges.size) + assertContentEquals(supplied, v.pendingChallenges.first()) + } + + @Test + fun startValidationDoesNothingIfAlreadyInProgress() { + val v = PathValidator(challengePayloadFactory = { ByteArray(8) { 0x07 } }) + v.recordPeerNewConnectionId(1L, 0L, ByteArray(8) { 0x01 }, ByteArray(16)) + v.recordPeerNewConnectionId(2L, 0L, ByteArray(8) { 0x02 }, ByteArray(16)) + assertEquals(PathMigrationResult.Started, v.tryStartValidation(0L, 100L)) + // Second trigger before resolution — must not pick a second CID. + assertEquals(PathMigrationResult.AlreadyInProgress, v.tryStartValidation(0L, 100L)) + // The second CID is still in the pool; only one challenge is in flight. + assertEquals(1, v.pendingChallenges.size) + } + + @Test + fun pathResponseWithMatchingPayloadValidatesAndQueuesRetire() { + val payload = ByteArray(8) { 0x55 } + val v = PathValidator(challengePayloadFactory = { payload.copyOf() }) + v.recordPeerNewConnectionId(1L, 0L, ByteArray(8) { 0xAA.toByte() }, ByteArray(16) { 0xBB.toByte() }) + v.tryStartValidation(0L, 100L) + val outcome = v.applyPathResponse(payload) + assertTrue(outcome is PathValidator.ValidationOutcome.Validated, "got $outcome") + assertEquals(1L, outcome.newSequence) + assertEquals(0L, outcome.retiredSequence) + assertContentEquals(ByteArray(8) { 0xAA.toByte() }, outcome.newConnectionIdBytes) + assertEquals(1L, v.activeCidSequence) + assertTrue(v.state is PathValidationState.Succeeded) + assertTrue(v.pendingRetireSequences.contains(0L), "old sequence must be queued for RETIRE_CONNECTION_ID") + assertEquals(1L, v.successfulValidations) + } + + @Test + fun pathResponseWithMismatchedPayloadIsIgnored() { + val v = PathValidator(challengePayloadFactory = { ByteArray(8) { 0x77 } }) + v.recordPeerNewConnectionId(1L, 0L, ByteArray(8) { 0x01 }, ByteArray(16)) + v.tryStartValidation(0L, 100L) + val outcome = v.applyPathResponse(ByteArray(8) { 0x11 }) + assertEquals(PathValidator.ValidationOutcome.PayloadMismatch, outcome) + assertTrue(v.state is PathValidationState.Validating, "still validating; mismatched response is dropped") + } + + @Test + fun pathResponseWithoutOutstandingChallengeIsIgnored() { + val v = PathValidator() + val outcome = v.applyPathResponse(ByteArray(8)) + assertEquals(PathValidator.ValidationOutcome.NotValidating, outcome) + } + + @Test + fun validationTimeoutAfter3PtoTransitionsToFailed() { + val v = PathValidator(challengePayloadFactory = { ByteArray(8) { 0x44 } }) + v.recordPeerNewConnectionId(1L, 0L, ByteArray(8) { 0x01 }, ByteArray(16)) + val pto = 1_000L + v.tryStartValidation(nowMillis = 0L, currentPtoMillis = pto) + + assertNull(v.checkValidationTimeout(nowMillis = pto * 2L), "still within budget") + val abandoned = v.checkValidationTimeout(nowMillis = pto * 4L) + assertTrue(abandoned != null, "validation must time out at 3*PTO") + assertEquals(1L, abandoned.newCidSequence) + assertTrue(v.state is PathValidationState.Failed) + assertEquals(1L, v.failedValidations) + } + + @Test + fun acknowledgeTerminalReturnsToIdle() { + val v = PathValidator(challengePayloadFactory = { ByteArray(8) { 0x33 } }) + v.recordPeerNewConnectionId(1L, 0L, ByteArray(8), ByteArray(16)) + v.tryStartValidation(0L, 100L) + v.applyPathResponse(ByteArray(8) { 0x33 }) + assertTrue(v.state is PathValidationState.Succeeded) + v.acknowledgeTerminal() + assertTrue(v.state is PathValidationState.Idle) + } + + @Test + fun staleOfferBelowWatermarkQueuesRetireWithoutStoring() { + val v = PathValidator() + // Seed watermark to 5 via a normal offer with retirePriorTo = 5. + v.recordPeerNewConnectionId(5L, 5L, ByteArray(8) { 0x05 }, ByteArray(16) { 0x50 }) + assertEquals(5L, v.retirePriorToWatermark) + // Now a reordered offer arrives with seq=3 and retirePriorTo=3 + // (the value the peer sent before raising it to 5). Even + // though seq=3 ≥ retirePriorTo=3 (so it isn't a frame error), + // the watermark we already raised says seq 3 has been retired. + // The validator clamps retirePriorTo to the watermark and + // detects that seq is now below the watermark — so the entry + // is queued for RETIRE_CONNECTION_ID rather than stored. + val r = v.recordPeerNewConnectionId(3L, 3L, ByteArray(8) { 0x03 }, ByteArray(16) { 0x30 }) + assertEquals(PathValidator.RecordResult.AlreadyRetired, r) + assertTrue(v.pendingRetireSequences.contains(3L), "stale offer must be retired immediately") + assertEquals(1, v.unusedCount(), "stale entry must not be stored in pool") + } +} diff --git a/quic/src/commonTest/kotlin/com/vitorpamplona/quic/connection/recovery/RecoveryTokenTest.kt b/quic/src/commonTest/kotlin/com/vitorpamplona/quic/connection/recovery/RecoveryTokenTest.kt index 1f9e5f5a8..b09badbe2 100644 --- a/quic/src/commonTest/kotlin/com/vitorpamplona/quic/connection/recovery/RecoveryTokenTest.kt +++ b/quic/src/commonTest/kotlin/com/vitorpamplona/quic/connection/recovery/RecoveryTokenTest.kt @@ -135,6 +135,8 @@ class RecoveryTokenTest { is RecoveryToken.ResetStream -> "rs:${it.streamId}:${it.errorCode}:${it.finalSize}" is RecoveryToken.StopSending -> "ss:${it.streamId}:${it.errorCode}" is RecoveryToken.NewConnectionId -> "ncid:${it.sequenceNumber}" + is RecoveryToken.PathChallenge -> "pc" + is RecoveryToken.RetireConnectionId -> "rcid:${it.sequenceNumber}" } } assertEquals(