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
14 changes: 11 additions & 3 deletions GraphcodeKit/Sources/Domain/CodespaceDialSchedule.swift
Original file line number Diff line number Diff line change
Expand Up @@ -6,20 +6,28 @@ import Foundation
/// human's per-user Codespaces rate limit — around a hundred while a codespace is
/// starting. Retrying on each reader's own clock spent that limit during every outage,
/// so all dialers share this one schedule, counted from the first failure: retry freely
/// for a minute, hold until the third, retry again until the fourth, then pause until a
/// human asks to reconnect.
/// for a minute, hold until the third, retry again until the fourth, then pause: one dial
/// every `slowRetryInterval` until the codespace answers, or at once when a human asks to
/// reconnect. The pause never ends in silence, because a codespace restarted from
/// outside graphcode can take longer than four minutes to come back, and its loops
/// should recover without a human finding them dead.
///
/// The daemon applies it in `CodespaceDialBreaker`; a terminal pane applies the same
/// numbers inside its shell loop (`SSHReconnectLoop`), which runs in another process.
public struct CodespaceDialSchedule: Equatable, Sendable {
public var freeRetryWindow: Int
public var holdUntil: Int
public var pauseAfter: Int
public var slowRetryInterval: Int

public init(freeRetryWindow: Int = 60, holdUntil: Int = 180, pauseAfter: Int = 240) {
public init(
freeRetryWindow: Int = 60, holdUntil: Int = 180, pauseAfter: Int = 240,
slowRetryInterval: Int = 60
) {
self.freeRetryWindow = freeRetryWindow
self.holdUntil = holdUntil
self.pauseAfter = pauseAfter
self.slowRetryInterval = slowRetryInterval
}

public static let standard = CodespaceDialSchedule()
Expand Down
32 changes: 18 additions & 14 deletions GraphcodeKit/Sources/Domain/SSHReconnectLoop.swift
Original file line number Diff line number Diff line change
Expand Up @@ -41,11 +41,12 @@ public enum SSHReconnectLoop {

/// A Codespace surface's loop: the same dials and exit handling, retried on
/// `CodespaceDialSchedule` instead of forever, because every gh run spends the human's
/// Codespaces rate limit (issue #480). Past `schedule.pauseAfter` it waits for Enter,
/// which also touches `pauseMarker` so `graphcoded`'s `CodespaceDialBreaker` resumes
/// the codespace's reads and ensures with it. The app touches the same marker when a
/// loop of this codespace is selected, which restarts the schedule and redials from any
/// wait — held or paused — within a second.
/// Codespaces rate limit (issue #480). Past `schedule.pauseAfter` it redials once per
/// `schedule.slowRetryInterval`, or at once on Enter, which also touches `pauseMarker`
/// so `graphcoded`'s `CodespaceDialBreaker` resumes the codespace's reads and ensures
/// with it. The app touches the same marker when a loop of this codespace is selected,
/// which restarts the schedule and redials from any wait — held or paused — within a
/// second.
///
/// The outage clock restarts only after a dial that lasted `upAfter`, longer than the
/// five minutes gh can spend waiting for a codespace to start before failing — a
Expand Down Expand Up @@ -74,13 +75,15 @@ public enum SSHReconnectLoop {
// does not count.
let asked = "{ [ -e \"$gc_stamp\" ] && [ \(marker) -nt \"$gc_stamp\" ]; }"
let enter = "mkdir -p \(directory) 2>/dev/null; touch \(marker) 2>/dev/null"
// On a tty the pause polls, so a selection can end it too. Off one, `read` blocks as
// it always did: bash 3.2's `read -t` answers 1 for a timeout and for end of input
// alike, and a closed stdin would spin. A pane always has a tty.
// On a tty the pause polls, so a selection can end it too, and runs out after the slow
// retry interval to redial on the same outage clock. Off one, `read` blocks as it
// always did: bash 3.2's `read -t` answers 1 for a timeout and for end of input alike,
// and a closed stdin would spin. A pane always has a tty.
let pause =
"if [ -t 0 ]; then while :; do \(asked) && break; "
+ "read -t 1 gc_line && { \(enter); break; }; done; "
+ "else read gc_line || exit 0; \(enter); fi; "
"gc_ask=; if [ -t 0 ]; then gc_left=\(schedule.slowRetryInterval); "
+ "while [ \"$gc_left\" -gt 0 ]; do \(asked) && { gc_ask=1; break; }; "
+ "read -t 1 gc_line && { \(enter); gc_ask=1; break; }; gc_left=$((gc_left - 1)); done; "
+ "else read gc_line || exit 0; \(enter); gc_ask=1; fi; "
let waitOrAsk =
"gc_wait_or_ask() { gc_left=$1; while [ \"$gc_left\" -gt 0 ]; do "
+ "\(asked) && return 0; sleep 1; gc_left=$((gc_left - 1)); done; return 1; }; "
Expand All @@ -93,9 +96,10 @@ public enum SSHReconnectLoop {
+ "while :; do gc_out=$(($(date +%s) - gc_down)); "
+ "if [ \"$gc_out\" -ge \(schedule.pauseAfter) ]; then "
+ #"printf '\033[1;33m── Codespace still unreachable (exit %s). Paused to save your "#
+ #"Codespaces API quota. Press Enter or select the loop to reconnect, Ctrl-C to "#
+ #"close. ──\033[0m\r\n' "$gc_rc"; "#
+ pause + "\(restart); "
+ #"Codespaces API quota; retrying every %ss. Press Enter or select the loop to "#
+ #"reconnect now, Ctrl-C to close. ──\033[0m\r\n' "$gc_rc" "#
+ "\(schedule.slowRetryInterval); "
+ pause + "[ -n \"$gc_ask\" ] && { \(restart); }; "
+ "else "
+ "if [ \"$gc_out\" -ge \(schedule.freeRetryWindow) ] "
+ "&& [ \"$gc_out\" -lt \(schedule.holdUntil) ]; then "
Expand Down
10 changes: 10 additions & 0 deletions GraphcodeKit/Sources/GraphStore.swift
Original file line number Diff line number Diff line change
Expand Up @@ -4863,6 +4863,16 @@ public actor GraphStore {
for node in graph.nodes where node.runsUnattended && !node.isResolved {
ensureSession(node)
}
// The children `pilotComposite` started live on the composite's sub-graph, not in
// `graph.nodes`, and a reboot kills their sessions just the same.
for composite in graph.nodes
where !composite.isResolved
&& (composite.pilotState == .piloted || composite.pilotState == .armed)
{
for child in composite.subGraph?.nodes ?? [] where child.runsUnattended && !child.isResolved {
ensureSession(child)
}
}
let finished = graph.nodes.filter {
$0.runsUnattended && $0.isResolved && $0.launchFailure == nil
}
Expand Down
19 changes: 16 additions & 3 deletions GraphcodeKit/Sources/ProjectRegistry.swift
Original file line number Diff line number Diff line change
Expand Up @@ -353,15 +353,20 @@ public actor ProjectRegistry {
/// dial is a bare `zmx get` — `remoteEnsureInvocation` keeps the hooks write and the
/// file delivery behind that check precisely so this can be cheap — multiplexed onto
/// the host's existing `ControlMaster` connection.
static let remoteLivenessSweepInterval: Duration = .seconds(60)
///
/// A codespace is swept on every `codespaceSweepTicks`th tick only: any of its dials can
/// fall back to gh and spend the human's Codespaces API quota (issue #480).
static let remoteLivenessSweepInterval: Duration = .seconds(30)
static let codespaceSweepTicks = 2

/// Generous next to the remote sweep's minute: on a healthy machine the condemned
/// Generous next to the remote sweep's interval: on a healthy machine the condemned
/// list is empty and a tick is one file read, but a tick that finds work spawns
/// processes, and a session that survived three confirmed kill attempts is not going
/// to die to a faster clock.
static let condemnedReapInterval: Duration = .seconds(300)

private var remoteSweeper: Task<Void, Never>?
private var remoteSweepTick = 0
private var condemnedReaper: Task<Void, Never>?

/// Started by the first remote project this daemon loads and left running: a store is
Expand Down Expand Up @@ -393,11 +398,19 @@ public actor ProjectRegistry {
/// `RemoteEnsureGate`: one ensure per node at a time, so a slow tick cannot pile a
/// second dial onto the same session.
private func sweepRemoteSessions() async {
for (path, store) in stores where RemoteProjectLocation.parse(projectPath: path) != nil {
remoteSweepTick += 1
for (path, store) in stores {
guard let location = RemoteProjectLocation.parse(projectPath: path),
Self.sweeps(location, onTick: remoteSweepTick)
else { continue }
await store.ensureUnattendedSessionsAlive()
}
}

static func sweeps(_ location: RemoteProjectLocation, onTick tick: Int) -> Bool {
!location.isCodespace || tick % codespaceSweepTicks == 0
}

/// Takes or drops the sleep assertion to match what is running right now
/// (`AwakeAssertion`), across every open project.
///
Expand Down
39 changes: 34 additions & 5 deletions GraphcodeKit/Sources/Sessions/CodespaceDialBreaker.swift
Original file line number Diff line number Diff line change
Expand Up @@ -10,19 +10,30 @@ import Foundation
/// newer than the outage clears it — holding or paused alike — so the next read or
/// ensure dials at once. A file rather than a daemon command because the pane is a shell
/// loop in the app's process, and a `stat` spends nothing.
///
/// Paused, one dial per `slowRetryInterval` still goes through, whichever reader or
/// ensure asks first. The one that reaches the codespace clears the outage for all of
/// them and calls `onRecovered`, which is how the finished loops on the codespace get
/// restored as well as the running ones.
public actor CodespaceDialBreaker {
static let shared = CodespaceDialBreaker()
static let shared = CodespaceDialBreaker(onRecovered: { location in
ZmxSessionLauncher.markRedialed(location)
})

private let schedule: CodespaceDialSchedule
private let markerDirectory: URL
private let onRecovered: @Sendable (RemoteProjectLocation) -> Void
private var downSince: [String: Date] = [:]
private var lastPausedDial: [String: Date] = [:]

init(
schedule: CodespaceDialSchedule = .standard,
markerDirectory: URL = CodespaceDialBreaker.defaultMarkerDirectory
markerDirectory: URL = CodespaceDialBreaker.defaultMarkerDirectory,
onRecovered: @escaping @Sendable (RemoteProjectLocation) -> Void = { _ in }
) {
self.schedule = schedule
self.markerDirectory = markerDirectory
self.onRecovered = onRecovered
}

public static var defaultMarkerDirectory: URL {
Expand Down Expand Up @@ -60,21 +71,39 @@ public actor CodespaceDialBreaker {
func permits(_ location: RemoteProjectLocation, now: Date = Date()) -> Bool {
guard location.isCodespace, let since = downSince[location.host] else { return true }
if reconnectRequested(for: location, after: since) {
downSince.removeValue(forKey: location.host)
clearOutage(location.host)
return true
}
switch schedule.verdict(secondsDown: now.timeIntervalSince(since)) {
case .dial: return true
case .hold: return false
case .paused:
let pausedAt = since.addingTimeInterval(TimeInterval(schedule.pauseAfter))
let last = lastPausedDial[location.host] ?? pausedAt
guard now.timeIntervalSince(last) >= TimeInterval(schedule.slowRetryInterval) else {
return false
}
lastPausedDial[location.host] = now
return true
}
return schedule.verdict(secondsDown: now.timeIntervalSince(since)) == .dial
}

func record(_ location: RemoteProjectLocation, reached: Bool, now: Date = Date()) {
guard location.isCodespace else { return }
if reached {
downSince.removeValue(forKey: location.host)
guard downSince[location.host] != nil else { return }
clearOutage(location.host)
onRecovered(location)
} else if downSince[location.host] == nil {
downSince[location.host] = now
}
}

private func clearOutage(_ host: String) {
downSince.removeValue(forKey: host)
lastPausedDial.removeValue(forKey: host)
}

private func reconnectRequested(for location: RemoteProjectLocation, after since: Date) -> Bool {
let marker = Self.reconnectMarker(for: location, in: markerDirectory)
guard
Expand Down
27 changes: 21 additions & 6 deletions GraphcodeKit/Sources/Sessions/ZmxSessionLauncher.swift
Original file line number Diff line number Diff line change
Expand Up @@ -2515,10 +2515,11 @@ public enum ZmxSessionLauncher {
/// 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).
/// On a codespace the probe runs only when a pane of that host has redialed since the
/// last probe that answered (`redialStamp`), or the codespace answered after an outage:
/// a codespace dial spends the human's API quota (issue #480). A plain ssh host is
/// probed on every sweep, multiplexed over its `ControlMaster`, so its finished loops
/// come back with no pane open.
///
/// `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.
Expand Down Expand Up @@ -2552,9 +2553,22 @@ public enum ZmxSessionLauncher {
.appendingPathComponent("\(location.host).redial")
}

/// Touches `redialStamp` from the daemon: a codespace that just answered after an
/// outage may have restarted under finished loops whose panes are closed, and nothing
/// else would ask `restoreRebootedRemote` to probe it.
static func markRedialed(_ location: RemoteProjectLocation) {
let stamp = redialStamp(for: location)
try? FileManager.default.createDirectory(
at: stamp.deletingLastPathComponent(), withIntermediateDirectories: true)
FileManager.default.createFile(atPath: stamp.path, contents: nil)
try? FileManager.default.setAttributes(
[.modificationDate: Date()], ofItemAtPath: stamp.path)
}

/// 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.
/// state `restoreRebootedRemote` keeps, and always yes for a plain ssh host. 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()

Expand All @@ -2570,6 +2584,7 @@ public enum ZmxSessionLauncher {
}

func panesRedialed(_ location: RemoteProjectLocation) -> Bool {
guard location.isCodespace else { return true }
guard
let touched =
(try? FileManager.default.attributesOfItem(
Expand Down
52 changes: 46 additions & 6 deletions graphcode/Tests/CodespaceDialScheduleTests.swift
Original file line number Diff line number Diff line change
@@ -1,10 +1,11 @@
import ComposableArchitecture
import Foundation
import Testing

@testable import GraphcodeKit

/// One outage schedule for every Codespace dialer (issue #480): retry freely for a
/// minute, hold until the third, retry until the fourth, then pause until a human asks.
/// minute, hold until the third, retry until the fourth, then pause to a slow retry.
/// Every gh run spends the human's Codespaces rate limit, so a dialer that retried on
/// its own clock forever spent it during every outage.
///
Expand Down Expand Up @@ -76,6 +77,7 @@ struct CodespaceDialScheduleTests {
#expect(schedule.verdict(secondsDown: 239) == .dial)
#expect(schedule.verdict(secondsDown: 240) == .paused)
#expect(schedule.verdict(secondsDown: 86_400) == .paused)
#expect(schedule.slowRetryInterval == 60)
}

// MARK: - The daemon's breaker
Expand All @@ -91,7 +93,45 @@ struct CodespaceDialScheduleTests {
#expect(await !breaker.permits(codespace, now: down.addingTimeInterval(90)))
#expect(await breaker.permits(codespace, now: down.addingTimeInterval(200)))
#expect(await !breaker.permits(codespace, now: down.addingTimeInterval(250)))
#expect(await !breaker.permits(codespace, now: down.addingTimeInterval(86_400)))
}

@Test
func aPausedCodespaceIsStillDialedOncePerSlowInterval() async throws {
// A codespace restarted from outside graphcode can take longer than the four minutes
// before the pause; its loops have to come back without a human selecting one.
let breaker = CodespaceDialBreaker(markerDirectory: try scratch())
let down = Date(timeIntervalSince1970: 1_000_000)
await breaker.record(codespace, reached: false, now: down)

#expect(await !breaker.permits(codespace, now: down.addingTimeInterval(299)))
#expect(await breaker.permits(codespace, now: down.addingTimeInterval(300)))
// One dial per interval for the whole codespace, not one per reader or loop.
#expect(await !breaker.permits(codespace, now: down.addingTimeInterval(300)))
#expect(await !breaker.permits(codespace, now: down.addingTimeInterval(359)))
#expect(await breaker.permits(codespace, now: down.addingTimeInterval(360)))
// A failed slow dial leaves the outage clock where it was.
await breaker.record(codespace, reached: false, now: down.addingTimeInterval(365))
#expect(await !breaker.permits(codespace, now: down.addingTimeInterval(400)))
#expect(await breaker.permits(codespace, now: down.addingTimeInterval(86_400)))
}

@Test
func aCodespaceThatAnswersAfterAnOutageAsksForTheRebootProbe() async throws {
let recovered = LockIsolated<[String]>([])
let breaker = CodespaceDialBreaker(
markerDirectory: try scratch(),
onRecovered: { location in recovered.withValue { $0.append(location.host) } })
let down = Date(timeIntervalSince1970: 1_000_000)

await breaker.record(codespace, reached: true, now: down)
#expect(recovered.value.isEmpty)

await breaker.record(codespace, reached: false, now: down)
await breaker.record(codespace, reached: true, now: down.addingTimeInterval(600))
await breaker.record(codespace, reached: true, now: down.addingTimeInterval(601))

#expect(recovered.value == [codespace.host])
#expect(await breaker.permits(codespace, now: down.addingTimeInterval(602)))
}

@Test
Expand Down Expand Up @@ -128,13 +168,13 @@ struct CodespaceDialScheduleTests {
FileManager.default.createFile(atPath: marker.path, contents: nil)
try FileManager.default.setAttributes(
[.modificationDate: down.addingTimeInterval(-60)], ofItemAtPath: marker.path)
#expect(await !breaker.permits(codespace, now: down.addingTimeInterval(300)))
#expect(await !breaker.permits(codespace, now: down.addingTimeInterval(250)))

try FileManager.default.setAttributes(
[.modificationDate: down.addingTimeInterval(299)], ofItemAtPath: marker.path)
#expect(await breaker.permits(codespace, now: down.addingTimeInterval(300)))
[.modificationDate: down.addingTimeInterval(249)], ofItemAtPath: marker.path)
#expect(await breaker.permits(codespace, now: down.addingTimeInterval(250)))
// Resumed, not a one-off: the next read goes through too.
#expect(await breaker.permits(codespace, now: down.addingTimeInterval(301)))
#expect(await breaker.permits(codespace, now: down.addingTimeInterval(251)))
}

@Test
Expand Down
25 changes: 25 additions & 0 deletions graphcode/Tests/CodespaceSelectionTests.swift
Original file line number Diff line number Diff line change
Expand Up @@ -212,4 +212,29 @@ struct CodespaceSelectionTests {

#expect(await waitFor(within: .seconds(5)) { lines(in: log) > pausedAt })
}

@Test
func aPausedPaneRedialsOnItsOwnAfterTheSlowInterval() async throws {
let directory = try scratch()
defer { try? FileManager.default.removeItem(at: directory) }
let log = directory.appendingPathComponent("dials")
let marker = directory.appendingPathComponent("space.reconnect")
let dial = "(echo dial >> \(RemoteProjectLocation.shellQuoted(log.path)); exit 1)"
let slow = CodespaceDialSchedule(
freeRetryWindow: 1, holdUntil: 3, pauseAfter: 4, slowRetryInterval: 3)
let pane = try startPane(
SSHReconnectLoop.codespaceScript(
connect: dial, reconnect: dial, pauseMarker: marker.path, schedule: slow),
onATTY: true)
defer { pane.process.terminate() }

#expect(await waitFor { lines(in: log) >= 3 })
try await Task.sleep(for: .seconds(5))
let pausedAt = lines(in: log)

// Nobody presses Enter or selects the loop.
#expect(await waitFor(within: .seconds(10)) { lines(in: log) >= pausedAt + 2 })
#expect(pane.process.isRunning)
#expect(!FileManager.default.fileExists(atPath: marker.path))
}
}
Loading
Loading