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 <noreply@anthropic.com>
This commit is contained in:
2026-03-06 08:48:19 +00:00
co-authored by Claude Opus 4.6
parent 6cfa40bad1
commit a0d9ea63f9
17 changed files with 878 additions and 712 deletions
+9 -3
View File
@@ -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,
},
],
+49 -8
View File
@@ -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<ServerWebSocket<WSData>, string>(); // ws → sidecar ID
const sidecarWebsocket = {
open(_ws: ServerWebSocket<WSData>) {
// Wait for registration message
},
message(ws: ServerWebSocket<WSData>, 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<WSData>) {
const id = sidecarConnections.get(ws);
if (id) {
unregisterSidecar(id);
sidecarConnections.delete(ws);
}
},
drain() {},
};
const handlers: Record<string, any> = {
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 () => {
+7 -2
View File
@@ -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<string, unknown>);
const authResult = await resolveAuth(
user.id,
account.authType,
account.email,
account.credentials as Record<string, unknown>,
);
if (!authResult.ok) throw BAD_REQUEST(authResult.error);
// Set status immediately so the UI reflects the queued state
+1 -1
View File
@@ -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';
+1 -1
View File
@@ -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();
+144 -104
View File
@@ -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<string, { term: import('node-pty').IPty, buffer: string, cols: number, rows: number, ws: import('ws').WebSocket | null, initConfig: object }>} */
/** @type {Map<string, { term: import('node-pty').IPty, buffer: string, cols: number, rows: number, initConfig: object }>} */
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();
+105 -185
View File
@@ -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<WSData>;
sidecar: WebSocket | null;
pendingMessages: string[];
sessionId: string;
unsubOutput: (() => void) | null;
unsubExit: (() => void) | null;
};
const HOST_SIDECAR_PORT = 5338;
const sessions = new Map<ServerWebSocket<WSData>, BridgeSession>();
let hostSidecarProcess: ReturnType<typeof import('bun').spawn> | null = null;
let idCounter = 0;
function nextId(): string {
return `pty_${Date.now()}_${++idCounter}`;
}
const sendOutput = (ws: ServerWebSocket<WSData>, data: string) => {
try {
@@ -35,105 +39,6 @@ const sendOutput = (ws: ServerWebSocket<WSData>, data: string) => {
}
};
const connectSidecar = async (): Promise<WebSocket> => {
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<WebSocket>((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<boolean> => {
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<WSData>, 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<WSData>) {
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');
}
};
+26 -14
View File
@@ -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<void> {
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<boolean> {
'`!model` — show current model\n' +
'`!model <id>` — 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<boolean> {
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<boolean> {
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<voi
}
await channel.send(
'I don\'t recognize your Discord account. To link it:\n' +
'1. Go to Officer Settings → Integrations → Discord\n' +
'2. Click "Link Discord" to get a pairing code\n' +
'3. Send the 6-character code to me here',
"I don't recognize your Discord account. To link it:\n" +
'1. Go to Officer Settings → Integrations → Discord\n' +
'2. Click "Link Discord" to get a pairing code\n' +
'3. Send the 6-character code to me here',
);
return;
}
// Handle commands
if (content.startsWith('!')) {
const handled = await handleCommand({ content, channel, userId: linked.user.id, email: linked.user.email, discordId });
const handled = await handleCommand({
content,
channel,
userId: linked.user.id,
email: linked.user.email,
discordId,
});
if (handled) return;
}
+1 -1
View File
@@ -1,6 +1,6 @@
import { logger } from '@@/api/pi/logger';
import type { MessageCost, PiEvent } from '@@/api/pi/types';
import * as sidecar from '@@/sidecar-client';
import * as sidecar from '@@/sidecar-registry';
type ClaudeCodeParams = {
userId: number;
+1 -1
View File
@@ -5,7 +5,7 @@ import { consumePairingCode } from '../pairing';
import { chunkMessage } from './chunker';
import { getTelegramBot } from './bot';
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';
+1 -1
View File
@@ -4,7 +4,7 @@ import { sendAndAwait, getSessionModel, setSessionModel } from '../send-and-awai
import { consumePairingCode } from '../pairing';
import { getWhatsAppClient } from './bot';
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';
-303
View File
@@ -1,303 +0,0 @@
import type {
SidecarCommand,
SidecarEvent,
SidecarState,
ClaudeSpawnParams,
ClaudeSpawnStreamingParams,
ClaudeCodeResult,
PiSpawnParams,
} from './sidecar/protocol';
import type { PiEvent } from './api/pi/types';
import type { Job, EnqueueParams } from './queue/types';
const SIDECAR_PORT = Number(process.env.SIDECAR_PORT ?? '5100');
const SIDECAR_URL = `ws://127.0.0.1:${SIDECAR_PORT}`;
const RECONNECT_DELAYS = [200, 500, 1000, 2000, 4000, 8000, 15000];
type PendingRequest = {
resolve: (value: SidecarEvent) => 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<string, PendingRequest>();
const eventHandlers = new Map<string, Set<EventHandler>>();
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<SidecarEvent> {
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<SidecarState> {
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<string> {
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<ClaudeCodeResult> {
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<void> {
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<void> {
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<Job> {
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<Job | null> {
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<Job[]> {
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<Job | null> {
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<boolean> {
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();
}
+315
View File
@@ -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<any>;
};
type PendingRequest = {
resolve: (value: any) => void;
reject: (error: Error) => void;
timer: Timer;
};
type EventHandler = (event: SidecarEvent | PtyEvent) => void;
// ── State ──
const sidecars = new Map<string, RegisteredSidecar>();
const pending = new Map<string, PendingRequest>();
const eventHandlers = new Map<string, Set<EventHandler>>();
let cachedState: SidecarState | null = null;
let idCounter = 0;
function nextId(): string {
return `sr_${Date.now()}_${++idCounter}`;
}
// ── Registration ──
export function registerSidecar(ws: ServerWebSocket<any>, 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<any> {
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<SidecarState> {
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<string> {
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<ClaudeCodeResult> {
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<void> {
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<void> {
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<Job> {
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<Job | null> {
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<Job[]> {
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<Job | null> {
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<PtyEvent> {
return sendCommand('terminal', cmd, timeoutMs);
}
export function isTerminalConnected(): boolean {
return findSidecarByCapability('terminal') !== undefined;
}
+135
View File
@@ -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,
};
}
+39 -88
View File
@@ -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<ServerWebSocket<unknown>>();
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<unknown>, 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<unknown>, 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<unknown>, 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<unknown>, 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<unknown>, 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<unknown>, 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<unknown>, 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<unknown>, 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<unknown>, 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<string, unknown>).type}` });
reply({
type: 'error',
id: (cmd as SidecarCommand).id,
error: `Unknown command type: ${(cmd as Record<string, unknown>).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);
+28
View File
@@ -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<string, string>;
};
// 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 };
@@ -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;