From 03baaeae77ffa21aa25494d17042144158d66d00 Mon Sep 17 00:00:00 2001 From: scgopi Date: Thu, 1 Oct 2026 19:57:50 -0700 Subject: [PATCH 1/5] Acknowledge node send before typing it; share one zmx ls per pass node send exited 75 ("may still have been applied") for messages that landed: the daemon typed the whole message (a zmx ls gate over every session, each chunk, the submit beat, Enter) before acknowledging, and on a CPU-starved machine that outlasted the CLI's 10s wait. - messageNode is acknowledged once the message is on the board and queued for its session; typing runs afterwards on a per-target chain that keeps send order, and follow-ups wait behind it. A failure after the ack is staged to the loop's memory and logged, not broadcast, so no other client takes it for its own verdict. - One zmx ls per presence pass (ZmxSessionLauncher.SessionListing), shared with sends and joined while in flight, instead of a full listing per node per read. Starting or killing a session invalidates it, and a "not alive" that would refuse a send or allow a resolution is confirmed against a fresh listing. - runZmx awaits exit through terminationHandler and reads on a GCD thread instead of blocking the cooperative pool. Co-Authored-By: Claude Opus 5.5 Signed-off-by: scgopi --- GraphcodeKit/Sources/GraphStore.swift | 128 +++++++-- GraphcodeKit/Sources/ProjectRegistry.swift | 4 + .../Sources/Sessions/CLISessionBackend.swift | 4 +- .../Sources/Sessions/RemoteGraphAccess.swift | 2 +- .../Sources/Sessions/ZmxSessionLauncher.swift | 191 +++++++++++--- graphcode-cli/Sources/main.swift | 10 +- .../Tests/GoalResolutionFollowUpTests.swift | 2 + graphcode/Tests/MessageAndSpawnTests.swift | 2 + graphcode/Tests/RemoteCLIShimTests.swift | 2 +- graphcode/Tests/RespawnOnSendTests.swift | 4 + .../Tests/SendAcknowledgementTests.swift | 248 ++++++++++++++++++ 11 files changed, 533 insertions(+), 64 deletions(-) create mode 100644 graphcode/Tests/SendAcknowledgementTests.swift diff --git a/GraphcodeKit/Sources/GraphStore.swift b/GraphcodeKit/Sources/GraphStore.swift index 3f09cd06..0ee01f03 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,15 @@ 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] = [] + if case .messageNode = command { + errors = await drainPendingErrors(broadcastErrors: broadcastErrors) + if errors.isEmpty { await broadcast() } + } + errors += await drainAndBroadcast(broadcastErrors: broadcastErrors) if let error = errors.first { return .rejected(message: error, graph: graph) } @@ -3679,10 +3690,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 @@ -3732,11 +3743,17 @@ public actor GraphStore { // 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). + // 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). Liveness is asked before the acknowledgement + // so a sender reporting to a finished parent still hears "staged" — one shared + // listing (`SessionListing`), 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 { + typeAfterAcknowledging(message, to: target, resolved: true) return } if MessageBus.deliverability(to: target) != nil { @@ -3746,29 +3763,79 @@ 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 } + 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. + private func typeAfterAcknowledging(_ message: String, to target: LoopNode, resolved: Bool) { + let previous = sessionTyping[target.id]?.task + let token = UUID() + let task = Task { [self] in + await previous?.value + await typeAdHocMessage(message, to: target.id, resolved: resolved) + if sessionTyping[target.id]?.token == token { + sessionTyping.removeValue(forKey: target.id) } - 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") + if pendingFollowUps.contains(where: { $0.nodeID == target.id }) { + await drainPendingFollowUps() + } + } + sessionTyping[target.id] = (token, task) + } + + /// 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 { + guard let target = graph.nodes[id: nodeID] else { return } + if resolved { + if await deliverToSession(target, message) { return } + stageUntyped(message, to: target, reason: "session-gone") return } + if await deliverToSession(target, message) { return } + // 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 } + } + stageUntyped(message, to: target, reason: "delivery-failed") + } + + private func stageUntyped(_ message: String, to target: LoopNode, reason: String) { + recordMemory(target.id, "while you were away: \(message)") + DaemonLog.shared.record( + "send-staged", [("node", target.id.uuidString), ("reason", reason)]) + 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 +4078,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 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..d70cb731 100644 --- a/GraphcodeKit/Sources/Sessions/ZmxSessionLauncher.swift +++ b/GraphcodeKit/Sources/Sessions/ZmxSessionLauncher.swift @@ -569,7 +569,11 @@ 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 } + // A shared listing can predate a session started outside this process — a pane a + // human just opened — so a refusal is confirmed against a listing of its own. + if await sessionTaskState(node) != .alive { + guard await sessionTaskState(node, freshListing: true) == .alive 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 @@ -737,13 +741,15 @@ public enum ZmxSessionLauncher { /// 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, freshListing: Bool = false) async + -> SessionTaskState + { guard ZmxLocator.isInstalled else { return .absent } - let result = runZmx(["ls"]) + let result = await SessionListing.shared.listing(fresh: freshListing) return sessionTaskState( lsStatus: result?.status, lsOutput: result?.output ?? "", sessionName: SurfaceRef(id: node.id, launchesClaudeCode: true).zmxSessionName) @@ -1038,11 +1044,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 +1088,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 +1117,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("\terr=") + && line.split(whereSeparator: \.isWhitespace).contains("name=\(name)") + } } + // A "no" decides whether a pane closing resolves the loop, and a shared listing can + // predate the session — so it is confirmed against a listing of its own. + if alive(in: await SessionListing.shared.listing()) { return true } + return alive(in: await SessionListing.shared.listing(fresh: true)) } private enum SessionNamedState { @@ -1128,34 +1142,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 @@ -2724,6 +2772,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 +2793,89 @@ 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 } + if let inFlight, joinBound.map({ inFlight.startedAt >= $0 }) ?? true { + 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/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..199c283c --- /dev/null +++ b/graphcode/Tests/SendAcknowledgementTests.swift @@ -0,0 +1,248 @@ +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 + 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) + } +} From 0da20eaa89af7452f8265b45973e63b6c09b2509 Mon Sep 17 00:00:00 2001 From: scgopi Date: Thu, 1 Oct 2026 20:10:45 -0700 Subject: [PATCH 2/5] Gate every local send on a listing of its own (#215) The shared listing can be older than a task's end, and a session whose task ended is a husk whose wrapper shell would take the keystrokes. The send gate now always takes a fresh zmx ls (typing runs after the ack, so this costs the sender nothing) and types only into a task zmx reports alive. A row with exit_code= but no ended= counts as ended too. Co-Authored-By: Claude Opus 5.5 Signed-off-by: scgopi --- .../Sources/Sessions/ZmxSessionLauncher.swift | 30 +++++++---- .../Tests/SendAcknowledgementTests.swift | 52 +++++++++++++++++++ 2 files changed, 71 insertions(+), 11 deletions(-) diff --git a/GraphcodeKit/Sources/Sessions/ZmxSessionLauncher.swift b/GraphcodeKit/Sources/Sessions/ZmxSessionLauncher.swift index d70cb731..a7f7ccb1 100644 --- a/GraphcodeKit/Sources/Sessions/ZmxSessionLauncher.swift +++ b/GraphcodeKit/Sources/Sessions/ZmxSessionLauncher.swift @@ -569,11 +569,7 @@ public enum ZmxSessionLauncher { return await sendRemote(text, to: node, at: remote) } guard ZmxLocator.isInstalled, !text.isEmpty else { return false } - // A shared listing can predate a session started outside this process — a pane a - // human just opened — so a refusal is confirmed against a listing of its own. - if await sessionTaskState(node) != .alive { - guard await sessionTaskState(node, freshListing: true) == .alive 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 @@ -603,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 @@ -745,11 +753,9 @@ public enum ZmxSessionLauncher { /// 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, freshListing: Bool = false) async - -> SessionTaskState - { + static func sessionTaskState(_ node: LoopNode) async -> SessionTaskState { guard ZmxLocator.isInstalled else { return .absent } - let result = await SessionListing.shared.listing(fresh: freshListing) + let result = await SessionListing.shared.listing() return sessionTaskState( lsStatus: result?.status, lsOutput: result?.output ?? "", sessionName: SurfaceRef(id: node.id, launchesClaudeCode: true).zmxSessionName) @@ -796,7 +802,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 } @@ -1126,7 +1134,7 @@ public enum ZmxSessionLauncher { 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("\terr=") + !line.contains("\tended=") && !line.contains("\texit_code=") && !line.contains("\terr=") && line.split(whereSeparator: \.isWhitespace).contains("name=\(name)") } } diff --git a/graphcode/Tests/SendAcknowledgementTests.swift b/graphcode/Tests/SendAcknowledgementTests.swift index 199c283c..93cb0563 100644 --- a/graphcode/Tests/SendAcknowledgementTests.swift +++ b/graphcode/Tests/SendAcknowledgementTests.swift @@ -246,3 +246,55 @@ struct RunCollectingOutputTests { #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) + } +} From 479572041cc81d1d1cb7a250071a3076be113ad3 Mon Sep 17 00:00:00 2001 From: scgopi Date: Thu, 1 Oct 2026 20:57:39 -0700 Subject: [PATCH 3/5] Address independent review of the send acknowledgement - A composite child's store lives for one command, so it types inline instead of on a chain no later send or parent could see (order kept). - Each post-ack typing is bounded by deliveryDeadline; a hung one is staged and logged, and the loop's follow-ups and wakes move on. - Every outcome is logged against the acknowledged request: send-typed, send-staged (delivery-failed, session-gone, deadline), send-dropped. - Lifecycle checks (start/terminate results, a resume that may have died, the first-pass kickoff, isSessionAlive) take fresh listings; a hung in-flight listing is joined for at most 30s. - A send whose drain changes nothing is broadcast once, not twice. - Tests: composite-child ordering with a slow transport, a hung typing, one broadcast per send, and a real zmx husk that send refuses. Co-Authored-By: Claude Opus 5.5 Signed-off-by: scgopi --- GraphcodeKit/Sources/GraphStore.swift | 109 +++++++++++++----- .../Sources/Sessions/ZmxSessionLauncher.swift | 41 ++++--- graphcode/Tests/MessageDeliveryTests.swift | 24 ++++ .../Tests/SendAcknowledgementTests.swift | 42 +++++++ graphcode/Tests/SubGraphAddressingTests.swift | 15 ++- 5 files changed, 185 insertions(+), 46 deletions(-) diff --git a/GraphcodeKit/Sources/GraphStore.swift b/GraphcodeKit/Sources/GraphStore.swift index 0ee01f03..d501140d 100644 --- a/GraphcodeKit/Sources/GraphStore.swift +++ b/GraphcodeKit/Sources/GraphStore.swift @@ -1344,11 +1344,16 @@ public actor GraphStore { // 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() } + if errors.isEmpty { + await broadcast() + acknowledged = graph + } } - errors += await drainAndBroadcast(broadcastErrors: broadcastErrors) + errors += await drainAndBroadcast( + broadcastErrors: broadcastErrors, unlessStillAt: acknowledged) if let error = errors.first { return .rejected(message: error, graph: graph) } @@ -3742,18 +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). - // 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). Liveness is asked before the acknowledgement - // so a sender reporting to a finished parent still hears "staged" — one shared - // listing (`SessionListing`), not a typing of the message. + // 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 { - typeAfterAcknowledging(message, to: target, resolved: true) + await typeAfterAcknowledging(message, to: target, resolved: true) return } if MessageBus.deliverability(to: target) != nil { @@ -3763,7 +3764,7 @@ public actor GraphStore { + "it will read it when it next wakes") return } - typeAfterAcknowledging(message, to: target, resolved: false) + await typeAfterAcknowledging(message, to: target, resolved: false) } /// Types a `node send` message into its target's session after the command that carried @@ -3778,13 +3779,29 @@ public actor GraphStore { /// /// 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. - private func typeAfterAcknowledging(_ message: String, to target: LoopNode, resolved: Bool) { + /// 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 typeAdHocMessage(message, to: target.id, resolved: resolved) + await typeLogged(message, to: target.id, resolved: resolved, context: context) if sessionTyping[target.id]?.token == token { sessionTyping.removeValue(forKey: target.id) } @@ -3795,6 +3812,41 @@ public actor GraphStore { 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 { @@ -3802,14 +3854,16 @@ public actor GraphStore { } } - private func typeAdHocMessage(_ message: String, to nodeID: UUID, resolved: Bool) async { - guard let target = graph.nodes[id: nodeID] else { return } + private func typeAdHocMessage( + _ message: String, to nodeID: UUID, resolved: Bool + ) async -> TypingOutcome { + guard let target = graph.nodes[id: nodeID] else { return .targetGone } if resolved { - if await deliverToSession(target, message) { return } - stageUntyped(message, to: target, reason: "session-gone") - return + if await deliverToSession(target, message) { return .typed } + stageUntyped(message, to: target) + return .staged(reason: "session-gone") } - if await deliverToSession(target, message) { return } + if await deliverToSession(target, message) { return .typed } // 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 @@ -3824,15 +3878,14 @@ public actor GraphStore { if target.runsUnattended, !target.isResolved { ensureSession(target) try? await Task.sleep(for: Self.respawnedSessionSettle) - if await deliverToSession(target, message) { return } + if await deliverToSession(target, message) { return .typed } } - stageUntyped(message, to: target, reason: "delivery-failed") + stageUntyped(message, to: target) + return .staged(reason: "delivery-failed") } - private func stageUntyped(_ message: String, to target: LoopNode, reason: String) { + private func stageUntyped(_ message: String, to target: LoopNode) { recordMemory(target.id, "while you were away: \(message)") - DaemonLog.shared.record( - "send-staged", [("node", target.id.uuidString), ("reason", reason)]) onAnnounceError?( "delivery to \(target.title)'s session failed — message staged to its memory; " + "it will read it when it next wakes") @@ -4579,7 +4632,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() @@ -4587,7 +4644,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/Sessions/ZmxSessionLauncher.swift b/GraphcodeKit/Sources/Sessions/ZmxSessionLauncher.swift index a7f7ccb1..e2270bc3 100644 --- a/GraphcodeKit/Sources/Sessions/ZmxSessionLauncher.swift +++ b/GraphcodeKit/Sources/Sessions/ZmxSessionLauncher.swift @@ -739,12 +739,17 @@ 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 @@ -753,9 +758,9 @@ public enum ZmxSessionLauncher { /// 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) async -> SessionTaskState { + static func sessionTaskState(_ node: LoopNode, fresh: Bool = false) async -> SessionTaskState { guard ZmxLocator.isInstalled else { return .absent } - let result = await SessionListing.shared.listing() + let result = await SessionListing.shared.listing(fresh: fresh) return sessionTaskState( lsStatus: result?.status, lsOutput: result?.output ?? "", sessionName: SurfaceRef(id: node.id, launchesClaudeCode: true).zmxSessionName) @@ -884,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? @@ -918,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) } @@ -929,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( @@ -946,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(()) @@ -1138,9 +1143,9 @@ public enum ZmxSessionLauncher { && line.split(whereSeparator: \.isWhitespace).contains("name=\(name)") } } - // A "no" decides whether a pane closing resolves the loop, and a shared listing can - // predate the session — so it is confirmed against a listing of its own. - if alive(in: await SessionListing.shared.listing()) { return true } + // 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)) } @@ -2633,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 } @@ -2748,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 @@ -2858,7 +2863,11 @@ extension ZmxSessionLauncher { 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 } - if let inFlight, joinBound.map({ inFlight.startedAt >= $0 }) ?? true { + // 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 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/SendAcknowledgementTests.swift b/graphcode/Tests/SendAcknowledgementTests.swift index 93cb0563..c3d58ab6 100644 --- a/graphcode/Tests/SendAcknowledgementTests.swift +++ b/graphcode/Tests/SendAcknowledgementTests.swift @@ -112,6 +112,48 @@ struct SendAcknowledgementTests { #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 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]>([]) 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 From d1674dc41c97c5b3bc70c9e335db331e268fed1b Mon Sep 17 00:00:00 2001 From: scgopi Date: Thu, 1 Oct 2026 21:07:53 -0700 Subject: [PATCH 4/5] Stop a deadline-abandoned typing from retrying or restaging withDeadline cancels the typing it gives up on but cannot stop it: the cancelled send read as a failure, so the abandoned task respawned an unattended loop, typed again behind the chain's next message, and staged the message a second time. It now returns as soon as it sees it was cancelled; typeLogged has already staged it once. Co-Authored-By: Claude Opus 5.5 Signed-off-by: scgopi --- GraphcodeKit/Sources/GraphStore.swift | 6 ++++ .../Tests/SendAcknowledgementTests.swift | 29 +++++++++++++++++++ 2 files changed, 35 insertions(+) diff --git a/GraphcodeKit/Sources/GraphStore.swift b/GraphcodeKit/Sources/GraphStore.swift index d501140d..f1c66f71 100644 --- a/GraphcodeKit/Sources/GraphStore.swift +++ b/GraphcodeKit/Sources/GraphStore.swift @@ -3858,12 +3858,17 @@ public actor GraphStore { _ 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 @@ -3879,6 +3884,7 @@ public actor GraphStore { ensureSession(target) try? await Task.sleep(for: Self.respawnedSessionSettle) if await deliverToSession(target, message) { return .typed } + if Task.isCancelled { return .staged(reason: "deadline") } } stageUntyped(message, to: target) return .staged(reason: "delivery-failed") diff --git a/graphcode/Tests/SendAcknowledgementTests.swift b/graphcode/Tests/SendAcknowledgementTests.swift index c3d58ab6..cf996871 100644 --- a/graphcode/Tests/SendAcknowledgementTests.swift +++ b/graphcode/Tests/SendAcknowledgementTests.swift @@ -138,6 +138,35 @@ struct SendAcknowledgementTests { #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 aSendIsBroadcastOnceWhenItsDrainChangesNothing() async { let broadcasts = LockIsolated(0) From 11ff1572680fe6fbefff3134d15a494c327e1041 Mon Sep 17 00:00:00 2001 From: scgopi Date: Thu, 1 Oct 2026 21:17:28 -0700 Subject: [PATCH 5/5] Check for an abandoned typing after the respawn settle too A deadline that fires during the respawn settle returned the sleep at once, and the retry then typed again inside the cancelled task. The check now follows the settle as well as each delivery attempt. Co-Authored-By: Claude Opus 5.5 Signed-off-by: scgopi --- GraphcodeKit/Sources/GraphStore.swift | 1 + .../Tests/SendAcknowledgementTests.swift | 24 +++++++++++++++++++ 2 files changed, 25 insertions(+) diff --git a/GraphcodeKit/Sources/GraphStore.swift b/GraphcodeKit/Sources/GraphStore.swift index f1c66f71..556629cb 100644 --- a/GraphcodeKit/Sources/GraphStore.swift +++ b/GraphcodeKit/Sources/GraphStore.swift @@ -3883,6 +3883,7 @@ public actor GraphStore { 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") } } diff --git a/graphcode/Tests/SendAcknowledgementTests.swift b/graphcode/Tests/SendAcknowledgementTests.swift index cf996871..eab6c1bd 100644 --- a/graphcode/Tests/SendAcknowledgementTests.swift +++ b/graphcode/Tests/SendAcknowledgementTests.swift @@ -167,6 +167,30 @@ struct SendAcknowledgementTests { #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)