diff --git a/PLAN.md b/PLAN.md index 333d84a..05b8e2c 100644 --- a/PLAN.md +++ b/PLAN.md @@ -95,13 +95,17 @@ Implemented and validated in Phase 2E: - Request handlers synchronously expose list/get/output/combined output/log files/file size/config. - Command handlers expose kill and clear over `pi.events`. - `/ps:kill` protocol kills are treated as intentional stops by reusing the same shared helper as `process stop`. +- Intentional stops suppress all lifecycle notifications, even when the process exits cleanly after a signal. - Log subscriptions return initial combined output and fan out live appended chunks to matching subscribers. +- Log subscriptions purge stale subscribers when a process ends or is removed by a clear operation. +- Protocol handlers ignore malformed payloads instead of throwing from shared event-bus listeners. +- Request/command reply callbacks are documented as an in-process protocol, not serializable IPC/RPC. - Session shutdown disposes protocol listeners and subscriptions before killing and cleaning up the manager. Latest validation: -- `pnpm typecheck` passes. - `pnpm lint` passes. -- `pnpm test` passes. +- `pnpm typecheck` passes. +- `pnpm test` passes with 213 tests. Current intentional gaps: - No settings/config loader yet. diff --git a/extensions/processes/handlers/commands.test.ts b/extensions/processes/handlers/commands.test.ts new file mode 100644 index 0000000..a12fcf7 --- /dev/null +++ b/extensions/processes/handlers/commands.test.ts @@ -0,0 +1,119 @@ +import { createEventBus } from "@earendil-works/pi-coding-agent"; +import { describe, expect, it, vi } from "vitest"; + +import type { ProcessManager } from "../../../src/manager"; +import { CHANNELS } from "../../../src/protocol"; +import type { KillResult, ProcessInfo } from "../../../src/types"; +import { createNotificationRegistry } from "../notifications/registry"; +import { registerCommandHandlers } from "./commands"; + +function makeInfo(overrides: Partial = {}): ProcessInfo { + return { + id: "proc_1", + name: "dev", + pid: 123, + command: "pnpm dev", + cwd: "/repo", + startTime: 1000, + endTime: 2000, + status: "killed", + exitCode: null, + success: false, + stdoutFile: "/tmp/stdout.log", + stderrFile: "/tmp/stderr.log", + endReason: "signal", + signal: null, + errorMessage: null, + ...overrides, + }; +} + +function flushPromises(): Promise { + return new Promise((resolve) => setTimeout(resolve, 0)); +} + +describe("registerCommandHandlers", () => { + it("kills processes as intentional stops", async () => { + const events = createEventBus(); + const registry = createNotificationRegistry(); + const result: KillResult = { ok: true, info: makeInfo() }; + let markedBeforeKill = false; + const manager = { + kill: vi.fn(async () => { + markedBeforeKill = registry.consumeIntentionalStop("proc_1"); + return result; + }), + } as unknown as ProcessManager; + const reply = vi.fn(); + + registerCommandHandlers(events, manager, registry); + events.emit(CHANNELS.COMMAND_KILL, { + id: "proc_1", + signal: "SIGKILL", + timeoutMs: 100, + reply, + }); + await flushPromises(); + + expect(markedBeforeKill).toBe(true); + expect(manager.kill).toHaveBeenCalledWith("proc_1", { + signal: "SIGKILL", + timeoutMs: 100, + }); + expect(reply).toHaveBeenCalledWith(result); + }); + + it("replies with error result when kill throws", async () => { + const events = createEventBus(); + const registry = createNotificationRegistry(); + const manager = { + kill: vi.fn(async () => { + throw new Error("boom"); + }), + } as unknown as ProcessManager; + const reply = vi.fn(); + + registerCommandHandlers(events, manager, registry); + events.emit(CHANNELS.COMMAND_KILL, { id: "proc_1", reply }); + await flushPromises(); + + expect(reply).toHaveBeenCalledWith( + expect.objectContaining({ + ok: false, + reason: "error", + info: expect.objectContaining({ id: "proc_1" }), + }), + ); + expect(registry.consumeIntentionalStop("proc_1")).toBe(false); + }); + + it("clears finished processes", () => { + const events = createEventBus(); + const registry = createNotificationRegistry(); + const manager = { + clearFinished: vi.fn(() => 2), + } as unknown as ProcessManager; + const reply = vi.fn(); + + registerCommandHandlers(events, manager, registry); + events.emit(CHANNELS.COMMAND_CLEAR, { reply }); + + expect(reply).toHaveBeenCalledWith(2); + }); + + it("disposes event listeners", async () => { + const events = createEventBus(); + const registry = createNotificationRegistry(); + const manager = { + kill: vi.fn(async () => ({ ok: true, info: makeInfo() }) as KillResult), + } as unknown as ProcessManager; + const reply = vi.fn(); + + const dispose = registerCommandHandlers(events, manager, registry); + dispose(); + events.emit(CHANNELS.COMMAND_KILL, { id: "proc_1", reply }); + await flushPromises(); + + expect(reply).not.toHaveBeenCalled(); + }); +}); diff --git a/extensions/processes/handlers/commands.ts b/extensions/processes/handlers/commands.ts new file mode 100644 index 0000000..837c5a9 --- /dev/null +++ b/extensions/processes/handlers/commands.ts @@ -0,0 +1,89 @@ +import type { EventBus } from "@earendil-works/pi-coding-agent"; + +import type { ProcessManager } from "../../../src/manager"; +import { + CHANNELS, + type CommandClearPayload, + type CommandKillPayload, +} from "../../../src/protocol"; +import type { KillResult } from "../../../src/types"; +import type { NotificationRegistry } from "../notifications/registry"; +import { killIntentionally } from "./kill-process"; + +export function registerCommandHandlers( + events: EventBus, + manager: ProcessManager, + notifications: NotificationRegistry, +): () => void { + const disposers = [ + events.on(CHANNELS.COMMAND_KILL, (payload) => { + const command = payload as CommandKillPayload; + if (!isCommandKillPayload(command)) return; + + void killIntentionally(manager, notifications, command.id, { + signal: command.signal, + timeoutMs: command.timeoutMs, + }).then(command.reply, () => { + command.reply(createKillErrorResult(command.id)); + }); + }), + events.on(CHANNELS.COMMAND_CLEAR, (payload) => { + const command = payload as CommandClearPayload; + if (!isCommandClearPayload(command)) return; + + command.reply(manager.clearFinished()); + }), + ]; + + return () => { + for (const dispose of disposers) dispose(); + }; +} + +function isCommandKillPayload( + payload: CommandKillPayload, +): payload is CommandKillPayload { + return ( + isRecord(payload) && typeof payload.id === "string" && isReply(payload) + ); +} + +function isCommandClearPayload( + payload: CommandClearPayload, +): payload is CommandClearPayload { + return isRecord(payload) && isReply(payload); +} + +function isReply( + payload: unknown, +): payload is { reply: (...args: never[]) => void } { + return isRecord(payload) && typeof payload.reply === "function"; +} + +function isRecord(payload: unknown): payload is Record { + return typeof payload === "object" && payload !== null; +} + +function createKillErrorResult(id: string): KillResult { + return { + ok: false, + reason: "error", + info: { + id, + name: "(unknown)", + pid: -1, + command: "", + cwd: "", + startTime: 0, + endTime: null, + status: "exited", + exitCode: null, + success: false, + stdoutFile: "", + stderrFile: "", + endReason: null, + signal: null, + errorMessage: "Failed to kill process", + }, + }; +} diff --git a/extensions/processes/handlers/kill-process.test.ts b/extensions/processes/handlers/kill-process.test.ts new file mode 100644 index 0000000..834b421 --- /dev/null +++ b/extensions/processes/handlers/kill-process.test.ts @@ -0,0 +1,107 @@ +import { describe, expect, it, vi } from "vitest"; + +import type { ProcessManager } from "../../../src/manager"; +import type { KillResult, ProcessInfo } from "../../../src/types"; +import { createNotificationRegistry } from "../notifications/registry"; +import { killIntentionally } from "./kill-process"; + +function makeInfo(overrides: Partial = {}): ProcessInfo { + return { + id: "proc_1", + name: "dev", + pid: 123, + command: "pnpm dev", + cwd: "/repo", + startTime: 1000, + endTime: 2000, + status: "killed", + exitCode: null, + success: false, + stdoutFile: "/tmp/stdout.log", + stderrFile: "/tmp/stderr.log", + endReason: "signal", + signal: null, + errorMessage: null, + ...overrides, + }; +} + +describe("killIntentionally", () => { + it("marks intentional stop before killing", async () => { + const registry = createNotificationRegistry(); + let markedBeforeKill = false; + const kill = vi.fn(async () => { + markedBeforeKill = registry.consumeIntentionalStop("proc_1"); + return { ok: true, info: makeInfo() } as KillResult; + }); + const manager = { kill } as unknown as ProcessManager; + + await killIntentionally(manager, registry, "proc_1", { signal: "SIGKILL" }); + + expect(markedBeforeKill).toBe(true); + expect(kill).toHaveBeenCalledWith("proc_1", { signal: "SIGKILL" }); + }); + + it("clears marker on not_found and error failures", async () => { + const cases: Array = [ + { ok: false, reason: "not_found", info: makeInfo() }, + { ok: false, reason: "error", info: makeInfo() }, + ]; + + for (const result of cases) { + const registry = createNotificationRegistry(); + const manager = { + kill: vi.fn(async () => result), + } as unknown as ProcessManager; + + await killIntentionally(manager, registry, "proc_1"); + + expect(registry.consumeIntentionalStop("proc_1")).toBe(false); + } + }); + + it("preserves marker on timeout failure", async () => { + const registry = createNotificationRegistry(); + const manager = { + kill: vi.fn( + async () => + ({ ok: false, reason: "timeout", info: makeInfo() }) as KillResult, + ), + } as unknown as ProcessManager; + + await killIntentionally(manager, registry, "proc_1"); + + expect(registry.consumeIntentionalStop("proc_1")).toBe(true); + }); + + it("clears marker for already-finished successful kill result", async () => { + const registry = createNotificationRegistry(); + const manager = { + kill: vi.fn( + async () => + ({ + ok: true, + info: makeInfo({ status: "exited", success: true }), + }) as KillResult, + ), + } as unknown as ProcessManager; + + await killIntentionally(manager, registry, "proc_1"); + + expect(registry.consumeIntentionalStop("proc_1")).toBe(false); + }); + + it("clears marker when kill throws", async () => { + const registry = createNotificationRegistry(); + const manager = { + kill: vi.fn(async () => { + throw new Error("boom"); + }), + } as unknown as ProcessManager; + + await expect( + killIntentionally(manager, registry, "proc_1"), + ).rejects.toThrow(/process stop failed/); + expect(registry.consumeIntentionalStop("proc_1")).toBe(false); + }); +}); diff --git a/extensions/processes/handlers/kill-process.ts b/extensions/processes/handlers/kill-process.ts new file mode 100644 index 0000000..390ce82 --- /dev/null +++ b/extensions/processes/handlers/kill-process.ts @@ -0,0 +1,31 @@ +import type { ProcessManager } from "../../../src/manager"; +import type { KillResult } from "../../../src/types"; +import { LIVE_STATUSES } from "../../../src/types"; +import type { NotificationRegistry } from "../notifications/registry"; + +export async function killIntentionally( + manager: ProcessManager, + notifications: NotificationRegistry, + id: string, + opts?: { signal?: NodeJS.Signals; timeoutMs?: number }, +): Promise { + notifications.markIntentionalStop(id); + + let result: KillResult; + try { + result = await manager.kill(id, opts); + } catch { + notifications.consumeIntentionalStop(id); + throw new Error(`process stop failed for ${id}`); + } + + if (!result.ok) { + if (result.reason === "not_found" || result.reason === "error") { + notifications.consumeIntentionalStop(id); + } + } else if (!LIVE_STATUSES.has(result.info.status)) { + notifications.consumeIntentionalStop(id); + } + + return result; +} diff --git a/extensions/processes/handlers/requests.test.ts b/extensions/processes/handlers/requests.test.ts new file mode 100644 index 0000000..768c300 --- /dev/null +++ b/extensions/processes/handlers/requests.test.ts @@ -0,0 +1,126 @@ +import { createEventBus } from "@earendil-works/pi-coding-agent"; +import { describe, expect, it, vi } from "vitest"; + +import type { ProcessManager } from "../../../src/manager"; +import { CHANNELS } from "../../../src/protocol"; +import type { ProcessInfo } from "../../../src/types"; +import { registerRequestHandlers } from "./requests"; + +function makeInfo(overrides: Partial = {}): ProcessInfo { + return { + id: "proc_1", + name: "dev", + pid: 123, + command: "pnpm dev", + cwd: "/repo", + startTime: 1000, + endTime: null, + status: "running", + exitCode: null, + success: null, + stdoutFile: "/tmp/stdout.log", + stderrFile: "/tmp/stderr.log", + endReason: null, + signal: null, + errorMessage: null, + ...overrides, + }; +} + +describe("registerRequestHandlers", () => { + it("replies to manager read requests", () => { + const events = createEventBus(); + const info = makeInfo(); + const manager = { + list: vi.fn(() => [info]), + get: vi.fn(() => info), + getOutput: vi.fn(() => ({ + stdout: ["out"], + stderr: ["err"], + status: "running", + })), + getCombinedOutput: vi.fn(() => [{ type: "stdout", text: "out" }]), + getLogFiles: vi.fn(() => ({ + stdoutFile: "/tmp/stdout.log", + stderrFile: "/tmp/stderr.log", + combinedFile: "/tmp/combined.log", + })), + getFileSize: vi.fn(() => ({ stdout: 1, stderr: 2 })), + } as unknown as ProcessManager; + + registerRequestHandlers(events, manager); + + const listReply = vi.fn(); + events.emit(CHANNELS.REQUEST_LIST, { reply: listReply }); + expect(listReply).toHaveBeenCalledWith([info]); + + const getReply = vi.fn(); + events.emit(CHANNELS.REQUEST_GET, { id: "proc_1", reply: getReply }); + expect(getReply).toHaveBeenCalledWith(info); + + const outputReply = vi.fn(); + events.emit(CHANNELS.REQUEST_OUTPUT, { + id: "proc_1", + tailLines: 5, + reply: outputReply, + }); + expect(manager.getOutput).toHaveBeenCalledWith("proc_1", 5); + expect(outputReply).toHaveBeenCalledWith({ + stdout: ["out"], + stderr: ["err"], + status: "running", + }); + + const combinedReply = vi.fn(); + events.emit(CHANNELS.REQUEST_COMBINED_OUTPUT, { + id: "proc_1", + tailLines: 3, + reply: combinedReply, + }); + expect(manager.getCombinedOutput).toHaveBeenCalledWith("proc_1", 3); + expect(combinedReply).toHaveBeenCalledWith([ + { type: "stdout", text: "out" }, + ]); + + const logsReply = vi.fn(); + events.emit(CHANNELS.REQUEST_LOG_FILES, { + id: "proc_1", + reply: logsReply, + }); + expect(logsReply).toHaveBeenCalledWith({ + stdoutFile: "/tmp/stdout.log", + stderrFile: "/tmp/stderr.log", + combinedFile: "/tmp/combined.log", + }); + + const sizeReply = vi.fn(); + events.emit(CHANNELS.REQUEST_FILE_SIZE, { + id: "proc_1", + reply: sizeReply, + }); + expect(sizeReply).toHaveBeenCalledWith({ stdout: 1, stderr: 2 }); + }); + + it("replies with temporary empty config", () => { + const events = createEventBus(); + const manager = {} as ProcessManager; + const reply = vi.fn(); + + registerRequestHandlers(events, manager); + events.emit(CHANNELS.REQUEST_CONFIG, { reply }); + + expect(reply).toHaveBeenCalledWith({}); + }); + + it("disposes event listeners", () => { + const events = createEventBus(); + const manager = { list: vi.fn(() => []) } as unknown as ProcessManager; + const reply = vi.fn(); + + const dispose = registerRequestHandlers(events, manager); + dispose(); + events.emit(CHANNELS.REQUEST_LIST, { reply }); + + expect(reply).not.toHaveBeenCalled(); + }); +}); diff --git a/extensions/processes/handlers/requests.ts b/extensions/processes/handlers/requests.ts new file mode 100644 index 0000000..b675a45 --- /dev/null +++ b/extensions/processes/handlers/requests.ts @@ -0,0 +1,108 @@ +import type { EventBus } from "@earendil-works/pi-coding-agent"; + +import type { ProcessManager } from "../../../src/manager"; +import { + CHANNELS, + type RequestCombinedOutputPayload, + type RequestConfigPayload, + type RequestFileSizePayload, + type RequestGetPayload, + type RequestListPayload, + type RequestLogFilesPayload, + type RequestOutputPayload, +} from "../../../src/protocol"; + +export function registerRequestHandlers( + events: EventBus, + manager: ProcessManager, +): () => void { + const disposers = [ + events.on(CHANNELS.REQUEST_LIST, (payload) => { + const request = payload as RequestListPayload; + if (!isRequestListPayload(request)) return; + + request.reply(manager.list()); + }), + events.on(CHANNELS.REQUEST_GET, (payload) => { + const request = payload as RequestGetPayload; + if (!isIdRequest(request)) return; + + request.reply(manager.get(request.id)); + }), + events.on(CHANNELS.REQUEST_OUTPUT, (payload) => { + const request = payload as RequestOutputPayload; + if (!isIdRequest(request)) return; + + request.reply(manager.getOutput(request.id, request.tailLines)); + }), + events.on(CHANNELS.REQUEST_COMBINED_OUTPUT, (payload) => { + const request = payload as RequestCombinedOutputPayload; + if (!isIdRequest(request)) return; + + request.reply(manager.getCombinedOutput(request.id, request.tailLines)); + }), + events.on(CHANNELS.REQUEST_LOG_FILES, (payload) => { + const request = payload as RequestLogFilesPayload; + if (!isIdRequest(request)) return; + + request.reply(manager.getLogFiles(request.id)); + }), + events.on(CHANNELS.REQUEST_FILE_SIZE, (payload) => { + const request = payload as RequestFileSizePayload; + if (!isIdRequest(request)) return; + + request.reply(manager.getFileSize(request.id)); + }), + events.on(CHANNELS.REQUEST_CONFIG, (payload) => { + const request = payload as RequestConfigPayload; + if (!isRequestConfigPayload(request)) return; + + // TODO(Phase 2F): return loaded process settings. + request.reply({}); + }), + ]; + + return () => { + for (const dispose of disposers) dispose(); + }; +} + +function isRequestListPayload( + payload: RequestListPayload, +): payload is RequestListPayload { + return isRecord(payload) && isReply(payload); +} + +function isRequestConfigPayload( + payload: RequestConfigPayload, +): payload is RequestConfigPayload { + return isRecord(payload) && isReply(payload); +} + +function isIdRequest( + payload: + | RequestGetPayload + | RequestOutputPayload + | RequestCombinedOutputPayload + | RequestLogFilesPayload + | RequestFileSizePayload, +): payload is + | RequestGetPayload + | RequestOutputPayload + | RequestCombinedOutputPayload + | RequestLogFilesPayload + | RequestFileSizePayload { + return ( + isRecord(payload) && typeof payload.id === "string" && isReply(payload) + ); +} + +function isReply( + payload: unknown, +): payload is { reply: (...args: never[]) => void } { + return isRecord(payload) && typeof payload.reply === "function"; +} + +function isRecord(payload: unknown): payload is Record { + return typeof payload === "object" && payload !== null; +} diff --git a/extensions/processes/handlers/subscriptions.test.ts b/extensions/processes/handlers/subscriptions.test.ts new file mode 100644 index 0000000..bbd6859 --- /dev/null +++ b/extensions/processes/handlers/subscriptions.test.ts @@ -0,0 +1,194 @@ +import { createEventBus } from "@earendil-works/pi-coding-agent"; +import { describe, expect, it, vi } from "vitest"; + +import type { ProcessManager } from "../../../src/manager"; +import { CHANNELS } from "../../../src/protocol"; +import type { ManagerEvent, ProcessInfo } from "../../../src/types"; +import { registerLogSubscriptions } from "./subscriptions"; + +function makeInfo(overrides: Partial = {}): ProcessInfo { + return { + id: "proc_1", + name: "dev", + pid: 123, + command: "pnpm dev", + cwd: "/repo", + startTime: 1000, + endTime: null, + status: "running", + exitCode: null, + success: null, + stdoutFile: "/tmp/stdout.log", + stderrFile: "/tmp/stderr.log", + endReason: null, + signal: null, + errorMessage: null, + ...overrides, + }; +} + +function createFakeManager() { + const listeners: Array<(event: ManagerEvent) => void> = []; + const manager = { + get: vi.fn((id: string) => (id === "proc_1" ? makeInfo({ id }) : null)), + getCombinedOutput: vi.fn(() => [{ type: "stdout", text: "initial" }]), + onEvent(listener: (event: ManagerEvent) => void): () => void { + listeners.push(listener); + return () => { + const index = listeners.indexOf(listener); + if (index >= 0) listeners.splice(index, 1); + }; + }, + } as unknown as ProcessManager; + + return { + manager, + emit(event: ManagerEvent): void { + for (const listener of listeners) listener(event); + }, + }; +} + +describe("registerLogSubscriptions", () => { + it("subscribes and returns initial combined output", () => { + const events = createEventBus(); + const fake = createFakeManager(); + const reply = vi.fn(); + + registerLogSubscriptions(events, fake.manager); + events.emit(CHANNELS.LOGS_SUBSCRIBE, { + subscriberId: "sub_1", + processId: "proc_1", + tailLines: 50, + reply, + }); + + expect(fake.manager.getCombinedOutput).toHaveBeenCalledWith("proc_1", 50); + expect(reply).toHaveBeenCalledWith({ + ok: true, + initialLines: [{ type: "stdout", text: "initial" }], + }); + }); + + it("rejects unknown processes", () => { + const events = createEventBus(); + const fake = createFakeManager(); + const reply = vi.fn(); + + registerLogSubscriptions(events, fake.manager); + events.emit(CHANNELS.LOGS_SUBSCRIBE, { + subscriberId: "sub_1", + processId: "missing", + reply, + }); + + expect(reply).toHaveBeenCalledWith({ + ok: false, + error: "Process not found", + }); + }); + + it("rejects unavailable process logs", () => { + const events = createEventBus(); + const fake = createFakeManager(); + const reply = vi.fn(); + + vi.mocked(fake.manager.getCombinedOutput).mockReturnValue(null); + + registerLogSubscriptions(events, fake.manager); + events.emit(CHANNELS.LOGS_SUBSCRIBE, { + subscriberId: "sub_1", + processId: "proc_1", + reply, + }); + + expect(reply).toHaveBeenCalledWith({ + ok: false, + error: "Process logs not found", + }); + }); + + it("emits chunks to matching subscribers only", () => { + const events = createEventBus(); + const fake = createFakeManager(); + const chunk = vi.fn(); + const appendedText = [{ type: "stderr" as const, text: "line" }]; + + registerLogSubscriptions(events, fake.manager); + events.on(CHANNELS.LOGS_CHUNK, chunk); + events.emit(CHANNELS.LOGS_SUBSCRIBE, { + subscriberId: "sub_1", + processId: "proc_1", + reply: vi.fn(), + }); + + fake.emit({ type: "process_output_changed", id: "other", appendedText }); + fake.emit({ type: "process_output_changed", id: "proc_1", appendedText }); + + expect(chunk).toHaveBeenCalledTimes(1); + expect(chunk).toHaveBeenCalledWith({ + subscriberId: "sub_1", + processId: "proc_1", + lines: appendedText, + }); + }); + + it("purges subscribers when processes end or disappear", () => { + const events = createEventBus(); + const fake = createFakeManager(); + const chunk = vi.fn(); + const appendedText = [{ type: "stdout" as const, text: "line" }]; + + registerLogSubscriptions(events, fake.manager); + events.on(CHANNELS.LOGS_CHUNK, chunk); + events.emit(CHANNELS.LOGS_SUBSCRIBE, { + subscriberId: "sub_1", + processId: "proc_1", + reply: vi.fn(), + }); + fake.emit({ type: "process_ended", info: makeInfo({ id: "proc_1" }) }); + fake.emit({ type: "process_output_changed", id: "proc_1", appendedText }); + + expect(chunk).not.toHaveBeenCalled(); + + events.emit(CHANNELS.LOGS_SUBSCRIBE, { + subscriberId: "sub_1", + processId: "proc_1", + reply: vi.fn(), + }); + vi.mocked(fake.manager.get).mockReturnValue(null); + fake.emit({ type: "processes_changed" }); + fake.emit({ type: "process_output_changed", id: "proc_1", appendedText }); + + expect(chunk).not.toHaveBeenCalled(); + }); + + it("unsubscribes and disposes subscriptions", () => { + const events = createEventBus(); + const fake = createFakeManager(); + const chunk = vi.fn(); + const appendedText = [{ type: "stdout" as const, text: "line" }]; + + const dispose = registerLogSubscriptions(events, fake.manager); + events.on(CHANNELS.LOGS_CHUNK, chunk); + events.emit(CHANNELS.LOGS_SUBSCRIBE, { + subscriberId: "sub_1", + processId: "proc_1", + reply: vi.fn(), + }); + events.emit(CHANNELS.LOGS_UNSUBSCRIBE, { subscriberId: "sub_1" }); + fake.emit({ type: "process_output_changed", id: "proc_1", appendedText }); + + expect(chunk).not.toHaveBeenCalled(); + + events.emit(CHANNELS.LOGS_SUBSCRIBE, { + subscriberId: "sub_1", + processId: "proc_1", + reply: vi.fn(), + }); + dispose(); + fake.emit({ type: "process_output_changed", id: "proc_1", appendedText }); + + expect(chunk).not.toHaveBeenCalled(); + }); +}); diff --git a/extensions/processes/handlers/subscriptions.ts b/extensions/processes/handlers/subscriptions.ts new file mode 100644 index 0000000..45258c8 --- /dev/null +++ b/extensions/processes/handlers/subscriptions.ts @@ -0,0 +1,133 @@ +import type { EventBus } from "@earendil-works/pi-coding-agent"; + +import type { ProcessManager } from "../../../src/manager"; +import { + CHANNELS, + type LogsSubscribePayload, + type LogsUnsubscribePayload, +} from "../../../src/protocol"; + +interface LogSubscriber { + subscriberId: string; + processId: string; +} + +export function registerLogSubscriptions( + events: EventBus, + manager: ProcessManager, +): () => void { + const subscribers = new Map(); + + const disposeSubscribe = events.on(CHANNELS.LOGS_SUBSCRIBE, (payload) => { + const request = payload as LogsSubscribePayload; + if (!isLogsSubscribePayload(request)) return; + + const processInfo = manager.get(request.processId); + + if (!processInfo) { + request.reply({ ok: false, error: "Process not found" }); + return; + } + + const initialLines = manager.getCombinedOutput( + request.processId, + request.tailLines ?? 100, + ); + + if (!initialLines) { + request.reply({ ok: false, error: "Process logs not found" }); + return; + } + + subscribers.set(request.subscriberId, { + subscriberId: request.subscriberId, + processId: request.processId, + }); + + request.reply({ ok: true, initialLines }); + }); + + const disposeUnsubscribe = events.on(CHANNELS.LOGS_UNSUBSCRIBE, (payload) => { + const request = payload as LogsUnsubscribePayload; + if (!isLogsUnsubscribePayload(request)) return; + + subscribers.delete(request.subscriberId); + }); + + const disposeManager = manager.onEvent((event) => { + if (event.type === "process_ended") { + removeSubscribersForProcess(subscribers, event.info.id); + return; + } + + if (event.type === "processes_changed") { + removeStaleSubscribers(subscribers, manager); + return; + } + + if (event.type !== "process_output_changed") return; + if (!event.appendedText || event.appendedText.length === 0) return; + + for (const subscriber of subscribers.values()) { + if (subscriber.processId !== event.id) continue; + + events.emit(CHANNELS.LOGS_CHUNK, { + subscriberId: subscriber.subscriberId, + processId: subscriber.processId, + lines: event.appendedText, + }); + } + }); + + return () => { + disposeSubscribe(); + disposeUnsubscribe(); + disposeManager(); + subscribers.clear(); + }; +} + +function removeSubscribersForProcess( + subscribers: Map, + processId: string, +): void { + for (const [subscriberId, subscriber] of subscribers.entries()) { + if (subscriber.processId === processId) subscribers.delete(subscriberId); + } +} + +function removeStaleSubscribers( + subscribers: Map, + manager: ProcessManager, +): void { + for (const [subscriberId, subscriber] of subscribers.entries()) { + if (!manager.get(subscriber.processId)) subscribers.delete(subscriberId); + } +} + +function isLogsSubscribePayload( + payload: LogsSubscribePayload, +): payload is LogsSubscribePayload { + return ( + isRecord(payload) && + typeof payload.subscriberId === "string" && + typeof payload.processId === "string" && + isReply(payload) + ); +} + +function isLogsUnsubscribePayload( + payload: LogsUnsubscribePayload, +): payload is LogsUnsubscribePayload { + return isRecord(payload) && typeof payload.subscriberId === "string"; +} + +function isReply( + payload: unknown, +): payload is { reply: (...args: never[]) => void } { + return isRecord(payload) && typeof payload.reply === "function"; +} + +function isRecord(payload: unknown): payload is Record { + return typeof payload === "object" && payload !== null; +} diff --git a/extensions/processes/hooks/cleanup.test.ts b/extensions/processes/hooks/cleanup.test.ts new file mode 100644 index 0000000..515a314 --- /dev/null +++ b/extensions/processes/hooks/cleanup.test.ts @@ -0,0 +1,49 @@ +import type { ExtensionAPI } from "@earendil-works/pi-coding-agent"; +import { describe, expect, it, vi } from "vitest"; + +import { createNotificationRegistry } from "../notifications/registry"; +import { registerCleanupHook } from "./cleanup"; + +describe("registerCleanupHook", () => { + it("runs disposers before manager cleanup and ignores duplicate shutdown", async () => { + let shutdown: () => Promise | void = () => { + throw new Error("shutdown handler was not registered"); + }; + const calls: string[] = []; + const pi = { + on: vi.fn((event: string, handler: () => Promise | void) => { + if (event === "session_shutdown") shutdown = handler; + }), + } as unknown as ExtensionAPI; + const manager = { + killAll: vi.fn(() => calls.push("killAll")), + cleanup: vi.fn(() => calls.push("cleanup")), + }; + const notificationService = { + dispose: vi.fn(() => calls.push("notificationService.dispose")), + }; + const notifications = createNotificationRegistry(); + const disposer = vi.fn(() => calls.push("disposer")); + + registerCleanupHook(pi, { + manager: manager as never, + notifications, + notificationService, + disposers: [disposer], + }); + + await shutdown(); + await shutdown(); + + expect(disposer).toHaveBeenCalledTimes(1); + expect(notificationService.dispose).toHaveBeenCalledTimes(1); + expect(manager.killAll).toHaveBeenCalledTimes(1); + expect(manager.cleanup).toHaveBeenCalledTimes(1); + expect(calls).toEqual([ + "disposer", + "notificationService.dispose", + "killAll", + "cleanup", + ]); + }); +}); diff --git a/extensions/processes/hooks/cleanup.ts b/extensions/processes/hooks/cleanup.ts index 7c3d0f4..883c75f 100644 --- a/extensions/processes/hooks/cleanup.ts +++ b/extensions/processes/hooks/cleanup.ts @@ -7,16 +7,32 @@ interface NotificationService { dispose(): void; } +type Disposer = () => void; + +interface CleanupHookDeps { + manager: ProcessManager; + notifications: NotificationRegistry; + notificationService: NotificationService; + disposers?: Disposer[]; +} + export function registerCleanupHook( pi: ExtensionAPI, - manager: ProcessManager, - notifications: NotificationRegistry, - notificationService: NotificationService, + deps: CleanupHookDeps, ): void { + let shuttingDown = false; + pi.on("session_shutdown", async () => { - notificationService.dispose(); - notifications.clear(); - manager.killAll(); - manager.cleanup(); + if (shuttingDown) return; + shuttingDown = true; + + for (const dispose of deps.disposers ?? []) { + dispose(); + } + + deps.notificationService.dispose(); + deps.notifications.clear(); + deps.manager.killAll(); + deps.manager.cleanup(); }); } diff --git a/extensions/processes/hooks/event-bridge.test.ts b/extensions/processes/hooks/event-bridge.test.ts new file mode 100644 index 0000000..55873e1 --- /dev/null +++ b/extensions/processes/hooks/event-bridge.test.ts @@ -0,0 +1,112 @@ +import { createEventBus } from "@earendil-works/pi-coding-agent"; +import { describe, expect, it, vi } from "vitest"; + +import type { ProcessManager } from "../../../src/manager"; +import { CHANNELS } from "../../../src/protocol"; +import type { ManagerEvent, ProcessInfo } from "../../../src/types"; +import { registerEventBridge } from "./event-bridge"; + +function makeInfo(overrides: Partial = {}): ProcessInfo { + return { + id: "proc_1", + name: "dev", + pid: 123, + command: "pnpm dev", + cwd: "/repo", + startTime: 1000, + endTime: null, + status: "running", + exitCode: null, + success: null, + stdoutFile: "/tmp/stdout.log", + stderrFile: "/tmp/stderr.log", + endReason: null, + signal: null, + errorMessage: null, + ...overrides, + }; +} + +function createFakeManager() { + const listeners: Array<(event: ManagerEvent) => void> = []; + + return { + manager: { + onEvent(listener: (event: ManagerEvent) => void): () => void { + listeners.push(listener); + return () => { + const index = listeners.indexOf(listener); + if (index >= 0) listeners.splice(index, 1); + }; + }, + } as unknown as ProcessManager, + emit(event: ManagerEvent): void { + for (const listener of listeners) listener(event); + }, + }; +} + +describe("registerEventBridge", () => { + it("bridges started events and changed notifications", () => { + const events = createEventBus(); + const fake = createFakeManager(); + const started = vi.fn(); + const changed = vi.fn(); + + events.on(CHANNELS.STARTED, started); + events.on(CHANNELS.CHANGED, changed); + registerEventBridge(events, fake.manager); + + const info = makeInfo(); + fake.emit({ type: "process_started", info }); + + expect(started).toHaveBeenCalledWith(info); + expect(changed).toHaveBeenCalledWith({ reason: "started" }); + }); + + it("bridges ended and output events", () => { + const events = createEventBus(); + const fake = createFakeManager(); + const ended = vi.fn(); + const output = vi.fn(); + + events.on(CHANNELS.ENDED, ended); + events.on(CHANNELS.OUTPUT_CHANGED, output); + registerEventBridge(events, fake.manager); + + const info = makeInfo({ status: "exited", success: true, exitCode: 0 }); + const appendedText = [{ type: "stdout" as const, text: "ready" }]; + fake.emit({ type: "process_ended", info }); + fake.emit({ type: "process_output_changed", id: "proc_1", appendedText }); + + expect(ended).toHaveBeenCalledWith(info); + expect(output).toHaveBeenCalledWith({ id: "proc_1", appendedText }); + }); + + it("bridges processes changed events", () => { + const events = createEventBus(); + const fake = createFakeManager(); + const changed = vi.fn(); + + events.on(CHANNELS.CHANGED, changed); + registerEventBridge(events, fake.manager); + + fake.emit({ type: "processes_changed" }); + + expect(changed).toHaveBeenCalledWith({ reason: "cleared" }); + }); + + it("disposes manager listener", () => { + const events = createEventBus(); + const fake = createFakeManager(); + const started = vi.fn(); + + events.on(CHANNELS.STARTED, started); + const dispose = registerEventBridge(events, fake.manager); + dispose(); + + fake.emit({ type: "process_started", info: makeInfo() }); + + expect(started).not.toHaveBeenCalled(); + }); +}); diff --git a/extensions/processes/hooks/event-bridge.ts b/extensions/processes/hooks/event-bridge.ts new file mode 100644 index 0000000..cd0db27 --- /dev/null +++ b/extensions/processes/hooks/event-bridge.ts @@ -0,0 +1,31 @@ +import type { EventBus } from "@earendil-works/pi-coding-agent"; + +import type { ProcessManager } from "../../../src/manager"; +import { CHANNELS } from "../../../src/protocol"; + +export function registerEventBridge( + events: EventBus, + manager: ProcessManager, +): () => void { + return manager.onEvent((event) => { + switch (event.type) { + case "process_started": + events.emit(CHANNELS.STARTED, event.info); + events.emit(CHANNELS.CHANGED, { reason: "started" }); + break; + case "process_ended": + events.emit(CHANNELS.ENDED, event.info); + events.emit(CHANNELS.CHANGED, { reason: "ended" }); + break; + case "process_output_changed": + events.emit(CHANNELS.OUTPUT_CHANGED, { + id: event.id, + appendedText: event.appendedText, + }); + break; + case "processes_changed": + events.emit(CHANNELS.CHANGED, { reason: "cleared" }); + break; + } + }); +} diff --git a/extensions/processes/index.ts b/extensions/processes/index.ts index 088fd51..c54401c 100644 --- a/extensions/processes/index.ts +++ b/extensions/processes/index.ts @@ -1,7 +1,11 @@ import type { ExtensionAPI } from "@earendil-works/pi-coding-agent"; import { getManager } from "../../src/get-manager"; +import { registerCommandHandlers } from "./handlers/commands"; +import { registerRequestHandlers } from "./handlers/requests"; +import { registerLogSubscriptions } from "./handlers/subscriptions"; import { registerCleanupHook } from "./hooks/cleanup"; +import { registerEventBridge } from "./hooks/event-bridge"; import { registerProcessNotificationRenderer } from "./message-renderer"; import { createNotificationRegistry, @@ -19,7 +23,19 @@ export default function processesExtension(pi: ExtensionAPI): void { getProcess: (id) => manager.get(id), }); + const disposers = [ + registerEventBridge(pi.events, manager), + registerRequestHandlers(pi.events, manager), + registerCommandHandlers(pi.events, manager, notifications), + registerLogSubscriptions(pi.events, manager), + ]; + registerProcessNotificationRenderer(pi); registerProcessTool(pi, manager, notifications); - registerCleanupHook(pi, manager, notifications, notificationService); + registerCleanupHook(pi, { + manager, + notifications, + notificationService, + disposers, + }); } diff --git a/extensions/processes/notifications/service.test.ts b/extensions/processes/notifications/service.test.ts index ec83a1c..b3c2e3f 100644 --- a/extensions/processes/notifications/service.test.ts +++ b/extensions/processes/notifications/service.test.ts @@ -157,6 +157,44 @@ describe("NotificationService", () => { service.dispose(); }); + it.each([ + { name: "successful", success: true, exitCode: 0 }, + { name: "failed", success: false, exitCode: 1 }, + ])("suppresses $name exit notification for intentional stop", async ({ + success, + exitCode, + }) => { + const fakeManager = createFakeManager(); + const fakePi = createFakePi(); + const registry = createNotificationRegistry(); + + registry.register("proc_1", { onSuccess: "turn", onFailure: "turn" }); + registry.markIntentionalStop("proc_1"); + + const service = createNotificationService({ + pi: fakePi as never, + manager: fakeManager as never, + registry, + getProcess: (id) => processes.get(id) ?? null, + }); + + fakeManager.emit({ + type: "process_ended", + info: makeInfo({ + id: "proc_1", + success, + exitCode, + endReason: "exit", + }), + }); + await flushQueuedMicrotasks(); + + expect(fakePi.sendMessage).not.toHaveBeenCalled(); + expect(registry.get("proc_1")).toBeNull(); + + service.dispose(); + }); + it("sends notification for killed process when not intentional", async () => { const fakeManager = createFakeManager(); const fakePi = createFakePi(); diff --git a/extensions/processes/notifications/service.ts b/extensions/processes/notifications/service.ts index 36c0083..f364e00 100644 --- a/extensions/processes/notifications/service.ts +++ b/extensions/processes/notifications/service.ts @@ -74,7 +74,7 @@ export function createNotificationService(deps: NotificationServiceDeps): { const isIntentionalStop = registry.consumeIntentionalStop(info.id); const kind = classifyProcessEnd(info); - if (isIntentionalStop && kind === "killed") { + if (isIntentionalStop) { cleanupMatcherState(info.id); registry.unregister(info.id); return; diff --git a/extensions/processes/tools/stop/index.test.ts b/extensions/processes/tools/stop/index.test.ts index ec8c89e..3c046b2 100644 --- a/extensions/processes/tools/stop/index.test.ts +++ b/extensions/processes/tools/stop/index.test.ts @@ -37,7 +37,7 @@ describe("executeStop", () => { await executeStop({ action: "stop", id: "proc_1" }, manager, registry); - expect(kill).toHaveBeenCalledWith("proc_1"); + expect(kill).toHaveBeenCalledWith("proc_1", undefined); }); it("marks intentional stop in registry before manager.kill is called", async () => { diff --git a/extensions/processes/tools/stop/index.ts b/extensions/processes/tools/stop/index.ts index 1f57fdb..8197022 100644 --- a/extensions/processes/tools/stop/index.ts +++ b/extensions/processes/tools/stop/index.ts @@ -1,6 +1,6 @@ import type { ProcessManager } from "../../../../src/manager"; import type { KillResult } from "../../../../src/types"; -import { LIVE_STATUSES } from "../../../../src/types"; +import { killIntentionally } from "../../handlers/kill-process"; import type { NotificationRegistry } from "../../notifications/registry"; import type { ProcessesParamsType } from "../schema"; @@ -18,23 +18,7 @@ export async function executeStop( throw new Error("process stop requires id"); } - notifications.markIntentionalStop(params.id); - - let result: KillResult; - try { - result = await manager.kill(params.id); - } catch { - notifications.consumeIntentionalStop(params.id); - throw new Error(`process stop failed for ${params.id}`); - } - - if (!result.ok) { - if (result.reason === "not_found" || result.reason === "error") { - notifications.consumeIntentionalStop(params.id); - } - } else if (!LIVE_STATUSES.has(result.info.status)) { - notifications.consumeIntentionalStop(params.id); - } + const result = await killIntentionally(manager, notifications, params.id); return { action: "stop", diff --git a/src/protocol.ts b/src/protocol.ts index ebc9e3a..e0031af 100644 --- a/src/protocol.ts +++ b/src/protocol.ts @@ -39,6 +39,8 @@ export type ProcessesChangedPayload = { }; // --- Request payloads (UI emits, core listens and calls reply synchronously) --- +// Reply callbacks make this an in-process protocol, not serializable IPC/RPC. +// Callers that need robustness should wrap requests with their own timeout. export interface RequestListPayload { reply: (processes: ProcessInfo[]) => void; @@ -103,6 +105,7 @@ export interface CommandClearPayload { export interface LogsSubscribePayload { subscriberId: string; processId: string; + tailLines?: number; reply: ( result: | {