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
13 changes: 6 additions & 7 deletions Sources/Cloud/VMClient.swift
Original file line number Diff line number Diff line change
Expand Up @@ -619,12 +619,6 @@ struct VMPublicationDomain: Equatable, Sendable {
}


struct VMSnapshotResult {
let id: String
let name: String?
let createdAt: Int64
}

/// One reflection read (`GET /api/vm/<id>/reflection[/<path>]`): the HTTP status and the
/// JSON body as sent. A 404 with `{error: "not_found", paths: […]}` is a normal result
/// (an unknown reflection path), so the CLI can print the paths that do exist.
Expand Down Expand Up @@ -777,7 +771,7 @@ actor VMClient {
/// the composition root.
@MainActor
static func bootstrap(auth: AuthCoordinator, session: URLSession = .shared, operations: CloudOperationRecorder? = nil) {
shared = VMClient(session: session, auth: auth, operations: operations)
shared = VMClient(session: session, auth: auth, checkpointRenames: SurfaceCatalog.shared.cloudRenameCoordinator, operations: operations)
}

/// Revoke endpoint credentials issued by the Cloud VM service during sign-out.
Expand Down Expand Up @@ -816,6 +810,7 @@ actor VMClient {

private let session: URLSession
private let auth: AuthCoordinator
private let checkpointRenames: CloudRenameCoordinator
private let telemetry: VMClientTelemetry
nonisolated let operations: CloudOperationRecorder?
private let machineCache: CloudMachineCache
Expand All @@ -824,13 +819,15 @@ actor VMClient {
init(
session: URLSession = .shared,
auth: AuthCoordinator,
checkpointRenames: CloudRenameCoordinator,
telemetry: VMClientTelemetry = .shared,
operations: CloudOperationRecorder? = nil,
machineCache: CloudMachineCache = CloudMachineCache(),
isDisabledByManagedPolicy: (@Sendable () -> Bool)? = nil
) {
self.session = session
self.auth = auth
self.checkpointRenames = checkpointRenames
self.telemetry = telemetry
self.operations = operations
self.machineCache = machineCache
Expand Down Expand Up @@ -1504,6 +1501,7 @@ actor VMClient {

func snapshot(id: String, name: String? = nil) async throws -> VMSnapshotResult {
return try await withOperation(.snapshot, foreground: true) {
try await checkpointRenames.waitForPendingRenames(on: .cloud(id))
var body: [String: Any] = [:]
if let name, !name.trimmingCharacters(in: .whitespacesAndNewlines).isEmpty {
body["name"] = name
Expand Down Expand Up @@ -1531,6 +1529,7 @@ actor VMClient {

func fork(id: String, name: String? = nil, idempotencyKey: String) async throws -> (snapshot: VMSnapshotResult?, vm: VMSummary) {
return try await withOperation(.fork, foreground: true) {
try await checkpointRenames.waitForPendingRenames(on: .cloud(id))
var body: [String: Any] = [:]
if let name, !name.trimmingCharacters(in: .whitespacesAndNewlines).isEmpty {
body["name"] = name
Expand Down
5 changes: 5 additions & 0 deletions Sources/Cloud/VMSnapshotResult.swift
Original file line number Diff line number Diff line change
@@ -0,0 +1,5 @@
struct VMSnapshotResult {
let id: String
let name: String?
let createdAt: Int64
}
104 changes: 104 additions & 0 deletions Sources/Surfaces/CloudRenameCoordinator.swift
Original file line number Diff line number Diff line change
@@ -0,0 +1,104 @@
import Foundation

/// Remote rename requests are process-wide because one daemon workspace or tab can be
/// projected into more than one local window. Keeping the lane here prevents two window
/// owners from sending the same remote identity out of order.
@MainActor
final class CloudRenameCoordinator {
struct Key: Hashable, Sendable {
enum Scope: String, Hashable, Sendable {
case workspace
case tab
case terminal
}

let machine: SurfaceMachineID
let scope: Scope
let remoteID: String

static func workspace(machine: SurfaceMachineID, id: String) -> Self {
Self(machine: machine, scope: .workspace, remoteID: id)
}

static func tab(machine: SurfaceMachineID, id: String) -> Self {
Self(machine: machine, scope: .tab, remoteID: id)
}

static func terminal(machine: SurfaceMachineID, id: String) -> Self {
Self(machine: machine, scope: .terminal, remoteID: id)
}
}

private struct Entry {
let generation: UInt64
let task: Task<Void, Error>
}

private struct PendingName {
let generation: UInt64
let value: String
let task: Task<Void, Error>
}

/// The daemon cursor is global to one machine, so all remote rename writes
/// share one lane. Identity keys remain separate for optimistic projection
/// reconciliation.
private var entries: [SurfaceMachineID: Entry] = [:]
private var pendingNames: [Key: PendingName] = [:]
private var nextGeneration: UInt64 = 0

func pendingName(for key: Key) -> String? {
pendingNames[key]?.value
}

/// Commits the names already chosen when a checkpoint was requested. Capture
/// one task per identity: a later edit supersedes an earlier failure for that
/// identity, but a failed workspace rename cannot hide behind a successful tab rename.
func waitForPendingRenames(on machine: SurfaceMachineID) async throws {
let pending = pendingNames.filter { $0.key.machine == machine }.values.sorted {
$0.generation < $1.generation
}
try Task.checkCancellation()
for name in pending {
try await name.task.value
try Task.checkCancellation()
}
}

/// Serializes every remote rename for one machine across every local window and
/// retains the newest optimistic name for each identity. A failed operation can
/// compensate its own local view; an older completion cannot clear a newer intent
/// or queue tail.
@discardableResult
func enqueue(
key: Key,
pendingName: String,
operation: @escaping @MainActor () async throws -> Void
) -> Task<Void, Error> {
let lane = key.machine
nextGeneration &+= 1
let generation = nextGeneration
let pendingGeneration = generation
let previous = entries[lane]?.task
let task = Task { @MainActor [weak self] in
defer { self?.finish(key: key, lane: lane, generation: generation, pendingGeneration: pendingGeneration) }
if let previous {
// A failed or cancelled rename must not strand later edits.
_ = try? await previous.value
}
try Task.checkCancellation()
try await operation()
}
entries[lane] = Entry(generation: generation, task: task)
pendingNames[key] = PendingName(generation: pendingGeneration, value: pendingName, task: task)
return task
}

private func finish(key: Key, lane: SurfaceMachineID, generation: UInt64, pendingGeneration: UInt64) {
if pendingNames[key]?.generation == pendingGeneration {
pendingNames[key] = nil
}
guard entries[lane]?.generation == generation else { return }
entries[lane] = nil
}
}
17 changes: 17 additions & 0 deletions Sources/Surfaces/CmuxTuiSurfaceProvider+Refresh.swift
Original file line number Diff line number Diff line change
@@ -1,6 +1,23 @@
import Foundation

extension CmuxTuiSurfaceProvider {
/// Suspended work cannot publish through a replacement provider.
func isRegisteredInCatalog() -> Bool {
guard let current = catalog.provider(for: machine) else { return false }
return ObjectIdentifier(current) == ObjectIdentifier(self)
}

/// A link acknowledgement can suspend between installing and publishing a graph.
/// Only the current graph may update catalog rows or restored local titles.
func canPublishCloudState(_ candidate: CloudVMState) -> Bool {
guard isRegisteredInCatalog(), candidate.machine == machine,
let current = cloudState else { return false }
if let cursor = current.cursor {
return candidate.cursor == cursor
}
return candidate == current
}

func refresh() async {
await refreshCurrentGraph(force: false)
}
Expand Down
17 changes: 8 additions & 9 deletions Sources/Surfaces/CmuxTuiSurfaceProviders.swift
Original file line number Diff line number Diff line change
Expand Up @@ -194,12 +194,6 @@ final class CmuxTuiSurfaceProvider: SurfaceProvider {
pendingRemoteRenames.removeAll()
acceptedCloudGenerations.removeAll()
}
/// Whether this provider is still registered for its machine. Suspended
/// network work must not write through a replacement provider.
func isRegisteredInCatalog() -> Bool {
guard let current = catalog.provider(for: machine) else { return false }
return ObjectIdentifier(current) == ObjectIdentifier(self)
}
func isCurrentLifecycleGeneration(_ generation: UInt64) -> Bool {
lifecycleGeneration == generation
}
Expand Down Expand Up @@ -416,7 +410,7 @@ final class CmuxTuiSurfaceProvider: SurfaceProvider {
}

@discardableResult
private func installSnapshotIfNewer(_ incoming: CloudVMState, requestVersion: UInt64? = nil) -> Bool {
func installSnapshotIfNewer(_ incoming: CloudVMState, requestVersion: UInt64? = nil) -> Bool {
guard acceptsIncomingGeneration(incoming.cursor) else {
#if DEBUG
cmuxDebugLog("cloud.state.snapshotIgnored machine=\(machineID) reason=old-generation")
Expand Down Expand Up @@ -548,12 +542,13 @@ final class CmuxTuiSurfaceProvider: SurfaceProvider {
/// Publishes the authoritative graph and every derived row in one catalog
/// transaction. Display and forwarded-port rows are machine capabilities, so
/// they join the daemon graph here without becoming a second session state.
private func publish(
func publish(
_ state: CloudVMState,
ports: [Int],
reconcileTitles: Bool = true,
observation: CloudVMStateObservation = .current
) {
guard canPublishCloudState(state) else { return }
var pool: [SurfaceResource] = []
// The control plane's resolved kind is authoritative. Freestyle snapshot
// ids are opaque and cannot tell us whether the machine has a desktop.
Expand Down Expand Up @@ -585,12 +580,13 @@ final class CmuxTuiSurfaceProvider: SurfaceProvider {
/// Applies a contiguous event to the catalog's canonical graph. Row-local changes rebuild
/// only their affected terminal, browser, or display rows. A topology change crosses a
/// relationship boundary and uses the authoritative complete publication path.
private func publishDelta(
func publishDelta(
_ state: CloudVMState,
impact: CloudVMStateDeltaImpact,
ports: [Int],
reconcileTitles: Bool
) {
guard canPublishCloudState(state) else { return }
if impact.requiresFullResourceRebuild {
publish(state, ports: ports, reconcileTitles: reconcileTitles)
return
Expand Down Expand Up @@ -1913,10 +1909,12 @@ final class CmuxTuiSurfaceProvider: SurfaceProvider {
if installSnapshotIfNewer(incoming) {
clearStateRecovery()
await link.setEventsCursor(incoming.cursor)
guard watchedLink === link, canPublishCloudState(incoming) else { return }
var subscriptionResumed = false
if let cursor = incoming.cursor {
subscriptionResumed = await link.resumeEventsSubscription(from: cursor)
}
guard watchedLink === link, canPublishCloudState(incoming) else { return }
if CloudVMEventFeedRecoveryDecision.shouldClearWarning(
snapshotCursor: incoming.cursor,
subscriptionResumed: subscriptionResumed
Expand Down Expand Up @@ -1985,6 +1983,7 @@ final class CmuxTuiSurfaceProvider: SurfaceProvider {
eventsFeedWarning = nil
clearStateRecovery()
await link.setEventsCursor(next.cursor)
guard watchedLink === link, canPublishCloudState(next) else { return }
info.linkState = .connected
info.linkError = nil
let titlesChanged = current.workspaces != next.workspaces || current.tabs != next.tabs
Expand Down
Loading
Loading