diff --git a/Packages/Shared/CmuxIrxTransport/Sources/CmuxIrxTransport/IrxEndpoint.swift b/Packages/Shared/CmuxIrxTransport/Sources/CmuxIrxTransport/IrxEndpoint.swift index ba129ebd0acd..369b0f6571fd 100644 --- a/Packages/Shared/CmuxIrxTransport/Sources/CmuxIrxTransport/IrxEndpoint.swift +++ b/Packages/Shared/CmuxIrxTransport/Sources/CmuxIrxTransport/IrxEndpoint.swift @@ -187,7 +187,17 @@ public actor IrxEndpointSupervisor { /// Queues a serialized installation and retries local failures independently /// of credential minting. Native installation emits its own outcome event. public func rotateCredentials(_ credentials: [IrxRelayCredential]) async { - guard !deactivated, configuration.pathMode != .directOnly else { return } + guard !deactivated, configuration.pathMode != .directOnly else { + journal.record("endpoint", "relay-rotation-skipped", [ + "reason": deactivated ? "deactivated" : "direct-only", + ]) + return + } + if relayInstaller == nil { + // The credentials are retained for the next bind, but nothing + // reaches the live driver; say so instead of rotating silently. + journal.record("endpoint", "relay-rotation-deferred", ["reason": "no-installer"]) + } desiredRelayCredentials = credentials desiredRelayOwnership = nil await relayInstaller?.replace(with: credentials) diff --git a/Packages/Shared/CmuxIrxTransport/Sources/CmuxIrxTransport/IrxRelayCredentialInstaller.swift b/Packages/Shared/CmuxIrxTransport/Sources/CmuxIrxTransport/IrxRelayCredentialInstaller.swift index c5ac62f4b2b0..b07d10c1e1f0 100644 --- a/Packages/Shared/CmuxIrxTransport/Sources/CmuxIrxTransport/IrxRelayCredentialInstaller.swift +++ b/Packages/Shared/CmuxIrxTransport/Sources/CmuxIrxTransport/IrxRelayCredentialInstaller.swift @@ -76,17 +76,27 @@ actor IrxRelayCredentialInstaller { var failures = 0 while !Task.isCancelled, !stopped, taskID == id { let observedRevision = revision + let unusable = desired.values.filter { !$0.isUsable(at: now()) }.count let pending = desired.values.filter { $0.isUsable(at: now()) && installed[$0.relayURL] != $0 }.sorted { $0.relayURL < $1.relayURL } + if unusable > 0 { + // Expired or not-yet-valid credentials silently shrink the + // fleet; make the drop visible next to the install outcomes. + journal.record("endpoint", "relay-credential-unusable", ["count": String(unusable)]) + } guard !pending.isEmpty else { return } var failed = false for credential in pending { guard !Task.isCancelled, !stopped, taskID == id else { return } guard desired[credential.relayURL] == credential, credential.isUsable(at: now()) else { continue } + // A started event with no matching outcome is the signature of + // a native install hung inside the driver. + journal.record("endpoint", "relay-credential-install-started", ["relay": credential.relayURL]) do { guard try await install(credential, ownership: ownership) else { + journal.record("endpoint", "relay-credential-install-superseded", ["relay": credential.relayURL]) if revision != observedRevision { break } return } diff --git a/Packages/Shared/CmuxIrxTransport/Sources/CmuxIrxTransport/V2/V2ControlDependencies.swift b/Packages/Shared/CmuxIrxTransport/Sources/CmuxIrxTransport/V2/V2ControlDependencies.swift index a592f0109ccb..001eb82ddf4b 100644 --- a/Packages/Shared/CmuxIrxTransport/Sources/CmuxIrxTransport/V2/V2ControlDependencies.swift +++ b/Packages/Shared/CmuxIrxTransport/Sources/CmuxIrxTransport/V2/V2ControlDependencies.swift @@ -9,6 +9,7 @@ public struct V2ControlDependencies: Sendable { let now: @Sendable () -> Date let sleep: @Sendable (TimeInterval) async throws -> Void let jitter: @Sendable () -> Double + let journal: IrxJournal? /// Injects real network/auth/signing effects or deterministic test replacements. /// - Parameters: @@ -19,6 +20,8 @@ public struct V2ControlDependencies: Sendable { /// - now: Wall clock for token expiry and server proofs. /// - sleep: Cancellable clock delay used only for deadlines and renewal/backoff. /// - jitter: A value in `0...1` to spread retries across clients. + /// - journal: Receives credential-lifecycle events so a silent renewal + /// stall is diagnosable from retained logs. Never carries token data. public init( connect: @escaping @Sendable (URLRequest) async throws -> any V2ControlSocket, http: @escaping @Sendable (URLRequest) async throws -> V2HTTPResponse, @@ -28,7 +31,8 @@ public struct V2ControlDependencies: Sendable { sleep: @escaping @Sendable (TimeInterval) async throws -> Void = { seconds in try await Task.sleep(for: .seconds(max(0, seconds))) }, - jitter: @escaping @Sendable () -> Double = { Double.random(in: 0...1) } + jitter: @escaping @Sendable () -> Double = { Double.random(in: 0...1) }, + journal: IrxJournal? = nil ) { self.connect = connect self.http = http @@ -37,5 +41,6 @@ public struct V2ControlDependencies: Sendable { self.now = now self.sleep = sleep self.jitter = jitter + self.journal = journal } } diff --git a/Packages/Shared/CmuxIrxTransport/Sources/CmuxIrxTransport/V2/V2ControlService+Connection.swift b/Packages/Shared/CmuxIrxTransport/Sources/CmuxIrxTransport/V2/V2ControlService+Connection.swift index 9bddabe77b86..e3d790adc750 100644 --- a/Packages/Shared/CmuxIrxTransport/Sources/CmuxIrxTransport/V2/V2ControlService+Connection.swift +++ b/Packages/Shared/CmuxIrxTransport/Sources/CmuxIrxTransport/V2/V2ControlService+Connection.swift @@ -23,6 +23,7 @@ extension V2ControlService { try await open(run: run) } catch { guard permitsHTTPRecovery(error) else { throw error } + journal("socket-open-failed", ["failure": mapFailure(error).diagnosticCode, "recovery": "http"]) let old = socket socket = nil socketID = nil @@ -32,6 +33,7 @@ extension V2ControlService { await old?.close() try assertCurrent(run) try await openHTTP(run: run) + journal("http-mode-entered", [:]) // Retry the preferred push channel while HTTP serves requests. try await dependencies.sleep(60) continue @@ -55,6 +57,7 @@ extension V2ControlService { await old?.close() guard runID == run else { return } if terminal(mapped) { + journal("run-stopped-terminal", ["failure": mapped.diagnosticCode]) status = .stopped runID = nil runTask = nil @@ -65,6 +68,11 @@ extension V2ControlService { status = .backingOff publish() let seconds = retryDelay(mapped, attempt: attempt) + journal("run-backing-off", [ + "failure": mapped.diagnosticCode, + "attempt": String(attempt), + "delay_s": String(Int(seconds)), + ]) attempt += 1 do { try await dependencies.sleep(seconds) } catch { return } @@ -120,6 +128,7 @@ extension V2ControlService { try await persist(run: run) status = .ready failure = nil + journal("session-ready", ["http_mode": String(httpMode)]) publish() requestDirectoryRefresh(run: run) scheduleMaintenance(run: run) @@ -231,6 +240,9 @@ extension V2ControlService { func socketFailed(_ error: any Error, run: UUID, connection: UUID) async { guard runID == run, socketID == connection else { return } + // Journaled from the service, not the snapshot consumer, so a stalled + // consumer cannot hide the socket's death from retained logs. + journal("socket-failed", ["failure": mapFailure(error).diagnosticCode]) failure = mapFailure(error) if case .socketClosed(1008, let reason) = failure, ["device_revoked", "team_access_revoked"].contains(reason ?? "") { @@ -284,13 +296,17 @@ extension V2ControlService { return false } + /// Only a deliberate revocation of this device's authority stops the run. + /// Every other failure, including identity or wire-shape rejections that + /// used to be permanent, retries at a bounded cadence: a stopped service + /// is invisible until relaunch, and a wedged Mac must keep announcing + /// itself to the backend rather than go silent. func terminal(_ error: V2ControlFailure) -> Bool { switch error { - case .scopeMismatch, .persistenceFailed, .capacityExceeded, .invalidWireData: return true - case .socketClosed(let code, _): return code == 1008 || code == 1009 + case .socketClosed(let code, let reason): + return code == 1008 && ["device_revoked", "team_access_revoked"].contains(reason ?? "") case .server(let response): - return [.teamAccessRevoked, .deviceRevoked, .identityMismatch, .environmentMismatch, .keyReplacementRequired, .endpointAlreadyOwned, .invalidDeviceProof].contains(response.code) - case .http(let status, _): return [400, 403, 404, 405, 413, 415].contains(status) + return [.teamAccessRevoked, .deviceRevoked].contains(response.code) default: return false } } @@ -305,6 +321,8 @@ extension V2ControlService { if let until = [cooldowns["session.open.v1"], cooldowns["session.open"]].compactMap({ $0 }).max() { delay = max(delay, until.timeIntervalSince(dependencies.now())) } - return delay + // No reconnect ever waits longer than the app-wide backoff ceiling, + // whatever a server retry-after or an accumulated cooldown says. + return min(delay, V2ControlService.maximumBackoff) } } diff --git a/Packages/Shared/CmuxIrxTransport/Sources/CmuxIrxTransport/V2/V2ControlService+HTTP.swift b/Packages/Shared/CmuxIrxTransport/Sources/CmuxIrxTransport/V2/V2ControlService+HTTP.swift index 572b32fe6f01..8068310a6697 100644 --- a/Packages/Shared/CmuxIrxTransport/Sources/CmuxIrxTransport/V2/V2ControlService+HTTP.swift +++ b/Packages/Shared/CmuxIrxTransport/Sources/CmuxIrxTransport/V2/V2ControlService+HTTP.swift @@ -63,7 +63,8 @@ extension V2ControlService { guard (200..<300).contains(response.status) else { let error = V2ControlFailure.http(status: response.status, retryAfter: response.retryAfter) if response.status == 429 { - cooldowns[operation(schema)] = dependencies.now().addingTimeInterval(max(1, response.retryAfter ?? 60)) + cooldowns[operation(schema)] = dependencies.now().addingTimeInterval( + min(max(1, response.retryAfter ?? 60), Self.maximumBackoff)) } record(error, schema: schema) throw error diff --git a/Packages/Shared/CmuxIrxTransport/Sources/CmuxIrxTransport/V2/V2ControlService+Maintenance.swift b/Packages/Shared/CmuxIrxTransport/Sources/CmuxIrxTransport/V2/V2ControlService+Maintenance.swift index 9c620b6512c7..2386d8407639 100644 --- a/Packages/Shared/CmuxIrxTransport/Sources/CmuxIrxTransport/V2/V2ControlService+Maintenance.swift +++ b/Packages/Shared/CmuxIrxTransport/Sources/CmuxIrxTransport/V2/V2ControlService+Maintenance.swift @@ -31,7 +31,14 @@ extension V2ControlService { } func scheduleMaintenance(run: UUID) { - guard runID == run, status == .ready else { return } + guard runID == run, status == .ready else { + journal("maintenance-not-scheduled", [ + "reason": runID == run ? "status" : "run-superseded", + "status": String(describing: status), + ]) + return + } + journal("maintenance-scheduled", ["http_mode": String(httpMode)]) #if DEBUG if let interval = verificationRenewalInterval, nextVerificationRenewalAt == nil { nextVerificationRenewalAt = dependencies.now().addingTimeInterval(interval) @@ -53,9 +60,38 @@ extension V2ControlService { let verificationDue = TimeInterval.infinity #endif let next = min(ticketDue, relayDue, directoryDue, verificationDue) + journal("maintenance-planned", [ + "sleep_s": String(Int(max(0, next - now))), + "ticket_in_s": String(Int(ticketDue - now)), + "relay_in_s": String(Int(relayDue - now)), + "directory_in_s": String(Int(directoryDue - now)), + // A renewal pushed past its refresh time by a cooldown is + // otherwise invisible; name the deferred schemas explicitly. + "deferred": deferredSchemas(now: now).joined(separator: ","), + ]) do { try await dependencies.sleep(max(0, next - now)) } - catch { return } - guard runID == run, !Task.isCancelled, (socket != nil || httpMode) else { return } + catch { + journal("maintenance-exited", ["reason": "sleep-cancelled"]) + return + } + guard runID == run, !Task.isCancelled else { + journal("maintenance-exited", ["reason": Task.isCancelled ? "cancelled" : "run-superseded"]) + return + } + guard socket != nil || httpMode else { + // The socket died in the window between renewal sleep and + // wake. Exiting here used to end renewals for good while + // status stayed .ready; instead hold at a bounded cadence + // until either the transport returns or the reconnect owner + // moves status off .ready and this loop ends normally. + journal("maintenance-idle", ["reason": "no-transport", "status": String(describing: status)]) + do { try await dependencies.sleep(30) } + catch { + journal("maintenance-exited", ["reason": "sleep-cancelled"]) + return + } + continue + } let deadline = dependencies.now().timeIntervalSince1970 + 0.1 #if DEBUG let forceVerification = verificationDue <= deadline @@ -87,15 +123,33 @@ extension V2ControlService { } } } + journal("maintenance-exited", [ + "reason": Task.isCancelled ? "cancelled" : (runID == run ? "status" : "run-superseded"), + "status": String(describing: status), + ]) } private func due(_ timestamp: Int?, schema: String, now: TimeInterval) -> TimeInterval { max(Double(timestamp ?? Int(now)), cooldowns[schema]?.timeIntervalSince1970 ?? 0, cooldowns[operation(schema)]?.timeIntervalSince1970 ?? 0) } + /// Schemas whose next run is later than their own refresh time because a + /// schema or operation cooldown dominates. + private func deferredSchemas(now: TimeInterval) -> [String] { + let wanted: [(String, Int?)] = [ + ("ticket.request.v1", cache.ticket?.refreshAfter), + ("relay.request.v1", cache.relayCredentials.map(\.refreshAfter).min()), + ("directory.request.v1", cache.directory.map { $0.permissionExpiresAt - 300 }), + ] + return wanted.compactMap { schema, refreshAfter in + due(refreshAfter, schema: schema, now: now) > Double(refreshAfter ?? Int(now)) ? schema : nil + } + } + private func maintenanceFailed(_ error: any Error, schema: String, run: UUID) { guard runID == run, !Task.isCancelled else { return } let mapped = mapFailure(error) + journal("refresh-failed", ["schema": schema, "failure": mapped.diagnosticCode]) // The receive owner already records server cooldowns. Do not count one // rejected schema twice or overwrite its long backoff with a short one. if case .server = mapped {} else { record(mapped, schema: schema) } diff --git a/Packages/Shared/CmuxIrxTransport/Sources/CmuxIrxTransport/V2/V2ControlService+Operations.swift b/Packages/Shared/CmuxIrxTransport/Sources/CmuxIrxTransport/V2/V2ControlService+Operations.swift index 7f697d9b4897..406b593df4d7 100644 --- a/Packages/Shared/CmuxIrxTransport/Sources/CmuxIrxTransport/V2/V2ControlService+Operations.swift +++ b/Packages/Shared/CmuxIrxTransport/Sources/CmuxIrxTransport/V2/V2ControlService+Operations.swift @@ -34,6 +34,12 @@ extension V2ControlService { try assertCurrent(run) cache.ticket = response.ticket failure = nil + let issued = Int(dependencies.now().timeIntervalSince1970) + journal("refresh-succeeded", [ + "schema": "ticket.request.v1", + "refresh_after_in_s": String(response.ticket.refreshAfter - issued), + "expires_in_s": String(response.ticket.expiresAt - issued), + ]) try await persist(run: run) return response.ticket } @@ -80,6 +86,13 @@ extension V2ControlService { try assertCurrent(run) cache.relayCredentials = response.credentials failure = nil + let issued = Int(dependencies.now().timeIntervalSince1970) + journal("refresh-succeeded", [ + "schema": "relay.request.v1", + "count": String(response.credentials.count), + "refresh_after_in_s": String((response.credentials.map(\.refreshAfter).min() ?? issued) - issued), + "expires_in_s": String((response.credentials.map(\.expiresAt).min() ?? issued) - issued), + ]) try await persist(run: run) return response.credentials } @@ -147,6 +160,12 @@ extension V2ControlService { ) cache.directory = directory failure = nil + journal("refresh-succeeded", [ + "schema": "directory.request.v1", + "revision": String(directory.revision), + "bindings": String(directory.devices.count), + "expires_in_s": String(directory.permissionExpiresAt - Int(dependencies.now().timeIntervalSince1970)), + ]) try await persist(run: run) return directory } diff --git a/Packages/Shared/CmuxIrxTransport/Sources/CmuxIrxTransport/V2/V2ControlService.swift b/Packages/Shared/CmuxIrxTransport/Sources/CmuxIrxTransport/V2/V2ControlService.swift index 4ab3e1b42eef..405dbca8bacd 100644 --- a/Packages/Shared/CmuxIrxTransport/Sources/CmuxIrxTransport/V2/V2ControlService.swift +++ b/Packages/Shared/CmuxIrxTransport/Sources/CmuxIrxTransport/V2/V2ControlService.swift @@ -167,7 +167,13 @@ public actor V2ControlService { try assertCurrent(run) let state = cache do { try await store.save(state) } - catch { throw V2ControlFailure.persistenceFailed } + catch { + // A failed save also skips publish, so downstream consumers never + // hear about state they would lose on relaunch. Without this event + // that outcome is indistinguishable from a renewal that never ran. + journal("persist-failed", ["error": String(describing: type(of: error))]) + throw V2ControlFailure.persistenceFailed + } try assertCurrent(run) publish() } @@ -305,19 +311,34 @@ public actor V2ControlService { func operation(_ schema: String) -> String { schema.split(separator: ".").dropLast().joined(separator: ".") } + /// The longest any cooldown or reconnect backoff may run, app-wide. A + /// credential or connection attempt deferred for hours is + /// indistinguishable from a dead Mac; half an hour bounds the damage of + /// any single misclassified failure while staying far under rate limits. + static let maximumBackoff: TimeInterval = 30 * 60 + func record(_ error: V2ControlFailure, schema: String) { failure = error if case .server(let response) = error { if response.code == .rateLimited { - cooldowns[operation(schema)] = dependencies.now().addingTimeInterval(max(1, Double(response.retryAfterMS ?? 60_000) / 1000)) + let delay = min(max(1, Double(response.retryAfterMS ?? 60_000) / 1000), Self.maximumBackoff) + cooldowns[operation(schema)] = dependencies.now().addingTimeInterval(delay) + journal("cooldown-set", ["schema": schema, "source": "rate_limited", "delay_s": String(Int(delay))]) } else if response.code == .clientUpgradeRequired { let attempt = retiredAttempts[schema, default: 0] - let delays: [TimeInterval] = [3600, 6 * 3600, 24 * 3600] - let delay = delays[min(attempt, delays.count - 1)] * (1 + 0.1 * dependencies.jitter()) + let delays: [TimeInterval] = [600, 1800, 1800] + let delay = min(delays[min(attempt, delays.count - 1)] * (1 + 0.1 * dependencies.jitter()), Self.maximumBackoff) retiredAttempts[schema] = attempt + 1 cooldowns[schema] = dependencies.now().addingTimeInterval(delay) + journal("cooldown-set", ["schema": schema, "source": "upgrade_required", "delay_s": String(Int(delay))]) } } publish() } + + /// Records one credential-lifecycle event. Attributes are schema names, + /// failure codes, and durations only, never tokens or credential bodies. + func journal(_ event: String, _ attributes: [String: String] = [:]) { + dependencies.journal?.record("v2-control", event, attributes) + } } diff --git a/Packages/Shared/CmuxIrxTransport/Tests/CmuxIrxTransportTests/V2/V2ControlServiceTests.swift b/Packages/Shared/CmuxIrxTransport/Tests/CmuxIrxTransportTests/V2/V2ControlServiceTests.swift index 5d54aa3ea25b..8e7f7f703c7c 100644 --- a/Packages/Shared/CmuxIrxTransport/Tests/CmuxIrxTransportTests/V2/V2ControlServiceTests.swift +++ b/Packages/Shared/CmuxIrxTransport/Tests/CmuxIrxTransportTests/V2/V2ControlServiceTests.swift @@ -14,15 +14,29 @@ import Testing ) } - private func service(backend: V2TestBackend, store: V2TestStateStore = V2TestStateStore()) throws -> V2ControlService { + private func service(backend: V2TestBackend, store: V2TestStateStore = V2TestStateStore(), journal: IrxJournal? = nil) throws -> V2ControlService { let fixedNow = now return V2ControlService( configuration: try V2ControlConfiguration(baseURL: URL(string: "https://control.example.com")!, device: device()), - dependencies: V2ControlDependencies(connect: { try await backend.connect($0) }, http: { try await backend.http($0) }, stackAccessToken: { _ in "existing-stack-session" }, sign: { _ in Data(repeating: 1, count: 64) }, now: { Date(timeIntervalSince1970: Double(fixedNow)) }, jitter: { 0.5 }), + dependencies: V2ControlDependencies(connect: { try await backend.connect($0) }, http: { try await backend.http($0) }, stackAccessToken: { _ in "existing-stack-session" }, sign: { _ in Data(repeating: 1, count: 64) }, now: { Date(timeIntervalSince1970: Double(fixedNow)) }, jitter: { 0.5 }, journal: journal), store: store ) } + private func events(_ journal: IrxJournal, _ event: String) -> [IrxJournalEvent] { + journal.tail(IrxJournal.ringCapacity).filter { $0.component == "v2-control" && $0.event == event } + } + + /// Journal writes trail the snapshot the test observed, so poll briefly. + private func journaled(_ journal: IrxJournal, _ event: String) async throws -> [IrxJournalEvent] { + for _ in 0..<200 { + let found = events(journal, event) + if !found.isEmpty { return found } + try await Task.sleep(for: .milliseconds(5)) + } + return [] + } + private func ready(_ service: V2ControlService) async throws -> V2ControlSnapshot { let events = await service.events() for await snapshot in events { @@ -236,12 +250,30 @@ import Testing _ = service } - @Test func permanentClosesStopWhileCapacityClosesBackOff() async throws { + @Test func onlyRevocationStopsEverythingElseRetriesWithinTheBackoffCeiling() async throws { let service = try service(backend: V2TestBackend(now: now)) #expect(await service.terminal(.socketClosed(code: 1008, reason: "device_revoked"))) - #expect(await service.terminal(.socketClosed(code: 1009, reason: "payload_too_large"))) + #expect(await service.terminal(.socketClosed(code: 1008, reason: "team_access_revoked"))) + // Policy closes without a revocation reason, wire-shape and + // persistence failures all retry now; a stopped service is silent + // until relaunch, which is how a wedged Mac disappears for hours. + #expect(await service.terminal(.socketClosed(code: 1008, reason: "policy_violation")) == false) + #expect(await service.terminal(.socketClosed(code: 1009, reason: "payload_too_large")) == false) + #expect(await service.terminal(.invalidWireData) == false) + #expect(await service.terminal(.persistenceFailed) == false) + #expect(await service.terminal(.capacityExceeded) == false) + #expect(await service.terminal(.scopeMismatch) == false) + #expect(await service.terminal(.http(status: 403, retryAfter: nil)) == false) #expect(await service.retryDelay(.socketClosed(code: 1013, reason: "slow_consumer"), attempt: 0) >= 60) - #expect(await service.terminal(.socketClosed(code: 1011, reason: "transport_error")) == false) + // A server retry-after or accumulated cooldown cannot push any + // reconnect past the app-wide 30-minute ceiling. + let capped = await service.retryDelay(.http(status: 429, retryAfter: 6 * 3600), attempt: 6) + #expect(capped <= 1800) + let cooldownCapped = await service.retryDelay( + .cooldown(schemaID: "session.open.v1", until: Date(timeIntervalSince1970: Double(now) + 24 * 3600)), + attempt: 0 + ) + #expect(cooldownCapped <= 1800) } @Test func paginationPinsRevisionAndRestartsAfterConcurrentChange() async throws { @@ -362,8 +394,8 @@ import Testing catch V2ControlFailure.server(let error) { #expect(error.code == .clientUpgradeRequired) } do { _ = try await service.refreshRelayCredentials(); Issue.record("Expected long cooldown") } catch V2ControlFailure.cooldown(_, let until) { - #expect(until.timeIntervalSince1970 - Double(now) >= 3600) - #expect(until.timeIntervalSince1970 - Double(now) <= 3960) + #expect(until.timeIntervalSince1970 - Double(now) >= 300) + #expect(until.timeIntervalSince1970 - Double(now) <= 1800) } await socket.rejectRelay(nil) await service.explicitRetry(schemaID: "relay.request.v1") @@ -371,4 +403,53 @@ import Testing #expect(await backend.sockets.count == 1) await service.stop() } + + @Test func credentialLifecycleIsJournaledFromReadyThroughRenewalAndShutdown() async throws { + let journal = IrxJournal(subsystem: "com.cmux.test", category: "v2-journal-test") + let backend = V2TestBackend(now: now) + let service = try service(backend: backend, journal: journal) + await service.start() + _ = try await ready(service) + #expect(try await journaled(journal, "session-ready").first?.attributes["http_mode"] == "false") + #expect(try await !journaled(journal, "maintenance-scheduled").isEmpty) + _ = try await service.refreshRelayCredentials() + _ = try await service.refreshAPITicket() + let schemas = events(journal, "refresh-succeeded").compactMap { $0.attributes["schema"] } + #expect(schemas.contains("relay.request.v1")) + #expect(schemas.contains("ticket.request.v1")) + #expect(schemas.contains("directory.request.v1")) + await service.stop() + // Stopping cancels the renewal sleep; the loop must say why it left. + let exits = try await journaled(journal, "maintenance-exited").compactMap { $0.attributes["reason"] } + #expect(exits.contains("sleep-cancelled") || exits.contains("cancelled") || exits.contains("run-superseded")) + } + + @Test func serverCooldownsAreJournaledWithTheirSource() async throws { + let journal = IrxJournal(subsystem: "com.cmux.test", category: "v2-journal-cooldown-test") + let rateLimitedBackend = V2TestBackend(now: now) + let rateLimitedService = try service(backend: rateLimitedBackend, journal: journal) + await rateLimitedService.start() + _ = try await ready(rateLimitedService) + await rateLimitedBackend.currentSocket().rejectRelay(.rateLimited) + do { _ = try await rateLimitedService.refreshRelayCredentials(); Issue.record("Expected relay rate limit") } + catch V2ControlFailure.server(let error) { #expect(error.code == .rateLimited) } + await rateLimitedService.stop() + let rateLimited = events(journal, "cooldown-set").first { $0.attributes["source"] == "rate_limited" } + #expect(rateLimited?.attributes["schema"] == "relay.request.v1") + #expect(Int(rateLimited?.attributes["delay_s"] ?? "") ?? -1 >= 1) + + let retiredBackend = V2TestBackend(now: now) + let retiredService = try service(backend: retiredBackend, journal: journal) + await retiredService.start() + _ = try await ready(retiredService) + await retiredBackend.currentSocket().rejectRelay(.clientUpgradeRequired) + do { _ = try await retiredService.refreshRelayCredentials(); Issue.record("Expected retired schema") } + catch V2ControlFailure.server(let error) { #expect(error.code == .clientUpgradeRequired) } + await retiredService.stop() + let retired = events(journal, "cooldown-set").first { $0.attributes["source"] == "upgrade_required" } + #expect(retired?.attributes["schema"] == "relay.request.v1") + let retiredDelay = Int(retired?.attributes["delay_s"] ?? "") ?? 0 + #expect(retiredDelay >= 300) + #expect(retiredDelay <= 1800) + } } diff --git a/Sources/Mobile/MobileHostIrxRuntime.swift b/Sources/Mobile/MobileHostIrxRuntime.swift index aff85b13cef3..a2057c6d85f5 100644 --- a/Sources/Mobile/MobileHostIrxRuntime.swift +++ b/Sources/Mobile/MobileHostIrxRuntime.swift @@ -40,7 +40,9 @@ final class MobileHostIrxRuntime: MobileHostPairingRuntime { nonisolated static func activationRetryDelay(after error: any Error, failureCount: Int, jitterUnitInterval: Double) -> TimeInterval { let ladder = min(5 * pow(2, Double(min(max(failureCount, 0), 16))), maximumActivationRetryDelay) let floor = TimeInterval(max(0, (error as? any CmxRetryAfterProviding)?.retryAfterSeconds ?? 0)) - let base = max(ladder, floor) + // A server retry-after may stretch one wait but never past the + // app-wide backoff ceiling; a host that waits hours is a dead host. + let base = min(max(ladder, floor), 30 * 60) return base + min(max(jitterUnitInterval, 0), 1) * base * 0.25 } @@ -87,6 +89,7 @@ final class MobileHostIrxRuntime: MobileHostPairingRuntime { private(set) var generationToken = UUID() private(set) var activationTask: Task? private var controlTask: Task? + private var renewalWatchdogTask: Task? private var endpointTask: Task? private var endpointRefreshPending = false private var relayAddressWatch: WatchHandle? @@ -94,6 +97,7 @@ final class MobileHostIrxRuntime: MobileHostPairingRuntime { private var permissionExpiryTask: Task? private var acceptLoop: Task? private var lastLoggedControlState: String? + private var lastAppliedControlSequence: UInt64 = 0 private var admission: V2InboundAdmissionAuthority? /// Compatibility publication for older iOS dialects. It shares the v2 /// signing key but has its own filtered authority and broker lifecycle. @@ -346,6 +350,7 @@ final class MobileHostIrxRuntime: MobileHostPairingRuntime { deviceMetadataTask?.cancel(); deviceMetadataTask = nil activationTask?.cancel(); activationTask = nil controlTask?.cancel(); controlTask = nil + renewalWatchdogTask?.cancel(); renewalWatchdogTask = nil endpointTask?.cancel(); endpointTask = nil endpointRefreshPending = false relayAddressWatch = nil; relayAddressWatchGeneration = nil @@ -459,7 +464,8 @@ final class MobileHostIrxRuntime: MobileHostPairingRuntime { sign: { data in guard await auth.isAuthenticatedTeamScopeCurrent(scope) else { throw V2ControlFailure.scopeMismatch } return try key.sign(data) - }) + }, + journal: Self.journal) let service = V2ControlService(configuration: try .init(baseURL: configuration.baseURL, device: device), dependencies: dependencies, store: store) listenerState.preferredPort = preferredPort @@ -507,6 +513,7 @@ final class MobileHostIrxRuntime: MobileHostPairingRuntime { await self?.apply(snapshot, token: token) } } + startRenewalWatchdog(token: token, service: service) await service.start() guard isCurrent(token), !Task.isCancelled else { await service.stop(); throw V2ControlFailure.stopped } Self.journal.record("v2-host", "control-started", ["cached": String(restored != nil)]) @@ -535,6 +542,10 @@ final class MobileHostIrxRuntime: MobileHostPairingRuntime { private func apply(_ snapshot: V2ControlSnapshot, token: UUID) async { guard isCurrent(token), let admission else { return } + // Marked on exit, not entry: an apply stuck in one of its awaits must + // stay visible to the renewal watchdog as an unapplied sequence. + let appliedSequence = snapshot.sequence + defer { lastAppliedControlSequence = max(lastAppliedControlSequence, appliedSequence) } let status = String(describing: snapshot.status) let failure = snapshot.failure?.diagnosticCode ?? "none" let state = status + ":" + failure @@ -606,15 +617,73 @@ final class MobileHostIrxRuntime: MobileHostPairingRuntime { guard isCurrent(token) else { return } schedulePermissionExpiry(token: token) // Installing credentials does not replace the endpoint or its admitted sessions. - if previousCredentials != snapshot.cache.relayCredentials, let supervisor = endpointSupervisor { - await supervisor.rotateCredentials(Self.credentials(snapshot.cache)) - guard isCurrent(token) else { return } + if previousCredentials != snapshot.cache.relayCredentials { + let now = Int(Date().timeIntervalSince1970) + Self.journal.record("v2-host", "credentials-received", [ + "count": String(snapshot.cache.relayCredentials.count), + "expires_in_s": String((snapshot.cache.relayCredentials.map(\.expiresAt).max() ?? now) - now), + "supervisor": String(endpointSupervisor != nil), + ]) + if let supervisor = endpointSupervisor { + await supervisor.rotateCredentials(Self.credentials(snapshot.cache)) + guard isCurrent(token) else { return } + } } requestEndpointReady(token: token) if activeDeviceCapabilities != deviceCapabilities { updateDeviceHostingMetadata() } publishIrxSettingsUpdate() } + /// Journals loudly when credential renewal has silently stopped: either the + /// service holds stale credentials past their refresh time, or the service + /// has newer snapshots than this runtime ever applied. Reads the service + /// directly, so a stalled `apply` pipeline cannot hide itself. Records + /// evidence only; recovery stays with the existing owners. + private func startRenewalWatchdog(token: UUID, service: V2ControlService) { + renewalWatchdogTask?.cancel() + renewalWatchdogTask = Task { @MainActor [weak self] in + let interval: TimeInterval = 300 + var pendingSinceLastTick: UInt64? + while !Task.isCancelled { + do { try await Task.sleep(for: .seconds(interval)) } + catch { return } + guard let self, self.isCurrent(token) else { return } + let snapshot = await service.snapshot() + guard self.isCurrent(token), !Task.isCancelled else { return } + let now = Int(Date().timeIntervalSince1970) + var overdue: [String: String] = [:] + if !snapshot.cache.authorityRevoked, Self.pathMode != .directOnly, + let refreshAfter = snapshot.cache.relayCredentials.map(\.refreshAfter).min(), + now - refreshAfter > Int(interval) { + overdue["relay_overdue_s"] = String(now - refreshAfter) + let expires = snapshot.cache.relayCredentials.map(\.expiresAt).max() ?? now + overdue["relay_expires_in_s"] = String(expires - now) + } + if !snapshot.cache.authorityRevoked, let ticket = snapshot.cache.ticket, + now - ticket.refreshAfter > Int(interval) { + overdue["ticket_overdue_s"] = String(now - ticket.refreshAfter) + } + if !overdue.isEmpty { + overdue["status"] = String(describing: snapshot.status) + overdue["failure"] = snapshot.failure?.diagnosticCode ?? "none" + Self.journal.record("v2-host", "credential-renewal-overdue", overdue) + } + // A sequence that was already pending one full tick ago and + // still has not been applied means the snapshot consumer is + // stuck, not merely busy. + let applied = self.lastAppliedControlSequence + if let pending = pendingSinceLastTick, applied < pending { + Self.journal.record("v2-host", "snapshot-apply-stalled", [ + "pending_sequence": String(pending), + "applied_sequence": String(applied), + "stalled_for_s": String(Int(interval)), + ]) + } + pendingSinceLastTick = snapshot.sequence > applied ? snapshot.sequence : nil + } + } + } + private func startLegacyCompatibility(token: UUID) { guard isCurrent(token), pairingEnabled(), legacyStartTask == nil, legacyAcceptorPeer == nil, let service = legacyService, let identity, @@ -685,7 +754,16 @@ final class MobileHostIrxRuntime: MobileHostPairingRuntime { guard endpointTask == nil else { endpointRefreshPending = true; return } guard let supervisor = endpointSupervisor, let cache = cachedState, !cache.authorityRevoked, - Self.pathMode == .directOnly || Self.credentials(cache).contains(where: { $0.isUsable(at: Date()) }) else { return } + Self.pathMode == .directOnly || Self.credentials(cache).contains(where: { $0.isUsable(at: Date()) }) else { + // Without this event, a host with only expired credentials skips + // endpoint readiness forever and logs nothing. + let reason = endpointSupervisor == nil ? "no-supervisor" + : cachedState == nil ? "no-cache" + : cachedState?.authorityRevoked == true ? "revoked" + : "no-usable-credential" + Self.journal.record("v2-host", "endpoint-ready-skipped", ["reason": reason]) + return + } endpointTask = Task { @MainActor [weak self] in defer { if let self, self.generationToken == token { diff --git a/cmuxTests/MobileHostIrxSettingsMappingTests.swift b/cmuxTests/MobileHostIrxSettingsMappingTests.swift index 12d45e48075e..93f4e8f2c62d 100644 --- a/cmuxTests/MobileHostIrxSettingsMappingTests.swift +++ b/cmuxTests/MobileHostIrxSettingsMappingTests.swift @@ -257,6 +257,15 @@ struct MobileHostIrxActivationRetryTests { ) #expect(delay == 75) } + + @Test func serverRetryAfterCannotPushTheDelayPastTheBackoffCeiling() { + let delay = MobileHostIrxRuntime.activationRetryDelay( + after: CmxRateLimitedError(retryAfterSeconds: 24 * 3600), + failureCount: 0, + jitterUnitInterval: 0 + ) + #expect(delay == 30 * 60) + } } @MainActor diff --git a/ios/cmuxPackage/Sources/cmuxFeature/MobileIrxRuntimeComposition+Lifecycle.swift b/ios/cmuxPackage/Sources/cmuxFeature/MobileIrxRuntimeComposition+Lifecycle.swift index 55267baa6caa..9b65238ccfa6 100644 --- a/ios/cmuxPackage/Sources/cmuxFeature/MobileIrxRuntimeComposition+Lifecycle.swift +++ b/ios/cmuxPackage/Sources/cmuxFeature/MobileIrxRuntimeComposition+Lifecycle.swift @@ -114,7 +114,8 @@ extension MobileIrxRuntimeComposition { sign: { data in guard await auth.isAuthenticatedTeamScopeCurrent(scope) else { throw CompositionError.scopeChanged } return try key.sign(data) - }) + }, + journal: journal) let service = V2ControlService(configuration: try V2ControlConfiguration( baseURL: configuration.baseURL, device: device), dependencies: dependencies, store: stateStore) try await assertScope(scope, epoch: currentEpoch)