From e8c060b6c3ada89b4795494a9e7e4672863d40e0 Mon Sep 17 00:00:00 2001 From: scgopi Date: Wed, 30 Sep 2026 22:21:45 -0700 Subject: [PATCH] Spool remote ensure payloads over stdin so big-memory loops reattach The daemon's remote ensure carried the CLI shim, the briefing (twice for Copilot) and the loop's wake digest inline in the ssh remote command, which the host runs as a single `zsh -c` argument. Linux refuses one argument over 128 KiB (MAX_ARG_STRLEN) with E2BIG before any of it runs, so loops with a large memory failed every ensure silently: no dial-log line on the host, and after a reboot their panes sat on "reboot wait-daemon" while siblings with smaller memories came back. The ensure now spools every delivered file from the dial's stdin (length and SHA-256 checked, 60s alarm), so its command line no longer grows with the loop's memory. The pane's inline delivery is deflated for margin. A failed ensure is recorded in graphcoded.log with the node, argv and stdin sizes and ssh's output. Co-Authored-By: Claude Opus 5.5 Signed-off-by: scgopi --- .../Sources/Sessions/RemoteGraphAccess.swift | 128 ++++++++++++++++-- .../Sources/Sessions/ZmxSessionLauncher.swift | 103 +++++++++++--- .../Tests/RemoteSessionLaunchTests.swift | 109 ++++++++++++++- .../Tests/RemoteSessionResumeTests.swift | 5 +- 4 files changed, 312 insertions(+), 33 deletions(-) diff --git a/GraphcodeKit/Sources/Sessions/RemoteGraphAccess.swift b/GraphcodeKit/Sources/Sessions/RemoteGraphAccess.swift index 58fe0dad..b9da814b 100644 --- a/GraphcodeKit/Sources/Sessions/RemoteGraphAccess.swift +++ b/GraphcodeKit/Sources/Sessions/RemoteGraphAccess.swift @@ -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 @@ -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" @@ -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=$?; " @@ -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: @@ -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 [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( @@ -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, @@ -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( @@ -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. @@ -1960,7 +1981,7 @@ 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 @@ -1968,7 +1989,8 @@ public enum ZmxSessionLauncher { // 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 @@ -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 @@ -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 @@ -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 diff --git a/graphcode/Tests/RemoteSessionLaunchTests.swift b/graphcode/Tests/RemoteSessionLaunchTests.swift index 6efd4705..6f56aaa4 100644 --- a/graphcode/Tests/RemoteSessionLaunchTests.swift +++ b/graphcode/Tests/RemoteSessionLaunchTests.swift @@ -1,3 +1,4 @@ +import ComposableArchitecture import Foundation import Testing @@ -374,8 +375,7 @@ struct RemoteSessionLaunchTests { let tokens = script.split(separator: " ") for token in tokens.reversed() { let raw = token.trimmingCharacters(in: CharacterSet(charactersIn: "'")) - guard let data = Data(base64Encoded: raw), - let manifest = try? JSONSerialization.jsonObject(with: data) as? [String: String] + guard let manifest = RemoteGraphAccess.manifest(fromInstallerArgument: raw) else { continue } return manifest.keys.sorted() } @@ -516,3 +516,108 @@ extension RemoteSessionLaunchTests { #expect(!files.keys.contains { $0.hasSuffix(NodeMemory.promptFileName) }) } } + +/// The daemon's ensure against the remote host's argument limit, and the stdin spool +/// that keeps it under (`RemotePayloadSpool`). +@Suite +struct RemoteEnsureArgumentLimitTests { + /// Linux's `MAX_ARG_STRLEN`: the most one argv string may hold, NUL included. sshd + /// hands the whole remote command to the login shell as the single `-c` argument, so + /// an ensure over it fails `execve` with E2BIG before its first fragment — the dial + /// log never hears of it, and after a reboot the loop's pane waits on graphcoded + /// forever while siblings with smaller memories come back. + private static let linuxArgumentLimit = 131_072 + + @Test(arguments: [CLISessionBackendKind.copilotCLI, .claudeCode]) + func anEnsureFitsOneLinuxArgumentWithAFullMemory(backend: CLISessionBackendKind) throws { + let codespace = RemoteProjectLocation( + user: nil, host: "curly-space-guide", port: nil, remotePath: "/workspaces/widget", + isCodespace: true) + let node = LoopNode( + title: "Refresh", loopType: .timeBased, + triggerPrompt: "/loop 1h refresh the factory", backend: backend) + defer { NodeMemory.remove(projectPath: codespace.projectPath, nodeID: node.id) } + for pass in 0.. empty.input.count) + } + + @Test + func aSpooledDeliveryLandsThroughTheDaemonsOwnPTY() async throws { + // The production path end to end, minus ssh: `runRemoteRetryingCollecting` writes + // the payload into a canonical-mode PTY, which drops any line past MAX_CANON — a + // payload well over that, and over Linux's argument limit, must land byte-exact. + let home = FileManager.default.temporaryDirectory + .appendingPathComponent("spool-\(UUID().uuidString)") + try FileManager.default.createDirectory(at: home, withIntermediateDirectories: true) + defer { try? FileManager.default.removeItem(at: home) } + let noise = (0..<2000).map { _ in + Data((0..<57).map { _ in UInt8.random(in: 0...255) }).base64EncodedString() + } + let text = "'é' $(x)\n" + noise.joined(separator: "\n") + let spool = RemotePayloadSpool() + let install = try #require( + RemoteGraphAccess.installerScript( + files: ["~/.graphcode/big.txt": text], neutered: false, spool: spool)) + let prelude = try #require(spool.prelude) + let script = + "export HOME=\(RemoteProjectLocation.shellQuoted(home.path)); " + + "\(prelude){ \(install); }; gc_rc=$?; \(spool.cleanup); exit $gc_rc" + + #expect(script.utf8.count < 16_384) + #expect(spool.input.count > Self.linuxArgumentLimit) + let result = await ZmxSessionLauncher.runRemoteRetryingCollecting( + ["/bin/sh", "-c", script], attempts: 1, standardInput: spool.input, + timeout: .seconds(60)) + #expect(result.succeeded, "\(result.output.suffix(300))") + let landed = try String( + contentsOf: home.appendingPathComponent(".graphcode/big.txt"), encoding: .utf8) + #expect(landed == text) + } + + @Test + func aFailedEnsureIsLoggedWithItsNodeAndSize() { + let node = LoopNode( + title: "Refresh", loopType: .timeBased, triggerPrompt: "/loop 1h refresh") + let lines = LockIsolated<[String]>([]) + let tap = DaemonLog.shared.tap { line in lines.withValue { $0.append(line) } } + defer { DaemonLog.shared.untap(tap) } + + ZmxSessionLauncher.recordEnsureFailure( + node: node, invocation: ["/usr/bin/gh", String(repeating: "x", count: 140_000)], + input: Data(count: 42), output: "zsh: argument list too long: zsh\r\n") + + let line = lines.value.first { $0.contains("event=remote-ensure-failed") } ?? "" + #expect(line.contains("node=\(node.id.uuidString)")) + #expect(line.contains("argvBytes=140000")) + #expect(line.contains("stdinBytes=42")) + #expect(line.contains("argument list too long")) + } +} diff --git a/graphcode/Tests/RemoteSessionResumeTests.swift b/graphcode/Tests/RemoteSessionResumeTests.swift index 984b8f27..d904bf6d 100644 --- a/graphcode/Tests/RemoteSessionResumeTests.swift +++ b/graphcode/Tests/RemoteSessionResumeTests.swift @@ -234,11 +234,10 @@ struct RemoteSessionResumeTests { script.components(separatedBy: ">> ").dropFirst() .allSatisfy { $0.hasPrefix(DialLog.logExpression) }) - // Not a manifest entry — the manifest is the one token that base64-decodes to JSON. + // Not a manifest entry — the manifest is the one token that decodes to JSON. let files = try #require( script.split(whereSeparator: { $0 == " " || $0 == "'" }) - .compactMap { Data(base64Encoded: String($0)) } - .compactMap { try? JSONSerialization.jsonObject(with: $0) as? [String: String] } + .compactMap { RemoteGraphAccess.manifest(fromInstallerArgument: String($0)) } .first) #expect(files[RemoteGraphAccess.cliInstallPath] != nil) #expect(files[RemoteGraphAccess.shimStampPath] == nil)