diff --git a/CLI/cmux.swift b/CLI/cmux.swift index 897cf8884d1d..5aed84843c6e 100644 --- a/CLI/cmux.swift +++ b/CLI/cmux.swift @@ -37979,6 +37979,12 @@ export default CMUXSessionRestore; /// remaining deadline. static let feedAttentionProbeTimeoutCapSeconds: TimeInterval = 1.0 + /// Send stamp that orders a hook's Feed frame against a pending blocking + /// request from the same agent (see `FeedCoordinator.supersedesPendingDecisions`). + static func feedHookSentAtMs() -> Int64 { + Int64((Date().timeIntervalSince1970 * 1000).rounded(.down)) + } + private func sendFeedTelemetry( client: SocketClient, source: String, @@ -38109,6 +38115,10 @@ export default CMUXSessionRestore; promptLength: feedPromptLength(from: parsedInput.object, compacted: true) ) event["_opencode_request_id"] = "\(source)-\(sessionId)-\(hookEventName)-\(Int(Date().timeIntervalSince1970 * 1000))" + if let agentID = firstString(in: fallbackObject, keys: ["agent_id", "agentId"]) { + event["agent_id"] = agentID + } + event["_hook_sent_at_ms"] = Self.feedHookSentAtMs() let frame: [String: Any] = [ "method": "feed.push", @@ -40711,7 +40721,7 @@ export default CMUXSessionRestore; } if isActionable { - try? waitForPriorAgentHookDeliveries( + let priorHooksDelivered = (try? waitForPriorAgentHookDeliveries( agent: source, client: activeClient, socketPassword: socketPassword, @@ -40720,7 +40730,14 @@ export default CMUXSessionRestore; max(0.01, clientDeadline.timeIntervalSinceNow) ), deadline: clientDeadline - ) + )) != nil + // Stamped only behind a completed barrier: every hook the agent + // published before this request (including this tool's own + // PreToolUse) was sent earlier, so a later-stamped hook from the + // same agent proves the decision was made elsewhere. + if priorHooksDelivered { + eventDict["_hook_sent_at_ms"] = Self.feedHookSentAtMs() + } let decisionWaitElapsed = max( 0, Date().timeIntervalSince(decisionWaitStartedAt) diff --git a/Sources/Feed/FeedCoordinator.swift b/Sources/Feed/FeedCoordinator.swift index 33a8ee832866..65f4a468ab95 100644 --- a/Sources/Feed/FeedCoordinator.swift +++ b/Sources/Feed/FeedCoordinator.swift @@ -155,6 +155,7 @@ final class FeedCoordinator: @unchecked Sendable { guard let store else { return nil } guard let item = store.ingestReturningItem(event) else { return nil } observeSemanticLifecycle(event) + retirePendingDecisionsSuperseded(by: event) if let ppid = event.ppid, ppid > 0 { armPidWatcher(ppid: ppid) } @@ -481,6 +482,49 @@ final class FeedCoordinator: @unchecked Sendable { func isAwaitingDecision(requestId: String) -> Bool { waiterRegistry.isAwaiting(requestId) } + /// Whether `event` proves its agent already moved past every earlier + /// blocking decision in the same agent context. + /// + /// Claude Code runs its PermissionRequest hook beside its own permission + /// dialog and auto-mode classifier. When the user answers in the terminal + /// or the classifier decides, Claude keeps the abandoned hook waiting until + /// the hook's own timeout, so the Feed request and its "Needs input" + /// sidebar overlay outlived the decision by up to two minutes while the + /// agent was visibly running again. Claude fires PreToolUse before the + /// permission check, and the blocking hook is stamped only after its + /// ordering barrier delivered every earlier hook, so a later-stamped tool, + /// prompt, or stop hook can only follow the decision. AskUserQuestion and + /// ExitPlanMode PreToolUse hooks announce a blocking prompt of their own. + static func supersedesPendingDecisions(_ event: WorkstreamEvent) -> Bool { + guard event.source == "claude", event.feedHookSentAtMs != nil else { return false } + switch event.hookEventName { + case .preToolUse: + return event.toolName != "AskUserQuestion" && event.toolName != "ExitPlanMode" + case .postToolUse, .postToolUseFailure, .userPromptSubmit, .stop, .sessionEnd: + return true + default: + return false + } + } + + /// Retires blocking requests that `event` proves were decided outside cmux: + /// the waiting hook returns no decision, the card expires, and the + /// needs-input overlay and banner clear, as when the user replies in Feed. + @MainActor + func retirePendingDecisionsSuperseded(by event: WorkstreamEvent) { + guard Self.supersedesPendingDecisions(event) else { return } + for (reply, itemID) in waiterRegistry.supersede(by: event) { + cancelNotification(requestId: reply.requestID) + concludeAttentionOnMain(reply.target) + notificationJournal.observeFeed(AgentFeedSemanticInput(event: reply.event, + agentKey: Self.lifecycleStatusKey(forSource: reply.event.source), + requestID: reply.requestID, resolvesRequest: true)) + clearSemanticFeedNotification(requestId: reply.requestID) + expireTimedOutItem(itemID) + waiterRegistry.cleanupStored(requestID: reply.requestID, groupID: reply.groupID) + } + } + private static func findItemId( for requestId: String, in items: [WorkstreamItem] diff --git a/Sources/Feed/FeedWaiterRegistry.swift b/Sources/Feed/FeedWaiterRegistry.swift index 6c77b433a07c..a695b48ca215 100644 --- a/Sources/Feed/FeedWaiterRegistry.swift +++ b/Sources/Feed/FeedWaiterRegistry.swift @@ -173,6 +173,35 @@ final class FeedWaiterRegistry: Sendable { } } + /// Retires every undecided request the agent provably moved past: one from + /// the same agent context (source, session, and subagent) whose hook was + /// sent before `event`. Requests or events without a send stamp never match. + func supersede(by event: WorkstreamEvent) -> [(Reply, UUID?)] { + guard let observedAtMs = event.feedHookSentAtMs else { return [] } + let session = FeedWorkstreamIdentifier.canonicalizedRawValue(agentID: event.source, rawValue: event.sessionId) + let agentID = event.feedAgentID + return groups.withLock { groups in + var superseded: [(Reply, UUID?)] = [] + for (requestID, var group) in groups { + guard group.decision == nil, group.terminalResult == nil, !group.cleanupClaimed, + group.event.source == event.source, + FeedWorkstreamIdentifier.canonicalizedRawValue( + agentID: group.event.source, rawValue: group.event.sessionId) == session, + group.event.feedAgentID == agentID, + let requestedAtMs = group.event.feedHookSentAtMs, + requestedAtMs < observedAtMs else { continue } + superseded.append((Reply(requestID: requestID, groupID: group.id, event: group.event, + target: group.target), group.itemID)) + group.target = nil + group.cleanupClaimed = true + group.terminalResult = .unavailable + groups[requestID] = group + for semaphore in group.subscribers.values { semaphore.signal() } + } + return superseded + } + } + func isAwaiting(_ requestID: String) -> Bool { groups.withLock { groups in guard let group = groups[requestID] else { return false } diff --git a/Sources/Feed/WorkstreamEvent+FeedIngress.swift b/Sources/Feed/WorkstreamEvent+FeedIngress.swift index 7058240dad2d..a3c69f51e11c 100644 --- a/Sources/Feed/WorkstreamEvent+FeedIngress.swift +++ b/Sources/Feed/WorkstreamEvent+FeedIngress.swift @@ -1,4 +1,5 @@ import CMUXAgentLaunch +import Foundation extension WorkstreamEvent { var feedIngressDeliveryKey: FeedIngressDeliveryKey { @@ -22,4 +23,25 @@ extension WorkstreamEvent { return .ordinary } } + + /// Wall-clock milliseconds at which the hook CLI sent this frame, on the + /// agent host's clock. A blocking request is stamped only after its + /// ordering barrier drained every hook the agent published before it. + var feedHookSentAtMs: Int64? { + guard let value = feedExtraFields?["_hook_sent_at_ms"] as? NSNumber, + value.int64Value >= 0 else { return nil } + return value.int64Value + } + + /// The agent's subagent identity; `nil` for the main agent thread. + var feedAgentID: String? { + guard let value = feedExtraFields?["agent_id"] as? String, + !value.isEmpty else { return nil } + return value + } + + private var feedExtraFields: [String: Any]? { + guard let data = extraFieldsJSON?.data(using: .utf8) else { return nil } + return (try? JSONSerialization.jsonObject(with: data)) as? [String: Any] + } } diff --git a/cmuxTests/ClaudeHookFeedTelemetrySwiftTests.swift b/cmuxTests/ClaudeHookFeedTelemetrySwiftTests.swift index dffc87a2b995..54a1a4497d96 100644 --- a/cmuxTests/ClaudeHookFeedTelemetrySwiftTests.swift +++ b/cmuxTests/ClaudeHookFeedTelemetrySwiftTests.swift @@ -93,6 +93,53 @@ struct ClaudeHookFeedTelemetrySwiftTests { #expect(result.status == 0, Comment(rawValue: result.stderr)) #expect(result.stdout == "{}\n") } + + // Feed retires a Claude permission request decided outside cmux when a + // later hook from the same agent arrives, which needs each telemetry + // frame's send stamp and subagent identity. + @Test func feedTelemetryCarriesSendStampAndAgentIdentity() throws { + let context = try FeedTelemetryTestContext(name: "sent-at") + defer { _ = context } + + let workspaceID = "11111111-1111-1111-1111-111111111111" + let surfaceID = "22222222-2222-2222-2222-222222222222" + let ttyName = "ttys-claude-sent-at" + let feedSeen = DispatchSemaphore(value: 0) + startServer( + listenerFD: context.listenerFD, + state: context.state, + workspaceID: workspaceID, + focusedSurfaceID: surfaceID, + ttyName: ttyName, + resolvedSurfaceID: surfaceID, + feedSeen: feedSeen + ) + + let cliPath = try BundledCLITestSupport.bundledCLIPath(for: BundledCLILinkageTests.self) + let startedAtMs = Int64(Date().timeIntervalSince1970 * 1000) + let result = runProcess( + executablePath: cliPath, + arguments: ["hooks", "claude", "session-start"], + environment: context.environment( + workspaceID: workspaceID, + surfaceID: surfaceID, + ttyName: ttyName + ), + standardInput: #"{"session_id":"claude-sent-at-session","source":"startup","cwd":"\#(context.root.path)","hook_event_name":"SessionStart","agent_id":"subagent-7"}"#, + timeout: 5 + ) + + #expect(result.timedOut == false, Comment(rawValue: result.stderr)) + #expect(result.status == 0, Comment(rawValue: result.stderr)) + #expect(feedSeen.wait(timeout: .now() + 5) == .success, "Expected feed.push, saw \(context.state.commandsSnapshot())") + let event = try #require( + context.state.feedEventsSnapshot().last { $0["hook_event_name"] as? String == "SessionStart" }, + "Expected SessionStart feed telemetry, saw \(context.state.commandsSnapshot())" + ) + let sentAtMs = try #require((event["_hook_sent_at_ms"] as? NSNumber)?.int64Value, "event=\(event)") + #expect(sentAtMs >= startedAtMs, "event=\(event)") + #expect(event["agent_id"] as? String == "subagent-7", "event=\(event)") + } } private final class FeedTelemetryTestContext { diff --git a/cmuxTests/FeedCoordinatorTests.swift b/cmuxTests/FeedCoordinatorTests.swift index 08f09b44de1c..cf1e4457081e 100644 --- a/cmuxTests/FeedCoordinatorTests.swift +++ b/cmuxTests/FeedCoordinatorTests.swift @@ -767,6 +767,92 @@ struct FeedCoordinatorTests { #expect(attention.events.first?.hookEventName == .permissionRequest) } + /// Claude Code keeps its PermissionRequest hook waiting after the user + /// answers the prompt in the terminal or the auto-mode classifier decides, + /// so the Feed request (and the "Needs input" overlay it owns) outlived the + /// decision until the hook timed out, beside the agent's own "Running". + /// A later hook from the same agent proves the decision was made elsewhere. + @Test func laterClaudeHookRetiresPermissionDecidedOutsideFeed() async { + defer { Self.resetFeedCoordinatorTestHooks() } + let requestId = "claude-decided-in-terminal-request" + let sessionId = "claude-decided-in-terminal-session" + let ingested = DispatchSemaphore(value: 0) + await MainActor.run { + FeedCoordinator.shared.install(store: WorkstreamStore(ringCapacity: 10)) + FeedCoordinatorTestHooks.afterBlockingEventIngested = { _, ingestedRequestId in + if ingestedRequestId == requestId { ingested.signal() } + } + } + + let permission = WorkstreamEvent( + sessionId: sessionId, + hookEventName: .permissionRequest, + source: "claude", + cwd: "/tmp", + toolName: "Write", + toolInputJSON: #"{"file_path":"/tmp/memory.md"}"#, + requestId: requestId, + extraFieldsJSON: #"{"_hook_sent_at_ms":2000}"# + ) + let done = DispatchSemaphore(value: 0) + let resultBox = IngestResultBox() + let startedAt = ContinuousClock.now + DispatchQueue.global(qos: .userInitiated).async { + resultBox.value = FeedCoordinator.shared.ingestBlocking(event: permission, waitTimeout: 5) + done.signal() + } + guard waitForFeedTestSignal(ingested, timeout: .now() + 2) == .success else { + Issue.record("the blocking PermissionRequest was never ingested") + return + } + + func deliver(_ event: WorkstreamEvent) { + // Acknowledged (non-decision) ingress commits before returning. + let delivered = DispatchSemaphore(value: 0) + DispatchQueue.global(qos: .userInitiated).async { + _ = FeedCoordinator.shared.ingestBlocking(event: event, waitTimeout: 1) + delivered.signal() + } + #expect(waitForFeedTestSignal(delivered, timeout: .now() + 2) == .success) + } + // The tool's own PreToolUse is sent before its permission request. + deliver(WorkstreamEvent( + sessionId: sessionId, hookEventName: .preToolUse, source: "claude", + toolName: "Write", extraFieldsJSON: #"{"_hook_sent_at_ms":1990}"# + )) + // A subagent's tool proves nothing about the main agent's prompt. + deliver(WorkstreamEvent( + sessionId: sessionId, hookEventName: .preToolUse, source: "claude", + toolName: "Bash", extraFieldsJSON: #"{"_hook_sent_at_ms":3000,"agent_id":"subagent-1"}"# + )) + #expect( + FeedCoordinator.shared.isAwaitingDecision(requestId: requestId), + "earlier or other-agent hooks must not retire a live permission request" + ) + + // The user answered in the terminal; Claude moved on to the next tool. + deliver(WorkstreamEvent( + sessionId: sessionId, hookEventName: .preToolUse, source: "claude", + toolName: "Bash", extraFieldsJSON: #"{"_hook_sent_at_ms":3000}"# + )) + guard waitForFeedTestSignal(done, timeout: .now() + 2) == .success else { + Issue.record("a later PreToolUse must retire the permission request decided in the terminal") + return + } + #expect(startedAt.duration(to: .now) < .seconds(4)) + guard case .unavailable = resultBox.value else { + Issue.record("the superseded hook must return no decision, got \(String(describing: resultBox.value))") + return + } + let status = await MainActor.run { + FeedCoordinator.shared.store.items.first { $0.kind == .permissionRequest }?.status + } + guard case .expired = status else { + Issue.record("the superseded permission card must stop being actionable") + return + } + } + @Test func blockingDecisionEventPredicateCoversEveryDecisionKind() { // The three blocking-decision kinds must all surface attention… #expect(FeedCoordinator.isBlockingDecisionEvent(.permissionRequest)) diff --git a/cmuxTests/FeedWaiterRegistryTests.swift b/cmuxTests/FeedWaiterRegistryTests.swift index b8c72dc95b64..ebed56d249d6 100644 --- a/cmuxTests/FeedWaiterRegistryTests.swift +++ b/cmuxTests/FeedWaiterRegistryTests.swift @@ -63,6 +63,53 @@ struct FeedWaiterRegistryTests { registry.cleanupStored(requestID: reply.requestID, groupID: reply.groupID) } + @Test func onlyALaterHookFromTheSameAgentSupersedesARequest() throws { + func stamped(_ hook: WorkstreamEvent.HookEventName, sentAt: Int?, agentID: String? = nil, + source: String = "claude", session: String = "session") -> WorkstreamEvent { + var extra: [String] = [] + if let sentAt { extra.append(#""_hook_sent_at_ms":\#(sentAt)"#) } + if let agentID { extra.append(#""agent_id":"\#(agentID)""#) } + return WorkstreamEvent(sessionId: session, hookEventName: hook, source: source, + toolName: "Tool", toolInputJSON: "{}", requestId: hook == .permissionRequest ? "request" : nil, + extraFieldsJSON: extra.isEmpty ? nil : "{\(extra.joined(separator: ","))}") + } + let registry = FeedWaiterRegistry() + let request = stamped(.permissionRequest, sentAt: 2_000) + let registration = try #require(registry.register(requestID: "request", event: request)) + registry.accepted(registration, event: request, item: item()) + + // The tool's own PreToolUse precedes its request; unstamped hooks and + // other sessions, sources, or subagents prove nothing. + #expect(registry.supersede(by: stamped(.preToolUse, sentAt: 1_990)).isEmpty) + #expect(registry.supersede(by: stamped(.preToolUse, sentAt: 2_000)).isEmpty) + #expect(registry.supersede(by: stamped(.preToolUse, sentAt: nil)).isEmpty) + #expect(registry.supersede(by: stamped(.preToolUse, sentAt: 3_000, session: "other")).isEmpty) + #expect(registry.supersede(by: stamped(.preToolUse, sentAt: 3_000, source: "codex")).isEmpty) + #expect(registry.supersede(by: stamped(.preToolUse, sentAt: 3_000, agentID: "subagent")).isEmpty) + #expect(registry.isAwaiting("request")) + + let superseded = registry.supersede(by: stamped(.preToolUse, sentAt: 3_000)) + #expect(superseded.map { $0.0.requestID } == ["request"]) + #expect(registration.semaphore.wait(timeout: .now()) == .success) + #expect(registry.supersede(by: stamped(.stop, sentAt: 4_000)).isEmpty) + let finished = registry.finish(registration) + #expect(!finished.shouldCancel) + guard case .unavailable = finished.outcome.result else { + Issue.record("A superseded request must return no decision") + return + } + } + + @Test func unstampedRequestIsNeverSuperseded() throws { + let registry = FeedWaiterRegistry() + let registration = try #require(registry.register(requestID: "request", event: event())) + registry.accepted(registration, event: event(), item: item()) + let later = WorkstreamEvent(sessionId: "session", hookEventName: .preToolUse, source: "claude", + toolName: "Bash", extraFieldsJSON: #"{"_hook_sent_at_ms":9000}"#) + #expect(registry.supersede(by: later).isEmpty) + #expect(registry.isAwaiting("request")) + } + @Test func lateFailureCannotReplaceDecisionBeforeStoreCommit() throws { let registry = FeedWaiterRegistry() let registration = try #require(registry.register(requestID: "request", event: event()))