413 lines
13 KiB
TypeScript
413 lines
13 KiB
TypeScript
/**
|
|
* Bridge Spawner
|
|
*
|
|
* Auto-starts inline HTTP bridges for detected CLI subscriptions. Each bridge
|
|
* exposes a `POST /api/generate` endpoint that the gateway can call as a regular
|
|
* external provider. Bridges run in-process to avoid the overhead of spawning
|
|
* separate Node processes — they listen on a dedicated port per subscription.
|
|
*/
|
|
|
|
import { execFile } from 'child_process';
|
|
import { createServer, type IncomingMessage, type Server, type ServerResponse } from 'http';
|
|
import { logger } from '../observability/logger.js';
|
|
import type { SubscriptionDescriptor, SubscriptionStatus } from './subscription-discovery.js';
|
|
|
|
interface RunningBridge {
|
|
descriptor: SubscriptionDescriptor;
|
|
server: Server;
|
|
port: number;
|
|
url: string;
|
|
startedAt: Date;
|
|
}
|
|
|
|
const runningBridges = new Map<string, RunningBridge>();
|
|
|
|
function extractCliContent(stdout: string): string {
|
|
let lastAgentMessage = '';
|
|
for (const line of stdout.split(/\r?\n/)) {
|
|
const jsonStart = line.indexOf('{');
|
|
if (jsonStart < 0) continue;
|
|
try {
|
|
const event = JSON.parse(line.slice(jsonStart)) as {
|
|
type?: string;
|
|
item?: { type?: string; text?: string };
|
|
};
|
|
if (event.type === 'item.completed' && event.item?.type === 'agent_message' && event.item.text) {
|
|
lastAgentMessage = event.item.text;
|
|
}
|
|
} catch {
|
|
// Non-JSON status lines from CLIs are ignored.
|
|
}
|
|
}
|
|
return (lastAgentMessage || stdout).trim();
|
|
}
|
|
|
|
/**
|
|
* Run a CLI tool with stdin-piped prompt, return stdout content.
|
|
* Generic implementation that all inline bridges share.
|
|
*/
|
|
async function runCli(
|
|
command: string,
|
|
args: readonly string[],
|
|
prompt: string,
|
|
timeoutMs: number = 300_000
|
|
): Promise<{ success: boolean; content?: string; error?: string }> {
|
|
return new Promise((resolve) => {
|
|
try {
|
|
const child = execFile(
|
|
command,
|
|
args as string[],
|
|
{ timeout: timeoutMs, maxBuffer: 10 * 1024 * 1024 },
|
|
(err, stdout) => {
|
|
if (err) {
|
|
resolve({ success: false, error: err.message.slice(0, 500) });
|
|
} else {
|
|
resolve({ success: true, content: extractCliContent(stdout) });
|
|
}
|
|
}
|
|
);
|
|
if (child.stdin) {
|
|
child.stdin.write(prompt);
|
|
child.stdin.end();
|
|
}
|
|
} catch (err) {
|
|
resolve({ success: false, error: err instanceof Error ? err.message : String(err) });
|
|
}
|
|
});
|
|
}
|
|
|
|
/**
|
|
* Build the CLI invocation for a given subscription.
|
|
*/
|
|
function buildCliInvocation(desc: SubscriptionDescriptor, model?: string): { cmd: string; args: string[] } {
|
|
switch (desc.bridgeImplementation) {
|
|
case 'inline-claude': {
|
|
const args = ['--print', '--output-format', 'text'];
|
|
if (model) args.push('--model', model);
|
|
return { cmd: 'claude', args };
|
|
}
|
|
case 'inline-copilot': {
|
|
// gh copilot suggest is interactive; we use the OpenAI-compatible copilot-api proxy if available.
|
|
return { cmd: 'gh', args: ['copilot', 'suggest', '--shell'] };
|
|
}
|
|
case 'inline-openai': {
|
|
// Generic OpenAI-compatible CLI (chatgpt-cli, gemini-cli with OpenAI compat)
|
|
return { cmd: desc.command, args: model ? ['--model', model] : [] };
|
|
}
|
|
case 'external-codex': {
|
|
// codex exec is the non-interactive CLI surface; prompt is read from stdin via "-".
|
|
const args = [
|
|
'exec',
|
|
'--sandbox',
|
|
'read-only',
|
|
'--skip-git-repo-check',
|
|
'--ephemeral',
|
|
'--color',
|
|
'never',
|
|
'--json',
|
|
];
|
|
args.push('-');
|
|
return { cmd: 'codex', args };
|
|
}
|
|
}
|
|
}
|
|
|
|
function sendJson(res: ServerResponse, statusCode: number, payload: Record<string, unknown>): void {
|
|
res.writeHead(statusCode);
|
|
res.end(JSON.stringify(payload));
|
|
}
|
|
|
|
function readJsonBody(req: IncomingMessage): Promise<Record<string, unknown>> {
|
|
return new Promise((resolve, reject) => {
|
|
let body = '';
|
|
req.on('data', (chunk) => {
|
|
body += chunk;
|
|
if (body.length > 2_000_000) {
|
|
req.destroy(new Error('request body too large'));
|
|
}
|
|
});
|
|
req.on('end', () => {
|
|
try {
|
|
resolve(JSON.parse(body || '{}') as Record<string, unknown>);
|
|
} catch (err) {
|
|
reject(err);
|
|
}
|
|
});
|
|
req.on('error', reject);
|
|
});
|
|
}
|
|
|
|
function contentToText(content: unknown): string {
|
|
if (typeof content === 'string') return content;
|
|
if (Array.isArray(content)) {
|
|
return content
|
|
.map((part) => {
|
|
if (typeof part === 'string') return part;
|
|
if (part && typeof part === 'object' && 'text' in part) {
|
|
return String((part as { text?: unknown }).text ?? '');
|
|
}
|
|
return '';
|
|
})
|
|
.filter(Boolean)
|
|
.join('\n');
|
|
}
|
|
return content == null ? '' : String(content);
|
|
}
|
|
|
|
function promptFromMessages(messages: readonly unknown[]): { prompt: string; system?: string } {
|
|
const system = messages
|
|
.filter((message): message is { role?: unknown; content?: unknown } => !!message && typeof message === 'object')
|
|
.filter((message) => message.role === 'system')
|
|
.map((message) => contentToText(message.content))
|
|
.filter(Boolean)
|
|
.join('\n\n');
|
|
|
|
const prompt = messages
|
|
.filter((message): message is { role?: unknown; content?: unknown } => !!message && typeof message === 'object')
|
|
.filter((message) => message.role !== 'system')
|
|
.map((message) => `${String(message.role ?? 'user')}: ${contentToText(message.content)}`)
|
|
.filter((line) => line.trim().length > 0)
|
|
.join('\n\n');
|
|
|
|
return { prompt, system: system || undefined };
|
|
}
|
|
|
|
function estimateTokens(text: string | undefined): number {
|
|
if (!text) return 0;
|
|
return Math.max(1, Math.ceil(text.length / 4));
|
|
}
|
|
|
|
function buildChatCompletionResponse(
|
|
desc: SubscriptionDescriptor,
|
|
model: string,
|
|
content: string,
|
|
prompt: string
|
|
): Record<string, unknown> {
|
|
const promptTokens = estimateTokens(prompt);
|
|
const completionTokens = estimateTokens(content);
|
|
return {
|
|
id: `chatcmpl-${desc.id}-${Date.now()}`,
|
|
object: 'chat.completion',
|
|
created: Math.floor(Date.now() / 1000),
|
|
model,
|
|
provider: desc.providerName,
|
|
success: true,
|
|
content,
|
|
choices: [
|
|
{
|
|
index: 0,
|
|
message: { role: 'assistant', content },
|
|
finish_reason: 'stop',
|
|
},
|
|
],
|
|
usage: {
|
|
prompt_tokens: promptTokens,
|
|
completion_tokens: completionTokens,
|
|
total_tokens: promptTokens + completionTokens,
|
|
},
|
|
};
|
|
}
|
|
|
|
/**
|
|
* Spawn an inline HTTP bridge for a subscription. Returns the URL the gateway
|
|
* should use to talk to it. Idempotent — calling twice returns the same bridge.
|
|
*/
|
|
export function spawnBridge(desc: SubscriptionDescriptor): Promise<RunningBridge> {
|
|
const existing = runningBridges.get(desc.id);
|
|
if (existing) {
|
|
return Promise.resolve(existing);
|
|
}
|
|
|
|
return new Promise((resolve, reject) => {
|
|
const server = createServer(async (req, res) => {
|
|
res.setHeader('Content-Type', 'application/json');
|
|
res.setHeader('Access-Control-Allow-Origin', '*');
|
|
res.setHeader('Access-Control-Allow-Methods', 'GET, POST, OPTIONS');
|
|
res.setHeader('Access-Control-Allow-Headers', 'Content-Type, Authorization');
|
|
|
|
if (req.method === 'OPTIONS') {
|
|
sendJson(res, 200, { ok: true });
|
|
return;
|
|
}
|
|
|
|
if (req.method === 'GET' && req.url === '/health') {
|
|
const current = runningBridges.get(desc.id);
|
|
sendJson(res, 200, {
|
|
status: 'ok',
|
|
subscription: desc.id,
|
|
label: desc.label,
|
|
command: desc.command,
|
|
provider: desc.providerName,
|
|
models: desc.models,
|
|
endpoints: ['/health', '/api/generate', '/v1/completion', '/v1/chat/completions'],
|
|
uptimeSeconds: current ? Math.floor((Date.now() - current.startedAt.getTime()) / 1000) : 0,
|
|
});
|
|
return;
|
|
}
|
|
|
|
if (req.method === 'POST' && req.url === '/v1/chat/completions') {
|
|
try {
|
|
const body = await readJsonBody(req);
|
|
const messages = body['messages'];
|
|
if (!Array.isArray(messages)) {
|
|
sendJson(res, 400, { error: 'messages array required' });
|
|
return;
|
|
}
|
|
|
|
const model = typeof body['model'] === 'string'
|
|
? body['model']
|
|
: desc.models[0]?.id ?? desc.id;
|
|
const { prompt, system } = promptFromMessages(messages);
|
|
if (!prompt) {
|
|
sendJson(res, 400, { error: 'messages content required' });
|
|
return;
|
|
}
|
|
|
|
const fullPrompt = system ? `${system}\n\n---\n\n${prompt}` : prompt;
|
|
const { cmd, args } = buildCliInvocation(desc, model);
|
|
const result = await runCli(cmd, args, fullPrompt);
|
|
if (result.success) {
|
|
sendJson(res, 200, buildChatCompletionResponse(desc, model, result.content ?? '', fullPrompt));
|
|
} else {
|
|
sendJson(res, 502, { success: false, error: result.error });
|
|
}
|
|
} catch (e) {
|
|
sendJson(res, 500, { error: e instanceof Error ? e.message : 'parse error' });
|
|
}
|
|
return;
|
|
}
|
|
|
|
if (req.method === 'POST' && (req.url === '/api/generate' || req.url === '/v1/completion')) {
|
|
try {
|
|
const body = await readJsonBody(req);
|
|
const prompt = typeof body['prompt'] === 'string' ? body['prompt'] : '';
|
|
const system = typeof body['system'] === 'string' ? body['system'] : undefined;
|
|
const model = typeof body['model'] === 'string'
|
|
? body['model']
|
|
: desc.models[0]?.id ?? desc.id;
|
|
if (!prompt) {
|
|
sendJson(res, 400, { error: 'prompt required' });
|
|
return;
|
|
}
|
|
const fullPrompt = system ? `${system}\n\n---\n\n${prompt}` : prompt;
|
|
const { cmd, args } = buildCliInvocation(desc, model);
|
|
const result = await runCli(cmd, args, fullPrompt);
|
|
if (result.success) {
|
|
sendJson(res, 200, {
|
|
success: true,
|
|
content: result.content,
|
|
response: result.content,
|
|
provider: desc.providerName,
|
|
model,
|
|
usage: {
|
|
prompt_tokens: estimateTokens(fullPrompt),
|
|
completion_tokens: estimateTokens(result.content),
|
|
total_tokens: estimateTokens(fullPrompt) + estimateTokens(result.content),
|
|
},
|
|
});
|
|
} else {
|
|
sendJson(res, 502, { success: false, error: result.error });
|
|
}
|
|
} catch (e) {
|
|
sendJson(res, 500, { error: e instanceof Error ? e.message : 'parse error' });
|
|
}
|
|
return;
|
|
}
|
|
|
|
sendJson(res, 404, { error: 'not found' });
|
|
});
|
|
|
|
server.on('error', (err) => {
|
|
// Port in use → assume an existing bridge is already running, treat as success
|
|
if ((err as NodeJS.ErrnoException).code === 'EADDRINUSE') {
|
|
logger.info(
|
|
{ subscription: desc.id, port: desc.bridgePort },
|
|
'Port already in use — assuming external bridge is healthy'
|
|
);
|
|
const url = `http://localhost:${desc.bridgePort}`;
|
|
const fakeBridge: RunningBridge = {
|
|
descriptor: desc,
|
|
server, // server failed to bind; OK to keep handle
|
|
port: desc.bridgePort,
|
|
url,
|
|
startedAt: new Date(),
|
|
};
|
|
runningBridges.set(desc.id, fakeBridge);
|
|
resolve(fakeBridge);
|
|
} else {
|
|
reject(err);
|
|
}
|
|
});
|
|
|
|
server.listen(desc.bridgePort, 'localhost', () => {
|
|
const url = `http://localhost:${desc.bridgePort}`;
|
|
const bridge: RunningBridge = {
|
|
descriptor: desc,
|
|
server,
|
|
port: desc.bridgePort,
|
|
url,
|
|
startedAt: new Date(),
|
|
};
|
|
runningBridges.set(desc.id, bridge);
|
|
// Set the env var so the existing external-providers logic finds the bridge
|
|
process.env[desc.bridgeEnvKey] = url;
|
|
logger.info(
|
|
{ subscription: desc.id, url, port: desc.bridgePort, envKey: desc.bridgeEnvKey },
|
|
'Inline subscription bridge started'
|
|
);
|
|
resolve(bridge);
|
|
});
|
|
});
|
|
}
|
|
|
|
/**
|
|
* Spawn bridges for every detected, authenticated subscription that doesn't
|
|
* already have a bridge URL configured. Returns the list of started bridges.
|
|
*/
|
|
export async function spawnDetectedBridges(
|
|
statuses: readonly SubscriptionStatus[]
|
|
): Promise<RunningBridge[]> {
|
|
const toSpawn = statuses.filter(
|
|
(s) => s.installed && s.authenticated !== false && !s.bridgeRunning
|
|
);
|
|
const results: RunningBridge[] = [];
|
|
for (const status of toSpawn) {
|
|
try {
|
|
const bridge = await spawnBridge(status.descriptor);
|
|
results.push(bridge);
|
|
} catch (err) {
|
|
logger.warn(
|
|
{ err, subscription: status.descriptor.id },
|
|
'Failed to spawn subscription bridge — continuing'
|
|
);
|
|
}
|
|
}
|
|
return results;
|
|
}
|
|
|
|
/**
|
|
* Snapshot of currently running in-process bridges. Used by the dashboard.
|
|
*/
|
|
export function getRunningBridges(): readonly RunningBridge[] {
|
|
return Array.from(runningBridges.values());
|
|
}
|
|
|
|
/**
|
|
* Stop all inline bridges (used during graceful shutdown).
|
|
*/
|
|
export async function stopAllBridges(): Promise<void> {
|
|
await Promise.all(
|
|
Array.from(runningBridges.values()).map(
|
|
(bridge) =>
|
|
new Promise<void>((resolve) => {
|
|
try {
|
|
bridge.server.close(() => resolve());
|
|
} catch {
|
|
resolve();
|
|
}
|
|
})
|
|
)
|
|
);
|
|
runningBridges.clear();
|
|
}
|