Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
4 changes: 4 additions & 0 deletions .github/workflows/tests.yml
Original file line number Diff line number Diff line change
Expand Up @@ -48,6 +48,7 @@ jobs:
test:
name: Test
runs-on: macos-26
timeout-minutes: 10

steps:
- name: Checkout code
Expand All @@ -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.
Expand Down
4 changes: 3 additions & 1 deletion Sources/SendspinKit/Audio/AudioEngine.swift
Original file line number Diff line number Diff line change
Expand Up @@ -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 {
Expand Down
2 changes: 1 addition & 1 deletion Sources/SendspinKit/Audio/AudioOutput.swift
Original file line number Diff line number Diff line change
Expand Up @@ -60,7 +60,7 @@ protocol AudioOutputPlatformMonitoring: Actor {
nonisolated var requiresActiveAudioSession: Bool { get }

/// Start a fresh, single-consumer observation stream.
func startMonitoring() -> AsyncStream<AudioOutputPlatformObservation>
func startMonitoring() async -> AsyncStream<AudioOutputPlatformObservation>

/// Remove all listeners and finish the current observation stream.
func stopMonitoring()
Expand Down
21 changes: 20 additions & 1 deletion Sources/SendspinKit/Client/SendspinClient+MultiServer.swift
Original file line number Diff line number Diff line change
Expand Up @@ -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()
Expand Down
74 changes: 53 additions & 21 deletions Sources/SendspinKit/Client/SendspinClient.swift
Original file line number Diff line number Diff line change
Expand Up @@ -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()
Expand Down Expand Up @@ -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()
Expand All @@ -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()
Expand All @@ -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 {
Expand Down Expand Up @@ -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.
Expand All @@ -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()
Expand All @@ -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
Expand Down Expand Up @@ -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
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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)))
}
Expand Down
73 changes: 56 additions & 17 deletions Sources/SendspinKit/Client/SendspinConnection+Outbound.swift
Original file line number Diff line number Diff line change
Expand Up @@ -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)
Expand All @@ -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
}
}
Expand All @@ -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
Expand All @@ -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
}
Expand Down
9 changes: 9 additions & 0 deletions Sources/SendspinKit/Client/SendspinConnection.swift
Original file line number Diff line number Diff line change
Expand Up @@ -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<Void, Never>] = []
/// 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
Expand Down
Loading
Loading