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
63 changes: 63 additions & 0 deletions Sources/Surfaces/CloudWorkspaceProjectionCoordinator.swift
Original file line number Diff line number Diff line change
Expand Up @@ -43,6 +43,7 @@ final class CloudWorkspaceProjectionCoordinator {
guard let self else { return }
defer { if self.tasks[machine]?.id == id { self.tasks[machine] = nil } }
guard let catalog else { return }
var budget = CloudWorkspaceReconcileBudget()
while self.requested.remove(machine) != nil {
guard !Task.isCancelled, self.localMutations[machine] == nil else { return }
await catalog.cloudPlacementCoordinator.waitForPendingMutations()
Expand All @@ -51,12 +52,44 @@ final class CloudWorkspaceProjectionCoordinator {
let state = catalog.cloudStates[machine],
catalog.cloudStateObservations[machine]?.freshness == .current,
catalog.cloudPlacementCoordinator.allowsNativeReconciliation(state) else { return }
let progress = CloudWorkspaceReconcileBudget.Mark(
state: state,
projectionVersion: catalog.projectionVersions[machine, default: 0],
bindings: self.environment.bindings().filter { $0.value.vmID == machine.rawValue }
)
guard budget.admit(progress) else {
self.reportNonConvergence(machine: machine, state: state, budget: budget)
return
}
await self.reconcile(state: state, catalog: catalog)
}
}
tasks[machine] = CloudWorkspaceProjectionTask(id: id, task: task)
}

/// Daemon generation last reported per machine. A persistent re-requester
/// would otherwise report once for every graph revision.
private var reportedNonConvergence: [SurfaceMachineID: String] = [:]

/// Reports non-convergence once per daemon generation, so a trigger outside
/// this loop that keeps restarting it cannot flood crash reporting.
private func reportNonConvergence(machine: SurfaceMachineID, state: CloudVMState, budget: CloudWorkspaceReconcileBudget) {
#if DEBUG
cmuxDebugLog("cloudWorkspace.projection.nonConvergent machine=\(machine.rawValue) passes=\(budget.passes) idle=\(budget.idlePasses)")
#endif
let generation = state.cursor?.generation ?? ""
guard reportedNonConvergence[machine] != generation else { return }
reportedNonConvergence[machine] = generation
sentryCaptureWarning(
"Cloud workspace projection did not converge",
category: "cloud.projection",
data: [
"passes": budget.passes, "idlePasses": budget.idlePasses,
"generation": generation, "revision": state.cursor?.revision ?? 0,
]
)
}

/// A bound mirror may not recreate a view that the accepted graph removed.
/// Unbound viewers retain their existing attachment-repair behavior.
func retainsProjection(_ projection: SurfaceProjection, in state: CloudVMState) -> Bool {
Expand Down Expand Up @@ -96,6 +129,7 @@ final class CloudWorkspaceProjectionCoordinator {
tasks.removeValue(forKey: machine)?.task.cancel()
requested.remove(machine)
localMutations[machine] = nil
reportedNonConvergence[machine] = nil
let bindings = environment.bindings()
failures = failures.filter { bindings[$0.key]?.vmID != machine.rawValue }
}
Expand Down Expand Up @@ -170,3 +204,32 @@ final class CloudWorkspaceProjectionCoordinator {
failures = failures.filter { live.contains($0.key) }
}
}

/// Bounds reconciliation of one accepted graph. A pass that changes nothing
/// (same graph, projections and bindings as the pass before) cannot make the
/// next one different, so a few in a row mean a consumer is requesting passes
/// without progress. The hard ceiling also stops a loop that rewrites the same
/// projections every pass, which looks like progress.
struct CloudWorkspaceReconcileBudget {
struct Mark: Equatable {
let state: CloudVMState
let projectionVersion: UInt64
let bindings: [UUID: WorkspaceCloudVMBinding]
}

static let maxIdlePasses = 3
static let maxPassesPerState = 64

private var last: Mark?
private(set) var passes = 0
private(set) var idlePasses = 0

/// Records the catalog before a pass; false means stop reconciling this graph.
mutating func admit(_ mark: Mark) -> Bool {
if last?.state != mark.state { passes = 0; idlePasses = 0 }
idlePasses = last == mark ? idlePasses + 1 : 0
passes += 1
last = mark
return idlePasses < Self.maxIdlePasses && passes <= Self.maxPassesPerState
}
}
36 changes: 19 additions & 17 deletions Sources/Surfaces/SurfaceCatalog.swift
Original file line number Diff line number Diff line change
Expand Up @@ -1245,16 +1245,21 @@ final class SurfaceCatalog {
}

func setRemotePlacement(for source: SurfaceProjection, workspaceID: String?, tabID: String?) {
// Only views whose coordinates change; an unchanged placement notifies nobody.
let matching = projections.filter {
$0.resource == source.resource && ($0.panelID == source.panelID
|| (tabID != nil && $0.remoteTabID == tabID))
&& ($0.remoteWorkspaceID != workspaceID || $0.remoteTabID != tabID)
}
guard !matching.isEmpty else { return }
var next = projections
for var projection in matching {
projections.remove(projection)
next.remove(projection)
projection.remoteWorkspaceID = workspaceID
projection.remoteTabID = tabID
projections.insert(projection)
next.insert(projection)
}
projections = next
notifyChange(for: source.resource.machine)
}

Expand All @@ -1264,23 +1269,20 @@ final class SurfaceCatalog {
@discardableResult
private func attachRemoteView(_ view: SurfaceRemoteView?, to projection: SurfaceProjection) -> SurfaceProjection {
guard let view else { return projection }
if view.isCloudDisplayMembershipView {
guard projection.remoteTabID == nil else { return projection }
projections.remove(projection)
var updated = projection
updated.remoteWorkspaceID = view.workspace.id
updated.remoteTabID = nil
projections.insert(updated)
reconcileCloudWorkspaceBinding(localWorkspaceID: updated.workspaceID)
notifyChange(for: updated.resource.machine)
return updated
}
guard projection.remoteTabID == nil || projection.remoteTabID == view.tabID else { return projection }
projections.remove(projection)
// A display membership view attaches only to a tabless preview.
let tabID = view.isCloudDisplayMembershipView ? nil : view.tabID
guard projection.remoteTabID == nil || projection.remoteTabID == tabID else { return projection }
var updated = projection
updated.remoteWorkspaceID = view.workspace.id
updated.remoteTabID = view.tabID
projections.insert(updated)
updated.remoteTabID = tabID
// Reconcile reprojects through here. Rewriting unchanged coordinates
// would request the next reconcile of this machine, forever.
guard updated != projection else { return projection }
// One assignment, so observers never see the pane briefly unprojected.
var next = projections
next.remove(projection)
next.insert(updated)
projections = next
reconcileCloudWorkspaceBinding(localWorkspaceID: updated.workspaceID)
notifyChange(for: updated.resource.machine)
return updated
Expand Down
83 changes: 83 additions & 0 deletions cmuxTests/CloudWorkspaceLiveProjectionTests.swift
Original file line number Diff line number Diff line change
Expand Up @@ -262,6 +262,89 @@ struct CloudWorkspaceLiveProjectionTests {
#expect(catalog.projections == [first.projection])
}

/// Reconcile reprojects through `project`, so a reuse that rewrites unchanged
/// coordinates would request its own next pass and spin the main actor.
@Test("Reusing a projection at its current placement changes nothing")
func reusingCurrentPlacementIsNoOp() async throws {
let fixture = boundWorkspaceFixture()
defer { fixture.workspace.teardownAllPanels() }
let catalog = fixture.catalog
catalog.register(CloudPlacementTestProvider(machine: machine))
install(try graph(["first": "a"], revision: 1), catalog: catalog)
await fixture.coordinator.waitForIdle()
let projection = try #require(catalog.projections.first { $0.workspaceID == fixture.workspace.id })
#expect(projection.remoteWorkspaceID == "a" && projection.remoteTabID == "first")
let view = try #require(try catalog.remoteView(for: projection.resource, tabID: projection.remoteTabID, workspaceID: "a"))
let version = catalog.projectionVersions[machine]

let reused = try await catalog.project(
projection.resource, into: .workspace(id: fixture.workspace.id, placement: .tab),
focus: false, reuseExisting: true, reuseInWorkspace: fixture.workspace.id, remoteView: view
)

#expect(reused.reused)
#expect(reused.projection == projection)
#expect(catalog.projectionVersions[machine] == version)
}

/// Any consumer that requests another pass without changing the graph (the
/// nightly b36a9b3 livelock) must end in a bounded number of passes.
@Test("Reconciling one graph stops when every pass requests another")
func reconcileOfOneGraphIsBounded() async throws {
let fixture = boundWorkspaceFixture()
defer { fixture.workspace.teardownAllPanels() }
let catalog = fixture.catalog
let coordinator = fixture.coordinator
let machine = self.machine
var passes = 0
coordinator.environment.applyLayout = { [unowned catalog, unowned coordinator] _, _, _ in
passes += 1
if passes < 1_000 { coordinator.request(machine: machine, catalog: catalog) }
}
catalog.register(CloudPlacementTestProvider(machine: machine))
install(try graph(["first": "a"], revision: 1), catalog: catalog)
await coordinator.waitForIdle()

#expect(passes > 0, "The fixture must reach the layout step")
#expect(passes <= CloudWorkspaceReconcileBudget.maxIdlePasses + 2)
#expect(catalog.projections.contains { $0.workspaceID == fixture.workspace.id && $0.remoteTabID == "first" })
}

/// The nightly b36a9b3 shape: every pass rewrote the same projection, which
/// advances the projection revision and so looks like progress.
@Test("Reconciling one graph stops when every pass rewrites projections and requests another")
func reconcileThatOnlyLooksBusyIsBounded() async throws {
let fixture = boundWorkspaceFixture()
defer { fixture.workspace.teardownAllPanels() }
let catalog = fixture.catalog
let coordinator = fixture.coordinator
let machine = self.machine
var passes = 0
coordinator.environment.applyLayout = { [unowned catalog, unowned coordinator] _, _, _ in
passes += 1
catalog.projectionVersions[machine, default: 0] &+= 1
if passes < 1_000 { coordinator.request(machine: machine, catalog: catalog) }
}
catalog.register(CloudPlacementTestProvider(machine: machine))
install(try graph(["first": "a"], revision: 1), catalog: catalog)
await coordinator.waitForIdle()

#expect(passes > 0, "The fixture must reach the layout step")
#expect(passes <= CloudWorkspaceReconcileBudget.maxPassesPerState)
}

@Test("A reconcile that keeps making progress is not cut short by the idle limit")
func progressingPassesStayAdmitted() throws {
let state = try graph(["first": "a"], revision: 1)
var budget = CloudWorkspaceReconcileBudget()
for version in 0..<UInt64(CloudWorkspaceReconcileBudget.maxPassesPerState) {
#expect(budget.admit(.init(state: state, projectionVersion: version, bindings: [:])))
}
#expect(!budget.admit(.init(state: state, projectionVersion: 1_000, bindings: [:])))
let next = try graph(["first": "a"], revision: 2)
#expect(budget.admit(.init(state: next, projectionVersion: 1_000, bindings: [:])), "A new graph starts a new budget")
}

@Test("Lifecycle cancellation is not retained as a projection failure")
func cancelledMaterializationIsNotAnError() async throws {
let live = LiveWorkspaceFixture()
Expand Down
Loading