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
17 changes: 13 additions & 4 deletions GraphcodeKit/Sources/Domain/SSHReconnectLoop.swift
Original file line number Diff line number Diff line change
Expand Up @@ -22,7 +22,12 @@ import Foundation
public enum SSHReconnectLoop {
public static let maxDelaySeconds = 15

public static func script(connect: String, reconnect: String) -> String {
/// `redialStamp` is touched before every redial: it is how `graphcoded` learns, for
/// free, that a host may have rebooted under a loop it restores
/// (`ZmxSessionLauncher.restoreRebootedRemote`).
public static func script(
connect: String, reconnect: String, redialStamp: String? = nil
) -> String {
let passExit = "; gc_rc=$?; [ \"$gc_rc\" -ne 255 ] && exit \"$gc_rc\""
return "trap 'exit 130' INT; "
+ connect + passExit + "; "
Expand All @@ -31,7 +36,7 @@ public enum SSHReconnectLoop {
+ #"Press Ctrl-C to stop. ──\033[0m\r\n' "$gc_rc" "$gc_delay"; "#
+ "sleep \"$gc_delay\"; gc_delay=$((gc_delay * 2)); "
+ "[ \"$gc_delay\" -gt \(maxDelaySeconds) ] && gc_delay=\(maxDelaySeconds); "
+ reconnect + passExit + "; done"
+ touching(redialStamp) + reconnect + passExit + "; done"
}

/// A Codespace surface's loop: the same dials and exit handling, retried on
Expand All @@ -52,7 +57,7 @@ public enum SSHReconnectLoop {
/// genuinely ended nonzero, which converges: the redial reattaches a live session, and
/// a gone one takes the reconnect script's session-ended branch to a clean exit 0.
public static func codespaceScript(
connect: String, reconnect: String, pauseMarker: String,
connect: String, reconnect: String, pauseMarker: String, redialStamp: String? = nil,
schedule: CodespaceDialSchedule = .standard, upAfter: Int = 330
) -> String {
let marker = quoted(pauseMarker)
Expand Down Expand Up @@ -100,7 +105,11 @@ public enum SSHReconnectLoop {
+ #"printf '\033[1;33m── Connection failed (exit %s). Retrying in %ss. "#
+ #"Press Ctrl-C to stop. ──\033[0m\r\n' "$gc_rc" "$gc_wait"; "#
+ "gc_wait_or_ask \"$gc_wait\" && { \(restart); }; fi; "
+ "gc_t=$(date +%s); " + reconnect + passExit + clock + "; done"
+ "gc_t=$(date +%s); " + touching(redialStamp) + reconnect + passExit + clock + "; done"
}

private static func touching(_ stamp: String?) -> String {
stamp.map { "touch \(quoted($0)) 2>/dev/null; " } ?? ""
}

/// `RemoteProjectLocation.shellQuoted`, repeated because this file also builds in the
Expand Down
36 changes: 33 additions & 3 deletions GraphcodeKit/Sources/GraphStore.swift
Original file line number Diff line number Diff line change
Expand Up @@ -52,6 +52,7 @@
private let onGraphEvent: (@Sendable (DaemonEvent) -> [UUID: DaemonWireEnvelope])?
private let onConnectionFailure: (@Sendable (UUID) -> Void)?
private let onEnsureSession: (@Sendable (LoopNode, String?) -> Void)?
private let onRestoreRebootedSessions: (@Sendable ([LoopNode], String) async -> Void)?
private let onTerminateSession: (@Sendable (LoopNode, String?) -> Void)?
/// Kills a loop's session and, for an unattended loop, relaunches it on the same
/// transcript. Awaited, unlike the two above: the answer is whether the old session
Expand Down Expand Up @@ -321,6 +322,7 @@
onGraphEvent: (@Sendable (DaemonEvent) -> [UUID: DaemonWireEnvelope])? = nil,
onConnectionFailure: (@Sendable (UUID) -> Void)? = nil,
onEnsureSession: (@Sendable (LoopNode, String?) -> Void)? = nil,
onRestoreRebootedSessions: (@Sendable ([LoopNode], String) async -> Void)? = nil,
onFindMissingProvider: (@Sendable (LoopNode, String?) async -> LaunchFailure?)? = nil,
onTerminateSession: (@Sendable (LoopNode, String?) -> Void)? = nil,
onRestartSession: (@Sendable (LoopNode, String?) async -> Bool)? = nil,
Expand Down Expand Up @@ -364,6 +366,7 @@
self.onGraphEvent = onGraphEvent
self.onConnectionFailure = onConnectionFailure
self.onEnsureSession = onEnsureSession
self.onRestoreRebootedSessions = onRestoreRebootedSessions
self.onFindMissingProvider = onFindMissingProvider
self.onTerminateSession = onTerminateSession
self.onRestartSession = onRestartSession
Expand Down Expand Up @@ -1287,7 +1290,7 @@
onRemoveMemory: onRemoveMemory,
onRefinePlaybook: onRefinePlaybook,
onRollbackPlaybook: onRollbackPlaybook,
onAnnounceError: effects.errors.append,

Check warning on line 1293 in GraphcodeKit/Sources/GraphStore.swift

View workflow job for this annotation

GitHub Actions / macos

converting non-Sendable function value to '@sendable (String) -> Void' may introduce data races
// The board's gate forwards like any other side effect: a loop inside a piloted
// composite is a real loop whose session got the standard briefing — teaching
// verbs the child store would refuse is exactly the incoherence the gate exists
Expand Down Expand Up @@ -2150,7 +2153,7 @@
pendingFollowUps.append(
PendingFollowUp(id: UUID(), nodeID: nodeID, text: message, watchedPostID: nil))
}
await drainAndBroadcast()

Check warning on line 2156 in GraphcodeKit/Sources/GraphStore.swift

View workflow job for this annotation

GitHub Actions / macos

result of call to 'drainAndBroadcast(broadcastErrors:)' is unused
}

/// A learned note into a node's memory log — `graphcode node memo`, the agent-written
Expand Down Expand Up @@ -2308,7 +2311,7 @@
let followUp = PendingFollowUp(id: UUID(), nodeID: nodeID, text: prompt, watchedPostID: nil)
pendingFollowUps.append(followUp)
goalFollowUps[nodeID] = followUp.id
await drainAndBroadcast()

Check warning on line 2314 in GraphcodeKit/Sources/GraphStore.swift

View workflow job for this annotation

GitHub Actions / macos

result of call to 'drainAndBroadcast(broadcastErrors:)' is unused
}

/// Opening a resolved loop whose session was ended brings its conversation back. Panes
Expand Down Expand Up @@ -2848,7 +2851,7 @@
return
}
if node.loopType == .composite {
await runInSubGraph(nodeID, .restartSessions, broadcastErrors: false)

Check warning on line 2854 in GraphcodeKit/Sources/GraphStore.swift

View workflow job for this annotation

GitHub Actions / macos

result of call to 'runInSubGraph(_:_:broadcastErrors:)' is unused
return
}
await restart([node])
Expand All @@ -2857,7 +2860,7 @@
private func restartSessions() async {
let live = graph.nodes.filter { !$0.isResolved }
for composite in live where composite.loopType == .composite {
await runInSubGraph(composite.id, .restartSessions, broadcastErrors: false)

Check warning on line 2863 in GraphcodeKit/Sources/GraphStore.swift

View workflow job for this annotation

GitHub Actions / macos

result of call to 'runInSubGraph(_:_:broadcastErrors:)' is unused
}
await restart(live.filter { $0.loopType != .composite })
}
Expand Down Expand Up @@ -2913,7 +2916,7 @@
// set below — a graph whose nodes have all stopped aggregates to `.idle`.
if node.loopType == .composite, let subGraph = node.subGraph {
for child in subGraph.nodes where !child.isResolved {
await runInSubGraph(

Check warning on line 2919 in GraphcodeKit/Sources/GraphStore.swift

View workflow job for this annotation

GitHub Actions / macos

result of call to 'runInSubGraph(_:_:broadcastErrors:)' is unused
node.id, .stopNode(child.id), broadcastErrors: false)
}
}
Expand Down Expand Up @@ -4023,7 +4026,7 @@
onRemoveMemory: onRemoveMemory,
onRefinePlaybook: onRefinePlaybook,
onRollbackPlaybook: onRollbackPlaybook,
onAnnounceError: effects.errors.append,

Check warning on line 4029 in GraphcodeKit/Sources/GraphStore.swift

View workflow job for this annotation

GitHub Actions / macos

converting non-Sendable function value to '@sendable (String) -> Void' may introduce data races
onMailroomEnabled: onMailroomEnabled,
goalCache: goalCache,
recurrence: effects.recurrence,
Expand All @@ -4039,7 +4042,7 @@
processRecurrence(effects.recurrence)
graph.nodes[id: ownerID]?.subGraph = await child.graph
rollUpComposite(ownerID)
await drainAndBroadcast()

Check warning on line 4045 in GraphcodeKit/Sources/GraphStore.swift

View workflow job for this annotation

GitHub Actions / macos

result of call to 'drainAndBroadcast(broadcastErrors:)' is unused
}

/// Applies the recurrence requests a child store handed up, in order — an update's
Expand Down Expand Up @@ -4110,7 +4113,7 @@
now.timeIntervalSince(node.createdAt) >= stallAfter
{
markStalled(nodeID)
await drainAndBroadcast()

Check warning on line 4116 in GraphcodeKit/Sources/GraphStore.swift

View workflow job for this annotation

GitHub Actions / macos

result of call to 'drainAndBroadcast(broadcastErrors:)' is unused
return
}

Expand Down Expand Up @@ -4474,16 +4477,43 @@
/// sweep would restart each poller's interval and a goal polled less often than the
/// sweep would never fire at all. The pollers are already running; they are in-memory
/// and a remote reboot doesn't touch them.
/// - **Resolved nodes are skipped whatever their loop type.** The load-time version
/// restarts a `.stopped` time-based node, which is defensible once at boot and wrong
/// every minute: a human who stopped a remote loop would watch it come back.
/// - **Resolved nodes are never ensured, whatever their loop type.** The load-time
/// version restarts a `.stopped` time-based node, which is defensible once at boot and
/// wrong every minute: a human who stopped a remote loop would watch it come back.
///
/// A resolved node's session is still brought back when the host's reboot killed it
/// (`onRestoreRebootedSessions`): a pane leaves every unattended loop to this daemon,
/// so a finished loop's open pane otherwise dialed "waiting for graphcoded" forever. The
/// restore is the conversation only (`rebootRestoreCopy`): no task, no poller or
/// heartbeat, no state change. Attended loops are not here — their pane restores them.
public func ensureUnattendedSessionsAlive() async {
for node in graph.nodes where node.runsUnattended && !node.isResolved {
ensureSession(node)
}
let finished = graph.nodes.filter {
$0.runsUnattended && $0.isResolved && $0.launchFailure == nil
}
if let onRestoreRebootedSessions, !finished.isEmpty {
let path = graph.project.path
let copies = finished.map(Self.rebootRestoreCopy)
Task.detached { await onRestoreRebootedSessions(copies, path) }
}
await broadcastIfTemplatesRefreshed()
}

/// A finished loop as its reboot restore launches it: the banked conversation resumed,
/// or, with nothing banked, a session opening on this note instead of the loop's task —
/// the same shape `resumeResolvedSession` gives a finished loop a human opens.
static func rebootRestoreCopy(of node: LoopNode) -> LoopNode {
var quiet = node
quiet.loopType = .sketch
quiet.attachments = []
quiet.firstInstruction =
"[graphcode] The machine this loop runs on restarted, and its earlier conversation "
+ "could not be resumed. This loop has finished; wait for the human's question."
return quiet
}

// MARK: - Broadcast

private func broadcast() async {
Expand Down
5 changes: 5 additions & 0 deletions GraphcodeKit/Sources/ProjectRegistry.swift
Original file line number Diff line number Diff line change
Expand Up @@ -59,6 +59,7 @@ public actor ProjectRegistry {
/// one named project — see `sidebarSubscribers`.
private var sidebarConnections: Set<UUID> = []
private let ensureSession: (@Sendable (LoopNode, String?) -> Void)?
private let restoreRebootedSessions: (@Sendable ([LoopNode], String) async -> Void)?
private let terminateSession: (@Sendable (LoopNode, String?) -> Void)?
private let restartSession: (@Sendable (LoopNode, String?) async -> Bool)?
private let startQuickChat:
Expand Down Expand Up @@ -121,6 +122,8 @@ public actor ProjectRegistry {
platformPaths: any PlatformPaths = CurrentPlatformPaths.value,
replayStore: DaemonReplayStore = DaemonReplayStore(),
ensureSession: (@Sendable (LoopNode, String?) -> Void)? = CLISessionBackend.ensureSession,
restoreRebootedSessions: (@Sendable ([LoopNode], String) async -> Void)? =
CLISessionBackend.restoreRebootedSessions,
terminateSession: (@Sendable (LoopNode, String?) -> Void)? =
CLISessionBackend.terminateSession,
restartSession: (@Sendable (LoopNode, String?) async -> Bool)? =
Expand Down Expand Up @@ -163,6 +166,7 @@ public actor ProjectRegistry {
quickChatStore = QuickChatStore(baseDirectory: persistenceDirectory)
self.replayStore = replayStore
self.ensureSession = ensureSession
self.restoreRebootedSessions = restoreRebootedSessions
self.terminateSession = terminateSession
self.restartSession = restartSession
self.evaluatePredicate = evaluatePredicate
Expand Down Expand Up @@ -1011,6 +1015,7 @@ public actor ProjectRegistry {
},
onConnectionFailure: onConnectionFailure,
onEnsureSession: ensureSession,
onRestoreRebootedSessions: restoreRebootedSessions,
onFindMissingProvider: { node, path in
await ProviderPath.missingProvider(for: node, projectPath: path)
},
Expand Down
5 changes: 5 additions & 0 deletions GraphcodeKit/Sources/Sessions/CLISessionBackend.swift
Original file line number Diff line number Diff line change
Expand Up @@ -287,6 +287,11 @@ extension CLISessionBackend {
Task.detached { await backend(for: node).launch(node, path) }
}

public static let restoreRebootedSessions: @Sendable ([LoopNode], String) async -> Void = {
nodes, path in
await ZmxSessionLauncher.restoreRebootedRemote(nodes, projectPath: path)
}

public static let terminateSession: @Sendable (LoopNode, String?) -> Void = { node, path in
Task.detached { await backend(for: node).terminate(node, path) }
}
Expand Down
130 changes: 124 additions & 6 deletions GraphcodeKit/Sources/Sessions/ZmxSessionLauncher.swift
Original file line number Diff line number Diff line change
Expand Up @@ -1635,7 +1635,7 @@ public enum ZmxSessionLauncher {
static func remoteEnsureInvocation(
forNode node: LoopNode, at location: RemoteProjectLocation,
settings: GraphcodeSettings = GraphcodeSettingsStore.load(),
bridgeState: RemoteBridgeWireState? = nil
bridgeState: RemoteBridgeWireState? = nil, onlyAfterReboot: Bool = false
) -> [String]? {
let shedPrompt = ShedPromptReport()
guard
Expand Down Expand Up @@ -1710,13 +1710,25 @@ public enum ZmxSessionLauncher {
let launch =
agentLabelCommand(zmxPath: "zmx", forNode: node)
.map { "\(create) && { \($0) || true; }" } ?? create
// The boot this session is alive in, recorded by the daemon as well as by a pane
// attach (`RemoteBootMarker`): it is how a pane, and `rebootProbeScript`, tell a
// session that died with the machine from one that ended.
let name = SurfaceRef(id: node.id, launchesClaudeCode: true).zmxSessionName
let markerWrite = RemoteBootMarker.writeFragment(forSessionName: name)
var missing = trustSeed + hooksWrite + "{ \(launch); } && { \(markerWrite); }"
if onlyAfterReboot {
let marker = RemoteBootMarker.markerExpression(forSessionName: name)
missing =
"\(RemoteBootMarker.captureFragment); gc_last=$(cat \(marker) 2>/dev/null); "
+ "if [ -n \"$gc_boot\" ] && [ -n \"$gc_last\" ] && [ \"$gc_boot\" != \"$gc_last\" ]; "
+ "then \(missing); fi"
}
let script =
"cd \(RemoteProjectLocation.shellQuoted(location.remotePath)) && { "
+ deliveryFragment(
delivery, ifSessionMissing: check,
bridgeStateGeneration: bridgeState.map(\.generation))
+ "\(check) >/dev/null 2>&1\(bank) || \(repair){ " + trustSeed + hooksWrite
+ "\(launch); }; }"
+ "\(check) >/dev/null 2>&1\(bank) && { \(markerWrite); } || \(repair){ \(missing); }; }"
return location.sshInvocation(remoteCommand: location.remoteLoginShellCommand(script))
}

Expand Down Expand Up @@ -2228,15 +2240,22 @@ public enum ZmxSessionLauncher {
static func remoteKillInvocation(
forNode node: LoopNode, at location: RemoteProjectLocation
) -> [String] {
let script = quotedCommand(["zmx"] + killArguments(forNode: node))
// A session ended on purpose must not read as one a reboot killed, to the pane or to
// `rebootProbeScript`.
let marker = RemoteBootMarker.markerExpression(
forSessionName: SurfaceRef(id: node.id, launchesClaudeCode: true).zmxSessionName)
let script =
"rm -f \(marker); " + quotedCommand(["zmx"] + killArguments(forNode: node))
return location.sshInvocation(remoteCommand: location.remoteLoginShellCommand(script))
}

static func killRemote(_ node: LoopNode, at location: RemoteProjectLocation) async {
_ = await runRemoteRetrying(remoteKillInvocation(forNode: node, at: location))
}

private static func startRemote(_ node: LoopNode, at location: RemoteProjectLocation) async {
private static func startRemote(
_ node: LoopNode, at location: RemoteProjectLocation, onlyAfterReboot: Bool = false
) async {
// A codespace that is down is redialed on the shared schedule, not on every sweep.
guard await CodespaceDialBreaker.shared.permits(location) else { return }
// A dial already in flight for this node is doing this job; a second one racing it
Expand Down Expand Up @@ -2294,7 +2313,7 @@ public enum ZmxSessionLauncher {
// as the local path: no UI here, the node's state stays honest, opening the loop
// retries.
if let ensure = remoteEnsureInvocation(
forNode: node, at: location, bridgeState: bridgeState
forNode: node, at: location, bridgeState: bridgeState, onlyAfterReboot: onlyAfterReboot
) {
if await runRemoteRetrying(ensure) {
await CodespaceDialBreaker.shared.record(location, reached: true)
Expand All @@ -2304,6 +2323,105 @@ public enum ZmxSessionLauncher {
await RemoteEnsureGate.shared.end(node.id, token: lease)
}

/// Brings back the sessions of finished loops that a reboot of their remote host killed
/// (`GraphStore.ensureUnattendedSessionsAlive`). One probe dial per host names the
/// sessions that are missing *and* were last seen alive in an earlier boot; only those
/// are dialed again, each behind the same boot gate.
///
/// The probe itself runs only when a pane of that host has redialed since the last
/// probe that answered (`redialStamp`): the one thing left dialing a finished loop's
/// host is its pane, and a healthy host has no pane redialing, so the sweep spends
/// nothing — a codespace dial spends the human's API quota (issue #480).
///
/// `nodes` are already the quiet copies the store made (`GraphStore.rebootRestoreCopy`):
/// the create resumes the banked conversation, or opens on a note, never on the task.
static func restoreRebootedRemote(
_ nodes: [LoopNode], projectPath: String, gate: RebootProbeGate = .shared
) async {
guard !nodes.isEmpty, let location = RemoteProjectLocation.parse(projectPath: projectPath)
else { return }
let asked = Date()
guard await gate.panesRedialed(location) else { return }
let name = { (node: LoopNode) in
SurfaceRef(id: node.id, launchesClaudeCode: true).zmxSessionName
}
let probe = location.sshInvocation(
remoteCommand: location.remoteLoginShellCommand(
rebootProbeScript(forSessionNames: nodes.map(name))))
let (succeeded, output) = await collectRemoteOutput(probe, location: location)
guard succeeded else { return }
await gate.probed(location, at: asked)
let rebooted = parseRebootProbe(output)
for node in nodes where rebooted.contains(name(node)) {
await startRemote(node, at: location, onlyAfterReboot: true)
}
}

/// Touched by a remote pane's reconnect loop before every redial
/// (`SSHReconnectLoop`), and read by `RebootProbeGate`. Per host, not per loop: one
/// probe answers for every loop on it.
public static func redialStamp(for location: RemoteProjectLocation) -> URL {
SupportDirectory.url.appendingPathComponent("remote-redials", isDirectory: true)
.appendingPathComponent("\(location.host).redial")
}

/// Whether a host's panes have redialed since its last answered probe — the only
/// state `restoreRebootedRemote` keeps. Stamps from before this daemon started count
/// once, so a pane left waiting across a daemon restart is still answered.
actor RebootProbeGate {
static let shared = RebootProbeGate()

private let stampFor: @Sendable (RemoteProjectLocation) -> URL
private var probedAt: [String: Date] = [:]

init(
stampFor: @escaping @Sendable (RemoteProjectLocation) -> URL = {
ZmxSessionLauncher.redialStamp(for: $0)
}
) {
self.stampFor = stampFor
}

func panesRedialed(_ location: RemoteProjectLocation) -> Bool {
guard
let touched =
(try? FileManager.default.attributesOfItem(
atPath: stampFor(location).path))?[.modificationDate] as? Date
else { return false }
guard let since = probedAt[location.host] else { return true }
return touched > since
}

/// Recorded only for a probe that answered: one that failed leaves the redial
/// pending, so the host is probed once it is back even if every pane has paused.
func probed(_ location: RemoteProjectLocation, at date: Date) {
probedAt[location.host] = date
}
}

/// Prints `rebooted <name>` for each session that is not running and whose boot marker
/// names a boot other than this one — the pane's reboot verdict, made for many sessions
/// in one shell. A host that cannot answer (no boot ID, `zmx` not up yet) prints nothing.
static func rebootProbeScript(forSessionNames names: [String]) -> String {
let quotedNames = names.map(RemoteProjectLocation.shellQuoted).joined(separator: " ")
return "\(RemoteBootMarker.captureFragment); [ -n \"$gc_boot\" ] || exit 0; "
+ "gc_ls=$(zmx ls 2>/dev/null) || exit 0; gc_tab=$(printf '\\t'); "
+ "for gc_n in \(quotedNames); do "
+ "gc_last=$(cat \"$HOME/.graphcode/boots/$gc_n\" 2>/dev/null); "
+ "[ -n \"$gc_last\" ] && [ \"$gc_last\" != \"$gc_boot\" ] || continue; "
+ "printf '%s\\n' \"$gc_ls\" | grep -v -e \"${gc_tab}ended=\" -e \"${gc_tab}err=\" "
+ "| grep -q \"name=$gc_n$gc_tab\" || printf 'rebooted %s\\n' \"$gc_n\"; done; exit 0"
}

static func parseRebootProbe(_ output: String) -> Set<String> {
Set(
output.split(whereSeparator: \.isNewline).compactMap { line in
let fields = line.split(whereSeparator: \.isWhitespace)
guard fields.count == 2, fields[0] == "rebooted" else { return nil }
return String(fields[1])
})
}

static func start(_ node: LoopNode, projectPath: String? = nil) async {
if let projectPath, let remote = RemoteProjectLocation.parse(projectPath: projectPath) {
await startRemote(node, at: remote)
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -45,14 +45,21 @@ extension GhosttyTerminalView {
remoteCommand: location.remoteLoginShellCommand(script), interactive: true)
let reconnect = location.sshCommandLine(
remoteCommand: location.remoteLoginShellCommand(reconnectScript), interactive: true)
let stamp = ZmxSessionLauncher.redialStamp(for: location)
try? FileManager.default.createDirectory(
at: stamp.deletingLastPathComponent(), withIntermediateDirectories: true)
guard location.isCodespace else {
return ["/bin/sh", "-c", SSHReconnectLoop.script(connect: connect, reconnect: reconnect)]
return [
"/bin/sh", "-c",
SSHReconnectLoop.script(connect: connect, reconnect: reconnect, redialStamp: stamp.path),
]
}
return [
"/bin/sh", "-c",
SSHReconnectLoop.codespaceScript(
connect: connect, reconnect: reconnect,
pauseMarker: CodespaceDialBreaker.reconnectMarker(for: location).path),
pauseMarker: CodespaceDialBreaker.reconnectMarker(for: location).path,
redialStamp: stamp.path),
]
}

Expand Down
Loading
Loading