返回源码地图

packages/coding-agent/src/modes/rpc/rpc-mode.ts

v1.0.0 · a13d35a742c6 · 10:JSONL响应disposition、事件和错误边界

完整原文供逐行核对;页面收录不代表每行都经过人工语义审核。MIT 许可见 许可证。

1/**
2 * RPC mode: Headless operation with JSON stdin/stdout protocol.
3 *
4 * Used for embedding the agent in other applications.
5 * Receives commands as JSON on stdin, outputs events and responses as JSON on stdout.
6 *
7 * Protocol:
8 * - Commands: JSON objects with `type` field, optional `id` for correlation
9 * - Responses: JSON objects with `type: "response"`, `command`, `success`, and optional `data`/`error`
10 * - Events: AgentSessionEvent objects streamed as they occur
11 * - Extension UI: Extension UI requests are emitted, client responds with extension_ui_response
12 */
13
14import * as crypto from "node:crypto";
15import type { AgentSessionRuntime } from "../../core/agent-session-runtime.ts";
16import type {
17 ExtensionUIContext,
18 ExtensionUIDialogOptions,
19 ExtensionWidgetOptions,
20 WorkingIndicatorOptions,
21} from "../../core/extensions/index.ts";
22import {
23 flushRawStdout,
24 takeOverStdout,
25 waitForRawStdoutBackpressure,
26 writeRawStdout,
27} from "../../core/output-guard.ts";
28import { killTrackedDetachedChildren } from "../../utils/shell.ts";
29import { type Theme, theme } from "../interactive/theme/theme.ts";
30import { toJsonEvent } from "../json-event.ts";
31import { attachJsonlLineReader, serializeJsonLine } from "./jsonl.ts";
32import type {
33 RpcCommand,
34 RpcExtensionUIRequest,
35 RpcExtensionUIResponse,
36 RpcResponse,
37 RpcSessionState,
38 RpcSlashCommand,
39} from "./rpc-types.ts";
40
41// Re-export types for consumers
42export type {
43 RpcCommand,
44 RpcExtensionUIRequest,
45 RpcExtensionUIResponse,
46 RpcResponse,
47 RpcSessionState,
48} from "./rpc-types.ts";
49
50/**
51 * Run in RPC mode.
52 * Listens for JSON commands on stdin, outputs events and responses on stdout.
53 */
54export async function runRpcMode(runtimeHost: AgentSessionRuntime): Promise<never> {
55 takeOverStdout();
56 let session = runtimeHost.session;
57 let unsubscribe: (() => void) | undefined;
58 let unsubscribeBackpressure: (() => void) | undefined;
59
60 const output = (obj: RpcResponse | RpcExtensionUIRequest | object) => {
61 writeRawStdout(serializeJsonLine(obj));
62 };
63
64 const success = <T extends RpcCommand["type"]>(
65 id: string | undefined,
66 command: T,
67 data?: object | null,
68 ): RpcResponse => {
69 if (data === undefined) {
70 return { id, type: "response", command, success: true } as RpcResponse;
71 }
72 return { id, type: "response", command, success: true, data } as RpcResponse;
73 };
74
75 const error = (id: string | undefined, command: string, message: string): RpcResponse => {
76 return { id, type: "response", command, success: false, error: message };
77 };
78
79 // Pending extension UI requests waiting for response
80 const pendingExtensionRequests = new Map<
81 string,
82 { resolve: (value: any) => void; reject: (error: Error) => void }
83 >();
84
85 // Shutdown request flag
86 let shutdownRequested = false;
87 let shuttingDown = false;
88 const signalCleanupHandlers: Array<() => void> = [];
89
90 /** Helper for dialog methods with signal/timeout support */
91 function createDialogPromise<T>(
92 opts: ExtensionUIDialogOptions | undefined,
93 defaultValue: T,
94 request: Record<string, unknown>,
95 parseResponse: (response: RpcExtensionUIResponse) => T,
96 ): Promise<T> {
97 if (opts?.signal?.aborted) return Promise.resolve(defaultValue);
98
99 const id = crypto.randomUUID();
100 return new Promise((resolve, reject) => {
101 let timeoutId: ReturnType<typeof setTimeout> | undefined;
102
103 const cleanup = () => {
104 if (timeoutId) clearTimeout(timeoutId);
105 opts?.signal?.removeEventListener("abort", onAbort);
106 pendingExtensionRequests.delete(id);
107 };
108
109 const onAbort = () => {
110 cleanup();
111 resolve(defaultValue);
112 };
113 opts?.signal?.addEventListener("abort", onAbort, { once: true });
114
115 if (opts?.timeout) {
116 timeoutId = setTimeout(() => {
117 cleanup();
118 resolve(defaultValue);
119 }, opts.timeout);
120 }
121
122 pendingExtensionRequests.set(id, {
123 resolve: (response: RpcExtensionUIResponse) => {
124 cleanup();
125 resolve(parseResponse(response));
126 },
127 reject,
128 });
129 output({ type: "extension_ui_request", id, ...request } as RpcExtensionUIRequest);
130 });
131 }
132
133 /**
134 * Create an extension UI context that uses the RPC protocol.
135 */
136 const createExtensionUIContext = (): ExtensionUIContext => ({
137 select: (title, options, opts) =>
138 createDialogPromise(opts, undefined, { method: "select", title, options, timeout: opts?.timeout }, (r) =>
139 "cancelled" in r && r.cancelled ? undefined : "value" in r ? r.value : undefined,
140 ),
141
142 confirm: (title, message, opts) =>
143 createDialogPromise(opts, false, { method: "confirm", title, message, timeout: opts?.timeout }, (r) =>
144 "cancelled" in r && r.cancelled ? false : "confirmed" in r ? r.confirmed : false,
145 ),
146
147 input: (title, placeholder, opts) =>
148 createDialogPromise(opts, undefined, { method: "input", title, placeholder, timeout: opts?.timeout }, (r) =>
149 "cancelled" in r && r.cancelled ? undefined : "value" in r ? r.value : undefined,
150 ),
151
152 notify(message: string, type?: "info" | "warning" | "error"): void {
153 // Fire and forget - no response needed
154 output({
155 type: "extension_ui_request",
156 id: crypto.randomUUID(),
157 method: "notify",
158 message,
159 notifyType: type,
160 } as RpcExtensionUIRequest);
161 },
162
163 onTerminalInput(): () => void {
164 // Raw terminal input not supported in RPC mode
165 return () => {};
166 },
167
168 setStatus(key: string, text: string | undefined): void {
169 // Fire and forget - no response needed
170 output({
171 type: "extension_ui_request",
172 id: crypto.randomUUID(),
173 method: "setStatus",
174 statusKey: key,
175 statusText: text,
176 } as RpcExtensionUIRequest);
177 },
178
179 setWorkingMessage(_message?: string): void {
180 // Working message not supported in RPC mode - requires TUI loader access
181 },
182
183 setWorkingVisible(_visible: boolean): void {
184 // Working visibility not supported in RPC mode - requires TUI loader access
185 },
186
187 setWorkingIndicator(_options?: WorkingIndicatorOptions): void {
188 // Working indicator customization not supported in RPC mode - requires TUI loader access
189 },
190
191 setHiddenThinkingLabel(_label?: string): void {
192 // Hidden thinking label not supported in RPC mode - requires TUI message rendering access
193 },
194
195 setWidget(key: string, content: unknown, options?: ExtensionWidgetOptions): void {
196 // Only support string arrays in RPC mode - factory functions are ignored
197 if (content === undefined || Array.isArray(content)) {
198 output({
199 type: "extension_ui_request",
200 id: crypto.randomUUID(),
201 method: "setWidget",
202 widgetKey: key,
203 widgetLines: content as string[] | undefined,
204 widgetPlacement: options?.placement,
205 } as RpcExtensionUIRequest);
206 }
207 // Component factories are not supported in RPC mode - would need TUI access
208 },
209
210 setFooter(_factory: unknown): void {
211 // Custom footer not supported in RPC mode - requires TUI access
212 },
213
214 setHeader(_factory: unknown): void {
215 // Custom header not supported in RPC mode - requires TUI access
216 },
217
218 setTitle(title: string): void {
219 // Fire and forget - host can implement terminal title control
220 output({
221 type: "extension_ui_request",
222 id: crypto.randomUUID(),
223 method: "setTitle",
224 title,
225 } as RpcExtensionUIRequest);
226 },
227
228 async custom() {
229 // Custom UI not supported in RPC mode
230 return undefined as never;
231 },
232
233 pasteToEditor(text: string): void {
234 // Paste handling not supported in RPC mode - falls back to setEditorText
235 this.setEditorText(text);
236 },
237
238 setEditorText(text: string): void {
239 // Fire and forget - host can implement editor control
240 output({
241 type: "extension_ui_request",
242 id: crypto.randomUUID(),
243 method: "set_editor_text",
244 text,
245 } as RpcExtensionUIRequest);
246 },
247
248 getEditorText(): string {
249 // Synchronous method can't wait for RPC response
250 // Host should track editor state locally if needed
251 return "";
252 },
253
254 async editor(title: string, prefill?: string): Promise<string | undefined> {
255 const id = crypto.randomUUID();
256 return new Promise((resolve, reject) => {
257 pendingExtensionRequests.set(id, {
258 resolve: (response: RpcExtensionUIResponse) => {
259 if ("cancelled" in response && response.cancelled) {
260 resolve(undefined);
261 } else if ("value" in response) {
262 resolve(response.value);
263 } else {
264 resolve(undefined);
265 }
266 },
267 reject,
268 });
269 output({ type: "extension_ui_request", id, method: "editor", title, prefill } as RpcExtensionUIRequest);
270 });
271 },
272
273 addAutocompleteProvider(): void {
274 // Autocomplete provider composition is not supported in RPC mode
275 },
276
277 setEditorComponent(): void {
278 // Custom editor components not supported in RPC mode
279 },
280
281 getEditorComponent() {
282 // Custom editor components not supported in RPC mode
283 return undefined;
284 },
285
286 get theme() {
287 return theme;
288 },
289
290 getAllThemes() {
291 return [];
292 },
293
294 getTheme(_name: string) {
295 return undefined;
296 },
297
298 setTheme(_theme: string | Theme) {
299 // Theme switching not supported in RPC mode
300 return { success: false, error: "Theme switching not supported in RPC mode" };
301 },
302
303 getToolsExpanded() {
304 // Tool expansion not supported in RPC mode - no TUI
305 return false;
306 },
307
308 setToolsExpanded(_expanded: boolean) {
309 // Tool expansion not supported in RPC mode - no TUI
310 },
311 });
312
313 runtimeHost.setRebindSession(async () => {
314 await rebindSession();
315 });
316
317 const rebindSession = async (): Promise<void> => {
318 session = runtimeHost.session;
319 await session.bindExtensions({
320 uiContext: createExtensionUIContext(),
321 mode: "rpc",
322 commandContextActions: {
323 waitForIdle: () => session.waitForIdle(),
324 newSession: async (options) => runtimeHost.newSession(options),
325 fork: async (entryId, forkOptions) => {
326 const result = await runtimeHost.fork(entryId, forkOptions);
327 return { cancelled: result.cancelled };
328 },
329 navigateTree: async (targetId, options) => {
330 const result = await session.navigateTree(targetId, {
331 summarize: options?.summarize,
332 customInstructions: options?.customInstructions,
333 replaceInstructions: options?.replaceInstructions,
334 label: options?.label,
335 });
336 return { cancelled: result.cancelled };
337 },
338 switchSession: async (sessionPath, options) => {
339 return runtimeHost.switchSession(sessionPath, options);
340 },
341 reload: async () => {
342 await session.reload();
343 },
344 },
345 shutdownHandler: () => {
346 shutdownRequested = true;
347 },
348 onError: (err) => {
349 output({ type: "extension_error", extensionPath: err.extensionPath, event: err.event, error: err.error });
350 },
351 });
352
353 unsubscribe?.();
354 unsubscribeBackpressure?.();
355 unsubscribe = session.subscribe((event) => {
356 output(toJsonEvent(event));
357 if (event.type === "agent_settled") {
358 void checkShutdownRequested();
359 }
360 });
361 unsubscribeBackpressure = session.agent.subscribe(async () => {
362 await waitForRawStdoutBackpressure();
363 });
364 };
365
366 const registerSignalHandlers = (): void => {
367 const signals: NodeJS.Signals[] = ["SIGTERM"];
368 if (process.platform !== "win32") {
369 signals.push("SIGHUP");
370 }
371
372 for (const signal of signals) {
373 const handler = () => {
374 killTrackedDetachedChildren();
375 void shutdown(signal === "SIGHUP" ? 129 : 143, signal);
376 };
377 process.on(signal, handler);
378 signalCleanupHandlers.push(() => process.off(signal, handler));
379 }
380 };
381
382 await rebindSession();
383 registerSignalHandlers();
384
385 // Handle a single command
386 const handleCommand = async (command: RpcCommand): Promise<RpcResponse | undefined> => {
387 const id = command.id;
388
389 switch (command.type) {
390 // =================================================================
391 // Prompting
392 // =================================================================
393
394 case "prompt": {
395 // Start prompt handling immediately, but emit the authoritative response only after
396 // prompt preflight succeeds. Queued and immediately handled prompts also count as success.
397 let preflightSucceeded = false;
398 void session
399 .prompt(command.message, {
400 images: command.images,
401 streamingBehavior: command.streamingBehavior,
402 source: "rpc",
403 preflightResult: (disposition) => {
404 preflightSucceeded = true;
405 output(success(id, "prompt", { disposition }));
406 },
407 })
408 .catch((e) => {
409 if (!preflightSucceeded) {
410 output(error(id, "prompt", e.message));
411 }
412 });
413 return undefined;
414 }
415
416 case "steer": {
417 const disposition = await session.steer(command.message, command.images, { source: "rpc" });
418 return success(id, "steer", { disposition });
419 }
420
421 case "follow_up": {
422 const disposition = await session.followUp(command.message, command.images, { source: "rpc" });
423 return success(id, "follow_up", { disposition });
424 }
425
426 case "abort": {
427 await session.abort();
428 return success(id, "abort");
429 }
430
431 case "clear_queue": {
432 return success(id, "clear_queue", session.clearQueue());
433 }
434
435 case "new_session": {
436 const options = command.parentSession ? { parentSession: command.parentSession } : undefined;
437 const result = await runtimeHost.newSession(options);
438 if (!result.cancelled) {
439 await rebindSession();
440 }
441 return success(id, "new_session", result);
442 }
443
444 // =================================================================
445 // State
446 // =================================================================
447
448 case "get_state": {
449 const state: RpcSessionState = {
450 model: session.model,
451 thinkingLevel: session.thinkingLevel,
452 isStreaming: session.isStreaming,
453 isCompacting: session.isCompacting,
454 steeringMode: session.steeringMode,
455 followUpMode: session.followUpMode,
456 sessionFile: session.sessionFile,
457 sessionId: session.sessionId,
458 sessionName: session.sessionName,
459 autoCompactionEnabled: session.autoCompactionEnabled,
460 messageCount: session.messages.length,
461 pendingMessageCount: session.pendingMessageCount,
462 };
463 return success(id, "get_state", state);
464 }
465
466 // =================================================================
467 // Model
468 // =================================================================
469
470 case "set_model": {
471 const models = session.modelRuntime.getAvailableSnapshot();
472 const model = models.find((m) => m.provider === command.provider && m.id === command.modelId);
473 if (!model) {
474 return error(id, "set_model", `Model not found: ${command.provider}/${command.modelId}`);
475 }
476 await session.setModel(model);
477 return success(id, "set_model", model);
478 }
479
480 case "cycle_model": {
481 const result = await session.cycleModel();
482 if (!result) {
483 return success(id, "cycle_model", null);
484 }
485 return success(id, "cycle_model", result);
486 }
487
488 case "get_available_models": {
489 const models = session.modelRuntime.getAvailableSnapshot();
490 return success(id, "get_available_models", { models });
491 }
492
493 // =================================================================
494 // Thinking
495 // =================================================================
496
497 case "set_thinking_level": {
498 session.setThinkingLevel(command.level);
499 return success(id, "set_thinking_level");
500 }
501
502 case "cycle_thinking_level": {
503 const level = session.cycleThinkingLevel();
504 if (!level) {
505 return success(id, "cycle_thinking_level", null);
506 }
507 return success(id, "cycle_thinking_level", { level });
508 }
509
510 case "get_available_thinking_levels": {
511 const levels = session.getAvailableThinkingLevels();
512 return success(id, "get_available_thinking_levels", { levels });
513 }
514
515 // =================================================================
516 // Queue Modes
517 // =================================================================
518
519 case "set_steering_mode": {
520 session.setSteeringMode(command.mode);
521 return success(id, "set_steering_mode");
522 }
523
524 case "set_follow_up_mode": {
525 session.setFollowUpMode(command.mode);
526 return success(id, "set_follow_up_mode");
527 }
528
529 // =================================================================
530 // Compaction
531 // =================================================================
532
533 case "compact": {
534 const result = await session.compact(command.customInstructions);
535 return success(id, "compact", result);
536 }
537
538 case "set_auto_compaction": {
539 session.setAutoCompactionEnabled(command.enabled);
540 return success(id, "set_auto_compaction");
541 }
542
543 // =================================================================
544 // Retry
545 // =================================================================
546
547 case "set_auto_retry": {
548 session.setAutoRetryEnabled(command.enabled);
549 return success(id, "set_auto_retry");
550 }
551
552 case "abort_retry": {
553 session.abortRetry();
554 return success(id, "abort_retry");
555 }
556
557 // =================================================================
558 // Bash
559 // =================================================================
560
561 case "bash": {
562 const eventResult = await session.extensionRunner.emitUserBash({
563 type: "user_bash",
564 command: command.command,
565 excludeFromContext: command.excludeFromContext ?? false,
566 cwd: session.sessionManager.getCwd(),
567 });
568
569 if (eventResult?.result) {
570 session.recordBashResult(command.command, eventResult.result, {
571 excludeFromContext: command.excludeFromContext,
572 });
573 return success(id, "bash", eventResult.result);
574 }
575
576 const result = await session.executeBash(command.command, undefined, {
577 excludeFromContext: command.excludeFromContext,
578 id,
579 operations: eventResult?.operations,
580 });
581 return success(id, "bash", result);
582 }
583
584 case "abort_bash": {
585 session.abortBash();
586 return success(id, "abort_bash");
587 }
588
589 // =================================================================
590 // Session
591 // =================================================================
592
593 case "get_session_stats": {
594 const stats = session.getSessionStats();
595 return success(id, "get_session_stats", stats);
596 }
597
598 case "export_html": {
599 const path = await session.exportToHtml(command.outputPath);
600 return success(id, "export_html", { path });
601 }
602
603 case "switch_session": {
604 const result = await runtimeHost.switchSession(command.sessionPath);
605 if (!result.cancelled) {
606 await rebindSession();
607 }
608 return success(id, "switch_session", result);
609 }
610
611 case "fork": {
612 const result = await runtimeHost.fork(command.entryId);
613 if (!result.cancelled) {
614 await rebindSession();
615 }
616 return success(id, "fork", { text: result.selectedText, cancelled: result.cancelled });
617 }
618
619 case "clone": {
620 const leafId = session.sessionManager.getLeafId();
621 if (!leafId) {
622 return error(id, "clone", "Cannot clone session: no current entry selected");
623 }
624 const result = await runtimeHost.fork(leafId, { position: "at" });
625 if (!result.cancelled) {
626 await rebindSession();
627 }
628 return success(id, "clone", { cancelled: result.cancelled });
629 }
630
631 case "get_fork_messages": {
632 const messages = session.getUserMessagesForForking();
633 return success(id, "get_fork_messages", { messages });
634 }
635
636 case "get_entries": {
637 const sessionManager = session.sessionManager;
638 let entries = sessionManager.getEntries();
639 if (command.since !== undefined) {
640 const sinceIndex = entries.findIndex((e) => e.id === command.since);
641 if (sinceIndex === -1) {
642 return error(id, "get_entries", `Entry not found: ${command.since}`);
643 }
644 entries = entries.slice(sinceIndex + 1);
645 }
646 return success(id, "get_entries", { entries, leafId: sessionManager.getLeafId() });
647 }
648
649 case "get_tree": {
650 const sessionManager = session.sessionManager;
651 return success(id, "get_tree", { tree: sessionManager.getTree(), leafId: sessionManager.getLeafId() });
652 }
653
654 case "get_last_assistant_text": {
655 const text = session.getLastAssistantText();
656 return success(id, "get_last_assistant_text", { text });
657 }
658
659 case "set_session_name": {
660 const name = command.name.trim();
661 if (!name) {
662 return error(id, "set_session_name", "Session name cannot be empty");
663 }
664 session.setSessionName(name);
665 return success(id, "set_session_name");
666 }
667
668 // =================================================================
669 // Messages
670 // =================================================================
671
672 case "get_messages": {
673 return success(id, "get_messages", { messages: session.messages });
674 }
675
676 // =================================================================
677 // Commands (available for invocation via prompt)
678 // =================================================================
679
680 case "get_commands": {
681 const commands: RpcSlashCommand[] = [];
682
683 for (const command of session.extensionRunner.getRegisteredCommands()) {
684 commands.push({
685 name: command.invocationName,
686 description: command.description,
687 source: "extension",
688 sourceInfo: command.sourceInfo,
689 });
690 }
691
692 for (const template of session.promptTemplates) {
693 commands.push({
694 name: template.name,
695 description: template.description,
696 source: "prompt",
697 sourceInfo: template.sourceInfo,
698 });
699 }
700
701 for (const skill of session.resourceLoader.getSkills().skills) {
702 commands.push({
703 name: `skill:${skill.name}`,
704 description: skill.description,
705 source: "skill",
706 sourceInfo: skill.sourceInfo,
707 });
708 }
709
710 return success(id, "get_commands", { commands });
711 }
712
713 default: {
714 const unknownCommand = command as { type: string };
715 return error(id, unknownCommand.type, `Unknown command: ${unknownCommand.type}`);
716 }
717 }
718 };
719
720 /**
721 * Check if shutdown was requested and perform shutdown if so.
722 * Called after handling each command when waiting for the next command.
723 */
724 let detachInput = () => {};
725
726 async function shutdown(exitCode = 0, signal?: NodeJS.Signals): Promise<never> {
727 if (shuttingDown) {
728 process.exit(exitCode);
729 }
730 shuttingDown = true;
731 for (const cleanup of signalCleanupHandlers) {
732 cleanup();
733 }
734 unsubscribe?.();
735 unsubscribeBackpressure?.();
736 await runtimeHost.dispose();
737 detachInput();
738 process.stdin.pause();
739 if (signal !== "SIGTERM") {
740 await flushRawStdout();
741 }
742 process.exit(exitCode);
743 }
744
745 async function checkShutdownRequested(): Promise<void> {
746 if (!shutdownRequested) return;
747 await shutdown();
748 }
749
750 const handleInputLine = async (line: string) => {
751 let parsed: unknown;
752 try {
753 parsed = JSON.parse(line);
754 } catch (parseError: unknown) {
755 output(
756 error(
757 undefined,
758 "parse",
759 `Failed to parse command: ${parseError instanceof Error ? parseError.message : String(parseError)}`,
760 ),
761 );
762 await waitForRawStdoutBackpressure();
763 return;
764 }
765
766 // Handle extension UI responses
767 if (
768 typeof parsed === "object" &&
769 parsed !== null &&
770 "type" in parsed &&
771 parsed.type === "extension_ui_response"
772 ) {
773 const response = parsed as RpcExtensionUIResponse;
774 const pending = pendingExtensionRequests.get(response.id);
775 if (pending) {
776 pendingExtensionRequests.delete(response.id);
777 pending.resolve(response);
778 }
779 return;
780 }
781
782 const command = parsed as RpcCommand;
783 try {
784 const response = await handleCommand(command);
785 if (response) {
786 output(response);
787 await waitForRawStdoutBackpressure();
788 }
789 await checkShutdownRequested();
790 } catch (commandError: unknown) {
791 output(
792 error(
793 command.id,
794 command.type,
795 commandError instanceof Error ? commandError.message : String(commandError),
796 ),
797 );
798 await waitForRawStdoutBackpressure();
799 }
800 };
801
802 const onInputEnd = () => {
803 void shutdown();
804 };
805 process.stdin.on("end", onInputEnd);
806
807 detachInput = (() => {
808 const detachJsonl = attachJsonlLineReader(process.stdin, (line) => {
809 void handleInputLine(line);
810 });
811 return () => {
812 detachJsonl();
813 process.stdin.off("end", onInputEnd);
814 };
815 })();
816
817 // Keep process alive forever
818 return new Promise(() => {});
819}