diff --git a/GraphcodeKit/Sources/GraphStore.swift b/GraphcodeKit/Sources/GraphStore.swift index 3f09cd06..556629cb 100644 --- a/GraphcodeKit/Sources/GraphStore.swift +++ b/GraphcodeKit/Sources/GraphStore.swift @@ -250,6 +250,9 @@ public actor GraphStore { } private var pendingFollowUps: [PendingFollowUp] = [] + /// Each target's `node send` messages still being typed, chained so they land in the + /// order they were sent — see `typeAfterAcknowledging`. + private var sessionTyping: [UUID: (token: UUID, task: Task)] = [:] private var pendingDeliveryAttempts: Set = [] private var completedTimedOutDeliveries: [UUID: Bool] = [:] /// `drainPendingFollowUps` runs across several awaits, and the presence poll that @@ -1337,7 +1340,20 @@ public actor GraphStore { // before anyone is told what the graph looks like. Cycle re-entries run before // hand-off deliveries because a re-entry *queues* one; nudges last, since an // update's memory record must exist before its session is told to go look. - let errors = await drainAndBroadcast(broadcastErrors: broadcastErrors) + // A `node send` is acknowledged before the drain, not after it: the drain types any + // follow-up an idle target can take now, and the CLI's ten-second wait for its + // verdict is not the place to spend that. + var errors: [String] = [] + var acknowledged: LoopGraph? + if case .messageNode = command { + errors = await drainPendingErrors(broadcastErrors: broadcastErrors) + if errors.isEmpty { + await broadcast() + acknowledged = graph + } + } + errors += await drainAndBroadcast( + broadcastErrors: broadcastErrors, unlessStillAt: acknowledged) if let error = errors.first { return .rejected(message: error, graph: graph) } @@ -3679,10 +3695,10 @@ public actor GraphStore { /// message edge (`MessageBus`, the target backend's `sendInput`), so there is one /// definition of "may this session be typed into", not two. /// - /// A failure is said out loud rather than swallowed: an `.errorOccurred` goes to - /// every connection, which the app shows as its error banner and the CLI prints — - /// the whole point of the message was that a peer be told something, and pretending - /// it landed is the one wrong answer. + /// A message that cannot be typed now is said out loud rather than swallowed: an + /// `.errorOccurred` goes to every connection, which the app shows as its error banner + /// and the CLI prints. One that can is acknowledged once it is on the board and queued + /// for its session, and typed after — see `typeAfterAcknowledging`. private func deliverAdHocMessage( to nodeID: UUID, text: String, from senderID: UUID?, followUp: Bool = false, mirror: Bool = true, watchedPostID: Int? = nil @@ -3731,12 +3747,14 @@ public actor GraphStore { // A follow-up question to a finished loop whose session is still up reaches it. The // graph calls a resolved loop "not live" so edges and wakes leave it alone, but a // human asking what it did is the point of keeping the session; the answer changes - // nothing about how it resolved (#346). + // nothing about how it resolved (#346). Liveness is asked before the acknowledgement + // so a sender reporting to a finished parent still hears "staged" — one `zmx ls`, + // not a typing of the message. if target.state == .succeeded || target.state == .failed, target.backend.capabilities.supportsMidSessionInput, - await onSessionAlive?(target, graph.project.path) == true, - await deliverToSession(target, message) + await onSessionAlive?(target, graph.project.path) == true { + await typeAfterAcknowledging(message, to: target, resolved: true) return } if MessageBus.deliverability(to: target) != nil { @@ -3746,29 +3764,138 @@ public actor GraphStore { + "it will read it when it next wakes") return } - guard await deliverToSession(target, message) else { - // The transport can also fail because the session died after the graph last - // looked — a goal loop whose agent exited on its very first turn had no session - // left to type into, and (before sessions that answer while dead stopped passing - // the send gate) even a "delivered" that nobody received (issue #215). An - // unattended loop is the daemon's to keep alive, so a failed delivery is the - // moment to do exactly that: the ensure is create-only and husk-aware, so it - // relaunches precisely the dead case, the settle is the fresh session's boot - // beat, and the retry lands the message that would otherwise have sat staged - // until a wake that a dead loop has no way to know about. Attended loops stay - // human-timed — a turn-based session is respawned by a human opening it, not by - // a message arriving. - if target.runsUnattended, !target.isResolved { - ensureSession(target) - try? await Task.sleep(for: Self.respawnedSessionSettle) - if await deliverToSession(target, message) { return } - } - recordMemory(nodeID, "while you were away: \(message)") - announceError( - "delivery to \(target.title)'s session failed — message staged to its memory; " - + "it will read it when it next wakes") + await typeAfterAcknowledging(message, to: target, resolved: false) + } + + /// Types a `node send` message into its target's session after the command that carried + /// it has been acknowledged. + /// + /// Typing used to happen inside the request: the send gate's `zmx ls`, one `zmx send` + /// per chunk, the submit beat, then Enter. On a CPU-starved machine that took 10–30s + /// against the CLI's ten-second wait, so a message that landed was reported as exit 75, + /// "may still have been applied". The acknowledgement now means what a sender can rely + /// on: the message is on the board and queued for this session, and if it cannot be + /// typed it is staged to the loop's memory for its next wake. + /// + /// A failure found after the acknowledgement is logged and staged, not announced: + /// `.errorOccurred` reaches every connection, and with the sender gone the first CLI to + /// be waiting on a verdict of its own would take it for one. Every outcome is logged + /// against the request that was acknowledged (`send-typed`, `send-staged`). + /// + /// Bounded by `deliveryDeadline`, as the follow-up drain is: a typing that hangs would + /// otherwise hold this target's chain, and every follow-up and wake behind it, forever. + /// + /// A composite's child is typed inline instead. Its store is built for one command and + /// discarded (`runInSubGraph`), so a chain there would be a new chain per send — no + /// order between them, and nothing the parent could wait on. Its parent does not + /// acknowledge early either, so inline costs nothing that was ever saved. + private func typeAfterAcknowledging( + _ message: String, to target: LoopNode, resolved: Bool + ) async { + let context = DaemonRequestContext.fields + guard subGraphDepth == 0 else { + await typeLogged(message, to: target.id, resolved: resolved, context: context) return } + let previous = sessionTyping[target.id]?.task + let token = UUID() + let task = Task { [self] in + await previous?.value + await typeLogged(message, to: target.id, resolved: resolved, context: context) + if sessionTyping[target.id]?.token == token { + sessionTyping.removeValue(forKey: target.id) + } + if pendingFollowUps.contains(where: { $0.nodeID == target.id }) { + await drainPendingFollowUps() + } + } + sessionTyping[target.id] = (token, task) + } + + private func typeLogged( + _ message: String, to nodeID: UUID, resolved: Bool, context: [(String, String)] + ) async { + let started = Date() + let outcome = await withDeadline(deliveryDeadline) { + await self.typeAdHocMessage(message, to: nodeID, resolved: resolved) + } + let fields = + context + [ + ("node", nodeID.uuidString), + ("typed_ms", DaemonLog.milliseconds(Date().timeIntervalSince(started))), + ] + switch outcome { + case .typed: + DaemonLog.shared.record("send-typed", fields) + case .staged(let reason): + DaemonLog.shared.record("send-staged", fields + [("reason", reason)]) + case .targetGone: + DaemonLog.shared.record("send-dropped", fields + [("reason", "target-deleted")]) + case nil: + // The abandoned typing may still land; a copy in memory is the cheaper mistake + // than a message neither typed nor kept. + if let target = graph.nodes[id: nodeID] { + stageUntyped(message, to: target) + } + DaemonLog.shared.record("send-staged", fields + [("reason", "deadline")]) + } + } + + private enum TypingOutcome: Sendable { + case typed + case staged(reason: String) + case targetGone + } + + /// Waits until every acknowledged `node send` has been typed or staged. + public func finishSessionTyping() async { + while let typing = sessionTyping.values.first { + await typing.task.value + } + } + + private func typeAdHocMessage( + _ message: String, to nodeID: UUID, resolved: Bool + ) async -> TypingOutcome { + guard let target = graph.nodes[id: nodeID] else { return .targetGone } + // Cancelled means `typeLogged`'s deadline gave up on this typing and staged it: a + // cancelled send reads as failed, and retrying or staging again from here would + // respawn the loop and type into it behind the chain's next message. + if resolved { + if await deliverToSession(target, message) { return .typed } + if Task.isCancelled { return .staged(reason: "deadline") } + stageUntyped(message, to: target) + return .staged(reason: "session-gone") + } + if await deliverToSession(target, message) { return .typed } + if Task.isCancelled { return .staged(reason: "deadline") } + // The transport can also fail because the session died after the graph last + // looked — a goal loop whose agent exited on its very first turn had no session + // left to type into, and (before sessions that answer while dead stopped passing + // the send gate) even a "delivered" that nobody received (issue #215). An + // unattended loop is the daemon's to keep alive, so a failed delivery is the + // moment to do exactly that: the ensure is create-only and husk-aware, so it + // relaunches precisely the dead case, the settle is the fresh session's boot + // beat, and the retry lands the message that would otherwise have sat staged + // until a wake that a dead loop has no way to know about. Attended loops stay + // human-timed — a turn-based session is respawned by a human opening it, not by + // a message arriving. + if target.runsUnattended, !target.isResolved { + ensureSession(target) + try? await Task.sleep(for: Self.respawnedSessionSettle) + if Task.isCancelled { return .staged(reason: "deadline") } + if await deliverToSession(target, message) { return .typed } + if Task.isCancelled { return .staged(reason: "deadline") } + } + stageUntyped(message, to: target) + return .staged(reason: "delivery-failed") + } + + private func stageUntyped(_ message: String, to target: LoopNode) { + recordMemory(target.id, "while you were away: \(message)") + onAnnounceError?( + "delivery to \(target.title)'s session failed — message staged to its memory; " + + "it will read it when it next wakes") } /// `GraphCommand.broadcastMessage`, as one operation over the whole tree rather than a @@ -4011,6 +4138,13 @@ public actor GraphStore { case nil: break } + // A `node send` still being typed into this session was sent first. + if sessionTyping[pending.nodeID] != nil { + remaining.append(staged(pending)) + drainDeferred = remaining + drainInFlight = nil + continue + } let presence: Presence if let known = readings[pending.nodeID] { presence = known @@ -4505,7 +4639,11 @@ public actor GraphStore { /// The same settle-then-tell sequence `handle` ends with, for the paths that mutate /// outside a command — goal polling resolves nodes and fires edges too, and an edge /// fired from a poll must not wait for the next unrelated command to be delivered. - private func drainAndBroadcast(broadcastErrors: Bool = true) async -> [String] { + /// `unlessStillAt` skips the closing broadcast when the graph is exactly the one + /// already broadcast — a `node send` acknowledged before its drain. + private func drainAndBroadcast( + broadcastErrors: Bool = true, unlessStillAt broadcasted: LoopGraph? = nil + ) async -> [String] { releaseHeldCompletions() let errors = await drainPendingErrors(broadcastErrors: broadcastErrors) await drainPendingMessages() @@ -4513,7 +4651,7 @@ public actor GraphStore { await drainPendingHandoffDeliveries() await drainPendingNudges() await drainPendingFollowUps() - if errors.isEmpty { + if errors.isEmpty, broadcasted == nil || graph != broadcasted { await broadcast() } return errors diff --git a/GraphcodeKit/Sources/ProjectRegistry.swift b/GraphcodeKit/Sources/ProjectRegistry.swift index 46bbf22c..f50c282b 100644 --- a/GraphcodeKit/Sources/ProjectRegistry.swift +++ b/GraphcodeKit/Sources/ProjectRegistry.swift @@ -436,10 +436,14 @@ public actor ProjectRegistry { /// Sequential rather than concurrent across projects: each store's poll already spawns /// one subprocess per loop, and firing every project's at once would turn a quiet /// background tick into a burst of them. + /// One pass for the shared session listing too, so the whole tick costs one `zmx ls` + /// however many loops it reads (`ZmxSessionLauncher.SessionListing`). private func pollPresence() async { + await ZmxSessionLauncher.SessionListing.shared.beginPass() for store in stores.values { await store.pollPresence() } + await ZmxSessionLauncher.SessionListing.shared.endPass() } // MARK: - Commands diff --git a/GraphcodeKit/Sources/Sessions/CLISessionBackend.swift b/GraphcodeKit/Sources/Sessions/CLISessionBackend.swift index f51b4d35..d8140d7a 100644 --- a/GraphcodeKit/Sources/Sessions/CLISessionBackend.swift +++ b/GraphcodeKit/Sources/Sessions/CLISessionBackend.swift @@ -304,7 +304,7 @@ extension CLISessionBackend { public static let attachedClients: @Sendable (LoopNode, String?) async -> Int? = { node, path in - ZmxSessionLauncher.attachedClients(node, projectPath: path) + await ZmxSessionLauncher.attachedClients(node, projectPath: path) } /// Returns whether an earlier conversation was resumed. @@ -358,7 +358,7 @@ extension CLISessionBackend { /// The liveness hook `GraphStore` is wired with — session-level like `terminate`, /// so it needs no per-backend adapter. public static let sessionAlive: @Sendable (LoopNode, String?) async -> Bool = { node, path in - ZmxSessionLauncher.isSessionAlive(node, projectPath: path) + await ZmxSessionLauncher.isSessionAlive(node, projectPath: path) } public static let readPresence: @Sendable (LoopNode, String?) async -> PresenceReading = { diff --git a/GraphcodeKit/Sources/Sessions/RemoteGraphAccess.swift b/GraphcodeKit/Sources/Sessions/RemoteGraphAccess.swift index b538a79f..58fe0dad 100644 --- a/GraphcodeKit/Sources/Sessions/RemoteGraphAccess.swift +++ b/GraphcodeKit/Sources/Sessions/RemoteGraphAccess.swift @@ -1292,7 +1292,7 @@ public enum RemoteGraphAccess { payload["followUp"] = True run_with_verdict(project, {"messageNode": payload}, "accepted — typed in when the loop next goes idle" - if follow_up else "delivered") + if follow_up else "accepted — typing it in now") else: run_with_verdict(project, {"memoNode": payload}, "noted") diff --git a/GraphcodeKit/Sources/Sessions/ZmxSessionLauncher.swift b/GraphcodeKit/Sources/Sessions/ZmxSessionLauncher.swift index fadf019c..e2270bc3 100644 --- a/GraphcodeKit/Sources/Sessions/ZmxSessionLauncher.swift +++ b/GraphcodeKit/Sources/Sessions/ZmxSessionLauncher.swift @@ -569,7 +569,7 @@ public enum ZmxSessionLauncher { return await sendRemote(text, to: node, at: remote) } guard ZmxLocator.isInstalled, !text.isEmpty else { return false } - guard await sessionExists(node) else { return false } + guard await sendGate(node) else { return false } // Typed as the writes `sendWrites` frames it into — plain keystrokes when it is // short, a bracketed paste when it is not, which is what keeps a long message's head // from being swallowed by the composer (`maxUnbracketedSendBytes`, issue #277). The @@ -599,6 +599,18 @@ public enum ZmxSessionLauncher { return delivered } + /// Issue #215's gate: whether a keystroke typed now reaches the node's agent rather than + /// the shell a finished task left behind. Always a listing taken for this send, never + /// the shared one: a task that ended after the shared listing was taken is a husk that + /// listing still calls alive. Typing runs after the request is acknowledged, so the + /// listing costs the sender nothing; it still joins no listing started before it. + static func sendGate(_ node: LoopNode, listing: SessionListing = .shared) async -> Bool { + let result = await listing.listing(fresh: true) + return sessionTaskState( + lsStatus: result?.status, lsOutput: result?.output ?? "", + sessionName: SurfaceRef(id: node.id, launchesClaudeCode: true).zmxSessionName) == .alive + } + /// Reads a session's presence, preferring what the backend reported over what we can /// infer. A live session with no label is reported as idle at `.heuristic` confidence /// rather than guessed at — docs/04-cli-backends.md asks for the fallback to be @@ -727,23 +739,28 @@ public enum ZmxSessionLauncher { /// swallows anything typed at it (issue #215's `node send` that reported "delivered" /// into a session whose `claude` had exited). Only a running task is a session a /// keystroke can reach. - static func sessionExists(_ node: LoopNode, projectPath: String? = nil) async -> Bool { + /// + /// `fresh` for a lifecycle decision — a start, a kill, a launch that may have died — + /// where the shared listing could predate the very change being checked. + static func sessionExists( + _ node: LoopNode, projectPath: String? = nil, fresh: Bool = false + ) async -> Bool { if let projectPath, let remote = RemoteProjectLocation.parse(projectPath: projectPath) { return await runRemoteRetrying( remoteStatusInvocation(forNode: node, label: "presence", at: remote)) } - return await sessionTaskState(node) == .alive + return await sessionTaskState(node, fresh: fresh) == .alive } /// What is actually inside the node's zmx session: a running task (`alive`), a /// completed one (`exited`, with the exit code when zmx caught it), or no session at - /// all (`absent`). One `zmx ls`, the same cost as the `zmx get` existence check it - /// replaces, and the one answer both the send gate and the create-only ensure need — - /// an ensure keyed on `zmx get` could never revive a husk, because the husk *is* the - /// session that check asks about. - static func sessionTaskState(_ node: LoopNode) async -> SessionTaskState { + /// all (`absent`). Read from the shared listing (`SessionListing`), not a `zmx ls` of + /// its own: every `ls` probes every session on the machine, so one per node made a + /// presence pass O(N²) probes, and on a CPU-starved machine a send's gate alone took + /// longer than the CLI waits for its acknowledgement. + static func sessionTaskState(_ node: LoopNode, fresh: Bool = false) async -> SessionTaskState { guard ZmxLocator.isInstalled else { return .absent } - let result = runZmx(["ls"]) + let result = await SessionListing.shared.listing(fresh: fresh) return sessionTaskState( lsStatus: result?.status, lsOutput: result?.output ?? "", sessionName: SurfaceRef(id: node.id, launchesClaudeCode: true).zmxSessionName) @@ -790,7 +807,9 @@ public enum ZmxSessionLauncher { // reporting a session it cannot reach — the daemon behind it is gone, which is as // absent as a missing row, and counts as neither alive nor exited. guard line.range(of: "\terr=") == nil else { return .absent } - guard line.range(of: "\tended=") != nil else { return .alive } + guard line.range(of: "\tended=") != nil || line.range(of: "\texit_code=") != nil else { + return .alive + } var exitCode: Int? if let range = line.range(of: "\texit_code=") { let digits = line[range.upperBound...].prefix { $0.isNumber } @@ -870,7 +889,7 @@ public enum ZmxSessionLauncher { guard ZmxLocator.isInstalled else { return .failure(.unavailable("zmx is not installed")) } - if await sessionExists(node, projectPath: projectPath) { + if await sessionExists(node, projectPath: projectPath, fresh: true) { return .success(.attached) } var spawnedProcess: Process? @@ -904,7 +923,7 @@ public enum ZmxSessionLauncher { } for delay in [100, 200, 400, 800, 1200] { try? await Task.sleep(for: .milliseconds(delay)) - if await sessionExists(node, projectPath: projectPath) { + if await sessionExists(node, projectPath: projectPath, fresh: true) { spawnedProcess = nil return .success(.started) } @@ -915,12 +934,12 @@ public enum ZmxSessionLauncher { } await kill(node, projectPath: projectPath) for delay in [100, 200, 400] { - if await sessionExists(node, projectPath: projectPath) { + if await sessionExists(node, projectPath: projectPath, fresh: true) { await kill(node, projectPath: projectPath) } try? await Task.sleep(for: .milliseconds(delay)) } - if await sessionExists(node, projectPath: projectPath) { + if await sessionExists(node, projectPath: projectPath, fresh: true) { return .failure(.failed("zmx session appeared after startup timeout")) } return .failure( @@ -932,12 +951,12 @@ public enum ZmxSessionLauncher { public static func terminateResult( _ node: LoopNode, projectPath: String? = nil ) async -> Result { - guard await sessionExists(node, projectPath: projectPath) else { + guard await sessionExists(node, projectPath: projectPath, fresh: true) else { SessionIDStore.remove(forNodeID: node.id) return .success(()) } await kill(node, projectPath: projectPath) - guard !(await sessionExists(node, projectPath: projectPath)) else { + guard !(await sessionExists(node, projectPath: projectPath, fresh: true)) else { return .failure(.failed("zmx session remained after terminate")) } return .success(()) @@ -1038,11 +1057,13 @@ public enum ZmxSessionLauncher { /// How many terminals are attached to a node's session, or `nil` when that cannot be /// told — a remote project, or `zmx` not answering. - static func attachedClients(_ node: LoopNode, projectPath: String? = nil) -> Int? { + static func attachedClients(_ node: LoopNode, projectPath: String? = nil) async -> Int? { if let projectPath, RemoteProjectLocation.parse(projectPath: projectPath) != nil { return nil } - guard ZmxLocator.isInstalled, let result = runZmx(["ls"]), result.status == 0 else { + guard ZmxLocator.isInstalled, let result = await SessionListing.shared.listing(), + result.status == 0 + else { return nil } let name = SurfaceRef(id: node.id, launchesClaudeCode: true).zmxSessionName @@ -1080,13 +1101,14 @@ public enum ZmxSessionLauncher { /// structural, which is what the condemned list's later reap is for. static func killConfirmingDeath(sessionNamed name: String) async -> Bool { for attempt in 1...killAttempts { - guard runZmx(["kill", name]) != nil else { return false } - if sessionNamedState(name) == .absent { return true } + guard await runZmx(["kill", name]) != nil else { return false } + await SessionListing.shared.noteChanged() + if await sessionNamedState(name) == .absent { return true } if attempt < killAttempts { try? await Task.sleep(for: .seconds(killRetrySeconds)) } } - return sessionNamedState(name) == .absent + return await sessionNamedState(name) == .absent } static let killAttempts = 3 @@ -1108,18 +1130,23 @@ public enum ZmxSessionLauncher { /// Whether the node's local session is alive and not a husk — the ensure's own /// create-or-run check (`aliveCheckCommand`), asked on its own. A remote session /// answers `false`: its liveness is read through presence over ssh (`GraphStore`). - static func isSessionAlive(_ node: LoopNode, projectPath: String? = nil) -> Bool { + static func isSessionAlive(_ node: LoopNode, projectPath: String? = nil) async -> Bool { if let projectPath, RemoteProjectLocation.parse(projectPath: projectPath) != nil { return false } - guard ZmxLocator.isInstalled, let result = runZmx(["ls"]), result.status == 0 else { - return false - } + guard ZmxLocator.isInstalled else { return false } let name = SurfaceRef(id: node.id, launchesClaudeCode: true).zmxSessionName - return result.output.split(separator: "\n").contains { line in - !line.contains("\tended=") && !line.contains("\terr=") - && line.split(whereSeparator: \.isWhitespace).contains("name=\(name)") + func alive(in result: ZmxResult?) -> Bool { + guard let result, result.status == 0 else { return false } + return result.output.split(separator: "\n").contains { line in + !line.contains("\tended=") && !line.contains("\texit_code=") && !line.contains("\terr=") + && line.split(whereSeparator: \.isWhitespace).contains("name=\(name)") + } } + // Fresh both ways: a "no" decides whether a pane closing resolves the loop, and a + // "yes" tells a sender their message to a finished loop will be typed in. The shared + // listing can predate the session starting, or its task ending. + return alive(in: await SessionListing.shared.listing(fresh: true)) } private enum SessionNamedState { @@ -1128,34 +1155,68 @@ public enum ZmxSessionLauncher { case unknown } - private static func sessionNamedState(_ name: String) -> SessionNamedState { - guard let result = runZmx(["ls"]), result.status == 0 else { return .unknown } + /// A listing of its own rather than the shared one: this is a kill's confirmation, and + /// only a listing taken after the kill can confirm it. + private static func sessionNamedState(_ name: String) async -> SessionNamedState { + guard let result = await runZmx(["ls"]), result.status == 0 else { return .unknown } let exists = result.output.split(separator: "\n").contains { line in line.split(whereSeparator: \.isWhitespace).contains("name=\(name)") } return exists ? .present : .absent } - private struct ZmxResult { + struct ZmxResult: Equatable, Sendable { let status: Int32 let output: String } /// Kill and confirmation must still work on a machine that has exhausted its PTYs, /// so these one-shot commands use pipes rather than `PTYProcessSession`. - private static func runZmx(_ arguments: [String]) -> ZmxResult? { + static func runZmx(_ arguments: [String]) async -> ZmxResult? { + await runCollectingOutput(ZmxLocator.binaryURL, arguments) + } + + /// Runs a process to its exit and returns its status and stdout, without parking a + /// cooperative-pool thread on either. `readDataToEndOfFile` and `waitUntilExit` both + /// block their thread, and with every presence read and send doing it at once on a + /// starved machine they held the daemon's pool for as long as `zmx` took to answer. + /// The read runs on a GCD thread and the exit arrives through `terminationHandler`, + /// set before `run` so it cannot be missed; the two are joined, not sequenced, + /// because reading only after the exit deadlocks a child that filled the pipe. + static func runCollectingOutput(_ executable: URL, _ arguments: [String]) async -> ZmxResult? { let process = Process() - process.executableURL = ZmxLocator.binaryURL + process.executableURL = executable process.arguments = arguments let output = Pipe() process.standardOutput = output process.standardError = FileHandle.nullDevice - do { try process.run() } catch { return nil } - let data = output.fileHandleForReading.readDataToEndOfFile() - process.waitUntilExit() - return ZmxResult( - status: process.terminationStatus, - output: String(data: data, encoding: .utf8) ?? "") + process.standardInput = FileHandle.nullDevice + let group = DispatchGroup() + let collected = CollectedOutput() + group.enter() + process.terminationHandler = { _ in group.leave() } + do { try process.run() } catch { + process.terminationHandler = nil + return nil + } + group.enter() + DispatchQueue.global(qos: .utility).async { + collected.data = output.fileHandleForReading.readDataToEndOfFile() + group.leave() + } + return await withCheckedContinuation { continuation in + group.notify(queue: .global(qos: .utility)) { + continuation.resume( + returning: ZmxResult( + status: process.terminationStatus, + output: String(data: collected.data, encoding: .utf8) ?? "")) + } + } + } + + /// The reader thread's result; the group orders the write before the read. + private final class CollectedOutput: @unchecked Sendable { + var data = Data() } /// `zmx get ` exits 0 when the session exists and 1 when it doesn't — raw @@ -2577,7 +2638,7 @@ public enum ZmxSessionLauncher { DialLog.record(session: name, dial: "first-pass", event: "already-served") return } - guard await sessionExists(node) else { + guard await sessionExists(node, fresh: true) else { DialLog.record(session: name, dial: "first-pass", event: "no-session") return } @@ -2692,7 +2753,7 @@ public enum ZmxSessionLauncher { /// count as the death they are. private static func sessionDiedImmediately(node: LoopNode) async -> Bool { try? await Task.sleep(for: .seconds(resumeSettleSeconds)) - return await sessionTaskState(node) != .alive + return await sessionTaskState(node, fresh: true) != .alive } /// `logFragment` rides inside the run branch, so an ensure whose check found the @@ -2724,6 +2785,7 @@ public enum ZmxSessionLauncher { try process.run() await Task.detached { process.waitUntilExit() }.value } catch {} + await SessionListing.shared.noteChanged() return #else // The stamp rides in the run branch, after the launch it describes and only if @@ -2744,7 +2806,93 @@ public enum ZmxSessionLauncher { workingDirectory: workingDirectory) else { return } _ = await session.waitUntilFinished() + await SessionListing.shared.noteChanged() #endif } } + +extension ZmxSessionLauncher { + /// One `zmx ls`, shared by everyone asking about local sessions at about the same time. + /// + /// A listing probes every session on the machine, one at a time — and the workspaces + /// share one zmx directory, so that is every workspace's sessions. Asked once per node, + /// a presence pass cost N listings of N probes each, and a send took a listing of its + /// own before typing a key; on a starved machine either took longer than the CLI waits. + /// + /// So a reader is handed a listing that is fresh enough rather than a new one: + /// - during a presence pass (`beginPass`/`endPass`), any listing started since the + /// pass began — one per pass, and a send arriving mid-pass shares it; + /// - otherwise, one started within `reuseWindow`; + /// - and a listing still being taken is joined rather than duplicated. + /// None started before the last session this process started or killed + /// (`noteChanged`) is reused, so a launch's own check never reads a listing older than + /// the launch. `passReuseLimit` bounds how stale a long pass's listing can get. + actor SessionListing { + static let shared = SessionListing() + + private let take: @Sendable () async -> ZmxResult? + private let reuseWindow: Duration + private let passReuseLimit: Duration + private let clock = ContinuousClock() + private var latest: (startedAt: ContinuousClock.Instant, result: ZmxResult?)? + private var inFlight: (startedAt: ContinuousClock.Instant, task: Task)? + private var changedAt: ContinuousClock.Instant? + private var passStartedAt: ContinuousClock.Instant? + private var openPasses = 0 + private(set) var listingsTaken = 0 + + init( + reuseWindow: Duration = .seconds(2), + passReuseLimit: Duration = .seconds(30), + take: @escaping @Sendable () async -> ZmxResult? = { await runZmx(["ls"]) } + ) { + self.reuseWindow = reuseWindow + self.passReuseLimit = passReuseLimit + self.take = take + } + + /// `fresh` asks for a listing started from now on — joined by later readers, never + /// satisfied by an earlier one. + func listing(fresh: Bool = false) async -> ZmxResult? { + let now = clock.now + // Nothing started before these may answer at all; a listing still being taken + // finishes after the reader asked, so only these bound joining it. + let floor = [changedAt, passStartedAt].compactMap { $0 }.max() + let joinBound = fresh ? now : floor + var reuseBound = passStartedAt.map { max($0, now - passReuseLimit) } ?? now - reuseWindow + if let floor { reuseBound = max(reuseBound, floor) } + if !fresh, let latest, latest.startedAt >= reuseBound { return latest.result } + // A listing that has run longer than a pass may reuse one is hung, not slow; it + // is left to its own caller rather than stalling everyone who joins it. + if let inFlight, + inFlight.startedAt >= max(joinBound ?? now - passReuseLimit, now - passReuseLimit) + { + return await inFlight.task.value + } + let startedAt = clock.now + let take = self.take + let task = Task { await take() } + inFlight = (startedAt, task) + listingsTaken += 1 + let result = await task.value + if inFlight?.startedAt == startedAt { inFlight = nil } + if latest.map({ $0.startedAt <= startedAt }) ?? true { latest = (startedAt, result) } + return result + } + + func noteChanged() { + changedAt = clock.now + } + + func beginPass() { + if openPasses == 0 { passStartedAt = clock.now } + openPasses += 1 + } + + func endPass() { + openPasses = max(0, openPasses - 1) + if openPasses == 0 { passStartedAt = nil } + } + } +} diff --git a/graphcode-cli/Sources/main.swift b/graphcode-cli/Sources/main.swift index f5a470ae..638e4d95 100644 --- a/graphcode-cli/Sources/main.swift +++ b/graphcode-cli/Sources/main.swift @@ -288,8 +288,10 @@ do { .graphCommand( projectPath: projectPath, command: .messageNode(nodeID, text: text, from: sender, followUp: followUp))) - // Delivery is judged and attempted before the daemon broadcasts, so the first event - // back is the verdict: an error means it did not land, the graph means it did. + // Deliverability is judged before the daemon broadcasts, so the first event back is + // the verdict: an error means it was staged to the loop's memory instead, the graph + // means it is on the board and queued for the session. The typing itself happens + // after the acknowledgement, and a failure there is staged to memory too. let verdict = try client.waitForEvent { event in switch event { case .graphChanged, .errorOccurred: return true @@ -299,7 +301,9 @@ do { if case .errorOccurred(let message) = verdict { fail(message) } // A follow-up to a busy loop is accepted, not typed: it's in the loop's memory now // and its session hears it when it next goes idle — "delivered" would overclaim. - print(followUp ? "accepted — typed in when the loop next goes idle" : "delivered") + print( + followUp + ? "accepted — typed in when the loop next goes idle" : "accepted — typing it in now") case .updateNode(let projectPath, let nodeID, let update): // Attributed like `node create`/`node send`: run from inside a loop, ZMX_SESSION diff --git a/graphcode/Tests/GoalResolutionFollowUpTests.swift b/graphcode/Tests/GoalResolutionFollowUpTests.swift index fdb2db47..afd88805 100644 --- a/graphcode/Tests/GoalResolutionFollowUpTests.swift +++ b/graphcode/Tests/GoalResolutionFollowUpTests.swift @@ -35,6 +35,7 @@ struct GoalResolutionFollowUpTests { let resolution = await store.graph.nodes[id: id]?.resolution await store.handle(.messageNode(id, text: "what did you change?", from: nil, followUp: false)) + await store.finishSessionTyping() #expect(delivered.value.contains { $0.contains("what did you change?") }) #expect(!memory.value.contains { $0.hasPrefix("while you were away") }) @@ -49,6 +50,7 @@ struct GoalResolutionFollowUpTests { let (store, id) = await finishedTarget(alive: false, delivered: delivered, memory: memory) await store.handle(.messageNode(id, text: "what did you change?", from: nil, followUp: false)) + await store.finishSessionTyping() #expect(!delivered.value.contains { $0.contains("what did you change?") }) #expect(memory.value.contains { $0.hasPrefix("while you were away") }) diff --git a/graphcode/Tests/MessageAndSpawnTests.swift b/graphcode/Tests/MessageAndSpawnTests.swift index 55db0ff2..cc124d61 100644 --- a/graphcode/Tests/MessageAndSpawnTests.swift +++ b/graphcode/Tests/MessageAndSpawnTests.swift @@ -248,6 +248,7 @@ struct MessageAndSpawnTests { await store.handle( .messageNode(target.id, text: "the API changed", from: sender.id, followUp: nil)) + await store.finishSessionTyping() #expect(delivered.value.count == 1) #expect(delivered.value[0].nodeTitle == "Review") @@ -261,6 +262,7 @@ struct MessageAndSpawnTests { let target = await store.graph.nodes[1] await store.handle(.messageNode(target.id, text: "wrap it up", from: nil, followUp: nil)) + await store.finishSessionTyping() #expect(delivered.value == [Delivery(nodeTitle: "Review", text: "[graphcode] wrap it up")]) } diff --git a/graphcode/Tests/MessageDeliveryTests.swift b/graphcode/Tests/MessageDeliveryTests.swift index 82b23873..cc0fa9a9 100644 --- a/graphcode/Tests/MessageDeliveryTests.swift +++ b/graphcode/Tests/MessageDeliveryTests.swift @@ -101,6 +101,30 @@ struct MessageDeliveryTests { return String(("[graphcode] \(tag): " + filler).prefix(bytes)) } + /// Issue #215 against a real husk: `zmx run … -d true` leaves the wrapper shell at its + /// prompt with the task ended, which is exactly the session a stale listing would call + /// alive. Nothing may be typed into it. + @Test + func aSessionWhoseTaskEndedIsNeverTypedInto() async throws { + let node = node() + let name = SurfaceRef(id: node.id, launchesClaudeCode: true).zmxSessionName + defer { Task { await kill(name) } } + await run("\(Self.quoted(Self.zmx)) run \(Self.quoted(name)) -d true >/dev/null 2>&1") + var husk = false + for _ in 0..<200 { + let listing = await run("\(Self.quoted(Self.zmx)) ls") + let row = listing.output.split(separator: "\n").first { $0.contains("name=\(name)") } + if row?.contains("ended=") == true { + husk = true + break + } + try? await Task.sleep(for: .milliseconds(100)) + } + try #require(husk) + + #expect(!(await ZmxSessionLauncher.send("[graphcode] are you there?", to: node))) + } + @Test func anOversizedMessageArrivesByteForByte() async throws { let node = node() diff --git a/graphcode/Tests/RemoteCLIShimTests.swift b/graphcode/Tests/RemoteCLIShimTests.swift index 7b627fe8..0c66904a 100644 --- a/graphcode/Tests/RemoteCLIShimTests.swift +++ b/graphcode/Tests/RemoteCLIShimTests.swift @@ -70,7 +70,7 @@ struct RemoteCLIShimTests { ["node", "send", Self.project, nodeID.uuidString, "the", "API", "changed"]) #expect(run.status == 0) - #expect(run.stdout.contains("delivered")) + #expect(run.stdout.contains("accepted — typing it in now")) #expect( run.commands.dropFirst().first == .graphCommand( diff --git a/graphcode/Tests/RespawnOnSendTests.swift b/graphcode/Tests/RespawnOnSendTests.swift index 0c30aad4..6cf38f4a 100644 --- a/graphcode/Tests/RespawnOnSendTests.swift +++ b/graphcode/Tests/RespawnOnSendTests.swift @@ -39,6 +39,7 @@ struct RespawnOnSendTests { let nodeID = graph.nodes[0].id await store.handle(.messageNode(nodeID, text: "wake up", from: nil, followUp: nil)) + await store.finishSessionTyping() #expect(deliveries.value == 2) #expect(ensured.value == 1) @@ -61,6 +62,7 @@ struct RespawnOnSendTests { let nodeID = graph.nodes[0].id await store.handle(.messageNode(nodeID, text: "wake up", from: nil, followUp: nil)) + await store.finishSessionTyping() #expect(deliveries.value == 2) #expect(remembered.value.contains { $0.contains("while you were away") }) @@ -85,6 +87,7 @@ struct RespawnOnSendTests { let nodeID = graph.nodes[0].id await store.handle(.messageNode(nodeID, text: "wake up", from: nil, followUp: nil)) + await store.finishSessionTyping() #expect(deliveries.value == 1) #expect(ensured.value == 0) @@ -102,6 +105,7 @@ struct RespawnOnSendTests { let nodeID = graph.nodes[0].id await store.handle(.messageNode(nodeID, text: "wake up", from: nil, followUp: nil)) + await store.finishSessionTyping() #expect(ensured.value == 0) } diff --git a/graphcode/Tests/SendAcknowledgementTests.swift b/graphcode/Tests/SendAcknowledgementTests.swift new file mode 100644 index 00000000..eab6c1bd --- /dev/null +++ b/graphcode/Tests/SendAcknowledgementTests.swift @@ -0,0 +1,395 @@ +import ComposableArchitecture +import Foundation +import Testing + +@testable import GraphcodeKit + +/// `node send` exiting 75 for a message that landed. The daemon used to type the whole +/// message — a `zmx ls` gate over every session, each chunk, the submit beat, Enter — +/// before it acknowledged the request, and on a starved machine that outlasted the CLI's +/// ten-second wait. The acknowledgement now comes once the message is on the board and +/// queued for its session; the typing follows. +@Suite +struct SendAcknowledgementTests { + private actor Gate { + private var opened = false + private var waiters: [CheckedContinuation] = [] + + func wait() async { + guard !opened else { return } + await withCheckedContinuation { waiters.append($0) } + } + + func open() { + opened = true + for waiter in waiters { waiter.resume() } + waiters = [] + } + } + + private static func graph() -> LoopGraph { + LoopGraph( + project: ProjectRef(path: "/tmp/send-ack", name: "send-ack"), + nodes: [ + LoopNode( + title: "Worker", loopType: .goalBased, goal: GoalSpec(summary: "work"), + state: .running) + ]) + } + + @Test + func aSendIsAcknowledgedBeforeItIsTyped() async { + let events = LockIsolated<[String]>([]) + let graph = Self.graph() + let store = GraphStore( + graph: graph, + onDeliverMessage: { _, text, _ in + try? await Task.sleep(for: .milliseconds(200)) + events.withValue { $0.append("typed \(text)") } + return true + }, + onAppendMemory: { _, _ in }) + + let result = await store.handle( + .messageNode(graph.nodes[0].id, text: "the API changed", from: nil, followUp: nil)) + events.withValue { $0.append("acknowledged") } + await store.finishSessionTyping() + + guard case .applied = result else { + Issue.record("expected the send to be acknowledged, got \(result)") + return + } + #expect(events.value == ["acknowledged", "typed [graphcode] the API changed"]) + } + + // Bounded: typed before the acknowledgement, the gate below is never opened. + @Test(.timeLimit(.minutes(1))) + func sendsToOneLoopAreTypedInTheOrderTheyWereSent() async { + let gate = Gate() + let typed = LockIsolated<[String]>([]) + let graph = Self.graph() + let store = GraphStore( + graph: graph, + onDeliverMessage: { _, text, _ in + if text.hasSuffix("first") { await gate.wait() } + typed.withValue { $0.append(text) } + return true + }, + onAppendMemory: { _, _ in }) + let nodeID = graph.nodes[0].id + + await store.handle(.messageNode(nodeID, text: "first", from: nil, followUp: nil)) + await store.handle(.messageNode(nodeID, text: "second", from: nil, followUp: nil)) + await gate.open() + await store.finishSessionTyping() + + #expect(typed.value == ["[graphcode] first", "[graphcode] second"]) + } + + // Bounded: typed before the acknowledgement, the gate below is never opened. + @Test(.timeLimit(.minutes(1))) + func aFollowUpWaitsBehindASendStillBeingTyped() async { + let gate = Gate() + let typed = LockIsolated<[String]>([]) + let graph = Self.graph() + let store = GraphStore( + graph: graph, + onDeliverMessage: { _, text, _ in + if text.hasSuffix("now") { await gate.wait() } + typed.withValue { $0.append(text) } + return true + }, + onReadPresence: { _, _ in PresenceReading(presence: .idle, confidence: .reported) }, + onAppendMemory: { _, _ in }) + let nodeID = graph.nodes[0].id + + await store.handle(.messageNode(nodeID, text: "now", from: nil, followUp: nil)) + await store.handle(.messageNode(nodeID, text: "later", from: nil, followUp: true)) + #expect(typed.value.isEmpty) + await gate.open() + await store.finishSessionTyping() + + #expect(typed.value == ["[graphcode] now", "[graphcode] later"]) + } + + @Test(.timeLimit(.minutes(1))) + func aTypingThatHangsIsStagedAtTheDeadlineAndFreesTheLoop() async { + let never = Gate() + let typed = LockIsolated<[String]>([]) + let remembered = LockIsolated<[String]>([]) + let graph = Self.graph() + let store = GraphStore( + graph: graph, + deliveryDeadline: .milliseconds(300), + onDeliverMessage: { _, text, _ in + if text.hasSuffix("now") { await never.wait() } + typed.withValue { $0.append(text) } + return true + }, + onReadPresence: { _, _ in PresenceReading(presence: .idle, confidence: .reported) }, + onAppendMemory: { _, entry in remembered.withValue { $0.append(entry) } }) + let nodeID = graph.nodes[0].id + + await store.handle(.messageNode(nodeID, text: "now", from: nil, followUp: nil)) + await store.handle(.messageNode(nodeID, text: "later", from: nil, followUp: true)) + await store.finishSessionTyping() + + #expect(remembered.value.contains("while you were away: [graphcode] now")) + #expect(typed.value == ["[graphcode] later"]) + } + + @Test + func aTypingAbandonedAtTheDeadlineIsNeitherRetriedNorStagedTwice() async { + let attempts = LockIsolated(0) + let ensured = LockIsolated(0) + let remembered = LockIsolated<[String]>([]) + let graph = Self.graph() + let store = GraphStore( + graph: graph, + deliveryDeadline: .milliseconds(300), + onEnsureSession: { _, _ in ensured.withValue { $0 += 1 } }, + onDeliverMessage: { _, _, _ in + attempts.withValue { $0 += 1 } + // Hangs until the deadline cancels it, then reports failure, as a cancelled + // `zmx send` does. + while !Task.isCancelled { try? await Task.sleep(for: .milliseconds(10)) } + return false + }, + onAppendMemory: { _, entry in remembered.withValue { $0.append(entry) } }) + + await store.handle(.messageNode(graph.nodes[0].id, text: "now", from: nil, followUp: nil)) + await store.finishSessionTyping() + // Past the respawn settle, where an abandoned typing would retry. + try? await Task.sleep(for: GraphStore.respawnedSessionSettle + .milliseconds(500)) + + #expect(attempts.value == 1) + #expect(ensured.value == 0) + #expect(remembered.value == ["while you were away: [graphcode] now"]) + } + + @Test + func aDeadlineDuringTheRespawnSettleDoesNotTypeAgain() async { + let attempts = LockIsolated(0) + let remembered = LockIsolated<[String]>([]) + let graph = Self.graph() + let store = GraphStore( + graph: graph, + deliveryDeadline: .milliseconds(300), + onEnsureSession: { _, _ in }, + onDeliverMessage: { _, _, _ in + // Fails at once, so the deadline lands in the respawn settle before the retry. + attempts.withValue { $0 += 1 } + return false + }, + onAppendMemory: { _, entry in remembered.withValue { $0.append(entry) } }) + + await store.handle(.messageNode(graph.nodes[0].id, text: "now", from: nil, followUp: nil)) + await store.finishSessionTyping() + try? await Task.sleep(for: GraphStore.respawnedSessionSettle + .milliseconds(500)) + + #expect(attempts.value == 1) + #expect(remembered.value == ["while you were away: [graphcode] now"]) + } + + @Test + func aSendIsBroadcastOnceWhenItsDrainChangesNothing() async { + let broadcasts = LockIsolated(0) + let graph = Self.graph() + let store = GraphStore( + graph: graph, + onGraphChanged: { _ in broadcasts.withValue { $0 += 1 } }, + onDeliverMessage: { _, _, _ in true }, + onAppendMemory: { _, _ in }) + + await store.handle(.messageNode(graph.nodes[0].id, text: "hi", from: nil, followUp: nil)) + await store.finishSessionTyping() + + #expect(broadcasts.value == 1) + } + + @Test + func aSendThatCannotBeTypedIsStagedWithoutAnErrorForOtherClients() async { + let remembered = LockIsolated<[String]>([]) + var graph = Self.graph() + graph.nodes[0].loopType = .turnBased + let store = GraphStore( + graph: graph, + onDeliverMessage: { _, _, _ in false }, + onAppendMemory: { _, entry in remembered.withValue { $0.append(entry) } }) + + let result = await store.handle( + .messageNode(graph.nodes[0].id, text: "wake up", from: nil, followUp: nil)) + await store.finishSessionTyping() + + guard case .applied = result else { + Issue.record("expected the send to be acknowledged, got \(result)") + return + } + #expect(remembered.value == ["while you were away: [graphcode] wake up"]) + // The failure is learned after the acknowledgement; an announced error would be the + // verdict of whichever command ran next, not this one. + let next = await store.handle(.memoNode(graph.nodes[0].id, text: "note", from: nil)) + guard case .applied = next else { + Issue.record("an unrelated command inherited the send's failure: \(next)") + return + } + } +} + +/// One `zmx ls` per presence pass, shared with sends, instead of one per node per read. +@Suite +struct SessionListingTests { + private static let listing = ZmxSessionLauncher.ZmxResult( + status: 0, output: "name=graphcode-A\tpid=1\tclients=0\n") + + @Test + func concurrentReadersShareOneListing() async { + let listing = ZmxSessionLauncher.SessionListing(reuseWindow: .zero) { + try? await Task.sleep(for: .milliseconds(100)) + return Self.listing + } + + let results = await withTaskGroup(of: ZmxSessionLauncher.ZmxResult?.self) { group in + for _ in 0..<20 { group.addTask { await listing.listing() } } + return await group.reduce(into: []) { $0.append($1) } + } + + #expect(results.allSatisfy { $0 == Self.listing }) + #expect(await listing.listingsTaken == 1) + } + + @Test + func aPassTakesOneListingHoweverManyReadsItMakes() async { + let listing = ZmxSessionLauncher.SessionListing(reuseWindow: .zero) { Self.listing } + + await listing.beginPass() + for _ in 0..<27 { _ = await listing.listing() } + await listing.endPass() + + #expect(await listing.listingsTaken == 1) + } + + @Test + func outsideAPassAListingOlderThanTheWindowIsRetaken() async { + let listing = ZmxSessionLauncher.SessionListing(reuseWindow: .zero) { Self.listing } + + _ = await listing.listing() + try? await Task.sleep(for: .milliseconds(5)) + _ = await listing.listing() + + #expect(await listing.listingsTaken == 2) + } + + @Test + func aSessionStartedOrKilledInvalidatesTheListing() async { + let listing = ZmxSessionLauncher.SessionListing(reuseWindow: .seconds(60)) { Self.listing } + + await listing.beginPass() + _ = await listing.listing() + await listing.noteChanged() + _ = await listing.listing() + _ = await listing.listing() + await listing.endPass() + + #expect(await listing.listingsTaken == 2) + } + + @Test + func aFreshListingIsNeverAnsweredByAnEarlierOne() async { + // A send that found its target missing confirms it against a listing of its own: the + // shared one may predate a session a human's pane has just started. + let listing = ZmxSessionLauncher.SessionListing(reuseWindow: .seconds(60)) { Self.listing } + + _ = await listing.listing() + _ = await listing.listing() + _ = await listing.listing(fresh: true) + + #expect(await listing.listingsTaken == 2) + } + + @Test + func aFailedListingIsSharedRatherThanRetriedPerNode() async { + let listing = ZmxSessionLauncher.SessionListing(reuseWindow: .zero) { nil } + + await listing.beginPass() + let first = await listing.listing() + let second = await listing.listing() + await listing.endPass() + + #expect(first == nil && second == nil) + #expect(await listing.listingsTaken == 1) + } +} + +@Suite +struct RunCollectingOutputTests { + @Test + func collectsMoreOutputThanAPipeHoldsAndTheExitStatus() async { + let result = await ZmxSessionLauncher.runCollectingOutput( + URL(fileURLWithPath: "/bin/sh"), + ["-c", "head -c 200000 /dev/zero | tr '\\000' x; exit 3"]) + + #expect(result?.status == 3) + #expect(result?.output.count == 200_000) + } + + @Test + func aMissingExecutableIsNil() async { + let result = await ZmxSessionLauncher.runCollectingOutput( + URL(fileURLWithPath: "/nonexistent/zmx"), ["ls"]) + + #expect(result == nil) + } +} + +/// Issue #215 through the shared listing: a husk — a session whose task has ended, its +/// wrapper shell still at the prompt — is never typed into. `zmx ls` marks one with +/// `ended=`/`exit_code=`, and the listing everyone else shares can be older than the end. +@Suite +struct HuskSendGateTests { + private static let node = LoopNode(title: "Worker") + private static var name: String { + SurfaceRef(id: node.id, launchesClaudeCode: true).zmxSessionName + } + private static var live: String { "name=\(name)\tpid=1\tclients=0\tcmd=claude\n" } + private static var husk: String { + "name=\(name)\tpid=1\tclients=0\tcmd=claude\tended=1790000000\texit_code=0\n" + } + + private static func listing(_ outputs: [String]) -> ZmxSessionLauncher.SessionListing { + let remaining = LockIsolated(outputs) + return ZmxSessionLauncher.SessionListing(reuseWindow: .seconds(60)) { + let output = remaining.withValue { $0.count > 1 ? $0.removeFirst() : $0[0] } + return ZmxSessionLauncher.ZmxResult(status: 0, output: output) + } + } + + @Test + func aLiveSessionPassesTheGate() async { + #expect(await ZmxSessionLauncher.sendGate(Self.node, listing: Self.listing([Self.live]))) + } + + @Test(arguments: [ + "\tended=1790000000\texit_code=0", "\tended=1790000000", "\texit_code=137", + ]) + func aHuskAtListingTimeIsNeverTypedInto(marker: String) async { + let row = "name=\(Self.name)\tpid=1\tclients=0\tcmd=claude\(marker)\n" + + #expect(!(await ZmxSessionLauncher.sendGate(Self.node, listing: Self.listing([row])))) + } + + @Test + func aSessionThatEndsAfterTheSharedListingIsNeverTypedInto() async { + let listing = Self.listing([Self.live, Self.husk]) + // A presence pass has the session alive in the shared listing; the task then ends. + await listing.beginPass() + let shared = await listing.listing() + #expect(shared?.output == Self.live) + + let typed = await ZmxSessionLauncher.sendGate(Self.node, listing: listing) + await listing.endPass() + + #expect(!typed) + #expect(await listing.listingsTaken == 2) + } +} diff --git a/graphcode/Tests/SubGraphAddressingTests.swift b/graphcode/Tests/SubGraphAddressingTests.swift index ac330f4a..7a24c611 100644 --- a/graphcode/Tests/SubGraphAddressingTests.swift +++ b/graphcode/Tests/SubGraphAddressingTests.swift @@ -74,19 +74,26 @@ struct SubGraphAddressingTests { @Test func aMessageAddressedToAChildLoopReachesItsTransport() async { - let delivered = LockIsolated<[UUID]>([]) + let delivered = LockIsolated<[String]>([]) let worker = LoopNode(title: "Worker", loopType: .turnBased, checkDescription: "?") let (store, _) = storeWithComposite( subNodes: [worker], - onDeliverMessage: { node, _, _ in - delivered.withValue { $0.append(node.id) } + onDeliverMessage: { node, text, _ in + // Slow on purpose: a child store lives for one command, so its sends must be + // typed before that command returns, and in the order they were sent. + if text.hasSuffix("inbox") { try? await Task.sleep(for: .milliseconds(50)) } + delivered.withValue { $0.append("\(node.id == worker.id) \(text)") } return true }) await store.handle( .messageNode(worker.id, text: "prioritize the inbox", from: nil, followUp: nil)) + await store.handle(.messageNode(worker.id, text: "then the backlog", from: nil, followUp: nil)) - #expect(delivered.value == [worker.id]) + #expect( + delivered.value == [ + "true [graphcode] prioritize the inbox", "true [graphcode] then the backlog", + ]) } @Test