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

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -41,6 +41,10 @@ public protocol CmxIrohSettingsControlling: AnyObject {

/// Erases the in-memory connection timeline and rotates its report session.
func clearIrohDiagnosticReport() async

/// The archived report from the previous process launch, if one exists.
/// Exports include it so a drop that preceded a relaunch stays diagnosable.
func irohPreviousLaunchDiagnosticReport() async -> DiagnosticReport?
}

public extension CmxIrohSettingsControlling {
Expand All @@ -56,6 +60,10 @@ public extension CmxIrohSettingsControlling {
.empty
}

func irohPreviousLaunchDiagnosticReport() async -> DiagnosticReport? {
nil
}

func exportIrohDiagnosticReport() async -> Data {
Data()
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -173,4 +173,18 @@ public enum DiagnosticEventCode: UInt16, Sendable, Codable, CaseIterable {
/// process-local session correlation ID. The event contains no peer or route
/// identity.
case transportSessionLifecycle = 51
/// The app's scene phase changed. `a` is ``DiagnosticAppLifecyclePhase``.
/// Session drops that follow a backgrounding within seconds are suspension
/// casualties, not network failures; this event makes that attributable.
case appLifecycleChanged = 52
/// Device reachability changed. `a` is 1 when a usable network path
/// exists, else 0. Correlates drops with WiFi/cellular transitions.
case reachabilityChanged = 53
}

/// Scene phase carried by ``DiagnosticEventCode/appLifecycleChanged``.
public enum DiagnosticAppLifecyclePhase: Int, Sendable, Codable, CaseIterable {
case background = 0
case active = 1
case inactive = 2
}
Original file line number Diff line number Diff line change
@@ -0,0 +1,61 @@
public import Foundation

/// Persists the most recent diagnostic snapshot across process launches.
///
/// The in-memory diagnostic ring dies with the process, so the events around
/// a connection drop were unrecoverable once the user relaunched before
/// exporting. The archive keeps exactly one previous-launch report on disk
/// (bounded, privacy-safe integers only, same vocabulary as the live ring) so
/// an export can include what happened before the current launch.
public struct DiagnosticReportArchive: Sendable {
/// Upper bound for a stored report file; a report of maximum event count
/// encodes far below this, so hitting it means corruption, not data.
public static let maximumFileByteCount = 1_024 * 1_024

public let fileURL: URL

public init(fileURL: URL) {
self.fileURL = fileURL
}

/// Default location inside Application Support.
public static func defaultArchive(
fileManager: FileManager = .default
) -> DiagnosticReportArchive? {
guard let base = fileManager.urls(
for: .applicationSupportDirectory,
in: .userDomainMask
).first else { return nil }
do {
try fileManager.createDirectory(at: base, withIntermediateDirectories: true)
} catch {
return nil
}
return DiagnosticReportArchive(
fileURL: base.appendingPathComponent("cmux-diagnostic-report.json")
)
}

/// Atomically replaces the stored report. Empty reports are not worth a
/// write and would only erase a more useful previous snapshot.
public func save(_ report: DiagnosticReport) {
guard !report.events.isEmpty else { return }
guard let data = try? JSONEncoder().encode(report),
data.count <= Self.maximumFileByteCount else { return }
try? data.write(to: fileURL, options: .atomic)
}

/// Loads the previous process's report, or `nil` when absent or invalid.
public func load() -> DiagnosticReport? {
guard let data = try? Data(contentsOf: fileURL),
data.count <= Self.maximumFileByteCount,
let report = try? JSONDecoder().decode(DiagnosticReport.self, from: data),
!report.events.isEmpty else { return nil }
return report
}

/// Removes the stored report (sign-out / account erase).
public func clear(fileManager: FileManager = .default) {
try? fileManager.removeItem(at: fileURL)
}
}
Original file line number Diff line number Diff line change
@@ -0,0 +1,62 @@
import Foundation
import Testing
@testable import CMUXMobileCore

@Suite
struct DiagnosticReportArchiveTests {
private func temporaryArchive() -> DiagnosticReportArchive {
DiagnosticReportArchive(
fileURL: FileManager.default.temporaryDirectory
.appendingPathComponent("cmuxdiag-archive-\(UUID().uuidString).json")
)
}

private func report(eventCount: Int) -> DiagnosticReport {
DiagnosticReport(
role: .mobileClient,
anchorWallNanos: 1,
anchorMonotonicNanos: 1,
buildStamp: "test",
events: (0 ..< eventCount).map {
DiagnosticEvent(code: .connect, tNanos: UInt64($0 + 1))
}
)
}

@Test
func savesAndReloadsAcrossInstances() {
let archive = temporaryArchive()
defer { archive.clear() }
archive.save(report(eventCount: 3))

let reloaded = DiagnosticReportArchive(fileURL: archive.fileURL).load()
#expect(reloaded?.events.count == 3)
#expect(reloaded?.role == .mobileClient)
}

@Test
func emptyReportDoesNotReplaceStoredSnapshot() {
let archive = temporaryArchive()
defer { archive.clear() }
archive.save(report(eventCount: 2))
archive.save(report(eventCount: 0))

#expect(archive.load()?.events.count == 2)
}

@Test
func clearRemovesTheSnapshot() {
let archive = temporaryArchive()
archive.save(report(eventCount: 1))
archive.clear()
#expect(archive.load() == nil)
}

@Test
func corruptFileLoadsAsNil() throws {
let archive = temporaryArchive()
defer { archive.clear() }
try Data("not json".utf8).write(to: archive.fileURL)
#expect(archive.load() == nil)
}
}
Original file line number Diff line number Diff line change
@@ -0,0 +1,92 @@
public import CMUXMobileCore
public import Foundation

/// Account-scoped cooldown honoring a broker's Retry-After directive across
/// endpoint activation attempts.
///
/// The credential coordinator already honors Retry-After inside one activation,
/// but a torn-down runtime discards that state, so every reconnect attempt
/// re-ran the full broker call set and kept the server's rate-limit window
/// exhausted. This value survives between activations: while active, callers
/// skip broker work entirely and surface the remaining delay to the reconnect
/// scheduler.
public struct CmxIrohBrokerCooldown: Equatable, Sendable {
public private(set) var accountID: String?
public private(set) var retryAt: Date?

public init() {}

/// Records a server directive, keeping the later of two overlapping floors.
/// A directive for a different account replaces the previous floor.
@discardableResult
public mutating func record(
accountID: String,
retryAfterSeconds: Int,
now: Date
) -> Date {
let proposed = now.addingTimeInterval(TimeInterval(max(1, retryAfterSeconds)))
if self.accountID != accountID {
self.accountID = accountID
retryAt = proposed
return proposed
}
if let retryAt, retryAt >= proposed { return retryAt }
retryAt = proposed
return proposed
}

/// Whole seconds until the floor expires, or `nil` when no floor applies.
/// Expired floors are cleared on read so state cannot go stale.
public mutating func remainingSeconds(accountID: String, now: Date) -> Int? {
guard self.accountID == accountID, let retryAt else { return nil }
let remaining = retryAt.timeIntervalSince(now)
guard remaining > 0 else {
clear()
return nil
}
return Int(remaining.rounded(.up))
}

/// Clears the floor, optionally only when it belongs to one account.
public mutating func clear(accountID: String? = nil) {
guard accountID == nil || self.accountID == accountID else { return }
self.accountID = nil
retryAt = nil
}
}

public extension CmxIrohBrokerCooldown {
/// Floor applied to a 429 whose response carried no Retry-After header.
static let defaultRateLimitedSeconds = 60

/// Seconds of cooldown one broker failure demands, or `nil` when the
/// error is not a rate-limit signal. Prefers the server's own Retry-After
/// directive; a bare 429 still arms a short default floor so a missing
/// header can never reopen the retry storm.
static func directiveSeconds(for error: any Error) -> Int? {
if let retryAfterSeconds = (error as? any CmxRetryAfterProviding)?
.retryAfterSeconds {
return retryAfterSeconds
}
if case let .rejected(statusCode, _)? = error as? CmxIrohTrustBrokerClientError,
statusCode == 429 {
return defaultRateLimitedSeconds
}
return nil
}
}

/// Thrown instead of a bare inactive-runtime error while a broker cooldown is
/// active, so the reconnect scheduler can adopt the server's floor instead of
/// its short transient backoff.
public struct CmxIrohBrokerCooldownError: CmxRetryAfterProviding, Equatable {
public let retryAfterSeconds: Int?

public init(retryAfterSeconds: Int) {
self.retryAfterSeconds = max(1, retryAfterSeconds)
}
}

extension CmxIrohBrokerCooldownError: DiagnosticFailureProviding {
public var diagnosticFailureKind: DiagnosticFailureKind { .policyUnavailable }
}
Original file line number Diff line number Diff line change
Expand Up @@ -77,6 +77,10 @@ extension CmxIrohClientRuntime {
}

func validateRelayFleet(_ fleet: [String]) throws {
// Without a verified managed fleet (relay policy unavailable) there is
// nothing to cross-check and no relay will be configured; activation
// continues on direct paths instead of failing closed here.
guard !managedRelayURLs.isEmpty else { return }
guard fleet.count == managedRelayURLs.count,
Set(fleet) == managedRelayURLs else {
throw CmxIrohClientRuntimeError.relayFleetMismatch
Expand All @@ -103,6 +107,6 @@ extension CmxIrohClientRuntime {
}

static func isConnectivity(_ error: any Error) -> Bool {
(error as? CmxIrohTrustBrokerClientError) == .connectivity
CmxIrohTrustBrokerClientError.preservesVerifiedPolicyDuringRefresh(error)
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -4,7 +4,8 @@ public import Foundation
extension CmxIrohClientRuntime {
func resolvePolicy(
expectedEndpointID: CmxIrohPeerIdentity,
revision: UInt64
revision: UInt64,
allowRegistrationDiscoveryFallback: Bool = false
) async throws -> ResolvedPolicy {
try await pendingRevocations.revokePending(
accountID: configuration.accountID,
Expand Down Expand Up @@ -47,13 +48,18 @@ extension CmxIrohClientRuntime {
pairingEnabled: false,
capabilities: configuration.capabilities
)
let offlineExpectation = try offlinePolicyCache.map { _ in
try CmxIrohClientOfflinePolicyExpectation(
accountID: configuration.accountID,
localBindingExpectation: expectation,
managedRelayURLs: managedRelayURLs
)
}
// Without a managed relay fleet (policy unavailable or direct-only)
// there is no relay bootstrap to cache offline; activation proceeds
// with direct paths instead of failing the expectation's fleet check.
let offlineExpectation: CmxIrohClientOfflinePolicyExpectation? =
try offlinePolicyCache.flatMap { _ in
guard !managedRelayURLs.isEmpty else { return nil }
return try CmxIrohClientOfflinePolicyExpectation(
accountID: configuration.accountID,
localBindingExpectation: expectation,
managedRelayURLs: managedRelayURLs
)
}
let signer = try CmxIrohRegistrationSigner(
identity: configuration.identity,
endpointID: expectedEndpointID.endpointID
Expand All @@ -63,6 +69,32 @@ extension CmxIrohClientRuntime {
do {
registration = try await broker.register(prepared: prepared, signer: signer)
} catch {
if allowRegistrationDiscoveryFallback,
Self.canResolveRegistrationFromDiscovery(error) {
let discovery = try await broker.discover()
try requireCurrent(revision)
guard discovery.routeContractVersion == payload.routeContractVersion else {
throw CmxIrohClientRuntimeError.routeContractMismatch
}
try validateRelayFleet(discovery.relayFleet)
let localMatches = discovery.bindings.filter(expectation.matches)
guard localMatches.count == 1,
let discovered = localMatches.first else {
throw CmxIrohClientRuntimeError.localBindingMissingFromDiscovery
}
return ResolvedPolicy(
registration: CmxIrohRegistrationResponse(
binding: discovered,
relay: .unavailable
),
discovery: discovery,
binding: discovered,
expectation: expectation,
offlineExpectation: offlineExpectation,
cachedTargetBindings: [],
cachedLANRendezvous: nil
)
}
guard Self.isConnectivity(error),
let cached = try await offlineBootstrap(
expectation: offlineExpectation,
Expand Down Expand Up @@ -123,6 +155,25 @@ extension CmxIrohClientRuntime {
)
}

private static func canResolveRegistrationFromDiscovery(_ error: any Error) -> Bool {
guard let brokerError = error as? CmxIrohTrustBrokerClientError else {
return false
}
switch brokerError {
case .rateLimited:
return true
case let .rejected(statusCode, _):
return statusCode == 429
case .connectivity,
.invalidBaseURL,
.missingAuthentication,
.invalidAuthentication,
.nonHTTPResponse,
.invalidResponse:
return false
}
}

func offlineBootstrap(
expectation: CmxIrohClientOfflinePolicyExpectation?,
confirmedLocalBinding: CmxIrohBrokerBinding?
Expand Down Expand Up @@ -199,7 +250,8 @@ extension CmxIrohClientRuntime {
automaticRefreshEnabled: automaticRelayCredentialRefreshEnabled,
credentialDidInstall: { [handleRelayCredential] response in
await handleRelayCredential(response, policy.binding)
}
},
rateLimitedDirective: handleRelayRateLimit
)
relayCoordinator = coordinator
}
Expand All @@ -213,7 +265,8 @@ extension CmxIrohClientRuntime {
bindingID: policy.binding.bindingID,
endpointIdentity: policy.binding.endpointID,
bootstrap: bootstrap,
waitForInitialCredential: requiresRelayReadiness
waitForInitialCredential: requiresRelayReadiness,
mintNotBefore: configuration.relayCredentialMintNotBefore
)
} catch {
if requiresRelayReadiness { throw error }
Expand Down
Loading