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
128 changes: 117 additions & 11 deletions GraphcodeKit/Sources/Sessions/RemoteGraphAccess.swift
Original file line number Diff line number Diff line change
Expand Up @@ -101,7 +101,8 @@ public enum RemoteGraphAccess {

/// A shell fragment that lands `files` (home-relative path → content) on the remote
/// host, or `nil` when there's nothing to send. One `python3 -c` with a base64 JSON
/// manifest rather than heredocs or scp: a single argument survives every quoting
/// manifest (`deflated` first — see there for the size bound it keeps the dial
/// under) rather than heredocs or scp: a single argument survives every quoting
/// layer between here and the remote shell, needs no extra ssh round-trip, and
/// content can't collide with a delimiter. Neutered because delivery must never block
/// the launch it precedes — a session without its briefing is the old behaviour, which
Expand Down Expand Up @@ -145,21 +146,29 @@ public enum RemoteGraphAccess {
/// brief is a pointer at a file that isn't there. That caller
/// (`ZmxSessionLauncher.remotePromptDelivery`) chains the launch behind this command's
/// exit status instead, so a failed delivery costs a retry rather than a blind pass.
///
/// With a `spool`, the manifest rides the dial's stdin instead of the command line
/// (`RemotePayloadSpool`) and the argument names the file it was spooled to.
public static func installerScript(
files: [String: String], receipt: (path: String, content: String)? = nil,
neutered: Bool = true
neutered: Bool = true, spool: RemotePayloadSpool? = nil
) -> String? {
guard !files.isEmpty else { return nil }
let manifest = files.mapValues { Data($0.utf8).base64EncodedString() }
guard let json = try? JSONSerialization.data(withJSONObject: manifest, options: [.sortedKeys])
guard !files.isEmpty,
let json = try? JSONSerialization.data(withJSONObject: files, options: [.sortedKeys])
else { return nil }
let payload = deflated(json)
let encoded = payload.data.base64EncodedString()
let source = "(open(a[1:]).read() if a[:1]==\"@\" else a)"
let decode =
payload.deflated
? "zlib.decompress(base64.b64decode(\(source)),-15)" : "base64.b64decode(\(source))"
let program =
"import base64,json,os,sys,tempfile; "
+ "m=json.loads(base64.b64decode(sys.argv[1])); "
"import base64,json,os,sys,tempfile,zlib; "
+ "a=sys.argv[1]; m=json.loads(\(decode)); "
+ "exec('def w(p,c):\\n"
+ " d=os.path.dirname(os.path.expanduser(p))\\n"
+ " os.makedirs(d,exist_ok=True)\\n"
+ " b=base64.b64decode(c)\\n"
+ " b=c.encode()\\n"
+ " if p.endswith(\"/bridge-state.json\"):\\n"
+ " fd,t=tempfile.mkstemp(dir=d)\\n"
+ " n=os.write(fd,b)\\n"
Expand All @@ -171,9 +180,14 @@ public enum RemoteGraphAccess {
+ "'); "
+ "[w(p,c) for p,c in sorted(m.items())]; "
+ "len(sys.argv)>2 and open(os.path.expanduser(sys.argv[2]),'w').write(sys.argv[3])"
var argv = ["python3", "-c", program, json.base64EncodedString()]
if let receipt { argv += [receipt.path, receipt.content] }
let install = argv.map(RemoteProjectLocation.shellQuoted).joined(separator: " ")
let argument =
spool.map { "\"@\($0.reference(for: encoded))\"" }
?? RemoteProjectLocation.shellQuoted(encoded)
var tail: [String] = []
if let receipt { tail = [receipt.path, receipt.content] }
let install =
(["python3", "-c", program].map(RemoteProjectLocation.shellQuoted) + [argument]
+ tail.map(RemoteProjectLocation.shellQuoted)).joined(separator: " ")
// stderr into a variable, stdout to `/dev/null` — `2>&1 >/dev/null` in that order,
// so the substitution keeps the diagnosis and drops the noise.
return "gc_di_err=$(\(install) 2>&1 >/dev/null); gc_di_rc=$?; "
Expand All @@ -196,6 +210,35 @@ public enum RemoteGraphAccess {
+ (neutered ? "true" : "[ \"$gc_di_rc\" -eq 0 ]")
}

/// The manifest as raw DEFLATE (`zlib.decompress(…, -15)` on the host), where this
/// platform's Foundation can make it. The whole remote command reaches the host's
/// login shell as one `-c` argument, and Linux refuses any single argument over
/// 128 KiB (`MAX_ARG_STRLEN`) with E2BIG before a byte of it runs. The shim alone was
/// ~80 KB once base64'd twice; add a briefing (twice, for Copilot) and a long-lived
/// loop's wake digest and the ensure crossed it — so exactly the loops with the most
/// memory were never restored after a reboot, and their dial log never said why.
static func deflated(_ json: Data) -> (data: Data, deflated: Bool) {
#if canImport(Darwin)
if let packed = try? (json as NSData).compressed(using: .zlib) as Data {
return (packed, true)
}
#endif
return (json, false)
}

/// The inverse of `installerScript`'s payload argument, for tests that assert on what a
/// delivery carries.
static func manifest(fromInstallerArgument argument: String) -> [String: String]? {
guard let data = Data(base64Encoded: argument) else { return nil }
var json = data
#if canImport(Darwin)
if let inflated = try? (data as NSData).decompressed(using: .zlib) as Data {
json = inflated
}
#endif
return (try? JSONSerialization.jsonObject(with: json)) as? [String: String]
}

/// How much of a failed delivery's stderr reaches the dial log — the **last** bytes,
/// not the first. A python traceback opens with frames and interpreter paths and ends
/// with the line that names the fault, so keeping the head throws away the answer:
Expand Down Expand Up @@ -1300,3 +1343,66 @@ public enum RemoteGraphAccess {
main(sys.argv[1:])
"""#
}

/// Payloads that ride a dial's stdin rather than its command line.
///
/// sshd hands the remote command to the login shell as one `-c` argument, and Linux
/// refuses any single argument over 128 KiB (`MAX_ARG_STRLEN`) with E2BIG before a byte
/// of it runs. The daemon's ensure carried the CLI shim, the briefing (twice, for
/// Copilot) and the loop's wake digest inline, so the loops with the most memory crossed
/// it: their ensure failed before the first dial-log fragment, every sweep, and after a
/// reboot their panes waited on graphcoded forever. Spooled, the command line holds
/// only a byte count and digest per payload, whatever the loop remembers.
///
/// The stdin is the daemon's PTY in canonical mode (`PTYProcessSession`), which drops a
/// line past `MAX_CANON` — 1024 bytes on macOS — so each payload travels as base64
/// wrapped at 76 columns, which the installer's decode skips over.
public final class RemotePayloadSpool {
private var chunks: [Data] = []

public init() {}

/// Registers `text` and returns the shell expression naming the file it lands in.
func reference(for text: String) -> String {
var wrapped = ""
var line = Substring(text)
while !line.isEmpty {
wrapped += String(line.prefix(76)) + "\n"
line = line.dropFirst(76)
}
chunks.append(Data(wrapped.utf8))
return "$gc_pl\(chunks.count - 1)"
}

/// Everything the dial must write to its stdin, in the order the prelude reads it.
public var input: Data { chunks.reduce(Data(), +) }

/// Reads each payload off stdin into its own temporary file, verifying length and
/// digest — a short or garbled read leaves the file empty, which the installer then
/// reports as a failed delivery in the host's dial log. Then stdin is closed, so
/// nothing later in the script can block on a stream that never ends. The alarm
/// bounds a dial whose input never arrives.
var prelude: String? {
guard !chunks.isEmpty else { return nil }
let names = chunks.indices.map { "gc_pl\($0)" }
let program = """
import hashlib,signal,sys
signal.alarm(60)
a=sys.argv[1:]
for i in range(0,len(a),3):
n=int(a[i+1]); d=sys.stdin.buffer.read(n)
ok=len(d)==n and hashlib.sha256(d).hexdigest()==a[i+2]
open(a[i],"wb").write(d if ok else b"")
"""
let specs = zip(names, chunks).map { name, chunk in
"\"$\(name)\" \(chunk.count) \(GraphcodeSHA256.hex(chunk))"
}
return names.map { "\($0)=$(mktemp)" }.joined(separator: " && ")
+ " && python3 -c \(RemoteProjectLocation.shellQuoted(program)) "
+ specs.joined(separator: " ") + " 2>/dev/null; exec </dev/null; "
}

var cleanup: String {
"rm -f " + chunks.indices.map { "\"$gc_pl\($0)\"" }.joined(separator: " ")
}
}
103 changes: 86 additions & 17 deletions GraphcodeKit/Sources/Sessions/ZmxSessionLauncher.swift
Original file line number Diff line number Diff line change
Expand Up @@ -1755,6 +1755,20 @@ public enum ZmxSessionLauncher {
settings: GraphcodeSettings = GraphcodeSettingsStore.load(),
bridgeState: RemoteBridgeWireState? = nil, onlyAfterReboot: Bool = false
) -> [String]? {
remoteEnsureDial(
forNode: node, at: location, settings: settings, bridgeState: bridgeState,
onlyAfterReboot: onlyAfterReboot)?.invocation
}

/// `remoteEnsureInvocation` with the stdin it must be fed: every delivered file rides
/// there (`RemotePayloadSpool`), so the command line stays the same few kilobytes
/// however much the loop remembers.
static func remoteEnsureDial(
forNode node: LoopNode, at location: RemoteProjectLocation,
settings: GraphcodeSettings = GraphcodeSettingsStore.load(),
bridgeState: RemoteBridgeWireState? = nil, onlyAfterReboot: Bool = false
) -> (invocation: [String], input: Data)? {
let spool = RemotePayloadSpool()
let shedPrompt = ShedPromptReport()
guard
let zmxArguments = arguments(
Expand All @@ -1770,7 +1784,7 @@ public enum ZmxSessionLauncher {
// prefixed when the prompt was typed in full, which is the ordinary case.
let launchCommand = remoteQuotedCommand(["zmx"] + zmxArguments)
let run =
remotePromptDelivery(shedPrompt, forNode: node)
remotePromptDelivery(shedPrompt, forNode: node, spool: spool)
.map { "\($0) && \(launchCommand)" } ?? launchCommand
// Copilot only, and remote only: an unattended Copilot queues its `--interactive`
// goal behind a per-session folder-trust dialog that nobody is present to answer,
Expand All @@ -1796,7 +1810,8 @@ public enum ZmxSessionLauncher {
.map { $0 + "; " } ?? ""
let delivery =
remoteDeliveryScript(
forNode: node, at: location, settings: settings, bridgeState: bridgeState
forNode: node, at: location, settings: settings, bridgeState: bridgeState,
spool: spool
)
.map { $0 + "; " } ?? ""
let create = remoteCreateScript(
Expand Down Expand Up @@ -1847,7 +1862,13 @@ public enum ZmxSessionLauncher {
delivery, ifSessionMissing: check,
bridgeStateGeneration: bridgeState.map(\.generation))
+ "\(check) >/dev/null 2>&1\(bank) && { \(markerWrite); } || \(repair){ \(missing); }; }"
return location.sshInvocation(remoteCommand: location.remoteLoginShellCommand(script))
let spooled =
spool.prelude.map { "\($0){ \(script); }; gc_rc=$?; \(spool.cleanup); exit $gc_rc" }
?? script
return (
location.sshInvocation(remoteCommand: location.remoteLoginShellCommand(spooled)),
spool.input
)
}

/// The delivery, run when the session is missing **or** the host's shim is out of date.
Expand Down Expand Up @@ -1960,15 +1981,16 @@ public enum ZmxSessionLauncher {
public static func remoteDeliveryScript(
forNode node: LoopNode?, backend: CLISessionBackendKind? = nil,
at location: RemoteProjectLocation, settings: GraphcodeSettings,
bridgeState: RemoteBridgeWireState? = nil
bridgeState: RemoteBridgeWireState? = nil, spool: RemotePayloadSpool? = nil
) -> String? {
_ = bridgeState
// The shim's receipt, written only once every file has landed — see
// `installerScript`. It is what lets a later ensure skip a delivery it doesn't need
// without ever claiming a shim the host never received.
return RemoteGraphAccess.installerScript(
files: remoteDeliveryFiles(forNode: node, backend: backend, at: location, settings: settings),
receipt: (path: RemoteGraphAccess.shimStampPath, content: RemoteGraphAccess.cliShimStamp))
receipt: (path: RemoteGraphAccess.shimStampPath, content: RemoteGraphAccess.cliShimStamp),
spool: spool)
}

/// The delivery for a prompt that moved to a file (issue #57), as its own command
Expand All @@ -1988,10 +2010,12 @@ public enum ZmxSessionLauncher {
/// with it. The node then stays honestly not-running and the next liveness sweep
/// retries, which is the same posture `startRemote` already takes on a dial that fails.
static func remotePromptDelivery(
_ shedPrompt: ShedPromptReport, forNode node: LoopNode
_ shedPrompt: ShedPromptReport, forNode node: LoopNode, spool: RemotePayloadSpool? = nil
) -> String? {
guard let path = shedPrompt.remotePath, let text = shedPrompt.text else { return nil }
guard let install = RemoteGraphAccess.installerScript(files: [path: text], neutered: false)
guard
let install = RemoteGraphAccess.installerScript(
files: [path: text], neutered: false, spool: spool)
else { return nil }
let name = SurfaceRef(id: node.id, launchesClaudeCode: true).zmxSessionName
// Logged on the remote host's dial log rather than swallowed: a launch that never
Expand Down Expand Up @@ -2262,25 +2286,45 @@ public enum ZmxSessionLauncher {
standardInput: Data? = nil,
timeout: Duration? = nil
) async -> Bool {
await runRemoteRetryingCollecting(
invocation, attempts: attempts, standardInput: standardInput, timeout: timeout
).succeeded
}

/// `runRemoteRetrying`, keeping what the last attempt printed — ssh's own stderr is
/// the only witness to a dial that failed before the remote script began.
///
/// The input is written off the calling task: a PTY write blocks until ssh reads it,
/// which it does only once the connection is up, and a codespace that is starting can
/// hold that for minutes. Once the child exits its slave end closes and a pending
/// write fails rather than waiting forever.
static func runRemoteRetryingCollecting(
_ invocation: [String],
attempts: Int = 3,
standardInput: Data? = nil,
timeout: Duration? = nil
) async -> (succeeded: Bool, output: String) {
RemoteProjectLocation.prepareControlSocketDirectory()
guard attempts > 0 else { return false }
guard attempts > 0 else { return (false, "") }
var output = ""
for attempt in 1...attempts {
guard !Task.isCancelled else { return false }
guard !Task.isCancelled else { return (false, output) }
guard
let session = try? PTYProcessSession(
executable: invocation[0], arguments: Array(invocation.dropFirst()))
else { return false }
else { return (false, output) }
if let standardInput {
session.sendInput(String(decoding: standardInput, as: UTF8.self) + "\n")
}
if await waitForRemoteProcess(session, timeout: timeout).succeeded {
return true
let text = String(decoding: standardInput, as: UTF8.self) + "\n"
DispatchQueue.global(qos: .utility).async { session.sendInput(text) }
}
let result = await waitForRemoteProcess(session, timeout: timeout)
if result.succeeded { return (true, result.output) }
output = result.output
if attempt < attempts, !Task.isCancelled {
try? await Task.sleep(for: .seconds(1 << (attempt - 1)))
}
}
return false
return (false, output)
}

/// The ceiling on a one-shot remote read, below `GraphStore`'s 45s presence deadline so
Expand Down Expand Up @@ -2430,17 +2474,42 @@ public enum ZmxSessionLauncher {
// and the run must share a shell. A failure after the retries is the same posture
// as the local path: no UI here, the node's state stays honest, opening the loop
// retries.
if let ensure = remoteEnsureInvocation(
if let ensure = remoteEnsureDial(
forNode: node, at: location, bridgeState: bridgeState, onlyAfterReboot: onlyAfterReboot
) {
if await runRemoteRetrying(ensure) {
let result = await runRemoteRetryingCollecting(
ensure.invocation, standardInput: ensure.input.isEmpty ? nil : ensure.input)
if result.succeeded {
await CodespaceDialBreaker.shared.record(location, reached: true)
await CodespaceSSHUser.shared.learnIfNeeded(location)
} else if !Task.isCancelled {
recordEnsureFailure(
node: node, invocation: ensure.invocation, input: ensure.input, output: result.output)
}
}
await RemoteEnsureGate.shared.end(node.id, token: lease)
}

/// A failed ensure, in this daemon's log. The host's dial log hears of an ensure only
/// once its script runs, so a dial refused before that — E2BIG, auth, a host that is
/// down — left no trace on either machine, and its loop's pane waited forever.
static func recordEnsureFailure(
node: LoopNode, invocation: [String], input: Data, output: String
) {
let tail = String(
output.replacingOccurrences(of: "\r", with: " ")
.replacingOccurrences(of: "\n", with: " ")
.suffix(200))
DaemonLog.shared.record(
"remote-ensure-failed",
[
("node", node.id.uuidString),
("argvBytes", String(invocation.map(\.utf8.count).max() ?? 0)),
("stdinBytes", String(input.count)),
("output", tail.trimmingCharacters(in: .whitespaces)),
])
}

/// 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
Expand Down
Loading
Loading