diff --git a/ProwlCLI/Commands/AgentsCommand.swift b/ProwlCLI/Commands/AgentsCommand.swift index b4ab3232..84865adc 100644 --- a/ProwlCLI/Commands/AgentsCommand.swift +++ b/ProwlCLI/Commands/AgentsCommand.swift @@ -6,8 +6,8 @@ import ProwlCLIShared struct AgentsCommand: ParsableCommand { static let configuration = CommandConfiguration( commandName: "agents", - abstract: "List detected agent panes.", - subcommands: [AgentsReadCommand.self] + abstract: "List or report agent state.", + subcommands: [AgentsReadCommand.self, AgentsSignalCommand.self] ) @OptionGroup var options: GlobalOptions diff --git a/ProwlCLI/Commands/AgentsSignalCommand.swift b/ProwlCLI/Commands/AgentsSignalCommand.swift new file mode 100644 index 00000000..642a3dea --- /dev/null +++ b/ProwlCLI/Commands/AgentsSignalCommand.swift @@ -0,0 +1,59 @@ +import ArgumentParser +import Foundation +import ProwlCLIShared + +extension AgentSignalEvent: ExpressibleByArgument {} + +struct AgentsSignalCommand: ParsableCommand { + static let configuration = CommandConfiguration( + commandName: "signal", + abstract: "Report an event for the agent in the caller's Prowl pane." + ) + + @Argument(help: "Event: turn-ended, needs-input, session-start, session-end, or progress.") + var event: AgentSignalEvent + + @Option(name: .long, help: "Progress from 0 through 100; omit for indeterminate progress.") + var progress: Int? + + @Option(name: .long, help: "Claimed producer origin metadata; does not upgrade trust.") + var origin: String? + + @Option(name: .customLong("session"), help: "Opaque agent session identifier.") + var sessionID: String? + + @Option(name: .long, help: "Short result or reason returned with the signal (maximum 4096 UTF-8 bytes).") + var detail: String? + + @OptionGroup var options: GlobalOptions + + mutating func run() throws { + try CLIExecution.run(command: "agents.signal", output: options.outputMode, colorEnabled: options.colorEnabled) { + let envelope = CommandEnvelope( + output: options.outputMode, + command: .agentsSignal( + AgentSignalInput( + event: event, + progress: progress, + origin: origin, + sessionID: sessionID, + detail: detail + )) + ) + try CLIRunner.execute(envelope) + } + } + + func validate() throws { + let input = AgentSignalInput( + event: event, + progress: progress, + origin: origin, + sessionID: sessionID, + detail: detail + ) + if let message = input.validationErrorMessage { + throw ValidationError(message) + } + } +} diff --git a/ProwlCLI/Output/OutputRenderer.swift b/ProwlCLI/Output/OutputRenderer.swift index 15d06830..79e4fc14 100644 --- a/ProwlCLI/Output/OutputRenderer.swift +++ b/ProwlCLI/Output/OutputRenderer.swift @@ -76,6 +76,14 @@ enum OutputRenderer { return } + if response.command == "agents.signal", + let data = response.data, + let payload = try? data.decode(as: AgentSignalCommandPayload.self) + { + print(agentSignalText(payload)) + return + } + if response.command == "profiles", let data = response.data, let payload = try? data.decode(as: ProfilesCommandPayload.self) @@ -275,6 +283,16 @@ enum OutputRenderer { }.joined(separator: "\n") } + static func agentSignalText(_ payload: AgentSignalCommandPayload) -> String { + let event: String + if payload.signal.event == .progress, let progress = payload.signal.progress { + event = "progress=\(progress)" + } else { + event = payload.signal.event.rawValue + } + return "Signaled \(event) for pane \(payload.pane.id)." + } + private static func renderProfiles(_ payload: ProfilesCommandPayload) -> String { guard !payload.profiles.isEmpty else { return "No Agent Profiles found." } return payload.profiles.map { profile in diff --git a/ProwlCLIContracts/Resources/cli-output-schema.json b/ProwlCLIContracts/Resources/cli-output-schema.json index 12e971c8..6063d514 100644 --- a/ProwlCLIContracts/Resources/cli-output-schema.json +++ b/ProwlCLIContracts/Resources/cli-output-schema.json @@ -16,6 +16,9 @@ { "$ref": "#/$defs/agentsReadResponse" }, + { + "$ref": "#/$defs/agentsSignalResponse" + }, { "$ref": "#/$defs/profilesResponse" }, @@ -674,6 +677,85 @@ } } }, + "agentSignalData": { + "type": "object", + "additionalProperties": false, + "required": [ + "pane", + "signal" + ], + "properties": { + "pane": { + "type": "object", + "additionalProperties": false, + "required": [ + "id", + "worktree_id" + ], + "properties": { + "id": { + "type": "string", + "format": "uuid" + }, + "worktree_id": { + "type": "string", + "minLength": 1 + } + } + }, + "signal": { + "type": "object", + "additionalProperties": false, + "required": [ + "event", + "source", + "confidence", + "at" + ], + "properties": { + "event": { + "enum": [ + "turn-ended", + "needs-input", + "session-start", + "session-end", + "progress" + ] + }, + "progress": { + "type": "integer", + "minimum": 0, + "maximum": 100 + }, + "source": { + "const": "cooperative_cli" + }, + "confidence": { + "const": "exact" + }, + "at": { + "type": "string", + "format": "date-time" + }, + "session_id": { + "type": "string", + "minLength": 1, + "maxLength": 256 + }, + "detail": { + "type": "string", + "minLength": 1, + "maxLength": 4096 + }, + "claimed_origin": { + "type": "string", + "minLength": 1, + "maxLength": 256 + } + } + } + } + }, "profilesData": { "type": "object", "additionalProperties": false, @@ -1660,6 +1742,58 @@ } ] }, + "agentsSignalResponse": { + "oneOf": [ + { + "type": "object", + "additionalProperties": false, + "required": [ + "ok", + "command", + "schema_version", + "data" + ], + "properties": { + "ok": { + "const": true + }, + "command": { + "const": "agents.signal" + }, + "schema_version": { + "const": "prowl.cli.agents.signal.v1" + }, + "data": { + "$ref": "#/$defs/agentSignalData" + } + } + }, + { + "type": "object", + "additionalProperties": false, + "required": [ + "ok", + "command", + "schema_version", + "error" + ], + "properties": { + "ok": { + "const": false + }, + "command": { + "const": "agents.signal" + }, + "schema_version": { + "const": "prowl.cli.agents.signal.v1" + }, + "error": { + "$ref": "#/$defs/error" + } + } + } + ] + }, "profilesResponse": { "oneOf": [ { diff --git a/ProwlCLITests/AgentSignalOutputRendererTests.swift b/ProwlCLITests/AgentSignalOutputRendererTests.swift new file mode 100644 index 00000000..b1308783 --- /dev/null +++ b/ProwlCLITests/AgentSignalOutputRendererTests.swift @@ -0,0 +1,49 @@ +import Foundation +import ProwlCLIShared +import XCTest + +@testable import prowl + +final class AgentSignalOutputRendererTests: XCTestCase { + func testTextReceiptNamesEventAndPaneWithoutEchoingDetail() { + let payload = AgentSignalCommandPayload( + pane: AgentSignalPanePayload(id: "pane-1", worktreeID: "wt-1"), + signal: AgentSignalPayload( + event: .turnEnded, + progress: nil, + source: "cooperative_cli", + confidence: "exact", + timestamp: "1970-01-01T00:16:40.000Z", + sessionID: "session-1", + detail: "sensitive result", + claimedOrigin: nil + ) + ) + + XCTAssertEqual( + OutputRenderer.agentSignalText(payload), + "Signaled turn-ended for pane pane-1." + ) + } + + func testProgressReceiptIncludesValue() { + let payload = AgentSignalCommandPayload( + pane: AgentSignalPanePayload(id: "pane-1", worktreeID: "wt-1"), + signal: AgentSignalPayload( + event: .progress, + progress: 75, + source: "cooperative_cli", + confidence: "exact", + timestamp: "1970-01-01T00:16:40.000Z", + sessionID: nil, + detail: nil, + claimedOrigin: nil + ) + ) + + XCTAssertEqual( + OutputRenderer.agentSignalText(payload), + "Signaled progress=75 for pane pane-1." + ) + } +} diff --git a/ProwlCLITests/AgentsCommandParsingTests.swift b/ProwlCLITests/AgentsCommandParsingTests.swift index 47a36046..e704ada8 100644 --- a/ProwlCLITests/AgentsCommandParsingTests.swift +++ b/ProwlCLITests/AgentsCommandParsingTests.swift @@ -24,4 +24,37 @@ final class AgentsCommandParsingTests: XCTestCase { XCTAssertThrowsError(try command.validateOutputMode()) } + + func testSignalParsesEventAndMetadata() throws { + let command = try AgentsSignalCommand.parse([ + "turn-ended", + "--origin", "manual-review", + "--session", "session-1", + "--detail", "Review complete", + "--json", + ]) + + XCTAssertEqual(command.event, .turnEnded) + XCTAssertNil(command.progress) + XCTAssertEqual(command.origin, "manual-review") + XCTAssertEqual(command.sessionID, "session-1") + XCTAssertEqual(command.detail, "Review complete") + XCTAssertTrue(command.options.json) + } + + func testSignalParsesIndeterminateAndNumericProgress() throws { + let indeterminate = try AgentsSignalCommand.parse(["progress"]) + let numeric = try AgentsSignalCommand.parse(["progress", "--progress", "75"]) + + XCTAssertEqual(indeterminate.event, .progress) + XCTAssertNil(indeterminate.progress) + XCTAssertEqual(numeric.progress, 75) + } + + func testSignalRejectsInvalidEventProgressAndBoundedText() { + XCTAssertThrowsError(try AgentsSignalCommand.parse(["complete"])) + XCTAssertThrowsError(try AgentsSignalCommand.parse(["turn-ended", "--progress", "1"])) + XCTAssertThrowsError(try AgentsSignalCommand.parse(["progress", "--progress", "101"])) + XCTAssertThrowsError(try AgentsSignalCommand.parse(["needs-input", "--detail", String(repeating: "x", count: 4_097)])) + } } diff --git a/ProwlCLITests/ProwlCLIIntegrationTests.swift b/ProwlCLITests/ProwlCLIIntegrationTests.swift index 5dca9cac..bcb7c18d 100644 --- a/ProwlCLITests/ProwlCLIIntegrationTests.swift +++ b/ProwlCLITests/ProwlCLIIntegrationTests.swift @@ -153,6 +153,95 @@ final class ProwlCLIIntegrationTests: XCTestCase { XCTAssertEqual(payload["command"] as? String, "agents") } + func testAgentsSignalRoundTripsOverSocketInTextMode() throws { + let socketPath = temporarySocketPath(suffix: "agents-signal-text") + let paneID = "6E1A2A10-D99F-4E3F-920C-D93AA3C05764" + let response = try CommandResponse( + ok: true, + command: "agents.signal", + schemaVersion: "prowl.cli.agents.signal.v1", + data: RawJSON( + encoding: AgentSignalCommandPayload( + pane: AgentSignalPanePayload(id: paneID, worktreeID: "/Projects/App"), + signal: AgentSignalPayload( + event: .turnEnded, + progress: nil, + source: "cooperative_cli", + confidence: "exact", + timestamp: "2026-08-22T12:00:00.000Z", + sessionID: "session-1", + detail: "Review complete", + claimedOrigin: "manual-review" + ) + )) + ) + + let (requestData, result) = try runWithMockServer( + socketPath: socketPath, + response: response, + args: [ + "agents", "signal", "turn-ended", + "--origin", "manual-review", + "--session", "session-1", + "--detail", "Review complete", + "--no-color", + ] + ) + + XCTAssertEqual(result.exitCode, 0) + XCTAssertEqual(result.stdout, "Signaled turn-ended for pane \(paneID).\n") + let envelope = try JSONDecoder().decode(CommandEnvelope.self, from: requestData) + guard case .agentsSignal(let input) = envelope.command else { + return XCTFail("Expected agents.signal command envelope") + } + XCTAssertEqual(input.event, .turnEnded) + XCTAssertEqual(input.origin, "manual-review") + XCTAssertEqual(input.sessionID, "session-1") + XCTAssertEqual(input.detail, "Review complete") + } + + func testAgentsSignalPreservesSchemaValidatedJSONResponse() throws { + let socketPath = temporarySocketPath(suffix: "agents-signal-json") + let response = try CommandResponse( + ok: true, + command: "agents.signal", + schemaVersion: "prowl.cli.agents.signal.v1", + data: RawJSON( + encoding: AgentSignalCommandPayload( + pane: AgentSignalPanePayload( + id: "6E1A2A10-D99F-4E3F-920C-D93AA3C05764", + worktreeID: "/Projects/App" + ), + signal: AgentSignalPayload( + event: .progress, + progress: 75, + source: "cooperative_cli", + confidence: "exact", + timestamp: "2026-08-22T12:00:00.000Z", + sessionID: nil, + detail: nil, + claimedOrigin: nil + ) + )) + ) + + let (requestData, result) = try runWithMockServer( + socketPath: socketPath, + response: response, + args: ["agents", "signal", "progress", "--progress", "75", "--json"] + ) + + XCTAssertEqual(result.exitCode, 0) + let request = try JSONDecoder().decode(CommandEnvelope.self, from: requestData) + guard case .agentsSignal(let input) = request.command else { + return XCTFail("Expected agents.signal command envelope") + } + XCTAssertEqual(input.progress, 75) + let rendered = try jsonObject(from: result.stdout) + XCTAssertEqual(rendered["command"] as? String, "agents.signal") + XCTAssertEqual(rendered["schema_version"] as? String, "prowl.cli.agents.signal.v1") + } + func testAgentsPayloadDetectionReasonRemainsBackwardCompatible() throws { let modernData = try JSONEncoder().encode( AgentsResponseData( diff --git a/docs-ai/013-prowl-cli/contracts/agents-signal.md b/docs-ai/013-prowl-cli/contracts/agents-signal.md new file mode 100644 index 00000000..afd543f6 --- /dev/null +++ b/docs-ai/013-prowl-cli/contracts/agents-signal.md @@ -0,0 +1,92 @@ +# `prowl agents signal` Contract + +Current version: `prowl.cli.agents.signal.v1`. + +```bash +prowl agents signal + [--progress <0...100>] + [--session ] + [--origin ] + [--detail ] + [--json|--no-color] +``` + +## Semantics + +The command reports an observation for the pane that spawned the calling `prowl` +process. It accepts no target selector. The app obtains the kernel peer PID from the +Unix socket, walks process ancestry against live pane shell PIDs, and either attributes the +signal to that exact pane or fails with `SOURCE_REQUIRED`. UI focus and +`PROWL_PANE_ID` are never fallback identity sources. + +Events: + +- `turn-ended` — the runtime ended one interaction turn; it does not mean an assigned task + or workflow step completed. +- `needs-input` — progress requires user/agent input. +- `session-start` / `session-end` — producer-reported session lifecycle. +- `progress` — indeterminate when `--progress` is absent, otherwise 0 through 100. + +S1 records every public invocation as `source: cooperative_cli`, `confidence: exact`. +Here `exact` means explicit channel plus exact caller-pane attribution; it does not make the +producer's business judgment authoritative. `--origin` is caller-authored metadata only and +cannot upgrade source/confidence or satisfy a future native-hook capability check. + +`--session` and `--origin` are non-empty, control-free UTF-8 up to 256 bytes. `--detail` +is non-empty, control-free UTF-8 up to 4096 bytes. Detail is a short result or reason returned +with the signal; it is not logged and does not change confidence. Large results use +`agents read`; workflow outputs use `workflow done -`. + +## Success response + +Text: + +```text +Signaled turn-ended for pane . +Signaled progress=75 for pane . +``` + +JSON: + +```json +{ + "ok": true, + "command": "agents.signal", + "schema_version": "prowl.cli.agents.signal.v1", + "data": { + "pane": { + "id": "6E1A2A10-D99F-4E3F-920C-D93AA3C05764", + "worktree_id": "/Projects/Prowl" + }, + "signal": { + "event": "turn-ended", + "source": "cooperative_cli", + "confidence": "exact", + "at": "2026-08-22T12:00:00.000Z", + "session_id": "session-1", + "detail": "Review complete", + "claimed_origin": "manual-review" + } + } +} +``` + +Optional fields are omitted rather than encoded as `null`. The executable schema is +`#/$defs/agentsSignalResponse` in +[`cli-output-schema.json`](../../../ProwlCLIContracts/Resources/cli-output-schema.json). + +## Errors + +- `INVALID_ARGUMENT` — invalid event/option combination, progress range, empty value, + control character, or UTF-8 byte limit. +- `SOURCE_REQUIRED` — the socket peer process cannot be attributed to a live Prowl pane + (including an external terminal or ancestry broken by tmux/detached wrappers). +- `AGENT_GONE` — the attributed pane closed before the signal was recorded. + +## Deferred paired completion + +`dispatch-complete` is deliberately not part of v1 S1. S2 ships one atomic paired path: +`create --profile --prompt` returns an opaque `dispatch_id`; the agent reports +`dispatch-complete --detail`; a bounded in-memory receipt survives pane closure but not app +restart; and `agents wait --dispatch` re-snapshots after observer overflow. Generic runtime +`turn-ended` never substitutes for that dispatch receipt or `workflow done`. diff --git a/docs-ai/013-prowl-cli/contracts/agents.md b/docs-ai/013-prowl-cli/contracts/agents.md index 20444307..9482ffd8 100644 --- a/docs-ai/013-prowl-cli/contracts/agents.md +++ b/docs-ai/013-prowl-cli/contracts/agents.md @@ -12,7 +12,9 @@ agent `type`/`name`, `status`, `raw_state`, optional `detection_reason`, `last_changed_at`, project/worktree/tab/pane metadata, and optional session attribution. Text output additionally shows a current-process `pN` handle. -Use `prowl agents read ` for a semantic agent snapshot; that command -has its own [contract](agents-read.md). The complete roster response schema is +Use `prowl agents read ` for a semantic agent snapshot. A process inside +a Prowl pane can report cooperative runtime events with `prowl agents signal`; these +commands have separate [read](agents-read.md) and [signal](agents-signal.md) contracts. +The complete roster response schema is `#/$defs/agentsResponse` in [`schema-bundle.json`](../../../ProwlCLIContracts/Resources/cli-output-schema.json). diff --git a/docs-ai/013-prowl-cli/contracts/architecture.md b/docs-ai/013-prowl-cli/contracts/architecture.md index ea2d933f..6f9c0f41 100644 --- a/docs-ai/013-prowl-cli/contracts/architecture.md +++ b/docs-ai/013-prowl-cli/contracts/architecture.md @@ -85,7 +85,7 @@ Phase-1 commands are **remote-control actions on running app state**. - `ReadCommandHandler` - `LifecycleCommandHandler` (`create`, `close`) - legacy `TabCommandHandler` / `PaneCommandHandler` during deprecation - - `AgentsCommandHandler`, `AgentReadCommandHandler`, and `HandoffCommandHandler` + - `AgentsCommandHandler`, `AgentReadCommandHandler`, `AgentSignalCommandHandler`, and `HandoffCommandHandler` - Shared services - `TargetResolver` - `TerminalCommandBridge` @@ -145,7 +145,7 @@ Resolution belongs to app runtime (state-aware), with CLI only enforcing selecto - Input normalization rules: `input.md` and `targeting.md` - Output contracts: one document per wire command, including `create.md`, `close.md`, - deprecated `tab.md` / `pane.md`, `agents.md`, and `handoff.md`. + deprecated `tab.md` / `pane.md`, `agents.md`, `agents-signal.md`, and `handoff.md`. - JSON schema validation source: the machine-readable bundle linked by `schema.md`. Every payload-bearing mock socket response is validated against that Draft 2020-12 diff --git a/docs-ai/013-prowl-cli/contracts/input.md b/docs-ai/013-prowl-cli/contracts/input.md index 9ab49eca..2b7fbba7 100644 --- a/docs-ai/013-prowl-cli/contracts/input.md +++ b/docs-ai/013-prowl-cli/contracts/input.md @@ -9,7 +9,7 @@ and the executable [schema bundle](schema.md). ```text prowl [path] prowl open [path] -prowl list | agents | profiles | focus | read | send | key | handoff | create | close +prowl list | agents [read|signal] | profiles | focus | read | send | key | handoff | create | close ``` Bare path forms (`/`, `./`, `../`, `~/`, `file://`, `.`, `..`) enter `open`. @@ -59,10 +59,25 @@ is a read-only global snapshot and accepts no target. `close` requires a pane-or shipped release. They keep their legacy parser/transport behavior while emitting a stderr warning; new automation must use the lifecycle grammar above. +## Agent signal grammar + +```bash +prowl agents signal + [--progress <0...100>] [--session ] + [--origin ] [--detail ] +``` + +`--progress` is valid only with `progress`; omitting it means indeterminate progress. +Session/origin are at most 256 UTF-8 bytes and detail is at most 4096. All are non-empty +and control-free when present. Parser and handler enforce the same shared validation. +See [agents-signal.md](agents-signal.md). + ## Command-specific exceptions - `agents read ` is a pane-only semantic snapshot, no selectors or focus fallback. +- `agents signal` accepts no selector. Its source is the caller pane resolved from the + socket peer process ancestry, never UI focus or `PROWL_PANE_ID`. - `handoff` defaults to the calling pane, not UI focus. - `list`, `agents`, and `profiles list` are global discovery commands with no target selector. - `open` consumes a path rather than a target. diff --git a/docs-ai/013-prowl-cli/contracts/schema.md b/docs-ai/013-prowl-cli/contracts/schema.md index 7ce6051e..e5c0303f 100644 --- a/docs-ai/013-prowl-cli/contracts/schema.md +++ b/docs-ai/013-prowl-cli/contracts/schema.md @@ -15,6 +15,7 @@ The bundle has one versioned success-or-error response schema for every wire com | `list` | `#/$defs/listResponse` | | `agents` | `#/$defs/agentsResponse` | | `agents.read` | `#/$defs/agentsReadResponse` | +| `agents.signal` | `#/$defs/agentsSignalResponse` | | `profiles` | `#/$defs/profilesResponse` | | `focus` | `#/$defs/focusResponse` | | `send` | `#/$defs/sendResponse` | diff --git a/docs-ai/063-agent-workflows/000-plan.md b/docs-ai/063-agent-workflows/000-plan.md index 9abe5bea..25a40237 100644 --- a/docs-ai/063-agent-workflows/000-plan.md +++ b/docs-ai/063-agent-workflows/000-plan.md @@ -157,19 +157,27 @@ reducer stream and typed so that disappearance is observable. The observer is de by 064-S1 in release R1 (it was first specified here) and consumed by this entry's B3: ```swift +struct AgentObservationSnapshot: Sendable { + let agent: ActiveAgentEntry? + let latestSignal: AgentSignal? + let revision: UInt64 +} enum ObservedAgentState: Sendable { - case snapshot(ActiveAgentEntry?) // always the first element, even when no agent is detected + case snapshot(AgentObservationSnapshot) // always first, including a normal shell case changed(ActiveAgentEntry) - case removed // agent process gone; entry removed (today: agentEntryRemoved) - case surfaceClosed // pane closed; stream finishes after this + case removed // published agent gone; pane may remain alive + case signal(AgentSignal) + case surfaceClosed // pane closed; stream finishes after this } -func observeAgentState(surfaceID: UUID) -> AsyncStream +func observeAgentState(surfaceID: UUID) -> AsyncThrowingStream ``` -Each subscriber gets its own bounded buffer (newest-wins coalescing, like -`TerminalEventCoalescer`); registration and snapshot capture happen in one main-actor -step so no change can fall between them, the snapshot precedes live events, cancellation -removes the subscriber, and `surfaceClosed` terminates the stream. `agents wait` maps `removed` / +Each subscriber gets its own bounded buffer. Registration and snapshot capture happen in +one main-actor step so no change can fall between them, the snapshot precedes live events, +cancellation removes the subscriber, and `surfaceClosed` terminates the stream. A slow +subscriber receives an explicit `bufferOverflow` error instead of silently losing signal or +lifecycle evidence; S2's `agents wait` re-subscribes and evaluates the newer snapshot before +surfacing an error. `agents wait` maps `removed` / `surfaceClosed` to a terminal `AGENT_GONE` error (not to `done`) unless `--until changed` was requested. The runner's watchdog likewise reads the role's *current* state first and schedules cancellable grace deadlines on the injected clock; it never relies on a later diff --git a/docs-ai/063-agent-workflows/release-plan.md b/docs-ai/063-agent-workflows/release-plan.md index 765690bc..2f79aa78 100644 --- a/docs-ai/063-agent-workflows/release-plan.md +++ b/docs-ai/063-agent-workflows/release-plan.md @@ -31,9 +31,9 @@ user-facing surface may merge before "their" release and stay dormant. Three rel | C0 | Merged | #709 | | A1 | Merged | #710 | | A1b | Merged | #713 | -| A2 | Implemented | #714 | -| S1 | Planned | **Next critical-path PR**: signal bus, multicast observer, `agents signal` | -| S2 | Planned | Follows S1: `agents wait` and honest heuristic fallback | +| A2 | Merged | #714 | +| S1 | In progress | `feat/agent-completion-signal-bus`: bus, multicast observer, `agents signal` | +| S2 | Planned | Follows S1: atomic dispatch ID → completion receipt → `agents wait` path and honest heuristic fallback | | S3 wave 1 | Planned | Follows A2 + S1: tier-A launch hooks | | 065-S0/K1 | Planned, parallel | Skill-target spike + bundled-skill registry | | 065-K2/K3 | Planned | Follow S0/K1 inside R1 | @@ -50,9 +50,9 @@ A2 completes 063's R1 implementation work. The orchestration critical path now m | 1 | **A1b** `PROWL_PANE_ID` per-pane environment variable (joins `PROWL_WORKTREE_PATH` / `PROWL_ROOT_PATH`) + `prowl-cli` skill self-identification rewrite | 063 | A1 | agents address their own pane deterministically (`--pane "$PROWL_PANE_ID"`) instead of guessing from `focused` | | 1 | **065-S0/K1** skill-target spike; `embed-skills` + `ProwlSkills` registry | 065 | — | skills ship in the bundle; D1 prerequisite | | 2 | **A2** profile launch boundary + `create tab\|pane --profile

--prompt -` + `profiles list` | 063 | A1 | CLI launches a profile with a kickoff prompt and gets the pane back | -| 2 | **S1** signal bus + `ObservedAgentState` multicast observer + `prowl agents signal` | 064 | — | layer-0 signals for every runtime | +| 2 | **S1** signal bus + `ObservedAgentState` multicast observer + `prowl agents signal` (`turn-ended`, needs-input/session/progress, bounded detail) | 064 | — | layer-0 signals for every runtime | | 2 | **065-K2** shared `SymlinkInstaller` + `prowl skills list\|install\|uninstall\|path` | 065 | 065-K1 | one command installs Prowl's skills into agent skill folders | -| 3 | **S2** `prowl agents wait` (`source`/`confidence`, `--include-screen`) + `agents` `signals` field + skill rubric | 064 | S1 | no hand-written polling; heuristic results are labelled | +| 3 | **S2** atomic dispatch pairing (`create` dispatch ID, cooperative `dispatch-complete --detail`, bounded receipt retention, `agents wait --dispatch` with overflow resnapshot) + `source`/`confidence`, `--include-screen`, `agents` `signals`, and skill rubric | 064 | S1 | no hand-written polling or stale completion; heuristic results are labelled | | 3 | **065-K3** Agent Skills section on Settings › Command Line Tool | 065 | 065-K2 | GUI users install skills without a terminal | | 4 | **S3 wave 1** launch-scoped hooks for tier-A runtimes (Claude Code, Codex `notify`, Copilot, Droid, Qoder, Pi, OMP, OpenCode) + self-check | 064 | A2, S1 | `agents wait` is deterministic for Prowl-launched agents | @@ -107,7 +107,10 @@ R3+: V2 / S5 rest; delete HANDOFF_RETIRED stubs ## Change log -- 2026-08-22 — A2 implemented in #714 after C0 #709, A1 #710, and A1b #713. The next R1 +- 2026-08-22 — S1 started on `feat/agent-completion-signal-bus`; owner review moved the + complete dispatch-ID issuance/receipt/wait protocol into S2, renamed the runtime edge to + `turn-ended`, retained bounded detail, and required explicit overflow resnapshot. +- 2026-08-22 — A2 merged in #714 after C0 #709, A1 #710, and A1b #713. The next R1 critical path is 064-S1 → S2 → S3 wave 1; 065-S0/K1 remains independent parallel work. - 2026-08-22 — A1 review: added **A1b** (`PROWL_PANE_ID`) to R1; `create pane` keeps an explicit anchor (no caller-pane default) and a background placement stays with A2. diff --git a/docs-ai/064-agent-completion-signals/000-plan.md b/docs-ai/064-agent-completion-signals/000-plan.md index 0c50cbf3..9b84f187 100644 --- a/docs-ai/064-agent-completion-signals/000-plan.md +++ b/docs-ai/064-agent-completion-signals/000-plan.md @@ -2,9 +2,9 @@ | | | | --- | --- | -| **Status** | Planned (research matrix recorded 2026-08-22) | +| **Status** | In progress — S1 implementation on `feat/agent-completion-signal-bus` | | **Anchor date** | 2026-08-22 | -| **Primary PRs** | TBD | +| **Primary PRs** | S1 TBD | | **Related** | [063 agent-workflows](../063-agent-workflows/000-plan.md) (consumer; defines the `ObservedAgentState` observer this entry feeds), [030 agent-status-detection](../030-agent-status-detection/000-plan.md), [045 native-agent-session-detection](../045-native-agent-session-detection/000-plan.md), [055 agent-profile-runtimes](../055-agent-profile-runtimes/000-plan.md), [059 agent-transcript-snapshots](../059-agent-transcript-snapshots/000-plan.md), [060 cli-targeting-and-contract-governance](../060-prowl-cli-targeting-and-contract-governance/000-plan.md), [#473](https://github.com/onevcat/Prowl/issues/473), [#676](https://github.com/onevcat/Prowl/issues/676), `docs/components/agent-detection.md`, `docs/components/cli.md` | ## Background @@ -38,9 +38,10 @@ at launch and have the agent report to Prowl through the bundled `prowl` binary. exit, OSC progress/notification sequences the CLI emits itself; 3. heuristic screen/process detection (existing). - Add `prowl agents signal ` so any agent (or a hook it runs) can report - `turn-complete` / `needs-input` / `session-start` / `session-end`, attributed by the - caller pane (a hook is a child of the agent process, so process ancestry still resolves - the pane). + `turn-ended` / `needs-input` / `session-start` / `session-end`, attributed by the caller + pane (a hook is a child of the agent process, so process ancestry still resolves the + pane). `turn-ended` deliberately means a runtime turn edge, not assigned-task or workflow + completion. - Add `prowl agents wait --until … [--timeout] [--min-confidence] [--include-screen]` that resolves on the bus and reports *what kind* of signal it got. - Make `prowl agents` honest about what each pane can offer (`signals` field) and make @@ -71,18 +72,22 @@ where ```swift struct AgentSignal: Sendable, Equatable { - enum Kind { case turnComplete, needsInput, sessionStart, sessionEnd, progress(Int?) } - enum Source { case cli, hook(runtime: AgentProfileRuntime, event: String), transcript, process, osc, screen } + enum Kind { case turnEnded, needsInput, sessionStart, sessionEnd, progress(Int?) } + enum Source { + case cooperativeCLI + case hook(runtime: AgentProfileRuntime, event: String) + case transcript, process, osc, screen + } enum Confidence { case exact, high, heuristic } - let kind: Kind; let source: Source; let confidence: Confidence; let at: Date - let sessionID: String?; let detail: String? // e.g. hook payload excerpt; never secrets + let kind: Kind; let source: Source; let confidence: Confidence; let timestamp: Date + let sessionID: String?; let detail: String?; let claimedOrigin: String? } ``` | Producer | Mechanism | Confidence | | --- | --- | --- | -| `prowl agents signal` | CLI handler, caller-pane attribution, optional `--origin hook:.` and `--session ` | exact | -| Launch-scoped hooks | adapter capability `signalHooks` renders the launch flag/config that makes the CLI run ` agents signal --event … --origin hook:…` on its native events (per-runtime syntax: research matrix) | exact | +| `prowl agents signal` | CLI handler, caller-pane attribution, optional bounded `--origin` (claimed metadata only), `--session`, and `--detail` | exact caller/channel attribution | +| Launch-scoped hooks | adapter capability `signalHooks` renders a Prowl-configured launch-scoped channel that reports native events through the bundled CLI; only validated capability upgrades provenance to `hook` (per-runtime syntax: research matrix) | exact channel attribution | | Transcript turn-end | 059's reader on the exact/high-attributed transcript, file-watch instead of polling | high/exact | | Process exit | existing `agentEntryRemoved` | exact | | OSC | existing progress/notification OSC handling in the Ghostty bridge, surfaced as signals | high | @@ -91,6 +96,9 @@ struct AgentSignal: Sendable, Equatable { Every producer writes to the same per-surface state; the reducer-side consumer (063 runner via `AppFeature`) and the CLI-side consumer (`agents wait` via the multicast observer) see identical events. Registration and snapshot capture stay one main-actor step. +Each subscriber is independently bounded. If it falls behind, state churn is recovered from +a new snapshot; signal or lifecycle overflow is explicit and S2's waiter re-subscribes before +surfacing an error. Critical events are never silently discarded. ### `prowl agents wait` @@ -110,6 +118,10 @@ prowl agents wait --until idle|blocked|changed|exit [--timeout 1…600] `--include-screen N`, a stable `detection`-source screen tail and, when available, the 059 result state — everything an orchestrating agent needs to judge a heuristic result in one call. +- Prowl-dispatched work uses an opaque `dispatch_id`, not timestamps, to exclude stale + completion. S2 ships `create` issuance, `dispatch-complete --detail`, bounded receipt + retention, and `agents wait --dispatch` atomically. Receipts survive pane closure but not + app restart; surface generation is only the unpaired fallback. - `removed` / `surfaceClosed` → `AGENT_GONE` (unless `--until exit`); timeout → `WAIT_TIMEOUT` with the last known status/source. The 600 s cap matches typical agent tool timeouts; the skill documents "re-arm on timeout". @@ -153,12 +165,12 @@ interleaves with 063's slices, is owned by the shared living | Slice | Depends | Contents / expectation | | --- | --- | --- | -| **S1** | — | Signal bus state + the `ObservedAgentState` multicast observer (snapshot / changed / removed / surfaceClosed / `.signal`; first specified in 063, delivered here so it ships first) + `prowl agents signal` (CLI four layers). Layer 0 works for every runtime immediately; 063-B3 later consumes the same observer. | -| **S2** | S1 | `prowl agents wait` + `agents` `signals` field + `--include-screen` + skill rubric. Route B usable; heuristic fallback honest. | -| **S3 wave 1** | 063-A2, S1, research matrix | Launch-scoped hook injection (adapter `signalHooks`, self-check) for tier A of the research matrix (flag/env per launch, live-verified): Claude Code `--settings`, Codex `-c notify=[…]` (turn-complete only; hook trust bypass is never passed), Copilot `--plugin-dir`, Droid `--settings`, Qoder `--settings`, Pi `-e`, OMP `--hook`, OpenCode `OPENCODE_CONFIG_CONTENT`. `agents wait` becomes deterministic for Prowl-launched agents on these runtimes. | +| **S1** | — | Signal bus state + the `ObservedAgentState` multicast observer (snapshot / changed / removed / surfaceClosed / `.signal`; first specified in 063, delivered here so it ships first) + `prowl agents signal` for `turn-ended`, `needs-input`, session, and progress events (CLI four layers, bounded detail). Layer 0 works for every runtime immediately; 063-B3 later consumes the same observer. | +| **S2** | S1 | One atomic paired-dispatch path: `create --profile --prompt` returns `dispatch_id`; cooperative `dispatch-complete --detail`; bounded non-destructive in-memory receipts; `prowl agents wait --dispatch` with automatic overflow resnapshot; `agents` `signals` field; `--include-screen`; skill rubric. Route B usable; heuristic fallback honest. | +| **S3 wave 1** | 063-A2, S1, research matrix | Launch-scoped hook injection (adapter `signalHooks`, self-check) for tier A of the research matrix (flag/env per launch, live-verified): Claude Code `--settings`, Codex `-c notify=[…]` (native `agent-turn-complete` maps to `turn-ended`; hook trust bypass is never passed), Copilot `--plugin-dir`, Droid `--settings`, Qoder `--settings`, Pi `-e`, OMP `--hook`, OpenCode `OPENCODE_CONFIG_CONTENT`. `agents wait` becomes deterministic for Prowl-launched agents on these runtimes. | | **S3 wave 2** | S3 wave 1, 053 dedicated homes | Tier B (`configDirOnly`: Gemini, Qwen, Grok, Cline, Kimi) for dedicated-home profiles only; tier C (Cursor, Amp: project files) is not attached. | | **S4** | S1 | Transcript file-watch and OSC producers — layer 2 without hooks. | -| **S5** | 063 C1 (part), S3/S4 + 063 V2 (rest) | 063's watchdog consumes exact signals (nudge on `turn-complete` without `done`, immediate attention on `needs-input`) — ships with 063-D2; later: 063 V2 observe mode (`expect.status` + `agents read` / hook `last_assistant_message`) and `on_attention: ask `. Recorded in 063 amendments. | +| **S5** | 063 C1 (part), S3/S4 + 063 V2 (rest) | 063's watchdog consumes exact signals (nudge on `turn-ended` without `done`, immediate attention on `needs-input`) — ships with 063-D2; later: 063 V2 observe mode (`expect.status` + `agents read` / hook `last_assistant_message`) and `on_attention: ask `. Recorded in 063 amendments. | ### Verification @@ -195,8 +207,8 @@ opencode; partial for qodercli/qwen/amp; docs/bundle for the rest). Key conclusi (tier A above); five more only through a Prowl-owned home (tier B, i.e. dedicated-home profiles); Cursor Agent and Amp only via project files (not attached). - Codex's hook system is trust-gated per command hash; per-launch `-c hooks.*` needs - `--dangerously-bypass-hook-trust`, which Prowl will **not** pass. Codex gets - `turn-complete` through the ungated `notify` config; its permission prompts stay + `--dangerously-bypass-hook-trust`, which Prowl will **not** pass. Codex gets the native + `agent-turn-complete` event (mapped to `turn-ended`) through ungated `notify`; its permission prompts stay heuristic/transcript-based. - Claude Code holds all hooks in interactive sessions until the workspace-trust dialog is accepted — the self-check grace must tolerate that, and a trust prompt is itself a @@ -220,6 +232,12 @@ opencode; partial for qodercli/qwen/amp; docs/bundle for the rest). Key conclusi ## Amendments +- Updated 2026-08-22 before S1 implementation: owner review separated runtime `turn-ended` + from S2's cooperative `dispatch-complete`; retained bounded `--detail`; made public origin + claimed metadata only; required explicit overflow/resnapshot; and moved the complete + dispatch-ID issuance/receipt/wait protocol into S2 so it cannot ship half-paired. The + implementation record is [001-action.md](001-action.md), with the authorized execution + checklist in [002-s1-work-note.md](002-s1-work-note.md). - Updated 2026-08-22: prerequisite 063-A2 is implemented in PR #714. Its typed synchronous launch boundary preserves the adapter-rendered invocation and launch-scoped surface environment that S3 will extend; A2 intentionally injects no hooks. With that dependency diff --git a/docs-ai/064-agent-completion-signals/001-action.md b/docs-ai/064-agent-completion-signals/001-action.md new file mode 100644 index 00000000..0e8232dd --- /dev/null +++ b/docs-ai/064-agent-completion-signals/001-action.md @@ -0,0 +1,83 @@ +# 064 — Agent Completion Signals: S1 Action + +## Status + +Implemented on `feat/agent-completion-signal-bus`; adversarial review and final post-review E2E are in progress. + +## Slice objective + +Ship the terminal-owned signal foundation as one vertical slice: + +- a per-surface canonical observation record; +- an independent typed multicast observer for agent entry, signal, and surface lifecycle changes; +- cooperative `prowl agents signal` ingress attributed to the caller pane; +- the governed CLI parser, wire contract, router/handler, schema, renderer, tests, manuals, and skill updates. + +S1 does not add `agents wait`, launch-scoped runtime hooks, workflow completion, or dispatch receipts. + +## Owner decisions fixed before implementation + +- Continue the 064 path before the 063 workflow runner. `prowl workflow done` remains the only command that completes a workflow step; agent signals are observation/control-plane evidence. +- Rename the runtime edge from ambiguous `turn-complete` to `turn-ended`. A runtime hook can prove that a turn ended, not that an assigned task completed. +- Reserve `dispatch-complete` for S2's paired dispatch protocol. S2 must ship the entire path atomically: `create` returns a `dispatch_id`, the agent reports `dispatch-complete --detail`, a bounded in-memory receipt survives pane closure (but not app restart), and `agents wait --dispatch` consumes it without destructive read semantics. +- Keep bounded `--detail` in S1 so a cooperative producer can attach a short result/reason without forcing another CLI command. Detail is caller-authored metadata, never logged, never raises confidence, and is not a replacement for large transcript/workflow output channels. +- Public `--origin` is only a claimed origin. It cannot mark a channel as verified, raise confidence, or satisfy future hook self-checks. S3 may upgrade provenance only through a Prowl-configured launch-scoped capability; this is a correctness boundary, not a heavyweight security boundary. +- Each observer has bounded buffering. State churn may be recovered by a newer snapshot. Signal/lifecycle loss is never silent: overflow terminates with an explicit internal error, and S2's waiter must re-subscribe/resnapshot before exposing a failure. +- S1 signal state is in-memory and surface-scoped. S2 owns the separate bounded dispatch-receipt retention needed to survive surface closure. + +## Implementation plan + +1. Add `AgentSignal`, `AgentObservationSnapshot`, `ObservedAgentState`, and an explicit observer overflow error. +2. Add terminal-manager-owned per-surface observation records and multicast subscription APIs without reusing the single-consumer `TerminalClient.events()` stream. +3. Wire published agent entry changes/removals and every surface teardown path into the observer. Remove false/duplicate agent-removal emissions while preserving the existing reducer event stream. +4. Add context-aware `agents.signal` handling using kernel peer PID → process ancestry → live shell-PID ownership. Never fall back to focus or `PROWL_PANE_ID`. +5. Add the CLI four layers and strict text/JSON/schema contracts for `turn-ended`, `needs-input`, `session-start`, `session-end`, and `progress`. +6. Update current manuals, 063/064 living plans, and the bundled `prowl-cli` skill. +7. Run focused RED/GREEN tests, complete app/CLI validation, live Debug socket checks, then adversarial review for at least three rounds. + +## Delivered implementation + +- Added terminal-owned `AgentObservationStore` records with atomic replay snapshots, + independent bounded subscribers, explicit overflow, latest-signal replay, cancellation + cleanup, and `removed` / `surfaceClosed` separation. +- Routed the existing detector's published entries into the store while correcting cleanup + so a warmed shell that never published an agent no longer emits a false removal. + Pane/tab/prune/close-all and split rollback converge on exactly-once surface cleanup. +- Added `AgentSignalCommandHandler` and the complete `agents.signal` CLI surface. The app + resolves only the socket peer's process ancestry; normal shell panes can signal without a + detected agent, while outside callers fail with `SOURCE_REQUIRED`. +- Added shared parser/handler validation, fractional ISO-8601 timestamps, strict wire payload, + text receipt, executable schema, raw socket fixtures, manuals, 063/064 plan amendments, and + `prowl-cli` skill guidance. + +## Verification record + +Observed passing before the first commit: + +```bash +make build-cli +make test-cli-smoke +make test-cli-integration # 87 tests +make check # format, strict swift-format lint, SwiftLint, 27 script tests +make test-app # xcresult verified: 2423 tests, zero failures +make build-app +``` + +Focused RED/GREEN coverage includes 9 observer/lifecycle tests, 5 handler validation tests, +parser/envelope/router/renderer/schema tests, caller ancestry, kernel `LOCAL_PEERPID`, app +composition, and an actual Unix socket request through `CLISocketServer` into the signal +handler. The full test run initially hung because that new socket test used blocking file-I/O +inside a detached Swift task that the generic executor scheduled on the OS main thread; a +sample identified the stack, the test client moved to a dedicated dispatch queue, and the +2423-test run then completed. + +An isolated Debug app was also launched on a custom socket. The bundled CLI successfully +exercised `list`, `agents`, `create`, and outside-caller rejection (`SOURCE_REQUIRED`). A +second app instance did not materialize Ghostty child surfaces (blank pane, no command/read), +so an inside-pane GUI signal could not be attested in that environment; the real Unix-socket +context test and app-composition observer test cover that path automatically. Final +post-review validation will repeat the Debug build and basic socket/E2E checks. + +## Review record + +To be completed after the required adversarial Claude Code rounds. diff --git a/docs-ai/064-agent-completion-signals/002-s1-work-note.md b/docs-ai/064-agent-completion-signals/002-s1-work-note.md new file mode 100644 index 00000000..e7095311 --- /dev/null +++ b/docs-ai/064-agent-completion-signals/002-s1-work-note.md @@ -0,0 +1,66 @@ +# 064-S1 — Implementation Work Note + +This note tracks the authorized S1 execution. Durable design decisions belong in +`000-plan.md`; delivered behavior and evidence are finalized in `001-action.md`. + +## Starting point + +- Branch: `feat/agent-completion-signal-bus` +- Base: `origin/main` at `cc800c6f` +- PR #714: merged; A2 is available for later S3 launch-hook injection +- Worktree: clean at branch creation +- Protected ignored path: `scripts/__pycache__/` (do not touch) +- Point-Free/TCA/Observation skill: not available in the current environment + +## Confirmed seams + +- `WorktreeTerminalManager.eventStream()` is strictly single-consumer; `AppFeature` owns its production subscription. S1 must use a separate multicast path. +- `WorktreeTerminalState` produces consumer-visible `ActiveAgentEntry` changes. `ActiveAgentsFeature` is a reducer projection and must not become observer truth. +- Existing agent cleanup emits removal unconditionally in some paths, including warmed panes that never published an agent entry. S1 must make `removed` mean a previously published entry actually disappeared. +- Surface teardown spans pane close, tab close, worktree prune/close-all, and split insertion rollback. Every path must produce at most one `removed` followed by `surfaceClosed`, then finish subscribers. +- `CLISocketServer` supplies a same-UID kernel peer PID in `CLICommandContext`; `CallerPaneResolver` maps its ancestry to live shell PIDs without focus/env fallback. +- CLI additions must update parser, shared input/payload, command envelope, router/handler/app wiring, executable schema, socket fixtures, text renderer, manuals, and skill in one change. + +## RED checklist + +- [x] Domain/wire validation: event/progress combinations, bounded session/origin/detail, NUL/control rejection. +- [x] Observer: snapshot first, normal shell snapshot, two concurrent subscribers, signal multicast, cancellation, overflow, removal without close, close ordering, already-closed surface. +- [x] Lifecycle regressions: never-published agent does not emit removal; pane/tab/prune/close-all/rollback converge on guarded cleanup. +- [x] Caller attribution: missing peer PID, ancestry miss, exact pane match, no focus fallback. +- [x] CLI: parsing, envelope coding, router context, handler success/failures, renderer, raw socket/schema round trip. + +## GREEN / validation checklist + +- [x] Focused app tests through `xcodebuild ... | xcsift -f toon`. +- [x] `make build-cli` +- [x] `make test-cli-smoke` +- [x] `make test-cli-integration` (87 tests) +- [x] `make check` +- [x] `make test` / `make test-app` (2423 xcresult tests, zero failures) +- [x] `make build-app` +- [x] Launch isolated Debug app/socket and exercise bundled CLI `list`, `agents`, `create`, plus outside-caller `SOURCE_REQUIRED`. +- [ ] Live caller-pane success from inside a second Debug Prowl pane (Ghostty child surfaces did not materialize in the concurrent-instance environment; covered by real socket + app composition tests). +- [ ] At least three Claude Code adversarial review rounds; fix every credible P0/P1 with tests, commit/push each accepted round, and record PR comments. +- [ ] Final Debug build and basic E2E after review. + +## Deferred S2 contract (must not be lost) + +S2 owns one atomic paired-dispatch slice: + +```text +create --profile --prompt + → return opaque dispatch_id + → cooperative dispatch-complete --detail + → bounded non-destructive in-memory receipt keyed by dispatch_id + → agents wait --dispatch (automatic overflow resnapshot) +``` + +The dispatch receipt survives agent/pane closure but not app restart. A later dispatch cannot +be satisfied by an older receipt. Internal surface generations remain only an unpaired +observation fallback. Full workflow output continues through `prowl workflow done -`; large +ad-hoc results continue through `agents read` or a future `agents wait --include-result`. + +## Progress log + +- 2026-08-22 — Preflight investigation and owner grill completed; branch created and the decisions above recorded before tests/code. +- 2026-08-23 — Observer and CLI RED/GREEN phases completed; full contract/manual/skill updates landed locally. CLI integration passed 87 tests and app xcresult passed 2423 tests. A full-suite deadlock in the newly added blocking socket test was sampled and fixed by moving test-client I/O to a dispatch queue. Initial isolated Debug socket validation completed with the second-instance Ghostty limitation recorded above. diff --git a/docs/components/agent-detection.md b/docs/components/agent-detection.md index a7afab95..02f3a6ae 100644 --- a/docs/components/agent-detection.md +++ b/docs/components/agent-detection.md @@ -104,6 +104,28 @@ showing their chrome keep the last trusted state instead of forcing idle. A **Done** pane becomes **Idle** the moment you focus it. +## Cooperative signal bus + +Detection remains heuristic UI state. Separately, every live pane now has a terminal-owned +multicast observation stream. A process inside the pane can report explicit runtime events: + +```bash +prowl agents signal turn-ended --detail "Review complete" +prowl agents signal needs-input +``` + +Prowl attributes the socket caller through process ancestry, not focus or +`PROWL_PANE_ID`. A signal can exist for an ordinary shell pane with no detected agent and +does not create or overwrite a detected-agent entry. `turn-ended` means one interaction +ended, not that a workflow step completed; only `prowl workflow done` will advance a +workflow. Future `agents wait` and launch-scoped hooks consume the same per-pane stream. + +Each observer receives an atomic current snapshot before live changes. Multiple observers +do not compete with the app's existing single-consumer terminal event stream. Agent removal +keeps the pane stream alive; pane closure emits `surfaceClosed` and finishes it. Bounded +overflow is explicit so future waiters can re-subscribe and resnapshot rather than silently +lose lifecycle or signal evidence. + ## How often it runs - **No polling** for cold panes that have not received recent input. diff --git a/docs/components/cli.md b/docs/components/cli.md index d9d5a80a..9c73c58e 100644 --- a/docs/components/cli.md +++ b/docs/components/cli.md @@ -4,7 +4,7 @@ > an agent) can list panes, read their screens, run commands and capture output, > send keystrokes, focus, and open/close tabs and panes programmatically. -**Keywords:** prowl cli, command line, prowl list, prowl agents, prowl agents read, prowl profiles list, prowl read, prowl send, prowl key, prowl focus, prowl create, prowl close, prowl open, prowl handoff, pane id, agent, profile, automation, json, capture, socket +**Keywords:** prowl cli, command line, prowl list, prowl agents, prowl agents read, prowl agents signal, prowl profiles list, prowl read, prowl send, prowl key, prowl focus, prowl create, prowl close, prowl open, prowl handoff, pane id, agent, profile, automation, json, capture, socket **Related:** [terminal](terminal.md) · [concepts](../concepts.md) · [active-agents](active-agents.md) · [agent-detection](agent-detection.md) · the bundled **`prowl-cli` skill** (`skills/prowl-cli/SKILL.md`) @@ -106,8 +106,8 @@ step that depends on knowing yourself inside the success branch. When it is unse matches nothing, stop rather than guess: `pane.cwd` only narrows the candidates — several panes usually share one cwd — and may stand in for you only when the match is unique; never assume the *focused* pane is you. Prowl itself never trusts the variable for -attribution; commands that need the calling pane (`handoff`) resolve it from the -caller's process ancestry. +attribution; commands that need the calling pane (`handoff`, `agents signal`) resolve it +from the caller's process ancestry. ## Commands @@ -227,6 +227,32 @@ no heading or added newline. For every other result state it exits non-zero with `SESSION_UNRESOLVED`, `RESULT_NOT_FOUND`, `RESULT_INCOMPLETE`, or `RESULT_TOO_LARGE`. Empty `agents` roster output remains `No agents found.`. +### `prowl agents signal ` + +Report a cooperative event for the agent in the **calling pane**: + +```bash +prowl agents signal turn-ended --detail "Review complete" +prowl agents signal needs-input --session session-1 --json +prowl agents signal progress --progress 75 +``` + +Events are `turn-ended`, `needs-input`, `session-start`, `session-end`, and `progress`. +`turn-ended` means one runtime interaction ended; it does not complete a workflow or prove +an assigned task is done. `--progress` accepts 0–100 and is valid only for `progress`. +Optional `--session` and claimed `--origin` are limited to 256 UTF-8 bytes; `--detail` +carries a short result/reason up to 4096 bytes. Values must be non-empty and control-free. + +The command accepts no target: Prowl attributes the kernel socket peer PID through process +ancestry to a live pane. It never uses UI focus or `PROWL_PANE_ID`; external terminals, +tmux/detached ancestry, and already-closed panes fail with `SOURCE_REQUIRED` or +`AGENT_GONE`. Public signals report `source=cooperative_cli`, `confidence=exact`; exact +means explicit channel and caller-pane attribution, not verified business completion. +Claimed origin never upgrades trust. JSON uses `prowl.cli.agents.signal.v1`. + +`dispatch-complete` and `agents wait` arrive together in the next slice with an opaque +`dispatch_id`; `workflow done` remains the only workflow-step completion command. + ### `prowl profiles list` Read-only snapshot of every configured Agent Profile, including disabled profiles, in @@ -517,6 +543,8 @@ artifacts and terminal excerpts do not appear in `git status`. | `PROFILE_NOT_FOUND` | No enabled Profile matches the UUID or exact name — re-run `profiles list`; disabled Profiles cannot launch. | | `PROFILE_NOT_UNIQUE` | Several enabled Profiles have the exact name — use the Profile UUID from `profiles list`. | | `AGENT_NOT_FOUND` / `AGENT_UNSUPPORTED` | `agents read` target no longer hosts an agent, or it is not Codex/Claude Code. Re-run `agents`. | +| `SOURCE_REQUIRED` | A caller-owned command such as `agents signal` or selector-free `handoff` could not map the socket peer ancestry to a Prowl pane. Run it inside the source pane without tmux/detached wrappers, or use an explicit selector where that command permits one. | +| `AGENT_GONE` | The caller pane closed before `agents signal` could record its event. | | `BLOCKER_UNREADABLE` | A blocked screen was detected but Prowl could not safely extract its current interaction text. Re-run `agents read` or inspect with `read`. | | `SESSION_UNRESOLVED` / `RESULT_NOT_FOUND` / `RESULT_INCOMPLETE` / `RESULT_TOO_LARGE` | `agents read --result-only` could not provide one trustworthy complete result. Drop `--result-only` to retain the live snapshot and inspect `.data.result`. | | `NO_ACTIVE_PANE` | No pane for focused-target; pass an explicit `--pane`. | @@ -561,6 +589,8 @@ fi - Use `prowl agents --json` for discovery, then `prowl agents read ` for a supported agent's status, blocker, and trustworthy result state; use `prowl list --json` when you need all panes, including ordinary shells. +- Use `prowl agents signal` only from the pane reporting the event; it never targets focus + or another pane and never substitutes for `workflow done`. - `--capture` needs shell integration; otherwise `read --wait-stable` or file redirection. - `open` is navigation, not a guaranteed new pane — use `create tab` or `create pane`. diff --git a/skills/prowl-cli/SKILL.md b/skills/prowl-cli/SKILL.md index 2ae1bd87..02854d0d 100644 --- a/skills/prowl-cli/SKILL.md +++ b/skills/prowl-cli/SKILL.md @@ -64,6 +64,15 @@ Pick targets by `pane.id`, `tab.id`, `worktree.id`/`name`/`path`, and `pane.cwd` For a currently active Codex or Claude Code agent, `prowl agents read p7 --json` returns an immediate semantic snapshot: `.data.agent.status`, `.data.blocker.text` when blocked, and `.data.result` — a result is trustworthy only when `.data.result.state == "complete"`. +Report an event only for the pane running the current process; `agents signal` has no target or focus fallback: + +```bash +prowl agents signal turn-ended --detail "Review complete" --json +prowl agents signal needs-input --session session-1 --json +``` + +The app attributes the socket peer PID through process ancestry. `turn-ended` means a runtime turn edge, not task/workflow completion; only `workflow done` completes a workflow step. Public `--origin` is claimed metadata and never upgrades trust. + ## Common Recipes Open a split beside yourself (or any positively identified anchor) and capture the new pane: @@ -138,6 +147,7 @@ Key fields by command: - `list` → `.data.items[]` with `.worktree.{id,name,path,root_path,kind}`, `.tab.{id,title,selected}`, `.pane.{id,title,cwd,focused,agent}`, `.task.status` (`running`|`idle`|null). - `agents` → `.data.agents[]` with `.status`, `.raw_state`, `.detection_reason`, `.type`, `.name`, `.pane.{id,focused,cwd}`, `.tab`, `.worktree`, `.project.{name,branch,path}`. - `agents read` → `.data.agent`, `.data.blocker.text`, `.data.result.{state,text}` — `pending`, `unavailable`, `missing`, `incomplete`, `too_large` carry no partial text. +- `agents signal` → `.data.pane.{id,worktree_id}`, `.data.signal.{event,source,confidence,at,session_id,detail,claimed_origin}`; optional fields are omitted. - `read` → `.data.text`, `.data.line_count`, `.data.truncated`, `.data.mode`, `.data.source`; `.data.stabilized` / `.data.waited_ms` with `--wait-stable`. - `send` → `.data.input`, `.data.wait.{exit_code,duration_ms}` when waiting, `.data.capture.{text,line_count,truncated}` with `--capture`. - `create tab` / `open` → `.data.target.{pane,tab,worktree}`; `create pane` → `.data.anchor`, `.data.direction`, `.data.target`; Profile launches also include `.data.launch.{profile_id,profile_name,agent}`. @@ -185,7 +195,7 @@ done - `TRANSPORT_FAILED`: the connection broke or the socket path is invalid (`ENOTSOCK`, too-long `PROWL_CLI_SOCKET`). - `TARGET_NOT_FOUND` / `TARGET_NOT_UNIQUE`: re-run `prowl list --json` and pass an explicit UUID or a current `pN`. - `PROFILE_NOT_FOUND` / `PROFILE_NOT_UNIQUE`: re-run `prowl profiles list --json`; choose an enabled Profile UUID. -- `NO_ACTIVE_PANE`: focused-pane targeting found nothing — pass `--pane`. `SOURCE_REQUIRED`: `handoff` was run outside a Prowl pane without a selector. +- `NO_ACTIVE_PANE`: focused-pane targeting found nothing — pass `--pane`. `SOURCE_REQUIRED`: a caller-owned command (`agents signal`, selector-free `handoff`) could not map process ancestry to a Prowl pane. `AGENT_GONE`: the signal's caller pane closed before recording. - `EMPTY_INPUT`, `INVALID_ARGUMENT`, `UNSUPPORTED_KEY`, `INVALID_REPEAT`: fix the arguments (`prowl --help`). - `CAPTURE_UNSUPPORTED`: drop `--capture` and use `read --wait-stable` or file redirection. `WAIT_TIMEOUT`: raise `--timeout` or use `--no-wait`. - `PATH_NOT_FOUND` / `PATH_NOT_DIRECTORY` / `PATH_NOT_ALLOWED`: fix the path given to `open` or `create tab --path`. @@ -210,4 +220,4 @@ Required sections are `## Objective`, `## Current State`, and `## Next Steps`; o ## Command Set -`list`, `agents`, `agents read`, `profiles list`, `read`, `send`, `key`, `focus`, `create tab`, `create pane`, `close`, `handoff to`, `handoff save`, and `open` (default). There is no CLI `quit`; close temporary tabs or panes with an explicit `close`. `tab create`, `tab close`, and `pane close` remain deprecated aliases for one release. +`list`, `agents`, `agents read`, `agents signal`, `profiles list`, `read`, `send`, `key`, `focus`, `create tab`, `create pane`, `close`, `handoff to`, `handoff save`, and `open` (default). There is no CLI `quit`; close temporary tabs or panes with an explicit `close`. `tab create`, `tab close`, and `pane close` remain deprecated aliases for one release. diff --git a/supacode/App/supacodeApp.swift b/supacode/App/supacodeApp.swift index eb8c8a51..f1190ee8 100644 --- a/supacode/App/supacodeApp.swift +++ b/supacode/App/supacodeApp.swift @@ -358,6 +358,9 @@ struct SupacodeApp: App { events: { terminalManager.eventStream() }, + observeAgentState: { surfaceID in + terminalManager.observeAgentState(surfaceID: surfaceID) + }, canvasFocusedWorktreeID: { terminalManager.canvasFocusedWorktreeID }, @@ -667,7 +670,8 @@ struct SupacodeApp: App { static func makeCLICommandRouter( appStore: StoreOf, terminalManager: WorktreeTerminalManager, - handoffRequestRegistry: HandoffRequestRegistry = HandoffRequestRegistry() + handoffRequestRegistry: HandoffRequestRegistry = HandoffRequestRegistry(), + agentSignalCallerResolver: AgentSignalCommandHandler.ResolveCaller? = nil ) -> CLICommandRouter { let listHandler = ListCommandHandler { @@ -700,6 +704,19 @@ struct SupacodeApp: App { screenDetectionsBySurfaceID: screenDetectionsBySurfaceID ) } + let resolveAgentSignalCaller: AgentSignalCommandHandler.ResolveCaller = + agentSignalCallerResolver ?? { callerProcessID in + CallerPaneResolver.pane( + forCallerProcess: callerProcessID, + paneByShellPID: terminalManager.paneByShellPID() + ) + } + let agentSignalHandler = AgentSignalCommandHandler( + resolveCaller: resolveAgentSignalCaller, + recordSignal: { caller, signal in + terminalManager.recordAgentSignal(signal, surfaceID: caller.surfaceID) + } + ) let agentReadHandler = AgentReadCommandHandler( snapshotProvider: { pane in await Self.makeAgentReadRuntimeSnapshot( @@ -987,6 +1004,7 @@ struct SupacodeApp: App { listHandler: listHandler, agentsHandler: agentsHandler, agentsReadHandler: agentReadHandler, + agentsSignalHandler: agentSignalHandler, profilesHandler: profilesHandler, focusHandler: focusHandler, sendHandler: sendHandler, diff --git a/supacode/CLIService/AgentSignalCommandHandler.swift b/supacode/CLIService/AgentSignalCommandHandler.swift new file mode 100644 index 00000000..bf3d59d6 --- /dev/null +++ b/supacode/CLIService/AgentSignalCommandHandler.swift @@ -0,0 +1,113 @@ +import Foundation + +@MainActor +final class AgentSignalCommandHandler: CommandHandler { + typealias ResolveCaller = @MainActor (pid_t) -> CallerPane? + typealias RecordSignal = @MainActor (CallerPane, AgentSignal) -> Bool + + private let resolveCaller: ResolveCaller + private let recordSignal: RecordSignal + private let now: @Sendable () -> Date + private let dateFormatter: ISO8601DateFormatter + + init( + resolveCaller: @escaping ResolveCaller, + recordSignal: @escaping RecordSignal, + now: @escaping @Sendable () -> Date = Date.init + ) { + self.resolveCaller = resolveCaller + self.recordSignal = recordSignal + self.now = now + self.dateFormatter = ISO8601DateFormatter() + self.dateFormatter.formatOptions = [.withInternetDateTime, .withFractionalSeconds] + self.dateFormatter.timeZone = TimeZone(secondsFromGMT: 0) + } + + func handle(envelope: CommandEnvelope) async -> CommandResponse { + await handle(envelope: envelope, context: CLICommandContext()) + } + + // swiftlint:disable async_without_await + func handle( + envelope: CommandEnvelope, + context: CLICommandContext + ) async -> CommandResponse { + guard case .agentsSignal(let input) = envelope.command else { + return failure(code: CLIErrorCode.invalidArgument, message: "Expected an agents.signal command.") + } + if let validationMessage = input.validationErrorMessage { + return failure(code: CLIErrorCode.invalidArgument, message: validationMessage) + } + guard let processID = context.callerProcessID, + let caller = resolveCaller(processID) + else { + return failure( + code: CLIErrorCode.sourceRequired, + message: "Run 'prowl agents signal' from inside the Prowl pane that is reporting the event." + ) + } + + let signal = AgentSignal( + kind: kind(for: input), + source: .cooperativeCLI, + confidence: .exact, + timestamp: now(), + sessionID: input.sessionID, + detail: input.detail, + claimedOrigin: input.origin + ) + guard recordSignal(caller, signal) else { + return failure( + code: CLIErrorCode.agentGone, + message: "The caller pane closed before its agent signal could be recorded." + ) + } + + let payload = AgentSignalCommandPayload( + pane: AgentSignalPanePayload( + id: caller.surfaceID.uuidString, + worktreeID: caller.worktreeID + ), + signal: AgentSignalPayload( + event: input.event, + progress: input.progress, + source: "cooperative_cli", + confidence: signal.confidence.rawValue, + timestamp: dateFormatter.string(from: signal.timestamp), + sessionID: signal.sessionID, + detail: signal.detail, + claimedOrigin: signal.claimedOrigin + ) + ) + do { + return try CommandResponse( + ok: true, + command: "agents.signal", + schemaVersion: "prowl.cli.agents.signal.v1", + data: RawJSON(encoding: payload) + ) + } catch { + return failure(code: CLIErrorCode.agentsFailed, message: "Failed to encode the agent signal receipt.") + } + } + // swiftlint:enable async_without_await + + private func kind(for input: AgentSignalInput) -> AgentSignal.Kind { + switch input.event { + case .turnEnded: .turnEnded + case .needsInput: .needsInput + case .sessionStart: .sessionStart + case .sessionEnd: .sessionEnd + case .progress: .progress(input.progress) + } + } + + private func failure(code: String, message: String) -> CommandResponse { + CommandResponse( + ok: false, + command: "agents.signal", + schemaVersion: "prowl.cli.agents.signal.v1", + error: CommandError(code: code, message: message) + ) + } +} diff --git a/supacode/CLIService/CLICommandRouter.swift b/supacode/CLIService/CLICommandRouter.swift index ee2385be..bc2534f3 100644 --- a/supacode/CLIService/CLICommandRouter.swift +++ b/supacode/CLIService/CLICommandRouter.swift @@ -9,6 +9,7 @@ final class CLICommandRouter { private let listHandler: any CommandHandler private let agentsHandler: any CommandHandler private let agentsReadHandler: any CommandHandler + private let agentsSignalHandler: any CommandHandler private let profilesHandler: any CommandHandler private let focusHandler: any CommandHandler private let sendHandler: any CommandHandler @@ -25,6 +26,7 @@ final class CLICommandRouter { listHandler: any CommandHandler = StubCommandHandler(command: "list"), agentsHandler: any CommandHandler = StubCommandHandler(command: "agents"), agentsReadHandler: any CommandHandler = StubCommandHandler(command: "agents.read"), + agentsSignalHandler: any CommandHandler = StubCommandHandler(command: "agents.signal"), profilesHandler: any CommandHandler = StubCommandHandler(command: "profiles"), focusHandler: any CommandHandler = StubCommandHandler(command: "focus"), sendHandler: any CommandHandler = StubCommandHandler(command: "send"), @@ -40,6 +42,7 @@ final class CLICommandRouter { self.listHandler = listHandler self.agentsHandler = agentsHandler self.agentsReadHandler = agentsReadHandler + self.agentsSignalHandler = agentsSignalHandler self.profilesHandler = profilesHandler self.focusHandler = focusHandler self.sendHandler = sendHandler @@ -62,6 +65,7 @@ final class CLICommandRouter { case .list: handler = listHandler case .agents: handler = agentsHandler case .agentsRead: handler = agentsReadHandler + case .agentsSignal: handler = agentsSignalHandler case .profiles: handler = profilesHandler case .focus: handler = focusHandler case .send: handler = sendHandler diff --git a/supacode/CLIService/Shared/AgentSignalCommandPayload.swift b/supacode/CLIService/Shared/AgentSignalCommandPayload.swift new file mode 100644 index 00000000..a77440ed --- /dev/null +++ b/supacode/CLIService/Shared/AgentSignalCommandPayload.swift @@ -0,0 +1,68 @@ +import Foundation + +public struct AgentSignalCommandPayload: Codable, Equatable, Sendable { + public let pane: AgentSignalPanePayload + public let signal: AgentSignalPayload + + public init(pane: AgentSignalPanePayload, signal: AgentSignalPayload) { + self.pane = pane + self.signal = signal + } +} + +public struct AgentSignalPanePayload: Codable, Equatable, Sendable { + public let id: String + public let worktreeID: String + + enum CodingKeys: String, CodingKey { + case id + case worktreeID = "worktree_id" + } + + public init(id: String, worktreeID: String) { + self.id = id + self.worktreeID = worktreeID + } +} + +public struct AgentSignalPayload: Codable, Equatable, Sendable { + public let event: AgentSignalEvent + public let progress: Int? + public let source: String + public let confidence: String + public let timestamp: String + public let sessionID: String? + public let detail: String? + public let claimedOrigin: String? + + enum CodingKeys: String, CodingKey { + case event + case progress + case source + case confidence + case timestamp = "at" + case sessionID = "session_id" + case detail + case claimedOrigin = "claimed_origin" + } + + public init( + event: AgentSignalEvent, + progress: Int?, + source: String, + confidence: String, + timestamp: String, + sessionID: String?, + detail: String?, + claimedOrigin: String? + ) { + self.event = event + self.progress = progress + self.source = source + self.confidence = confidence + self.timestamp = timestamp + self.sessionID = sessionID + self.detail = detail + self.claimedOrigin = claimedOrigin + } +} diff --git a/supacode/CLIService/Shared/CommandEnvelope.swift b/supacode/CLIService/Shared/CommandEnvelope.swift index b4829f4b..d73c3ee5 100644 --- a/supacode/CLIService/Shared/CommandEnvelope.swift +++ b/supacode/CLIService/Shared/CommandEnvelope.swift @@ -18,6 +18,7 @@ public enum Command: Codable, Sendable { case list(ListInput) case agents(AgentsInput) case agentsRead(AgentReadInput) + case agentsSignal(AgentSignalInput) case profiles(ProfilesInput) case focus(FocusInput) case send(SendInput) @@ -35,6 +36,7 @@ public enum Command: Codable, Sendable { case .list: "list" case .agents: "agents" case .agentsRead: "agents.read" + case .agentsSignal: "agents.signal" case .profiles: "profiles" case .focus: "focus" case .send: "send" diff --git a/supacode/CLIService/Shared/ErrorCodes.swift b/supacode/CLIService/Shared/ErrorCodes.swift index a4d6ff2b..63b21f73 100644 --- a/supacode/CLIService/Shared/ErrorCodes.swift +++ b/supacode/CLIService/Shared/ErrorCodes.swift @@ -25,6 +25,7 @@ public enum CLIErrorCode { public static let agentNotFound = "AGENT_NOT_FOUND" public static let agentUnsupported = "AGENT_UNSUPPORTED" public static let agentReadFailed = "AGENT_READ_FAILED" + public static let agentGone = "AGENT_GONE" public static let blockerUnreadable = "BLOCKER_UNREADABLE" public static let sessionUnresolved = "SESSION_UNRESOLVED" public static let resultNotFound = "RESULT_NOT_FOUND" diff --git a/supacode/CLIService/Shared/InputModels.swift b/supacode/CLIService/Shared/InputModels.swift index 02b59460..4ec4f2f1 100644 --- a/supacode/CLIService/Shared/InputModels.swift +++ b/supacode/CLIService/Shared/InputModels.swift @@ -30,6 +30,88 @@ public struct AgentsInput: Codable, Sendable { public init() {} } +public enum AgentSignalEvent: String, Codable, CaseIterable, Sendable { + case turnEnded = "turn-ended" + case needsInput = "needs-input" + case sessionStart = "session-start" + case sessionEnd = "session-end" + case progress +} + +public struct AgentSignalInput: Codable, Sendable { + public static let maximumSessionIDBytes = 256 + public static let maximumOriginBytes = 256 + public static let maximumDetailBytes = 4 * 1_024 + + public let event: AgentSignalEvent + public let progress: Int? + public let origin: String? + public let sessionID: String? + public let detail: String? + + enum CodingKeys: String, CodingKey { + case event + case progress + case origin + case sessionID = "session_id" + case detail + } + + public init( + event: AgentSignalEvent, + progress: Int? = nil, + origin: String? = nil, + sessionID: String? = nil, + detail: String? = nil + ) { + self.event = event + self.progress = progress + self.origin = origin + self.sessionID = sessionID + self.detail = detail + } + + public var validationErrorMessage: String? { + if event != .progress, progress != nil { + return "--progress is only valid with the 'progress' event." + } + if let progress, !(0...100).contains(progress) { + return "--progress must be between 0 and 100." + } + if let message = Self.validateText( + sessionID, + name: "--session", + maximumBytes: Self.maximumSessionIDBytes + ) { + return message + } + if let message = Self.validateText( + origin, + name: "--origin", + maximumBytes: Self.maximumOriginBytes + ) { + return message + } + return Self.validateText( + detail, + name: "--detail", + maximumBytes: Self.maximumDetailBytes + ) + } + + private static func validateText(_ value: String?, name: String, maximumBytes: Int) -> String? { + guard let value else { return nil } + guard !value.isEmpty else { return "\(name) must not be empty." } + guard value.utf8.count <= maximumBytes else { + return "\(name) must be at most \(maximumBytes) UTF-8 bytes." + } + guard !value.unicodeScalars.contains(where: CharacterSet.controlCharacters.contains) else { + return "\(name) must not contain control characters." + } + return nil + } +} + public struct ProfilesInput: Codable, Sendable { public init() {} } diff --git a/supacode/Clients/Terminal/TerminalClient.swift b/supacode/Clients/Terminal/TerminalClient.swift index 716e91e7..29f899c7 100644 --- a/supacode/Clients/Terminal/TerminalClient.swift +++ b/supacode/Clients/Terminal/TerminalClient.swift @@ -11,6 +11,8 @@ struct TerminalClient { var launchAgentProfile: @MainActor @Sendable (Worktree, AgentProfileLaunchRequest) -> Result var events: @MainActor @Sendable () -> AsyncStream + /// Per-surface multicast stream. Independent from the single-consumer event stream. + var observeAgentState: @MainActor @Sendable (UUID) -> AgentObservationStream var canvasFocusedWorktreeID: @MainActor @Sendable () -> Worktree.ID? /// Active surface in the selected tab. Lets the reducer capture the target /// synchronously before an async dispatch races against AppKit focus reshuffle @@ -108,6 +110,7 @@ extension TerminalClient: DependencyKey { createTabInDirectory: { _, _ in fatalError("TerminalClient.createTabInDirectory not configured") }, launchAgentProfile: { _, _ in fatalError("TerminalClient.launchAgentProfile not configured") }, events: { fatalError("TerminalClient.events not configured") }, + observeAgentState: { _ in fatalError("TerminalClient.observeAgentState not configured") }, canvasFocusedWorktreeID: { nil }, selectedSurfaceID: { _ in nil }, handoffSourceContext: { _ in nil }, @@ -125,6 +128,7 @@ extension TerminalClient: DependencyKey { createTabInDirectory: { _, _ in nil }, launchAgentProfile: { _, _ in .failure(.tabCreationFailed) }, events: { AsyncStream { $0.finish() } }, + observeAgentState: { _ in AgentObservationStream { $0.finish() } }, canvasFocusedWorktreeID: { nil }, selectedSurfaceID: { _ in nil }, handoffSourceContext: { _ in nil }, diff --git a/supacode/Domain/AgentDetection/AgentSignal.swift b/supacode/Domain/AgentDetection/AgentSignal.swift new file mode 100644 index 00000000..73c8d305 --- /dev/null +++ b/supacode/Domain/AgentDetection/AgentSignal.swift @@ -0,0 +1,57 @@ +import Foundation + +struct AgentSignal: Equatable, Sendable { + enum Kind: Equatable, Sendable { + case turnEnded + case needsInput + case sessionStart + case sessionEnd + case progress(Int?) + } + + enum Source: Equatable, Sendable { + case cooperativeCLI + case hook(runtime: AgentProfileRuntime, event: String) + case transcript + case process + case osc + case screen + } + + enum Confidence: String, Equatable, Sendable { + case exact + case high + case heuristic + } + + let kind: Kind + let source: Source + let confidence: Confidence + let timestamp: Date + let sessionID: String? + let detail: String? + /// Caller-authored provenance hint. It never upgrades `source` or `confidence`. + let claimedOrigin: String? +} + +struct AgentObservationSnapshot: Equatable, Sendable { + let agent: ActiveAgentEntry? + let latestSignal: AgentSignal? + /// Monotonic within one live surface. Consumers use it to recognize a newer + /// resubscription snapshot after an overflow; it is not a persisted cursor. + let revision: UInt64 +} + +enum ObservedAgentState: Equatable, Sendable { + case snapshot(AgentObservationSnapshot) + case changed(ActiveAgentEntry) + case removed + case signal(AgentSignal) + case surfaceClosed +} + +enum AgentObservationError: Error, Equatable, Sendable { + case bufferOverflow +} + +typealias AgentObservationStream = AsyncThrowingStream diff --git a/supacode/Features/Terminal/BusinessLogic/AgentObservationStore.swift b/supacode/Features/Terminal/BusinessLogic/AgentObservationStore.swift new file mode 100644 index 00000000..99dc46bd --- /dev/null +++ b/supacode/Features/Terminal/BusinessLogic/AgentObservationStore.swift @@ -0,0 +1,132 @@ +import Foundation + +/// Terminal-owned publication state for one multicast observer per surface. +/// The detector remains the producer of agent entries; this store is the +/// canonical replay/stream boundary shared by reducers and CLI observers. +@MainActor +final class AgentObservationStore { + private struct SurfaceRecord { + var agent: ActiveAgentEntry? + var latestSignal: AgentSignal? + var revision: UInt64 = 0 + var subscribers: [UUID: AgentObservationStream.Continuation] = [:] + + var snapshot: AgentObservationSnapshot { + AgentObservationSnapshot( + agent: agent, + latestSignal: latestSignal, + revision: revision + ) + } + } + + private var records: [UUID: SurfaceRecord] = [:] + private let bufferCapacity: Int + + init(bufferCapacity: Int) { + self.bufferCapacity = max(1, bufferCapacity) + } + + func observe(surfaceID: UUID, isLive: Bool) -> AgentObservationStream { + guard isLive else { + return AgentObservationStream(bufferingPolicy: .unbounded) { continuation in + continuation.yield( + .snapshot( + AgentObservationSnapshot(agent: nil, latestSignal: nil, revision: 0) + )) + continuation.yield(.surfaceClosed) + continuation.finish() + } + } + + let subscriberID = UUID() + var continuation: AgentObservationStream.Continuation? + let stream = AgentObservationStream( + bufferingPolicy: .bufferingOldest(bufferCapacity) + ) { continuation = $0 } + guard let continuation else { return stream } + + continuation.onTermination = { @Sendable [weak self] _ in + Task { @MainActor [weak self] in + self?.removeSubscriber(subscriberID, surfaceID: surfaceID) + } + } + + var record = records[surfaceID] ?? SurfaceRecord() + record.subscribers[subscriberID] = continuation + let snapshot = record.snapshot + records[surfaceID] = record + continuation.yield(.snapshot(snapshot)) + return stream + } + + func publishAgentChanged(_ entry: ActiveAgentEntry) { + var record = records[entry.surfaceID] ?? SurfaceRecord() + guard record.agent != entry else { return } + record.agent = entry + record.revision &+= 1 + records[entry.surfaceID] = record + publish(.changed(entry), surfaceID: entry.surfaceID) + } + + func publishAgentRemoved(surfaceID: UUID) { + guard var record = records[surfaceID], record.agent != nil else { return } + record.agent = nil + record.revision &+= 1 + records[surfaceID] = record + publish(.removed, surfaceID: surfaceID) + } + + func publishSignal(_ signal: AgentSignal, surfaceID: UUID) { + var record = records[surfaceID] ?? SurfaceRecord() + record.latestSignal = signal + record.revision &+= 1 + records[surfaceID] = record + publish(.signal(signal), surfaceID: surfaceID) + } + + func publishSurfaceClosed(surfaceID: UUID) { + guard records[surfaceID] != nil else { return } + publishAgentRemoved(surfaceID: surfaceID) + guard let record = records.removeValue(forKey: surfaceID) else { return } + for continuation in record.subscribers.values { + switch continuation.yield(.surfaceClosed) { + case .enqueued: + continuation.finish() + case .dropped, .terminated: + continuation.finish(throwing: AgentObservationError.bufferOverflow) + @unknown default: + continuation.finish(throwing: AgentObservationError.bufferOverflow) + } + } + } + + func subscriberCount(surfaceID: UUID) -> Int { + records[surfaceID]?.subscribers.count ?? 0 + } + + private func publish(_ event: ObservedAgentState, surfaceID: UUID) { + guard var record = records[surfaceID] else { return } + for (subscriberID, continuation) in record.subscribers { + switch continuation.yield(event) { + case .enqueued: + continue + case .dropped: + continuation.finish(throwing: AgentObservationError.bufferOverflow) + record.subscribers.removeValue(forKey: subscriberID) + case .terminated: + record.subscribers.removeValue(forKey: subscriberID) + @unknown default: + continuation.finish(throwing: AgentObservationError.bufferOverflow) + record.subscribers.removeValue(forKey: subscriberID) + } + } + records[surfaceID] = record + } + + private func removeSubscriber(_ subscriberID: UUID, surfaceID: UUID) { + guard var record = records[surfaceID] else { return } + record.subscribers.removeValue(forKey: subscriberID) + records[surfaceID] = record + } +} diff --git a/supacode/Features/Terminal/BusinessLogic/WorktreeTerminalManager.swift b/supacode/Features/Terminal/BusinessLogic/WorktreeTerminalManager.swift index 7530c971..ed5fa0c3 100644 --- a/supacode/Features/Terminal/BusinessLogic/WorktreeTerminalManager.swift +++ b/supacode/Features/Terminal/BusinessLogic/WorktreeTerminalManager.swift @@ -12,6 +12,7 @@ final class WorktreeTerminalManager { private let runtime: GhosttyRuntime? private let layoutPersistence: TerminalLayoutPersistenceClient private let targetHandleRegistry = TerminalTargetHandleRegistry() + @ObservationIgnored private let agentObservationStore: AgentObservationStore private var states: [Worktree.ID: WorktreeTerminalState] = [:] private var notificationsEnabled = true private var commandFinishedNotificationEnabled = true @@ -34,11 +35,13 @@ final class WorktreeTerminalManager { init( runtime: GhosttyRuntime, preferredFontSize: Float32? = nil, - layoutPersistence: TerminalLayoutPersistenceClient = .liveValue + layoutPersistence: TerminalLayoutPersistenceClient = .liveValue, + agentObservationBufferCapacity: Int = 64 ) { self.runtime = runtime self.layoutPersistence = layoutPersistence self.preferredFontSize = preferredFontSize + self.agentObservationStore = AgentObservationStore(bufferCapacity: agentObservationBufferCapacity) baselineFontSize = runtime.defaultFontSize() } @@ -223,6 +226,23 @@ final class WorktreeTerminalManager { } } + /// Independent per-surface multicast observation. This deliberately does not + /// reuse `eventStream()`, whose single production subscriber is `AppFeature`. + func observeAgentState(surfaceID: UUID) -> AgentObservationStream { + agentObservationStore.observe(surfaceID: surfaceID, isLive: containsSurface(surfaceID)) + } + + @discardableResult + func recordAgentSignal(_ signal: AgentSignal, surfaceID: UUID) -> Bool { + guard containsSurface(surfaceID) else { return false } + agentObservationStore.publishSignal(signal, surfaceID: surfaceID) + return true + } + + func agentObservationSubscriberCount(surfaceID: UUID) -> Int { + agentObservationStore.subscriberCount(surfaceID: surfaceID) + } + func eventStream() -> AsyncStream { eventContinuation?.finish() let (stream, continuation) = AsyncStream.makeStream( @@ -304,10 +324,17 @@ final class WorktreeTerminalManager { self?.emit(.taskStatusChanged(worktreeID: worktree.id, status: status)) } state.onAgentEntryChanged = { [weak self] entry in - self?.emit(.agentEntryChanged(entry)) + guard let self else { return } + agentObservationStore.publishAgentChanged(entry) + emit(.agentEntryChanged(entry)) } state.onAgentEntryRemoved = { [weak self] id in - self?.emit(.agentEntryRemoved(id)) + guard let self else { return } + agentObservationStore.publishAgentRemoved(surfaceID: id) + emit(.agentEntryRemoved(id)) + } + state.onSurfaceClosed = { [weak self] surfaceID in + self?.agentObservationStore.publishSurfaceClosed(surfaceID: surfaceID) } state.onRunScriptStatusChanged = { [weak self] isRunning in self?.emit(.runScriptStatusChanged(worktreeID: worktree.id, isRunning: isRunning)) @@ -451,6 +478,10 @@ final class WorktreeTerminalManager { activeWorktreeStates.first { $0.surfaceView(for: tabId) != nil } } + private func containsSurface(_ surfaceID: UUID) -> Bool { + states.values.contains { $0.surfaceView(for: surfaceID) != nil } + } + @discardableResult func broadcastCommittedText( _ text: String, @@ -745,6 +776,7 @@ final class WorktreeTerminalManager { self.runtime = nil self.layoutPersistence = .liveValue self.preferredFontSize = nil + self.agentObservationStore = AgentObservationStore(bufferCapacity: 64) self.baselineFontSize = 13 } #endif diff --git a/supacode/Features/Terminal/Models/WorktreeTerminalState+AgentDetection.swift b/supacode/Features/Terminal/Models/WorktreeTerminalState+AgentDetection.swift index b835a08b..5c3fba3f 100644 --- a/supacode/Features/Terminal/Models/WorktreeTerminalState+AgentDetection.swift +++ b/supacode/Features/Terminal/Models/WorktreeTerminalState+AgentDetection.swift @@ -283,7 +283,7 @@ extension WorktreeTerminalState { surfaceAgentStates[surfaceID] = PaneAgentState(lastChangedAt: Date()) lastWorkingAtBySurface.removeValue(forKey: surfaceID) lastAgentScreenScanBySurface.removeValue(forKey: surfaceID) - lastEmittedAgentEntriesBySurface.removeValue(forKey: surfaceID) + let hadPublishedEntry = lastEmittedAgentEntriesBySurface.removeValue(forKey: surfaceID) != nil // The launch identity lives exactly as long as the launched agent // (docs-ai 053/006): once the pane is a bare shell again, a manually // started agent is the user's own — default home, default account — and @@ -291,7 +291,9 @@ extension WorktreeTerminalState { launchProfilesBySurface.removeValue(forKey: surfaceID) lastAgentEntryEmitAtBySurface.removeValue(forKey: surfaceID) pendingAgentEntryBySurface.removeValue(forKey: surfaceID) - onAgentEntryRemoved?(surfaceID) + if hadPublishedEntry { + onAgentEntryRemoved?(surfaceID) + } if let tabId = tabId(containing: surfaceID) { updateTabAgentBusyState(for: tabId) } @@ -362,10 +364,12 @@ extension WorktreeTerminalState { now: Date = Date() ) { guard let entry = activeAgentEntry(surfaceID: surfaceID, tabId: tabId, state: state) else { - lastEmittedAgentEntriesBySurface.removeValue(forKey: surfaceID) + let hadPublishedEntry = lastEmittedAgentEntriesBySurface.removeValue(forKey: surfaceID) != nil lastAgentEntryEmitAtBySurface.removeValue(forKey: surfaceID) pendingAgentEntryBySurface.removeValue(forKey: surfaceID) - onAgentEntryRemoved?(surfaceID) + if hadPublishedEntry { + onAgentEntryRemoved?(surfaceID) + } return } if let previous = lastEmittedAgentEntriesBySurface[surfaceID] { @@ -464,17 +468,19 @@ extension WorktreeTerminalState { lastWorkingAtBySurface.removeValue(forKey: surfaceId) lastAgentDetectionDiagnosticsBySurface.removeValue(forKey: surfaceId) lastAgentScreenScanBySurface.removeValue(forKey: surfaceId) - lastEmittedAgentEntriesBySurface.removeValue(forKey: surfaceId) + let hadPublishedEntry = lastEmittedAgentEntriesBySurface.removeValue(forKey: surfaceId) != nil lastAgentEntryEmitAtBySurface.removeValue(forKey: surfaceId) pendingAgentEntryBySurface.removeValue(forKey: surfaceId) - onAgentEntryRemoved?(surfaceId) + if hadPublishedEntry { + onAgentEntryRemoved?(surfaceId) + } } func cleanupAllAgentDetectionState() { for task in agentDetectionTasks.values { task.cancel() } - let removedIDs = Array(surfaceAgentStates.keys) + let removedIDs = Array(lastEmittedAgentEntriesBySurface.keys) agentDetectionTasks.removeAll() agentDetectionSchedules.removeAll() surfaceAgentStates.removeAll() diff --git a/supacode/Features/Terminal/Models/WorktreeTerminalState+Surfaces.swift b/supacode/Features/Terminal/Models/WorktreeTerminalState+Surfaces.swift index 4ee8b05c..25b80891 100644 --- a/supacode/Features/Terminal/Models/WorktreeTerminalState+Surfaces.swift +++ b/supacode/Features/Terminal/Models/WorktreeTerminalState+Surfaces.swift @@ -137,10 +137,7 @@ extension WorktreeTerminalState { return .success(newSurface.id) } catch { newSurface.closeSurface() - surfaces.removeValue(forKey: newSurface.id) - surfaceRunningStartedAtById.removeValue(forKey: newSurface.id) - cleanupCommandDetectorState(forSurfaceId: newSurface.id) - cleanupAgentDetectionState(forSurfaceId: newSurface.id) + forgetSurface(newSurface.id) return .failure(.insertionFailed) } } @@ -328,9 +325,11 @@ extension WorktreeTerminalState { for tab in tabManager.tabs { unregisterTargetHandle(for: tab.id) } - for surface in surfaces.values { - unregisterTargetHandle(for: surface.id) + // `closeSurface()` may synchronously re-enter `forgetSurface`; snapshot the + // values and let the guarded cleanup seam make each lifecycle exactly once. + for surface in Array(surfaces.values) { surface.closeSurface() + forgetSurface(surface.id) } surfaces.removeAll() launchProfilesBySurface.removeAll() @@ -631,6 +630,7 @@ extension WorktreeTerminalState { /// without dropping them here the worktree's unseen indicator (bell + Dock /// badge) would stay lit until the user manually dismisses everything. func forgetSurface(_ surfaceID: UUID) { + guard surfaces[surfaceID] != nil else { return } unregisterTargetHandle(for: surfaceID) surfaces.removeValue(forKey: surfaceID) launchProfilesBySurface.removeValue(forKey: surfaceID) @@ -645,6 +645,7 @@ extension WorktreeTerminalState { let previousHasUnseen = hasUnseenNotification notifications = Self.prunedNotifications(from: notifications, removingSurfaceID: surfaceID) emitNotificationIndicatorIfNeeded(previousHasUnseen: previousHasUnseen) + onSurfaceClosed?(surfaceID) } func removeTree(for tabId: TerminalTabID) { diff --git a/supacode/Features/Terminal/Models/WorktreeTerminalState.swift b/supacode/Features/Terminal/Models/WorktreeTerminalState.swift index 32ed6af0..991753fc 100644 --- a/supacode/Features/Terminal/Models/WorktreeTerminalState.swift +++ b/supacode/Features/Terminal/Models/WorktreeTerminalState.swift @@ -266,6 +266,8 @@ final class WorktreeTerminalState { var onTaskStatusChanged: ((WorktreeTaskStatus) -> Void)? var onAgentEntryChanged: ((ActiveAgentEntry) -> Void)? var onAgentEntryRemoved: ((ActiveAgentEntry.ID) -> Void)? + /// Emitted exactly once after agent cleanup for each torn-down surface. + var onSurfaceClosed: ((UUID) -> Void)? var onRunScriptStatusChanged: ((Bool) -> Void)? var onCommandPaletteToggle: (() -> Void)? var onSetupScriptConsumed: (() -> Void)? diff --git a/supacodeTests/AgentObservationTests.swift b/supacodeTests/AgentObservationTests.swift new file mode 100644 index 00000000..b604ce74 --- /dev/null +++ b/supacodeTests/AgentObservationTests.swift @@ -0,0 +1,232 @@ +import Foundation +import GhosttyKit +import Testing + +@testable import supacode + +@MainActor +struct AgentObservationTests { + @Test func liveShellStartsWithAtomicEmptySnapshotAndReplaysLatestSignal() async throws { + let fixture = makeFixture() + let firstStream = fixture.manager.observeAgentState(surfaceID: fixture.surfaceID) + var firstIterator = firstStream.makeAsyncIterator() + + let first = try await firstIterator.next() + guard case .snapshot(let initial) = first else { + Issue.record("Expected snapshot first") + return + } + #expect(initial.agent == nil) + #expect(initial.latestSignal == nil) + #expect(initial.revision == 0) + + let signal = makeSignal() + #expect(fixture.manager.recordAgentSignal(signal, surfaceID: fixture.surfaceID)) + #expect(try await firstIterator.next() == .signal(signal)) + + let secondStream = fixture.manager.observeAgentState(surfaceID: fixture.surfaceID) + var secondIterator = secondStream.makeAsyncIterator() + guard case .snapshot(let replay) = try await secondIterator.next() else { + Issue.record("Expected replay snapshot first") + return + } + #expect(replay.agent == nil) + #expect(replay.latestSignal == signal) + #expect(replay.revision == 1) + } + + @Test func signalIsMulticastToIndependentSubscribers() async throws { + let fixture = makeFixture() + var first = fixture.manager.observeAgentState(surfaceID: fixture.surfaceID).makeAsyncIterator() + var second = fixture.manager.observeAgentState(surfaceID: fixture.surfaceID).makeAsyncIterator() + _ = try await first.next() + _ = try await second.next() + + let signal = makeSignal(detail: "Review complete") + #expect(fixture.manager.recordAgentSignal(signal, surfaceID: fixture.surfaceID)) + + #expect(try await first.next() == .signal(signal)) + #expect(try await second.next() == .signal(signal)) + } + + @Test func publishedAgentRemovalPrecedesSurfaceClosureAndFinishesStream() async throws { + let fixture = makeFixture() + var iterator = fixture.manager.observeAgentState(surfaceID: fixture.surfaceID).makeAsyncIterator() + _ = try await iterator.next() + + fixture.state.emitAgentEntry( + surfaceID: fixture.surfaceID, + tabId: fixture.tabID, + state: PaneAgentState(detectedAgent: .claude, state: .working) + ) + guard case .changed(let entry) = try await iterator.next() else { + Issue.record("Expected changed event") + return + } + #expect(entry.surfaceID == fixture.surfaceID) + + #expect(fixture.state.closeSurface(id: fixture.surfaceID, confirmation: .skip)) + #expect(try await iterator.next() == .removed) + #expect(try await iterator.next() == .surfaceClosed) + #expect(try await iterator.next() == nil) + } + + @Test func shellWithoutPublishedAgentClosesWithoutFalseRemoval() async throws { + let fixture = makeFixture() + fixture.state.wakeAgentDetection(forSurfaceID: fixture.surfaceID) + var iterator = fixture.manager.observeAgentState(surfaceID: fixture.surfaceID).makeAsyncIterator() + _ = try await iterator.next() + + #expect(fixture.state.closeSurface(id: fixture.surfaceID, confirmation: .skip)) + #expect(try await iterator.next() == .surfaceClosed) + #expect(try await iterator.next() == nil) + } + + @Test func agentRemovalDoesNotCloseObserverAndLaterSignalStillArrives() async throws { + let fixture = makeFixture() + var iterator = fixture.manager.observeAgentState(surfaceID: fixture.surfaceID).makeAsyncIterator() + _ = try await iterator.next() + fixture.state.emitAgentEntry( + surfaceID: fixture.surfaceID, + tabId: fixture.tabID, + state: PaneAgentState(detectedAgent: .claude, state: .working) + ) + _ = try await iterator.next() + + fixture.state.emitAgentEntry( + surfaceID: fixture.surfaceID, + tabId: fixture.tabID, + state: PaneAgentState() + ) + #expect(try await iterator.next() == .removed) + + let signal = makeSignal() + #expect(fixture.manager.recordAgentSignal(signal, surfaceID: fixture.surfaceID)) + #expect(try await iterator.next() == .signal(signal)) + } + + @Test func closeAllAndPruneFinishEverySurfaceObserver() async throws { + let closeAllFixture = makeFixture() + var closeAllIterator = closeAllFixture.manager + .observeAgentState(surfaceID: closeAllFixture.surfaceID) + .makeAsyncIterator() + _ = try await closeAllIterator.next() + + closeAllFixture.state.closeAllSurfaces() + #expect(try await closeAllIterator.next() == .surfaceClosed) + #expect(try await closeAllIterator.next() == nil) + + let pruneFixture = makeFixture() + var pruneIterator = pruneFixture.manager + .observeAgentState(surfaceID: pruneFixture.surfaceID) + .makeAsyncIterator() + _ = try await pruneIterator.next() + + pruneFixture.manager.prune(keeping: []) + #expect(try await pruneIterator.next() == .surfaceClosed) + #expect(try await pruneIterator.next() == nil) + } + + @Test func cancellationRemovesOnlyTheCancelledSubscriber() async throws { + let fixture = makeFixture() + let firstStream = fixture.manager.observeAgentState(surfaceID: fixture.surfaceID) + let secondStream = fixture.manager.observeAgentState(surfaceID: fixture.surfaceID) + #expect(fixture.manager.agentObservationSubscriberCount(surfaceID: fixture.surfaceID) == 2) + + let task = Task { + for try await _ in firstStream {} + } + await Task.yield() + task.cancel() + _ = await task.result + for _ in 0..<10 where fixture.manager.agentObservationSubscriberCount(surfaceID: fixture.surfaceID) == 2 { + await Task.yield() + } + + #expect(fixture.manager.agentObservationSubscriberCount(surfaceID: fixture.surfaceID) == 1) + var secondIterator = secondStream.makeAsyncIterator() + _ = try await secondIterator.next() + let signal = makeSignal() + #expect(fixture.manager.recordAgentSignal(signal, surfaceID: fixture.surfaceID)) + #expect(try await secondIterator.next() == .signal(signal)) + } + + @Test func closedSurfaceProducesSnapshotThenSurfaceClosed() async throws { + let manager = WorktreeTerminalManager(runtime: GhosttyRuntime()) + var iterator = manager.observeAgentState(surfaceID: UUID()).makeAsyncIterator() + + guard case .snapshot(let snapshot) = try await iterator.next() else { + Issue.record("Expected snapshot first") + return + } + #expect(snapshot.agent == nil) + #expect(snapshot.latestSignal == nil) + #expect(try await iterator.next() == .surfaceClosed) + #expect(try await iterator.next() == nil) + } + + @Test func boundedOverflowIsExplicitAndRecoverableByResubscription() async throws { + let fixture = makeFixture(bufferCapacity: 1) + let stream = fixture.manager.observeAgentState(surfaceID: fixture.surfaceID) + let signal = makeSignal() + + // Leave the initial snapshot buffered so the next critical event overflows. + #expect(fixture.manager.recordAgentSignal(signal, surfaceID: fixture.surfaceID)) + + var iterator = stream.makeAsyncIterator() + guard case .snapshot = try await iterator.next() else { + Issue.record("Expected the protected initial snapshot") + return + } + do { + _ = try await iterator.next() + Issue.record("Expected explicit buffer overflow") + } catch let error as AgentObservationError { + #expect(error == .bufferOverflow) + } + + var replacement = fixture.manager.observeAgentState(surfaceID: fixture.surfaceID).makeAsyncIterator() + guard case .snapshot(let snapshot) = try await replacement.next() else { + Issue.record("Expected resubscription snapshot") + return + } + #expect(snapshot.latestSignal == signal) + } + + private struct Fixture { + let manager: WorktreeTerminalManager + let state: WorktreeTerminalState + let tabID: TerminalTabID + let surfaceID: UUID + } + + private func makeFixture(bufferCapacity: Int = 64) -> Fixture { + let manager = WorktreeTerminalManager( + runtime: GhosttyRuntime(), + agentObservationBufferCapacity: bufferCapacity + ) + let worktree = Worktree( + id: "/tmp/agent-observation", + name: "agent-observation", + detail: "", + workingDirectory: URL(fileURLWithPath: "/tmp/agent-observation"), + repositoryRootURL: URL(fileURLWithPath: "/tmp/agent-observation") + ) + let state = manager.state(for: worktree) + let tabID = state.createTab()! + let surfaceID = state.focusedSurfaceId(in: tabID)! + return Fixture(manager: manager, state: state, tabID: tabID, surfaceID: surfaceID) + } + + private func makeSignal(detail: String? = nil) -> AgentSignal { + AgentSignal( + kind: .turnEnded, + source: .cooperativeCLI, + confidence: .exact, + timestamp: Date(timeIntervalSince1970: 1_000), + sessionID: "session-1", + detail: detail, + claimedOrigin: nil + ) + } +} diff --git a/supacodeTests/CLIAgentSignalCommandHandlerTests.swift b/supacodeTests/CLIAgentSignalCommandHandlerTests.swift new file mode 100644 index 00000000..d70efc13 --- /dev/null +++ b/supacodeTests/CLIAgentSignalCommandHandlerTests.swift @@ -0,0 +1,144 @@ +import Foundation +import Testing + +@testable import supacode + +@MainActor +struct CLIAgentSignalCommandHandlerTests { + @Test func recordsSignalForExactCallerPaneAndReturnsReceipt() async throws { + let pane = CallerPane(worktreeID: "/tmp/repo", surfaceID: UUID()) + var recorded: (CallerPane, AgentSignal)? + let handler = AgentSignalCommandHandler( + resolveCaller: { processID in + #expect(processID == 42) + return pane + }, + recordSignal: { caller, signal in + recorded = (caller, signal) + return true + }, + now: { Date(timeIntervalSince1970: 1_000) } + ) + let input = AgentSignalInput( + event: .turnEnded, + origin: "manual-review", + sessionID: "session-1", + detail: "Review complete" + ) + + let response = await handler.handle( + envelope: CommandEnvelope(output: .json, command: .agentsSignal(input)), + context: CLICommandContext(callerProcessID: 42) + ) + + #expect(response.ok) + #expect(response.command == "agents.signal") + #expect(response.schemaVersion == "prowl.cli.agents.signal.v1") + let payload = try #require(try response.data?.decode(as: AgentSignalCommandPayload.self)) + #expect(payload.pane.id == pane.surfaceID.uuidString) + #expect(payload.pane.worktreeID == pane.worktreeID) + #expect(payload.signal.event == .turnEnded) + #expect(payload.signal.progress == nil) + #expect(payload.signal.source == "cooperative_cli") + #expect(payload.signal.confidence == "exact") + #expect(payload.signal.timestamp == "1970-01-01T00:16:40.000Z") + #expect(payload.signal.sessionID == "session-1") + #expect(payload.signal.detail == "Review complete") + #expect(payload.signal.claimedOrigin == "manual-review") + #expect(recorded?.0 == pane) + #expect(recorded?.1.kind == .turnEnded) + #expect(recorded?.1.source == .cooperativeCLI) + #expect(recorded?.1.confidence == .exact) + #expect(recorded?.1.timestamp == Date(timeIntervalSince1970: 1_000)) + } + + @Test func supportsIndeterminateAndBoundedProgress() async throws { + var signals: [AgentSignal] = [] + let pane = CallerPane(worktreeID: "wt", surfaceID: UUID()) + let handler = AgentSignalCommandHandler( + resolveCaller: { _ in pane }, + recordSignal: { _, signal in + signals.append(signal) + return true + } + ) + + for value in [nil, 0, 100] as [Int?] { + let response = await handler.handle( + envelope: CommandEnvelope( + output: .json, + command: .agentsSignal(AgentSignalInput(event: .progress, progress: value)) + ), + context: CLICommandContext(callerProcessID: 7) + ) + #expect(response.ok) + } + + #expect(signals.map(\.kind) == [.progress(nil), .progress(0), .progress(100)]) + } + + @Test func rejectsMissingCallerWithoutRecording() async { + var didRecord = false + let handler = AgentSignalCommandHandler( + resolveCaller: { _ in nil }, + recordSignal: { _, _ in + didRecord = true + return true + } + ) + + let response = await handler.handle( + envelope: CommandEnvelope(output: .json, command: .agentsSignal(AgentSignalInput(event: .needsInput))), + context: CLICommandContext(callerProcessID: nil) + ) + + #expect(response.ok == false) + #expect(response.error?.code == CLIErrorCode.sourceRequired) + #expect(didRecord == false) + } + + @Test func rejectsClosedCallerPane() async { + let handler = AgentSignalCommandHandler( + resolveCaller: { _ in CallerPane(worktreeID: "wt", surfaceID: UUID()) }, + recordSignal: { _, _ in false } + ) + + let response = await handler.handle( + envelope: CommandEnvelope(output: .json, command: .agentsSignal(AgentSignalInput(event: .sessionEnd))), + context: CLICommandContext(callerProcessID: 7) + ) + + #expect(response.ok == false) + #expect(response.error?.code == CLIErrorCode.agentGone) + } + + @Test func validatesWireInputBeforeCallerResolution() async { + var resolved = false + let handler = AgentSignalCommandHandler( + resolveCaller: { _ in + resolved = true + return nil + }, + recordSignal: { _, _ in true } + ) + let invalidInputs = [ + AgentSignalInput(event: .turnEnded, progress: 1), + AgentSignalInput(event: .progress, progress: -1), + AgentSignalInput(event: .progress, progress: 101), + AgentSignalInput(event: .needsInput, sessionID: ""), + AgentSignalInput(event: .needsInput, origin: "bad\norigin"), + AgentSignalInput(event: .needsInput, detail: "bad\u{0}detail"), + AgentSignalInput(event: .needsInput, detail: String(repeating: "x", count: 4_097)), + ] + + for input in invalidInputs { + let response = await handler.handle( + envelope: CommandEnvelope(output: .json, command: .agentsSignal(input)), + context: CLICommandContext(callerProcessID: 7) + ) + #expect(response.ok == false) + #expect(response.error?.code == CLIErrorCode.invalidArgument) + } + #expect(resolved == false) + } +} diff --git a/supacodeTests/CLICommandEnvelopeTests.swift b/supacodeTests/CLICommandEnvelopeTests.swift index c71e3953..5d4116cf 100644 --- a/supacodeTests/CLICommandEnvelopeTests.swift +++ b/supacodeTests/CLICommandEnvelopeTests.swift @@ -71,6 +71,31 @@ struct CLICommandEnvelopeTests { } } + @Test func envelopeAgentsSignalRoundTrips() throws { + let envelope = CommandEnvelope( + output: .json, + command: .agentsSignal( + AgentSignalInput( + event: .needsInput, + origin: "manual", + sessionID: "session-1", + detail: "Approval required" + )) + ) + let data = try JSONEncoder().encode(envelope) + let decoded = try JSONDecoder().decode(CommandEnvelope.self, from: data) + + if case .agentsSignal(let input) = decoded.command { + #expect(input.event == .needsInput) + #expect(input.progress == nil) + #expect(input.origin == "manual") + #expect(input.sessionID == "session-1") + #expect(input.detail == "Approval required") + } else { + Issue.record("Expected .agentsSignal command") + } + } + @Test func envelopeAgentsReadRoundTrips() throws { let envelope = CommandEnvelope( output: .text, @@ -265,6 +290,7 @@ struct CLICommandEnvelopeTests { (.list(ListInput()), "list"), (.agents(AgentsInput()), "agents"), (.agentsRead(AgentReadInput(pane: "p7")), "agents.read"), + (.agentsSignal(AgentSignalInput(event: .turnEnded)), "agents.signal"), (.focus(FocusInput()), "focus"), (.send(SendInput(text: "x")), "send"), (.key(KeyInput(rawToken: "tab", token: "tab")), "key"), @@ -285,6 +311,7 @@ struct CLICommandEnvelopeTests { .list(ListInput()), .agents(AgentsInput()), .agentsRead(AgentReadInput(pane: "p7")), + .agentsSignal(AgentSignalInput(event: .sessionStart)), .focus(FocusInput()), .send(SendInput(text: "test")), .key(KeyInput(rawToken: "enter", token: "enter")), diff --git a/supacodeTests/CLICommandRouterTests.swift b/supacodeTests/CLICommandRouterTests.swift index 17e7d566..c7117d32 100644 --- a/supacodeTests/CLICommandRouterTests.swift +++ b/supacodeTests/CLICommandRouterTests.swift @@ -67,6 +67,24 @@ struct CLICommandRouterTests { #expect(response.error?.code == "NOT_IMPLEMENTED") } + @MainActor + @Test func routerDispatchesAgentsSignalWithCallerContext() async { + let handler = ContextRecordingCommandHandler(command: "agents.signal") + let router = CLICommandRouter(agentsSignalHandler: handler) + let envelope = CommandEnvelope( + output: .json, + command: .agentsSignal(AgentSignalInput(event: .turnEnded)) + ) + + let response = await router.route( + envelope, + context: CLICommandContext(callerProcessID: 42) + ) + + #expect(response.command == "agents.signal") + #expect(handler.callerProcessID == 42) + } + @MainActor @Test func routerDispatchesProfilesToProfilesHandler() async { let router = CLICommandRouter() @@ -162,6 +180,7 @@ struct CLICommandRouterTests { .list(ListInput()), .agents(AgentsInput()), .agentsRead(AgentReadInput(pane: "p7")), + .agentsSignal(AgentSignalInput(event: .turnEnded)), .profiles(ProfilesInput()), .focus(FocusInput()), .send(SendInput(text: "x")), @@ -190,3 +209,27 @@ private struct MockCommandHandler: CommandHandler { response } } + +@MainActor +private final class ContextRecordingCommandHandler: CommandHandler { + let command: String + var callerProcessID: pid_t? + + init(command: String) { + self.command = command + } + + func handle(envelope: CommandEnvelope) async -> CommandResponse { + await handle(envelope: envelope, context: CLICommandContext()) + } + + // swiftlint:disable:next async_without_await + func handle(envelope: CommandEnvelope, context: CLICommandContext) async -> CommandResponse { + callerProcessID = context.callerProcessID + return CommandResponse( + ok: true, + command: command, + schemaVersion: "prowl.cli.\(command).v1" + ) + } +} diff --git a/supacodeTests/CLISocketServerTests.swift b/supacodeTests/CLISocketServerTests.swift index 7e520d41..2c204fca 100644 --- a/supacodeTests/CLISocketServerTests.swift +++ b/supacodeTests/CLISocketServerTests.swift @@ -88,6 +88,64 @@ struct CLISocketServerTests { #expect(!CLISocketServer.isAllowedPeerUID(502, currentUID: 501)) } + @Test func socketRoundTripThreadsKernelPeerPIDIntoAgentSignalHandler() async throws { + let socketPath = temporarySocketPath(suffix: "signal-context") + let pane = CallerPane(worktreeID: "wt", surfaceID: UUID()) + var recordedSignal: AgentSignal? + let handler = AgentSignalCommandHandler( + resolveCaller: { processID in + #expect(processID == getpid()) + return pane + }, + recordSignal: { caller, signal in + #expect(caller == pane) + recordedSignal = signal + return true + }, + now: { Date(timeIntervalSince1970: 1_000) } + ) + let server = CLISocketServer( + router: CLICommandRouter(agentsSignalHandler: handler), + socketPath: socketPath + ) + try server.start() + defer { server.stop() } + let envelope = CommandEnvelope( + output: .json, + command: .agentsSignal(AgentSignalInput(event: .turnEnded, detail: "socket result")) + ) + + let requestData = try JSONEncoder().encode(envelope) + let responseData = try await withCheckedThrowingContinuation { continuation in + DispatchQueue.global(qos: .userInitiated).async { + continuation.resume( + with: Result { + try Self.send(requestData: requestData, socketPath: socketPath) + } + ) + } + } + let response = try JSONDecoder().decode(CommandResponse.self, from: responseData) + + #expect(response.ok) + #expect(response.command == "agents.signal") + #expect(recordedSignal?.kind == .turnEnded) + #expect(recordedSignal?.detail == "socket result") + } + + #if canImport(Darwin) + @Test func kernelReportsPeerPIDForLocalSocket() throws { + var descriptors: [Int32] = [-1, -1] + #expect(socketpair(AF_UNIX, SOCK_STREAM, 0, &descriptors) == 0) + defer { + close(descriptors[0]) + close(descriptors[1]) + } + + #expect(CLISocketServer.peerProcessID(descriptors[0]) == getpid()) + } + #endif + private func temporarySocketPath(suffix: String) -> String { URL(fileURLWithPath: "/tmp", isDirectory: true) .appending(path: "prowl-cli-tests-\(UUID().uuidString.prefix(8))", directoryHint: .isDirectory) @@ -97,6 +155,73 @@ struct CLISocketServerTests { .path(percentEncoded: false) } + nonisolated private static func send(requestData: Data, socketPath: String) throws -> Data { + let socketFD = socket(AF_UNIX, SOCK_STREAM, 0) + guard socketFD >= 0 else { throw CLIServiceError.socketCreationFailed } + defer { close(socketFD) } + + let connected = withSocketAddress(socketPath) { address in + withUnsafePointer(to: address) { pointer in + pointer.withMemoryRebound(to: sockaddr.self, capacity: 1) { socketPointer in + connect(socketFD, socketPointer, socklen_t(MemoryLayout.size)) + } + } + } + guard connected == 0 else { throw TestSocketClientError.connectFailed } + + var requestLength = UInt32(requestData.count).bigEndian + try withUnsafeBytes(of: &requestLength) { try writeAll(socketFD, buffer: $0) } + try requestData.withUnsafeBytes { try writeAll(socketFD, buffer: $0) } + + let lengthData = try readExact(socketFD, count: 4) + let responseLength = lengthData.withUnsafeBytes { UInt32(bigEndian: $0.load(as: UInt32.self)) } + return try readExact(socketFD, count: Int(responseLength)) + } + + nonisolated private static func writeAll(_ fileDescriptor: Int32, buffer: UnsafeRawBufferPointer) throws { + var offset = 0 + while offset < buffer.count { + let written = Darwin.write( + fileDescriptor, + buffer.baseAddress!.advanced(by: offset), + buffer.count - offset + ) + guard written > 0 else { throw CLIServiceError.writeFailed } + offset += written + } + } + + nonisolated private static func readExact(_ fileDescriptor: Int32, count: Int) throws -> Data { + var data = Data(count: count) + var offset = 0 + while offset < count { + let readCount = data.withUnsafeMutableBytes { buffer in + Darwin.read(fileDescriptor, buffer.baseAddress!.advanced(by: offset), count - offset) + } + guard readCount > 0 else { throw CLIServiceError.readFailed } + offset += readCount + } + return data + } + + nonisolated private static func withSocketAddress( + _ socketPath: String, + _ body: (sockaddr_un) throws -> Result + ) rethrows -> Result { + var address = sockaddr_un() + address.sun_family = sa_family_t(AF_UNIX) + let pathBytes = Array(socketPath.utf8) + let maxLength = MemoryLayout.size(ofValue: address.sun_path) - 1 + precondition(pathBytes.count <= maxLength) + withUnsafeMutableBytes(of: &address.sun_path) { buffer in + for index in pathBytes.indices { + buffer[index] = pathBytes[index] + } + buffer[pathBytes.count] = 0 + } + return try body(address) + } + private func fileMode(at path: String) -> mode_t? { var statValue = stat() guard stat(path, &statValue) == 0 else { return nil } @@ -168,3 +293,7 @@ struct CLISocketServerTests { return try body(address) } } + +private enum TestSocketClientError: Error { + case connectFailed +} diff --git a/supacodeTests/CallerPaneResolverTests.swift b/supacodeTests/CallerPaneResolverTests.swift new file mode 100644 index 00000000..93ef298f --- /dev/null +++ b/supacodeTests/CallerPaneResolverTests.swift @@ -0,0 +1,64 @@ +import Foundation +import Testing + +@testable import supacode + +struct CallerPaneResolverTests { + @Test func resolvesDirectAndNestedCallerAncestry() { + let pane = CallerPane(worktreeID: "wt", surfaceID: UUID()) + let parents: [pid_t: pid_t] = [400: 300, 300: 200, 200: 100] + + #expect( + CallerPaneResolver.pane( + forCallerProcess: 100, + paneByShellPID: [100: pane], + parentProcessID: { parents[$0] } + ) == pane + ) + #expect( + CallerPaneResolver.pane( + forCallerProcess: 400, + paneByShellPID: [100: pane], + parentProcessID: { parents[$0] } + ) == pane + ) + } + + @Test func unresolvedAndCyclicAncestryNeverGuess() { + let focusedButUnrelated = CallerPane(worktreeID: "focused", surfaceID: UUID()) + + #expect( + CallerPaneResolver.pane( + forCallerProcess: 400, + paneByShellPID: [100: focusedButUnrelated], + parentProcessID: { _ in nil } + ) == nil + ) + #expect( + CallerPaneResolver.pane( + forCallerProcess: 400, + paneByShellPID: [100: focusedButUnrelated], + parentProcessID: { $0 } + ) == nil + ) + } + + @Test func ancestryWalkIsBounded() { + let pane = CallerPane(worktreeID: "wt", surfaceID: UUID()) + + #expect( + CallerPaneResolver.pane( + forCallerProcess: 100, + paneByShellPID: [67: pane], + parentProcessID: { $0 - 1 } + ) == nil + ) + #expect( + CallerPaneResolver.pane( + forCallerProcess: 100, + paneByShellPID: [69: pane], + parentProcessID: { $0 - 1 } + ) == pane + ) + } +} diff --git a/supacodeTests/SupacodeAppCLITests.swift b/supacodeTests/SupacodeAppCLITests.swift index 0ac1fede..c0af3bc1 100644 --- a/supacodeTests/SupacodeAppCLITests.swift +++ b/supacodeTests/SupacodeAppCLITests.swift @@ -37,6 +37,50 @@ struct SupacodeAppCLITests { #expect(readResponse.error?.code != "NOT_IMPLEMENTED") } + @Test func cliRouterWiresCallerAttributedAgentSignalIntoTerminalObserver() async throws { + let store = Store(initialState: AppFeature.State()) { + AppFeature() + } + let terminalManager = WorktreeTerminalManager(runtime: GhosttyRuntime()) + let worktree = Worktree( + id: "/tmp/app-cli-signal", + name: "app-cli-signal", + detail: "", + workingDirectory: URL(fileURLWithPath: "/tmp/app-cli-signal"), + repositoryRootURL: URL(fileURLWithPath: "/tmp/app-cli-signal") + ) + let state = terminalManager.state(for: worktree) + let tabID = try #require(state.createTab()) + let surfaceID = try #require(state.focusedSurfaceId(in: tabID)) + var observer = terminalManager.observeAgentState(surfaceID: surfaceID).makeAsyncIterator() + _ = try await observer.next() + let router = SupacodeApp.makeCLICommandRouter( + appStore: store, + terminalManager: terminalManager, + agentSignalCallerResolver: { processID in + #expect(processID == 42) + return CallerPane(worktreeID: worktree.id, surfaceID: surfaceID) + } + ) + + let response = await router.route( + CommandEnvelope( + output: .json, + command: .agentsSignal(AgentSignalInput(event: .needsInput, detail: "Approval required")) + ), + context: CLICommandContext(callerProcessID: 42) + ) + + #expect(response.ok) + guard case .signal(let signal) = try await observer.next() else { + Issue.record("Expected signal on terminal observer") + return + } + #expect(signal.kind == .needsInput) + #expect(signal.detail == "Approval required") + #expect(signal.source == .cooperativeCLI) + } + @Test func resolveCLITerminalWorktreeBuildsSyntheticRunnableFolderWorktree() { let repository = Repository( id: "/Users/test/PlainFolder",