Skip to content
Closed
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
26 changes: 14 additions & 12 deletions Sources/Mobile/MobileBrowserStreamCoordinator.swift
Original file line number Diff line number Diff line change
@@ -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,
Expand All @@ -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)
Expand Down Expand Up @@ -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
Expand All @@ -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(
Expand All @@ -81,15 +83,15 @@ 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
}

@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
Expand All @@ -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
}
Expand Down
58 changes: 27 additions & 31 deletions Sources/Mobile/MobileBrowserStreamSession.swift
Original file line number Diff line number Diff line change
Expand Up @@ -218,7 +218,7 @@ final class MobileBrowserStreamSession {
scheduleDeadline(after: pacing.ackStallTimeout)
return
case .idle:
scheduleIdleReconciliation()
scheduleDeadline(after: Self.idleReconcileInterval, reconcile: true)
return
}
}
Expand Down Expand Up @@ -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<Void, Never> {
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()
}
}

Expand Down
140 changes: 48 additions & 92 deletions Sources/Mobile/MobileSimulatorStreamCoordinator.swift
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand All @@ -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)
}

Expand All @@ -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
}

Expand All @@ -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()
}
Expand All @@ -191,16 +125,23 @@ final class MobileSimulatorStreamCoordinator {
func requestFrameReplay(connectionID: UUID, panelIDStrings: Set<String>) {
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(
Expand All @@ -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
)
}

Expand Down
Loading
Loading