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

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -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)
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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:
Expand All @@ -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,
Expand All @@ -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
Expand All @@ -37,5 +41,6 @@ public struct V2ControlDependencies: Sendable {
self.now = now
self.sleep = sleep
self.jitter = jitter
self.journal = journal
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand All @@ -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
Expand All @@ -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
Expand All @@ -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 }
Expand Down Expand Up @@ -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)
Expand Down Expand Up @@ -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 ?? "") {
Expand Down Expand Up @@ -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
}
}
Expand All @@ -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)
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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)
Expand All @@ -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
Expand Down Expand Up @@ -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) }
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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
}
Expand Down Expand Up @@ -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
}
Expand Down Expand Up @@ -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
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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()
}
Expand Down Expand Up @@ -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)
}
}
Loading
Loading