Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
200 changes: 169 additions & 31 deletions GraphcodeKit/Sources/GraphStore.swift
Original file line number Diff line number Diff line change
Expand Up @@ -250,6 +250,9 @@
}

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<Void, Never>)] = [:]
private var pendingDeliveryAttempts: Set<UUID> = []
private var completedTimedOutDeliveries: [UUID: Bool] = [:]
/// `drainPendingFollowUps` runs across several awaits, and the presence poll that
Expand Down Expand Up @@ -1337,7 +1340,20 @@
// 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)
}
Expand Down Expand Up @@ -1484,7 +1500,7 @@
onRemoveMemory: onRemoveMemory,
onRefinePlaybook: onRefinePlaybook,
onRollbackPlaybook: onRollbackPlaybook,
onAnnounceError: effects.errors.append,

Check warning on line 1503 in GraphcodeKit/Sources/GraphStore.swift

View workflow job for this annotation

GitHub Actions / macos

converting non-Sendable function value to '@sendable (String) -> Void' may introduce data races
// The board's gate forwards like any other side effect: a loop inside a piloted
// composite is a real loop whose session got the standard briefing — teaching
// verbs the child store would refuse is exactly the incoherence the gate exists
Expand Down Expand Up @@ -2347,7 +2363,7 @@
pendingFollowUps.append(
PendingFollowUp(id: UUID(), nodeID: nodeID, text: message, watchedPostID: nil))
}
await drainAndBroadcast()

Check warning on line 2366 in GraphcodeKit/Sources/GraphStore.swift

View workflow job for this annotation

GitHub Actions / macos

result of call to 'drainAndBroadcast(broadcastErrors:unlessStillAt:)' is unused
}

/// A learned note into a node's memory log — `graphcode node memo`, the agent-written
Expand Down Expand Up @@ -2505,7 +2521,7 @@
let followUp = PendingFollowUp(id: UUID(), nodeID: nodeID, text: prompt, watchedPostID: nil)
pendingFollowUps.append(followUp)
goalFollowUps[nodeID] = followUp.id
await drainAndBroadcast()

Check warning on line 2524 in GraphcodeKit/Sources/GraphStore.swift

View workflow job for this annotation

GitHub Actions / macos

result of call to 'drainAndBroadcast(broadcastErrors:unlessStillAt:)' is unused
}

/// Opening a resolved loop whose session was ended brings its conversation back. Panes
Expand Down Expand Up @@ -3064,7 +3080,7 @@
return
}
if node.loopType == .composite {
await runInSubGraph(nodeID, .restartSessions, broadcastErrors: false)

Check warning on line 3083 in GraphcodeKit/Sources/GraphStore.swift

View workflow job for this annotation

GitHub Actions / macos

result of call to 'runInSubGraph(_:_:broadcastErrors:)' is unused
return
}
await restart([node])
Expand All @@ -3073,7 +3089,7 @@
private func restartSessions() async {
let live = graph.nodes.filter { !$0.isResolved }
for composite in live where composite.loopType == .composite {
await runInSubGraph(composite.id, .restartSessions, broadcastErrors: false)

Check warning on line 3092 in GraphcodeKit/Sources/GraphStore.swift

View workflow job for this annotation

GitHub Actions / macos

result of call to 'runInSubGraph(_:_:broadcastErrors:)' is unused
}
await restart(live.filter { $0.loopType != .composite })
}
Expand Down Expand Up @@ -3129,7 +3145,7 @@
// set below — a graph whose nodes have all stopped aggregates to `.idle`.
if node.loopType == .composite, let subGraph = node.subGraph {
for child in subGraph.nodes where !child.isResolved {
await runInSubGraph(

Check warning on line 3148 in GraphcodeKit/Sources/GraphStore.swift

View workflow job for this annotation

GitHub Actions / macos

result of call to 'runInSubGraph(_:_:broadcastErrors:)' is unused
node.id, .stopNode(child.id), broadcastErrors: false)
}
}
Expand Down Expand Up @@ -3679,10 +3695,10 @@
/// 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
Expand Down Expand Up @@ -3731,12 +3747,14 @@
// 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 {
Expand All @@ -3746,29 +3764,138 @@
+ "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
Expand Down Expand Up @@ -4011,6 +4138,13 @@
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
Expand Down Expand Up @@ -4239,7 +4373,7 @@
onRemoveMemory: onRemoveMemory,
onRefinePlaybook: onRefinePlaybook,
onRollbackPlaybook: onRollbackPlaybook,
onAnnounceError: effects.errors.append,

Check warning on line 4376 in GraphcodeKit/Sources/GraphStore.swift

View workflow job for this annotation

GitHub Actions / macos

converting non-Sendable function value to '@sendable (String) -> Void' may introduce data races
onMailroomEnabled: onMailroomEnabled,
goalCache: goalCache,
recurrence: effects.recurrence,
Expand All @@ -4255,7 +4389,7 @@
processRecurrence(effects.recurrence)
graph.nodes[id: ownerID]?.subGraph = await child.graph
rollUpComposite(ownerID)
await drainAndBroadcast()

Check warning on line 4392 in GraphcodeKit/Sources/GraphStore.swift

View workflow job for this annotation

GitHub Actions / macos

result of call to 'drainAndBroadcast(broadcastErrors:unlessStillAt:)' is unused
}

/// Applies the recurrence requests a child store handed up, in order — an update's
Expand Down Expand Up @@ -4326,7 +4460,7 @@
now.timeIntervalSince(node.createdAt) >= stallAfter
{
markStalled(nodeID)
await drainAndBroadcast()

Check warning on line 4463 in GraphcodeKit/Sources/GraphStore.swift

View workflow job for this annotation

GitHub Actions / macos

result of call to 'drainAndBroadcast(broadcastErrors:unlessStillAt:)' is unused
return
}

Expand Down Expand Up @@ -4505,15 +4639,19 @@
/// 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()
await drainPendingCycleReentries()
await drainPendingHandoffDeliveries()
await drainPendingNudges()
await drainPendingFollowUps()
if errors.isEmpty {
if errors.isEmpty, broadcasted == nil || graph != broadcasted {
await broadcast()
}
return errors
Expand Down
4 changes: 4 additions & 0 deletions GraphcodeKit/Sources/ProjectRegistry.swift
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
4 changes: 2 additions & 2 deletions GraphcodeKit/Sources/Sessions/CLISessionBackend.swift
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Expand Down Expand Up @@ -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 = {
Expand Down
2 changes: 1 addition & 1 deletion GraphcodeKit/Sources/Sessions/RemoteGraphAccess.swift
Original file line number Diff line number Diff line change
Expand Up @@ -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")

Expand Down
Loading
Loading