Merge master into PR #153 (unified Session Manager)

Resolve the 4 conflicted files toward master's merged #146 work while
keeping PR #153's genuinely-new additions:

- app.js: keep the full Escape chain (closeSessionManager +
  closeCommandPalette + closeShortcutOverlay).
- index.html: keep master's Command Palette modal markup alongside the
  PR's Session Manager modal + header button.
- styles.css: keep master's Command Palette + COD-157 shortcut CSS AND
  the PR's COD-130 session-row kebab-menu CSS (both inserted at the same
  spot — reunited each with its own closing brace).
- terminal-ui.js: resolve _buildHistoryItem's main-row click handler to
  master's options.onActivate contract with a liveness + claudeSessionId
  -aware resume default, preserving the PR's two-shape/badges/kebab body.
- panels-ui.js: the PR's pre-#146 Session Manager block auto-merged as a
  duplicate AFTER master's fixed block (last-key-wins regression) — drop
  it, keep master's implementation plus the PR's new
  _onSessionListMaybeChanged.

Backend projectKey plumbing and the SSE live-refresh listeners in app.js
merge additively and are kept as-is.
This commit is contained in:
Codeman maintainer
2026-07-13 00:34:43 +02:00
87 changed files with 13392 additions and 661 deletions
+142 -1
View File
@@ -11,15 +11,30 @@ import { join, resolve } from 'node:path';
import { homedir } from 'node:os';
import type { ApiResponse, CaseInfo } from '../../types.js';
import { ApiErrorCode, createErrorResponse, getErrorMessage } from '../../types.js';
import { CreateCaseSchema, LinkCaseSchema, CaseOrderSchema } from '../schemas.js';
import {
CreateCaseSchema,
LinkCaseSchema,
CaseOrderSchema,
RemoteCaseLinkSchema,
RemoteHostSchema,
} from '../schemas.js';
import { generateClaudeMd } from '../../templates/claude-md.js';
import { writeHooksConfig } from '../../hooks-config.js';
import { CASES_DIR, SETTINGS_PATH, validatePathWithinBase, parseBody, readJsonConfig } from '../route-helpers.js';
import { SseEvent } from '../sse-events.js';
import type { EventPort, ConfigPort } from '../ports/index.js';
import { dataPath, getDataDir } from '../../config/instance.js';
import {
checkRemoteTmuxAvailable,
readRemoteCases,
readRemoteHosts,
remoteDisplayPath,
writeRemoteCases,
writeRemoteHosts,
} from '../../remote-hosts.js';
const LINKED_CASES_FILE = dataPath('linked-cases.json');
const CODEMAN_CONFIG_DIR = getDataDir();
const SAFE_CASE_NAME = /^[a-zA-Z0-9_-]+$/;
/** Read and parse linked-cases.json, returning empty object on missing/invalid file. */
@@ -53,6 +68,7 @@ export function registerCaseRoutes(app: FastifyInstance, ctx: EventPort & Config
name: e.name,
path: join(CASES_DIR, e.name),
hasClaudeMd: existsSync(join(CASES_DIR, e.name, 'CLAUDE.md')),
location: 'local',
});
}
}
@@ -69,10 +85,39 @@ export function registerCaseRoutes(app: FastifyInstance, ctx: EventPort & Config
name,
path,
hasClaudeMd: existsSync(join(path, 'CLAUDE.md')),
linked: true,
location: 'linked-local',
});
}
}
// Get remote cases
const remoteHosts = await readRemoteHosts(CODEMAN_CONFIG_DIR);
const remoteHostMap = new Map(remoteHosts.map((host) => [host.id, host]));
for (const remoteCase of await readRemoteCases(CODEMAN_CONFIG_DIR)) {
const host = remoteHostMap.get(remoteCase.hostId);
if (!host || !SAFE_CASE_NAME.test(remoteCase.name)) continue;
existingNames.add(remoteCase.name);
const remoteCaseInfo: CaseInfo = {
name: remoteCase.name,
path: remoteDisplayPath({ username: host.username, host: host.host, path: remoteCase.remotePath }),
hasClaudeMd: false,
location: 'remote',
remote: {
hostId: host.id,
host: host.host,
username: host.username,
path: remoteCase.remotePath,
},
};
const existingIndex = cases.findIndex((item) => item.name === remoteCase.name);
if (existingIndex === -1) {
cases.push(remoteCaseInfo);
} else {
cases[existingIndex] = remoteCaseInfo;
}
}
// Sort by persisted caseOrder from settings.json
const settings = await readJsonConfig<Record<string, unknown>>(SETTINGS_PATH, 'settings', {});
const caseOrder = Array.isArray(settings.caseOrder) ? (settings.caseOrder as string[]) : [];
@@ -120,6 +165,73 @@ export function registerCaseRoutes(app: FastifyInstance, ctx: EventPort & Config
}
});
app.get('/api/remote-hosts', async () => readRemoteHosts(CODEMAN_CONFIG_DIR));
app.post('/api/remote-hosts', async (req): Promise<ApiResponse<{ host: unknown }>> => {
const host = parseBody(RemoteHostSchema, req.body);
const hosts = await readRemoteHosts(CODEMAN_CONFIG_DIR);
if (hosts.some((item) => item.id === host.id)) {
return createErrorResponse(ApiErrorCode.ALREADY_EXISTS, 'Remote host already exists');
}
await writeRemoteHosts(CODEMAN_CONFIG_DIR, [...hosts, host]);
return { success: true, data: { host } };
});
app.put('/api/remote-hosts/:id', async (req): Promise<ApiResponse<{ host: unknown }>> => {
const { id } = req.params as { id: string };
const host = parseBody(RemoteHostSchema, { ...(req.body as object), id });
const hosts = await readRemoteHosts(CODEMAN_CONFIG_DIR);
const index = hosts.findIndex((item) => item.id === id);
if (index === -1) return createErrorResponse(ApiErrorCode.NOT_FOUND, 'Remote host not found');
const next = [...hosts];
next[index] = host;
await writeRemoteHosts(CODEMAN_CONFIG_DIR, next);
return { success: true, data: { host } };
});
app.delete('/api/remote-hosts/:id', async (req): Promise<ApiResponse<{ id: string }>> => {
const { id } = req.params as { id: string };
const cases = await readRemoteCases(CODEMAN_CONFIG_DIR);
if (cases.some((item) => item.hostId === id)) {
return createErrorResponse(ApiErrorCode.OPERATION_FAILED, 'Remote host is still used by remote cases');
}
const hosts = await readRemoteHosts(CODEMAN_CONFIG_DIR);
await writeRemoteHosts(
CODEMAN_CONFIG_DIR,
hosts.filter((item) => item.id !== id)
);
return { success: true, data: { id } };
});
app.post('/api/cases/remote-link', async (req): Promise<ApiResponse<{ case: unknown }>> => {
const remoteCase = { ...parseBody(RemoteCaseLinkSchema, req.body), type: 'remote' as const };
const hosts = await readRemoteHosts(CODEMAN_CONFIG_DIR);
const host = hosts.find((item) => item.id === remoteCase.hostId);
if (!host) return createErrorResponse(ApiErrorCode.NOT_FOUND, 'Remote host not found');
const linkedCases = await readLinkedCases();
const remoteCases = await readRemoteCases(CODEMAN_CONFIG_DIR);
if (
remoteCases.some((item) => item.name === remoteCase.name) ||
linkedCases[remoteCase.name] ||
existsSync(join(CASES_DIR, remoteCase.name))
) {
return createErrorResponse(ApiErrorCode.ALREADY_EXISTS, 'Case already exists');
}
// Courtesy validation: tmux is a hard prerequisite for durable remote sessions.
// Verify it up-front so linking surfaces a clear error now instead of a dead pane
// at first launch (also confirms the SSH connection actually works).
const tmuxCheck = await checkRemoteTmuxAvailable(host);
if (!tmuxCheck.ok) {
return createErrorResponse(ApiErrorCode.OPERATION_FAILED, tmuxCheck.error || 'remote host is missing tmux');
}
await writeRemoteCases(CODEMAN_CONFIG_DIR, [...remoteCases, remoteCase]);
ctx.broadcast(SseEvent.CaseLinked, { name: remoteCase.name, path: remoteCase.remotePath, type: 'remote' });
return { success: true, data: { case: remoteCase } };
});
// Link an existing folder as a case
app.post('/api/cases/link', async (req): Promise<ApiResponse<{ case: { name: string; path: string } }>> => {
const { name, path: folderPath } = parseBody(LinkCaseSchema, req.body, 'Invalid request body');
@@ -173,6 +285,16 @@ export function registerCaseRoutes(app: FastifyInstance, ctx: EventPort & Config
return createErrorResponse(ApiErrorCode.INVALID_INPUT, 'Invalid case name');
}
const remoteCases = await readRemoteCases(CODEMAN_CONFIG_DIR);
if (remoteCases.some((item) => item.name === name)) {
await writeRemoteCases(
CODEMAN_CONFIG_DIR,
remoteCases.filter((item) => item.name !== name)
);
ctx.broadcast(SseEvent.CaseDeleted, { name, type: 'remote-unlinked' });
return { success: true, data: { name } };
}
// Check linked cases first — unlink only, don't delete the actual directory
const linkedCases = await readLinkedCases();
if (linkedCases[name]) {
@@ -233,6 +355,25 @@ export function registerCaseRoutes(app: FastifyInstance, ctx: EventPort & Config
return createErrorResponse(ApiErrorCode.INVALID_INPUT, 'Invalid case name');
}
const remoteCases = await readRemoteCases(CODEMAN_CONFIG_DIR);
const remoteCase = remoteCases.find((item) => item.name === name);
if (remoteCase) {
const host = (await readRemoteHosts(CODEMAN_CONFIG_DIR)).find((item) => item.id === remoteCase.hostId);
if (!host) return createErrorResponse(ApiErrorCode.NOT_FOUND, 'Remote host not found');
return {
name,
path: remoteDisplayPath({ username: host.username, host: host.host, path: remoteCase.remotePath }),
hasClaudeMd: false,
location: 'remote',
remote: {
hostId: host.id,
host: host.host,
username: host.username,
path: remoteCase.remotePath,
},
};
}
const casePath = await resolveCasePath(name);
if (!existsSync(casePath)) {
+80
View File
@@ -0,0 +1,80 @@
/**
* @fileoverview Cron Jobs routes.
*
* CRUD + enable/disable + Run Now + run history for `CronJob`s. These are
* separate from the legacy `/api/scheduled` (ScheduledRun) endpoints — see
* docs/cron-discovery.md §0.
*/
import { FastifyInstance } from 'fastify';
import { ApiErrorCode, createErrorResponse } from '../../types.js';
import { CronJobSchema, CronJobUpdateSchema, CronJobEnabledSchema } from '../schemas.js';
import { parseBody } from '../route-helpers.js';
import type { CronPort } from '../ports/index.js';
export function registerCronRoutes(app: FastifyInstance, ctx: CronPort): void {
// ── Jobs ────────────────────────────────────────────────────────────────
app.get('/api/cron/jobs', async () => {
return ctx.cron.listJobs();
});
app.post('/api/cron/jobs', async (req) => {
// No custom errorMessage: surface the schema's field-specific messages
// (e.g. "runAt is required for a one-time schedule").
const body = parseBody(CronJobSchema, req.body);
return { job: ctx.cron.createJob(body) };
});
app.get('/api/cron/jobs/:id', async (req) => {
const { id } = req.params as { id: string };
const job = ctx.cron.getJob(id);
if (!job) return createErrorResponse(ApiErrorCode.NOT_FOUND, 'Cron job not found');
return job;
});
app.put('/api/cron/jobs/:id', async (req) => {
const { id } = req.params as { id: string };
const body = parseBody(CronJobUpdateSchema, req.body);
const job = ctx.cron.updateJob(id, body);
if (!job) return createErrorResponse(ApiErrorCode.NOT_FOUND, 'Cron job not found');
return { job };
});
app.delete('/api/cron/jobs/:id', async (req) => {
const { id } = req.params as { id: string };
if (!ctx.cron.deleteJob(id)) {
return createErrorResponse(ApiErrorCode.NOT_FOUND, 'Cron job not found');
}
return {};
});
app.put('/api/cron/jobs/:id/enabled', async (req) => {
const { id } = req.params as { id: string };
const { enabled } = parseBody(CronJobEnabledSchema, req.body, 'Invalid request body');
const job = ctx.cron.setEnabled(id, enabled);
if (!job) return createErrorResponse(ApiErrorCode.NOT_FOUND, 'Cron job not found');
return { job };
});
// ── Run Now ──────────────────────────────────────────────────────────────
app.post('/api/cron/jobs/:id/run', async (req) => {
const { id } = req.params as { id: string };
const job = ctx.cron.getJob(id);
if (!job) return createErrorResponse(ApiErrorCode.NOT_FOUND, 'Cron job not found');
const run = await ctx.cron.runNow(id);
return { run, activeAgents: ctx.cron.countActiveAgents(job.agentType, job.id) };
});
// ── Run history ──────────────────────────────────────────────────────────
app.get('/api/cron/jobs/:id/runs', async (req) => {
const { id } = req.params as { id: string };
return ctx.cron.listRuns(id);
});
app.get('/api/cron/runs', async () => {
return ctx.cron.listRuns();
});
}
+1
View File
@@ -7,6 +7,7 @@ export { registerTeamRoutes } from './team-routes.js';
export { registerMuxRoutes } from './mux-routes.js';
export { registerFileRoutes } from './file-routes.js';
export { registerScheduledRoutes } from './scheduled-routes.js';
export { registerCronRoutes } from './cron-routes.js';
export { registerSystemRoutes } from './system-routes.js';
export { registerHookEventRoutes } from './hook-event-routes.js';
export { registerStatusTelemetryRoutes } from './status-telemetry-routes.js';
+4
View File
@@ -246,6 +246,10 @@ export function registerRespawnRoutes(
}
}
// Re-attach listener wiring if a prior PTY exit detached it (the wiring exit
// handler removes ALL session listeners; idempotent — no-op while still attached).
await ctx.setupSessionListeners(session);
// Start interactive session
await session.startInteractive();
getLifecycleLog().log({
+548 -76
View File
@@ -5,7 +5,7 @@
*/
import { FastifyInstance } from 'fastify';
import { join, dirname, extname } from 'node:path';
import { join, dirname, extname, basename } from 'node:path';
import { homedir } from 'node:os';
import { existsSync, statSync, mkdirSync, writeFileSync } from 'node:fs';
import { execFile } from 'node:child_process';
@@ -34,6 +34,7 @@ import {
FlickerFilterSchema,
QuickRunSchema,
QuickStartSchema,
InteractiveStartSchema,
} from '../schemas.js';
import {
autoConfigureRalph,
@@ -54,6 +55,7 @@ import {
} from '../../hooks-config.js';
import { generateClaudeMd } from '../../templates/claude-md.js';
import { imageWatcher } from '../../image-watcher.js';
import { convertHeicToJpeg } from '../heic-jpeg-converter.js';
import { getLifecycleLog } from '../../session-lifecycle-log.js';
import {
mergeUnifiedSessions,
@@ -70,10 +72,13 @@ import { RunSummaryTracker } from '../../run-summary.js';
import { MAX_INPUT_LENGTH, MAX_SESSION_NAME_LENGTH } from '../../config/terminal-limits.js';
import { MAX_PASTE_IMAGE_BYTES } from '../../config/buffer-limits.js';
import { dataPath } from '../../config/instance.js';
import { dataPath, getDataDir } from '../../config/instance.js';
import { checkRemoteTmuxAvailable, readRemoteCases, readRemoteHosts, toSessionRemote } from '../../remote-hosts.js';
import { LRUMap } from '../../utils/lru-map.js';
// Path to linked-cases registry (same file used by case-routes resolveCasePath)
const LINKED_CASES_FILE = dataPath('linked-cases.json');
const CODEMAN_CONFIG_DIR = getDataDir();
// Pre-compiled regex for terminal buffer cleaning (avoids per-request compilation)
// eslint-disable-next-line no-control-regex
@@ -198,6 +203,15 @@ export function imageMagicMatchesExt(data: Buffer, ext: string): boolean {
return u32be(0) === 0x52494646 && u32be(8) === 0x57454250;
case '.bmp':
return data[0] === 0x42 && data[1] === 0x4d;
case '.heic':
case '.heif': {
// ISO Base Media File Format: size + "ftyp" + major brand. The brand
// list matches heic-decode's own isHeic() — accepting more brands here
// would only route bytes into a conversion that always throws.
if (u32be(4) !== 0x66747970) return false;
const brand = data.subarray(8, 12).toString('ascii');
return ['heic', 'heix', 'hevc', 'hevx', 'mif1', 'msf1'].includes(brand);
}
default:
return false;
}
@@ -620,6 +634,14 @@ export function registerSessionRoutes(
app.post('/api/sessions/:id/interactive', async (req) => {
const { id } = req.params as { id: string };
// Body is optional (auto-reattach callers send none) — same idiom as /interactive-respawn.
const bodyResult = req.body
? InteractiveStartSchema.safeParse(req.body)
: { success: true as const, data: {} as { clearBreaker?: boolean } };
if (!bodyResult.success) {
return createErrorResponse(ApiErrorCode.INVALID_INPUT, 'Invalid request body');
}
const { clearBreaker } = bodyResult.data;
const session = findSessionOrFail(ctx, id);
if (session.isBusy()) {
@@ -642,6 +664,20 @@ export function registerSessionRoutes(
}
}
// COD-118: ONLY an explicit user-initiated restart (body {clearBreaker:true})
// clears a tripped PTY-exit circuit breaker. This endpoint is ALSO the frontend's
// automatic re-attach path (selectSession auto-POSTs it for any pid===null
// session), so an unconditional reset here would re-arm the exact crash loop
// the breaker exists to stop — auto-reattach sends no body and must not clear.
if (clearBreaker) {
session.resetRespawnBreaker();
}
// Re-attach listener wiring if a prior PTY exit detached it: the wiring exit
// handler removes ALL session listeners (incl. respawnBreakerTripped), and only
// session-create/boot-recovery paths ran setupSessionListeners before this fix —
// without this, a re-attached session's SSE/terminal/trip events go unobserved.
// setupSessionListeners is idempotent (no-op while refs are still attached).
await ctx.setupSessionListeners(session);
await session.startInteractive();
getLifecycleLog().log({
event: 'started',
@@ -669,6 +705,8 @@ export function registerSessionRoutes(
}
try {
// Re-attach listener wiring if a prior PTY exit detached it (see /interactive).
await ctx.setupSessionListeners(session);
await session.startShell();
getLifecycleLog().log({
event: 'started',
@@ -872,6 +910,14 @@ export function registerSessionRoutes(
const { id } = req.params as { id: string };
const session = findSessionOrFail(ctx, id);
// Codex sessions don't write to ~/.claude/projects — their transcripts
// live in ~/.codex/sessions/**. Branch to a Codex-specific reader so the
// response-viewer works for Codex panes too.
if (session.mode === 'codex') {
const codexQuery = req.query as { context?: string };
return await readCodexLastResponse(session, codexQuery.context === 'full');
}
// Scan ~/.claude/projects/*/ for the transcript file
const projectsDir = join(process.env.HOME || '/tmp', '.claude', 'projects');
@@ -985,15 +1031,346 @@ export function registerSessionRoutes(
};
});
function isCodexInjectedContext(text: string): boolean {
return (
/^# AGENTS\.md instructions\b/i.test(text) ||
/^<environment_context\b/i.test(text) ||
/^<turn_aborted\b/i.test(text) ||
/^<codex_internal_context\b/i.test(text) ||
/^<recommended_plugins\b/i.test(text) ||
/^<user_instructions\b/i.test(text) ||
/^# Options\b/i.test(text)
);
}
// ── Codex response-viewer support ───────────────────────────────────────────────────────
// Read the rollout's session_meta identity fields (plus turn_context cwd as
// a fallback when the huge session_meta line got truncated by the head read).
function readCodexRolloutMeta(head: string): { cwd?: string; originator?: string } {
let cwd: string | undefined;
let originator: string | undefined;
for (const line of head.split('\n')) {
if (!line) continue;
try {
const entry = JSON.parse(line) as {
type?: string;
payload?: { cwd?: string; originator?: string };
};
if (entry.type === 'session_meta') {
cwd ??= entry.payload?.cwd;
originator ??= entry.payload?.originator;
} else if (entry.type === 'turn_context') {
cwd ??= entry.payload?.cwd;
}
} catch {
// Malformed or truncated head line — keep scanning.
}
if (cwd && originator) break;
}
return { cwd, originator };
}
// The pane's last Enter (Session.codexLastSubmitAt) correlated against
// ~/.codex/history.jsonl, which logs every submitted user message as
// {session_id, ts}. This identifies the thread the pane is ACTUALLY on and
// is the only signal that survives /resume, /new and /fork typed inside the
// codex TUI itself. An entry is credited to this pane only when its Enter is
// the closest among all codex panes, so a menu keystroke in another pane
// can't steal the attribution.
const codexHistoryPinCache = new LRUMap<string, { submitAt: number; threadId: string }>({ maxSize: 1024 });
async function resolveCodexThreadFromHistory(
session: { id: string; codexLastSubmitAt?: number },
codexHome: string
): Promise<string | null> {
const submitAt = session.codexLastSubmitAt || 0;
if (!submitAt) return null;
const cached = codexHistoryPinCache.get(session.id);
if (cached && cached.submitAt === submitAt) return cached.threadId;
const histPath = join(codexHome, 'history.jsonl');
const st = await fs.stat(histPath).catch(() => null);
if (!st || st.size === 0) return null;
const tail = await readFileTail(histPath, Buffer.alloc(65536), st.size);
if (!tail) return null;
const WINDOW_MS = 15_000;
const otherSubmits: number[] = [];
for (const s of ctx.sessions.values()) {
if (s.id !== session.id && s.mode === 'codex' && s.codexLastSubmitAt) {
otherSubmits.push(s.codexLastSubmitAt);
}
}
let best: { threadId: string; dist: number } | undefined;
for (const line of tail.split('\n')) {
if (!line) continue;
let e: { session_id?: string; ts?: number };
try {
e = JSON.parse(line);
} catch {
continue; // first tail line may be cut mid-JSON
}
if (!e.session_id || typeof e.ts !== 'number') continue;
const tsMs = e.ts * 1000; // history timestamps are unix seconds
const dist = Math.abs(tsMs - submitAt);
if (dist > WINDOW_MS) continue;
if (otherSubmits.some((o) => Math.abs(tsMs - o) < dist)) continue; // another pane is closer
if (!best || dist < best.dist) best = { threadId: e.session_id, dist };
}
if (!best) return null;
codexHistoryPinCache.set(session.id, { submitAt, threadId: best.threadId });
return best.threadId;
}
// Locate THIS pane's rollout, in order of confidence:
// 0. history match — the thread the pane last submitted a message to
// (see resolveCodexThreadFromHistory); tracks the pane through
// /resume //new //fork typed inside the TUI.
// 1. originator match — Codeman spawns codex panes with
// CODEX_INTERNAL_ORIGINATOR_OVERRIDE=codeman_<sessionId>, which codex
// writes into session_meta.originator of every rollout it creates
// (including new files after /new in the same pane; newest match wins).
// 2. resume-id match — resumed rollouts keep their ORIGINAL session_meta
// (codex appends without rewriting it), so originator matching can't
// see them; but the rollout uuid is in the filename and we know the id.
// 3. legacy cwd+mtime heuristic — panes started before this feature, or
// TUI-resumed threads before their first tracked submit. Case-blind
// cwd compare (codex records the launch-time case, /mnt paths vary)
// and rollouts claimed by OTHER codeman panes are excluded.
async function findActiveCodexFile(session: {
id: string;
workingDir: string;
codexLastSubmitAt?: number;
codexConfig?: { resumeSessionId?: string };
}): Promise<string | null> {
const codexHome = process.env.CODEX_HOME || join(process.env.HOME || '/tmp', '.codex');
const sessionsDir = join(codexHome, 'sessions');
const files: Array<{ path: string; mtimeMs: number }> = [];
const walk = async (dir: string): Promise<void> => {
let entries: import('node:fs').Dirent[];
try {
entries = await fs.readdir(dir, { withFileTypes: true });
} catch {
return;
}
for (const entry of entries) {
const fullPath = join(dir, entry.name);
if (entry.isDirectory()) {
await walk(fullPath);
continue;
}
if (!entry.isFile() || !entry.name.endsWith('.jsonl')) continue;
const st = await fs.stat(fullPath).catch(() => null);
if (!st || st.size < 100) continue;
files.push({ path: fullPath, mtimeMs: st.mtimeMs });
}
};
await walk(sessionsDir);
files.sort((a, b) => b.mtimeMs - a.mtimeMs);
const historyThreadId = await resolveCodexThreadFromHistory(session, codexHome);
if (historyThreadId) {
const hit = files.find((f) => basename(f.path).endsWith(`-${historyThreadId}.jsonl`));
if (hit) return hit.path;
}
const rawResumeId = session.codexConfig?.resumeSessionId;
const resumeId =
rawResumeId && /^[a-f0-9]{8}-[a-f0-9]{4}-[a-f0-9]{4}-[a-f0-9]{4}-[a-f0-9]{12}$/.test(rawResumeId)
? rawResumeId
: undefined;
const idMatch = resumeId ? files.find((f) => basename(f.path).endsWith(`-${resumeId}.jsonl`)) : undefined;
// Scan newest-first for our originator; anything strictly older than the
// id match can never beat it, so the head reads stop there (mtime ties are
// still scanned — a /new rollout may land in the same clock tick). The
// 128 KiB head budget covers the session_meta line, which embeds full
// base_instructions (observed max ~22 KiB on codex 0.144).
const originator = `codeman_${session.id}`;
const wantCwd = session.workingDir.toLowerCase();
const headBuf = Buffer.alloc(131072);
let cwdFallback: { path: string; mtimeMs: number } | undefined;
for (const f of files) {
if (idMatch && f.mtimeMs < idMatch.mtimeMs) break;
const meta = await readCodexRolloutMetaCached(f.path, headBuf);
if (!meta) continue;
if (meta.originator === originator) return f.path; // newest-first → first hit wins
if (
!cwdFallback &&
!idMatch &&
meta.cwd?.toLowerCase() === wantCwd &&
// A rollout stamped by another codeman pane belongs to that pane.
!(meta.originator?.startsWith('codeman_') && meta.originator !== originator)
) {
cwdFallback = f;
}
}
return idMatch?.path ?? cwdFallback?.path ?? null;
}
// session_meta is written once when codex creates the rollout and never
// rewritten (verified: resume appends without touching it), so the parsed
// identity of a given path can be cached forever. This turns the per-request
// scan into stat calls plus head reads for new files only.
const codexRolloutMetaCache = new LRUMap<string, { cwd?: string; originator?: string }>({ maxSize: 4096 });
async function readCodexRolloutMetaCached(
filePath: string,
headBuf: Buffer
): Promise<{ cwd?: string; originator?: string } | null> {
const cached = codexRolloutMetaCache.get(filePath);
if (cached) return cached;
const head = await readFileHead(filePath, headBuf);
if (!head) return null;
const meta = readCodexRolloutMeta(head);
// Don't cache a still-incomplete head: a rollout being created may not
// have flushed session_meta/turn_context yet.
if (!meta.cwd && !meta.originator) return meta;
codexRolloutMetaCache.set(filePath, meta);
return meta;
}
function extractCodexBlockText(content: unknown, kinds: string[]): string {
if (typeof content === 'string') return content;
if (!Array.isArray(content)) return '';
return content
.filter(
(b): b is { type: string; text: string } =>
!!b &&
typeof b === 'object' &&
kinds.includes((b as { type?: string }).type || '') &&
typeof (b as { text?: string }).text === 'string'
)
.map((b) => b.text)
.join('\n\n');
}
// Single pass over a Codex rollout: track the last assistant message (for the
// default eye view) and, when `full`, the whole user/assistant thread.
//
// User turns come from event_msg/user_message when available: codex emits one
// per REAL user input, and injected context (AGENTS.md, environment_context,
// compaction summaries, …) never appears there — so no filtering heuristics.
// response_item user rows duplicate those inputs mixed with the injections;
// they are kept only as a fallback for old rollouts without event_msg rows.
async function readCodexLastResponse(
session: { id: string; workingDir: string; codexConfig?: { resumeSessionId?: string } },
full: boolean
): Promise<{
text: string;
timestamp: string;
messages?: Array<{ role: string; text: string; timestamp?: string }>;
}> {
const empty = full ? { text: '', timestamp: '', messages: [] } : { text: '', timestamp: '' };
const filePath = await findActiveCodexFile(session);
if (!filePath) return empty;
let content: string;
try {
content = await fs.readFile(filePath, 'utf8');
} catch {
return empty;
}
let lastText = '';
let lastTimestamp = '';
const messages: Array<{ role: string; text: string; timestamp?: string; legacyUser?: boolean }> = [];
// Multiset of event-sourced user texts: a real input appears BOTH as an
// event_msg and as a response_item row, so each event text cancels exactly
// one legacy twin. Legacy rows without an event twin (turns written by an
// older codex appending to the same rollout) survive — a file-wide boolean
// would wrongly drop them.
const eventUserTexts = new Map<string, number>();
for (const line of content.split('\n')) {
if (!line) continue;
let entry: {
timestamp?: string;
type?: string;
payload?: {
type?: string;
role?: string;
content?: unknown;
message?: unknown;
images?: unknown;
local_images?: unknown;
};
};
try {
entry = JSON.parse(line);
} catch {
continue;
}
if (full && entry.type === 'event_msg' && entry.payload?.type === 'user_message') {
let text = typeof entry.payload.message === 'string' ? entry.payload.message.trim() : '';
if (text && isCodexInjectedContext(text)) continue;
if (text) eventUserTexts.set(text, (eventUserTexts.get(text) || 0) + 1);
// Image-only (or image+text) inputs: the text field alone would make
// the turn vanish, so surface a placeholder.
const imageCount =
(Array.isArray(entry.payload.images) ? entry.payload.images.length : 0) +
(Array.isArray(entry.payload.local_images) ? entry.payload.local_images.length : 0);
if (imageCount > 0) text = text ? `${text}\n\n*[image ×${imageCount}]*` : `*[image ×${imageCount}]*`;
if (text) messages.push({ role: 'user', text, timestamp: entry.timestamp });
continue;
}
if (entry.type !== 'response_item' || entry.payload?.type !== 'message') continue;
const role = entry.payload?.role;
if (role === 'assistant') {
const text = extractCodexBlockText(entry.payload?.content, ['output_text', 'text']);
if (text) {
lastText = text;
lastTimestamp = entry.timestamp || '';
if (full) messages.push({ role: 'assistant', text, timestamp: entry.timestamp });
}
} else if (role === 'user' && full) {
const text = extractCodexBlockText(entry.payload?.content, ['input_text', 'text']).trim();
// Drop Codex's injected context turns (AGENTS.md, environment_context, …)
// so the thread shows real user prompts only.
if (text && !isCodexInjectedContext(text)) {
messages.push({ role: 'user', text, timestamp: entry.timestamp, legacyUser: true });
}
}
}
const thread = messages
.filter((m) => {
if (!m.legacyUser) return true;
const n = eventUserTexts.get(m.text) || 0;
if (n > 0) {
eventUserTexts.set(m.text, n - 1);
return false; // duplicate of an event_msg row already in the thread
}
return true;
})
.map(({ role, text, timestamp }) => ({ role, text, timestamp }));
return full
? { text: lastText, timestamp: lastTimestamp, messages: thread }
: { text: lastText, timestamp: lastTimestamp };
}
// ========== Get Terminal Buffer ==========
// Query params:
// tail=<bytes> - Only return last N bytes (faster initial load)
// full=1 - Full page reload: replay the entire tmux scrollback (COD-47)
app.get('/api/sessions/:id/terminal', async (req) => {
const { id } = req.params as { id: string };
const query = req.query as { tail?: string };
const query = req.query as { tail?: string; full?: string };
const session = findSessionOrFail(ctx, id);
// `full=1` is the EXPLICIT full-reload signal (COD-47): the browser reloaded
// the page and wants the whole scroll history back, so we capture the ENTIRE
// tmux scrollback and the user gets back history that scrolled off Codeman's
// byte buffer. Requests WITHOUT it — tab switches (`tail=`) and the legacy
// no-param callers (response-viewer fallback, clearTerminal refresh) — keep
// the fast visible-frame capture.
const tailBytes = query.tail ? parseInt(query.tail, 10) : 0;
const isFullReload = query.full === '1' || query.full === 'true';
const { tmuxHistoryLimit, terminalBufferMaxBytes } = await ctx.getTerminalHistoryConfig();
// Prepend the live tmux pane buffer so tab-switch replay shows the current
// on-screen frame, not just the accumulated byte history. This matters for
// TUI modes (codex/opencode) that repaint only their latest frame: the
@@ -1004,24 +1381,57 @@ export function registerSessionRoutes(
const muxName = session.muxName;
const liveMuxBuffer =
muxName && typeof ctx.mux.captureActivePaneBuffer === 'function'
? ctx.mux.captureActivePaneBuffer(muxName)
? ctx.mux.captureActivePaneBuffer(
muxName,
isFullReload
? { fullHistory: true, historyLimitLines: tmuxHistoryLimit, maxCaptureBytes: terminalBufferMaxBytes }
: undefined
)
: null;
const rawBuffer =
liveMuxBuffer !== null && liveMuxBuffer.length > 0
? session.terminalBufferLength > 0
const hasLiveMuxBuffer = liveMuxBuffer !== null && liveMuxBuffer.length > 0;
const source: 'history' | 'mux-visible' | 'mux-full-history' = hasLiveMuxBuffer
? isFullReload
? 'mux-full-history'
: 'mux-visible'
: 'history';
let rawBuffer: string;
if (liveMuxBuffer !== null && liveMuxBuffer.length > 0) {
// Full-history capture is the RENDERED form of everything already in the
// byte buffer (up to tmux eviction) — return it alone. Prepending the byte
// history would replay the whole conversation twice: `\x1b[2J` clears only
// the viewport, not xterm scrollback. The history+clear+frame concat stays
// for the visible-frame path, where the single pane frame lacks history.
rawBuffer = isFullReload
? liveMuxBuffer
: session.terminalBufferLength > 0
? `${session.terminalBuffer}\x1b[H\x1b[2J${liveMuxBuffer}`
: liveMuxBuffer
: session.terminalBuffer;
const tailBytes = query.tail ? parseInt(query.tail, 10) : 0;
: liveMuxBuffer;
} else {
rawBuffer = session.terminalBuffer;
}
const fullSize = rawBuffer.length;
let truncated = false;
let cleanBuffer: string;
// Cap the payload EARLY — before the regex normalization passes below run
// over it. A full-history tmux capture can be tens of MB of scrollback;
// normalizing all of it would stall the event loop only to discard most
// bytes anyway. Keep the most RECENT bytes (slice from the end) and align
// to a line boundary so we never start mid-ANSI-escape.
if (terminalBufferMaxBytes > 0 && rawBuffer.length > terminalBufferMaxBytes) {
rawBuffer = rawBuffer.slice(-terminalBufferMaxBytes);
truncated = true;
const capNewline = rawBuffer.indexOf('\n');
if (capNewline > 0 && capNewline < 4096) {
rawBuffer = rawBuffer.slice(capNewline + 1);
}
}
// Strip redundant Ink spinner/status redraws BEFORE tailing.
// During long thinking phases, Ink rewrites the same rows thousands of times
// (500KB+). Without stripping, tail mode returns only spinner frames and
// the terminal appears empty when switching tabs.
let strippedBuffer = stripInkRedrawBloat(rawBuffer);
let strippedBuffer = session.mode === 'shell' ? rawBuffer : stripInkRedrawBloat(rawBuffer);
// Strip alt-screen toggles and scrollback-erase from Codex/Claude byte
// streams. xterm.js obeys them by switching to its scrollback-less alt
@@ -1069,6 +1479,7 @@ export function registerSessionRoutes(
status: session.status,
fullSize,
truncated,
source,
};
});
@@ -1269,81 +1680,125 @@ export function registerSessionRoutes(
effort,
} = parseBody(QuickStartSchema, req.body);
// Check OpenCode availability if requested
if (mode === 'opencode') {
const { isOpenCodeAvailable } = await import('../../utils/opencode-cli-resolver.js');
if (!isOpenCodeAvailable()) {
// Resolve the remote case FIRST — the CLI executes on the REMOTE host over ssh,
// so the LOCAL availability gates below (isCodexAvailable() etc.) don't apply and
// would wrongly reject a machine that hasn't got the CLI installed locally.
let remote = undefined;
let casePath: string | null = null;
const remoteCases = await readRemoteCases(CODEMAN_CONFIG_DIR);
const remoteCase = remoteCases.find((item) => item.name === caseName);
if (remoteCase) {
const host = (await readRemoteHosts(CODEMAN_CONFIG_DIR)).find((item) => item.id === remoteCase.hostId);
if (!host) return createErrorResponse(ApiErrorCode.NOT_FOUND, 'Remote host not found');
// Per-session config that is applied to the LOCAL tmux/CLI wrapper (env vars via
// tmux setenv, effort/model CLI args, codex/gemini/opencode config) does NOT
// cross ssh, so it would silently no-op. Reject rather than pretend it worked —
// remote command/env customization goes through the per-host command override.
if (
(envOverrides && Object.keys(envOverrides).length > 0) ||
effort ||
codexConfig ||
geminiConfig ||
openCodeConfig
) {
return createErrorResponse(
ApiErrorCode.OPERATION_FAILED,
'OpenCode CLI not found. Install with: curl -fsSL https://opencode.ai/install | bash'
ApiErrorCode.INVALID_INPUT,
'envOverrides, effort, and per-CLI config are not supported for remote cases (they do not cross ssh). Configure the remote command via the host command override instead.'
);
}
}
// Check Codex availability if requested
if (mode === 'codex') {
const { isCodexAvailable } = await import('../../utils/codex-cli-resolver.js');
if (!isCodexAvailable()) {
return createErrorResponse(
ApiErrorCode.OPERATION_FAILED,
'Codex CLI not found. Install with: npm install -g @openai/codex'
);
// tmux is a hard prerequisite on the remote host (the agent runs inside a remote
// tmux server so it survives ssh drops). Probe before spawning so a missing tmux
// surfaces a clear, structured error instead of a dead "tmux: command not found" pane.
const tmuxCheck = await checkRemoteTmuxAvailable(host);
if (!tmuxCheck.ok) {
return createErrorResponse(ApiErrorCode.OPERATION_FAILED, tmuxCheck.error || 'remote host is missing tmux');
}
}
// Check Gemini availability if requested
if (mode === 'gemini') {
const { isGeminiAvailable } = await import('../../utils/gemini-cli-resolver.js');
if (!isGeminiAvailable()) {
return createErrorResponse(
ApiErrorCode.OPERATION_FAILED,
'Gemini CLI not found. Install with: npm install -g @google/gemini-cli'
);
casePath = remoteCase.remotePath;
remote = toSessionRemote(host, remoteCase);
} else {
// Check OpenCode availability if requested
if (mode === 'opencode') {
const { isOpenCodeAvailable } = await import('../../utils/opencode-cli-resolver.js');
if (!isOpenCodeAvailable()) {
return createErrorResponse(
ApiErrorCode.OPERATION_FAILED,
'OpenCode CLI not found. Install with: curl -fsSL https://opencode.ai/install | bash'
);
}
}
}
// Resolve case path: check linked-cases registry first, then fall back to CASES_DIR.
// This mirrors the behaviour of resolveCasePath() in case-routes so that linked
// external project directories are honoured by quick-start just like regular case routes.
let linkedCases: Record<string, string> = {};
try {
const raw = await fs.readFile(LINKED_CASES_FILE, 'utf-8');
linkedCases = JSON.parse(raw);
} catch {
// File missing or unparseable — treat as empty registry
}
const linkedCasePath = linkedCases[caseName];
const casePath = linkedCasePath || validatePathWithinBase(caseName, CASES_DIR);
if (!casePath) {
return createErrorResponse(ApiErrorCode.INVALID_INPUT, 'Invalid case path');
}
// Check Codex availability if requested
if (mode === 'codex') {
const { isCodexAvailable } = await import('../../utils/codex-cli-resolver.js');
if (!isCodexAvailable()) {
return createErrorResponse(
ApiErrorCode.OPERATION_FAILED,
'Codex CLI not found. Install with: npm install -g @openai/codex'
);
}
}
// Create case folder and CLAUDE.md if it doesn't exist (only for non-linked cases)
if (!existsSync(casePath)) {
// Check Gemini availability if requested
if (mode === 'gemini') {
const { isGeminiAvailable } = await import('../../utils/gemini-cli-resolver.js');
if (!isGeminiAvailable()) {
return createErrorResponse(
ApiErrorCode.OPERATION_FAILED,
'Gemini CLI not found. Install with: npm install -g @google/gemini-cli'
);
}
}
// Resolve case path: check linked-cases registry first, then fall back to CASES_DIR.
// This mirrors the behaviour of resolveCasePath() in case-routes so that linked
// external project directories are honoured by quick-start just like regular case routes.
let linkedCases: Record<string, string> = {};
try {
mkdirSync(casePath, { recursive: true });
mkdirSync(join(casePath, 'src'), { recursive: true });
const raw = await fs.readFile(LINKED_CASES_FILE, 'utf-8');
linkedCases = JSON.parse(raw);
} catch {
// File missing or unparseable — treat as empty registry
}
casePath = linkedCases[caseName] || validatePathWithinBase(caseName, CASES_DIR);
if (!casePath) {
return createErrorResponse(ApiErrorCode.INVALID_INPUT, 'Invalid case path');
}
}
// By this point casePath is guaranteed non-null: for remote cases it was set from remoteCase.remotePath,
// for local cases the !casePath guard above returned early. TypeScript can't narrow across the if/else.
const resolvedCasePath = casePath as string;
// Create case folder and CLAUDE.md if it doesn't exist (only for non-linked, non-remote cases)
if (!remote && !existsSync(resolvedCasePath)) {
try {
mkdirSync(resolvedCasePath, { recursive: true });
mkdirSync(join(resolvedCasePath, 'src'), { recursive: true });
// Read settings to get custom template path
const templatePath = await ctx.getDefaultClaudeMdPath();
const claudeMd = generateClaudeMd(caseName, '', templatePath);
writeFileSync(join(casePath, 'CLAUDE.md'), claudeMd);
writeFileSync(join(resolvedCasePath, 'CLAUDE.md'), claudeMd);
// Write .claude/settings.local.json with hooks for desktop notifications
// (Claude-specific — OpenCode, Codex, and Gemini use their own systems)
if (mode !== 'opencode' && mode !== 'codex' && mode !== 'gemini') {
await writeHooksConfig(casePath);
await writeHooksConfig(resolvedCasePath);
}
ctx.broadcast(SseEvent.CaseCreated, { name: caseName, path: casePath });
ctx.broadcast(SseEvent.CaseCreated, { name: caseName, path: resolvedCasePath });
} catch (err) {
return createErrorResponse(ApiErrorCode.OPERATION_FAILED, `Failed to create case: ${getErrorMessage(err)}`);
}
} else if (mode !== 'opencode') {
} else if (!remote && mode !== 'opencode') {
// COD-91 self-heal for an EXISTING case: refresh a pre-secret hooks block so the
// now-unconditional hook-secret gate keeps accepting its hook events. No-op when
// the hooks aren't ours or already carry the secret.
await refreshStaleHookSecret(casePath).catch(() => {});
// the hooks aren't ours or already carry the secret. Skipped for remote cases —
// resolvedCasePath is a REMOTE path that doesn't exist on the local filesystem.
await refreshStaleHookSecret(resolvedCasePath).catch(() => {});
}
// Strip stale disk entries for keys this request is actively setting (Claude only —
@@ -1352,10 +1807,11 @@ export function registerSessionRoutes(
mode !== 'opencode' &&
mode !== 'codex' &&
mode !== 'gemini' &&
!remote &&
envOverrides &&
Object.keys(envOverrides).length > 0
) {
await stripCaseEnvKeys(casePath, Object.keys(envOverrides));
await stripCaseEnvKeys(resolvedCasePath, Object.keys(envOverrides));
}
// Create a new session with the case as working directory
@@ -1375,7 +1831,7 @@ export function registerSessionRoutes(
const qsClaudeModeConfig = await ctx.getClaudeModeConfig();
const qsTerminalHistoryConfig = await ctx.getTerminalHistoryConfig();
const session = new Session({
workingDir: casePath,
workingDir: resolvedCasePath,
mux: ctx.mux,
useMux: true,
mode: mode,
@@ -1388,13 +1844,14 @@ export function registerSessionRoutes(
geminiConfig: mode === 'gemini' ? geminiConfig : undefined,
envOverrides,
effort,
remote,
tmuxHistoryLimit: qsTerminalHistoryConfig.tmuxHistoryLimit,
});
// Auto-detect completion phrase from CLAUDE.md BEFORE broadcasting
// so the initial state already has the phrase configured (only if globally enabled)
if (mode === 'claude' && ctx.store.getConfig().ralphEnabled) {
autoConfigureRalph(session, casePath, ctx);
if (mode === 'claude' && !remote && ctx.store.getConfig().ralphEnabled) {
autoConfigureRalph(session, resolvedCasePath, ctx);
if (!session.ralphTracker.enabled) {
session.ralphTracker.enable();
session.ralphTracker.enableAutoEnable(); // Allow re-enabling on restart
@@ -1463,7 +1920,7 @@ export function registerSessionRoutes(
return {
sessionId: session.id,
casePath,
casePath: resolvedCasePath,
caseName,
};
} catch (err) {
@@ -1875,7 +2332,7 @@ export function registerSessionRoutes(
// Paste Image (clipboard / drag-drop upload)
// ═══════════════════════════════════════════════════════════════
const ALLOWED_IMAGE_EXTS = new Set(['.png', '.jpg', '.jpeg', '.gif', '.webp', '.bmp']);
const ALLOWED_IMAGE_EXTS = new Set(['.png', '.jpg', '.jpeg', '.gif', '.webp', '.bmp', '.heic', '.heif']);
// The per-file size cap (MAX_PASTE_IMAGE_BYTES) is enforced by @fastify/multipart (registered in server.ts).
app.post('/api/sessions/:id/paste-image', async (req, reply) => {
@@ -1967,7 +2424,7 @@ export function registerSessionRoutes(
const origExt = extname(part.filename).toLowerCase();
if (ALLOWED_IMAGE_EXTS.has(origExt)) ext = origExt;
}
const mimeMatch = (part.mimetype || '').toLowerCase().match(/^image\/(png|jpeg|jpg|webp|gif|bmp)$/);
const mimeMatch = (part.mimetype || '').toLowerCase().match(/^image\/(png|jpeg|jpg|webp|gif|bmp|heic|heif)$/);
if (mimeMatch) {
const map: Record<string, string> = {
png: '.png',
@@ -1976,6 +2433,8 @@ export function registerSessionRoutes(
webp: '.webp',
gif: '.gif',
bmp: '.bmp',
heic: '.heic',
heif: '.heif',
};
ext = map[mimeMatch[1]] ?? ext;
}
@@ -1988,14 +2447,27 @@ export function registerSessionRoutes(
);
}
// Sniff actual bytes — filename and Content-Type are both attacker-supplied.
// Polyglot HTML/PNG would otherwise pass and serve back with image/png MIME.
if (!imageMagicMatchesExt(imageBytes, ext)) {
// Diagnostic: on some Android galleries (e.g. MIUI) a WebP/HEIF is
// mislabeled as image/jpeg, so the declared ext passes the allowlist but
// the magic bytes do not. Log the real header so format mismatches can be
// pinned down without a reproduce-and-guess loop. The client now
// re-encodes images to JPEG/PNG before upload, so this should be rare.
// Route HEIC on the raw bytes, NOT the declared ext/mime: on some Android
// galleries (e.g. MIUI) a HEIF comes back mislabeled as image/jpeg, and
// browsers that cannot decode HEIF upload the original file as-is — so a
// HEIC payload can arrive under any declared type. Filename and
// Content-Type are attacker-supplied anyway; only the bytes are trusted.
if (imageMagicMatchesExt(imageBytes, '.heic')) {
try {
imageBytes = await convertHeicToJpeg(imageBytes);
ext = '.jpg';
} catch (err: unknown) {
console.warn(
`[paste-image] HEIC conversion failed: filename=${JSON.stringify(part.filename)} mime=${JSON.stringify(part.mimetype)} error=${getErrorMessage(err)}`
);
reply.code(415);
return createErrorResponse(ApiErrorCode.INVALID_INPUT, 'Could not convert HEIC image to JPEG');
}
} else if (!imageMagicMatchesExt(imageBytes, ext)) {
// Sniff actual bytes — a polyglot HTML/PNG would otherwise pass and
// serve back with image/png MIME. Log the real header so format
// mismatches can be pinned down without a reproduce-and-guess loop. The
// client re-encodes images to JPEG/PNG before upload, so this is rare.
console.warn(
`[paste-image] magic mismatch: filename=${JSON.stringify(part.filename)} mime=${JSON.stringify(part.mimetype)} declaredExt=${ext} magic=${imageBytes.subarray(0, 12).toString('hex')}`
);
+210 -166
View File
@@ -34,6 +34,7 @@ import type { WebSocket } from 'ws';
import type { SessionPort } from '../ports/session-port.js';
import { MAX_INPUT_LENGTH } from '../../config/terminal-limits.js';
import { isAllowedRequestHost, isAllowedRequestOrigin, type HostPolicy } from '../network-auth-policy.js';
import { WsConnectionRegistry } from '../ws-connection-registry.js';
/** Micro-batch interval for terminal output (ms). Short enough for low latency,
* long enough to group Ink's rapid cursor-up redraw sequences into single frames. */
@@ -59,197 +60,240 @@ const DEC_2026_END = '\x1b[?2026l';
/** Max concurrent WS connections per session. Prevents listener/bandwidth multiplication. */
const MAX_WS_PER_SESSION = 5;
/** Track active WS connections per session for connection limiting. */
const sessionWsCount = new Map<string, number>();
/**
* Track live WS connections per session, keyed by clientId (COD-137).
* Replaces a bare counter that over-counted across the async-close gap on
* reconnect (spurious 4008). A same-`cid` reconnect supersedes its own socket
* (reclaims the slot) instead of consuming a new one; cid-less upgrades are
* admitted anonymously up to the cap. See ws-connection-registry.ts.
*/
const sessionWsRegistry = new WsConnectionRegistry<WebSocket>(MAX_WS_PER_SESSION);
export function registerWsRoutes(app: FastifyInstance, ctx: SessionPort, getHostPolicy: () => HostPolicy): void {
app.get<{ Params: { id: string } }>('/ws/sessions/:id/terminal', { websocket: true }, (socket: WebSocket, req) => {
// Reject cross-site WebSocket hijacking (CSWSH) and DNS-rebinding before doing
// anything: the upgrade must come from an allowed Host and (when the browser
// sends one — it always does for WS) a same-site Origin. Writing to this socket
// injects keystrokes into a --dangerously-skip-permissions agent, so this gate
// matters even on the default no-password install. See security review H5.
const policy = getHostPolicy();
if (!isAllowedRequestHost(req.headers.host, policy) || !isAllowedRequestOrigin(req.headers.origin, policy)) {
socket.close(4003, 'Forbidden');
return;
}
app.get<{ Params: { id: string }; Querystring: { cid?: string } }>(
'/ws/sessions/:id/terminal',
{ websocket: true },
(socket: WebSocket, req) => {
// Reject cross-site WebSocket hijacking (CSWSH) and DNS-rebinding before doing
// anything: the upgrade must come from an allowed Host and (when the browser
// sends one — it always does for WS) a same-site Origin. Writing to this socket
// injects keystrokes into a --dangerously-skip-permissions agent, so this gate
// matters even on the default no-password install. See security review H5.
const policy = getHostPolicy();
if (!isAllowedRequestHost(req.headers.host, policy) || !isAllowedRequestOrigin(req.headers.origin, policy)) {
socket.close(4003, 'Forbidden');
return;
}
const { id } = req.params;
const session = ctx.sessions.get(id);
const { id } = req.params;
const session = ctx.sessions.get(id);
if (!session) {
socket.close(4004, 'Session not found');
return;
}
if (!session) {
socket.close(4004, 'Session not found');
return;
}
// Enforce per-session connection limit
const currentCount = sessionWsCount.get(id) ?? 0;
if (currentCount >= MAX_WS_PER_SESSION) {
socket.close(4008, 'Too many connections');
return;
}
sessionWsCount.set(id, currentCount + 1);
// Structured transport logging — surfaces WS open/close/timeout churn so the
// tunnel-flap behavior (COD-134) is observable in the server logs. Fastify is
// configured logger:false, so we log via console (→ journald under systemd).
// Swallow socket errors — cleanup happens in 'close'
socket.on('error', () => {});
// Enforce per-session connection limit, scoped by clientId. A same-cid
// reconnect reclaims its own slot (registry evicts the stale socket), so a
// drop+reconnect burst can no longer over-count across the async-close gap
// and trip a spurious 4008. cid-less upgrades are admitted anonymously.
const cid = typeof req.query?.cid === 'string' && req.query.cid.length > 0 ? req.query.cid : null;
const { admitted, evictedSocket } = sessionWsRegistry.register(id, cid, socket);
if (!admitted) {
console.warn('[ws] terminal rejected: too many connections', {
sessionId: id,
wsCount: sessionWsRegistry.liveCount(id),
});
socket.close(4008, 'Too many connections');
return;
}
if (evictedSocket) {
// Same client reconnected; retire the stale socket so it doesn't linger.
console.info('[ws] terminal superseded by reconnect', { sessionId: id });
try {
evictedSocket.close(4010, 'Superseded by reconnect');
} catch {
/* socket may already be closing */
}
}
console.info('[ws] terminal open', { sessionId: id, wsCount: sessionWsRegistry.liveCount(id) });
// Per-connection micro-batch state
let batchChunks: string[] = [];
let batchSize = 0;
let batchTimer: ReturnType<typeof setTimeout> | null = null;
// Eagerly free the slot on error/terminate — don't wait for the async
// 'close' (idempotent with the 'close' handler below). This is what kills
// the reconnect over-count: the slot is released the instant the socket dies.
socket.on('error', () => {
sessionWsRegistry.unregister(id, socket);
});
const flushBatch = () => {
batchTimer = null;
if (batchChunks.length === 0 || socket.readyState !== 1) {
// Per-connection micro-batch state
let batchChunks: string[] = [];
let batchSize = 0;
let batchTimer: ReturnType<typeof setTimeout> | null = null;
const flushBatch = () => {
batchTimer = null;
if (batchChunks.length === 0 || socket.readyState !== 1) {
batchChunks = [];
batchSize = 0;
return;
}
const data = batchChunks.join('');
batchChunks = [];
batchSize = 0;
return;
}
const data = batchChunks.join('');
batchChunks = [];
batchSize = 0;
socket.send(`{"t":"o","d":${JSON.stringify(DEC_2026_START + data + DEC_2026_END)}}`);
};
socket.send(`{"t":"o","d":${JSON.stringify(DEC_2026_START + data + DEC_2026_END)}}`);
};
// Per-connection desktop sizing claim — registered on the first
// desktop-typed resize and released on socket close, so Session.resize()
// can ignore small-viewport resizes only while a desktop is actually
// connected (see Session._desktopSizeClaims).
const sizingToken = Symbol('ws-desktop-sizing');
let holdsDesktopClaim = false;
// Per-connection desktop sizing claim — registered on the first
// desktop-typed resize and released on socket close, so Session.resize()
// can ignore small-viewport resizes only while a desktop is actually
// connected (see Session._desktopSizeClaims).
const sizingToken = Symbol('ws-desktop-sizing');
let holdsDesktopClaim = false;
// Attach message handler synchronously BEFORE any async work
// (@fastify/websocket requirement to avoid dropped messages).
socket.on('message', (raw) => {
try {
const msg = JSON.parse(String(raw));
if (msg.t === 'i' && typeof msg.d === 'string') {
if (msg.d.length > MAX_INPUT_LENGTH) return;
// Reliable delivery: when the frame carries a clientId + seq, apply it
// exactly once (skip a duplicate redelivery) but ACK it regardless so
// the client can drop it from its durable queue. Frames without seq
// (legacy/other tools) are applied as-is — no behavior change.
const cid = typeof msg.cid === 'string' ? msg.cid : null;
const seq = Number.isInteger(msg.seq) ? (msg.seq as number) : null;
const apply = cid && seq !== null ? session.shouldApplyInput(cid, seq) : true;
if (apply) {
// Typed input from a claim-holding desktop keeps the claim "hot"
// and re-asserts the desktop layout after a mobile override.
if (holdsDesktopClaim) session.noteDesktopActivity();
session.write(msg.d);
// Attach message handler synchronously BEFORE any async work
// (@fastify/websocket requirement to avoid dropped messages).
socket.on('message', (raw) => {
try {
const msg = JSON.parse(String(raw));
if (msg.t === 'i' && typeof msg.d === 'string') {
if (msg.d.length > MAX_INPUT_LENGTH) return;
// Reliable delivery: when the frame carries a clientId + seq, apply it
// exactly once (skip a duplicate redelivery) but ACK it regardless so
// the client can drop it from its durable queue. Frames without seq
// (legacy/other tools) are applied as-is — no behavior change.
const cid = typeof msg.cid === 'string' ? msg.cid : null;
const seq = Number.isInteger(msg.seq) ? (msg.seq as number) : null;
const apply = cid && seq !== null ? session.shouldApplyInput(cid, seq) : true;
if (apply) {
// Typed input from a claim-holding desktop keeps the claim "hot"
// and re-asserts the desktop layout after a mobile override.
if (holdsDesktopClaim) session.noteDesktopActivity();
session.write(msg.d);
}
if (seq !== null && socket.readyState === 1) {
socket.send(`{"t":"ia","seq":${seq}}`);
}
} else if (
msg.t === 'z' &&
Number.isInteger(msg.c) &&
Number.isInteger(msg.r) &&
msg.c >= 1 &&
msg.c <= 500 &&
msg.r >= 1 &&
msg.r <= 200
) {
const viewportType = msg.v === 'mobile' || msg.v === 'tablet' || msg.v === 'desktop' ? msg.v : undefined;
if (viewportType === 'desktop') {
session.claimDesktopSizing(sizingToken);
holdsDesktopClaim = true;
} else if (viewportType) {
// The connection's viewport can change (e.g. browser window
// narrowed past the tablet breakpoint) — drop a stale claim.
session.releaseDesktopSizing(sizingToken);
holdsDesktopClaim = false;
}
const force = msg.f === true;
session.resize(msg.c, msg.r, { viewportType, force });
}
if (seq !== null && socket.readyState === 1) {
socket.send(`{"t":"ia","seq":${seq}}`);
}
} else if (
msg.t === 'z' &&
Number.isInteger(msg.c) &&
Number.isInteger(msg.r) &&
msg.c >= 1 &&
msg.c <= 500 &&
msg.r >= 1 &&
msg.r <= 200
) {
const viewportType = msg.v === 'mobile' || msg.v === 'tablet' || msg.v === 'desktop' ? msg.v : undefined;
if (viewportType === 'desktop') {
session.claimDesktopSizing(sizingToken);
holdsDesktopClaim = true;
} else if (viewportType) {
// The connection's viewport can change (e.g. browser window
// narrowed past the tablet breakpoint) — drop a stale claim.
session.releaseDesktopSizing(sizingToken);
holdsDesktopClaim = false;
}
const force = msg.f === true;
session.resize(msg.c, msg.r, { viewportType, force });
} catch {
// Ignore malformed messages
}
} catch {
// Ignore malformed messages
}
});
});
// Terminal output -> micro-batched WS send
const onTerminal = (data: string) => {
if (socket.readyState !== 1) return;
batchChunks.push(data);
batchSize += data.length;
// Terminal output -> micro-batched WS send
const onTerminal = (data: string) => {
if (socket.readyState !== 1) return;
batchChunks.push(data);
batchSize += data.length;
// Flush immediately for large batches (responsiveness during bulk output)
if (batchSize > WS_BATCH_FLUSH_THRESHOLD) {
if (batchTimer) {
clearTimeout(batchTimer);
// Flush immediately for large batches (responsiveness during bulk output)
if (batchSize > WS_BATCH_FLUSH_THRESHOLD) {
if (batchTimer) {
clearTimeout(batchTimer);
}
flushBatch();
return;
}
flushBatch();
return;
}
// Start timer if not already running
if (!batchTimer) {
batchTimer = setTimeout(flushBatch, WS_BATCH_INTERVAL_MS);
}
};
// Start timer if not already running
if (!batchTimer) {
batchTimer = setTimeout(flushBatch, WS_BATCH_INTERVAL_MS);
}
};
const onClearTerminal = () => {
if (socket.readyState === 1) {
socket.send('{"t":"c"}');
}
};
const onClearTerminal = () => {
if (socket.readyState === 1) {
socket.send('{"t":"c"}');
}
};
const onNeedsRefresh = () => {
if (socket.readyState === 1) {
socket.send('{"t":"r"}');
}
};
const onNeedsRefresh = () => {
if (socket.readyState === 1) {
socket.send('{"t":"r"}');
}
};
// Close WS when session exits (deleted, respawned, or crashed) — prevents
// orphaned listeners and stale writes to a dead PTY.
const onSessionExit = () => {
socket.close(4009, 'Session terminated');
};
// Close WS when session exits (deleted, respawned, or crashed) — prevents
// orphaned listeners and stale writes to a dead PTY.
const onSessionExit = () => {
socket.close(4009, 'Session terminated');
};
session.on('terminal', onTerminal);
session.on('clearTerminal', onClearTerminal);
session.on('needsRefresh', onNeedsRefresh);
session.on('exit', onSessionExit);
session.on('terminal', onTerminal);
session.on('clearTerminal', onClearTerminal);
session.on('needsRefresh', onNeedsRefresh);
session.on('exit', onSessionExit);
// Heartbeat: detect stale connections (especially through tunnels where
// TCP RST can take minutes to propagate).
let pongTimeout: ReturnType<typeof setTimeout> | null = null;
// Heartbeat: detect stale connections (especially through tunnels where
// TCP RST can take minutes to propagate).
let pongTimeout: ReturnType<typeof setTimeout> | null = null;
socket.on('pong', () => {
if (pongTimeout) {
clearTimeout(pongTimeout);
pongTimeout = null;
}
});
socket.on('pong', () => {
if (pongTimeout) {
clearTimeout(pongTimeout);
pongTimeout = null;
}
});
const pingInterval = setInterval(() => {
if (socket.readyState !== 1) return;
socket.ping();
pongTimeout = setTimeout(() => {
socket.terminate();
}, WS_PONG_TIMEOUT_MS);
}, WS_PING_INTERVAL_MS);
const pingInterval = setInterval(() => {
if (socket.readyState !== 1) return;
socket.ping();
pongTimeout = setTimeout(() => {
console.warn('[ws] terminal ping timeout — terminating', { sessionId: id });
// Free the slot eagerly — terminate()'s 'close' may lag, and a client
// reconnecting after a stale-connection drop must not be over-counted.
sessionWsRegistry.unregister(id, socket);
socket.terminate();
}, WS_PONG_TIMEOUT_MS);
}, WS_PING_INTERVAL_MS);
socket.on('close', () => {
clearInterval(pingInterval);
if (pongTimeout) clearTimeout(pongTimeout);
if (batchTimer) clearTimeout(batchTimer);
batchChunks = [];
session.off('terminal', onTerminal);
session.off('clearTerminal', onClearTerminal);
session.off('needsRefresh', onNeedsRefresh);
session.off('exit', onSessionExit);
session.releaseDesktopSizing(sizingToken);
socket.on('close', (code: number, reason: Buffer) => {
clearInterval(pingInterval);
if (pongTimeout) clearTimeout(pongTimeout);
if (batchTimer) clearTimeout(batchTimer);
batchChunks = [];
session.off('terminal', onTerminal);
session.off('clearTerminal', onClearTerminal);
session.off('needsRefresh', onNeedsRefresh);
session.off('exit', onSessionExit);
session.releaseDesktopSizing(sizingToken);
// Decrement per-session connection count
const count = sessionWsCount.get(id) ?? 1;
if (count <= 1) {
sessionWsCount.delete(id);
} else {
sessionWsCount.set(id, count - 1);
}
});
});
// Release this socket's slot. Idempotent and identity-matched: if this
// socket was already superseded (a same-cid reconnect took its slot) or
// eagerly unregistered on terminate/error, this is a no-op and the
// reconnected socket keeps the slot.
sessionWsRegistry.unregister(id, socket);
console.info('[ws] terminal close', {
sessionId: id,
code,
reason: String(reason),
wsCount: sessionWsRegistry.liveCount(id),
});
});
}
);
}