From a0d9ea63f9fd2925a0155c59dfeff0e346761d2e Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Andr=C3=A9=20Padez?= Date: Fri, 6 Mar 2026 08:48:19 +0000 Subject: [PATCH] sidecar self-registration: sidecars connect to API server instead of vice versa Flips the connection model so sidecars register themselves with the API server via WebSocket at /api/sidecar/register, enabling dynamic discovery, location independence, and automatic reconnection from either side. - Add registration protocol types and PTY command/event types - Create sidecar-registry.ts (replaces sidecar-client.ts) as passive registry - Create sidecar connector (connect.ts) with exponential backoff reconnect - Convert process sidecar from WS server to WS client - Convert PTY sidecar from WS server to multiplexed WS client - Simplify terminal bridge to thin adapter using registry - Add PTY sidecar as PM2-managed process - Update all consumer imports Co-Authored-By: Claude Opus 4.6 --- ecosystem.config.cjs | 12 +- src/server.tsx | 57 +++- src/servers/api/email/accounts.ts | 9 +- src/servers/api/pi/websocket.ts | 2 +- src/servers/api/queue/queue.ts | 2 +- src/servers/api/terminal/pty-sidecar.mjs | 248 +++++++++------ src/servers/api/terminal/websocket.ts | 290 +++++++---------- src/servers/channels/discord/handler.ts | 40 ++- src/servers/channels/send-claude-code.ts | 2 +- src/servers/channels/telegram/handler.ts | 2 +- src/servers/channels/whatsapp/handler.ts | 2 +- src/servers/sidecar-client.ts | 303 ------------------ src/servers/sidecar-registry.ts | 315 +++++++++++++++++++ src/servers/sidecar/connect.ts | 135 ++++++++ src/servers/sidecar/index.ts | 127 +++----- src/servers/sidecar/protocol.ts | 28 ++ src/servers/sidecar/registration-protocol.ts | 16 + 17 files changed, 878 insertions(+), 712 deletions(-) delete mode 100644 src/servers/sidecar-client.ts create mode 100644 src/servers/sidecar-registry.ts create mode 100644 src/servers/sidecar/connect.ts create mode 100644 src/servers/sidecar/registration-protocol.ts diff --git a/ecosystem.config.cjs b/ecosystem.config.cjs index 758bb23e..bf151e32 100644 --- a/ecosystem.config.cjs +++ b/ecosystem.config.cjs @@ -1,5 +1,11 @@ module.exports = { apps: [ + { + name: 'officer', + script: 'bun', + args: 'start', + watch: false, + }, { name: 'officer-sidecar', script: 'bun', @@ -7,9 +13,9 @@ module.exports = { watch: false, }, { - name: 'officer', - script: 'bun', - args: 'start', + name: 'officer-pty', + script: 'node', + args: 'src/servers/api/terminal/pty-sidecar.mjs', watch: false, }, ], diff --git a/src/server.tsx b/src/server.tsx index ffde24b5..f45d6ae5 100644 --- a/src/server.tsx +++ b/src/server.tsx @@ -4,7 +4,7 @@ import { serve } from 'bun'; import { honoServer } from './servers/hono'; import { verify } from './servers/jwt'; import { isTokenBlacklisted } from 'officerdb'; -import { terminalWebsocket, initTerminalSidecars } from './servers/api/terminal/websocket'; +import { terminalWebsocket } from './servers/api/terminal/websocket'; import { piWebsocket } from './servers/api/pi/websocket'; import { cliampWebsocket } from './servers/api/cliamp/websocket'; import { cliampAudioWebsocket } from './servers/api/cliamp/audio-ws'; @@ -12,7 +12,8 @@ import { desktopWebsocket } from './servers/api/desktop/websocket'; import { findEntryByProxyId, touchEntry } from './servers/api/dev-server/router'; import officerWeb from './apps/officer-web/index.html'; import { startBrowserRelay } from './servers/api/browser/relay'; -import { initSidecarClient } from './servers/sidecar-client'; +import { registerSidecar, unregisterSidecar, handleSidecarMessage } from './servers/sidecar-registry'; +import type { SidecarRegistration } from './servers/sidecar/registration-protocol'; import { toShellUsername } from './servers/data-path'; const { PORT = '5000' } = process.env; @@ -22,7 +23,7 @@ type WSData = { email: string; username: string; role: string; - provider: 'terminal' | 'pi' | 'dev-server' | 'cliamp' | 'cliamp-audio' | 'desktop'; + provider: 'terminal' | 'pi' | 'dev-server' | 'cliamp' | 'cliamp-audio' | 'desktop' | 'sidecar'; sandboxed: boolean; sessionId?: string; cwd?: string; @@ -36,12 +37,50 @@ type WSData = { wsToken?: string; }; +// Sidecar registration WebSocket handler +const sidecarConnections = new Map, string>(); // ws → sidecar ID + +const sidecarWebsocket = { + open(_ws: ServerWebSocket) { + // Wait for registration message + }, + message(ws: ServerWebSocket, raw: string | Buffer) { + try { + const data = typeof raw === 'string' ? raw : raw.toString(); + const msg = JSON.parse(data); + + if (msg.type === 'register') { + const id = registerSidecar(ws, msg as SidecarRegistration); + sidecarConnections.set(ws, id); + ws.send(JSON.stringify({ type: 'registered', id })); + return; + } + + const id = sidecarConnections.get(ws); + if (id) { + handleSidecarMessage(id, msg); + } + } catch { + // skip malformed messages + } + }, + close(ws: ServerWebSocket) { + const id = sidecarConnections.get(ws); + if (id) { + unregisterSidecar(id); + sidecarConnections.delete(ws); + } + }, + drain() {}, +}; + const handlers: Record = { terminal: terminalWebsocket, pi: piWebsocket, cliamp: cliampWebsocket, 'cliamp-audio': cliampAudioWebsocket, desktop: desktopWebsocket, + sidecar: sidecarWebsocket, }; // Dev-server WebSocket proxy: bridges client WS ↔ upstream dev server WS (for HMR etc.) @@ -187,6 +226,12 @@ const server = serve({ if (req.headers.get('upgrade') === 'websocket') return upgradeDevServerWs(req, server); return honoServer.fetch(req, server); }, + '/api/sidecar/register': (req: Request, server: any) => { + const ok = server.upgrade(req, { + data: { provider: 'sidecar', userId: 0, email: '', username: '', role: '', sandboxed: false }, + }); + if (!ok) return new Response('Upgrade failed', { status: 500 }); + }, '/api/terminal/ws': (req, server) => upgradeWs(req, server, 'terminal'), '/api/pi/chat/ws': (req, server) => upgradeWs(req, server, 'pi'), '/api/cliamp/ws': (req, server) => upgradeWs(req, server, 'cliamp'), @@ -229,11 +274,7 @@ try { console.error('[browser-relay] failed to start:', err instanceof Error ? err.message : err); } -// Connect to process sidecar (owns proxy, Claude Code, Pi, queue) -initSidecarClient(); - - -void initTerminalSidecars(); +// Sidecars connect to us via /api/sidecar/register — no init needed // Ensure PulseAudio is running with virtual sink for cliamp audio streaming (async () => { diff --git a/src/servers/api/email/accounts.ts b/src/servers/api/email/accounts.ts index 56b119f6..8a40401f 100644 --- a/src/servers/api/email/accounts.ts +++ b/src/servers/api/email/accounts.ts @@ -9,7 +9,7 @@ import { updateEmailAccountStatus, } from 'officerdb'; import { validateImapConnection } from './imap-validate'; -import * as sidecar from '../../sidecar-client'; +import * as sidecar from '../../sidecar-registry'; type CreateAccountBody = { provider: string; @@ -137,7 +137,12 @@ accountsRouter.post('/:id/sync', async (ctx) => { if (account.status === 'syncing') throw BAD_REQUEST('Account is already syncing'); // Resolve auth before enqueueing - const authResult = await resolveAuth(user.id, account.authType, account.email, account.credentials as Record); + const authResult = await resolveAuth( + user.id, + account.authType, + account.email, + account.credentials as Record, + ); if (!authResult.ok) throw BAD_REQUEST(authResult.error); // Set status immediately so the UI reflects the queued state diff --git a/src/servers/api/pi/websocket.ts b/src/servers/api/pi/websocket.ts index d2bf9e90..c255683f 100644 --- a/src/servers/api/pi/websocket.ts +++ b/src/servers/api/pi/websocket.ts @@ -5,7 +5,7 @@ import { sessionManager } from './session-manager'; import * as storage from './storage'; import * as piBridge from './pi-bridge'; import { sendClaudeCodeStreaming } from '@@/channels/send-claude-code'; -import * as sidecar from '@@/sidecar-client'; +import * as sidecar from '@@/sidecar-registry'; import { join } from 'path'; import { getHomeDirForRole } from '../../../servers/data-path'; import { getUserSettings } from 'officerdb'; diff --git a/src/servers/api/queue/queue.ts b/src/servers/api/queue/queue.ts index 6fc727e3..c0fa7759 100644 --- a/src/servers/api/queue/queue.ts +++ b/src/servers/api/queue/queue.ts @@ -1,5 +1,5 @@ import { createRouter } from '../../create-router'; -import * as sidecar from '../../sidecar-client'; +import * as sidecar from '../../sidecar-registry'; import { NOT_FOUND } from '../../custom-errors'; export const queueRouter = createRouter(); diff --git a/src/servers/api/terminal/pty-sidecar.mjs b/src/servers/api/terminal/pty-sidecar.mjs index 45c776be..b1d22de9 100644 --- a/src/servers/api/terminal/pty-sidecar.mjs +++ b/src/servers/api/terminal/pty-sidecar.mjs @@ -1,13 +1,21 @@ -// Ignore SIGINT — sudo/pty child processes may propagate it -process.on('SIGINT', () => {}); +// Graceful shutdown +process.on('SIGINT', () => { + console.log('[pty-sidecar] shutting down...'); + for (const [id, session] of sessions) { + try { session.term.kill(); } catch { /* ignore */ } + } + sessions.clear(); + if (ws) { try { ws.close(); } catch { /* ignore */ } } + process.exit(0); +}); +process.on('SIGTERM', () => process.emit('SIGINT')); -import http from 'node:http'; import { existsSync } from 'node:fs'; import { cp, mkdir } from 'node:fs/promises'; import { join, dirname } from 'node:path'; import { fileURLToPath } from 'node:url'; import { execFile } from 'node:child_process'; -import { WebSocketServer } from 'ws'; +import WebSocket from 'ws'; import * as pty from 'node-pty'; const run = (cmd, args, opts = {}) => @@ -20,24 +28,24 @@ const __dirname = dirname(fileURLToPath(import.meta.url)); const templateDir = join(__dirname, 'templates'); -const port = Number(process.env.TERMINAL_PTY_PORT ?? '5337'); -const host = process.env.TERMINAL_PTY_HOST ?? '127.0.0.1'; +import 'dotenv/config'; + +const API_URL = process.env.API_URL ?? `ws://127.0.0.1:${process.env.PORT ?? '5000'}`; +const REGISTER_URL = `${API_URL}/api/sidecar/register`; const BUFFER_MAX = 50 * 1024; +const RECONNECT_DELAYS = [200, 500, 1000, 2000, 4000, 8000, 15000]; -/** @type {Map} */ +/** @type {Map} */ const sessions = new Map(); -const server = http.createServer((req, res) => { - res.writeHead(200, { 'Content-Type': 'text/plain' }); - res.end('terminal-sidecar'); -}); - -const wss = new WebSocketServer({ server }); +// ── Helpers ── const sendJson = (ws, msg) => { try { - ws.send(JSON.stringify(msg)); + if (ws.readyState === WebSocket.OPEN) { + ws.send(JSON.stringify(msg)); + } } catch { // ignore } @@ -76,47 +84,26 @@ const appendBuffer = (session, data) => { } }; -wss.on('connection', (ws) => { - let currentSessionId = null; +// ── Command handler ── - ws.on('message', async (data) => { - let msg; - try { - msg = JSON.parse(typeof data === 'string' ? data : data.toString()); - } catch { - return; - } - - if (msg.type === 'init') { - const sessionId = msg.sessionId; +async function handleCommand(ws, msg) { + switch (msg.type) { + case 'pty:init': { + const { sessionId, config } = msg; if (!sessionId) return; - currentSessionId = sessionId; const existing = sessions.get(sessionId); - - console.log(`[sidecar] init sessionId=${sessionId} existing=${!!existing} total=${sessions.size}`); + console.log(`[pty-sidecar] init sessionId=${sessionId} existing=${!!existing} total=${sessions.size}`); if (existing) { - // Evict old WS if still attached - if (existing.ws && existing.ws !== ws) { - sendJson(existing.ws, { type: 'detached' }); - try { - existing.ws.close(); - } catch { - // ignore - } - } - - existing.ws = ws; - // Replay buffer if (existing.buffer.length > 0) { - sendJson(ws, { type: 'output', data: existing.buffer }); + sendJson(ws, { type: 'pty:output', sessionId, data: existing.buffer }); } // Resize PTY to new client dimensions - const cols = msg.cols ?? existing.cols; - const rows = msg.rows ?? existing.rows; + const cols = config.cols ?? existing.cols; + const rows = config.rows ?? existing.rows; if (cols > 0 && rows > 0 && (cols !== existing.cols || rows !== existing.rows)) { existing.cols = cols; existing.rows = rows; @@ -127,44 +114,40 @@ wss.on('connection', (ws) => { } } + sendJson(ws, { type: 'pty:ready', id: msg.id, sessionId }); return; } // New session — spawn PTY - const shell = msg.shell ?? { command: '/bin/bash', args: ['-i'] }; - const cwd = msg.cwd ?? process.cwd(); - const homeDir = msg.homeDir ?? process.cwd(); - const userLabel = msg.userLabel ?? 'officer'; - const username = msg.username ?? null; - const cols = msg.cols ?? 80; - const rows = msg.rows ?? 24; - const isHost = !!msg.host; + const shell = config.shell ?? { command: '/bin/bash', args: ['-i'] }; + const cwd = config.cwd ?? process.cwd(); + const homeDir = config.homeDir ?? process.cwd(); + const userLabel = config.userLabel ?? 'officer'; + const username = config.username ?? null; + const cols = config.cols ?? 80; + const rows = config.rows ?? 24; + const isHost = !!config.host; let spawnCommand; let spawnArgs; let ptyEnv; if (isHost) { - // Host session — spawn shell directly as current user spawnCommand = shell.command; spawnArgs = shell.args ?? []; - ptyEnv = { ...process.env, TERM: 'xterm-256color', ...(msg.env ?? {}) }; + ptyEnv = { ...process.env, TERM: 'xterm-256color', ...(config.env ?? {}) }; } else if (username) { - // User session — spawn via sudo -u as the target Linux user spawnCommand = 'sudo'; spawnArgs = ['-u', username, '-i', '/bin/zsh']; - ptyEnv = { - TERM: 'xterm-256color', - }; + ptyEnv = { TERM: 'xterm-256color' }; } else { - // Fallback — direct spawn with custom env (legacy) spawnCommand = shell.command; spawnArgs = shell.args ?? []; try { await ensureUserFiles(homeDir); } catch (err) { - console.error('[sidecar] ensureUserFiles failed:', err); + console.error('[pty-sidecar] ensureUserFiles failed:', err); } ptyEnv = { @@ -180,8 +163,6 @@ wss.on('connection', (ws) => { }; } - // For username sessions, don't set cwd — sudo -u -i will cd to the user's home. - // node-pty does chdir before exec, so it would fail if the service user can't access the dir. const ptyCwd = username ? undefined : cwd; let term; @@ -195,8 +176,8 @@ wss.on('connection', (ws) => { }); } catch (err) { const message = err instanceof Error ? err.message : 'Failed to start terminal'; - sendJson(ws, { type: 'output', data: `\r\n[Terminal error] ${message}\r\n` }); - sendJson(ws, { type: 'exit' }); + sendJson(ws, { type: 'pty:output', sessionId, data: `\r\n[Terminal error] ${message}\r\n` }); + sendJson(ws, { type: 'pty:exit', sessionId, exitCode: 1 }); return; } @@ -205,66 +186,125 @@ wss.on('connection', (ws) => { buffer: '', cols, rows, - ws, initConfig: { shell, cwd, homeDir, userLabel }, }; sessions.set(sessionId, session); term.onData((output) => { appendBuffer(session, output); - if (session.ws) { - sendJson(session.ws, { type: 'output', data: output }); - } + sendJson(ws, { type: 'pty:output', sessionId, data: output }); }); term.onExit(({ exitCode, signal }) => { - console.log(`[sidecar] session ${sessionId} exited code=${exitCode} signal=${signal}`); - if (session.ws) { - sendJson(session.ws, { type: 'exit' }); - } + console.log(`[pty-sidecar] session ${sessionId} exited code=${exitCode} signal=${signal}`); + sendJson(ws, { type: 'pty:exit', sessionId, exitCode, signal }); sessions.delete(sessionId); }); + sendJson(ws, { type: 'pty:ready', id: msg.id, sessionId }); return; } - // Route other messages to current session - if (!currentSessionId) return; - const session = sessions.get(currentSessionId); - if (!session) return; - - switch (msg.type) { - case 'input': + case 'pty:input': { + const session = sessions.get(msg.sessionId); + if (session) { session.term.write(msg.data ?? ''); - break; - case 'resize': - if (msg.cols > 0 && msg.rows > 0) { - session.cols = msg.cols; - session.rows = msg.rows; - try { - session.term.resize(msg.cols, msg.rows); - } catch { - // PTY may have already exited - } + } + break; + } + + case 'pty:resize': { + const session = sessions.get(msg.sessionId); + if (session && msg.cols > 0 && msg.rows > 0) { + session.cols = msg.cols; + session.rows = msg.rows; + try { + session.term.resize(msg.cols, msg.rows); + } catch { + // PTY may have already exited } - break; - case 'cwd': - if (msg.path) session.term.write(`cd ${JSON.stringify(msg.path)}\r`); - break; + } + break; + } + + case 'pty:close': { + const session = sessions.get(msg.sessionId); + if (session) { + try { + session.term.kill(); + } catch { + // ignore + } + sessions.delete(msg.sessionId); + } + break; + } + } +} + +// ── Connect to API server with reconnect ── + +let ws = null; +let reconnectAttempt = 0; +let reconnectTimer = null; + +console.log(`[pty-sidecar] starting, connecting to ${REGISTER_URL}`); + +function connect() { + if (ws && (ws.readyState === WebSocket.CONNECTING || ws.readyState === WebSocket.OPEN)) return; + + try { + ws = new WebSocket(REGISTER_URL); + } catch (err) { + console.error(`[pty-sidecar] failed to create WebSocket:`, err.message ?? err); + scheduleReconnect(); + return; + } + + ws.on('open', () => { + reconnectAttempt = 0; + console.log('[pty-sidecar] connected, sending registration...'); + sendJson(ws, { type: 'register', name: 'pty', capabilities: ['terminal'] }); + }); + + ws.on('message', (data) => { + try { + const msg = JSON.parse(typeof data === 'string' ? data : data.toString()); + + if (msg.type === 'registered') { + console.log(`[pty-sidecar] registered with API server (id=${msg.id})`); + return; + } + + handleCommand(ws, msg); + } catch { + // skip malformed messages } }); ws.on('close', () => { - // Detach WS from session — do NOT kill PTY - if (currentSessionId) { - const session = sessions.get(currentSessionId); - if (session && session.ws === ws) { - session.ws = null; - } - } + console.log('[pty-sidecar] disconnected from API server'); + ws = null; + scheduleReconnect(); }); -}); -server.listen(port, host, () => { - console.log(`[terminal-sidecar] listening on ${host}:${port}`); -}); + ws.on('error', (err) => { + if (reconnectAttempt <= 1) { + console.error(`[pty-sidecar] connection error: ${err.message ?? err}`); + } + // onclose will fire after this + }); +} + +function scheduleReconnect() { + if (reconnectTimer) return; + const delay = RECONNECT_DELAYS[Math.min(reconnectAttempt, RECONNECT_DELAYS.length - 1)]; + reconnectAttempt++; + reconnectTimer = setTimeout(() => { + reconnectTimer = null; + connect(); + }, delay); +} + +// Start connecting +connect(); diff --git a/src/servers/api/terminal/websocket.ts b/src/servers/api/terminal/websocket.ts index 17a01b9e..d687c4c4 100644 --- a/src/servers/api/terminal/websocket.ts +++ b/src/servers/api/terminal/websocket.ts @@ -1,8 +1,9 @@ import type { ServerWebSocket } from 'bun'; import { mkdirSync } from 'node:fs'; import { dirname, join } from 'node:path'; -import { fileURLToPath } from 'node:url'; import { getHomeDir } from '@@/data-path'; +import { sendPtyCommand, sendPtyCommandAsync, on, isTerminalConnected } from '@@/sidecar-registry'; +import type { PtyInitConfig } from '../../sidecar/protocol'; type WSData = { userId: number; @@ -15,17 +16,20 @@ type WSData = { cols?: number; rows?: number; }; + type BridgeSession = { client: ServerWebSocket; - sidecar: WebSocket | null; - pendingMessages: string[]; + sessionId: string; + unsubOutput: (() => void) | null; + unsubExit: (() => void) | null; }; -const HOST_SIDECAR_PORT = 5338; - const sessions = new Map, BridgeSession>(); -let hostSidecarProcess: ReturnType | null = null; +let idCounter = 0; +function nextId(): string { + return `pty_${Date.now()}_${++idCounter}`; +} const sendOutput = (ws: ServerWebSocket, data: string) => { try { @@ -35,105 +39,6 @@ const sendOutput = (ws: ServerWebSocket, data: string) => { } }; -const connectSidecar = async (): Promise => { - const delays = [200, 300, 500, 800, 1200, 1600, 2000]; - let lastError: Error | null = null; - - for (const delay of delays) { - try { - const ws = await new Promise((resolve, reject) => { - const socket = new WebSocket(`ws://127.0.0.1:${HOST_SIDECAR_PORT}`); - const timeout = setTimeout(() => { - try { - socket.close(); - } catch { - // ignore - } - reject(new Error('Terminal sidecar timeout')); - }, 2000); - - socket.addEventListener('open', () => { - clearTimeout(timeout); - resolve(socket); - }); - socket.addEventListener('error', () => { - clearTimeout(timeout); - reject(new Error('Terminal sidecar connection failed')); - }); - }); - - return ws; - } catch (err) { - lastError = err instanceof Error ? err : new Error('Terminal sidecar connection failed'); - await new Promise((resolve) => setTimeout(resolve, delay)); - } - } - - throw lastError ?? new Error('Terminal sidecar connection failed'); -}; - -const sidecarAlive = async (): Promise => { - try { - const res = await fetch(`http://127.0.0.1:${HOST_SIDECAR_PORT}`, { signal: AbortSignal.timeout(500) }); - return res.ok; - } catch { - return false; - } -}; - -const killSidecarOnPort = (port: number) => { - try { - const result = Bun.spawnSync({ cmd: ['fuser', '-k', `${port}/tcp`], stdout: 'ignore', stderr: 'ignore' }); - if (result.exitCode === 0) console.log(`[terminal] killed stale sidecar on port ${port}`); - } catch { - // fuser not available or failed - } -}; - -const startHostSidecar = async () => { - if (hostSidecarProcess) { - hostSidecarProcess.kill(); - await hostSidecarProcess.exited.catch(() => {}); - hostSidecarProcess = null; - } - - killSidecarOnPort(HOST_SIDECAR_PORT); - await new Promise((resolve) => setTimeout(resolve, 200)); - - const sidecarPath = fileURLToPath(new URL('./pty-sidecar.mjs', import.meta.url)); - hostSidecarProcess = Bun.spawn({ - cmd: ['node', sidecarPath], - env: { ...process.env, TERMINAL_PTY_PORT: String(HOST_SIDECAR_PORT) }, - stdout: 'inherit', - stderr: 'inherit', - }); - console.log(`[terminal] host sidecar started on port ${HOST_SIDECAR_PORT}`); -}; - -export const ensureHostSidecar = async () => { - const alive = await sidecarAlive(); - if (!alive) { - await startHostSidecar(); - for (let i = 0; i < 10; i++) { - await new Promise((r) => setTimeout(r, 300)); - if (await sidecarAlive()) return; - } - throw new Error('Host sidecar failed to start'); - } -}; - -export const initTerminalSidecars = async () => { - // Sidecar is managed by pm2 — wait for it to be available - for (let i = 0; i < 15; i++) { - if (await sidecarAlive()) { - console.log(`[terminal] host sidecar already running on port ${HOST_SIDECAR_PORT}`); - return; - } - await new Promise((r) => setTimeout(r, 500)); - } - console.warn(`[terminal] host sidecar not detected on port ${HOST_SIDECAR_PORT} — terminals will retry on connect`); -}; - const resolveCwd = (home: string, cwd?: string) => { if (!cwd || cwd === '~') return home; if (cwd.startsWith('~/')) return join(home, cwd.slice(2)); @@ -150,102 +55,126 @@ export const terminalWebsocket = { `[terminal] open: email=${email} username=${username} role=${role} sandboxed=${sandboxed} isHost=${isHost}`, ); - const session: BridgeSession = { client: ws, sidecar: null, pendingMessages: [] }; - sessions.set(ws, session); - - let sidecar: WebSocket | null = null; - try { - sidecar = await connectSidecar(); - console.log('[terminal] sidecar connected'); - } catch (err) { - const message = err instanceof Error ? err.message : 'Failed to connect terminal sidecar'; - console.error('[terminal] sidecar connection failed:', message); - sendOutput(ws, `\r\n[Terminal error] ${message}\r\n`); - sessions.delete(ws); + if (!isTerminalConnected()) { + sendOutput(ws, '\r\n[Terminal error] PTY sidecar is not connected\r\n'); return; } - session.sidecar = sidecar; + const sessionId = ws.data.sessionId ?? (isHost ? `host-${ws.data.userId}` : `default-${ws.data.userId}`); - sidecar.addEventListener('message', (ev) => { - try { - if (typeof ev.data === 'string') { - ws.send(ev.data); - } else { - ws.send(new TextDecoder().decode(ev.data)); - } - } catch { - // ws already closed - } - }); + // Build PTY init config + let config: PtyInitConfig; if (isHost) { - // Super Admin host terminal — spawn as the service user directly - sidecar.send( - JSON.stringify({ - type: 'init', - host: true, - sessionId: ws.data.sessionId ?? `host-${ws.data.userId}`, - shell: { command: process.env.SHELL ?? '/bin/zsh', args: ['-i'] }, - cwd: resolveCwd(process.env.HOME!, ws.data.cwd), - homeDir: process.env.HOME, - userLabel: email, - cols: ws.data.cols, - rows: ws.data.rows, - }), - ); + config = { + sessionId, + host: true, + shell: { command: process.env.SHELL ?? '/bin/zsh', args: ['-i'] }, + cwd: resolveCwd(process.env.HOME!, ws.data.cwd), + homeDir: process.env.HOME!, + userLabel: email, + cols: ws.data.cols, + rows: ws.data.rows, + }; } else { - // User terminal — spawn as the target Linux user via sudo -u const homeDir = getHomeDir(email); mkdirSync(dirname(homeDir), { recursive: true }); mkdirSync(homeDir, { recursive: true }); - sidecar.send( - JSON.stringify({ - type: 'init', - username, - sessionId: ws.data.sessionId ?? `default-${ws.data.userId}`, - cwd: resolveCwd(homeDir, ws.data.cwd), - homeDir, - userLabel: email, - cols: ws.data.cols, - rows: ws.data.rows, - }), - ); + config = { + sessionId, + username, + cwd: resolveCwd(homeDir, ws.data.cwd), + homeDir, + userLabel: email, + cols: ws.data.cols, + rows: ws.data.rows, + }; } - for (const msg of session.pendingMessages) sidecar.send(msg); - session.pendingMessages = []; + // Subscribe to events for this session + const unsubOutput = on('pty:output', (msg) => { + if (msg.type === 'pty:output' && msg.sessionId === sessionId) { + try { + ws.send(JSON.stringify({ type: 'output', data: msg.data })); + } catch { + // ws already closed + } + } + }); + + const unsubExit = on('pty:exit', (msg) => { + if (msg.type === 'pty:exit' && msg.sessionId === sessionId) { + try { + ws.send(JSON.stringify({ type: 'exit' })); + } catch { + // ws already closed + } + } + }); + + const session: BridgeSession = { client: ws, sessionId, unsubOutput, unsubExit }; + sessions.set(ws, session); + + // Send init command to PTY sidecar + try { + await sendPtyCommandAsync({ type: 'pty:init', id: nextId(), sessionId, config }); + } catch (err) { + const message = err instanceof Error ? err.message : 'Failed to initialize terminal'; + sendOutput(ws, `\r\n[Terminal error] ${message}\r\n`); + unsubOutput(); + unsubExit(); + sessions.delete(ws); + } }, message(ws: ServerWebSocket, raw: string | Buffer) { const session = sessions.get(ws); if (!session) return; - const payload = typeof raw === 'string' ? raw : raw.toString(); - - if (!session.sidecar || session.sidecar.readyState !== WebSocket.OPEN) { - session.pendingMessages.push(payload); - return; - } - try { - session.sidecar.send(payload); + const payload = typeof raw === 'string' ? raw : raw.toString(); + const msg = JSON.parse(payload); + + switch (msg.type) { + case 'input': + sendPtyCommand({ type: 'pty:input', id: nextId(), sessionId: session.sessionId, data: msg.data ?? '' }); + break; + case 'resize': + if (msg.cols > 0 && msg.rows > 0) { + sendPtyCommand({ + type: 'pty:resize', + id: nextId(), + sessionId: session.sessionId, + cols: msg.cols, + rows: msg.rows, + }); + } + break; + case 'cwd': + if (msg.path) { + sendPtyCommand({ + type: 'pty:input', + id: nextId(), + sessionId: session.sessionId, + data: `cd ${JSON.stringify(msg.path)}\r`, + }); + } + break; + } } catch { - // ignore + // ignore malformed messages } }, close(ws: ServerWebSocket) { const session = sessions.get(ws); - if (session?.sidecar) { - try { - session.sidecar.close(); - } catch { - // ignore - } + if (session) { + session.unsubOutput?.(); + session.unsubExit?.(); + // Don't kill PTY — it can be reattached + sessions.delete(ws); } - sessions.delete(ws); }, drain() {}, @@ -253,8 +182,8 @@ export const terminalWebsocket = { export const broadcastPanelRefresh = (email: string) => { const msg = JSON.stringify({ type: 'panel-refresh' }); - for (const [ws, session] of sessions) { - if (ws.data.email === email && session.sidecar) { + for (const [ws] of sessions) { + if (ws.data.email === email) { try { ws.send(msg); } catch { @@ -263,12 +192,3 @@ export const broadcastPanelRefresh = (email: string) => { } } }; - -export const stopAllSidecars = async () => { - if (hostSidecarProcess) { - hostSidecarProcess.kill(); - await hostSidecarProcess.exited.catch(() => {}); - hostSidecarProcess = null; - console.log('[terminal] host sidecar stopped'); - } -}; diff --git a/src/servers/channels/discord/handler.ts b/src/servers/channels/discord/handler.ts index 634dae64..ae89c2c2 100644 --- a/src/servers/channels/discord/handler.ts +++ b/src/servers/channels/discord/handler.ts @@ -4,7 +4,7 @@ import { sendAndAwait, getSessionModel, setSessionModel } from '../send-and-awai import { consumePairingCode } from '../pairing'; import { chunkMessage } from './chunker'; import { listPiModels } from '@@/api/pi/list-models'; -import { enqueueJob } from '../../sidecar-client'; +import { enqueueJob } from '../../sidecar-registry'; import { readJob } from '@@/queue/storage'; import { openEmailDb } from '@@/api/email/email-db'; import type { ModelInfo } from '@@/api/pi/types'; @@ -95,9 +95,9 @@ async function handleEmailSync(ctx: CommandContext): Promise { return; } - const newest = db.query( - 'SELECT from_name, from_address, subject FROM emails WHERE deleted = 0 ORDER BY date DESC LIMIT ?', - ).all(Math.min(newCount, 20)) as Array<{ from_name: string | null; from_address: string; subject: string }>; + const newest = db + .query('SELECT from_name, from_address, subject FROM emails WHERE deleted = 0 ORDER BY date DESC LIMIT ?') + .all(Math.min(newCount, 20)) as Array<{ from_name: string | null; from_address: string; subject: string }>; db.close(); const lines = newest.map((e) => { @@ -128,9 +128,7 @@ async function handleCommand(ctx: CommandContext): Promise { '`!model` — show current model\n' + '`!model ` — switch model\n' + '`!models` — list available models', - email: - '**Email:**\n' + - '`!email sync` — sync Gmail and show new emails', + email: '**Email:**\n' + '`!email sync` — sync Gmail and show new emails', }; if (lower === '!help' || lower.startsWith('!help ')) { @@ -140,7 +138,11 @@ async function handleCommand(ctx: CommandContext): Promise { return true; } if (topic) { - await channel.send(`Unknown topic: \`${topic}\`\nAvailable: ${Object.keys(helpSections).map((k) => `\`${k}\``).join(', ')}`); + await channel.send( + `Unknown topic: \`${topic}\`\nAvailable: ${Object.keys(helpSections) + .map((k) => `\`${k}\``) + .join(', ')}`, + ); return true; } const full = Object.values(helpSections).join('\n\n'); @@ -172,7 +174,11 @@ async function handleCommand(ctx: CommandContext): Promise { if (lower === '!model') { const current = getSessionModel('discord', ctx.userId, discordId); - await channel.send(current ? `Current model: **${current}**` : 'No active session yet — the default model will be used on your next message.'); + await channel.send( + current + ? `Current model: **${current}**` + : 'No active session yet — the default model will be used on your next message.', + ); return true; } @@ -229,17 +235,23 @@ export async function handleDiscordMessage(message: DiscordMessage): Promise void; - reject: (error: Error) => void; - timer: Timer; -}; - -type EventHandler = (event: SidecarEvent) => void; - -let ws: WebSocket | null = null; -let connected = false; -let reconnectAttempt = 0; -let reconnectTimer: Timer | null = null; -const pending = new Map(); -const eventHandlers = new Map>(); -let cachedState: SidecarState | null = null; - -// ── Connection management ── - -function connect() { - if (ws && (ws.readyState === WebSocket.CONNECTING || ws.readyState === WebSocket.OPEN)) return; - - try { - ws = new WebSocket(SIDECAR_URL); - } catch { - scheduleReconnect(); - return; - } - - ws.onopen = () => { - connected = true; - reconnectAttempt = 0; - console.log('[sidecar-client] connected'); - - // Sync state on connect - sendCommand({ type: 'state:sync', id: nextId() }).then((res) => { - if (res.type === 'state:sync') { - cachedState = res.state; - console.log('[sidecar-client] state synced'); - } - }).catch(() => { /* best effort */ }); - }; - - ws.onmessage = (event) => { - try { - const msg = JSON.parse(event.data as string) as SidecarEvent; - - // Check if this is a response to a pending request - if ('id' in msg && msg.id && pending.has(msg.id)) { - const req = pending.get(msg.id)!; - pending.delete(msg.id); - clearTimeout(req.timer); - req.resolve(msg); - return; - } - - // Otherwise dispatch as event - dispatchEvent(msg); - } catch { - // skip malformed messages - } - }; - - ws.onclose = () => { - connected = false; - ws = null; - rejectAllPending('WebSocket disconnected'); - scheduleReconnect(); - }; - - ws.onerror = () => { - // onclose will fire after this - }; -} - -function scheduleReconnect() { - if (reconnectTimer) return; - const delay = RECONNECT_DELAYS[Math.min(reconnectAttempt, RECONNECT_DELAYS.length - 1)]!; - reconnectAttempt++; - reconnectTimer = setTimeout(() => { - reconnectTimer = null; - connect(); - }, delay); -} - -function rejectAllPending(reason: string) { - for (const [id, req] of pending) { - clearTimeout(req.timer); - req.reject(new Error(reason)); - } - pending.clear(); -} - -// ── Event dispatch ── - -function dispatchEvent(msg: SidecarEvent) { - const handlers = eventHandlers.get(msg.type); - if (handlers) { - for (const handler of handlers) { - try { handler(msg); } catch { /* ignore */ } - } - } -} - -export function on(eventType: string, handler: EventHandler): () => void { - if (!eventHandlers.has(eventType)) { - eventHandlers.set(eventType, new Set()); - } - eventHandlers.get(eventType)!.add(handler); - return () => { - eventHandlers.get(eventType)?.delete(handler); - }; -} - -// ── Command sending ── - -let idCounter = 0; -function nextId(): string { - return `sc_${Date.now()}_${++idCounter}`; -} - -const DEFAULT_TIMEOUT_MS = 30_000; -const LONG_TIMEOUT_MS = 6 * 60 * 1000; // 6 min for claude spawn - -function sendCommand(cmd: SidecarCommand, timeoutMs = DEFAULT_TIMEOUT_MS): Promise { - return new Promise((resolve, reject) => { - if (!ws || ws.readyState !== WebSocket.OPEN) { - reject(new Error('Sidecar not connected')); - return; - } - - const timer = setTimeout(() => { - pending.delete(cmd.id); - reject(new Error(`Sidecar command ${cmd.type} timed out`)); - }, timeoutMs); - - pending.set(cmd.id, { resolve, reject, timer }); - ws.send(JSON.stringify(cmd)); - }); -} - -function sendFire(cmd: SidecarCommand): void { - if (ws && ws.readyState === WebSocket.OPEN) { - ws.send(JSON.stringify(cmd)); - } -} - -// ── Public API ── - -export function isConnected(): boolean { - return connected; -} - -export function getCachedState(): SidecarState | null { - return cachedState; -} - -export async function syncState(): Promise { - const res = await sendCommand({ type: 'state:sync', id: nextId() }); - if (res.type === 'state:sync') { - cachedState = res.state; - return res.state; - } - throw new Error('Unexpected response'); -} - -export async function getProxySecret(): Promise { - if (cachedState?.proxySecret) return cachedState.proxySecret; - const res = await sendCommand({ type: 'proxy:secret', id: nextId() }); - if (res.type === 'proxy:secret') return res.secret; - throw new Error('Failed to get proxy secret'); -} - -export function getProxySecretSync(): string { - return cachedState?.proxySecret ?? ''; -} - -// ── Claude Code ── - -export async function spawnClaude(params: ClaudeSpawnParams): Promise { - const res = await sendCommand({ type: 'claude:spawn', id: nextId(), params }, LONG_TIMEOUT_MS); - if (res.type === 'claude:result') return res.result; - if (res.type === 'claude:error') throw new Error(res.error); - throw new Error('Unexpected response'); -} - -export async function spawnClaudeStreaming(params: ClaudeSpawnStreamingParams): Promise { - const res = await sendCommand({ type: 'claude:spawn-streaming', id: nextId(), params }); - if (res.type === 'claude:spawned') return; - if (res.type === 'claude:error') throw new Error(res.error); - throw new Error('Unexpected response'); -} - -export function killClaude(sessionKey: string): void { - sendFire({ type: 'claude:kill', id: nextId(), sessionKey }); -} - -export function clearClaudeSession(sessionKey: string): void { - sendFire({ type: 'claude:clear-session', id: nextId(), sessionKey }); -} - -export function onClaudeEvent(handler: (sessionKey: string, event: PiEvent) => void): () => void { - return on('claude:event', (msg) => { - if (msg.type === 'claude:event') { - handler(msg.sessionKey, msg.event); - } - }); -} - -// ── Pi ── - -export async function spawnPi(params: PiSpawnParams): Promise { - const res = await sendCommand({ type: 'pi:spawn', id: nextId(), params }); - if (res.type === 'pi:spawned') return; - if (res.type === 'pi:error') throw new Error(res.error); - throw new Error('Unexpected response'); -} - -export function sendPiPrompt(sessionId: string, prompt: string, requestId: string): void { - sendFire({ type: 'pi:prompt', id: nextId(), sessionId, prompt, requestId }); -} - -export function abortPi(sessionId: string, requestId: string): void { - sendFire({ type: 'pi:abort', id: nextId(), sessionId, requestId }); -} - -export function killPi(sessionId: string): void { - sendFire({ type: 'pi:kill', id: nextId(), sessionId }); -} - -export function setPiThinking(sessionId: string, level: string): void { - sendFire({ type: 'pi:set-thinking', id: nextId(), sessionId, level }); -} - -export function onPiEvent(handler: (sessionId: string, event: PiEvent) => void): () => void { - return on('pi:event', (msg) => { - if (msg.type === 'pi:event') { - handler(msg.sessionId, msg.event); - } - }); -} - -// ── Queue ── - -export async function enqueueJob(params: EnqueueParams): Promise { - const res = await sendCommand({ type: 'queue:enqueue', id: nextId(), params }); - if (res.type === 'queue:enqueued') return res.job; - if (res.type === 'queue:error') throw new Error(res.error); - throw new Error('Unexpected response'); -} - -export async function cancelJob(jobId: string): Promise { - const res = await sendCommand({ type: 'queue:cancel', id: nextId(), jobId }); - if (res.type === 'queue:cancelled') return res.job; - if (res.type === 'queue:error') throw new Error(res.error); - throw new Error('Unexpected response'); -} - -export async function listJobs(): Promise { - const res = await sendCommand({ type: 'queue:list', id: nextId() }); - if (res.type === 'queue:list') return res.jobs; - throw new Error('Unexpected response'); -} - -export async function getJob(jobId: string): Promise { - const res = await sendCommand({ type: 'queue:get', id: nextId(), jobId }); - if (res.type === 'queue:get') return res.job; - throw new Error('Unexpected response'); -} - -// ── Health check ── - -export async function isSidecarAlive(): Promise { - try { - const res = await fetch(`http://127.0.0.1:${SIDECAR_PORT}/`, { signal: AbortSignal.timeout(1000) }); - const text = await res.text(); - return text === 'process-sidecar'; - } catch { - return false; - } -} - -// ── Init ── - -export function initSidecarClient(): void { - connect(); -} diff --git a/src/servers/sidecar-registry.ts b/src/servers/sidecar-registry.ts new file mode 100644 index 00000000..6434d30c --- /dev/null +++ b/src/servers/sidecar-registry.ts @@ -0,0 +1,315 @@ +import type { ServerWebSocket } from 'bun'; +import type { + SidecarCommand, + SidecarEvent, + SidecarState, + ClaudeSpawnParams, + ClaudeSpawnStreamingParams, + ClaudeCodeResult, + PiSpawnParams, + PtyCommand, + PtyEvent, +} from './sidecar/protocol'; +import type { SidecarRegistration } from './sidecar/registration-protocol'; +import type { PiEvent } from './api/pi/types'; +import type { Job, EnqueueParams } from './queue/types'; + +// ── Types ── + +type RegisteredSidecar = { + id: string; + name: string; + capabilities: string[]; + ws: ServerWebSocket; +}; + +type PendingRequest = { + resolve: (value: any) => void; + reject: (error: Error) => void; + timer: Timer; +}; + +type EventHandler = (event: SidecarEvent | PtyEvent) => void; + +// ── State ── + +const sidecars = new Map(); +const pending = new Map(); +const eventHandlers = new Map>(); +let cachedState: SidecarState | null = null; +let idCounter = 0; + +function nextId(): string { + return `sr_${Date.now()}_${++idCounter}`; +} + +// ── Registration ── + +export function registerSidecar(ws: ServerWebSocket, registration: SidecarRegistration): string { + // If a sidecar with the same name is already registered, unregister it first + for (const [id, sc] of sidecars) { + if (sc.name === registration.name) { + console.log(`[registry] replacing existing sidecar "${registration.name}" (id=${id})`); + unregisterSidecar(id); + break; + } + } + + const id = `sc_${registration.name}_${Date.now()}`; + sidecars.set(id, { + id, + name: registration.name, + capabilities: registration.capabilities, + ws, + }); + console.log( + `[registry] registered sidecar "${registration.name}" (id=${id}, capabilities=[${registration.capabilities.join(', ')}])`, + ); + return id; +} + +export function unregisterSidecar(id: string): void { + const sc = sidecars.get(id); + if (!sc) return; + sidecars.delete(id); + console.log(`[registry] unregistered sidecar "${sc.name}" (id=${id})`); + + // Reject all pending requests for this sidecar + for (const [reqId, req] of pending) { + clearTimeout(req.timer); + req.reject(new Error(`Sidecar "${sc.name}" disconnected`)); + pending.delete(reqId); + } + + // Clear cached state if the process sidecar disconnects + if (sc.capabilities.includes('proxy')) { + cachedState = null; + } +} + +export function handleSidecarMessage(id: string, msg: SidecarEvent | PtyEvent): void { + // Check if this is a response to a pending request + if ('id' in msg && msg.id && pending.has(msg.id)) { + const req = pending.get(msg.id)!; + pending.delete(msg.id); + clearTimeout(req.timer); + req.resolve(msg); + return; + } + + // Otherwise dispatch as event + dispatchEvent(msg); +} + +// ── Lookup ── + +function findSidecarByCapability(cap: string): RegisteredSidecar | undefined { + for (const sc of sidecars.values()) { + if (sc.capabilities.includes(cap)) return sc; + } + return undefined; +} + +function requireSidecar(cap: string): RegisteredSidecar { + const sc = findSidecarByCapability(cap); + if (!sc) throw new Error(`No sidecar with capability "${cap}" is connected`); + return sc; +} + +// ── Event dispatch ── + +function dispatchEvent(msg: SidecarEvent | PtyEvent) { + const handlers = eventHandlers.get(msg.type); + if (handlers) { + for (const handler of handlers) { + try { + handler(msg); + } catch { + /* ignore */ + } + } + } +} + +export function on(eventType: string, handler: EventHandler): () => void { + if (!eventHandlers.has(eventType)) { + eventHandlers.set(eventType, new Set()); + } + eventHandlers.get(eventType)!.add(handler); + return () => { + eventHandlers.get(eventType)?.delete(handler); + }; +} + +// ── Command sending ── + +const DEFAULT_TIMEOUT_MS = 30_000; +const LONG_TIMEOUT_MS = 6 * 60 * 1000; + +function sendCommand(cap: string, cmd: SidecarCommand | PtyCommand, timeoutMs = DEFAULT_TIMEOUT_MS): Promise { + return new Promise((resolve, reject) => { + const sc = findSidecarByCapability(cap); + if (!sc) { + reject(new Error(`No sidecar with capability "${cap}" is connected`)); + return; + } + + const timer = setTimeout(() => { + pending.delete((cmd as any).id); + reject(new Error(`Sidecar command ${cmd.type} timed out`)); + }, timeoutMs); + + pending.set((cmd as any).id, { resolve, reject, timer }); + sc.ws.send(JSON.stringify(cmd)); + }); +} + +function sendFire(cap: string, cmd: SidecarCommand | PtyCommand): void { + const sc = findSidecarByCapability(cap); + if (sc) { + sc.ws.send(JSON.stringify(cmd)); + } +} + +// ── Public API (same signatures as sidecar-client.ts) ── + +export function isConnected(): boolean { + return findSidecarByCapability('proxy') !== undefined; +} + +export function getCachedState(): SidecarState | null { + return cachedState; +} + +export async function syncState(): Promise { + const res = await sendCommand('proxy', { type: 'state:sync', id: nextId() }); + if (res.type === 'state:sync') { + cachedState = res.state; + return res.state; + } + throw new Error('Unexpected response'); +} + +export async function getProxySecret(): Promise { + if (cachedState?.proxySecret) return cachedState.proxySecret; + const res = await sendCommand('proxy', { type: 'proxy:secret', id: nextId() }); + if (res.type === 'proxy:secret') return res.secret; + throw new Error('Failed to get proxy secret'); +} + +export function getProxySecretSync(): string { + return cachedState?.proxySecret ?? ''; +} + +// ── Claude Code ── + +export async function spawnClaude(params: ClaudeSpawnParams): Promise { + const res = await sendCommand('claude', { type: 'claude:spawn', id: nextId(), params }, LONG_TIMEOUT_MS); + if (res.type === 'claude:result') return res.result; + if (res.type === 'claude:error') throw new Error(res.error); + throw new Error('Unexpected response'); +} + +export async function spawnClaudeStreaming(params: ClaudeSpawnStreamingParams): Promise { + const res = await sendCommand('claude', { type: 'claude:spawn-streaming', id: nextId(), params }); + if (res.type === 'claude:spawned') return; + if (res.type === 'claude:error') throw new Error(res.error); + throw new Error('Unexpected response'); +} + +export function killClaude(sessionKey: string): void { + sendFire('claude', { type: 'claude:kill', id: nextId(), sessionKey }); +} + +export function clearClaudeSession(sessionKey: string): void { + sendFire('claude', { type: 'claude:clear-session', id: nextId(), sessionKey }); +} + +export function onClaudeEvent(handler: (sessionKey: string, event: PiEvent) => void): () => void { + return on('claude:event', (msg) => { + if (msg.type === 'claude:event') { + handler( + (msg as SidecarEvent & { type: 'claude:event' }).sessionKey, + (msg as SidecarEvent & { type: 'claude:event' }).event, + ); + } + }); +} + +// ── Pi ── + +export async function spawnPi(params: PiSpawnParams): Promise { + const res = await sendCommand('pi', { type: 'pi:spawn', id: nextId(), params }); + if (res.type === 'pi:spawned') return; + if (res.type === 'pi:error') throw new Error(res.error); + throw new Error('Unexpected response'); +} + +export function sendPiPrompt(sessionId: string, prompt: string, requestId: string): void { + sendFire('pi', { type: 'pi:prompt', id: nextId(), sessionId, prompt, requestId }); +} + +export function abortPi(sessionId: string, requestId: string): void { + sendFire('pi', { type: 'pi:abort', id: nextId(), sessionId, requestId }); +} + +export function killPi(sessionId: string): void { + sendFire('pi', { type: 'pi:kill', id: nextId(), sessionId }); +} + +export function setPiThinking(sessionId: string, level: string): void { + sendFire('pi', { type: 'pi:set-thinking', id: nextId(), sessionId, level }); +} + +export function onPiEvent(handler: (sessionId: string, event: PiEvent) => void): () => void { + return on('pi:event', (msg) => { + if (msg.type === 'pi:event') { + handler( + (msg as SidecarEvent & { type: 'pi:event' }).sessionId, + (msg as SidecarEvent & { type: 'pi:event' }).event, + ); + } + }); +} + +// ── Queue ── + +export async function enqueueJob(params: EnqueueParams): Promise { + const res = await sendCommand('queue', { type: 'queue:enqueue', id: nextId(), params }); + if (res.type === 'queue:enqueued') return res.job; + if (res.type === 'queue:error') throw new Error(res.error); + throw new Error('Unexpected response'); +} + +export async function cancelJob(jobId: string): Promise { + const res = await sendCommand('queue', { type: 'queue:cancel', id: nextId(), jobId }); + if (res.type === 'queue:cancelled') return res.job; + if (res.type === 'queue:error') throw new Error(res.error); + throw new Error('Unexpected response'); +} + +export async function listJobs(): Promise { + const res = await sendCommand('queue', { type: 'queue:list', id: nextId() }); + if (res.type === 'queue:list') return res.jobs; + throw new Error('Unexpected response'); +} + +export async function getJob(jobId: string): Promise { + const res = await sendCommand('queue', { type: 'queue:get', id: nextId(), jobId }); + if (res.type === 'queue:get') return res.job; + throw new Error('Unexpected response'); +} + +// ── Terminal (PTY sidecar) ── + +export function sendPtyCommand(cmd: PtyCommand): void { + sendFire('terminal', cmd); +} + +export async function sendPtyCommandAsync(cmd: PtyCommand, timeoutMs = DEFAULT_TIMEOUT_MS): Promise { + return sendCommand('terminal', cmd, timeoutMs); +} + +export function isTerminalConnected(): boolean { + return findSidecarByCapability('terminal') !== undefined; +} diff --git a/src/servers/sidecar/connect.ts b/src/servers/sidecar/connect.ts new file mode 100644 index 00000000..81edaf63 --- /dev/null +++ b/src/servers/sidecar/connect.ts @@ -0,0 +1,135 @@ +import type { SidecarCommand, SidecarEvent, PtyCommand, PtyEvent } from './protocol'; +import type { SidecarRegistration, RegistrationAck } from './registration-protocol'; + +type AnyCommand = SidecarCommand | PtyCommand; +type AnyEvent = SidecarEvent | PtyEvent; + +type SidecarConnectorConfig = { + apiUrl: string; // ws://127.0.0.1:5000/api/sidecar/register + name: string; + capabilities: string[]; + onCommand: (cmd: AnyCommand, reply: (msg: AnyEvent) => void) => void; + onConnected?: () => void; + onDisconnected?: () => void; +}; + +type SidecarConnection = { + send: (msg: AnyEvent) => void; + destroy: () => void; + isConnected: () => boolean; +}; + +const RECONNECT_DELAYS = [200, 500, 1000, 2000, 4000, 8000, 15000]; + +export function createSidecarConnector(config: SidecarConnectorConfig): SidecarConnection { + let ws: WebSocket | null = null; + let connected = false; + let reconnectAttempt = 0; + let reconnectTimer: Timer | null = null; + let destroyed = false; + let sidecarId: string | null = null; + + function connect() { + if (destroyed) return; + if (ws && (ws.readyState === WebSocket.CONNECTING || ws.readyState === WebSocket.OPEN)) return; + + try { + ws = new WebSocket(config.apiUrl); + } catch { + scheduleReconnect(); + return; + } + + ws.onopen = () => { + reconnectAttempt = 0; + + // Send registration + const registration: SidecarRegistration = { + type: 'register', + name: config.name, + capabilities: config.capabilities, + }; + ws!.send(JSON.stringify(registration)); + }; + + ws.onmessage = (event) => { + try { + const msg = JSON.parse(event.data as string); + + // Handle registration ack + if (msg.type === 'registered') { + const ack = msg as RegistrationAck; + sidecarId = ack.id; + connected = true; + console.log(`[sidecar] registered with API server (id=${sidecarId})`); + config.onConnected?.(); + return; + } + + // Handle commands from API server + const reply = (response: AnyEvent) => { + if (ws && ws.readyState === WebSocket.OPEN) { + ws.send(JSON.stringify(response)); + } + }; + config.onCommand(msg as AnyCommand, reply); + } catch { + // skip malformed messages + } + }; + + ws.onclose = () => { + const wasConnected = connected; + connected = false; + sidecarId = null; + ws = null; + if (wasConnected) { + console.log('[sidecar] disconnected from API server'); + config.onDisconnected?.(); + } + scheduleReconnect(); + }; + + ws.onerror = () => { + // onclose will fire after this + }; + } + + function scheduleReconnect() { + if (destroyed || reconnectTimer) return; + const delay = RECONNECT_DELAYS[Math.min(reconnectAttempt, RECONNECT_DELAYS.length - 1)]!; + reconnectAttempt++; + reconnectTimer = setTimeout(() => { + reconnectTimer = null; + connect(); + }, delay); + } + + function send(msg: AnyEvent): void { + if (ws && ws.readyState === WebSocket.OPEN) { + ws.send(JSON.stringify(msg)); + } + } + + function destroy(): void { + destroyed = true; + if (reconnectTimer) { + clearTimeout(reconnectTimer); + reconnectTimer = null; + } + if (ws) { + ws.close(); + ws = null; + } + connected = false; + } + + // Start connecting + connect(); + + return { + send, + destroy, + isConnected: () => connected, + }; +} diff --git a/src/servers/sidecar/index.ts b/src/servers/sidecar/index.ts index b3e32eba..c1554096 100644 --- a/src/servers/sidecar/index.ts +++ b/src/servers/sidecar/index.ts @@ -1,4 +1,3 @@ -import type { ServerWebSocket } from 'bun'; import type { SidecarCommand, SidecarEvent, SidecarState } from './protocol'; import { loadState, flushAndSave, acquireLock, releaseLock, getState } from './state'; import { startAnthropicProxy, getProxySecret, ensureProxySecret } from './proxy'; @@ -6,8 +5,9 @@ import * as claudeManager from './claude-manager'; import * as piManager from './pi-manager'; import * as queueRunner from './queue-runner'; import { initEmailCron, stopEmailCron } from './email-cron'; +import { createSidecarConnector } from './connect'; -const PORT = Number(process.env.SIDECAR_PORT ?? '5100'); +const API_URL = process.env.API_URL ?? `ws://127.0.0.1:${process.env.PORT ?? '5000'}`; const startedAt = Date.now(); // ── Startup ── @@ -35,24 +35,7 @@ queueRunner.initQueue().catch((err) => { // Start email sync cron // initEmailCron(); // TODO: re-enable after initial sync testing -// ── WebSocket connections ── - -const clients = new Set>(); - -function broadcast(msg: SidecarEvent) { - const data = JSON.stringify(msg); - for (const ws of clients) { - if (ws.readyState === 1) { - ws.send(data); - } - } -} - -function reply(ws: ServerWebSocket, msg: SidecarEvent) { - if (ws.readyState === 1) { - ws.send(JSON.stringify(msg)); - } -} +// ── State ── function buildState(): SidecarState { return { @@ -65,18 +48,20 @@ function buildState(): SidecarState { // ── Command handlers ── -async function handleCommand(ws: ServerWebSocket, cmd: SidecarCommand) { +type ReplyFn = (msg: SidecarEvent) => void; + +async function handleCommand(cmd: SidecarCommand, reply: ReplyFn) { switch (cmd.type) { case 'ping': - reply(ws, { type: 'pong', id: cmd.id }); + reply({ type: 'pong', id: cmd.id }); break; case 'state:sync': - reply(ws, { type: 'state:sync', id: cmd.id, state: buildState() }); + reply({ type: 'state:sync', id: cmd.id, state: buildState() }); break; case 'proxy:secret': - reply(ws, { type: 'proxy:secret', id: cmd.id, secret: getProxySecret() }); + reply({ type: 'proxy:secret', id: cmd.id, secret: getProxySecret() }); break; // ── Claude Code ── @@ -84,22 +69,22 @@ async function handleCommand(ws: ServerWebSocket, cmd: SidecarCommand) case 'claude:spawn': { try { const result = await claudeManager.spawnClaude(cmd.params); - reply(ws, { type: 'claude:result', id: cmd.id, result }); + reply({ type: 'claude:result', id: cmd.id, result }); } catch (err) { - reply(ws, { type: 'claude:error', id: cmd.id, error: err instanceof Error ? err.message : String(err) }); + reply({ type: 'claude:error', id: cmd.id, error: err instanceof Error ? err.message : String(err) }); } break; } case 'claude:spawn-streaming': { - reply(ws, { type: 'claude:spawned', id: cmd.id, sessionKey: cmd.params.sessionKey }); + reply({ type: 'claude:spawned', id: cmd.id, sessionKey: cmd.params.sessionKey }); const onEvent = (event: import('../api/pi/types').PiEvent) => { - broadcast({ type: 'claude:event', sessionKey: cmd.params.sessionKey, event }); + connection.send({ type: 'claude:event', sessionKey: cmd.params.sessionKey, event }); }; claudeManager.spawnClaudeStreaming(cmd.params, onEvent).catch((err) => { - broadcast({ + connection.send({ type: 'claude:event', sessionKey: cmd.params.sessionKey, event: { type: 'error', message: err instanceof Error ? err.message : String(err) }, @@ -110,12 +95,12 @@ async function handleCommand(ws: ServerWebSocket, cmd: SidecarCommand) case 'claude:kill': claudeManager.killClaudeSession(cmd.sessionKey); - reply(ws, { type: 'claude:killed', id: cmd.id }); + reply({ type: 'claude:killed', id: cmd.id }); break; case 'claude:clear-session': claudeManager.clearSession(cmd.sessionKey); - reply(ws, { type: 'claude:session-cleared', id: cmd.id }); + reply({ type: 'claude:session-cleared', id: cmd.id }); break; // ── Pi ── @@ -123,7 +108,7 @@ async function handleCommand(ws: ServerWebSocket, cmd: SidecarCommand) case 'pi:spawn': { try { const onEvent = (event: import('../api/pi/types').PiEvent) => { - broadcast({ type: 'pi:event', sessionId: cmd.params.sessionId, event }); + connection.send({ type: 'pi:event', sessionId: cmd.params.sessionId, event }); }; await piManager.spawnPi({ @@ -138,9 +123,9 @@ async function handleCommand(ws: ServerWebSocket, cmd: SidecarCommand) onEvent, }); - reply(ws, { type: 'pi:spawned', id: cmd.id, sessionId: cmd.params.sessionId }); + reply({ type: 'pi:spawned', id: cmd.id, sessionId: cmd.params.sessionId }); } catch (err) { - reply(ws, { type: 'pi:error', id: cmd.id, error: err instanceof Error ? err.message : String(err) }); + reply({ type: 'pi:error', id: cmd.id, error: err instanceof Error ? err.message : String(err) }); } break; } @@ -155,7 +140,7 @@ async function handleCommand(ws: ServerWebSocket, cmd: SidecarCommand) case 'pi:kill': piManager.killPiSession(cmd.sessionId); - reply(ws, { type: 'pi:killed', id: cmd.id }); + reply({ type: 'pi:killed', id: cmd.id }); break; case 'pi:set-thinking': @@ -167,9 +152,9 @@ async function handleCommand(ws: ServerWebSocket, cmd: SidecarCommand) case 'queue:enqueue': { try { const job = await queueRunner.enqueue(cmd.params); - reply(ws, { type: 'queue:enqueued', id: cmd.id, job }); + reply({ type: 'queue:enqueued', id: cmd.id, job }); } catch (err) { - reply(ws, { type: 'queue:error', id: cmd.id, error: err instanceof Error ? err.message : String(err) }); + reply({ type: 'queue:error', id: cmd.id, error: err instanceof Error ? err.message : String(err) }); } break; } @@ -177,85 +162,51 @@ async function handleCommand(ws: ServerWebSocket, cmd: SidecarCommand) case 'queue:cancel': { try { const job = await queueRunner.cancelJob(cmd.jobId); - reply(ws, { type: 'queue:cancelled', id: cmd.id, job }); + reply({ type: 'queue:cancelled', id: cmd.id, job }); } catch (err) { - reply(ws, { type: 'queue:error', id: cmd.id, error: err instanceof Error ? err.message : String(err) }); + reply({ type: 'queue:error', id: cmd.id, error: err instanceof Error ? err.message : String(err) }); } break; } case 'queue:list': { const jobs = await queueRunner.listAllJobs(); - reply(ws, { type: 'queue:list', id: cmd.id, jobs }); + reply({ type: 'queue:list', id: cmd.id, jobs }); break; } case 'queue:get': { const job = await queueRunner.readJob(cmd.jobId); - reply(ws, { type: 'queue:get', id: cmd.id, job }); + reply({ type: 'queue:get', id: cmd.id, job }); break; } default: - reply(ws, { type: 'error', id: (cmd as SidecarCommand).id, error: `Unknown command type: ${(cmd as Record).type}` }); + reply({ + type: 'error', + id: (cmd as SidecarCommand).id, + error: `Unknown command type: ${(cmd as Record).type}`, + }); } } -// ── Server ── +// ── Connect to API server ── -const server = Bun.serve({ - port: PORT, - hostname: '127.0.0.1', - - fetch(req, server) { - if (req.headers.get('upgrade') === 'websocket') { - const ok = server.upgrade(req); - if (!ok) return new Response('Upgrade failed', { status: 500 }); - return; - } - - const url = new URL(req.url); - if (url.pathname === '/') return new Response('process-sidecar'); - if (url.pathname === '/health') { - return new Response(JSON.stringify({ - status: 'ok', - uptime: Date.now() - startedAt, - piSessions: piManager.getAllSessions().length, - claudeSessions: claudeManager.getActiveSessionKeys().length, - }), { headers: { 'Content-Type': 'application/json' } }); - } - return new Response('Not found', { status: 404 }); - }, - - websocket: { - open(ws) { - clients.add(ws); - console.log(`[sidecar] client connected (${clients.size} total)`); - }, - message(ws, raw) { - try { - const data = typeof raw === 'string' ? raw : raw.toString(); - const cmd = JSON.parse(data) as SidecarCommand; - handleCommand(ws, cmd); - } catch (err) { - reply(ws, { type: 'error', error: `Invalid message: ${err instanceof Error ? err.message : String(err)}` }); - } - }, - close(ws) { - clients.delete(ws); - console.log(`[sidecar] client disconnected (${clients.size} total)`); - }, - drain() {}, +const connection = createSidecarConnector({ + apiUrl: `${API_URL}/api/sidecar/register`, + name: 'process', + capabilities: ['claude', 'pi', 'queue', 'proxy'], + onCommand(cmd, reply) { + handleCommand(cmd as SidecarCommand, reply as ReplyFn); }, }); -console.log(`[sidecar] listening on 127.0.0.1:${PORT}`); - // ── Graceful shutdown ── async function shutdown(signal: string) { console.log(`[sidecar] ${signal} received, saving state...`); stopEmailCron(); + connection.destroy(); await flushAndSave(); releaseLock(); process.exit(0); diff --git a/src/servers/sidecar/protocol.ts b/src/servers/sidecar/protocol.ts index 72756cba..172875e1 100644 --- a/src/servers/sidecar/protocol.ts +++ b/src/servers/sidecar/protocol.ts @@ -113,3 +113,31 @@ export type PiSpawnParams = { model: string; sessionFile?: string; }; + +// ── PTY types ── + +export type PtyInitConfig = { + sessionId: string; + shell?: { command: string; args?: string[] }; + cwd?: string; + homeDir?: string; + userLabel?: string; + username?: string; + host?: boolean; + cols?: number; + rows?: number; + env?: Record; +}; + +// PTY commands (API → PTY sidecar) +export type PtyCommand = + | { type: 'pty:init'; id: string; sessionId: string; config: PtyInitConfig } + | { type: 'pty:input'; id: string; sessionId: string; data: string } + | { type: 'pty:resize'; id: string; sessionId: string; cols: number; rows: number } + | { type: 'pty:close'; id: string; sessionId: string }; + +// PTY events (PTY sidecar → API) +export type PtyEvent = + | { type: 'pty:ready'; id: string; sessionId: string } + | { type: 'pty:output'; sessionId: string; data: string } + | { type: 'pty:exit'; sessionId: string; exitCode: number; signal?: number }; diff --git a/src/servers/sidecar/registration-protocol.ts b/src/servers/sidecar/registration-protocol.ts new file mode 100644 index 00000000..d30e4418 --- /dev/null +++ b/src/servers/sidecar/registration-protocol.ts @@ -0,0 +1,16 @@ +// ── Sidecar self-registration protocol ── + +// Sidecar → API server on connect +export type SidecarRegistration = { + type: 'register'; + name: string; // 'process' | 'pty' | custom + capabilities: string[]; // ['claude', 'pi', 'queue', 'proxy'] or ['terminal'] +}; + +// API server → sidecar ack +export type RegistrationAck = { + type: 'registered'; + id: string; // server-assigned ID +}; + +export type RegistrationMessage = SidecarRegistration | RegistrationAck;