diff --git a/.github/workflows/tests.yml b/.github/workflows/tests.yml index 0ac0c98..6dadfb8 100644 --- a/.github/workflows/tests.yml +++ b/.github/workflows/tests.yml @@ -48,6 +48,7 @@ jobs: test: name: Test runs-on: macos-26 + timeout-minutes: 10 steps: - name: Checkout code @@ -57,6 +58,9 @@ jobs: run: swift --version - name: Run all tests with coverage + timeout-minutes: 5 + env: + NSUnbufferedIO: "YES" # Warnings-as-errors keeps the package warning-free for downstream Xcode # consumers, which surface path-dependency warnings local build logs filter # out. Toolchain bumps that add new warnings must fail here. diff --git a/Sources/SendspinKit/Audio/AudioEngine.swift b/Sources/SendspinKit/Audio/AudioEngine.swift index 46729d2..b06476d 100644 --- a/Sources/SendspinKit/Audio/AudioEngine.swift +++ b/Sources/SendspinKit/Audio/AudioEngine.swift @@ -891,7 +891,9 @@ actor AudioEngine { startupReleaseDeferredChunks.removeAll(keepingCapacity: true) startupReleaseInProgress = false cancelStartupDeadline() - signalStartupCoordinator(.stateChanged) + // No self-wake: the restored chunks are identical to the ones that just + // failed the scan, so `.stateChanged` here spins forever. Only a fresh + // arrival changes the outcome, and the chunk path signals the coordinator. return } if candidate.index > 0 { diff --git a/Sources/SendspinKit/Audio/AudioOutput.swift b/Sources/SendspinKit/Audio/AudioOutput.swift index 70f1b1d..8c7a189 100644 --- a/Sources/SendspinKit/Audio/AudioOutput.swift +++ b/Sources/SendspinKit/Audio/AudioOutput.swift @@ -60,7 +60,7 @@ protocol AudioOutputPlatformMonitoring: Actor { nonisolated var requiresActiveAudioSession: Bool { get } /// Start a fresh, single-consumer observation stream. - func startMonitoring() -> AsyncStream + func startMonitoring() async -> AsyncStream /// Remove all listeners and finish the current observation stream. func stopMonitoring() diff --git a/Sources/SendspinKit/Client/SendspinClient+MultiServer.swift b/Sources/SendspinKit/Client/SendspinClient+MultiServer.swift index 39fb8a0..156bf06 100644 --- a/Sources/SendspinKit/Client/SendspinClient+MultiServer.swift +++ b/Sources/SendspinKit/Client/SendspinClient+MultiServer.swift @@ -70,11 +70,30 @@ extension SendspinClient { case .keepExisting: await HandshakeDriver.reject(outcome, reason: .concurrentAttempt, on: transport) case .acceptIncoming: + // Promotion is a session transition: claim a fresh epoch so a parked + // connect/accept at the older epoch cannot install over this winner. + guard sessionEpoch == arbitrationEpoch else { + await transport.disconnect() + return + } + sessionEpoch += 1 + let promotionEpoch = sessionEpoch if let incumbent = retireSession() { await incumbent.disconnect(reason: .anotherServer) } + // The incumbent teardown suspends; a disconnect may have landed. + guard sessionEpoch == promotionEpoch else { + await transport.disconnect() + return + } updateConnectionState(.connecting) - await setupConnection(with: transport, outcome: outcome, negotiation: negotiation) + await setupConnection( + with: transport, + outcome: outcome, + negotiation: negotiation, + runtimeConfiguration: runtimeConfiguration, + setupEpoch: promotionEpoch + ) } } catch { await transport.disconnect() diff --git a/Sources/SendspinKit/Client/SendspinClient.swift b/Sources/SendspinKit/Client/SendspinClient.swift index 7c06e7b..d875040 100644 --- a/Sources/SendspinKit/Client/SendspinClient.swift +++ b/Sources/SendspinKit/Client/SendspinClient.swift @@ -447,10 +447,8 @@ public final class SendspinClient { throw error } - // Re-validate: the guard above ran before a network-length suspension, during - // which an inbound server can win arbitration and be promoted, or the caller can - // call `disconnect()`. Either bumps the epoch, and neither is visible in - // `connection` — a cancelled dial leaves it nil, exactly as an untouched one does. + // Re-validate after the network suspension: a promotion or a disconnect bumps + // the epoch, and neither is visible in `connection` (nil either way). guard sessionEpoch == dialEpoch else { Log.client.warning("The session changed while dialing \(url); abandoning this dial") await transport.disconnect() @@ -481,7 +479,13 @@ public final class SendspinClient { await transport.disconnect() throw SendspinClientError.alreadyConnected } - await setupConnection(with: transport, outcome: outcome, negotiation: negotiation) + await setupConnection( + with: transport, + outcome: outcome, + negotiation: negotiation, + runtimeConfiguration: runtimeConfiguration, + setupEpoch: dialEpoch + ) try requireOpen() } catch { await transport.disconnect() @@ -507,6 +511,10 @@ public final class SendspinClient { defer { pendingTransports.removeValue(forKey: pendingID) } if connectionState == .disconnected { connectionState = .connecting + // Claim the epoch before the first suspension: the dial window holds no + // `connection`, so only the epoch tracks caller intent. + sessionEpoch += 1 + let acceptEpoch = sessionEpoch await preparePairingConfiguration() do { let negotiation = try await makeSessionFormatNegotiation() @@ -527,10 +535,29 @@ public final class SendspinClient { phaseTimeout: handshakeTimeout ) try requireOpen() - await setupConnection(with: transport, outcome: outcome, negotiation: negotiation) + guard sessionEpoch == acceptEpoch else { + // A disconnect/close or a promoted competitor invalidated this claim. + await transport.disconnect() + throw SendspinClientError.alreadyConnected + } + await setupConnection( + with: transport, + outcome: outcome, + negotiation: negotiation, + runtimeConfiguration: runtimeConfiguration, + setupEpoch: acceptEpoch + ) + try requireOpen() } catch { await transport.disconnect() - updateConnectionState(.disconnected) + // Only a candidate still owning this epoch may reset the visible + // state; a stale failure must not clobber a replacement session. + if sessionEpoch == acceptEpoch, connectionState == .connecting { + updateConnectionState(.disconnected) + } + if isTerminated { + throw TerminatedError() + } throw error } } else { @@ -561,26 +588,29 @@ public final class SendspinClient { return retired } - /// - Parameter preReadHello: When non-nil, the `client/hello` was already sent - /// and the `server/hello` already consumed during competing-connection - /// arbitration. In that case we process the hello directly instead of sending - /// another `client/hello`, and the message loop resumes the transport's stream - /// from the (buffered) frames that follow. + /// Install an admitted session on `transport` and start it. /// - /// This setup path is intentionally non-throwing: all genuine dial/handshake - /// failures are handled before a transport reaches this point, so callers do not - /// need duplicate rollback logic after they set `.connecting`. + /// Caller already resolved the pairing runtime snapshot and claimed its + /// epoch: never install for an epoch that lost its claim. Non-throwing. @MainActor // swiftlint:disable:next function_body_length func setupConnection( with transport: any SendspinTransport, outcome: consuming HandshakeDriver.Result, - negotiation: SessionFormatNegotiation + negotiation: SessionFormatNegotiation, + runtimeConfiguration: PairingManagementConfiguration, + setupEpoch: Int ) async { guard !isTerminated else { await transport.disconnect() return } + // Claim the epoch re-check: an interleaved `disconnect()` or promoted + // competitor bumps the epoch; this candidate must not install for it. + guard sessionEpoch == setupEpoch else { + await transport.disconnect() + return + } // A new connection is a new session: drop any server-reported state carried // over from a prior connection (notably one lost without an explicit // disconnect) before the first server/state update is applied. @@ -592,14 +622,16 @@ public final class SendspinClient { // Retire the old session synchronously (token + identity guards both // reject its late events from this point), then await its teardown. - // - // `oldConnection` is nil for every current caller, so this await does not run. If - // that changes, the nil-`connection` window makes `disconnect()` a silent no-op and - // can orphan a live connection. Keep the rest of this suspension-free. + // `oldConnection` is nil for current callers. Re-check the epoch after + // that suspension: teardown can take arbitrarily long. let oldConnection = retireSession() if let oldConnection { await oldConnection.shutdown() } + guard sessionEpoch == setupEpoch else { + await transport.disconnect() + return + } // Build the SendspinConnection with configuration from this facade let validity = SessionValidityToken() @@ -624,7 +656,6 @@ public final class SendspinClient { let outcomeServerStaticPublicKey = outcome.serverStaticPublicKey let outcomeSuite = outcome.suite let sessionChannel = outcome.takeChannel() - let runtimeConfiguration = await pairingRuntimeConfiguration() #if DEBUG let nonceBOverride = nonceBOverride let pairingHandshakeHashOverride = pairingHandshakeHashOverride @@ -713,6 +744,7 @@ public final class SendspinClient { clock: clockSync, engine: audioEngine ) + // No suspension occurs between the re-check above and this install. connection = newConnection currentOutputFormatStatus = nil diff --git a/Sources/SendspinKit/Client/SendspinConnection+MessageHandling.swift b/Sources/SendspinKit/Client/SendspinConnection+MessageHandling.swift index f2d3e29..28dc533 100644 --- a/Sources/SendspinKit/Client/SendspinConnection+MessageHandling.swift +++ b/Sources/SendspinKit/Client/SendspinConnection+MessageHandling.swift @@ -767,7 +767,13 @@ extension SendspinConnection { } func pairingAttemptTimedOut() async { + // A stale wake (its handle was cancelled and replaced by a newer attempt + // or teardown) must not detach the newer handle or abort the fresh attempt. + guard !Task.isCancelled else { return } guard pendingPairingPsk != nil || dynamicPairingAttempt != nil || staticPairingAttempt != nil else { return } + // Detach this task's handle before clear: clearPairingAttempt cancels the + // owned task, which would self-cancel the abort send below. + pairingAttemptTask = nil clearPairingAttempt(reason: .attemptTimeout) try? await sendWrapped(PairAbortMessage(payload: PairAbortPayload(reason: .attemptTimeout))) } diff --git a/Sources/SendspinKit/Client/SendspinConnection+Outbound.swift b/Sources/SendspinKit/Client/SendspinConnection+Outbound.swift index 5afeb91..6e5ba66 100644 --- a/Sources/SendspinKit/Client/SendspinConnection+Outbound.swift +++ b/Sources/SendspinKit/Client/SendspinConnection+Outbound.swift @@ -3,11 +3,62 @@ import Foundation extension SendspinConnection { // MARK: - Outbound sends + /// Park until the outbound slot is free, then take it. FIFO, no busy spin. + private func acquireOutboundSlot() async { + if outboundInFlight { + await withCheckedContinuation { outboundWaiters.append($0) } + } + outboundInFlight = true + } + + /// Free the slot or hand it to the queued head (inFlight stays true so a + /// fresh sender cannot steal it between the release and the wake). + private func releaseOutboundSlot() { + if outboundWaiters.isEmpty { + outboundInFlight = false + } else { + outboundWaiters.removeFirst().resume() + } + } + + /// A send failure burns nonces, so the channel is crypto-dead. Latch the + /// failure before the async teardown: the deferred slot release then chains + /// queued senders into the latch, which rejects them without encrypting. + private func failOutbound() async { + outboundFailed = true + if !shuttingDown { + shuttingDown = true + if disconnectReason == nil { + disconnectReason = .connectionLost(nil) + } + } + await transport.disconnect() + } + func sendWrapped(_ message: some Codable & Sendable, bypassRehandshakeGate: Bool = false) async throws { - // Check the gate before encryption; a post-encryption re-check would consume a nonce before throwing. + await acquireOutboundSlot() + defer { releaseOutboundSlot() } + + guard !outboundFailed else { + throw SendspinClientError.sendFailed("outbound channel is dead") + } + // Gate check comes after acquisition: a sender that parked during the + // exchange must not encrypt under pre-swap keys. guard bypassRehandshakeGate || !rehandshakeInProgress else { throw SendspinClientError.handshakeIncomplete } + guard lifecycle == .running || lifecycle == .shuttingDown else { + // `.shuttingDown` permits the intentional goodbye; everything else + // on a stopped connection is rejected. + throw SendspinClientError.notConnected + } + if Task.isCancelled { + // The pairing timeout handler detaches its own task handle before + // clearing, so a cancelled sender here is always abandoned work: + // don't burn a nonce for a frame nothing will carry. + throw CancellationError() + } + let data = try SendspinEncoding.makeEncoder().encode(message) var plaintext = Data([NoiseFrameType.json]) plaintext.append(data) @@ -16,17 +67,7 @@ extension SendspinConnection { try await transport.sendBinary(frame) } } catch { - // A failed send is terminal: encryption already consumed AEAD nonces, - // so the peer can never decrypt a later frame — the session is - // cryptographically dead, not merely degraded. Tear down (unless a - // teardown is already driving this send) and surface the error. - if !shuttingDown { - shuttingDown = true - if disconnectReason == nil { - disconnectReason = .connectionLost(nil) - } - await transport.disconnect() - } + await failOutbound() throw error } } @@ -35,10 +76,6 @@ extension SendspinConnection { /// Send a facade-initiated protocol message, wrapping transport errors in /// the public typed ``SendspinClientError/sendFailed(_:)``. - /// - /// All outbound protocol I/O flows through this actor — the facade holds - /// no send path of its own — so public API sends serialize with the - /// handshake/time/state/goodbye sequencing this actor owns. func send(clientMessage message: some Codable & Sendable) async throws { guard lifecycle == .running, !rehandshakeInProgress else { throw SendspinClientError.handshakeIncomplete @@ -63,7 +100,9 @@ extension SendspinConnection { while clientStateDirty { clientStateDirty = false let payload = currentClientStatePayload() - try await sendWrapped(ClientStateMessage(payload: payload)) + // Forward the rehandshake bypass: handleServerActivate publishes the + // post-swap full state while rehandshakeInProgress is still true. + try await sendWrapped(ClientStateMessage(payload: payload), bypassRehandshakeGate: bypassRehandshakeGate) if payload.player != nil { playerStateSent = true } diff --git a/Sources/SendspinKit/Client/SendspinConnection.swift b/Sources/SendspinKit/Client/SendspinConnection.swift index 676a646..5ceba51 100644 --- a/Sources/SendspinKit/Client/SendspinConnection.swift +++ b/Sources/SendspinKit/Client/SendspinConnection.swift @@ -102,6 +102,15 @@ actor SendspinConnection { var clientStateSendInFlight = false var clientStateDirty = false + /// Outbound whole-message fence: nonces burn per fragment up front and the + /// fragments must reach the transport with nothing interleaved, so one + /// message sends at a time; the rest park on `outboundWaiters` (FIFO). + var outboundInFlight = false + var outboundWaiters: [CheckedContinuation] = [] + /// Set when an outbound send fails: a burned nonce makes the channel + /// crypto-dead, so queued and later senders must fail without encrypting. + var outboundFailed = false + /// Server info var currentServerId: String? var serverName: String diff --git a/Sources/SendspinKit/Transport/NoiseSessionEstablisher.swift b/Sources/SendspinKit/Transport/NoiseSessionEstablisher.swift index 44aa3d6..ef4d8bc 100644 --- a/Sources/SendspinKit/Transport/NoiseSessionEstablisher.swift +++ b/Sources/SendspinKit/Transport/NoiseSessionEstablisher.swift @@ -178,39 +178,42 @@ enum NoiseSessionEstablisher { ) } - /// Pull the next frame, requiring a text frame within `timeout`. The timeout is - /// a watchdog that *disconnects the transport*: a parked `nextFrame()` pull is - /// released by `disconnect()` finishing the frame stream, never by cancellation - /// (the FrameInbox contract). + /// Both timeout and caller cancellation must disconnect: cancellation alone + /// cannot release a parked `FrameInbox` pull. A fired watchdog wins over an + /// arriving frame because the transport is already being closed. private static func nextTextFrame( from transport: any SendspinTransport, timeout: Duration ) async throws -> Data { - let watchdog = Task { () -> Bool in - do { - try await Task.sleep(for: timeout) - } catch { - return false // cancelled: the frame arrived in time + try await withTaskCancellationHandler(operation: { + let watchdog = Task { () -> Bool in + do { + try await Task.sleep(for: timeout) + } catch { + return false // cancelled: the frame arrived in time + } + await transport.disconnect() + return true + } + let frame = await transport.nextFrame() + watchdog.cancel() + let timedOut = await watchdog.value + // The deadline wins a race with an arriving frame: the watchdog already + // disconnected, so proceeding would run the next phase on a dead connection. + if timedOut { + throw HandshakeError.timeout } - await transport.disconnect() - return true - } - let frame = await transport.nextFrame() - watchdog.cancel() - let timedOut = await watchdog.value - // The deadline wins a race with an arriving frame: the watchdog already - // disconnected, so proceeding would run the next phase on a dead connection. - if timedOut { - throw HandshakeError.timeout - } - switch frame { - case let .text(text): - return Data(text.utf8) - case .binary: - throw HandshakeError.malformed - case nil: - throw HandshakeError.transportClosed - } + switch frame { + case let .text(text): + return Data(text.utf8) + case .binary: + throw HandshakeError.malformed + case nil: + throw HandshakeError.transportClosed + } + }, onCancel: { + Task { await transport.disconnect() } + }) } } diff --git a/Tests/SendspinKitTests/Audio/AudioStartupReleaseTests.swift b/Tests/SendspinKitTests/Audio/AudioStartupReleaseTests.swift index 85e57f3..481a956 100644 --- a/Tests/SendspinKitTests/Audio/AudioStartupReleaseTests.swift +++ b/Tests/SendspinKitTests/Audio/AudioStartupReleaseTests.swift @@ -348,6 +348,49 @@ struct AudioStartupReleaseTests { await engine.shutdown() } + /// With no self-wake in the all-stale branch, the coordinator parks after the stale + /// chunks drain instead of re-scanning identical data; a new arrival re-enters it. + @Test("an all-stale startup buffer parks without self-woken re-evaluation") + func allStaleBufferParksWithoutSelfWake() async throws { + let clock = StubClock(anchorToNow: true) + let output = SpyAudioOutput() + let scheduler = AudioScheduler(clockSync: clock) + let engine = AudioEngine(output: output, scheduler: scheduler, clock: clock, enableStartupBuffering: true) + let format = try AudioFormatSpec(codec: .pcm, channels: 2, sampleRate: 48_000, bitDepth: 16) + await engine.start() + await engine.commands.enqueue(.streamStart(format, codecHeader: nil)) + + let staleCount = 3 + for index in 0 ..< staleCount { + await engine.commands.enqueue( + .chunk(Data(repeating: UInt8(index), count: 100), ts: -1_000_000 + Int64(index) * 20_000) + ) + } + #expect( + await waitUntil { await engine.appliedCommandKinds().count(where: { $0 == .chunk }) == staleCount }, + "the stale chunks should have reached the engine" + ) + #expect(await !output.recordedCalls.contains("startPrepared()"), "nothing is viable yet") + + // Without parking, the self-signal would keep burning evaluations on the same + // buffer; a tight spread here means the coordinator went back to waiting. + let afterDrain = await engine.startupReleaseEvaluations + try? await Task.sleep(for: .milliseconds(100)) + let afterQuiet = await engine.startupReleaseEvaluations + #expect( + afterQuiet - afterDrain <= staleCount, + "the all-stale buffer kept re-evaluating (\(afterDrain) -> \(afterQuiet))" + ) + + // A single viable chunk arrival must restart the stalled release. + await engine.commands.enqueue(.chunk(Data(repeating: 0xAA, count: 100), ts: 1_500_000)) + #expect( + await waitUntil(timeout: .seconds(3)) { await engine.startupReleaseCommits == 1 }, + "a fresh viable chunk must commit a release" + ) + await engine.shutdown() + } + @Test("startup release remains single-flight while a deadline probe is suspended") func startupReleaseRemainsSingleFlightWhileDeadlineProbeIsSuspended() async throws { let clock = StubClock(anchorToNow: true) diff --git a/Tests/SendspinKitTests/Client/DynamicPairingTranscriptTests.swift b/Tests/SendspinKitTests/Client/DynamicPairingTranscriptTests.swift index 5f1a1d3..59943d4 100644 --- a/Tests/SendspinKitTests/Client/DynamicPairingTranscriptTests.swift +++ b/Tests/SendspinKitTests/Client/DynamicPairingTranscriptTests.swift @@ -533,6 +533,7 @@ struct DynamicPairingTimeoutTests { @Test("attempt timeout uses the exact attempt_timeout reason") func attemptTimeout() async throws { let session = try await makeDynamicTestSession(attemptTimeout: .milliseconds(100)) + await session.server.transport.setHonorCancellationSends(true) try await activateDynamic(session.server) _ = try await waitForClientMessage(session.server, type: ClientPairInitMessage.typeString) let abortData = try await waitForClientMessage(session.server, type: PairAbortMessage.typeString) diff --git a/Tests/SendspinKitTests/Client/FacadeLifecycleTests.swift b/Tests/SendspinKitTests/Client/FacadeLifecycleTests.swift new file mode 100644 index 0000000..7cbd542 --- /dev/null +++ b/Tests/SendspinKitTests/Client/FacadeLifecycleTests.swift @@ -0,0 +1,122 @@ +import Foundation +@testable import SendspinKit +import Testing + +@MainActor +struct FacadeLifecycleTests { + @Test("a disconnect while an accept handshakes prevents a later install") + func acceptCancelledByDisconnectCannotInstall() async throws { + let client = try makeTestClient() + let transport = MockTransport() + let server = MockNoiseServer(transport: transport, psk: .sentinel) + + let accepted = Task { try? await client.acceptConnection(transport) } + #expect(await waitUntil { await transport.hasSentFrames }) + + await client.disconnect(reason: .userRequest) + // Complete the handshake after the disconnect: the accept's epoch guard must win. + try await server.establishSession(activities: [.playback], activeRoles: [.playerV1]) + _ = await accepted.value + + #expect(client.connection == nil) + #expect(await transport.disconnectCalled) + #expect(client.connectionState == .disconnected) + } + + @Test("a competing promotion supersedes a parked primary accept") + func competingPromotionBumpsPrimaryEpoch() async throws { + let client = try makeTestClient() + + let primary = MockTransport() + let primaryServer = MockNoiseServer(transport: primary, psk: .sentinel) + let primaryAccept = Task { try? await client.acceptConnection(primary) } + #expect(await waitUntil { await primary.hasSentFrames }) + + let competitor = MockTransport() + let competitorServer = MockNoiseServer(transport: competitor, psk: .sentinel) + let competitorAccept = Task { try? await client.acceptConnection(competitor) } + try await competitorServer.establishSession(activities: [.playback], activeRoles: [.playerV1]) + _ = await competitorAccept.value + #expect(await waitUntil { await MainActor.run { client.connectionState == .connected } }) + let promoted = client.connection + + try await primaryServer.establishSession(activities: [.playback], activeRoles: [.playerV1]) + _ = await primaryAccept.value + + #expect(client.connection === promoted) + #expect(await primary.disconnectCalled) + #expect(await competitor.disconnectCalled == false) + await client.disconnect() + } + + @Test("a disconnect during promotion teardown prevents the install") + func disconnectDuringPromotionTeardownPreventsInstall() async throws { + let client = try makeTestClient() + + let incumbent = MockTransport() + let incumbentServer = MockNoiseServer(transport: incumbent, psk: .sentinel) + async let incumbentAccepted: Void = client.acceptConnection(incumbent) + try await incumbentServer.establishSession(activities: [], activeRoles: []) + try await incumbentAccepted + #expect(await waitUntil { await MainActor.run { client.connectionState == .connected } }) + + let candidate = MockTransport() + let candidateServer = MockNoiseServer(transport: candidate, psk: .sentinel) + let replacementAccept = Task { try? await client.acceptConnection(candidate) } + #expect(await waitUntil { await candidate.hasSentFrames }) + + // Park the incumbent's next send so promotion's teardown stalls. + await incumbent.enableGoodbyeGate() + try await candidateServer.establishSession(activities: [.playback], activeRoles: [.playerV1]) + #expect(await waitUntil { await incumbent.isGoodbyeGateWaiting }) + + await client.disconnect(reason: .userRequest) + await incumbent.releaseGoodbyeGate() + _ = await replacementAccept.value + + #expect(client.connection == nil) + #expect(await candidate.disconnectCalled) + #expect(client.connectionState == .disconnected) + } + + @Test("a failing parked accept cannot clobber a replacement session") + func abandonedFailingAcceptCannotClobberReplacement() async throws { + let client = try makeTestClient() + + let first = MockTransport() + let firstAccept = Task { try? await client.acceptConnection(first) } + #expect(await waitUntil { await first.hasSentFrames }) + + let replacement = MockTransport() + let replacementServer = MockNoiseServer(transport: replacement, psk: .sentinel) + let replacementAccept = Task { try? await client.acceptConnection(replacement) } + try await replacementServer.establishSession(activities: [.playback], activeRoles: [.playerV1]) + _ = await replacementAccept.value + #expect(await waitUntil { await MainActor.run { client.connectionState == .connected } }) + let installed = client.connection + + // Fail the parked accept's handshake; its catch must not touch the winner. + await first.finishStreams() + _ = await firstAccept.value + + #expect(client.connection === installed) + #expect(client.connectionState == .connected) + #expect(await replacement.disconnectCalled == false) + await client.disconnect() + } + + @Test("close() during a paused accept terminates without installing") + func closeDuringPausedAcceptTerminates() async throws { + let client = try makeTestClient() + let transport = MockTransport() + + let accepted = Task { try await client.acceptConnection(transport) } + #expect(await waitUntil { await transport.hasSentFrames }) + + await client.close() + await #expect(throws: TerminatedError.self) { try await accepted.value } + + #expect(client.connection == nil) + #expect(client.connectionState == .disconnected) + } +} diff --git a/Tests/SendspinKitTests/Client/RehandshakeTests.swift b/Tests/SendspinKitTests/Client/RehandshakeTests.swift index b0bb43b..db84543 100644 --- a/Tests/SendspinKitTests/Client/RehandshakeTests.swift +++ b/Tests/SendspinKitTests/Client/RehandshakeTests.swift @@ -228,21 +228,27 @@ struct RehandshakeTests { ) } + /// Timeout cleanup must leave the abort sender uncancelled; the + /// cancellation-aware mock rejects a send that arrives pre-cancelled. @Test("Pairing attempt timeout aborts without persistence") func pairingAttemptTimeoutAbortsWithoutPersistence() async throws { let session = try await makePairableSession(pairingAttemptTimeout: .milliseconds(100)) let server = session.server + await server.transport.setHonorCancellationSends(true) try await rehandshake(server, to: session.pairingPsk) try await server.sendJSON( #"{"type":"server/activate","payload":{"activities":["pairing"],"active_roles":[],"pairing":{"method":"pairing_psk"}}}"# ) - #expect(await waitUntil { - await server.clientJSONMessages(ofType: PairAbortMessage.typeString).count == 1 - }, "the pending attempt must expire") + #expect( + await waitUntil(timeout: .seconds(3)) { + await server.clientJSONMessages(ofType: PairAbortMessage.typeString).count == 1 + }, + "the pending attempt must expire AND the intentional pair/abort must reach the server" + ) let abort = try #require(await server.clientJSONMessages(ofType: PairAbortMessage.typeString).first) #expect(try JSONDecoder().decode(PairAbortMessage.self, from: abort).payload.reason == .attemptTimeout) - #expect(await session.store.listRecords().allSatisfy { $0.serverId == nil }) + #expect(await session.store.listRecords().allSatisfy { $0.serverId == nil }, "a timed-out attempt must not persist") #expect(await session.client.connectionState == .connected) await session.client.disconnect() } @@ -333,6 +339,105 @@ struct RehandshakeTests { await session.client.disconnect() } + /// A sender parked on the outbound queue before a rehandshake re-checks the + /// gate once woken: the gate closes on message-1 receipt, ahead of the reply's + /// key swap, so the woken sender is rejected while the swap is still pending. + @Test("a queued sender woken under the closed rehandshake gate is rejected before the key swap") + func queuedSenderWokenUnderClosedGateIsRejected() async throws { + let longTermPsk = Psk.generate() + let session = try await makePairableSession(seededLongTermPsk: longTermPsk) + let server = session.server + let transport = server.transport + let connection = try #require(await MainActor.run { session.client.connection }) + // start() returns before messageLoop installs the clock-sync sampler; + // wait for the handle first so cancel/join drains it deterministically. + #expect(await waitUntil { await connection.clockSyncTask != nil }, "clock-sync task handle must appear before cancel") + await connection.clockSyncTask?.cancel() + await connection.clockSyncTask?.value + + // #1 takes the outbound slot and parks mid-fragment on the gate. + await transport.enableGoodbyeGate() + let first = Task { () -> Result in + do { + try await connection.send(clientMessage: + OutboundTestMessage( + type: OutboundTestMessageType.padded, + note: String(repeating: "g", count: NoiseChannel.maxSinglePayload + 2_000) + ) + ) + return .success(()) + } catch { + return .failure(error) + } + } + #expect(await waitUntil { await transport.isGoodbyeGateWaiting }) + + // #2 queues behind #1 (it has NOT acquired the slot, so it has not checked + // the gate — it will only do so once woken). + let second = Task { () -> Result in + do { + try await connection.send(clientMessage: OutboundTestMessage(type: OutboundTestMessageType.small, note: "queued-before-rekey")) + return .success(()) + } catch { + return .failure(error) + } + } + #expect(await waitUntil { await connection.outboundWaiters.count == 1 }, "the second sender must queue") + + // beginRehandshake only injects message 1; the gate closes when the + // connection consumes it. Wait for that before releasing the sender, and + // release unconditionally so failure cannot wedge the parked send. + try await server.beginRehandshake(to: longTermPsk, pskCategoryOverride: .longTerm) + #expect(await waitUntil { await connection.isRehandshakeInProgress }) + await transport.releaseGoodbyeGate() + + let firstResult = await first.value + #expect((try? firstResult.get()) != nil, "the first send must complete") + + // #2 wakes under the (now closed) rehandshake gate and is rejected. + let secondResult = await second.value + let secondWasRejected: Bool = { + guard case .failure = secondResult else { return false } + return true + }() + #expect(secondWasRejected, "a sender that parked before the re-key swap must be gate-rejected when woken after it") + + // The reply sent bypass under the old keys; the swap lands on the client. + #expect(await waitUntil { await server.rehandshakeComplete }) + + // Neither queued-before-rekey send may have hit the wire under the old or + // new keys: the small one was gate-rejected (no encrypt); the padded one + // completed before the handshake and is legitimate old-key traffic. + let wireTypes = await server.decryptedMessages.compactMap(typeOfDecryptedJSON) + #expect(!wireTypes.contains(OutboundTestMessageType.small), "a gate-rejected send must not reach the wire at all") + + try await server.sendJSON(#"{"type":"server/hello","payload":{"name":"Test Server"}}"#) + #expect(await waitUntil { await server.clientJSONMessages(ofType: ClientHelloMessage.typeString).count == 1 }) + await session.client.disconnect() + } + + @Test("Post-swap activate publishes the full client/state under the new keys") + func postSwapActivationPublishesClientState() async throws { + let longTermPsk = Psk.generate() + let session = try await makePairableSession(seededLongTermPsk: longTermPsk) + let server = session.server + + try await rehandshake(server, to: longTermPsk, pskCategory: .longTerm) + + // Between the swap and the post-swap activate, publishClientState must be + // rejected by the gate (the rehandshakeInProgress bypass is not yet set). + #expect(await server.clientJSONMessages(ofType: ClientStateMessage.typeString).isEmpty) + + try await server.sendActivation(activities: [], activeRoles: []) + #expect( + await waitUntil(timeout: .seconds(3)) { + await server.clientJSONMessages(ofType: ClientStateMessage.typeString).count == 1 + }, + "the completed rehandshake must publish the post-swap full client/state under the new keys" + ) + await session.client.disconnect() + } + @Test("Cancelling pairing discards the attempt and ignores a late finalize") func cancellingPairingDiscardsLateFinalize() async throws { // On a Pairing PSK session the only admissible activity set is ['pairing'], diff --git a/Tests/SendspinKitTests/Client/SendspinConnectionTests.swift b/Tests/SendspinKitTests/Client/SendspinConnectionTests.swift index 3a3a2c0..6dde728 100644 --- a/Tests/SendspinKitTests/Client/SendspinConnectionTests.swift +++ b/Tests/SendspinKitTests/Client/SendspinConnectionTests.swift @@ -1456,6 +1456,364 @@ struct SendspinConnectionSessionTests { "start() after an idle shutdown must be a no-op (no client/hello)" ) } + + // MARK: - Outbound whole-message serialization + + /// Core regression: a fragmented message parks mid-send with all fragment + /// nonces already consumed; the peer must still decrypt both messages in send + /// order with no AEAD gap. + @Test("fragmented outbound message reaches the peer before a concurrent message (no nonce gap)") + func outboundMessageFragmentsSerializeBeforeConcurrentMessage() async throws { + let transport = MockTransport() + let connection = try await makeConnectionWithTransport(transport) + // start() returns before messageLoop installs the clock-sync sampler; + // wait for the handle first so cancel/join drains it deterministically. + #expect(await waitUntil { await connection.clockSyncTask != nil }, "clock-sync task handle must appear before cancel") + await connection.clockSyncTask?.cancel() + await connection.clockSyncTask?.value + // Drain any residual in-flight send so no unrelated sender races the test. + #expect(await waitUntil { await !(connection.outboundInFlight) }, "initial clock samples must drain") + + await transport.enableGoodbyeGate() + let big = OutboundTestMessage( + type: OutboundTestMessageType.padded, + note: String(repeating: "a", count: NoiseChannel.maxSinglePayload + 2_000) + ) + async let bigSend: Void = connection.send(clientMessage: big) + + // The first fragment of the fragmented message parks mid-send; every + // fragment's nonces are already consumed by encryptMessage. + #expect(await waitUntil { await transport.isGoodbyeGateWaiting }, "the first fragment must park on the transport gate") + + let small = OutboundTestMessage(type: OutboundTestMessageType.small, note: "after") + async let smallSend: Void = connection.send(clientMessage: small) + // The concurrent sender must be QUEUED on the outbound slot, not + // encrypting — the exact interleaving that used to burn nonce n+2 + // before n+1. + #expect( + await waitUntil { await connection.outboundWaiters.count == 1 }, + "the second sender must park on the outbound slot" + ) + + await transport.releaseGoodbyeGate() + try await bigSend + try await smallSend + + // The peer must have decrypted both messages with contiguous AEAD nonces: + // a gap makes decryptFrame throw, so the message would never appear here. + let server = try #require(await connectionReadbacks.server(for: transport)) + #expect( + await waitUntil(timeout: .seconds(3)) { + await server.decryptedMessages.contains { typeOfDecryptedJSON($0) == OutboundTestMessageType.small } + }, + "the fragmented message must decrypt at the peer (no nonce gap)" + ) + let types = await server.decryptedMessages.compactMap(typeOfDecryptedJSON) + let bigIndex = try #require(types.lastIndex(of: OutboundTestMessageType.padded)) + let smallIndex = try #require(types.lastIndex(of: OutboundTestMessageType.small)) + #expect(bigIndex < smallIndex, "the fragmented message must complete before the concurrent one") + await transport.finishStreams() + await connection.shutdown() + } + + /// Three concurrent senders: the fragmented message is admitted first, then + /// two single-frame messages queue behind it. All three must decrypt in the + /// order they were enqueued (FIFO). + @Test("three concurrent outbound messages decrypt in FIFO order") + func outboundQueueDeliveryIsFIFO() async throws { + let transport = MockTransport() + let connection = try await makeConnectionWithTransport(transport) + // start() returns before messageLoop installs the clock-sync sampler; + // wait for the handle first so cancel/join drains it deterministically. + #expect(await waitUntil { await connection.clockSyncTask != nil }, "clock-sync task handle must appear before cancel") + await connection.clockSyncTask?.cancel() + await connection.clockSyncTask?.value + #expect(await waitUntil { await !(connection.outboundInFlight) }) + + await transport.enableGoodbyeGate() + let big = OutboundTestMessage( + type: OutboundTestMessageType.padded, + note: String(repeating: "b", count: NoiseChannel.maxSinglePayload + 2_000) + ) + async let first: Void = connection.send(clientMessage: big) + #expect(await waitUntil { await transport.isGoodbyeGateWaiting }) + + // Sequence: wait for the FIRST queued sender before launching the second, + // so both are deterministically in the queue when the gate opens. + let second = Task { () -> Result in + do { + try await connection.send(clientMessage: OutboundTestMessage(type: OutboundTestMessageType.one, note: "one")) + return .success(()) + } catch { + return .failure(error) + } + } + #expect(await waitUntil { await connection.outboundWaiters.count == 1 }, "the first single-frame sender must queue") + + let third = Task { () -> Result in + do { + try await connection.send(clientMessage: OutboundTestMessage(type: OutboundTestMessageType.two, note: "two")) + return .success(()) + } catch { + return .failure(error) + } + } + #expect(await waitUntil { await connection.outboundWaiters.count == 2 }, "both single-frame senders must queue") + + await transport.releaseGoodbyeGate() + try await first + let secondResult = await second.value + _ = try? secondResult.get() + let thirdResult = await third.value + _ = try? thirdResult.get() + + // Wait for the peer readback of all three before snapshotting order. + let server = try #require(await connectionReadbacks.server(for: transport)) + #expect( + await waitUntil(timeout: .seconds(3)) { + await server.decryptedMessages.count >= 3 + }, + "the peer must read all three messages" + ) + let types = await server.decryptedMessages.compactMap(typeOfDecryptedJSON) + let ours = types.filter { OutboundTestMessageType.all.contains($0) } + #expect( + ours == [OutboundTestMessageType.padded, OutboundTestMessageType.one, OutboundTestMessageType.two], + "messages must reach the peer in FIFO enqueue order" + ) + await transport.finishStreams() + await connection.shutdown() + } + + /// A queued sender cancelled while parked must not encrypt a frame (no nonce + /// burn), must not wedge the chain, and the slot must stay usable afterwards. + @Test("a queued outbound send cancelled while parked must not burn a nonce or wedge the chain") + func queuedCancellationDoesNotBurnNonceOrWedgeChain() async throws { + let transport = MockTransport() + let connection = try await makeConnectionWithTransport(transport) + // start() returns before messageLoop installs the clock-sync sampler; + // wait for the handle first so cancel/join drains it deterministically. + #expect(await waitUntil { await connection.clockSyncTask != nil }, "clock-sync task handle must appear before cancel") + await connection.clockSyncTask?.cancel() + await connection.clockSyncTask?.value + #expect(await waitUntil { await !(connection.outboundInFlight) }) + + await transport.enableGoodbyeGate() + let big = OutboundTestMessage( + type: OutboundTestMessageType.padded, + note: String(repeating: "c", count: NoiseChannel.maxSinglePayload + 2_000) + ) + let bigTask = Task { () -> Result in + do { + try await connection.send(clientMessage: big) + return .success(()) + } catch { + return .failure(error) + } + } + #expect(await waitUntil { await transport.isGoodbyeGateWaiting }) + + let smallTask = Task { () -> Result in + do { + try await connection.send(clientMessage: OutboundTestMessage(type: OutboundTestMessageType.small, note: "cancelled")) + return .success(()) + } catch { + return .failure(error) + } + } + // WAIT for the small sender to be parked BEFORE cancelling, so it is + // cancelled while parked (wasCancelledBeforeWaiting == false) and the + // parked-cancel path deterministically drops it. + #expect(await waitUntil { await connection.outboundWaiters.count == 1 }, "the small sender must queue") + + smallTask.cancel() + await transport.releaseGoodbyeGate() + + let bigResult = await bigTask.value + #expect((try? bigResult.get()) != nil, "the fragmented message must still send") + + let smallResult = await smallTask.value + guard case .failure = smallResult else { + Issue.record("a cancelled queued send must fail") + return + } + #expect(await connection.outboundWaiters.isEmpty, "the cancelled waiter must release the slot") + + // The cancelled message never encrypted, so it must not appear at the peer. + let server = try #require(await connectionReadbacks.server(for: transport)) + let wireTypes = await server.decryptedMessages.compactMap(typeOfDecryptedJSON) + #expect(!wireTypes.contains(OutboundTestMessageType.small), "a cancelled send must not reach the wire") + + // The chain is still live: a follow-up send goes out under the next nonce. + try await connection.send(clientMessage: OutboundTestMessage(type: OutboundTestMessageType.one, note: "after-cancel")) + #expect( + await waitUntil(timeout: .seconds(3)) { + await server.decryptedMessages.contains { typeOfDecryptedJSON($0) == OutboundTestMessageType.one } + }, + "the chain must stay usable after a cancelled waiter" + ) + await transport.finishStreams() + await connection.shutdown() + } + + /// The cancellation policy is uniform: a send reaching a cancelled task — + /// including one already-cancelled at entry — is rejected before encrypting, + /// so no nonce is burned for a frame nothing will carry. + @Test("a send already-cancelled at entry is rejected without burning a nonce or wedging the chain") + func cancelledEntryIsRejectedWithoutBurningNonceOrWedgingChain() async throws { + let transport = MockTransport() + let connection = try await makeConnectionWithTransport(transport) + // start() returns before messageLoop installs the clock-sync sampler; + // wait for the handle first so cancel/join drains it deterministically. + #expect(await waitUntil { await connection.clockSyncTask != nil }, "clock-sync task handle must appear before cancel") + await connection.clockSyncTask?.cancel() + await connection.clockSyncTask?.value + #expect(await waitUntil { await !(connection.outboundInFlight) }) + + await transport.enableGoodbyeGate() + // Make the mock transport mirror the real one: a pre-cancelled sender that + // somehow reached the transport must also be rejected, not delivered. + await transport.setHonorCancellationSends(true) + let big = OutboundTestMessage( + type: OutboundTestMessageType.padded, + note: String(repeating: "e", count: NoiseChannel.maxSinglePayload + 2_000) + ) + let bigTask = Task { () -> Result in + do { + try await connection.send(clientMessage: big) + return .success(()) + } catch { + return .failure(error) + } + } + #expect(await waitUntil { await transport.isGoodbyeGateWaiting }) + + // The task cancels ITSELF before sending; with the uniform policy the + // send must fail without encrypting, and the chain must stay live. + let cancelledAtEntryTask = Task { () -> Result in + withUnsafeCurrentTask { $0?.cancel() } + do { + try await connection.send(clientMessage: OutboundTestMessage(type: OutboundTestMessageType.small, note: "pre-cancelled")) + return .success(()) + } catch { + return .failure(error) + } + } + #expect(await waitUntil { await connection.outboundWaiters.count == 1 }, "the pre-cancelled sender must queue") + + // Release the parked first sender; the pre-cancelled sender wakes and must be rejected. + await transport.releaseGoodbyeGate() + let bigResult = await bigTask.value + #expect((try? bigResult.get()) != nil, "the fragmented message must still send") + let smallResult = await cancelledAtEntryTask.value + guard case .failure = smallResult else { + Issue.record("a send already-cancelled at entry must fail under the uniform policy") + return + } + #expect(await connection.outboundWaiters.isEmpty, "the rejected waiter must release the slot") + + // The cancelled message never encrypted, so it must not appear at the peer. + let server = try #require(await connectionReadbacks.server(for: transport)) + let wireTypes = await server.decryptedMessages.compactMap(typeOfDecryptedJSON) + #expect(!wireTypes.contains(OutboundTestMessageType.small), "a pre-cancelled send must not reach the wire") + + // The chain is still live after the rejection. + try await connection.send(clientMessage: OutboundTestMessage(type: OutboundTestMessageType.one, note: "after-reject")) + #expect( + await waitUntil(timeout: .seconds(3)) { + await server.decryptedMessages.contains { typeOfDecryptedJSON($0) == OutboundTestMessageType.one } + }, + "the chain must stay usable after a rejected pre-cancelled sender" + ) + await transport.finishStreams() + await connection.shutdown() + } + + /// A failed fragmented send is terminal (nonces burned): latch so queued and + /// later senders fail without encrypting, and nothing hangs. + @Test("outbound failure latches: queued and later senders send nothing after the error") + func outboundFailureLatchStopsFurtherSends() async throws { + let transport = MockTransport() + let connection = try await makeConnectionWithTransport(transport) + // start() returns before messageLoop installs the clock-sync sampler; + // wait for the handle first so cancel/join drains it deterministically. + #expect(await waitUntil { await connection.clockSyncTask != nil }, "clock-sync task handle must appear before cancel") + await connection.clockSyncTask?.cancel() + await connection.clockSyncTask?.value + #expect(await waitUntil { await !(connection.outboundInFlight) }) + + await transport.enableGoodbyeGate() + let big = OutboundTestMessage( + type: OutboundTestMessageType.padded, + note: String(repeating: "d", count: NoiseChannel.maxSinglePayload + 2_000) + ) + let bigTask = Task { () -> Result in + do { + try await connection.send(clientMessage: big) + return .success(()) + } catch { + return .failure(error) + } + } + #expect(await waitUntil { await transport.isGoodbyeGateWaiting }) + + let queuedTask = Task { () -> Result in + do { + try await connection.send(clientMessage: OutboundTestMessage(type: OutboundTestMessageType.small, note: "queued")) + return .success(()) + } catch { + return .failure(error) + } + } + #expect(await waitUntil { await connection.outboundWaiters.count == 1 }, "the queued sender must park") + + // Arm failure so the SECOND fragment (post-release) throws; the parked + // first fragment already passed the fail check, so it sends once. + let framesBeforeRelease = await transport.sentBinaryMessages.count + await transport.setShouldFailOnSend(true) + await transport.releaseGoodbyeGate() + + let bigResult = await bigTask.value + guard case .failure = bigResult else { + Issue.record("a send that burned nonces then failed must surface the error") + return + } + let queuedResult = await queuedTask.value + guard case .failure = queuedResult else { + Issue.record("a queued sender behind a failed send must fail (channel is crypto-dead)") + return + } + + // A LATER sender must also fail without encrypting or sending: nothing + // beyond the failure point may reach the wire. + let laterResult = await Task { + do { + try await connection.send(clientMessage: OutboundTestMessage(type: OutboundTestMessageType.one, note: "later")) + return Result.success(()) + } catch { + return .failure(error) + } + }.value + guard case .failure = laterResult else { + Issue.record("a later send after an outbound failure must fail (channel is crypto-dead)") + return + } + + #expect(await connection.outboundWaiters.isEmpty, "no sender may remain parked") + #expect(await connection.outboundFailed, "the outbound-failure state must be latched") + #expect(await transport.disconnectCalled, "a burned-nonce send must tear the session down") + + // Only the first fragment ever reached the wire (it parked BEFORE the + // fail arm and sent once on release); the failing second fragment and + // every queued/later send must not have produced a frame. + let sentCount = await transport.sentBinaryMessages.count + #expect( + sentCount == framesBeforeRelease + 1, + "only the first fragment may reach the wire after the failure point; got \(sentCount) frames" + ) + await transport.finishStreams() + await connection.shutdown() + } } // MARK: - Connection factories (shared by both suites) @@ -1576,6 +1934,35 @@ private func audioChunkFrame(index: Int = 0, baseTimestamp: Int64 = 1_000_000) - return frame } +/// Wire discriminators for the outbound-serialization tests. The `padded` type +/// carries a payload large enough to force fragmentation, so the connection +/// exercises the multi-frame path with one wait per fragment. +enum OutboundTestMessageType: String, Codable { + case padded + case small + case one + case two + + static let all: [OutboundTestMessageType] = [.padded, .small, .one, .two] +} + +/// A minimal custom control message the connection's `send(clientMessage:)` +/// accepts and encrypts directly (not a `client/state` snapshot, so the state +/// coalescing gate cannot serialize the test away). +struct OutboundTestMessage: Codable, Sendable { + let type: OutboundTestMessageType + let note: String +} + +/// Message type of a decrypted `[json type][payload]` plaintext, or nil for +/// non-JSON (audio) frames. +func typeOfDecryptedJSON(_ message: Data) -> OutboundTestMessageType? { + guard message.first == NoiseFrameType.json, + let decoded = try? JSONDecoder().decode(OutboundTestMessage.self, from: Data(message.dropFirst())) + else { return nil } + return decoded.type +} + /// Encode a `stream/start` carrying a player format. `codec` is a raw wire string /// (not `AudioCodec`) so callers can deliberately exercise the unsupported-codec /// path — see `streamStartUnknownCodec_emitsClientStateError`. diff --git a/Tests/SendspinKitTests/Client/StaticPairingTranscriptTests.swift b/Tests/SendspinKitTests/Client/StaticPairingTranscriptTests.swift index 0902eba..282f410 100644 --- a/Tests/SendspinKitTests/Client/StaticPairingTranscriptTests.swift +++ b/Tests/SendspinKitTests/Client/StaticPairingTranscriptTests.swift @@ -209,10 +209,13 @@ struct StaticPairingWindowTests { @Test("a rejected static activation cancels its attempt and allows a fresh activation") func rejectedStaticActivationCleansUpAttempt() async throws { - let session = try await makeStaticTestSession(attemptTimeout: .milliseconds(100)) + let session = try await makeStaticTestSession() try await session.client.openPairingWindow() try await activateStatic(session.server) _ = try await waitForStaticClientMessage(session.server, type: ClientPairInitMessage.typeString) + let connection = try #require(await MainActor.run { session.client.connection }) + let attemptTask = try #require(await connection.pairingAttemptTask) + #expect(!attemptTask.isCancelled) let runtime = try #require(await MainActor.run { session.client.pairingConfiguration?.runtime }) let pairingPsk = await runtime.snapshot().pairingPsk @@ -227,7 +230,13 @@ struct StaticPairingWindowTests { try await activateStatic(session.server) let firstAbort = try await waitForStaticClientMessage(session.server, type: PairAbortMessage.typeString) #expect(try JSONDecoder().decode(PairAbortMessage.self, from: firstAbort).payload.reason == .methodNotSupported) - try await Task.sleep(for: .milliseconds(150)) + #expect(attemptTask.isCancelled) + #expect(await connection.pairingAttemptTask == nil) + // Clean up even when the cancellation assertion fails. + attemptTask.cancel() + if case .timedOut = await observeTask(attemptTask, timeout: .seconds(2)) { + Issue.record("cancelled pairing timer did not finish") + } #expect(await session.server.clientJSONMessages(ofType: PairAbortMessage.typeString).count == 1) await runtime.update(PairingManagementConfiguration( @@ -240,7 +249,15 @@ struct StaticPairingWindowTests { )) try await session.client.openPairingWindow() try await activateStatic(session.server) - _ = try await waitForStaticClientMessage(session.server, type: ClientPairInitMessage.typeString) + let refreshedInit = try await waitForStaticClientMessage( + session.server, + type: ClientPairInitMessage.typeString, + count: 2 + ) + let refreshedPairInit = try JSONDecoder().decode(ClientPairInitMessage.self, from: refreshedInit) + #expect(refreshedPairInit.payload.commitB == nil) + let replacementTask = try #require(await connection.pairingAttemptTask) + #expect(!replacementTask.isCancelled) #expect(await MainActor.run { session.client.connectionState == .connected }) await session.client.disconnect() } @@ -440,6 +457,7 @@ struct PairingCancellationTests { @Test("static attempt timeout sends attempt_timeout") func attemptTimeout() async throws { let session = try await makeStaticTestSession(attemptTimeout: .milliseconds(100)) + await session.server.transport.setHonorCancellationSends(true) try await activateStatic(session.server) try await session.client.openPairingWindow() _ = try await waitForStaticClientMessage(session.server, type: ClientPairInitMessage.typeString) diff --git a/Tests/SendspinKitTests/Helpers/MockTransport.swift b/Tests/SendspinKitTests/Helpers/MockTransport.swift index 0eb3d3b..8ff8fe7 100644 --- a/Tests/SendspinKitTests/Helpers/MockTransport.swift +++ b/Tests/SendspinKitTests/Helpers/MockTransport.swift @@ -43,6 +43,11 @@ actor MockTransport: ClientDialingTransport { private var goodbyeGateEnabled = false private var goodbyeGateContinuation: CheckedContinuation? + /// Opt-in mirror of `NWWebSocketTransport`: senders that reach the + /// transport already-cancelled are rejected. Defaults off so existing + /// tests that queue non-cancelled senders stay unchanged. + private var honorCancellationSends = false + // MARK: - SendspinTransport conformance func connect() async throws {} @@ -55,6 +60,9 @@ actor MockTransport: ClientDialingTransport { if shouldFailOnSend { throw MockTransportError.simulatedFailure } + if honorCancellationSends, Task.isCancelled { + throw CancellationError() + } await parkNextOutboundFrameIfArmed() sentTextMessages.append(Data(text.utf8)) outbox.yield(.text(text)) @@ -64,6 +72,9 @@ actor MockTransport: ClientDialingTransport { if shouldFailOnSend { throw MockTransportError.simulatedFailure } + if honorCancellationSends, Task.isCancelled { + throw CancellationError() + } await parkNextOutboundFrameIfArmed() sentBinaryMessages.append(data) outbox.yield(.binary(data)) @@ -154,6 +165,11 @@ actor MockTransport: ClientDialingTransport { goodbyeGateEnabled = true } + /// Opt in to rejecting sends that arrive here already-cancelled. + func setHonorCancellationSends(_ value: Bool) { + honorCancellationSends = value + } + /// Whether an outbound frame is currently parked on the gate. var isGoodbyeGateWaiting: Bool { goodbyeGateContinuation != nil diff --git a/Tests/SendspinKitTests/Models/AudioFormatSpecTests.swift b/Tests/SendspinKitTests/Models/AudioFormatSpecTests.swift index 112ac10..8726be3 100644 --- a/Tests/SendspinKitTests/Models/AudioFormatSpecTests.swift +++ b/Tests/SendspinKitTests/Models/AudioFormatSpecTests.swift @@ -715,9 +715,9 @@ private actor BlockingStartAudioOutputPlatformMonitor: AudioOutputPlatformMonito private(set) var stopCount = 0 private(set) var activeListenerCount = 0 - func startMonitoring() -> AsyncStream { + func startMonitoring() async -> AsyncStream { startCount += 1 - gate.enterAndWaitForCancellationThenRelease() + await gate.enterAndWaitForCancellationThenRelease() activeListenerCount = 1 return AsyncStream { _ in } } @@ -734,9 +734,7 @@ private actor BlockingStartAudioOutputPlatformMonitor: AudioOutputPlatformMonito } nonisolated func waitUntilStartWasCancelled() async { - while !gate.hasObservedCancellation { - await Task.yield() - } + await gate.waitForCancellationObserved() } nonisolated func releaseStart() { @@ -745,38 +743,71 @@ private actor BlockingStartAudioOutputPlatformMonitor: AudioOutputPlatformMonito } private final class BlockingStartGate: @unchecked Sendable { - private let condition = NSCondition() + private let lock = NSLock() private var entered = false private var cancellationObserved = false private var released = false - - func enterAndWaitForCancellationThenRelease() { - condition.lock() - entered = true - condition.broadcast() - while !released { - if Task.isCancelled { - cancellationObserved = true - condition.broadcast() + private var waiter: CheckedContinuation? + private var cancelWaiters: [CheckedContinuation] = [] + + func enterAndWaitForCancellationThenRelease() async { + // Suspends until releaseStart(); cancellation is observed but does not + // resume, preserving the start-after-cancel race. + await withTaskCancellationHandler { + await withCheckedContinuation { continuation in + lock.lock() + entered = true + if released { + lock.unlock() + continuation.resume() + return + } + if Task.isCancelled { + cancellationObserved = true + } + waiter = continuation + lock.unlock() + } + } onCancel: { + lock.lock() + cancellationObserved = true + let observers = cancelWaiters + cancelWaiters.removeAll() + lock.unlock() + for observer in observers { + observer.resume() } - _ = condition.wait(until: Date(timeIntervalSinceNow: 0.01)) } - condition.unlock() } var hasEntered: Bool { - condition.withLock { entered } + lock.withLock { entered } } var hasObservedCancellation: Bool { - condition.withLock { cancellationObserved } + lock.withLock { cancellationObserved } } func release() { - condition.lock() + lock.lock() released = true - condition.broadcast() - condition.unlock() + let waiter = waiter + self.waiter = nil + lock.unlock() + waiter?.resume() + } + + func waitForCancellationObserved() async { + await withCheckedContinuation { continuation in + lock.lock() + if cancellationObserved { + lock.unlock() + continuation.resume() + return + } + cancelWaiters.append(continuation) + lock.unlock() + } } } diff --git a/Tests/SendspinKitTests/Transport/NoiseSessionEstablisherTests.swift b/Tests/SendspinKitTests/Transport/NoiseSessionEstablisherTests.swift index e559848..47c4e7f 100644 --- a/Tests/SendspinKitTests/Transport/NoiseSessionEstablisherTests.swift +++ b/Tests/SendspinKitTests/Transport/NoiseSessionEstablisherTests.swift @@ -212,6 +212,43 @@ struct NoiseSessionEstablisherTests { } } + /// Cancelling the caller must not strand the parked `FrameInbox` pull: only + /// `disconnect()` releases it, so a regression that ignores cancellation + /// parks until the phase watchdog fires. + @Test("Cancelling after client/init returns promptly and disconnects the transport") + func cancellationAfterClientInitReturnsPromptly() async { + let transport = MockTransport() + let clientSide = establishmentTask( + transport: transport, + phaseTimeout: NoiseSessionEstablisher.defaultPhaseTimeout + ) + _ = await transport.nextSentFrame() // client/init went out; the server never replies + clientSide.cancel() + + let outcome = TestBox?>(nil) + let observer = Task { + do { + _ = try await clientSide.value + await outcome.set(.success(())) + } catch { + await outcome.set(.failure(error)) + } + } + let returned = await waitUntil(timeout: .seconds(2)) { await outcome.value != nil } + let disconnectedOnCancellation = await transport.disconnectCalled + + // Release and join even when the cancellation assertion fails. + await transport.disconnect() + _ = try? await clientSide.value + await observer.value + + #expect(returned, "a cancelled establishment must return promptly, not park until the phase watchdog") + #expect( + disconnectedOnCancellation, + "cancellation must disconnect the transport, which is what releases the parked pull" + ) + } + @Test("client/init carries the identity, core version, and chosen suite") func clientInitContents() async throws { let transport = MockTransport()