File size: 5,498 Bytes
cd8bd0a
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
/**
 * ACP (Agent Client Protocol) — Process Spawner & Manager
 *
 * Spawns CLI agents as child processes and manages their lifecycle.
 * Communication happens via stdin/stdout (JSON-RPC style) or piped HTTP.
 *
 * This module provides a "CLI-as-backend" transport: instead of intercepting
 * HTTP API calls, OmniRoute spawns the CLI directly and feeds prompts through
 * its native interface.
 */

import { spawn, ChildProcess } from "child_process";
import { EventEmitter } from "events";

export interface AcpSession {
  /** Unique session ID */
  id: string;
  /** Agent ID (e.g., "codex", "claude") */
  agentId: string;
  /** Child process handle */
  process: ChildProcess;
  /** Whether the process is alive */
  alive: boolean;
  /** Accumulated stdout buffer */
  stdoutBuffer: string;
  /** Accumulated stderr buffer */
  stderrBuffer: string;
  /** Created timestamp */
  createdAt: Date;
}

/**
 * ACP Session Manager
 *
 * Manages the lifecycle of CLI agent processes.
 * Each session represents one running CLI agent instance.
 */
export class AcpManager extends EventEmitter {
  private sessions: Map<string, AcpSession> = new Map();

  /**
   * Spawn a new CLI agent process.
   */
  spawn(
    agentId: string,
    binary: string,
    args: string[] = [],
    env: Record<string, string> = {}
  ): AcpSession {
    const ALLOWED_AGENTS = ["claude", "codex", "gemini", "qwen"];
    if (!ALLOWED_AGENTS.includes(agentId)) {
      throw new Error(`Unknown agent: ${agentId}`);
    }

    const sessionId = `acp-${agentId}-${Date.now()}-${crypto.randomUUID().slice(0, 8)}`;

    const child = spawn(binary, args, {
      stdio: ["pipe", "pipe", "pipe"],
      env: { ...process.env, ...env },
      shell: false,
    });

    const session: AcpSession = {
      id: sessionId,
      agentId,
      process: child,
      alive: true,
      stdoutBuffer: "",
      stderrBuffer: "",
      createdAt: new Date(),
    };

    child.stdout?.on("data", (chunk: Buffer) => {
      session.stdoutBuffer += chunk.toString();
      this.emit("stdout", { sessionId, data: chunk.toString() });
    });

    child.stderr?.on("data", (chunk: Buffer) => {
      session.stderrBuffer += chunk.toString();
      this.emit("stderr", { sessionId, data: chunk.toString() });
    });

    child.on("exit", (code, signal) => {
      session.alive = false;
      this.emit("exit", { sessionId, code, signal });
    });

    child.on("error", (err) => {
      session.alive = false;
      this.emit("error", { sessionId, error: err });
    });

    this.sessions.set(sessionId, session);
    return session;
  }

  /**
   * Send input to a running session's stdin.
   */
  sendInput(sessionId: string, input: string): boolean {
    const session = this.sessions.get(sessionId);
    if (!session?.alive || !session.process.stdin?.writable) return false;

    session.process.stdin.write(input);
    return true;
  }

  /**
   * Send a prompt to a CLI agent and collect the response.
   * This is a higher-level method that handles the send/receive cycle.
   */
  async sendPrompt(sessionId: string, prompt: string, timeoutMs: number = 120000): Promise<string> {
    const session = this.sessions.get(sessionId);
    if (!session?.alive) throw new Error(`Session ${sessionId} is not alive`);

    // Clear buffer before sending
    session.stdoutBuffer = "";

    // Send prompt
    this.sendInput(sessionId, prompt + "\n");

    // Wait for response (collect until process goes idle or timeout)
    return new Promise((resolve, reject) => {
      const timer = setTimeout(() => {
        reject(new Error(`ACP timeout after ${timeoutMs}ms`));
      }, timeoutMs);

      let idleTimer: ReturnType<typeof setTimeout>;

      const onData = ({ sessionId: sid }: { sessionId: string }) => {
        if (sid !== sessionId) return;
        // Reset idle timer on new data
        clearTimeout(idleTimer);
        idleTimer = setTimeout(() => {
          clearTimeout(timer);
          this.removeListener("stdout", onData);
          this.removeListener("exit", onExit);
          resolve(session.stdoutBuffer);
        }, 2000); // 2s idle = response complete
      };

      const onExit = ({ sessionId: sid }: { sessionId: string }) => {
        if (sid !== sessionId) return;
        clearTimeout(timer);
        clearTimeout(idleTimer);
        this.removeListener("stdout", onData);
        this.removeListener("exit", onExit);
        resolve(session.stdoutBuffer);
      };

      this.on("stdout", onData);
      this.on("exit", onExit);
    });
  }

  /**
   * Kill a session and clean up.
   */
  kill(sessionId: string): boolean {
    const session = this.sessions.get(sessionId);
    if (!session) return false;

    if (session.alive) {
      session.process.kill("SIGTERM");
      // Force kill after 5s
      setTimeout(() => {
        if (session.alive) {
          session.process.kill("SIGKILL");
        }
      }, 5000);
    }

    this.sessions.delete(sessionId);
    return true;
  }

  /**
   * Get all active sessions.
   */
  getActiveSessions(): AcpSession[] {
    return Array.from(this.sessions.values()).filter((s) => s.alive);
  }

  /**
   * Get a specific session.
   */
  getSession(sessionId: string): AcpSession | undefined {
    return this.sessions.get(sessionId);
  }

  /**
   * Kill all sessions.
   */
  killAll(): void {
    for (const [id] of this.sessions) {
      this.kill(id);
    }
  }
}

// Singleton manager instance
export const acpManager = new AcpManager();