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
21 changes: 18 additions & 3 deletions Sources/AgentRestoreEvidenceObservation.swift
Original file line number Diff line number Diff line change
Expand Up @@ -10,11 +10,26 @@ struct AgentRestoreEvidenceObservation: Sendable {
@Sendable
#endif
nonisolated func wait(process: AgentPIDProcessIdentity?, paths: [String]) async {
if let process, AgentPIDProcessIdentity(pid: process.pid) != process { return }
let observation = AgentRestoreEvidenceSubscription(process: process, paths: paths)
await wait(processes: process.map { [$0] } ?? [], paths: paths)
}

/// Waits for any supplied process generation or watched path to change.
/// Every PID is checked before and after registration so a reused PID can
/// never make a stale owner look live.
#if compiler(>=6.2)
@concurrent
#else
@Sendable
#endif
nonisolated func wait(
processes: [AgentPIDProcessIdentity],
paths: [String]
) async {
guard processes.allSatisfy({ AgentPIDProcessIdentity(pid: $0.pid) == $0 }) else { return }
let observation = AgentRestoreEvidenceSubscription(processes: processes, paths: paths)
defer { observation.cancel() }
// Close the registration race without ever accepting a reused PID.
if let process, AgentPIDProcessIdentity(pid: process.pid) != process { return }
guard processes.allSatisfy({ AgentPIDProcessIdentity(pid: $0.pid) == $0 }) else { return }
for await _ in observation.events { return }
}
}
13 changes: 11 additions & 2 deletions Sources/AgentRestoreEvidenceSubscription.swift
Original file line number Diff line number Diff line change
Expand Up @@ -9,12 +9,21 @@ final class AgentRestoreEvidenceSubscription: @unchecked Sendable {
private let continuation: AsyncStream<Void>.Continuation
private let sources: [any DispatchSourceProtocol]

init(process: AgentPIDProcessIdentity?, paths: [String], deadline: DispatchTime = .now() + .seconds(8)) {
convenience init(process: AgentPIDProcessIdentity?, paths: [String], deadline: DispatchTime = .now() + .seconds(8)) {
self.init(processes: process.map { [$0] } ?? [], paths: paths, deadline: deadline)
}

init(
processes: [AgentPIDProcessIdentity],
paths: [String],
deadline: DispatchTime = .now() + .seconds(8)
) {
let (events, continuation) = AsyncStream<Void>.makeStream(bufferingPolicy: .bufferingNewest(1))
self.events = events
self.continuation = continuation
var sources: [any DispatchSourceProtocol] = []
if let process {
var subscribedPIDs = Set<pid_t>()
for process in processes where subscribedPIDs.insert(process.pid).inserted {
let source = DispatchSource.makeProcessSource(
identifier: process.pid, eventMask: .exit, queue: .global(qos: .utility)
)
Expand Down
43 changes: 36 additions & 7 deletions Sources/DeferredAgentResumeAdmissionOwner.swift
Original file line number Diff line number Diff line change
Expand Up @@ -21,7 +21,10 @@ protocol DeferredAgentResumeAdmissionOwner: AnyObject {
/// Fresh evidence for this owner's pending requests.
var deferredAgentResumeIndexProvider: @MainActor @Sendable () async -> SharedLiveAgentIndexRefreshOutcome { get }
/// Waits for the next evidence event; cancellation releases the wait.
var deferredAgentResumeEvidenceWait: @Sendable () async -> Void { get }
/// The process identities are the owners observed by the just-completed
/// scan. Waiting on their kernel exit events closes the gap where a
/// SessionEnd hook cannot write because the old cmux socket is gone.
var deferredAgentResumeEvidenceWait: @Sendable ([AgentPIDProcessIdentity]) async -> Void { get }
}

extension DeferredAgentResumeAdmissionOwner {
Expand All @@ -39,11 +42,17 @@ extension DeferredAgentResumeAdmissionOwner {
deferredAgentResumeIndexTask = Task { @MainActor [weak self] in
while !Task.isCancelled {
let outcome = await refresh()
guard !Task.isCancelled else { return }
// Do not hold the owner across the evidence wait: teardown must
// not be delayed by a restore that is still observing evidence.
guard !Task.isCancelled,
self?.applyDeferredAgentResumeOutcome(outcome) == true else { return }
await waitForEvidence()
let processIdentities: [AgentPIDProcessIdentity]? = {
guard let self,
self.applyDeferredAgentResumeOutcome(outcome) else { return nil }
guard case .index(let index) = outcome else { return [] }
return self.deferredAgentResumeOwnerProcessIdentities(using: index)
}()
guard !Task.isCancelled, let processIdentities else { return }
await waitForEvidence(processIdentities)
}
}
}
Expand All @@ -52,15 +61,35 @@ extension DeferredAgentResumeAdmissionOwner {
{ await SharedLiveAgentIndex.shared.indexForOwnershipDecision() }
}

var deferredAgentResumeEvidenceWait: @Sendable () async -> Void {
{
var deferredAgentResumeEvidenceWait: @Sendable ([AgentPIDProcessIdentity]) async -> Void {
{ processIdentities in
await AgentRestoreEvidenceObservation().wait(
process: nil,
processes: processIdentities,
paths: [RestorableAgentKind.claude.hookStoreFileURL().deletingLastPathComponent().path]
)
}
}

/// Returns the process generations that kept a staged restore pending in
/// the completed scan. A missing identity is intentionally not synthesized:
/// the next refresh must decide whether the record is stale or unavailable.
func deferredAgentResumeOwnerProcessIdentities(
using index: RestorableAgentSessionIndex
) -> [AgentPIDProcessIdentity] {
var identities = Set<AgentPIDProcessIdentity>()
for restore in deferredAgentResumeRestoresByPanelId.values {
guard let kind = restore.restorableAgent?.kind.rawValue ?? restore.resumeBinding?.kind,
let sessionID = restore.restorableAgent?.sessionId ?? restore.resumeBinding?.checkpointId,
let owner = index.liveSessionOwner(
kind: kind,
sessionID: sessionID,
revalidateProcessEvidence: false
) else { continue }
identities.insert(owner.processIdentity)
}
return Array(identities)
}

/// Applies one index outcome to every pending restore.
/// - Returns: Whether restores remain pending and the loop should keep observing evidence.
func applyDeferredAgentResumeOutcome(_ outcome: SharedLiveAgentIndexRefreshOutcome) -> Bool {
Expand Down
16 changes: 11 additions & 5 deletions Sources/DockSplitStore+DeferredAgentRestoreAdmission.swift
Original file line number Diff line number Diff line change
Expand Up @@ -39,17 +39,23 @@ extension DockSplitStore {
continue
}
let currentResumeBinding: SurfaceResumeBindingSnapshot?
let observedResumeBinding: SurfaceResumeBindingSnapshot?
if let capturedBinding = restore.resumeBinding {
guard let currentBinding = surfaceResumeBindingsByPanelId[panelId],
currentBinding.isAgentHookBinding,
currentBinding.isSameManagedSession(as: capturedBinding),
currentBinding.autoResume == true else {
currentBinding.isSameManagedSession(as: capturedBinding) else {
cancelDeferredAgentResumeRestore(panelId: panelId, restore: restore)
continue
}
currentResumeBinding = currentBinding
observedResumeBinding = currentBinding
// The staged snapshot owns this launch. A late SessionEnd from
// the previous cmux can retire the live binding in place while
// leaving the same session identity; do not turn that race into
// a silent shell or rebuild the launch from autoResume=false.
currentResumeBinding = capturedBinding
} else {
currentResumeBinding = nil
observedResumeBinding = nil
}
if restore.remoteResumeCommandEmbedded {
// The attach command was embedded in the terminal's initial
Expand All @@ -58,8 +64,8 @@ extension DockSplitStore {
// cwd, and launch flavor) so a changed resume payload can
// never execute from the stale terminal configuration.
guard let capturedBinding = restore.resumeBinding,
let currentResumeBinding,
capturedBinding == currentResumeBinding else {
let observedResumeBinding,
capturedBinding == observedResumeBinding else {
cancelDeferredAgentResumeRestore(panelId: panelId, restore: restore)
continue
}
Expand Down
16 changes: 11 additions & 5 deletions Sources/Workspace+DeferredAgentRestoreAdmission.swift
Original file line number Diff line number Diff line change
Expand Up @@ -39,17 +39,23 @@ extension Workspace {
continue
}
let currentResumeBinding: SurfaceResumeBindingSnapshot?
let observedResumeBinding: SurfaceResumeBindingSnapshot?
if let capturedBinding = restore.resumeBinding {
guard let currentBinding = surfaceResumeBindingsByPanelId[panelId],
currentBinding.isAgentHookBinding,
currentBinding.isSameManagedSession(as: capturedBinding),
currentBinding.autoResume == true else {
currentBinding.isSameManagedSession(as: capturedBinding) else {
cancelDeferredAgentResumeRestore(panelId: panelId, restore: restore)
continue
}
currentResumeBinding = currentBinding
observedResumeBinding = currentBinding
// The staged snapshot owns this launch. A late SessionEnd from
// the previous cmux can retire the live binding in place while
// leaving the same session identity; do not turn that race into
// a silent shell or rebuild the launch from autoResume=false.
currentResumeBinding = capturedBinding
} else {
currentResumeBinding = nil
observedResumeBinding = nil
}
if restore.remoteResumeCommandEmbedded {
// The attach command was embedded in the terminal's initial
Expand All @@ -58,8 +64,8 @@ extension Workspace {
// cwd, and launch flavor) so a changed resume payload can
// never execute from the stale terminal configuration.
guard let capturedBinding = restore.resumeBinding,
let currentResumeBinding,
capturedBinding == currentResumeBinding else {
let observedResumeBinding,
capturedBinding == observedResumeBinding else {
cancelDeferredAgentResumeRestore(panelId: panelId, restore: restore)
continue
}
Expand Down
4 changes: 4 additions & 0 deletions cmux.xcodeproj/project.pbxproj
Original file line number Diff line number Diff line change
Expand Up @@ -164,6 +164,7 @@
15042AA9E67B074D4C03CA18 /* AgentRestoreCodexEvidence.swift in Sources */ = {isa = PBXBuildFile; fileRef = 922A32567AAA6DD33BCE05CF /* AgentRestoreCodexEvidence.swift */; };
5166F999E9205E6CEB90AA79 /* AgentRestoreEvidenceObservation.swift in Sources */ = {isa = PBXBuildFile; fileRef = 34AE0530BC772C9E4DC2BE33 /* AgentRestoreEvidenceObservation.swift */; };
A64E0D245603951ACCAEB957 /* AgentRestoreEvidenceSubscription.swift in Sources */ = {isa = PBXBuildFile; fileRef = AC44C6882243B4D699AE73FA /* AgentRestoreEvidenceSubscription.swift */; };
6229C7F75EEE378D764DDB47 /* AgentRestoreIssue12775Tests.swift in Sources */ = {isa = PBXBuildFile; fileRef = 9BE0A210ACEDFB6FB21BBE0E /* AgentRestoreIssue12775Tests.swift */; };
C41D15BA2CF8651980082256 /* AgentRestoreLiveOwnerAdmissionTests+Fixture.swift in Sources */ = {isa = PBXBuildFile; fileRef = F0E91730F780E2461753BB8A /* AgentRestoreLiveOwnerAdmissionTests+Fixture.swift */; };
110430011104300111043001 /* AgentRestoreLiveOwnerAdmissionTests.swift in Sources */ = {isa = PBXBuildFile; fileRef = 110430021104300211043002 /* AgentRestoreLiveOwnerAdmissionTests.swift */; };
110432021104320211043202 /* AgentRestoreLiveOwnerNotice.swift in Sources */ = {isa = PBXBuildFile; fileRef = 110432011104320111043201 /* AgentRestoreLiveOwnerNotice.swift */; };
Expand Down Expand Up @@ -4269,6 +4270,7 @@
922A32567AAA6DD33BCE05CF /* AgentRestoreCodexEvidence.swift */ = {isa = PBXFileReference; lastKnownFileType = sourcecode.swift; path = "AgentRestoreCodexEvidence.swift"; sourceTree = "<group>"; };
34AE0530BC772C9E4DC2BE33 /* AgentRestoreEvidenceObservation.swift */ = {isa = PBXFileReference; lastKnownFileType = sourcecode.swift; path = "AgentRestoreEvidenceObservation.swift"; sourceTree = "<group>"; };
AC44C6882243B4D699AE73FA /* AgentRestoreEvidenceSubscription.swift */ = {isa = PBXFileReference; lastKnownFileType = sourcecode.swift; path = "AgentRestoreEvidenceSubscription.swift"; sourceTree = "<group>"; };
9BE0A210ACEDFB6FB21BBE0E /* AgentRestoreIssue12775Tests.swift */ = {isa = PBXFileReference; lastKnownFileType = sourcecode.swift; path = "AgentRestoreIssue12775Tests.swift"; sourceTree = "<group>"; };
F0E91730F780E2461753BB8A /* AgentRestoreLiveOwnerAdmissionTests+Fixture.swift */ = {isa = PBXFileReference; lastKnownFileType = sourcecode.swift; path = "AgentRestoreLiveOwnerAdmissionTests+Fixture.swift"; sourceTree = "<group>"; };
110430021104300211043002 /* AgentRestoreLiveOwnerAdmissionTests.swift */ = {isa = PBXFileReference; lastKnownFileType = sourcecode.swift; path = AgentRestoreLiveOwnerAdmissionTests.swift; sourceTree = "<group>"; };
110432011104320111043201 /* AgentRestoreLiveOwnerNotice.swift */ = {isa = PBXFileReference; lastKnownFileType = sourcecode.swift; path = AgentRestoreLiveOwnerNotice.swift; sourceTree = "<group>"; };
Expand Down Expand Up @@ -12056,6 +12058,7 @@
511EAB11A81C099E4F5A8ECF /* DeferredAgentResumeAdmissionOwnerTests.swift */,
C7E3FFB9B84F6DE45250D5D5 /* CodexSessionStartDeadTurnTests.swift */,
0B0073F9140EA86D5ADEE233 /* CodexStaleTurnRestoreIntentTests.swift */,
9BE0A210ACEDFB6FB21BBE0E /* AgentRestoreIssue12775Tests.swift */,
25DB476025F922037411EA98 /* CommandPaletteAuthCommandTests.swift */,
F0677DB157F30F725FECA0EC /* CMUXCLICodexUnavailableAdmissionTests.swift */,
D745F48BE9C06D2ED4C2E3DC /* WorkspaceCreateIdempotencyIsolationTests.swift */,
Expand Down Expand Up @@ -15609,6 +15612,7 @@
762CA407531356D761E008F5 /* AgentQuitOwnershipTests.swift in Sources */,
F6F0F8FCF48DD40A3DB48F42 /* AgentQuitTerminationCoordinatorTests.swift in Sources */,
93840005AA11BB22CC33EE05 /* AgentRelayTTYOwnershipTests.swift in Sources */,
6229C7F75EEE378D764DDB47 /* AgentRestoreIssue12775Tests.swift in Sources */,
C41D15BA2CF8651980082256 /* AgentRestoreLiveOwnerAdmissionTests+Fixture.swift in Sources */,
110430011104300111043001 /* AgentRestoreLiveOwnerAdmissionTests.swift in Sources */,
32870528BA7DA7C68021A6EA /* AgentRestoreRecoveryRegressionTests.swift in Sources */,
Expand Down
Loading