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/IrxJournal.swift b/Packages/Shared/CmuxIrxTransport/Sources/CmuxIrxTransport/IrxJournal.swift index 8c8478645e40..19b78c5abc72 100644 --- a/Packages/Shared/CmuxIrxTransport/Sources/CmuxIrxTransport/IrxJournal.swift +++ b/Packages/Shared/CmuxIrxTransport/Sources/CmuxIrxTransport/IrxJournal.swift @@ -53,6 +53,7 @@ public final class IrxJournal: @unchecked Sendable { private var fileHandle: FileHandle? private var ring: [IrxJournalEvent] = [] private var counters: [String: Int] = [:] + private var taps: [UUID: @Sendable (IrxJournalEvent) -> Void] = [:] private var terminalTraceWindowInitialized = false private var terminalTraceWindowStartMs: UInt64 = 0 private var terminalTraceEventsInWindow = 0 @@ -129,6 +130,27 @@ public final class IrxJournal: @unchecked Sendable { if let fileHandle { try? fileHandle.write(contentsOf: Data((rendered + "\n").utf8)) } + let observers = Array(taps.values) + lock.unlock() + // Delivered outside the lock so a tap can never deadlock the journal. + // Taps receive the redacted entry and must not block. + for observer in observers { observer(entry) } + } + + /// Registers an observer for every subsequent redacted event. The tap runs + /// synchronously on the recording thread, so implementations only enqueue. + /// - Returns: A token for ``removeTap(_:)``. + public func addTap(_ tap: @escaping @Sendable (IrxJournalEvent) -> Void) -> UUID { + let id = UUID() + lock.lock() + taps[id] = tap + lock.unlock() + return id + } + + public func removeTap(_ id: UUID) { + lock.lock() + taps.removeValue(forKey: id) lock.unlock() } diff --git a/Packages/Shared/CmuxIrxTransport/Sources/CmuxIrxTransport/IrxJournalUploader.swift b/Packages/Shared/CmuxIrxTransport/Sources/CmuxIrxTransport/IrxJournalUploader.swift new file mode 100644 index 000000000000..ee02b0431771 --- /dev/null +++ b/Packages/Shared/CmuxIrxTransport/Sources/CmuxIrxTransport/IrxJournalUploader.swift @@ -0,0 +1,227 @@ +public import Foundation + +/// Ships the credential-lifecycle slice of the transport journal to the +/// authenticated `/api/observability/transport` route, so wedge evidence +/// survives the Mac's short unified-log retention. Data-plane chatter never +/// leaves the device: only allowlisted components pass, minus an explicit +/// event denylist for the periodic ones. +public actor IrxJournalUploader { + public struct ClientMetadata: Sendable { + public var platform: String + public var clientChannel: String + public var appVersion: String? + public var buildNumber: String? + public var bundleIdentifier: String? + public var osVersion: String? + /// 12-hex endpoint prefix; matches client journal and server sink rows. + public var endpoint: String? + public var deviceId: String? + public var buildTag: String? + + public init( + platform: String, clientChannel: String, appVersion: String? = nil, + buildNumber: String? = nil, bundleIdentifier: String? = nil, osVersion: String? = nil, + endpoint: String? = nil, deviceId: String? = nil, buildTag: String? = nil + ) { + self.platform = platform + self.clientChannel = clientChannel + self.appVersion = appVersion + self.buildNumber = buildNumber + self.bundleIdentifier = bundleIdentifier + self.osVersion = osVersion + self.endpoint = endpoint + self.deviceId = deviceId + self.buildTag = buildTag + } + } + + /// Components worth durable retention: the credential-renewal pipeline, + /// connection lifecycle, and admission. Everything else stays local. + public static let exportedComponents: Set = [ + "v2-control", "v2-host", "v2-lifecycle", "endpoint", "admission", + "host-runtime", "engine", "broker", "control-plane", "legacy-dialect", + "connection", "client-runtime", "registry", "host-events", "device-list", + ] + /// Periodic events inside exported components that would dominate volume + /// without adding wedge evidence. + public static let deniedEvents: Set = [ + "pong-sent", "ponged", "pong", "hint-update", "discovered", "acked", + "directory", "lane-accepted", + ] + + public static let maximumBatch = 100 + private static let bufferCapacity = 500 + private static let flushThreshold = 50 + + private let endpoint: URL + private let metadata: ClientMetadata + private let token: @Sendable (_ forceRefresh: Bool) async throws -> String + private let transport: @Sendable (URLRequest) async throws -> Int + private let flushInterval: TimeInterval + private var buffer: [[String: Any]] = [] + private var flushTask: Task? + private var stopped = false + public private(set) var droppedCount = 0 + public private(set) var uploadedCount = 0 + + /// - Parameters: + /// - endpoint: The complete transport-journal ingest URL. + /// - metadata: Attribution attached to every exported event. + /// - token: Bearer credential provider; `true` forces a refresh. + /// - transport: Executes the request and returns the HTTP status. + /// - flushInterval: Idle time before a partial batch ships. + public init( + endpoint: URL, + metadata: ClientMetadata, + token: @escaping @Sendable (_ forceRefresh: Bool) async throws -> String, + transport: @escaping @Sendable (URLRequest) async throws -> Int, + flushInterval: TimeInterval = 30 + ) { + self.endpoint = endpoint + self.metadata = metadata + self.token = token + self.transport = transport + self.flushInterval = flushInterval + } + + /// Journal-tap entry point; synchronous and non-blocking by contract. + public nonisolated func offer(_ event: IrxJournalEvent) { + guard Self.exportedComponents.contains(event.component), + !Self.deniedEvents.contains(event.event) else { return } + Task { await self.enqueue(event) } + } + + public func stop() { + stopped = true + flushTask?.cancel() + flushTask = nil + buffer.removeAll() + } + + /// Ships everything currently buffered; used by tests and shutdown paths. + public func flushNow() async { + flushTask?.cancel() + flushTask = nil + await flush() + } + + private func enqueue(_ event: IrxJournalEvent) async { + guard !stopped else { return } + buffer.append(wire(event)) + if buffer.count > Self.bufferCapacity { + droppedCount += buffer.count - Self.bufferCapacity + buffer.removeFirst(buffer.count - Self.bufferCapacity) + } + if buffer.count >= Self.flushThreshold { + flushTask?.cancel() + flushTask = nil + await flush() + return + } + guard flushTask == nil else { return } + let interval = flushInterval + flushTask = Task { [weak self] in + do { try await Task.sleep(for: .seconds(interval)) } + catch { return } + await self?.scheduledFlush() + } + } + + private func scheduledFlush() async { + flushTask = nil + await flush() + } + + private func flush() async { + guard !stopped, !buffer.isEmpty else { return } + let batch = Array(buffer.prefix(Self.maximumBatch)) + buffer.removeFirst(batch.count) + guard let body = try? JSONSerialization.data(withJSONObject: ["batch": batch]) else { + droppedCount += batch.count + return + } + var status = await post(body: body, forceToken: false) + guard !stopped else { return } + if status == 401 { + status = await post(body: body, forceToken: true) + guard !stopped else { return } + } + switch status { + case 200..<300: + uploadedCount += batch.count + case 400, 404, 413: + // The server rejected the batch shape, or this backend does not + // serve the route yet; retrying this batch cannot succeed. Later + // batches still try, so the lane comes up when the route ships. + droppedCount += batch.count + default: + // Auth outage, rate limit, or transport failure: retain for the + // next flush, bounded by the buffer capacity. + buffer.insert(contentsOf: batch, at: 0) + if buffer.count > Self.bufferCapacity { + droppedCount += buffer.count - Self.bufferCapacity + buffer.removeLast(buffer.count - Self.bufferCapacity) + } + } + if !buffer.isEmpty, flushTask == nil, !stopped { + let interval = flushInterval + flushTask = Task { [weak self] in + do { try await Task.sleep(for: .seconds(interval)) } + catch { return } + await self?.scheduledFlush() + } + } + } + + private func post(body: Data, forceToken: Bool) async -> Int { + guard !stopped else { return -1 } + guard let credential = try? await token(forceToken) else { return -1 } + guard !stopped else { return -1 } + var request = URLRequest(url: endpoint) + request.httpMethod = "POST" + request.timeoutInterval = 15 + request.setValue("application/json", forHTTPHeaderField: "Content-Type") + request.setValue("Bearer " + credential, forHTTPHeaderField: "Authorization") + request.httpBody = body + let status = (try? await transport(request)) ?? -1 + guard !stopped else { return -1 } + return status + } + + private let timestampFormatter: ISO8601DateFormatter = { + let formatter = ISO8601DateFormatter() + formatter.formatOptions = [.withInternetDateTime, .withFractionalSeconds] + return formatter + }() + + private func wire(_ event: IrxJournalEvent) -> [String: Any] { + var value: [String: Any] = [ + "timestamp": timestampFormatter.string(from: event.wallTime), + "monoMs": event.monotonicMs, + "component": event.component, + "event": event.event, + "platform": metadata.platform, + "clientChannel": metadata.clientChannel, + ] + if let appVersion = metadata.appVersion { value["appVersion"] = appVersion } + if let buildNumber = metadata.buildNumber { value["buildNumber"] = buildNumber } + if let bundleIdentifier = metadata.bundleIdentifier { value["bundleIdentifier"] = bundleIdentifier } + if let osVersion = metadata.osVersion { value["osVersion"] = osVersion } + if let endpoint = metadata.endpoint { value["endpoint"] = endpoint } + if let deviceId = metadata.deviceId { value["deviceId"] = deviceId } + if let buildTag = metadata.buildTag { value["buildTag"] = buildTag } + if !event.attributes.isEmpty { + // Server-side caps: 16 keys, snake-case keys, 160-char values. + var attributes: [String: String] = [:] + for (key, item) in event.attributes.sorted(by: { $0.key < $1.key }).prefix(16) { + let normalized = key.lowercased().replacingOccurrences( + of: "[^a-z0-9_]", with: "_", options: .regularExpression) + guard !normalized.isEmpty, normalized.count <= 32 else { continue } + guard !item.isEmpty else { continue } + attributes[normalized] = String(item.prefix(160)) + } + if !attributes.isEmpty { value["attributes"] = attributes } + } + return value + } +} diff --git a/Packages/Shared/CmuxIrxTransport/Sources/CmuxIrxTransport/IrxRelayCredentialInstaller.swift b/Packages/Shared/CmuxIrxTransport/Sources/CmuxIrxTransport/IrxRelayCredentialInstaller.swift index c5ac62f4b2b0..08f55bc0ea3a 100644 --- a/Packages/Shared/CmuxIrxTransport/Sources/CmuxIrxTransport/IrxRelayCredentialInstaller.swift +++ b/Packages/Shared/CmuxIrxTransport/Sources/CmuxIrxTransport/IrxRelayCredentialInstaller.swift @@ -76,17 +76,36 @@ actor IrxRelayCredentialInstaller { var failures = 0 while !Task.isCancelled, !stopped, taskID == id { let observedRevision = revision - let pending = desired.values.filter { - $0.isUsable(at: now()) && installed[$0.relayURL] != $0 - }.sorted { $0.relayURL < $1.relayURL } + let observedAt = now() + var unusable = 0 + var pending: [IrxRelayCredential] = [] + for credential in desired.values { + guard credential.isUsable(at: observedAt) else { + unusable += 1 + continue + } + if installed[credential.relayURL] != credential { + pending.append(credential) + } + } + pending.sort { $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..d2028b26548e 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 ?? "") { diff --git a/Packages/Shared/CmuxIrxTransport/Sources/CmuxIrxTransport/V2/V2ControlService+Maintenance.swift b/Packages/Shared/CmuxIrxTransport/Sources/CmuxIrxTransport/V2/V2ControlService+Maintenance.swift index 9c620b6512c7..bc563d2601e1 100644 --- a/Packages/Shared/CmuxIrxTransport/Sources/CmuxIrxTransport/V2/V2ControlService+Maintenance.swift +++ b/Packages/Shared/CmuxIrxTransport/Sources/CmuxIrxTransport/V2/V2ControlService+Maintenance.swift @@ -31,7 +31,15 @@ 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)]) + scheduleRenewalHealthCheck(run: run) #if DEBUG if let interval = verificationRenewalInterval, nextVerificationRenewalAt == nil { nextVerificationRenewalAt = dependencies.now().addingTimeInterval(interval) @@ -41,21 +49,101 @@ extension V2ControlService { renewalTask = Task { [weak self] in await self?.maintain(run: run) } } + func scheduleRenewalHealthCheck(run: UUID) { + guard runID == run, status == .ready else { return } + let now = dependencies.now().timeIntervalSince1970 + let refreshAfter = [ + cache.ticket?.refreshAfter, + cache.relayCredentials.map(\.refreshAfter).min(), + ].compactMap { $0 }.map { TimeInterval($0) }.min() + guard let refreshAfter else { + renewalHealthTask?.cancel() + renewalHealthTask = nil + return + } + armRenewalHealthCheck(run: run, delay: max(0, refreshAfter + 300 - now)) + } + + private func armRenewalHealthCheck(run: UUID, delay: TimeInterval) { + renewalHealthTask?.cancel() + let dependencies = dependencies + renewalHealthTask = Task { [weak self] in + do { try await dependencies.sleep(delay) } + catch { return } + await self?.renewalHealthCheckFired(run: run) + } + } + + private func renewalHealthCheckFired(run: UUID) { + renewalHealthTask = nil + guard runID == run, status == .ready, !Task.isCancelled else { return } + let now = Int(dependencies.now().timeIntervalSince1970) + var overdue: [String: String] = [:] + if !cache.authorityRevoked, + let refreshAfter = cache.relayCredentials.map(\.refreshAfter).min(), + now - refreshAfter > 300 { + overdue["relay_overdue_s"] = String(now - refreshAfter) + let expires = cache.relayCredentials.map(\.expiresAt).max() ?? now + overdue["relay_expires_in_s"] = String(expires - now) + } + if !cache.authorityRevoked, let ticket = cache.ticket, + now - ticket.refreshAfter > 300 { + overdue["ticket_overdue_s"] = String(now - ticket.refreshAfter) + } + if overdue.isEmpty { + scheduleRenewalHealthCheck(run: run) + } else { + overdue["status"] = String(describing: status) + overdue["failure"] = failure?.diagnosticCode ?? "none" + journal("credential-renewal-overdue", overdue) + armRenewalHealthCheck(run: run, delay: 300) + } + } + private func maintain(run: UUID) async { while runID == run, status == .ready, !Task.isCancelled { let now = dependencies.now().timeIntervalSince1970 - let ticketDue = due(cache.ticket?.refreshAfter, schema: "ticket.request.v1", now: now) - let relayDue = due(cache.relayCredentials.map(\.refreshAfter).min(), schema: "relay.request.v1", now: now) - let directoryDue = due(cache.directory.map { $0.permissionExpiresAt - 300 }, schema: "directory.request.v1", now: now) + let ticketRefreshAfter = cache.ticket?.refreshAfter + let relayRefreshAfter = cache.relayCredentials.map(\.refreshAfter).min() + let directoryRefreshAfter = cache.directory.map { $0.permissionExpiresAt - 300 } + let ticketDue = due(ticketRefreshAfter, schema: "ticket.request.v1", now: now) + let relayDue = due(relayRefreshAfter, schema: "relay.request.v1", now: now) + let directoryDue = due(directoryRefreshAfter, schema: "directory.request.v1", now: now) #if DEBUG let verificationDue = nextVerificationRenewalAt?.timeIntervalSince1970 ?? .infinity #else 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, + ticket: (ticketRefreshAfter, ticketDue), + relay: (relayRefreshAfter, relayDue), + directory: (directoryRefreshAfter, directoryDue) + ).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 { + // Nothing restarts maintenance until the next reconnect or + // foreground. Renewals stop here while status stays .ready. + journal("maintenance-exited", ["reason": "no-transport", "status": String(describing: status)]) + return + } let deadline = dependencies.now().timeIntervalSince1970 + 0.1 #if DEBUG let forceVerification = verificationDue <= deadline @@ -87,15 +175,38 @@ 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, + ticket: (Int?, TimeInterval), + relay: (Int?, TimeInterval), + directory: (Int?, TimeInterval) + ) -> [String] { + let wanted: [(String, Int?, TimeInterval)] = [ + ("ticket.request.v1", ticket.0, ticket.1), + ("relay.request.v1", relay.0, relay.1), + ("directory.request.v1", directory.0, directory.1), + ] + return wanted.compactMap { schema, refreshAfter, nextDue in + nextDue > 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..239c016ce8fe 100644 --- a/Packages/Shared/CmuxIrxTransport/Sources/CmuxIrxTransport/V2/V2ControlService+Operations.swift +++ b/Packages/Shared/CmuxIrxTransport/Sources/CmuxIrxTransport/V2/V2ControlService+Operations.swift @@ -35,6 +35,12 @@ extension V2ControlService { cache.ticket = response.ticket failure = nil try await persist(run: run) + 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), + ]) return response.ticket } @@ -81,6 +87,13 @@ extension V2ControlService { cache.relayCredentials = response.credentials failure = nil try await persist(run: run) + 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), + ]) return response.credentials } @@ -148,6 +161,12 @@ extension V2ControlService { cache.directory = directory failure = nil try await persist(run: run) + 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)), + ]) return directory } throw V2ControlFailure.server(V2ErrorResponse(code: .revisionConflict, requestID: "directory-refresh", retryable: true, retryAfterMS: 1000, schemaID: .errorV1)) diff --git a/Packages/Shared/CmuxIrxTransport/Sources/CmuxIrxTransport/V2/V2ControlService.swift b/Packages/Shared/CmuxIrxTransport/Sources/CmuxIrxTransport/V2/V2ControlService.swift index 4ab3e1b42eef..d63c7b73886e 100644 --- a/Packages/Shared/CmuxIrxTransport/Sources/CmuxIrxTransport/V2/V2ControlService.swift +++ b/Packages/Shared/CmuxIrxTransport/Sources/CmuxIrxTransport/V2/V2ControlService.swift @@ -22,6 +22,10 @@ public actor V2ControlService { var socket: (any V2ControlSocket)? var receiveTask: Task? var renewalTask: Task? + var renewalHealthTask: Task? + var applyWatchdogTask: Task? + var pendingApplySequence: UInt64? + var acknowledgedApplySequence: UInt64 = 0 var directoryTask: Task? var directorySyncTask: Task? var ticketTask: Task? @@ -103,6 +107,8 @@ public actor V2ControlService { guard runID == nil, !Task.isCancelled else { return } let id = UUID() runID = id + pendingApplySequence = nil + acknowledgedApplySequence = sequence status = .connecting publish() runTask = Task { [weak self] in await self?.run(id) } @@ -120,6 +126,9 @@ public actor V2ControlService { receiveTask?.cancel() receiveTask = nil cancelMaintenance() + applyWatchdogTask?.cancel() + applyWatchdogTask = nil + pendingApplySequence = nil finishAll(throwing: V2ControlFailure.stopped) status = .stopped publish() @@ -155,10 +164,48 @@ public actor V2ControlService { sequence &+= 1 let value = snapshot() for observer in observers.values { observer.yield(value) } + guard status != .stopped, !observers.isEmpty, let run = runID else { return } + guard pendingApplySequence == nil else { return } + pendingApplySequence = value.sequence + armApplyWatchdog(run: run, sequence: value.sequence) } private func removeObserver(_ id: UUID) { observers.removeValue(forKey: id) } + /// Acknowledges that the platform consumer finished applying a snapshot. + /// + /// The acknowledgement is the boundary between transport publication and + /// endpoint/UI state. A missing acknowledgement is retained as evidence by + /// the transport owner instead of being inferred by a second polling owner. + /// - Parameter sequence: The snapshot sequence that was fully applied. + public func acknowledgeApplied(sequence: UInt64) { + acknowledgedApplySequence = max(acknowledgedApplySequence, sequence) + guard let pendingApplySequence, sequence >= pendingApplySequence else { return } + self.pendingApplySequence = nil + applyWatchdogTask?.cancel() + applyWatchdogTask = nil + } + + private func armApplyWatchdog(run: UUID, sequence: UInt64) { + applyWatchdogTask?.cancel() + let dependencies = dependencies + applyWatchdogTask = Task { [weak self] in + do { try await dependencies.sleep(300) } + catch { return } + await self?.applyWatchdogFired(run: run, sequence: sequence) + } + } + + private func applyWatchdogFired(run: UUID, sequence: UInt64) { + guard runID == run, pendingApplySequence == sequence, status != .stopped else { return } + journal("snapshot-apply-stalled", [ + "pending_sequence": String(sequence), + "applied_sequence": String(acknowledgedApplySequence), + "stalled_for_s": "300", + ]) + armApplyWatchdog(run: run, sequence: sequence) + } + func assertCurrent(_ run: UUID) throws { guard runID == run, !Task.isCancelled else { throw V2ControlFailure.stopped } } @@ -167,9 +214,16 @@ 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() + if status == .ready { scheduleRenewalHealthCheck(run: run) } } func cancelMaintenance() { @@ -194,6 +248,8 @@ public actor V2ControlService { relayTask = nil authTask?.cancel() authTask = nil + renewalHealthTask?.cancel() + renewalHealthTask = nil #if DEBUG nextVerificationRenewalAt = nil #endif @@ -309,15 +365,24 @@ public actor V2ControlService { 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 = max(1, Double(response.retryAfterMS ?? 60_000) / 1000) + 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()) 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/IrxJournalUploaderTests.swift b/Packages/Shared/CmuxIrxTransport/Tests/CmuxIrxTransportTests/IrxJournalUploaderTests.swift new file mode 100644 index 000000000000..f7eeaa9ee1e6 --- /dev/null +++ b/Packages/Shared/CmuxIrxTransport/Tests/CmuxIrxTransportTests/IrxJournalUploaderTests.swift @@ -0,0 +1,179 @@ +import Foundation +import Testing +@testable import CmuxIrxTransport + +@Suite(.timeLimit(.minutes(1))) struct IrxJournalUploaderTests { + private final class Recorder: @unchecked Sendable { + private let lock = NSLock() + private var requests: [URLRequest] = [] + private var statuses: [Int] + private var tokens: [String] = [] + init(statuses: [Int]) { self.statuses = statuses } + func record(_ request: URLRequest) -> Int { + lock.lock(); defer { lock.unlock() } + requests.append(request) + return statuses.isEmpty ? 200 : statuses.removeFirst() + } + func recordToken(force: Bool) -> String { + lock.lock(); defer { lock.unlock() } + let token = force ? "fresh-token" : "cached-token" + tokens.append(token) + return token + } + var sent: [URLRequest] { lock.lock(); defer { lock.unlock() }; return requests } + var issuedTokens: [String] { lock.lock(); defer { lock.unlock() }; return tokens } + } + + private func event( + component: String = "v2-control", event name: String = "refresh-failed", + attributes: [String: String] = ["schema": "relay.request.v1"] + ) -> IrxJournalEvent { + IrxJournalEvent( + wallTime: Date(timeIntervalSince1970: 1_789_000_000), monotonicMs: 12_345, + component: component, event: name, attributes: attributes + ) + } + + private func uploader(_ recorder: Recorder, flushInterval: TimeInterval = 600) -> IrxJournalUploader { + IrxJournalUploader( + endpoint: URL(string: "https://api.example.com/api/observability/transport")!, + metadata: IrxJournalUploader.ClientMetadata( + platform: "mac", clientChannel: "nightly", appVersion: "1.2", + endpoint: "10533b43db35", deviceId: "device-1", buildTag: "default" + ), + token: { force in recorder.recordToken(force: force) }, + transport: { request in recorder.record(request) }, + flushInterval: flushInterval + ) + } + + /// `offer` hops onto the actor asynchronously; poll until ready or timeout. + private func drain(_ uploader: IrxJournalUploader, until ready: @Sendable () -> Bool) async throws { + let clock = ContinuousClock() + let deadline = clock.now.advanced(by: .seconds(10)) + while !ready() { + await uploader.flushNow() + if ready() { return } + guard clock.now < deadline else { + Issue.record("Timed out waiting for journal uploader readiness") + return + } + try await Task.sleep(for: .milliseconds(5)) + } + } + + @Test func exportsAllowlistedEventsWithMetadataAndBoundedAttributes() async throws { + let recorder = Recorder(statuses: [200]) + let uploader = uploader(recorder) + uploader.offer(event(attributes: ["schema": "relay.request.v1", "empty": ""])) + // Denied component and denied periodic event never reach the wire. + uploader.offer(event(component: "terminal-trace")) + uploader.offer(event(component: "control-plane", event: "pong-sent")) + try await drain(uploader) { !recorder.sent.isEmpty } + let sent = recorder.sent + #expect(sent.count == 1) + let request = try #require(sent.first) + #expect(request.url?.path == "/api/observability/transport") + #expect(request.value(forHTTPHeaderField: "Authorization") == "Bearer cached-token") + let body = try JSONSerialization.jsonObject(with: #require(request.httpBody)) as? [String: Any] + let batch = try #require(body?["batch"] as? [[String: Any]]) + #expect(batch.count == 1) + #expect(batch.first?["component"] as? String == "v2-control") + #expect(batch.first?["event"] as? String == "refresh-failed") + #expect(batch.first?["platform"] as? String == "mac") + #expect(batch.first?["clientChannel"] as? String == "nightly") + #expect(batch.first?["endpoint"] as? String == "10533b43db35") + #expect(batch.first?["buildTag"] as? String == "default") + let attributes = batch.first?["attributes"] as? [String: String] + #expect(attributes?["schema"] == "relay.request.v1") + #expect(attributes?["empty"] == nil) + #expect(await uploader.uploadedCount == 1) + } + + @Test func unauthorizedUploadRetriesOnceWithAForcedToken() async throws { + let recorder = Recorder(statuses: [401, 200]) + let uploader = uploader(recorder) + uploader.offer(event()) + try await drain(uploader) { recorder.sent.count == 2 } + #expect(recorder.sent.count == 2) + #expect(recorder.issuedTokens == ["cached-token", "fresh-token"]) + #expect(recorder.sent.last?.value(forHTTPHeaderField: "Authorization") == "Bearer fresh-token") + #expect(await uploader.uploadedCount == 1) + } + + @Test func stopPreventsAnInFlightUploadFromSending() async throws { + let started = AsyncStream.makeStream() + let release = AsyncStream.makeStream() + let returned = AsyncStream.makeStream() + defer { + started.continuation.finish() + release.continuation.finish() + returned.continuation.finish() + } + + let recorder = Recorder(statuses: [200]) + let uploader = IrxJournalUploader( + endpoint: URL(string: "https://api.example.com/api/observability/transport")!, + metadata: IrxJournalUploader.ClientMetadata( + platform: "mac", clientChannel: "nightly", appVersion: "1.2", + endpoint: "10533b43db35", deviceId: "device-1", buildTag: "default" + ), + token: { _ in + started.continuation.yield(()) + for await _ in release.stream { break } + returned.continuation.yield(()) + return "token" + }, + transport: { request in recorder.record(request) }, + flushInterval: 600 + ) + + for _ in 0..<50 { + uploader.offer(event()) + } + var startedIterator = started.stream.makeAsyncIterator() + #expect(await startedIterator.next() != nil) + await uploader.stop() + release.continuation.yield(()) + var returnedIterator = returned.stream.makeAsyncIterator() + #expect(await returnedIterator.next() != nil) + await uploader.flushNow() + + #expect(recorder.sent.isEmpty) + #expect(await uploader.uploadedCount == 0) + } + + @Test func transientFailureRetainsTheBatchAndRejectionDropsIt() async throws { + let recorder = Recorder(statuses: [503, 200]) + let uploader = uploader(recorder) + uploader.offer(event()) + try await drain(uploader) { recorder.sent.count == 1 } + #expect(await uploader.uploadedCount == 0) + await uploader.flushNow() + #expect(await uploader.uploadedCount == 1) + + let rejecting = Recorder(statuses: [400]) + let dropper = self.uploader(rejecting) + dropper.offer(event()) + try await drain(dropper) { rejecting.sent.count == 1 } + #expect(await dropper.droppedCount == 1) + await dropper.flushNow() + #expect(rejecting.sent.count == 1) + } + + @Test func journalTapDeliversRedactedEventsToTheUploader() async throws { + let journal = IrxJournal(subsystem: "com.cmux.test", category: "uploader-tap-test") + let recorder = Recorder(statuses: [200]) + let uploader = uploader(recorder) + let tap = journal.addTap { [weak uploader] entry in uploader?.offer(entry) } + journal.record("v2-control", "cooldown-set", ["schema": "relay.request.v1", "source": "rate_limited"]) + journal.record("terminal-trace", "host_received") + try await drain(uploader) { !recorder.sent.isEmpty } + journal.removeTap(tap) + let request = try #require(recorder.sent.first) + let body = try JSONSerialization.jsonObject(with: #require(request.httpBody)) as? [String: Any] + let batch = try #require(body?["batch"] as? [[String: Any]]) + #expect(batch.count == 1) + #expect(batch.first?["event"] as? String == "cooldown-set") + } +} diff --git a/Packages/Shared/CmuxIrxTransport/Tests/CmuxIrxTransportTests/V2/V2ControlServiceTests.swift b/Packages/Shared/CmuxIrxTransport/Tests/CmuxIrxTransportTests/V2/V2ControlServiceTests.swift index 5d54aa3ea25b..5117e6f7adbc 100644 --- a/Packages/Shared/CmuxIrxTransport/Tests/CmuxIrxTransportTests/V2/V2ControlServiceTests.swift +++ b/Packages/Shared/CmuxIrxTransport/Tests/CmuxIrxTransportTests/V2/V2ControlServiceTests.swift @@ -2,6 +2,37 @@ import Foundation import Testing @testable import CmuxIrxTransport +private final class SleepRecorder: @unchecked Sendable { + private let lock = NSLock() + private var recordedDurations: [TimeInterval] = [] + let requests: AsyncStream + private let continuation: AsyncStream.Continuation + + init() { + let pair = AsyncStream.makeStream() + requests = pair.stream + continuation = pair.continuation + } + + func sleep(_ duration: TimeInterval) async throws { + record(duration) + continuation.yield(duration) + try await Task.sleep(for: .seconds(3600)) + } + + private func record(_ duration: TimeInterval) { + lock.lock() + recordedDurations.append(duration) + lock.unlock() + } + + func durations() -> [TimeInterval] { + lock.lock() + defer { lock.unlock() } + return recordedDurations + } +} + @Suite(.timeLimit(.minutes(1))) struct V2ControlServiceTests { private let now = 1_789_000_000 @@ -14,15 +45,36 @@ import Testing ) } - private func service(backend: V2TestBackend, store: V2TestStateStore = V2TestStateStore()) throws -> V2ControlService { + private func service( + backend: V2TestBackend, + store: V2TestStateStore = V2TestStateStore(), + journal: IrxJournal? = nil, + sleep: @escaping @Sendable (TimeInterval) async throws -> Void = { seconds in + try await Task.sleep(for: .seconds(max(0, seconds))) + } + ) 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)) }, sleep: sleep, 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 { @@ -32,6 +84,36 @@ import Testing throw V2ControlFailure.stopped } + @Test func applyWatchdogStaysAnchoredToTheFirstUnacknowledgedSnapshot() async throws { + let backend = V2TestBackend(now: now) + let sleeps = SleepRecorder() + let service = try service( + backend: backend, + sleep: { seconds in try await sleeps.sleep(seconds) } + ) + let observer = await service.events() + var sleepIterator = sleeps.requests.makeAsyncIterator() + await service.start() + let firstSleep = try #require(await sleepIterator.next()) + #expect(firstSleep == 300) + + var readySnapshot: V2ControlSnapshot? + for await snapshot in observer { + if snapshot.status == .ready { + readySnapshot = snapshot + break + } + if snapshot.status == .stopped, let failure = snapshot.failure { throw failure } + } + let observedReady = try #require(readySnapshot) + #expect(observedReady.sequence > 1) + + let watchdogSleeps = sleeps.durations().filter { $0 == 300 } + #expect(watchdogSleeps.count == 1) + + await service.stop() + } + @Test func enrollmentThenResumeUsesOneRegistrationAndSignedTicket() async throws { let backend = V2TestBackend(now: now) let service = try service(backend: backend) @@ -371,4 +453,51 @@ 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") + #expect(Int(retired?.attributes["delay_s"] ?? "") ?? 0 >= 3600) + } } diff --git a/Sources/Mobile/MobileHostIrxRuntime.swift b/Sources/Mobile/MobileHostIrxRuntime.swift index aff85b13cef3..9bd0112a18a4 100644 --- a/Sources/Mobile/MobileHostIrxRuntime.swift +++ b/Sources/Mobile/MobileHostIrxRuntime.swift @@ -87,6 +87,8 @@ final class MobileHostIrxRuntime: MobileHostPairingRuntime { private(set) var generationToken = UUID() private(set) var activationTask: Task? private var controlTask: Task? + private var journalUploader: IrxJournalUploader? + private var journalTapID: UUID? private var endpointTask: Task? private var endpointRefreshPending = false private var relayAddressWatch: WatchHandle? @@ -346,6 +348,11 @@ final class MobileHostIrxRuntime: MobileHostPairingRuntime { deviceMetadataTask?.cancel(); deviceMetadataTask = nil activationTask?.cancel(); activationTask = nil controlTask?.cancel(); controlTask = nil + if let journalTapID { Self.journal.removeTap(journalTapID) } + journalTapID = nil + let oldUploader = journalUploader + journalUploader = nil + Task { await oldUploader?.stop() } endpointTask?.cancel(); endpointTask = nil endpointRefreshPending = false relayAddressWatch = nil; relayAddressWatchGeneration = nil @@ -459,7 +466,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 @@ -504,9 +512,12 @@ final class MobileHostIrxRuntime: MobileHostPairingRuntime { controlTask = Task { @MainActor [weak self] in for await snapshot in await service.events() { guard !Task.isCancelled else { return } - await self?.apply(snapshot, token: token) + guard let self else { return } + await self.apply(snapshot, token: token) + await service.acknowledgeApplied(sequence: snapshot.sequence) } } + startJournalUpload(auth: auth, scope: scope, identity: tuple, endpointID: key.endpointID) 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)]) @@ -606,15 +617,65 @@ 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() } + /// Ships the credential-lifecycle journal slice to the transport + /// observability route, so the next renewal wedge is diagnosable from the + /// sink after local log retention has erased the failure window. + private func startJournalUpload( + auth: AuthCoordinator, scope: AuthenticatedTeamScope, identity: V2Identity, endpointID: String + ) { + if let journalTapID { Self.journal.removeTap(journalTapID) } + let oldUploader = journalUploader + Task { await oldUploader?.stop() } + let channel: String + switch BuildFlavor.current { + case .dev: channel = "dev" + case .nightly: channel = "nightly" + case .stable: channel = "production" + case .rc: channel = "unknown" + } + let bundle = Bundle.main + let uploader = IrxJournalUploader( + endpoint: AuthEnvironment.vmAPIBaseURL.appendingPathComponent("api/observability/transport"), + metadata: IrxJournalUploader.ClientMetadata( + platform: "mac", + clientChannel: channel, + appVersion: bundle.object(forInfoDictionaryKey: "CFBundleShortVersionString") as? String, + buildNumber: bundle.object(forInfoDictionaryKey: "CFBundleVersion") as? String, + bundleIdentifier: bundle.bundleIdentifier, + osVersion: ProcessInfo.processInfo.operatingSystemVersionString, + endpoint: String(endpointID.prefix(12)), + deviceId: identity.deviceID, + buildTag: identity.buildTag + ), + token: { force in try await Self.accessToken(auth: auth, scope: scope, force: force) }, + transport: { request in + let (_, response) = try await URLSession.shared.data(for: request) + return (response as? HTTPURLResponse)?.statusCode ?? -1 + } + ) + journalUploader = uploader + journalTapID = Self.journal.addTap { [weak uploader] event in + uploader?.offer(event) + } + } + private func startLegacyCompatibility(token: UUID) { guard isCurrent(token), pairingEnabled(), legacyStartTask == nil, legacyAcceptorPeer == nil, let service = legacyService, let identity, @@ -685,7 +746,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/ios/cmuxPackage/Sources/cmuxFeature/MobileIrxRuntimeComposition+Lifecycle.swift b/ios/cmuxPackage/Sources/cmuxFeature/MobileIrxRuntimeComposition+Lifecycle.swift index 55267baa6caa..80800f57171e 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) @@ -122,7 +123,9 @@ extension MobileIrxRuntimeComposition { controlTask = Task { [weak self] in for await snapshot in await service.events() { guard !Task.isCancelled else { return } - await self?.apply(snapshot, scope: scope, epoch: currentEpoch) + guard let self else { return } + await self.apply(snapshot, scope: scope, epoch: currentEpoch) + await service.acknowledgeApplied(sequence: snapshot.sequence) } } await service.start() diff --git a/web/app/api/observability/transport/route.ts b/web/app/api/observability/transport/route.ts new file mode 100644 index 000000000000..06a66bebae2d --- /dev/null +++ b/web/app/api/observability/transport/route.ts @@ -0,0 +1,140 @@ +import { checkRateLimit as checkVercelRateLimit } from "@vercel/firewall"; + +import { readBoundedJsonObject } from "../../../../services/apns/routePolicy"; +import { + emitTransportJournalEvents, + MAX_TRANSPORT_JOURNAL_BATCH_EVENTS, + MAX_TRANSPORT_JOURNAL_REQUEST_BYTES, + parseTransportJournalEvent, + type TransportJournalEvent, +} from "../../../../services/observability/transportJournal"; +import { reportMissingRateLimitRule } from "../../../../services/rateLimitObservability"; +import { forceFlushTraces, setSpanAttributes, withApiRouteSpan } from "../../../../services/telemetry"; +import { verifyRequest } from "../../../../services/vms/auth"; +import { jsonResponse } from "../../../../services/vms/routeHelpers"; + +const ROUTE = "/api/observability/transport"; + +export type TransportJournalRouteDependencies = { + readonly verifyRequest: ( + request: Request, + options: { readonly allowCookie: false }, + ) => Promise<{ readonly id: string } | null>; + readonly checkRateLimit: typeof checkVercelRateLimit; + readonly emitEvents: (userId: string, batch: readonly TransportJournalEvent[]) => Promise; + readonly flushTraces: (timeoutMs?: number) => Promise; +}; + +const defaultDependencies: TransportJournalRouteDependencies = { + verifyRequest, + checkRateLimit: checkVercelRateLimit, + emitEvents: emitTransportJournalEvents, + flushTraces: forceFlushTraces, +}; + +export const POST = makeTransportJournalHandler(); + +export function makeTransportJournalHandler( + dependencies: TransportJournalRouteDependencies = defaultDependencies, +) { + return async function POST(request: Request): Promise { + return withApiRouteSpan( + request, + ROUTE, + { "cmux.subsystem": "transport-journal" }, + async (span) => { + const rateLimitResponse = await enforceRateLimit(request, dependencies); + if (rateLimitResponse) return rateLimitResponse; + + let user: { readonly id: string } | null; + try { + user = await dependencies.verifyRequest(request, { allowCookie: false }); + } catch { + return jsonResponse({ error: "auth_unavailable" }, 503, { "cache-control": "no-store" }); + } + if (!user) { + return jsonResponse({ error: "unauthorized" }, 401, { "cache-control": "no-store" }); + } + + const body = await readBoundedJsonObject(request, MAX_TRANSPORT_JOURNAL_REQUEST_BYTES); + if (!body.ok) { + return jsonResponse( + { error: body.error }, + body.error === "request_too_large" ? 413 : 400, + { "cache-control": "no-store" }, + ); + } + if (!Array.isArray(body.value.batch)) { + return jsonResponse({ error: "missing_batch" }, 400, { "cache-control": "no-store" }); + } + if (body.value.batch.length > MAX_TRANSPORT_JOURNAL_BATCH_EVENTS) { + return jsonResponse({ error: "batch_too_large" }, 400, { "cache-control": "no-store" }); + } + + const accepted = body.value.batch + .map(parseTransportJournalEvent) + .filter((entry): entry is TransportJournalEvent => entry !== null); + if (accepted.length !== body.value.batch.length) { + return jsonResponse({ error: "invalid_event" }, 400, { "cache-control": "no-store" }); + } + if (accepted.length === 0) { + return jsonResponse({ ok: true, accepted: 0 }, 200, { "cache-control": "no-store" }); + } + + setSpanAttributes(span, { + "cmux.user_id": user.id, + "cmux.transport.event_count": accepted.length, + }); + try { + await dependencies.emitEvents(user.id, accepted); + } catch { + return jsonResponse({ error: "observability_unavailable" }, 503, { "cache-control": "no-store" }); + } + // Serverless instances can be torn down after the response, and the + // batch is already accepted once emission succeeds; an ambiguous + // flush must not turn into a client retry that duplicates spans. + try { + await dependencies.flushTraces(1_000); + } catch { + // Best effort after emission. + } + return jsonResponse({ ok: true, accepted: accepted.length }, 200, { "cache-control": "no-store" }); + }, + { priority: true }, + ); + }; +} + +async function enforceRateLimit( + request: Request, + dependencies: TransportJournalRouteDependencies, +): Promise { + // Shares the client-observability firewall rule with the mobile-network + // route so this lane works without a separate Vercel rule rollout. + const rateLimitId = process.env.CMUX_MOBILE_OBSERVABILITY_RATE_LIMIT_ID?.trim(); + if (process.env.VERCEL !== "1") return null; + if (!rateLimitId) { + void reportMissingRateLimitRule({ route: ROUTE, reason: "unset" }); + return jsonResponse({ error: "observability_unavailable" }, 503, { "cache-control": "no-store" }); + } + try { + const { error, rateLimited } = await dependencies.checkRateLimit(rateLimitId, { request }); + if (rateLimited || error === "blocked") { + return jsonResponse( + { error: "rate_limited" }, + 429, + { "cache-control": "no-store", "retry-after": "60" }, + ); + } + if (error === "not-found") { + void reportMissingRateLimitRule({ route: ROUTE, reason: "not-found" }); + return jsonResponse({ error: "observability_unavailable" }, 503, { "cache-control": "no-store" }); + } + if (error) { + return jsonResponse({ error: "observability_unavailable" }, 503, { "cache-control": "no-store" }); + } + return null; + } catch { + return jsonResponse({ error: "observability_unavailable" }, 503, { "cache-control": "no-store" }); + } +} diff --git a/web/services/observability/transportJournal.ts b/web/services/observability/transportJournal.ts new file mode 100644 index 000000000000..514debe72828 --- /dev/null +++ b/web/services/observability/transportJournal.ts @@ -0,0 +1,178 @@ +import { SpanStatusCode } from "@opentelemetry/api"; + +import { withSpan } from "../telemetry"; + +export const MAX_TRANSPORT_JOURNAL_REQUEST_BYTES = 64 * 1_024; +export const MAX_TRANSPORT_JOURNAL_BATCH_EVENTS = 100; + +// The client-side irx journal components whose events may be exported. A +// bounded set keeps sink cardinality owned by this file rather than by +// whatever a future client build emits. Chatty data-plane components +// (keepalive, terminal-trace, host-surface-lanes) are deliberately absent. +const components = new Set([ + "v2-control", "v2-host", "v2-lifecycle", "endpoint", "admission", + "host-runtime", "engine", "broker", "control-plane", "legacy-dialect", + "connection", "client-runtime", "registry", "host-events", "host-lanes", + "device-list", +]); +const platforms = new Set(["mac", "ios"]); +const channels = new Set(["dev", "nightly", "production", "unknown"]); + +const EVENT_PATTERN = /^[a-z0-9][a-z0-9_-]{0,47}$/; +const ATTRIBUTE_KEY_PATTERN = /^[a-z0-9][a-z0-9_]{0,31}$/; +const ENDPOINT_PATTERN = /^[0-9a-f]{12}$/; +const MAX_ATTRIBUTES = 16; +const MAX_ATTRIBUTE_VALUE_LENGTH = 160; +const MAX_STRING_LENGTH = 120; +// Events whose presence is itself the alarm; they set span error status so +// the standard error monitors see a wedge without a bespoke query. +const FAILURE_EVENT_PATTERN = /(-failed|-overdue|-stalled|-terminal)$/; + +export type TransportJournalEvent = { + readonly timestamp: string; + readonly monoMs: number; + readonly component: string; + readonly event: string; + readonly platform: "mac" | "ios"; + readonly clientChannel?: string; + readonly appVersion?: string; + readonly buildNumber?: string; + readonly bundleIdentifier?: string; + readonly osVersion?: string; + /** 12-hex endpoint prefix, matching client journals and iroh-v2 sink rows. */ + readonly endpoint?: string; + readonly deviceId?: string; + readonly buildTag?: string; + readonly attributes?: Readonly>; +}; + +function boundedString(value: unknown, maxLength = MAX_STRING_LENGTH): string | null { + if (typeof value !== "string" || value.length === 0 || value.length > maxLength) return null; + return value; +} + +function isRecord(value: unknown): value is Record { + return typeof value === "object" && value !== null && !Array.isArray(value); +} + +function parseTransportJournalCore( + record: Record, +): Pick | null { + const timestamp = boundedString(record.timestamp, 40); + if (!timestamp || Number.isNaN(Date.parse(timestamp))) return null; + if (typeof record.monoMs !== "number" || !Number.isFinite(record.monoMs) || record.monoMs < 0) return null; + const component = boundedString(record.component, 32); + if (!component || !components.has(component)) return null; + const event = boundedString(record.event, 48); + if (!event || !EVENT_PATTERN.test(event)) return null; + const platform = boundedString(record.platform, 8); + if (!platform || !platforms.has(platform)) return null; + return { + timestamp, + monoMs: Math.floor(record.monoMs), + component, + event, + platform: platform as "mac" | "ios", + }; +} + +type ParsedTransportJournalMetadata = { + clientChannel?: string; + appVersion?: string; + buildNumber?: string; + bundleIdentifier?: string; + osVersion?: string; + endpoint?: string; + deviceId?: string; + buildTag?: string; + attributes?: Readonly>; +}; + +function parseTransportJournalMetadata(record: Record): ParsedTransportJournalMetadata | null { + const metadata: ParsedTransportJournalMetadata = {}; + if (record.clientChannel !== undefined) { + const channel = boundedString(record.clientChannel, 16); + if (!channel || !channels.has(channel)) return null; + metadata.clientChannel = channel; + } + for (const key of ["appVersion", "buildNumber", "bundleIdentifier", "osVersion", "deviceId", "buildTag"] as const) { + if (record[key] === undefined) continue; + const parsed = boundedString(record[key]); + if (!parsed) return null; + metadata[key] = parsed; + } + if (record.endpoint !== undefined) { + const endpoint = boundedString(record.endpoint, 12); + if (!endpoint || !ENDPOINT_PATTERN.test(endpoint)) return null; + metadata.endpoint = endpoint; + } + if (record.attributes !== undefined) { + const attributes = parseTransportJournalAttributes(record.attributes); + if (!attributes) return null; + metadata.attributes = attributes; + } + return metadata; +} + +function parseTransportJournalAttributes(value: unknown): Readonly> | null { + if (!isRecord(value)) return null; + const entries = Object.entries(value); + if (entries.length > MAX_ATTRIBUTES) return null; + const attributes: Record = {}; + for (const [key, item] of entries) { + if (!ATTRIBUTE_KEY_PATTERN.test(key)) return null; + const parsed = boundedString(item, MAX_ATTRIBUTE_VALUE_LENGTH); + if (parsed === null) return null; + attributes[key] = parsed; + } + return attributes; +} + +/** Validates one client-submitted journal event; null rejects the batch. */ +export function parseTransportJournalEvent(value: unknown): TransportJournalEvent | null { + if (!isRecord(value)) return null; + const core = parseTransportJournalCore(value); + if (!core) return null; + const metadata = parseTransportJournalMetadata(value); + return metadata ? { ...core, ...metadata } : null; +} + +/** + * Emits one span per exported client transport-journal event. The span name + * is fixed and the component/event live in attributes, so cardinality stays + * bounded and one query covers the whole credential-renewal pipeline. + */ +export async function emitTransportJournalEvents( + userId: string, + batch: readonly TransportJournalEvent[], +): Promise { + await Promise.all(batch.map((entry) => withSpan( + "cmux-transport-journal", + "cmux.transport.journal", + { + "cmux.subsystem": "transport-journal", + "cmux.observation.source": "client", + "cmux.user_id": userId, + "cmux.client.channel": entry.clientChannel, + "cmux.transport.platform": entry.platform, + "cmux.transport.component": entry.component, + "cmux.transport.event": entry.event, + "cmux.transport.occurred_at": entry.timestamp, + "cmux.transport.mono_ms": entry.monoMs, + "cmux.device.endpoint": entry.endpoint, + "cmux.device.id": entry.deviceId, + "cmux.device.build_tag": entry.buildTag, + "cmux.transport.app_version": entry.appVersion, + "cmux.transport.build_number": entry.buildNumber, + "cmux.transport.bundle_identifier": entry.bundleIdentifier, + "cmux.transport.os_version": entry.osVersion, + ...Object.fromEntries(Object.entries(entry.attributes ?? {}) + .map(([key, value]) => [`cmux.transport.attr.${key}`, value])), + }, + (span) => { + if (FAILURE_EVENT_PATTERN.test(entry.event)) { + span.setStatus({ code: SpanStatusCode.ERROR, message: `${entry.component}/${entry.event}` }); + } + }, + ))); +} diff --git a/web/tests/transport-journal-observability-route.test.ts b/web/tests/transport-journal-observability-route.test.ts new file mode 100644 index 000000000000..f1b93e302588 --- /dev/null +++ b/web/tests/transport-journal-observability-route.test.ts @@ -0,0 +1,135 @@ +import { afterAll, beforeEach, describe, expect, mock, test } from "bun:test"; +import type { checkRateLimit as checkVercelRateLimit } from "@vercel/firewall"; + +import { makeTransportJournalHandler } from "../app/api/observability/transport/route"; +import { + parseTransportJournalEvent, + type TransportJournalEvent, +} from "../services/observability/transportJournal"; + +const originalVercel = process.env.VERCEL; +const originalRuleId = process.env.CMUX_MOBILE_OBSERVABILITY_RATE_LIMIT_ID; + +let authenticatedUser: { readonly id: string } | null = { id: "user-9" }; +let emitError: unknown = null; +let rateLimitResult: Awaited> = { rateLimited: false }; +const emitted: Array<{ readonly userId: string; readonly batch: readonly TransportJournalEvent[] }> = []; + +const verifyRequest = mock(async () => authenticatedUser); +const POST = makeTransportJournalHandler({ + verifyRequest, + checkRateLimit: async () => rateLimitResult, + emitEvents: async (userId, batch) => { + if (emitError) throw emitError; + emitted.push({ userId, batch }); + }, + flushTraces: async () => true, +}); + +beforeEach(() => { + delete process.env.VERCEL; + process.env.CMUX_MOBILE_OBSERVABILITY_RATE_LIMIT_ID = "client-observability-test"; + authenticatedUser = { id: "user-9" }; + emitError = null; + rateLimitResult = { rateLimited: false }; + emitted.length = 0; + verifyRequest.mockClear(); +}); + +afterAll(() => { + restoreEnv("VERCEL", originalVercel); + restoreEnv("CMUX_MOBILE_OBSERVABILITY_RATE_LIMIT_ID", originalRuleId); +}); + +function journalEvent(overrides: Record = {}): Record { + return { + timestamp: "2026-09-28T03:00:50.168Z", + monoMs: 91_670_979, + component: "v2-control", + event: "cooldown-set", + platform: "mac", + clientChannel: "nightly", + endpoint: "10533b43db35", + deviceId: "device-9", + buildTag: "default", + attributes: { schema: "relay.request.v1", source: "rate_limited", delay_s: "60" }, + ...overrides, + }; +} + +function journalRequest(batch: unknown): Request { + return new Request("https://cmux.test/api/observability/transport", { + method: "POST", + headers: { "content-type": "application/json", authorization: "Bearer token" }, + body: JSON.stringify({ batch }), + }); +} + +describe("transport journal observability route", () => { + test("attributes an accepted batch to the authenticated user", async () => { + const response = await POST(journalRequest([ + journalEvent(), + journalEvent({ event: "credential-renewal-overdue", component: "v2-host" }), + ])); + expect(response.status).toBe(200); + expect(await response.json()).toEqual({ ok: true, accepted: 2 }); + expect(emitted).toHaveLength(1); + expect(emitted[0]?.userId).toBe("user-9"); + expect(emitted[0]?.batch[0]?.endpoint).toBe("10533b43db35"); + expect(emitted[0]?.batch[1]?.event).toBe("credential-renewal-overdue"); + }); + + test("rejects a batch containing an event outside the component allowlist", async () => { + const response = await POST(journalRequest([ + journalEvent(), + journalEvent({ component: "terminal-trace" }), + ])); + expect(response.status).toBe(400); + expect(await response.json()).toEqual({ error: "invalid_event" }); + expect(emitted).toHaveLength(0); + }); + + test("requires authentication", async () => { + authenticatedUser = null; + const response = await POST(journalRequest([journalEvent()])); + expect(response.status).toBe(401); + expect(emitted).toHaveLength(0); + }); + + test("reports emission failure without accepting the batch", async () => { + emitError = new Error("sink down"); + const response = await POST(journalRequest([journalEvent()])); + expect(response.status).toBe(503); + expect(await response.json()).toEqual({ error: "observability_unavailable" }); + }); +}); + +describe("transport journal event validation", () => { + test("accepts a complete event and floors the monotonic clock", () => { + const parsed = parseTransportJournalEvent(journalEvent({ monoMs: 12.7 })); + expect(parsed?.monoMs).toBe(12); + expect(parsed?.attributes?.schema).toBe("relay.request.v1"); + }); + + test.each([ + ["unknown component", journalEvent({ component: "not-a-component" })], + ["uppercase event", journalEvent({ event: "Cooldown-Set" })], + ["bad endpoint", journalEvent({ endpoint: "10533B43DB35" })], + ["bad channel", journalEvent({ clientChannel: "beta" })], + ["attribute key with dots", journalEvent({ attributes: { "a.b": "x" } })], + ["oversized attribute value", journalEvent({ attributes: { schema: "x".repeat(161) } })], + ["too many attributes", journalEvent({ + attributes: Object.fromEntries(Array.from({ length: 17 }, (_, index) => [`k${index}`, "v"])), + })], + ["missing timestamp", journalEvent({ timestamp: undefined })], + ["unparseable timestamp", journalEvent({ timestamp: "not-a-date" })], + ["negative monotonic", journalEvent({ monoMs: -1 })], + ])("rejects %s", (_name, value) => { + expect(parseTransportJournalEvent(value)).toBeNull(); + }); +}); + +function restoreEnv(key: string, value: string | undefined): void { + if (value === undefined) delete process.env[key]; + else process.env[key] = value; +}