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 @@ -25,6 +25,7 @@ public actor CmxConnectivityEngine {
private var endpointGeneration: UInt64?
private var localIdentity: CmxIrohPeerIdentity?
private var routeRevision: UInt64?
private var routeContent: CmxConnectivityRouteContent?
private var endpointEventTask: Task<Void, Never>?
private var routeSyncOperation: RouteSyncOperation?
private var peers: [CmxConnectivityPeerID: CmxConnectivityPeerSession] = [:]
Expand Down Expand Up @@ -240,10 +241,22 @@ public actor CmxConnectivityEngine {
}

/// Records the last route revision installed atomically by the composition root.
public func didInstallRouteRevision(_ revision: UInt64) async {
guard routeRevision != revision else { return }
await invalidateAllPeers(failure: .superseded)
///
/// Peers whose material route content is unchanged keep their live
/// sessions; every other peer is invalidated before the new revision
/// becomes visible.
public func didInstallRouteRevision(
_ revision: UInt64,
routes: CmxIrohDiscoveryResponse
) async {
let content = CmxConnectivityRouteContent(snapshot: routes)
guard routeRevision != revision else {
routeContent = content
return
}
Comment on lines +253 to +256

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

P2: A live peer can survive the first snapshot installed after the engine recorded a revision without route content. performRouteSync can leave routeContent nil for an unchanged response, but this same-revision branch skips the documented fail-closed invalidation; compare the prior content and invalidate, or invalidate all when the baseline is missing, before storing it.

Prompt for AI agents
Check if this issue is valid — if so, understand the root cause and fix it. At Packages/Shared/CmuxIrohTransport/Sources/CmuxIrohTransport/CmxConnectivityEngine.swift, line 253:

<comment>A live peer can survive the first snapshot installed after the engine recorded a revision without route content. `performRouteSync` can leave `routeContent` nil for an unchanged response, but this same-revision branch skips the documented fail-closed invalidation; compare the prior content and invalidate, or invalidate all when the baseline is missing, before storing it.</comment>

<file context>
@@ -240,10 +241,22 @@ public actor CmxConnectivityEngine {
+        routes: CmxIrohDiscoveryResponse
+    ) async {
+        let content = CmxConnectivityRouteContent(snapshot: routes)
+        guard routeRevision != revision else {
+            routeContent = content
+            return
</file context>
Suggested change
guard routeRevision != revision else {
routeContent = content
return
}
guard routeRevision != revision else {
if let previous = routeContent {
if previous != content {
await invalidatePeersSuperseded(by: content)
}
} else {
await invalidateAllPeers(failure: .superseded)
}
routeContent = content
return
}

await invalidatePeersSuperseded(by: content)
routeRevision = revision
routeContent = content
publishSnapshot()
}

Expand Down Expand Up @@ -678,10 +691,36 @@ public actor CmxConnectivityEngine {
throw CmxConnectivityEngineError.superseded
}
}
let content = response.snapshot.map(CmxConnectivityRouteContent.init)
if routeRevision != response.revision {
await invalidateAllPeers(failure: .superseded)
await invalidatePeersSuperseded(by: content)
routeRevision = response.revision
routeContent = content
publishSnapshot()
} else if let content {
routeContent = content
}
}

/// Invalidates peers whose authoritative route material changed.
///
/// A missing baseline or replacement fails closed and tears down every
/// peer, preserving the pre-content-tracking behavior.
private func invalidatePeersSuperseded(
by content: CmxConnectivityRouteContent?
) async {
guard let previous = routeContent,
let content,
previous.account == content.account else {
await invalidateAllPeers(failure: .superseded)
return
}
for (peerID, peer) in peers {
guard let previousRoute = previous.peerRoute(for: peerID),
previousRoute == content.peerRoute(for: peerID) else {
await peer.invalidate(failure: .superseded)
continue
}
}
}

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -200,6 +200,16 @@ actor CmxConnectivityPeerSession {
continue redial
}

// The dead-on-arrival probe suspends this actor. A concurrent
// caller that dialed in that window may have installed first;
// installing over it would leak its session and double-record
// an established lifecycle for the same peer.
if let installed = activeConnection {

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

P2: A concurrent invalidation, remote close, or replacement dial can run while the redundant session is being closed, but this returns the previously captured installed.session without revalidating the slot. Re-read the active ID (and liveness) after the await and retry the dial when the winner changed, otherwise callers can receive a closed session.

Prompt for AI agents
Check if this issue is valid — if so, understand the root cause and fix it. At Packages/Shared/CmuxIrohTransport/Sources/CmuxIrohTransport/CmxConnectivityPeerSession.swift, line 207:

<comment>A concurrent invalidation, remote close, or replacement dial can run while the redundant session is being closed, but this returns the previously captured `installed.session` without revalidating the slot. Re-read the active ID (and liveness) after the await and retry the dial when the winner changed, otherwise callers can receive a closed session.</comment>

<file context>
@@ -200,6 +200,16 @@ actor CmxConnectivityPeerSession {
+            // caller that dialed in that window may have installed first;
+            // installing over it would leak its session and double-record
+            // an established lifecycle for the same peer.
+            if let installed = activeConnection {
+                if installed.id != pending.id {
+                    await connected.close()
</file context>

if installed.id != pending.id {
await connected.close()
}
return installed.session
}
install(
connected,
id: pending.id,
Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,65 @@
/// Authoritative route material whose change requires live session teardown.
///
/// Volatile freshness fields are excluded on purpose: `last_seen_at`, path
/// hints, direct ports, and display names move on every registration
/// heartbeat and shape only the next dial, never the trust of an already
/// admitted connection. Comparing this content lets a route revision bump
/// keep healthy sessions whose routes did not materially change.
struct CmxConnectivityRouteContent: Equatable, Sendable {
/// Trust material shared by every route in one account snapshot.
struct AccountMaterial: Equatable, Sendable {
let relayFleet: [String]
let lanRendezvous: CmxIrohLANRendezvous
let grantVerificationKeys: CmxIrohGrantVerificationKeySet
}

/// Admission-relevant material of one broker binding.
struct BindingMaterial: Equatable, Sendable {
let bindingID: String
let appInstanceID: String
let tag: String
let platform: CmxIrohPlatform
let identityGeneration: Int
let pairingEnabled: Bool
let capabilities: [String]

init(binding: CmxIrohBrokerBinding) {
bindingID = binding.bindingID
appInstanceID = binding.appInstanceID
tag = binding.tag
platform = binding.platform
identityGeneration = binding.identityGeneration
pairingEnabled = binding.pairingEnabled
capabilities = binding.capabilities

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

P2: A capability-order-only revision will invalidate the live peer session even though the advertised capability set is unchanged. Store a canonical ordering here so route-content equality matches the set semantics used by the admission policy.

Prompt for AI agents
Check if this issue is valid — if so, understand the root cause and fix it. At Packages/Shared/CmuxIrohTransport/Sources/CmuxIrohTransport/CmxConnectivityRouteContent.swift, line 33:

<comment>A capability-order-only revision will invalidate the live peer session even though the advertised capability set is unchanged. Store a canonical ordering here so route-content equality matches the set semantics used by the admission policy.</comment>

<file context>
@@ -0,0 +1,65 @@
+            platform = binding.platform
+            identityGeneration = binding.identityGeneration
+            pairingEnabled = binding.pairingEnabled
+            capabilities = binding.capabilities
+        }
+    }
</file context>

}
}

let account: AccountMaterial
private let peerRoutes: [CmxConnectivityPeerID: [BindingMaterial]]

init(snapshot: CmxIrohDiscoveryResponse) {
account = AccountMaterial(
relayFleet: snapshot.relayFleet,

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

P2: A revision that only reorders the same relay fleet will be treated as an account-material change, causing unnecessary teardown of all live sessions. Canonicalize the fleet order (or compare it as a set) when constructing route content so equivalent snapshots remain equivalent.

Prompt for AI agents
Check if this issue is valid — if so, understand the root cause and fix it. At Packages/Shared/CmuxIrohTransport/Sources/CmuxIrohTransport/CmxConnectivityRouteContent.swift, line 42:

<comment>A revision that only reorders the same relay fleet will be treated as an account-material change, causing unnecessary teardown of all live sessions. Canonicalize the fleet order (or compare it as a set) when constructing route content so equivalent snapshots remain equivalent.</comment>

<file context>
@@ -0,0 +1,65 @@
+
+    init(snapshot: CmxIrohDiscoveryResponse) {
+        account = AccountMaterial(
+            relayFleet: snapshot.relayFleet,
+            lanRendezvous: snapshot.lanRendezvous,
+            grantVerificationKeys: snapshot.grantVerificationKeys
</file context>
Suggested change
relayFleet: snapshot.relayFleet,
relayFleet: snapshot.relayFleet.sorted(),

lanRendezvous: snapshot.lanRendezvous,
grantVerificationKeys: snapshot.grantVerificationKeys

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

P2: A revision that only reorders equivalent grant verification keys will be treated as changed account trust material and tear down all live sessions. Canonicalize the key array by kid (or compare key sets order-independently) when building route content.

Prompt for AI agents
Check if this issue is valid — if so, understand the root cause and fix it. At Packages/Shared/CmuxIrohTransport/Sources/CmuxIrohTransport/CmxConnectivityRouteContent.swift, line 44:

<comment>A revision that only reorders equivalent grant verification keys will be treated as changed account trust material and tear down all live sessions. Canonicalize the key array by `kid` (or compare key sets order-independently) when building route content.</comment>

<file context>
@@ -0,0 +1,65 @@
+        account = AccountMaterial(
+            relayFleet: snapshot.relayFleet,
+            lanRendezvous: snapshot.lanRendezvous,
+            grantVerificationKeys: snapshot.grantVerificationKeys
+        )
+        var routes: [CmxConnectivityPeerID: [BindingMaterial]] = [:]
</file context>

)
var routes: [CmxConnectivityPeerID: [BindingMaterial]] = [:]
for binding in snapshot.bindings {
let peerID = CmxConnectivityPeerID(
identity: binding.endpointID,
deviceID: binding.deviceID
)
routes[peerID, default: []].append(BindingMaterial(binding: binding))
}
peerRoutes = routes.mapValues { bindings in
bindings.sorted { $0.bindingID < $1.bindingID }
}
}

/// Returns the material route for one peer, or nil when unrouted.
func peerRoute(
for peerID: CmxConnectivityPeerID
) -> [BindingMaterial]? {
peerRoutes[peerID]
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -88,7 +88,10 @@ extension CmxIrohClientRuntime {
try requireCurrent(revision)
guard published else { return .failed(.superseded) }
if let routeRevision = discovery.revision {
await connectivityEngine.didInstallRouteRevision(routeRevision)
await connectivityEngine.didInstallRouteRevision(
routeRevision,
routes: discovery
)
}
liveDiscoveryGeneration &+= 1
return .refreshed
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -376,7 +376,10 @@ public actor CmxIrohClientRuntime {
guard published else {
return .failed(.superseded)
}
await connectivityEngine.didInstallRouteRevision(discoveredRevision)
await connectivityEngine.didInstallRouteRevision(
discoveredRevision,
routes: discovery
)
liveDiscoveryGeneration &+= 1
return .refreshed
} catch {
Expand Down Expand Up @@ -520,7 +523,8 @@ public actor CmxIrohClientRuntime {
if published {
if let routeRevision = discovery.revision {
await connectivityEngine.didInstallRouteRevision(
routeRevision
routeRevision,
routes: discovery
)
}
liveDiscoveryGeneration &+= 1
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -459,7 +459,10 @@ extension CmxIrohHostRuntime {
await handleRoute(policy.binding, policy.routePathHints)
try requireCurrent(revision)
if let routeRevision = discovery.revision {
await connectivityEngine.didInstallRouteRevision(routeRevision)
await connectivityEngine.didInstallRouteRevision(
routeRevision,
routes: discovery
)
}
scheduleLANPublication(
binding: policy.binding,
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -149,7 +149,10 @@ extension CmxIrohHostRuntime {
try requireCurrent(revision)
await handleRoute(metadata, discovered.pathHints)
try requireCurrent(revision)
await connectivityEngine.didInstallRouteRevision(discoveredRevision)
await connectivityEngine.didInstallRouteRevision(
discoveredRevision,

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

P2: Overlapping host reconciliations can install an older discovery snapshot after a newer one because this call forwards the fetched revision without a monotonic check. That can roll back the engine's installed route revision and tear down or retain sessions based on stale route content; an engine-level monotonic install/coalescing guard would keep older completions from winning.

Prompt for AI agents
Check if this issue is valid — if so, understand the root cause and fix it. At Packages/Shared/CmuxIrohTransport/Sources/CmuxIrohTransport/CmxIrohHostRuntime+PublicAPI.swift, line 153:

<comment>Overlapping host reconciliations can install an older discovery snapshot after a newer one because this call forwards the fetched revision without a monotonic check. That can roll back the engine's installed route revision and tear down or retain sessions based on stale route content; an engine-level monotonic install/coalescing guard would keep older completions from winning.</comment>

<file context>
@@ -149,7 +149,10 @@ extension CmxIrohHostRuntime {
             try requireCurrent(revision)
-            await connectivityEngine.didInstallRouteRevision(discoveredRevision)
+            await connectivityEngine.didInstallRouteRevision(
+                discoveredRevision,
+                routes: discovery
+            )
</file context>

routes: discovery
)
scheduleLANPublication(
binding: metadata,
rendezvous: discovery.lanRendezvous,
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -321,7 +321,10 @@ public actor CmxIrohHostRuntime {
await handleBinding(registration, discovery, publishedPolicy.attestation)
try requireCurrent(revision)
if let routeRevision = discovery.revision {
await connectivityEngine.didInstallRouteRevision(routeRevision)
await connectivityEngine.didInstallRouteRevision(
routeRevision,
routes: discovery
)
}
scheduleRegistrationRenewal(
binding: registration.binding,
Expand Down
Loading