Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
17 commits
Select commit Hold shift + click to select a range
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 @@ -464,7 +464,10 @@ actor CmxConnectivityPeerSession {
activeConnection.closureTask?.cancel()
activeConnection.pathObservationTask?.cancel()
activeConnection.pathEventObservationTask?.cancel()
await activeConnection.pathEventObservationTask?.value
// Path-event diagnostics are observational. The session is already
// closed on this callback, so waiting for a cancelled observer here
// would retain control ownership and serialize the next dial behind
// an event stream that may not finish promptly.
await recordSessionClosure(
.remoteClosed,
active: activeConnection,
Expand Down Expand Up @@ -497,7 +500,10 @@ actor CmxConnectivityPeerSession {
activeConnection.pathObservationTask?.cancel()
activeConnection.pathEventObservationTask?.cancel()
await activeConnection.session.close()
await activeConnection.pathEventObservationTask?.value
// Path-event diagnostics are observational. They can outlive the
// physical session close while Iroh drains its event stream, but
// control ownership must be released as soon as the session itself
// is closed so a foreground handoff can admit the next owner.
await recordSessionClosure(
reason,
active: activeConnection,
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -107,6 +107,79 @@ struct CmxConnectivityPeerSessionTests {
await peer.releaseControl(ownerID: secondOwner)
}

@Test
func releaseDoesNotWaitForPathEventObserverToFinish() async throws {
let request = try Self.request()
let peerID = try CmxConnectivityPeerID(request: request)
let firstSession = TestConnectivitySession(
continuityID: 15,
keepsPathEventStreamOpen: true
)
let secondSession = TestConnectivitySession(continuityID: 16)
let builder = SequencedConnectivitySessionBuilder(
sessions: [firstSession, secondSession]
)
let peer = CmxConnectivityPeerSession(
peerID: peerID,
buildSession: { request in
try await builder.build(request)
},
diagnosticLog: DiagnosticLog(capacity: 16, role: .mobileClient)
)
let firstOwner = UUID()

_ = try await peer.acquireControl(for: request, ownerID: firstOwner)
try await Self.waitUntil { await firstSession.hasPathEventObserver() }

let release = Task {
await peer.releaseControl(ownerID: firstOwner)
}
try await Self.waitUntil { await firstSession.closeCount() == 1 }

let nextOwner = UUID()
let next = Task {
try await peer.acquireControl(for: request, ownerID: nextOwner)
}
try await Self.waitUntil { await builder.callCount() == 2 }
_ = try await next.value
await release.value
await peer.releaseControl(ownerID: nextOwner)
}

@Test
func remoteCloseDoesNotWaitForPathEventObserverToFinish() async throws {
let request = try Self.request()
let peerID = try CmxConnectivityPeerID(request: request)
let firstSession = TestConnectivitySession(
continuityID: 17,
keepsPathEventStreamOpen: true
)
let secondSession = TestConnectivitySession(continuityID: 18)
let builder = SequencedConnectivitySessionBuilder(
sessions: [firstSession, secondSession]
)
let peer = CmxConnectivityPeerSession(
peerID: peerID,
buildSession: { request in
try await builder.build(request)
},
diagnosticLog: DiagnosticLog(capacity: 16, role: .mobileClient)
)
let firstOwner = UUID()

_ = try await peer.acquireControl(for: request, ownerID: firstOwner)
try await Self.waitUntil { await firstSession.hasPathEventObserver() }
await firstSession.finishRemotely(failure: .connectionClosed)

let nextOwner = UUID()
let next = Task {
try await peer.acquireControl(for: request, ownerID: nextOwner)
}
try await Self.waitUntil { await builder.callCount() == 2 }
_ = try await next.value
await peer.releaseControl(ownerID: nextOwner)
}

@Test
func cancelledControlWaiterCannotBlockTheNextOwner() async throws {
let request = try Self.request()
Expand Down Expand Up @@ -779,6 +852,7 @@ private actor TestConnectivitySession: CmxConnectivitySession {
private let continuityID: UInt64
private let gatesCloseAttribution: Bool
private let keepsSelectedPathStreamOpen: Bool
private let keepsPathEventStreamOpen: Bool
private var closed = false
private var closes = 0
private var closeFailure = DiagnosticFailureKind.connectionClosed
Expand All @@ -795,17 +869,21 @@ private actor TestConnectivitySession: CmxConnectivitySession {
private var selectedPath = CmxIrohObservedConnectionPath.direct
private var selectedPathContinuation:
AsyncStream<CmxIrohObservedConnectionPath>.Continuation?
private var pathEventContinuation:
AsyncStream<CmxIrohConnectionPathEvent>.Continuation?

init(
continuityID: UInt64,
gatesCloseAttribution: Bool = false,
keepsSelectedPathStreamOpen: Bool = false,
keepsPathEventStreamOpen: Bool = false,
gatesFirstIsClosedCheck: Bool = false,
gatesFirstClose: Bool = false
) {
self.continuityID = continuityID
self.gatesCloseAttribution = gatesCloseAttribution
self.keepsSelectedPathStreamOpen = keepsSelectedPathStreamOpen
self.keepsPathEventStreamOpen = keepsPathEventStreamOpen
isClosedGatePending = gatesFirstIsClosedCheck
closeGatePending = gatesFirstClose
}
Expand Down Expand Up @@ -910,9 +988,17 @@ private actor TestConnectivitySession: CmxConnectivitySession {
}

func observedPathEvents() -> AsyncStream<CmxIrohConnectionPathEvent> {
AsyncStream { continuation in
continuation.finish()
let pair = AsyncStream<CmxIrohConnectionPathEvent>.makeStream()
guard keepsPathEventStreamOpen else {
pair.continuation.finish()
return pair.stream
}
pathEventContinuation = pair.continuation
return pair.stream
}

func hasPathEventObserver() -> Bool {
pathEventContinuation != nil
}

func close() async {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -2163,12 +2163,14 @@ public final class MobileShellComposite: MobileTerminalOutputSinking {
pairedMacAliasIDsByRepresentativeID = [:]
pairedMacs = []
pairedMacLoadState = .notLoaded
pairedMacLoadGeneration &+= 1
hiddenComputers = []
hasHiddenComputers = false
resetTerminalThemes()
// Likewise drop the registry-backed device tree so a shared device never
// shows the previous user's team devices after sign-out.
registryDevices = []
registryDevicesLoadGeneration &+= 1
// Reset the in-memory restoring flags; hasKnownPairedMac stays driven by
// the hide path. On a real account switch the next reconnect's no-mac
// branch clears the hint. Bump the reconnect generation so any in-flight
Expand Down Expand Up @@ -2286,10 +2288,12 @@ public final class MobileShellComposite: MobileTerminalOutputSinking {
pairedMacAliasIDsByRepresentativeID = [:]
pairedMacs = []
pairedMacLoadState = .notLoaded
pairedMacLoadGeneration &+= 1
hiddenMacDeviceIDsByScope = [:]
hiddenComputers = []
hasHiddenComputers = false
registryDevices = []
registryDevicesLoadGeneration &+= 1
teamScopeCleanupTask?.cancel()
teamScopeCleanupTask = Task {
if let refresher {
Expand Down Expand Up @@ -3449,6 +3453,9 @@ public final class MobileShellComposite: MobileTerminalOutputSinking {
var hasStoredUsableTailscaleAuthorization = false
/// Load status for ``pairedMacs`` in the current signed-in account/team scope.
public internal(set) var pairedMacLoadState: PairedMacLoadState = .notLoaded
/// Monotonic token so overlapping same-scope loads cannot publish an older
/// snapshot after a newer refresh has started.
private var pairedMacLoadGeneration: UInt64 = 0
/// Visible representative id to all stored ids for that logical paired Mac.
public private(set) var pairedMacAliasIDsByRepresentativeID: [String: [String]] = [:]
/// Cached device-local hidden ids keyed by signed-in account/team scope.
Expand Down Expand Up @@ -3565,6 +3572,8 @@ public final class MobileShellComposite: MobileTerminalOutputSinking {
/// known paired Macs, so the tree degrades to the same hosts the switcher
/// shows rather than going blank.
public internal(set) var registryDevices: [RegistryDevice] = []
/// Monotonic token so overlapping registry requests are latest-wins.
private var registryDevicesLoadGeneration: UInt64 = 0

/// The cmux device id of the Mac the live connection currently targets, or
/// `nil` when not connected. Used by the device tree to mark which device row
Expand Down Expand Up @@ -3614,10 +3623,13 @@ public final class MobileShellComposite: MobileTerminalOutputSinking {
/// leads with the host the user is on. Mirrors ``loadPairedMacs()``: signed
/// out yields an empty list.
public func loadRegistryDevices() async {
registryDevicesLoadGeneration &+= 1
let loadGeneration = registryDevicesLoadGeneration
let startedAt = appDiagnosticNow()
recordAppEvent(.deviceRegistryLoadStarted)
guard let deviceRegistry,
let scope = await currentScopeSnapshot() else {
guard loadGeneration == registryDevicesLoadGeneration else { return }
registryDevices = []
recordAppEvent(
.deviceRegistryLoadFailed,
Expand All @@ -3639,7 +3651,8 @@ public final class MobileShellComposite: MobileTerminalOutputSinking {
// requesting user still being current (mirroring the `.ok` path):
// a stale 401 from a signed-out session that lands after a
// different user signed in must not blank the new user's tree.
if await isScopeCurrent(scope) {
if loadGeneration == registryDevicesLoadGeneration,
await isScopeCurrent(scope) {
registryDevices = []
}
recordAppEvent(
Expand All @@ -3662,10 +3675,12 @@ public final class MobileShellComposite: MobileTerminalOutputSinking {
// are still in the same signed-in account/team scope, so a slow load can
// never repopulate another scope's devices after sign-out, account switch,
// or same-account team switch.
guard await isScopeCurrent(scope) else { return }
guard loadGeneration == registryDevicesLoadGeneration,
await isScopeCurrent(scope) else { return }
let connectedID = connectedMacDeviceID
let hiddenIDs = await hiddenMacDeviceIDs(scope: scope)
guard await isScopeCurrent(scope) else { return }
guard loadGeneration == registryDevicesLoadGeneration,
await isScopeCurrent(scope) else { return }
let compatible = compatibleRegistryDevices(loaded)
registryDevices = compatible
.compactMap { device in
Expand Down Expand Up @@ -3941,6 +3956,8 @@ public final class MobileShellComposite: MobileTerminalOutputSinking {
/// back to the unscoped all-users query, so a shared device never exposes
/// another user's Macs in the switcher.
public func loadPairedMacs() async {
pairedMacLoadGeneration &+= 1
let loadGeneration = pairedMacLoadGeneration
// The demo-content paired-Mac decorator reads the account's
// demonstration flag lazily on every load, so any load can reveal the
// Demo Mac row. Re-evaluate activation at the same moment: the flag
Expand All @@ -3952,6 +3969,7 @@ public final class MobileShellComposite: MobileTerminalOutputSinking {
recordAppEvent(.computerListRefreshStarted)
guard let pairedMacStore,
let scope = await currentScopeSnapshot() else {
guard loadGeneration == pairedMacLoadGeneration else { return }
storedPairedMacs = []
clearStoredPairedMacCache()
pairedMacAliasIDsByRepresentativeID = [:]
Expand All @@ -3966,13 +3984,15 @@ public final class MobileShellComposite: MobileTerminalOutputSinking {
)
return
}
guard loadGeneration == pairedMacLoadGeneration else { return }
pairedMacLoadState = .notLoaded
let loaded: [MobilePairedMac]
do {
loaded = try await pairedMacStore.loadAll(stackUserID: scope.userID, teamID: scope.teamID)
} catch {
mobileShellLog.error("paired mac store loadAll failed: \(String(describing: error), privacy: .public)")
if await isScopeCurrent(scope) {
if loadGeneration == pairedMacLoadGeneration,
await isScopeCurrent(scope) {
pairedMacLoadState = .failed
hiddenComputers = []
hasHiddenComputers = false
Expand All @@ -3992,7 +4012,8 @@ public final class MobileShellComposite: MobileTerminalOutputSinking {
// The await above suspended the main actor; a sign-out, user switch, or
// same-account team switch may have run meanwhile. Discard unless the
// captured account/team scope is still current.
guard await isScopeCurrent(scope) else {
guard loadGeneration == pairedMacLoadGeneration,
await isScopeCurrent(scope) else {
return
}
migrateLegacyWorkspaceComputerPriority(loadedMacs: loaded)
Expand All @@ -4006,7 +4027,8 @@ public final class MobileShellComposite: MobileTerminalOutputSinking {
from: loaded,
hiddenIDs: hiddenIDs
)
guard await isScopeCurrent(scope) else {
guard loadGeneration == pairedMacLoadGeneration,
await isScopeCurrent(scope) else {
return
}
installStoredPairedMacCache(loaded, scope: scope)
Expand Down Expand Up @@ -13976,8 +13998,29 @@ public final class MobileShellComposite: MobileTerminalOutputSinking {
deadline.schedule(deadline: .now() + .nanoseconds(Int(clamping: timeoutNanoseconds)))
deadline.setEventHandler { probe.cancel() }
deadline.resume()
let ack = await probe.value
var ack = await probe.value
deadline.cancel()

if case .failed = ack {
// A probe timeout is weaker evidence than a failed subscription
// round-trip. The keepalive lane can still be healthy while one
// control request stalls during Iroh path migration. Give the
// idempotent repair two fresh, independently bounded attempts
// before promoting this suspicion to a session replacement. The
// host-side operation is idempotent and preserves the live reader.
for _ in 0..<2 {
let retry = await requestTerminalEventSubscription(
client: client,
reason: "liveness_probe_retry",
topics: topics,
timeoutNanoseconds: timeoutNanoseconds
)
ack = retry
if retry.isSubscribed {
break
}
}
}
return ack
}

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -412,6 +412,32 @@ import Testing
#expect(store.registryDevices.map(\.deviceId) == ["device-b"])
}

@Test func staleSameScopeRegistryLoadCannotReplaceNewerSnapshot() async throws {
let registry = SequencedDeviceRegistry(
outcomes: [
.ok([Self.registryDevice(id: "old-device")]),
.ok([Self.registryDevice(id: "new-device")]),
]
)
let store = MobileShellComposite(
isSignedIn: true,
deviceRegistry: registry,
identityProvider: StaticIdentityProvider(userID: "user-1"),
teamIDProvider: { "team-a" }
)

let oldLoad = Task { await store.loadRegistryDevices() }
await registry.waitUntilCall(1)
let newLoad = Task { await store.loadRegistryDevices() }
await registry.waitUntilCall(2)

await registry.releaseFirstCall()
await oldLoad.value
await newLoad.value

#expect(store.registryDevices.map(\.deviceId) == ["new-device"])
}

@Test func teamChangeDoesNotStartACompetingStoredMacReconnect() async throws {
let team = MutableTeamID("team-a")
let pairedStore = DelayedTeamPairedMacStore(
Expand Down Expand Up @@ -503,7 +529,11 @@ import Testing
platform: "mac",
displayName: id,
lastSeenAt: Date(timeIntervalSince1970: 2),
instances: []
instances: [RegistryAppInstance(
tag: "default",
routes: [],
lastSeenAt: Date(timeIntervalSince1970: 2)
)]
)
}

Expand Down Expand Up @@ -1286,6 +1316,7 @@ import Testing
#expect(route?.0 == "100.71.210.41")
#expect(route?.1 == CmxMobileDefaults.defaultHostPort)
}

}

private func hostPortRoute(
Expand Down
Loading