diff --git a/PLAN.md b/PLAN.md index a877836..7b37ad6 100644 --- a/PLAN.md +++ b/PLAN.md @@ -137,7 +137,7 @@ Implemented and validated in Phase 2F: - Unit tests for config, build-sections, apply-setting-change, background-blocker, and i18n. Current intentional gaps: -- No logs/dock UI extensions yet. +- No dock UI extension yet (`extensions/processes-dock/` is the next phase). - Keybindings are managed by Pi's built-in KeybindingsManager; not in extension config. - `clear` and `write` tool actions are deferred. The agent can `read` log file paths returned by `list` or `output` for full-log access. - `package.json` still references `./skills/pi-processes`, but the local `skills/` directory is absent. Either restore the skill later or remove the `pi.skills`/`files` entries during cleanup. @@ -1013,6 +1013,10 @@ After Phase 2F: ## Phase 3: Logs Extension (`extensions/processes-logs/`) +### Status + +Phase 3 is complete. The logs extension is registered in `package.json`, `/ps:logs` opens a tabbed overlay backed by the log subscription protocol, and the extension has unit tests for the protocol client, completions, and the `LogFileViewer` (render, scroll, search, stream filter, follow, notify-match highlighting). The `LogOverlayComponent` itself is covered indirectly through its dependencies plus a manual/e2e pass, since there is no TUI test harness. + ### Goal Own `/ps:logs` only. The logs extension is focused on process output: selecting a process, viewing live logs, scrolling, searching, stream filtering, and follow mode. It does not own process-management controls such as kill or clear. @@ -1133,6 +1137,20 @@ function connectToProcessLogs( ## Phase 3 bis: Notification Event Fanout and Log-Match Highlighting +### Status + +Phase 3 bis is complete. Notification flow is now event-driven: `NotificationService` emits a protocol-safe payload on `CHANNELS.NOTIFICATION`, a core delivery listener converts the payload back into a persisted `ad-process:notification` custom message with attention-derived send options, and the logs overlay consumes the same channel to highlight matched log lines. Lifecycle and log-match notifications behave exactly as before; only the delivery path changed. + +Implemented: +- `CHANNELS.NOTIFICATION` channel and `ProcessProtocolNotificationPayload` protocol type in `src/protocol/notifications.ts` (re-exported from `src/protocol/index.ts`). +- Extension notification types now re-export the protocol types so the two cannot drift. +- `NotificationService` takes `events: EventBus` and emits payloads instead of calling `pi.sendMessage`. +- `attentionToSendOptions` moved to `extensions/processes/notification-sender.ts` and is reused by the delivery listener. +- `extensions/processes/handlers/notifications.ts` delivery listener with malformed-payload guarding; wired into `extensions/processes/index.ts` disposers (disposed before the manager is killed on `session_shutdown`). +- `/ps:logs` listens for `kind === "log_match"` notifications, stores per-process match markers (capped at 100 per process), forwards markers to the focused viewer, and replays stored markers on tab switch. +- `LogFileViewer` underlines notify-match lines with lower priority than manual search (current search match > search match > notify match > stream/stderr color). +- Tests: service tests rewritten to assert on emitted events; delivery handler tests; viewer notify-highlight and priority tests. + ### Goal Make process notifications observable through `pi.events` before they are delivered to Pi messages. This lets the main extension keep owning agent/user notification delivery while UI extensions can react to the same events without importing notification internals. diff --git a/extensions/processes-logs/components/log-file-viewer.test.ts b/extensions/processes-logs/components/log-file-viewer.test.ts index 6758670..5931c57 100644 --- a/extensions/processes-logs/components/log-file-viewer.test.ts +++ b/extensions/processes-logs/components/log-file-viewer.test.ts @@ -1,5 +1,5 @@ import type { Theme } from "@earendil-works/pi-coding-agent"; -import { describe, expect, it } from "vitest"; +import { describe, expect, it, vi } from "vitest"; import { LogFileViewer } from "./log-file-viewer"; function makeTheme(): Theme { @@ -8,7 +8,7 @@ function makeTheme(): Theme { bg: (_color: string, text: string) => text, bold: (text: string) => `[b]${text}[/b]`, italic: (text: string) => text, - underline: (text: string) => text, + underline: (text: string) => `[u]${text}[/u]`, inverse: (text: string) => `[inv]${text}[/inv]`, strikethrough: (text: string) => text, } as unknown as Theme; @@ -140,4 +140,55 @@ describe("LogFileViewer", () => { const lines = trimLines(viewer.render(20, 5)); expect(lines[4]).toBe("line 14"); }); + + it("highlights notify log-match lines with lower priority than search", () => { + const underline = vi.fn((text: string) => text); + const theme = { + ...makeTheme(), + underline, + } as unknown as Theme; + const viewer = new LogFileViewer( + [ + { type: "stdout", text: "match" }, + { type: "stdout", text: "other" }, + ], + theme, + { followEnabled: false, maxBufferLines: 10 }, + ); + + viewer.addNotifyMatch({ line: "match" }); + expect(viewer.getNotifyMatchCount()).toBe(1); + + viewer.render(100, 2); + // Notify match line is underlined... + expect(underline).toHaveBeenCalledTimes(1); + expect(underline.mock.calls[0][0].trim()).toBe("match"); + }); + + it("search matches take priority over notify matches", () => { + const underline = vi.fn((text: string) => text); + const inverse = vi.fn((text: string) => text); + const bold = vi.fn((text: string) => text); + const theme = { + ...makeTheme(), + underline, + inverse, + bold, + } as unknown as Theme; + const viewer = new LogFileViewer( + [{ type: "stdout", text: "match" }], + theme, + { followEnabled: false, maxBufferLines: 10 }, + ); + + viewer.addNotifyMatch({ line: "match" }); + viewer.setSearch("match"); + + viewer.render(100, 1); + // Current search match is bold+inverse, not underline-only. + expect(inverse).toHaveBeenCalledTimes(1); + expect(inverse.mock.calls[0][0].trim()).toBe("match"); + expect(bold).toHaveBeenCalled(); + expect(underline).not.toHaveBeenCalled(); + }); }); diff --git a/extensions/processes-logs/components/log-file-viewer.ts b/extensions/processes-logs/components/log-file-viewer.ts index 1101401..edd2c41 100644 --- a/extensions/processes-logs/components/log-file-viewer.ts +++ b/extensions/processes-logs/components/log-file-viewer.ts @@ -18,6 +18,7 @@ export class LogFileViewer { private searchMatches: number[] = []; private searchCurrentMatch = -1; private centerTarget: number | null = null; + private readonly notifyLines = new Set(); constructor( initialLines: ProcessLogLine[], @@ -104,6 +105,23 @@ export class LogFileViewer { this.searchCurrentMatch = -1; } + /** + * Records a notify log-match marker. Lines whose text equals the matched + * line are highlighted distinctly from manual search matches and with lower + * priority (search current match > search match > notify match > stream). + */ + addNotifyMatch(match: { line: string }): void { + if (match.line) this.notifyLines.add(match.line); + } + + clearNotifyMatches(): void { + this.notifyLines.clear(); + } + + getNotifyMatchCount(): number { + return this.notifyLines.size; + } + nextMatch(): void { if (this.searchMatches.length === 0) return; this.searchCurrentMatch = @@ -166,6 +184,9 @@ export class LogFileViewer { if (matchSet.has(visibleIndex)) { return truncateToWidth(this.theme.fg("warning", text), width); } + if (this.notifyLines.has(line.text)) { + return truncateToWidth(this.theme.underline(text), width); + } if (line.type === "stderr") { return truncateToWidth(this.theme.fg("warning", text), width); } diff --git a/extensions/processes-logs/components/log-overlay-component.ts b/extensions/processes-logs/components/log-overlay-component.ts index 025fd58..41a3ecf 100644 --- a/extensions/processes-logs/components/log-overlay-component.ts +++ b/extensions/processes-logs/components/log-overlay-component.ts @@ -13,6 +13,7 @@ import { CHANNELS, type ProcessesChangedPayload, type ProcessProtocolConfig, + type ProcessProtocolNotificationPayload, } from "../../../src/protocol"; import { LIVE_STATUSES, type ProcessInfo } from "../../../src/types"; import { formatRuntime, truncateCmd } from "../../../src/utils/format"; @@ -42,6 +43,15 @@ const MIN_OVERLAY_WIDTH = 80; const MIN_OVERLAY_HEIGHT = 12; const OVERLAY_FRACTION = 0.9; const MAX_TAB_NAME = 12; +const MAX_NOTIFY_MARKERS_PER_PROCESS = 100; + +interface NotifyMatchMark { + pattern: string; + line: string; + stream: "stdout" | "stderr"; + matcherIndex: number; + timestamp: number; +} export class LogOverlayComponent implements Component { private processes: ProcessInfo[] = []; @@ -55,6 +65,8 @@ export class LogOverlayComponent implements Component { private hasSeenRunningProcess = false; private readonly disposers: Array<() => void> = []; private disposed = false; + /** Per-process notify log-match markers, capped per process. */ + private readonly notifyMarkers = new Map(); constructor(private readonly opts: LogOverlayOptions) { this.configureSearchInput(); @@ -65,6 +77,35 @@ export class LogOverlayComponent implements Component { if (isChangedPayload(payload)) this.handleProcessesChanged(payload); }), ); + this.disposers.push( + opts.events.on(CHANNELS.NOTIFICATION, (payload) => { + this.handleNotification(payload); + }), + ); + } + + private handleNotification(payload: unknown): void { + if (!isLogMatchNotification(payload)) return; + const mark: NotifyMatchMark = { + pattern: payload.logMatch.pattern, + line: payload.logMatch.line, + stream: payload.logMatch.stream, + matcherIndex: payload.logMatch.matcherIndex, + timestamp: payload.timestamp, + }; + + const list = this.notifyMarkers.get(payload.processId) ?? []; + list.push(mark); + if (list.length > MAX_NOTIFY_MARKERS_PER_PROCESS) { + list.splice(0, list.length - MAX_NOTIFY_MARKERS_PER_PROCESS); + } + this.notifyMarkers.set(payload.processId, list); + + const selected = this.selectedProcess(); + if (selected && selected.id === payload.processId) { + this.viewer?.addNotifyMatch(mark); + this.opts.tui.requestRender(); + } } render(width: number): string[] { @@ -338,6 +379,8 @@ export class LogOverlayComponent implements Component { followEnabled: this.opts.config.follow.enabledByDefault, maxBufferLines: this.opts.config.output.maxOutputLines, }); + const stored = this.notifyMarkers.get(selected.id); + if (stored) for (const mark of stored) this.viewer.addNotifyMatch(mark); connection.onChunk((lines: ProcessLogLine[]) => { this.viewer?.appendLines(lines); this.opts.tui.requestRender(); @@ -583,3 +626,23 @@ function isChangedPayload( payload.reason === "cleared") ); } + +function isLogMatchNotification( + payload: unknown, +): payload is ProcessProtocolNotificationPayload & { + kind: "log_match"; + logMatch: NonNullable; +} { + if (!isRecord(payload)) return false; + if (payload.kind !== "log_match") return false; + if (typeof payload.processId !== "string") return false; + if (typeof payload.timestamp !== "number") return false; + const logMatch = payload.logMatch; + if (!isRecord(logMatch)) return false; + if (typeof logMatch.pattern !== "string") return false; + if (typeof logMatch.line !== "string") return false; + if (logMatch.stream !== "stdout" && logMatch.stream !== "stderr") + return false; + if (typeof logMatch.matcherIndex !== "number") return false; + return true; +} diff --git a/extensions/processes/handlers/notifications.test.ts b/extensions/processes/handlers/notifications.test.ts new file mode 100644 index 0000000..82da080 --- /dev/null +++ b/extensions/processes/handlers/notifications.test.ts @@ -0,0 +1,86 @@ +import { createEventBus } from "@earendil-works/pi-coding-agent"; +import { describe, expect, it, vi } from "vitest"; +import type { ProcessProtocolNotificationPayload } from "../../../src/protocol"; +import { CHANNELS } from "../../../src/protocol"; +import { MESSAGE_TYPE_PROCESS_NOTIFICATION } from "../constants"; +import { registerNotificationDelivery } from "./notifications"; + +function makePayload( + overrides: Partial = {}, +): ProcessProtocolNotificationPayload { + return { + kind: "failure", + processId: "proc_1", + processName: "dev", + command: "pnpm dev", + timestamp: 123, + summary: "Process failed.", + status: "exited", + exitCode: 1, + endReason: "exit", + signal: null, + attention: "turn", + ...overrides, + }; +} + +function piWithSendMessage(sendMessage: ReturnType) { + return { sendMessage } as never; +} + +describe("registerNotificationDelivery", () => { + it("sends a displayed custom message with attention-derived options", () => { + const events = createEventBus(); + const sendMessage = vi.fn(); + registerNotificationDelivery(events, piWithSendMessage(sendMessage)); + + events.emit(CHANNELS.NOTIFICATION, makePayload({ attention: "turn" })); + + expect(sendMessage).toHaveBeenCalledTimes(1); + const [message, options] = sendMessage.mock.calls[0]; + expect(message.customType).toBe(MESSAGE_TYPE_PROCESS_NOTIFICATION); + expect(message.display).toBe(true); + expect(message.details.attention).toBe("turn"); + expect(options.triggerTurn).toBe(true); + expect(options.deliverAs).toBe("steer"); + }); + + it("maps context attention to a non-turn steer message", () => { + const events = createEventBus(); + const sendMessage = vi.fn(); + registerNotificationDelivery(events, piWithSendMessage(sendMessage)); + + events.emit(CHANNELS.NOTIFICATION, makePayload({ attention: "context" })); + + const [, options] = sendMessage.mock.calls[0]; + expect(options.triggerTurn).toBe(false); + expect(options.deliverAs).toBe("steer"); + }); + + it("ignores malformed payloads", () => { + const events = createEventBus(); + const sendMessage = vi.fn(); + registerNotificationDelivery(events, piWithSendMessage(sendMessage)); + + events.emit(CHANNELS.NOTIFICATION, null); + events.emit(CHANNELS.NOTIFICATION, {}); + events.emit(CHANNELS.NOTIFICATION, { kind: "success" }); + events.emit(CHANNELS.NOTIFICATION, { ...makePayload(), attention: 5 }); + + expect(sendMessage).not.toHaveBeenCalled(); + }); + + it("stops delivering after the disposer is called", () => { + const events = createEventBus(); + const sendMessage = vi.fn(); + const dispose = registerNotificationDelivery( + events, + piWithSendMessage(sendMessage), + ); + + dispose(); + events.emit(CHANNELS.NOTIFICATION, makePayload()); + + expect(sendMessage).not.toHaveBeenCalled(); + }); +}); diff --git a/extensions/processes/handlers/notifications.ts b/extensions/processes/handlers/notifications.ts new file mode 100644 index 0000000..af28880 --- /dev/null +++ b/extensions/processes/handlers/notifications.ts @@ -0,0 +1,48 @@ +import type { EventBus, ExtensionAPI } from "@earendil-works/pi-coding-agent"; + +import { + CHANNELS, + type ProcessProtocolNotificationPayload, +} from "../../../src/protocol"; +import { isRecord } from "../../../src/utils/is-record"; +import { + attentionToSendOptions, + sendProcessNotificationMessage, +} from "../notification-sender"; + +/** + * Delivers notification events emitted on {@link CHANNELS.NOTIFICATION} to Pi as + * persisted custom messages. This is the core side of the notification fanout: + * the NotificationService emits language-neutral payloads, and this listener + * converts each payload into a displayed `ad-process:notification` message with + * the attention-derived send options. UI extensions observe the same channel + * for display concerns (e.g. log-match highlighting) without importing this + * module. + * + * Returns a disposer that removes the listener; it must be called on + * `session_shutdown` before the manager is killed. + */ +export function registerNotificationDelivery( + events: EventBus, + pi: ExtensionAPI, +): () => void { + return events.on(CHANNELS.NOTIFICATION, (payload: unknown) => { + if (!isNotificationPayload(payload)) return; + const options = attentionToSendOptions(payload.attention); + sendProcessNotificationMessage(pi, payload, options); + }); +} + +function isNotificationPayload( + payload: unknown, +): payload is ProcessProtocolNotificationPayload { + if (!isRecord(payload)) return false; + if (typeof payload.kind !== "string") return false; + if (typeof payload.processId !== "string") return false; + if (typeof payload.processName !== "string") return false; + if (typeof payload.command !== "string") return false; + if (typeof payload.timestamp !== "number") return false; + if (typeof payload.summary !== "string") return false; + if (typeof payload.attention !== "string") return false; + return true; +} diff --git a/extensions/processes/index.ts b/extensions/processes/index.ts index b39b2d6..92a5013 100644 --- a/extensions/processes/index.ts +++ b/extensions/processes/index.ts @@ -2,6 +2,7 @@ import type { ExtensionAPI } from "@earendil-works/pi-coding-agent"; import { getManager } from "../../src/get-manager"; import { configLoader } from "./config"; import { registerCommandHandlers } from "./handlers/commands"; +import { registerNotificationDelivery } from "./handlers/notifications"; import { registerRequestHandlers } from "./handlers/requests"; import { registerLogSubscriptions } from "./handlers/subscriptions"; import { registerBackgroundBlocker } from "./hooks/background-blocker"; @@ -33,7 +34,7 @@ export default async function processesExtension( }); const notifications = createNotificationRegistry(); const notificationService = createNotificationService({ - pi, + events: pi.events, manager, registry: notifications, getProcess: (id) => manager.get(id), @@ -46,6 +47,7 @@ export default async function processesExtension( registerRequestHandlers(pi.events, manager, getConfig), registerCommandHandlers(pi.events, manager, notifications), registerLogSubscriptions(pi.events, manager), + registerNotificationDelivery(pi.events, pi), ]; registerBackgroundBlocker( diff --git a/extensions/processes/notification-sender.ts b/extensions/processes/notification-sender.ts index 8aae4fa..ad3bd68 100644 --- a/extensions/processes/notification-sender.ts +++ b/extensions/processes/notification-sender.ts @@ -2,13 +2,30 @@ import type { ExtensionAPI } from "@earendil-works/pi-coding-agent"; import { MESSAGE_TYPE_PROCESS_NOTIFICATION } from "./constants"; import { buildProcessNotificationContent } from "./notifications/render-content"; -import type { ProcessNotificationDetails } from "./notifications/types"; +import type { + Attention, + ProcessNotificationDetails, +} from "./notifications/types"; export interface ProcessNotificationSendOptions { triggerTurn: boolean; deliverAs: "steer" | "followUp" | "nextTurn"; } +/** Maps a notification attention level to Pi send-message options. */ +export function attentionToSendOptions( + attention: Attention, +): ProcessNotificationSendOptions { + switch (attention) { + case "turn": + return { triggerTurn: true, deliverAs: "steer" }; + case "context": + return { triggerTurn: false, deliverAs: "steer" }; + case "ignore": + return { triggerTurn: false, deliverAs: "steer" }; + } +} + export function sendProcessNotificationMessage( pi: ExtensionAPI, details: ProcessNotificationDetails, diff --git a/extensions/processes/notifications/service.test.ts b/extensions/processes/notifications/service.test.ts index b3c2e3f..faf2973 100644 --- a/extensions/processes/notifications/service.test.ts +++ b/extensions/processes/notifications/service.test.ts @@ -1,5 +1,7 @@ -import { describe, expect, it, vi } from "vitest"; - +import { createEventBus } from "@earendil-works/pi-coding-agent"; +import { describe, expect, it } from "vitest"; +import type { ProcessProtocolNotificationPayload } from "../../../src/protocol"; +import { CHANNELS } from "../../../src/protocol"; import type { ProcessInfo } from "../../../src/types"; import { flushQueuedMicrotasks } from "../../../tests/utils/async"; @@ -27,6 +29,8 @@ function makeInfo(overrides: Partial = {}): ProcessInfo { }; } +const processes = new Map(); + function createFakeManager() { const listeners: Array<(event: unknown) => void> = []; @@ -49,58 +53,65 @@ function createFakeManager() { }; } -const processes = new Map(); - -function createFakePi() { - return { - sendMessage: vi.fn(), - }; +/** + * Real in-memory event bus with a spy listener on CHANNELS.NOTIFICATION so + * tests assert on the payloads the service fans out rather than on Pi's + * sendMessage (now handled by the delivery listener). + */ +function createNotificationSpy(): { + events: ReturnType; + emitted: ProcessProtocolNotificationPayload[]; +} { + const events = createEventBus(); + const emitted: ProcessProtocolNotificationPayload[] = []; + events.on(CHANNELS.NOTIFICATION, (payload: unknown) => { + emitted.push(payload as ProcessProtocolNotificationPayload); + }); + return { events, emitted }; } describe("NotificationService", () => { - it("sends a turn notification for a failed process with default config", async () => { + it("emits a turn notification for a failed process with default config", async () => { const fakeManager = createFakeManager(); - const fakePi = createFakePi(); + const spy = createNotificationSpy(); const registry = createNotificationRegistry(); registry.register("proc_1", {}); const service = createNotificationService({ - pi: fakePi as never, + events: spy.events, manager: fakeManager as never, registry, getProcess: (id) => processes.get(id) ?? null, }); - const info = makeInfo({ - id: "proc_1", - success: false, - exitCode: 1, - endReason: "exit", + fakeManager.emit({ + type: "process_ended", + info: makeInfo({ + id: "proc_1", + success: false, + exitCode: 1, + endReason: "exit", + }), }); - - fakeManager.emit({ type: "process_ended", info }); await flushQueuedMicrotasks(); - expect(fakePi.sendMessage).toHaveBeenCalledTimes(1); - const [message, options] = fakePi.sendMessage.mock.calls[0]; - expect(message.customType).toBe("ad-process:notification"); - expect(message.display).toBe(true); - expect(options.triggerTurn).toBe(true); - expect(options.deliverAs).toBe("steer"); + expect(spy.emitted).toHaveLength(1); + expect(spy.emitted[0].attention).toBe("turn"); + expect(spy.emitted[0].processId).toBe("proc_1"); service.dispose(); }); - it("sends a context notification for a successful process with default config", async () => { + it("emits a context notification for a successful process with default config", async () => { const fakeManager = createFakeManager(); - const fakePi = createFakePi(); + const spy = createNotificationSpy(); const registry = createNotificationRegistry(); registry.register("proc_1", {}); const service = createNotificationService({ - pi: fakePi as never, + events: spy.events, manager: fakeManager as never, registry, getProcess: (id) => processes.get(id) ?? null, @@ -117,24 +128,22 @@ describe("NotificationService", () => { }); await flushQueuedMicrotasks(); - expect(fakePi.sendMessage).toHaveBeenCalledTimes(1); - const [, options] = fakePi.sendMessage.mock.calls[0]; - expect(options.triggerTurn).toBe(false); - expect(options.deliverAs).toBe("steer"); + expect(spy.emitted).toHaveLength(1); + expect(spy.emitted[0].attention).toBe("context"); service.dispose(); }); it("suppresses killed notification for intentional stop", async () => { const fakeManager = createFakeManager(); - const fakePi = createFakePi(); + const spy = createNotificationSpy(); const registry = createNotificationRegistry(); registry.register("proc_1", { onKilled: "ignore" }); registry.markIntentionalStop("proc_1"); const service = createNotificationService({ - pi: fakePi as never, + events: spy.events, manager: fakeManager as never, registry, getProcess: (id) => processes.get(id) ?? null, @@ -152,7 +161,7 @@ describe("NotificationService", () => { }); await flushQueuedMicrotasks(); - expect(fakePi.sendMessage).not.toHaveBeenCalled(); + expect(spy.emitted).toHaveLength(0); service.dispose(); }); @@ -165,14 +174,14 @@ describe("NotificationService", () => { exitCode, }) => { const fakeManager = createFakeManager(); - const fakePi = createFakePi(); + const spy = createNotificationSpy(); const registry = createNotificationRegistry(); registry.register("proc_1", { onSuccess: "turn", onFailure: "turn" }); registry.markIntentionalStop("proc_1"); const service = createNotificationService({ - pi: fakePi as never, + events: spy.events, manager: fakeManager as never, registry, getProcess: (id) => processes.get(id) ?? null, @@ -189,23 +198,22 @@ describe("NotificationService", () => { }); await flushQueuedMicrotasks(); - expect(fakePi.sendMessage).not.toHaveBeenCalled(); + expect(spy.emitted).toHaveLength(0); expect(registry.get("proc_1")).toBeNull(); service.dispose(); }); - it("sends notification for killed process when not intentional", async () => { + it("does not emit for killed process with default config when not intentional", async () => { const fakeManager = createFakeManager(); - const fakePi = createFakePi(); + const spy = createNotificationSpy(); const registry = createNotificationRegistry(); - // Default onKilled is "ignore", but killed is not forced display. - // So with default config and non-intentional kill, no message is sent. + // Default onKilled is "ignore", and killed is not forced display. registry.register("proc_1", {}); const service = createNotificationService({ - pi: fakePi as never, + events: spy.events, manager: fakeManager as never, registry, getProcess: (id) => processes.get(id) ?? null, @@ -223,21 +231,20 @@ describe("NotificationService", () => { }); await flushQueuedMicrotasks(); - // killed with ignore attention and not forced display = no message - expect(fakePi.sendMessage).not.toHaveBeenCalled(); + expect(spy.emitted).toHaveLength(0); service.dispose(); }); - it("forces display for crash even when attention is ignore", async () => { + it("forces emit for crash even when attention is ignore", async () => { const fakeManager = createFakeManager(); - const fakePi = createFakePi(); + const spy = createNotificationSpy(); const registry = createNotificationRegistry(); registry.register("proc_1", { onFailure: "ignore" }); const service = createNotificationService({ - pi: fakePi as never, + events: spy.events, manager: fakeManager as never, registry, getProcess: (id) => processes.get(id) ?? null, @@ -256,25 +263,22 @@ describe("NotificationService", () => { await flushQueuedMicrotasks(); // Crash forces display: attention becomes "context" - expect(fakePi.sendMessage).toHaveBeenCalledTimes(1); - const [message, options] = fakePi.sendMessage.mock.calls[0]; - expect(message.display).toBe(true); - expect(options.triggerTurn).toBe(false); - expect(options.deliverAs).toBe("steer"); - expect(message.details.kind).toBe("crash"); + expect(spy.emitted).toHaveLength(1); + expect(spy.emitted[0].kind).toBe("crash"); + expect(spy.emitted[0].attention).toBe("context"); service.dispose(); }); - it("forces display for timeout even when attention is ignore", async () => { + it("forces emit for timeout even when attention is ignore", async () => { const fakeManager = createFakeManager(); - const fakePi = createFakePi(); + const spy = createNotificationSpy(); const registry = createNotificationRegistry(); registry.register("proc_1", { onFailure: "ignore" }); const service = createNotificationService({ - pi: fakePi as never, + events: spy.events, manager: fakeManager as never, registry, getProcess: (id) => processes.get(id) ?? null, @@ -292,17 +296,16 @@ describe("NotificationService", () => { }); await flushQueuedMicrotasks(); - expect(fakePi.sendMessage).toHaveBeenCalledTimes(1); - const [message] = fakePi.sendMessage.mock.calls[0]; - expect(message.display).toBe(true); - expect(message.details.kind).toBe("timeout"); + expect(spy.emitted).toHaveLength(1); + expect(spy.emitted[0].kind).toBe("timeout"); + expect(spy.emitted[0].attention).toBe("context"); service.dispose(); }); - it("sends log match notification on output changed", () => { + it("emits log match notification on output changed", () => { const fakeManager = createFakeManager(); - const fakePi = createFakePi(); + const spy = createNotificationSpy(); const registry = createNotificationRegistry(); processes.set("proc_1", makeInfo({ id: "proc_1" })); @@ -311,7 +314,7 @@ describe("NotificationService", () => { }); const service = createNotificationService({ - pi: fakePi as never, + events: spy.events, manager: fakeManager as never, registry, getProcess: (id) => processes.get(id) ?? null, @@ -323,19 +326,18 @@ describe("NotificationService", () => { appendedText: [{ type: "stdout", text: "Server ready on port 3000" }], }); - expect(fakePi.sendMessage).toHaveBeenCalledTimes(1); - const [message, options] = fakePi.sendMessage.mock.calls[0]; - expect(message.details.kind).toBe("log_match"); - expect(options.triggerTurn).toBe(true); - expect(options.deliverAs).toBe("steer"); + expect(spy.emitted).toHaveLength(1); + expect(spy.emitted[0].kind).toBe("log_match"); + expect(spy.emitted[0].attention).toBe("turn"); + expect(spy.emitted[0].logMatch?.pattern).toBe("ready"); processes.delete("proc_1"); service.dispose(); }); - it("does not send log match notification when no appended text", () => { + it("does not emit log match notification when no appended text", () => { const fakeManager = createFakeManager(); - const fakePi = createFakePi(); + const spy = createNotificationSpy(); const registry = createNotificationRegistry(); registry.register("proc_1", { @@ -343,7 +345,7 @@ describe("NotificationService", () => { }); const service = createNotificationService({ - pi: fakePi as never, + events: spy.events, manager: fakeManager as never, registry, getProcess: (id) => processes.get(id) ?? null, @@ -355,18 +357,18 @@ describe("NotificationService", () => { appendedText: undefined, }); - expect(fakePi.sendMessage).not.toHaveBeenCalled(); + expect(spy.emitted).toHaveLength(0); service.dispose(); }); - it("does not send log match notification when no config registered", () => { + it("does not emit log match notification when no config registered", () => { const fakeManager = createFakeManager(); - const fakePi = createFakePi(); + const spy = createNotificationSpy(); const registry = createNotificationRegistry(); const service = createNotificationService({ - pi: fakePi as never, + events: spy.events, manager: fakeManager as never, registry, getProcess: (id) => processes.get(id) ?? null, @@ -378,21 +380,21 @@ describe("NotificationService", () => { appendedText: [{ type: "stdout", text: "ready" }], }); - expect(fakePi.sendMessage).not.toHaveBeenCalled(); + expect(spy.emitted).toHaveLength(0); service.dispose(); }); describe("disposal", () => { - it("does not send after dispose()", async () => { + it("does not emit after dispose()", async () => { const fakeManager = createFakeManager(); - const fakePi = createFakePi(); + const spy = createNotificationSpy(); const registry = createNotificationRegistry(); registry.register("proc_1", {}); const service = createNotificationService({ - pi: fakePi as never, + events: spy.events, manager: fakeManager as never, registry, getProcess: (id) => processes.get(id) ?? null, @@ -400,7 +402,6 @@ describe("NotificationService", () => { service.dispose(); - // Emit events after disposal fakeManager.emit({ type: "process_ended", info: makeInfo({ @@ -418,24 +419,23 @@ describe("NotificationService", () => { appendedText: [{ type: "stdout", text: "ready" }], }); - expect(fakePi.sendMessage).not.toHaveBeenCalled(); + expect(spy.emitted).toHaveLength(0); }); - it("does not send for events emitted during disposal sequence", async () => { + it("does not emit for events emitted during disposal sequence", async () => { const fakeManager = createFakeManager(); - const fakePi = createFakePi(); + const spy = createNotificationSpy(); const registry = createNotificationRegistry(); registry.register("proc_1", {}); const service = createNotificationService({ - pi: fakePi as never, + events: spy.events, manager: fakeManager as never, registry, getProcess: (id) => processes.get(id) ?? null, }); - // Dispose and then emit service.dispose(); fakeManager.emit({ @@ -449,16 +449,16 @@ describe("NotificationService", () => { }); await flushQueuedMicrotasks(); - expect(fakePi.sendMessage).not.toHaveBeenCalled(); + expect(spy.emitted).toHaveLength(0); }); it("dispose is idempotent", () => { const fakeManager = createFakeManager(); - const fakePi = createFakePi(); + const spy = createNotificationSpy(); const registry = createNotificationRegistry(); const service = createNotificationService({ - pi: fakePi as never, + events: spy.events, manager: fakeManager as never, registry, getProcess: (id) => processes.get(id) ?? null, @@ -467,18 +467,17 @@ describe("NotificationService", () => { service.dispose(); service.dispose(); - // No errors thrown - expect(fakePi.sendMessage).not.toHaveBeenCalled(); + expect(spy.emitted).toHaveLength(0); }); }); - it("sends default failure notification for unregistered process with no config", async () => { + it("emits default failure notification for unregistered process with no config", async () => { const fakeManager = createFakeManager(); - const fakePi = createFakePi(); + const spy = createNotificationSpy(); const registry = createNotificationRegistry(); const service = createNotificationService({ - pi: fakePi as never, + events: spy.events, manager: fakeManager as never, registry, getProcess: (id) => processes.get(id) ?? null, @@ -496,21 +495,19 @@ describe("NotificationService", () => { await flushQueuedMicrotasks(); // No config registered, but defaults resolve by kind: failure -> turn - expect(fakePi.sendMessage).toHaveBeenCalledTimes(1); - const [, options] = fakePi.sendMessage.mock.calls[0]; - expect(options.triggerTurn).toBe(true); - expect(options.deliverAs).toBe("steer"); + expect(spy.emitted).toHaveLength(1); + expect(spy.emitted[0].attention).toBe("turn"); service.dispose(); }); - it("sends context notification for unregistered successful process", async () => { + it("emits context notification for unregistered successful process", async () => { const fakeManager = createFakeManager(); - const fakePi = createFakePi(); + const spy = createNotificationSpy(); const registry = createNotificationRegistry(); const service = createNotificationService({ - pi: fakePi as never, + events: spy.events, manager: fakeManager as never, registry, getProcess: (id) => processes.get(id) ?? null, @@ -527,24 +524,21 @@ describe("NotificationService", () => { }); await flushQueuedMicrotasks(); - // No config registered, defaults resolve by kind: success -> context - expect(fakePi.sendMessage).toHaveBeenCalledTimes(1); - const [, options] = fakePi.sendMessage.mock.calls[0]; - expect(options.triggerTurn).toBe(false); - expect(options.deliverAs).toBe("steer"); + expect(spy.emitted).toHaveLength(1); + expect(spy.emitted[0].attention).toBe("context"); service.dispose(); }); it("unregisters process from registry after process_ended", async () => { const fakeManager = createFakeManager(); - const fakePi = createFakePi(); + const spy = createNotificationSpy(); const registry = createNotificationRegistry(); registry.register("proc_1", {}); const service = createNotificationService({ - pi: fakePi as never, + events: spy.events, manager: fakeManager as never, registry, getProcess: (id) => processes.get(id) ?? null, @@ -561,18 +555,19 @@ describe("NotificationService", () => { }); await flushQueuedMicrotasks(); + expect(spy.emitted).toHaveLength(1); expect(registry.get("proc_1")).toBeNull(); service.dispose(); }); - it("handles processes_changed events without sending notifications", () => { + it("handles processes_changed events without emitting notifications", () => { const fakeManager = createFakeManager(); - const fakePi = createFakePi(); + const spy = createNotificationSpy(); const registry = createNotificationRegistry(); const service = createNotificationService({ - pi: fakePi as never, + events: spy.events, manager: fakeManager as never, registry, getProcess: (id) => processes.get(id) ?? null, @@ -580,20 +575,20 @@ describe("NotificationService", () => { fakeManager.emit({ type: "processes_changed" }); - expect(fakePi.sendMessage).not.toHaveBeenCalled(); + expect(spy.emitted).toHaveLength(0); service.dispose(); }); it("uses custom onSuccess attention from config", async () => { const fakeManager = createFakeManager(); - const fakePi = createFakePi(); + const spy = createNotificationSpy(); const registry = createNotificationRegistry(); registry.register("proc_1", { onSuccess: "turn" }); const service = createNotificationService({ - pi: fakePi as never, + events: spy.events, manager: fakeManager as never, registry, getProcess: (id) => processes.get(id) ?? null, @@ -610,21 +605,19 @@ describe("NotificationService", () => { }); await flushQueuedMicrotasks(); - expect(fakePi.sendMessage).toHaveBeenCalledTimes(1); - const [, options] = fakePi.sendMessage.mock.calls[0]; - expect(options.triggerTurn).toBe(true); - expect(options.deliverAs).toBe("steer"); + expect(spy.emitted).toHaveLength(1); + expect(spy.emitted[0].attention).toBe("turn"); service.dispose(); }); - it("handles process_started event without sending notification", () => { + it("handles process_started event without emitting notification", () => { const fakeManager = createFakeManager(); - const fakePi = createFakePi(); + const spy = createNotificationSpy(); const registry = createNotificationRegistry(); const service = createNotificationService({ - pi: fakePi as never, + events: spy.events, manager: fakeManager as never, registry, getProcess: (id) => processes.get(id) ?? null, @@ -635,7 +628,7 @@ describe("NotificationService", () => { info: makeInfo(), }); - expect(fakePi.sendMessage).not.toHaveBeenCalled(); + expect(spy.emitted).toHaveLength(0); service.dispose(); }); @@ -643,18 +636,16 @@ describe("NotificationService", () => { describe("deferred process_ended and config registration race", () => { it("downgrades forced failure display to context when onFailure:ignore is registered before microtask", async () => { const fakeManager = createFakeManager(); - const fakePi = createFakePi(); + const spy = createNotificationSpy(); const registry = createNotificationRegistry(); const service = createNotificationService({ - pi: fakePi as never, + events: spy.events, manager: fakeManager as never, registry, getProcess: (id) => processes.get(id) ?? null, }); - // Simulate: manager.start() emits process_ended synchronously - // (e.g. missing_pid), then executeStart() registers config. fakeManager.emit({ type: "process_ended", info: makeInfo({ @@ -671,22 +662,20 @@ describe("NotificationService", () => { await flushQueuedMicrotasks(); // failure is forced display, so ignore is upgraded to context - expect(fakePi.sendMessage).toHaveBeenCalledTimes(1); - const [message, options] = fakePi.sendMessage.mock.calls[0]; - expect(message.display).toBe(true); - expect(options.triggerTurn).toBe(false); - expect(options.deliverAs).toBe("steer"); + expect(spy.emitted).toHaveLength(1); + expect(spy.emitted[0].kind).toBe("crash"); + expect(spy.emitted[0].attention).toBe("context"); service.dispose(); }); it("applies onKilled:context registered after emit but before microtask", async () => { const fakeManager = createFakeManager(); - const fakePi = createFakePi(); + const spy = createNotificationSpy(); const registry = createNotificationRegistry(); const service = createNotificationService({ - pi: fakePi as never, + events: spy.events, manager: fakeManager as never, registry, getProcess: (id) => processes.get(id) ?? null, @@ -708,21 +697,19 @@ describe("NotificationService", () => { await flushQueuedMicrotasks(); - expect(fakePi.sendMessage).toHaveBeenCalledTimes(1); - const [, options] = fakePi.sendMessage.mock.calls[0]; - expect(options.triggerTurn).toBe(false); - expect(options.deliverAs).toBe("steer"); + expect(spy.emitted).toHaveLength(1); + expect(spy.emitted[0].attention).toBe("context"); service.dispose(); }); - it("does not send if disposed before microtask fires", async () => { + it("does not emit if disposed before microtask fires", async () => { const fakeManager = createFakeManager(); - const fakePi = createFakePi(); + const spy = createNotificationSpy(); const registry = createNotificationRegistry(); const service = createNotificationService({ - pi: fakePi as never, + events: spy.events, manager: fakeManager as never, registry, getProcess: (id) => processes.get(id) ?? null, @@ -745,7 +732,7 @@ describe("NotificationService", () => { await flushQueuedMicrotasks(); - expect(fakePi.sendMessage).not.toHaveBeenCalled(); + expect(spy.emitted).toHaveLength(0); }); }); }); diff --git a/extensions/processes/notifications/service.ts b/extensions/processes/notifications/service.ts index 839f9c4..57014f5 100644 --- a/extensions/processes/notifications/service.ts +++ b/extensions/processes/notifications/service.ts @@ -1,8 +1,8 @@ -import type { ExtensionAPI } from "@earendil-works/pi-coding-agent"; +import type { EventBus } from "@earendil-works/pi-coding-agent"; import type { ProcessManager } from "../../../src/manager"; +import { CHANNELS } from "../../../src/protocol"; import type { ManagerEvent, ProcessInfo } from "../../../src/types"; -import { sendProcessNotificationMessage } from "../notification-sender"; import { classifyProcessEnd } from "./classify"; import { type CompiledLogMatcher, @@ -31,7 +31,14 @@ const DEFAULT_ATTENTION: Record< }; export interface NotificationServiceDeps { - pi: ExtensionAPI; + /** + * Event bus used to fan out notification events on CHANNELS.NOTIFICATION. + * The service never calls pi.sendMessage directly; a core delivery listener + * converts the emitted payload into a persisted custom message. This keeps + * notification flow event-driven and lets UI extensions observe the same + * events (e.g. for log-match highlighting). + */ + events: EventBus; manager: ProcessManager; registry: NotificationRegistry; getProcess: (id: string) => ProcessInfo | null; @@ -46,7 +53,7 @@ interface ProcessMatcherState { export function createNotificationService(deps: NotificationServiceDeps): { dispose: () => void; } { - const { pi, manager, registry, getProcess } = deps; + const { events, manager, registry, getProcess } = deps; let disposed = false; const matcherStates = new Map(); @@ -98,9 +105,7 @@ export function createNotificationService(deps: NotificationServiceDeps): { attention === "ignore" && shouldForceDisplay ? "context" : attention; const details = buildLifecycleDetails(info, kind, effectiveAttention); - const sendOptions = attentionToSendOptions(effectiveAttention); - - sendProcessNotificationMessage(pi, details, sendOptions); + events.emit(CHANNELS.NOTIFICATION, details); cleanupMatcherState(info.id); registry.unregister(info.id); @@ -135,8 +140,7 @@ export function createNotificationService(deps: NotificationServiceDeps): { attention, now, ); - const sendOptions = attentionToSendOptions(attention); - sendProcessNotificationMessage(pi, details, sendOptions); + events.emit(CHANNELS.NOTIFICATION, details); } } @@ -237,20 +241,6 @@ export function createNotificationService(deps: NotificationServiceDeps): { }; } - function attentionToSendOptions(attention: Attention): { - triggerTurn: boolean; - deliverAs: "steer" | "followUp" | "nextTurn"; - } { - switch (attention) { - case "turn": - return { triggerTurn: true, deliverAs: "steer" }; - case "context": - return { triggerTurn: false, deliverAs: "steer" }; - case "ignore": - return { triggerTurn: false, deliverAs: "steer" }; - } - } - function syncMatcherState( processId: string, watchState: WatchState, diff --git a/extensions/processes/notifications/types.ts b/extensions/processes/notifications/types.ts index d07274b..e7f1b6d 100644 --- a/extensions/processes/notifications/types.ts +++ b/extensions/processes/notifications/types.ts @@ -1,38 +1,11 @@ -import type { - ProcessEndReason, - ProcessSignalInfo, - ProcessStatus, -} from "../../../src/types"; - -export type Attention = "turn" | "context" | "ignore"; - -export type ProcessNotificationKind = - | "success" - | "failure" - | "crash" - | "killed" - | "timeout" - | "log_match"; - -export interface ProcessNotificationLogMatchDetails { - pattern: string; - mode: "literal" | "regex"; - stream: "stdout" | "stderr"; - line: string; - matcherIndex: number; -} - -export interface ProcessNotificationDetails { - kind: ProcessNotificationKind; - processId: string; - processName: string; - command: string; - timestamp: number; - summary: string; - status?: ProcessStatus; - exitCode?: number | null; - endReason?: ProcessEndReason | null; - signal?: ProcessSignalInfo | null; - logMatch?: ProcessNotificationLogMatchDetails; - attention: Attention; -} +// Re-export the protocol-safe notification types so the core extension and the +// protocol layer cannot drift apart. The canonical shape lives in +// `src/protocol/notifications.ts`; these names keep existing import sites +// stable while guaranteeing structural compatibility with the events emitted on +// CHANNELS.NOTIFICATION. +export type { + ProcessProtocolAttention as Attention, + ProcessProtocolNotificationKind as ProcessNotificationKind, + ProcessProtocolNotificationLogMatch as ProcessNotificationLogMatchDetails, + ProcessProtocolNotificationPayload as ProcessNotificationDetails, +} from "../../../src/protocol"; diff --git a/src/protocol/channels.ts b/src/protocol/channels.ts index 742c3d0..bf3e804 100644 --- a/src/protocol/channels.ts +++ b/src/protocol/channels.ts @@ -22,4 +22,7 @@ export const CHANNELS = { LOGS_SUBSCRIBE: "processes:logs:subscribe", LOGS_UNSUBSCRIBE: "processes:logs:unsubscribe", LOGS_CHUNK: "processes:logs:chunk", + + // Notification fanout (core emits, UI + core delivery listen) + NOTIFICATION: "processes:notification", } as const; diff --git a/src/protocol/index.ts b/src/protocol/index.ts index c87cc6e..7847db4 100644 --- a/src/protocol/index.ts +++ b/src/protocol/index.ts @@ -11,6 +11,12 @@ export type { LogsSubscribePayload, LogsUnsubscribePayload, } from "./logs"; +export type { + ProcessProtocolAttention, + ProcessProtocolNotificationKind, + ProcessProtocolNotificationLogMatch, + ProcessProtocolNotificationPayload, +} from "./notifications"; export type { ProcessProtocolConfig, RequestCombinedOutputPayload, diff --git a/src/protocol/notifications.ts b/src/protocol/notifications.ts new file mode 100644 index 0000000..49fc333 --- /dev/null +++ b/src/protocol/notifications.ts @@ -0,0 +1,48 @@ +import type { + ProcessEndReason, + ProcessSignalInfo, + ProcessStatus, +} from "../types"; + +/** + * Notification event payload broadcast on {@link CHANNELS.NOTIFICATION}. + * + * This is the protocol-safe mirror of the core extension's + * `ProcessNotificationDetails`. It is intentionally free of Pi imports so UI + * extensions (logs, dock) can observe notification events without importing + * core notification internals. The core extension emits this payload; the core + * delivery listener converts it back into a persisted custom message, while UI + * extensions use it for highlighting (e.g. log-match markers). + */ +export type ProcessProtocolAttention = "turn" | "context" | "ignore"; + +export type ProcessProtocolNotificationKind = + | "success" + | "failure" + | "crash" + | "killed" + | "timeout" + | "log_match"; + +export interface ProcessProtocolNotificationLogMatch { + pattern: string; + mode: "literal" | "regex"; + stream: "stdout" | "stderr"; + line: string; + matcherIndex: number; +} + +export interface ProcessProtocolNotificationPayload { + kind: ProcessProtocolNotificationKind; + processId: string; + processName: string; + command: string; + timestamp: number; + summary: string; + status?: ProcessStatus; + exitCode?: number | null; + endReason?: ProcessEndReason | null; + signal?: ProcessSignalInfo | null; + logMatch?: ProcessProtocolNotificationLogMatch; + attention: ProcessProtocolAttention; +}