diff --git a/Sources/Mobile/MobileBrowserStreamCoordinator.swift b/Sources/Mobile/MobileBrowserStreamCoordinator.swift index 139088f20582..2dbb77946955 100644 --- a/Sources/Mobile/MobileBrowserStreamCoordinator.swift +++ b/Sources/Mobile/MobileBrowserStreamCoordinator.swift @@ -1,14 +1,16 @@ import CMUXMobileCore import Foundation +/// Keys one per-connection, per-panel mobile stream session; shared by the +/// browser and simulator stream coordinators. +struct MobileStreamSessionKey: Hashable { + let connectionID: UUID + let panelID: UUID +} + @MainActor final class MobileBrowserStreamCoordinator { - private struct SessionKey: Hashable { - let connectionID: UUID - let panelID: UUID - } - - private var sessions: [SessionKey: MobileBrowserStreamSession] = [:] + private var sessions: [MobileStreamSessionKey: MobileBrowserStreamSession] = [:] func start( connectionID: UUID, @@ -18,7 +20,7 @@ final class MobileBrowserStreamCoordinator { guard let connection = MobileHostConnectionRegistry.shared.connection(id: connectionID) else { return nil } - let key = SessionKey(connectionID: connectionID, panelID: panel.id) + let key = MobileStreamSessionKey(connectionID: connectionID, panelID: panel.id) let replacedExisting = sessions[key] != nil if let previous = sessions.removeValue(forKey: key) { await previous.stop(sendClosed: false) @@ -49,7 +51,7 @@ final class MobileBrowserStreamCoordinator { } func hasStream(connectionID: UUID, panelID: UUID) -> Bool { - sessions[SessionKey(connectionID: connectionID, panelID: panelID)] != nil + sessions[MobileStreamSessionKey(connectionID: connectionID, panelID: panelID)] != nil } @discardableResult @@ -68,7 +70,7 @@ final class MobileBrowserStreamCoordinator { @discardableResult func stop(connectionID: UUID, panelID: UUID) async -> Bool { - let key = SessionKey(connectionID: connectionID, panelID: panelID) + let key = MobileStreamSessionKey(connectionID: connectionID, panelID: panelID) guard let session = sessions.removeValue(forKey: key) else { return false } await session.stop(sendClosed: false) MobileHostIrohRuntime.hostDiagnosticLog.record(DiagnosticEvent( @@ -81,7 +83,7 @@ final class MobileBrowserStreamCoordinator { @discardableResult func acknowledge(connectionID: UUID, panelID: UUID, sequence: UInt64) -> Bool { - let key = SessionKey(connectionID: connectionID, panelID: panelID) + let key = MobileStreamSessionKey(connectionID: connectionID, panelID: panelID) guard let session = sessions[key] else { return false } session.acknowledge(sequence: sequence) return true @@ -89,7 +91,7 @@ final class MobileBrowserStreamCoordinator { @discardableResult func noteInputReplayed(connectionID: UUID, panelID: UUID) -> Bool { - let key = SessionKey(connectionID: connectionID, panelID: panelID) + let key = MobileStreamSessionKey(connectionID: connectionID, panelID: panelID) guard let session = sessions[key] else { return false } session.noteInputReplayed() return true @@ -103,7 +105,7 @@ final class MobileBrowserStreamCoordinator { } } - private func sessionEnded(key: SessionKey, sessionID: UUID) { + private func sessionEnded(key: MobileStreamSessionKey, sessionID: UUID) { guard sessions[key]?.id == sessionID else { return } sessions[key] = nil } diff --git a/Sources/Mobile/MobileBrowserStreamSession.swift b/Sources/Mobile/MobileBrowserStreamSession.swift index e9a3a60c8902..fcb13a3a2e87 100644 --- a/Sources/Mobile/MobileBrowserStreamSession.swift +++ b/Sources/Mobile/MobileBrowserStreamSession.swift @@ -218,7 +218,7 @@ final class MobileBrowserStreamSession { scheduleDeadline(after: pacing.ackStallTimeout) return case .idle: - scheduleIdleReconciliation() + scheduleDeadline(after: Self.idleReconcileInterval, reconcile: true) return } } @@ -328,53 +328,49 @@ final class MobileBrowserStreamSession { } } - private func scheduleDeadline(after interval: TimeInterval) { - guard !isStopped else { return } - deadlineTask?.cancel() + /// Runs `body` after a bounded, cancellable sleep. Every deadline in this + /// session (cadence/settle, idle reconciliation, state coalescing) goes + /// through this one task shape; storing the returned task and cancelling + /// it on a fresh signal is the caller's job. + private func scheduleAfter( + _ interval: TimeInterval, + _ body: @escaping @MainActor () async -> Void + ) -> Task { let clock = clock - deadlineTask = Task { @MainActor [weak self, clock] in + return Task { @MainActor [clock] in do { - // Bounded, cancellable cadence/settle deadline; new signals cancel it. try await clock.sleep(for: max(0, interval)) guard !Task.isCancelled else { return } - self?.requestDrive() + await body() } catch {} } } - /// Schedules one lossless reconciliation frame after a quiet interval. - /// - /// An idle stream's correctness otherwise hangs entirely on page-driven - /// dirty signals, which are silently lost when WebKit suspends - /// `requestAnimationFrame` for the occluded offscreen host. This bounds - /// how long the phone can display a stale (or blank first) frame to - /// `idleReconcileInterval`. Any real dirty signal cancels it via - /// `requestDrive`, and a fresh one is armed when the stream idles again. - private func scheduleIdleReconciliation() { + /// Arms the drive deadline. With `reconcile`, expiry also requests one + /// lossless settle frame: an idle stream's dirty signals are silently + /// lost when WebKit suspends `requestAnimationFrame` for the occluded + /// offscreen host, so this bounds how long the phone can display a stale + /// (or blank first) frame to `idleReconcileInterval`. Any real dirty + /// signal cancels the deadline via `requestDrive`, and `drive()` arms a + /// fresh one when the stream idles again. + private func scheduleDeadline(after interval: TimeInterval, reconcile: Bool = false) { guard !isStopped else { return } deadlineTask?.cancel() - let clock = clock - deadlineTask = Task { @MainActor [weak self, clock] in - do { - try await clock.sleep(for: Self.idleReconcileInterval) - guard !Task.isCancelled, let self, !self.isStopped else { return } + deadlineTask = scheduleAfter(interval) { [weak self] in + guard let self, !self.isStopped else { return } + if reconcile { self.pacing.requestSettleReconciliation() - self.requestDrive() - } catch {} + } + self.requestDrive() } } private func scheduleStateEmission() { guard !isStopped else { return } stateTask?.cancel() - let clock = clock - stateTask = Task { @MainActor [weak self, clock] in - do { - // Bounded, cancellable coalescing delay for bursty WebKit state KVO. - try await clock.sleep(for: 0.016) - guard !Task.isCancelled else { return } - await self?.emitStateIfChanged() - } catch {} + // Coalesces bursty WebKit state KVO behind one bounded delay. + stateTask = scheduleAfter(0.016) { [weak self] in + await self?.emitStateIfChanged() } } diff --git a/Sources/Mobile/MobileSimulatorStreamCoordinator.swift b/Sources/Mobile/MobileSimulatorStreamCoordinator.swift index b6ca0277b36b..70af1f9c31e0 100644 --- a/Sources/Mobile/MobileSimulatorStreamCoordinator.swift +++ b/Sources/Mobile/MobileSimulatorStreamCoordinator.swift @@ -9,65 +9,38 @@ final class MobileSimulatorStreamCoordinator { case unavailable } - private struct SessionKey: Hashable { - let connectionID: UUID - let panelID: UUID - } - - private var sessions: [SessionKey: MobileSimulatorStreamSession] = [:] + private var sessions: [MobileStreamSessionKey: MobileSimulatorStreamSession] = [:] private var ownersByPanelID: [UUID: UUID] = [:] private var workspaceIDsByPanelID: [UUID: UUID] = [:] private var lastFramesByPanelID: [UUID: MobileSimulatorFrameEvent] = [:] private let wireEncoder = MobileSimulatorWireEncoder() func start(connectionID: UUID, panel: SimulatorPanel, workspaceID: UUID) async -> StartResult { - MobileSimulatorDiagnostics.recordStream( - panelID: panel.id, - state: .startRequested, - ownership: MobileSimulatorDiagnostics.ownershipState( - ownerConnectionID: ownersByPanelID[panel.id], - currentConnectionID: connectionID - ), - activeSessions: sessions.count + recordStream( + panel.id, + .startRequested, + ownership: currentOwnership(panelID: panel.id, connectionID: connectionID) ) guard let connection = MobileHostConnectionRegistry.shared.connection(id: connectionID) else { - MobileSimulatorDiagnostics.recordStream( - panelID: panel.id, - state: .startFailed, - ownership: .unknown, - activeSessions: sessions.count - ) + recordStream(panel.id, .startFailed, ownership: .unknown) return .unavailable } workspaceIDsByPanelID[panel.id] = workspaceID if let owner = ownersByPanelID[panel.id], owner != connectionID { guard let descriptor = descriptor(panel: panel, currentConnectionID: connectionID) else { - MobileSimulatorDiagnostics.recordStream( - panelID: panel.id, - state: .startFailed, - ownership: .otherConnection, - activeSessions: sessions.count - ) + recordStream(panel.id, .startFailed, ownership: .otherConnection) return .unavailable } - MobileSimulatorDiagnostics.recordStream( - panelID: panel.id, - state: .locked, - ownership: .otherConnection, - activeSessions: sessions.count - ) + recordStream(panel.id, .locked, ownership: .otherConnection) return .locked(descriptor) } - let key = SessionKey(connectionID: connectionID, panelID: panel.id) + let key = MobileStreamSessionKey(connectionID: connectionID, panelID: panel.id) if let previous = sessions.removeValue(forKey: key) { await previous.stop(sendClosed: false) } - let previousOwnership = MobileSimulatorDiagnostics.ownershipState( - ownerConnectionID: ownersByPanelID[panel.id], - currentConnectionID: connectionID - ) + let previousOwnership = currentOwnership(panelID: panel.id, connectionID: connectionID) ownersByPanelID[panel.id] = connectionID MobileSimulatorDiagnostics.recordOwnership( panelID: panel.id, @@ -94,20 +67,10 @@ final class MobileSimulatorStreamCoordinator { session.start() pruneCachesForClosedPanels() guard let descriptor = descriptor(panel: panel, currentConnectionID: connectionID) else { - MobileSimulatorDiagnostics.recordStream( - panelID: panel.id, - state: .startFailed, - ownership: .currentConnection, - activeSessions: sessions.count - ) + recordStream(panel.id, .startFailed, ownership: .currentConnection) return .unavailable } - MobileSimulatorDiagnostics.recordStream( - panelID: panel.id, - state: .started, - ownership: .currentConnection, - activeSessions: sessions.count - ) + recordStream(panel.id, .started, ownership: .currentConnection) return .started(descriptor) } @@ -133,33 +96,17 @@ final class MobileSimulatorStreamCoordinator { @discardableResult func stop(connectionID: UUID, panelID: UUID) async -> Bool { - let key = SessionKey(connectionID: connectionID, panelID: panelID) + let key = MobileStreamSessionKey(connectionID: connectionID, panelID: panelID) guard let session = sessions.removeValue(forKey: key) else { return false } - MobileSimulatorDiagnostics.recordStream( - panelID: panelID, - state: .stopRequested, - ownership: MobileSimulatorDiagnostics.ownershipState( - ownerConnectionID: ownersByPanelID[panelID], - currentConnectionID: connectionID - ), + recordStream( + panelID, + .stopRequested, + ownership: currentOwnership(panelID: panelID, connectionID: connectionID), activeSessions: sessions.count + 1 ) await session.stop(sendClosed: false) - if ownersByPanelID[panelID] == connectionID { - ownersByPanelID[panelID] = nil - MobileSimulatorDiagnostics.recordOwnership( - panelID: panelID, - ownership: .unowned, - previousOwnership: .currentConnection - ) - } + releaseOwnership(key, recording: .stopped) pruneCachesForClosedPanels() - MobileSimulatorDiagnostics.recordStream( - panelID: panelID, - state: .stopped, - ownership: .unowned, - activeSessions: sessions.count - ) return true } @@ -168,20 +115,7 @@ final class MobileSimulatorStreamCoordinator { for (key, session) in matching { sessions[key] = nil await session.stop(sendClosed: false) - if ownersByPanelID[key.panelID] == connectionID { - ownersByPanelID[key.panelID] = nil - MobileSimulatorDiagnostics.recordOwnership( - panelID: key.panelID, - ownership: .unowned, - previousOwnership: .currentConnection - ) - } - MobileSimulatorDiagnostics.recordStream( - panelID: key.panelID, - state: .closed, - ownership: .unowned, - activeSessions: sessions.count - ) + releaseOwnership(key, recording: .closed) } pruneCachesForClosedPanels() } @@ -191,16 +125,23 @@ final class MobileSimulatorStreamCoordinator { func requestFrameReplay(connectionID: UUID, panelIDStrings: Set) { for panelIDString in panelIDStrings { guard let panelID = UUID(uuidString: panelIDString), - let session = sessions[SessionKey(connectionID: connectionID, panelID: panelID)] else { + let session = sessions[MobileStreamSessionKey(connectionID: connectionID, panelID: panelID)] else { continue } session.requestFrameReplay() } } - private func sessionEnded(key: SessionKey, sessionID: UUID) { + private func sessionEnded(key: MobileStreamSessionKey, sessionID: UUID) { guard sessions[key]?.id == sessionID else { return } sessions[key] = nil + releaseOwnership(key, recording: .closed) + pruneCachesForClosedPanels() + } + + /// Releases `key`'s panel ownership when its connection still holds it, + /// then records the stream lifecycle transition for the removed session. + private func releaseOwnership(_ key: MobileStreamSessionKey, recording state: DiagnosticSimulatorStreamLifecycle) { if ownersByPanelID[key.panelID] == key.connectionID { ownersByPanelID[key.panelID] = nil MobileSimulatorDiagnostics.recordOwnership( @@ -209,12 +150,27 @@ final class MobileSimulatorStreamCoordinator { previousOwnership: .currentConnection ) } - pruneCachesForClosedPanels() + recordStream(key.panelID, state, ownership: .unowned) + } + + private func recordStream( + _ panelID: UUID, + _ state: DiagnosticSimulatorStreamLifecycle, + ownership: DiagnosticSimulatorOwnershipState, + activeSessions: Int? = nil + ) { MobileSimulatorDiagnostics.recordStream( - panelID: key.panelID, - state: .closed, - ownership: .unowned, - activeSessions: sessions.count + panelID: panelID, + state: state, + ownership: ownership, + activeSessions: activeSessions ?? sessions.count + ) + } + + private func currentOwnership(panelID: UUID, connectionID: UUID) -> DiagnosticSimulatorOwnershipState { + MobileSimulatorDiagnostics.ownershipState( + ownerConnectionID: ownersByPanelID[panelID], + currentConnectionID: connectionID ) } diff --git a/Sources/Mobile/MobileSimulatorStreamSession.swift b/Sources/Mobile/MobileSimulatorStreamSession.swift index 9ca54b5330fd..3f00e6bee25f 100644 --- a/Sources/Mobile/MobileSimulatorStreamSession.swift +++ b/Sources/Mobile/MobileSimulatorStreamSession.swift @@ -101,11 +101,7 @@ final class MobileSimulatorStreamSession { && !changedTopics.contains("simulator.state") { return } - MobileSimulatorDiagnostics.recordFrame( - panelID: self.panelID, - state: .subscriptionReasserted, - sequence: self.lastSentSequence - ) + self.recordFrame(.subscriptionReasserted, sequence: self.lastSentSequence) self.emitState() self.requestFrameSend() } @@ -173,11 +169,11 @@ final class MobileSimulatorStreamSession { _ = refresh.detachedReader?.setFramePublicationHandler(nil) guard let reader = refresh.attachedReader else { if refresh.isMissing { - MobileSimulatorDiagnostics.recordFrame(panelID: panelID, state: .readerMissing) + recordFrame(.readerMissing) } return } - MobileSimulatorDiagnostics.recordFrame(panelID: panelID, state: .readerAttached) + recordFrame(.readerAttached) _ = reader.setFramePublicationHandler { [weak self] in Task { @MainActor [weak self] in self?.requestFrameSend() @@ -218,12 +214,7 @@ final class MobileSimulatorStreamSession { guard let event = await nextFrameEvent(allowCachedReplay: shouldReplay) else { return } let payloadBytes = event.dataBase64.utf8.count guard let payload = wireEncoder.object(event) else { - MobileSimulatorDiagnostics.recordFrame( - panelID: panelID, - state: .encodeFailed, - sequence: event.sequence, - payloadBytes: payloadBytes - ) + recordFrame(.encodeFailed, sequence: event.sequence, payloadBytes: payloadBytes) continue } let delivered = await connection.sendEvent(topic: "simulator.frame", payload: payload) @@ -233,20 +224,10 @@ final class MobileSimulatorStreamSession { // healthy. Keep the control session alive; a later frame // publication or subscription reassertion will retry the // current latest frame because `lastSentSequence` is unchanged. - MobileSimulatorDiagnostics.recordFrame( - panelID: panelID, - state: .refused, - sequence: event.sequence, - payloadBytes: payloadBytes - ) + recordFrame(.refused, sequence: event.sequence, payloadBytes: payloadBytes) return } - MobileSimulatorDiagnostics.recordFrame( - panelID: panelID, - state: .sent, - sequence: event.sequence, - payloadBytes: payloadBytes - ) + recordFrame(.sent, sequence: event.sequence, payloadBytes: payloadBytes) lastSentSequence = event.sequence cachedFrame = event onFrame(panelID, event) @@ -256,7 +237,7 @@ final class MobileSimulatorStreamSession { private func nextFrameEvent(allowCachedReplay: Bool = false) async -> MobileSimulatorFrameEvent? { guard let reader = readerAttachment.reader else { if allowCachedReplay { return cachedFrame } - MobileSimulatorDiagnostics.recordFrame(panelID: panelID, state: .readerMissing) + recordFrame(.readerMissing) return nil } guard reader.hasPublishedFrame(after: lastSentSequence) else { @@ -265,12 +246,7 @@ final class MobileSimulatorStreamSession { guard let frame = await reader.copyLatestFrame(after: lastSentSequence) else { return allowCachedReplay ? cachedFrame : nil } - MobileSimulatorDiagnostics.recordFrame( - panelID: panelID, - state: .copied, - sequence: frame.sequence, - payloadBytes: frame.data.count - ) + recordFrame(.copied, sequence: frame.sequence, payloadBytes: frame.data.count) let format: MobileSimulatorFrameFormat switch frame.format { case .jpeg: @@ -295,15 +271,27 @@ final class MobileSimulatorStreamSession { guard let self, !self.isStopped else { return } guard let payload = self.wireEncoder.object(cachedFrame) else { return } let delivered = await self.connection.sendEvent(topic: "simulator.frame", payload: payload) - MobileSimulatorDiagnostics.recordFrame( - panelID: self.panelID, - state: delivered ? .cachedSent : .refused, + self.recordFrame( + delivered ? .cachedSent : .refused, sequence: cachedFrame.sequence, payloadBytes: cachedFrame.dataBase64.utf8.count ) } } + private func recordFrame( + _ state: DiagnosticSimulatorFrameLifecycle, + sequence: UInt64? = nil, + payloadBytes: Int? = nil + ) { + MobileSimulatorDiagnostics.recordFrame( + panelID: panelID, + state: state, + sequence: sequence, + payloadBytes: payloadBytes + ) + } + private func emitState() { guard !isStopped, let descriptor = descriptorProvider(connectionID),