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)