P2: découverte & reprise des sessions Claude
Arboretum découvre désormais toutes les sessions Claude de la machine (scan ~/.claude/projects + registre ~/.claude/sessions), distingue vivantes/mortes par pid+procStart, et permet de reprendre une morte (--resume dans son cwd d'origine) ou forker une vivante sans la corrompre. - shared: SessionSummary enrichi (source, claudeSessionId, pid, resumable, attachable, registryStatus) — additif, PROTOCOL_VERSION inchangé ; types REST resume/fork. - db: migration id:2 (claude_session_id, resumed_from). - core: jsonl-discovery (parseur tolérant, scan asynchrone non bloquant), session-registry (vivacité pid+procStart), discovery-service (cache + refresh périodique + diff/broadcast), pty-manager (resume/fork + capture du claudeSessionId via le registre). - routes: /sessions/:id/resume (garde-fou 409 anti-corruption sur session vivante) et /fork ; GET fusionné managées + découvertes ; relais WS. - web: badges managed/discovered + busy/idle/waiting, actions conditionnelles (Open/Observe/Kill vs Fork/View vs Resume/Fork), vue read-only des sessions externes, i18n EN/FR. - tests: jsonl-discovery, session-registry, discovery-service + resume/fork (130 verts) ; acceptation E2E acceptance-p2.mjs (sans quota) ALL GREEN. Conforme aux verdicts S1 (resume dans cwd d'origine, vivacité pid+procStart) et S4 (munge cwd, parseur tête+queue, priorité de titre).
This commit is contained in:
@@ -9,6 +9,7 @@ import type { Config } from './config.js';
|
||||
import type { Db } from './db/index.js';
|
||||
import { AuthService, LoginRateLimiter, type AuthContext } from './auth/service.js';
|
||||
import { PtyManager } from './core/pty-manager.js';
|
||||
import { DiscoveryService } from './core/discovery-service.js';
|
||||
import { registerAuthRoutes } from './routes/auth.js';
|
||||
import { registerSessionRoutes } from './routes/sessions.js';
|
||||
import { registerWsGateway } from './ws/gateway.js';
|
||||
@@ -26,13 +27,19 @@ export interface AppBundle {
|
||||
app: FastifyInstance;
|
||||
auth: AuthService;
|
||||
manager: PtyManager;
|
||||
discovery: DiscoveryService;
|
||||
}
|
||||
|
||||
export function buildApp(config: Config, db: Db, serverVersion: string): AppBundle {
|
||||
const app = Fastify({ logger: { level: process.env.ARBORETUM_LOG ?? 'info' } });
|
||||
const auth = new AuthService(db);
|
||||
const limiter = new LoginRateLimiter();
|
||||
const manager = new PtyManager(db);
|
||||
const manager = new PtyManager(db, config.claudeSessionsDir);
|
||||
const discovery = new DiscoveryService({
|
||||
ptyManager: manager,
|
||||
projectsDir: config.claudeProjectsDir,
|
||||
sessionsDir: config.claudeSessionsDir,
|
||||
});
|
||||
|
||||
void app.register(fastifyCookie);
|
||||
void app.register(fastifyWebsocket, {
|
||||
@@ -72,11 +79,11 @@ export function buildApp(config: Config, db: Db, serverVersion: string): AppBund
|
||||
});
|
||||
|
||||
registerAuthRoutes(app, auth, limiter, serverVersion);
|
||||
registerSessionRoutes(app, manager);
|
||||
registerSessionRoutes(app, manager, discovery);
|
||||
// La route websocket doit être déclarée APRÈS le chargement du plugin (contexte
|
||||
// encapsulé) — sinon le handler reçoit la signature REST (request, reply).
|
||||
void app.register(async (scoped) => {
|
||||
registerWsGateway(scoped, manager, serverVersion);
|
||||
registerWsGateway(scoped, manager, discovery, serverVersion);
|
||||
});
|
||||
|
||||
// SPA buildée embarquée dans le paquet npm (public/) — absente en dev (vite dev sert le front)
|
||||
@@ -91,5 +98,5 @@ export function buildApp(config: Config, db: Db, serverVersion: string): AppBund
|
||||
});
|
||||
}
|
||||
|
||||
return { app, auth, manager };
|
||||
return { app, auth, manager, discovery };
|
||||
}
|
||||
|
||||
@@ -11,6 +11,10 @@ export interface Config {
|
||||
/** origins supplémentaires autorisées (ex. https://machine.tailnet.ts.net) */
|
||||
allowedOrigins: string[];
|
||||
printToken: boolean;
|
||||
/** ~/.claude/projects (transcripts JSONL) — surchargeable via --claude-home (tests). */
|
||||
claudeProjectsDir: string;
|
||||
/** ~/.claude/sessions (registre des sessions CLI vivantes). */
|
||||
claudeSessionsDir: string;
|
||||
}
|
||||
|
||||
export function loadConfig(argv = process.argv.slice(2)): Config {
|
||||
@@ -23,6 +27,8 @@ export function loadConfig(argv = process.argv.slice(2)): Config {
|
||||
'allow-origin': { type: 'string', multiple: true },
|
||||
'print-token': { type: 'boolean', default: false },
|
||||
'i-know-this-exposes-a-terminal': { type: 'boolean', default: false },
|
||||
// racine de l'install Claude (~/.claude par défaut) — surchargée par les tests d'acceptation.
|
||||
'claude-home': { type: 'string' },
|
||||
},
|
||||
strict: true,
|
||||
});
|
||||
@@ -39,6 +45,7 @@ export function loadConfig(argv = process.argv.slice(2)): Config {
|
||||
|
||||
const dataDir = join(process.env.XDG_DATA_HOME ?? join(homedir(), '.local', 'share'), 'arboretum');
|
||||
mkdirSync(dataDir, { recursive: true });
|
||||
const claudeHome = values['claude-home'] ?? join(homedir(), '.claude');
|
||||
return {
|
||||
port: Number(values.port),
|
||||
bind,
|
||||
@@ -46,5 +53,7 @@ export function loadConfig(argv = process.argv.slice(2)): Config {
|
||||
dataDir,
|
||||
allowedOrigins: values['allow-origin'] ?? [],
|
||||
printToken: values['print-token'] ?? false,
|
||||
claudeProjectsDir: join(claudeHome, 'projects'),
|
||||
claudeSessionsDir: join(claudeHome, 'sessions'),
|
||||
};
|
||||
}
|
||||
|
||||
@@ -6,6 +6,12 @@ export interface SpawnSpec {
|
||||
env: NodeJS.ProcessEnv;
|
||||
}
|
||||
|
||||
export interface SpawnOptions {
|
||||
command: 'claude' | 'bash';
|
||||
/** reprise d'une session existante (P2) : `--resume <id>`, `--fork-session` si fork. */
|
||||
resume?: { claudeSessionId: string; fork?: boolean };
|
||||
}
|
||||
|
||||
let cachedClaudeBin: string | null = null;
|
||||
|
||||
export function resolveClaudeBin(): string {
|
||||
@@ -21,14 +27,20 @@ export function resolveClaudeBin(): string {
|
||||
}
|
||||
|
||||
/** Module volontairement abstrait : le plan B « BYO API key / Agent SDK » se brancherait ici. */
|
||||
export function buildSpawnSpec(command: 'claude' | 'bash'): SpawnSpec {
|
||||
export function buildSpawnSpec(opts: SpawnOptions): SpawnSpec {
|
||||
const env: NodeJS.ProcessEnv = {
|
||||
...process.env,
|
||||
TERM: 'xterm-256color',
|
||||
COLORTERM: 'truecolor',
|
||||
};
|
||||
if (command === 'bash') {
|
||||
if (opts.command === 'bash') {
|
||||
return { file: 'bash', args: ['--norc'], env };
|
||||
}
|
||||
return { file: resolveClaudeBin(), args: [], env };
|
||||
const args: string[] = [];
|
||||
if (opts.resume) {
|
||||
// `--resume` doit toujours s'exécuter dans le cwd d'origine (garanti par l'appelant — spike S1).
|
||||
args.push('--resume', opts.resume.claudeSessionId);
|
||||
if (opts.resume.fork) args.push('--fork-session');
|
||||
}
|
||||
return { file: resolveClaudeBin(), args, env };
|
||||
}
|
||||
|
||||
137
packages/server/src/core/discovery-service.ts
Normal file
137
packages/server/src/core/discovery-service.ts
Normal file
@@ -0,0 +1,137 @@
|
||||
// Service de découverte : second producteur de SessionSummary (à côté du PtyManager).
|
||||
// Scanne ~/.claude/projects (JSONL) + ~/.claude/sessions (registre), corrèle par claudeSessionId,
|
||||
// calcule la vivacité, met en cache et notifie les changements. Lecture seule du disque.
|
||||
import { EventEmitter } from 'node:events';
|
||||
import { homedir } from 'node:os';
|
||||
import { join } from 'node:path';
|
||||
import type { SessionSummary } from '@arboretum/shared';
|
||||
import { scanProjects, type DiscoveredJsonl } from './jsonl-discovery.js';
|
||||
import { readRegistry, type RegistryEntry } from './session-registry.js';
|
||||
import type { PtyManager } from './pty-manager.js';
|
||||
|
||||
const DEFAULT_REFRESH_MS = 10_000;
|
||||
|
||||
export interface DiscoveryServiceEvents {
|
||||
/** session découverte nouvelle ou modifiée (relayée en session_update par la gateway). */
|
||||
discovery_update: [SessionSummary];
|
||||
}
|
||||
|
||||
export interface DiscoveryOptions {
|
||||
ptyManager: PtyManager;
|
||||
projectsDir?: string;
|
||||
sessionsDir?: string;
|
||||
refreshMs?: number;
|
||||
}
|
||||
|
||||
export class DiscoveryService extends EventEmitter<DiscoveryServiceEvents> {
|
||||
private readonly projectsDir: string;
|
||||
private readonly sessionsDir: string;
|
||||
private readonly ptyManager: PtyManager;
|
||||
private readonly refreshMs: number;
|
||||
private timer: NodeJS.Timeout | null = null;
|
||||
private cache: SessionSummary[] = [];
|
||||
private byId = new Map<string, DiscoveredJsonl>();
|
||||
private prevJson = new Map<string, string>();
|
||||
|
||||
constructor(opts: DiscoveryOptions) {
|
||||
super();
|
||||
this.ptyManager = opts.ptyManager;
|
||||
this.projectsDir = opts.projectsDir ?? join(homedir(), '.claude', 'projects');
|
||||
this.sessionsDir = opts.sessionsDir ?? join(homedir(), '.claude', 'sessions');
|
||||
this.refreshMs = opts.refreshMs ?? DEFAULT_REFRESH_MS;
|
||||
}
|
||||
|
||||
start(): void {
|
||||
if (this.timer) return;
|
||||
void this.refresh(); // premier scan asynchrone : ne bloque pas le boot
|
||||
this.timer = setInterval(() => void this.refresh(), this.refreshMs);
|
||||
this.timer.unref(); // ne maintient pas le process en vie
|
||||
}
|
||||
|
||||
stop(): void {
|
||||
if (this.timer) {
|
||||
clearInterval(this.timer);
|
||||
this.timer = null;
|
||||
}
|
||||
}
|
||||
|
||||
/** Cache courant (jamais de scan synchrone dans le chemin chaud). */
|
||||
list(): SessionSummary[] {
|
||||
return this.cache;
|
||||
}
|
||||
|
||||
/** Session découverte (avec son cwd d'origine lu sur disque) — pour resume/fork. null si absente. */
|
||||
getDiscovered(claudeSessionId: string): DiscoveredJsonl | null {
|
||||
return this.byId.get(claudeSessionId) ?? null;
|
||||
}
|
||||
|
||||
/**
|
||||
* Vivacité FRAÎCHE d'une session (relit le registre, ne se fie pas au cache) — garde-fou
|
||||
* anti-corruption : la route /resume doit refuser une session devenue vivante depuis le dernier scan.
|
||||
*/
|
||||
isClaudeSessionLive(claudeSessionId: string): boolean {
|
||||
const entry = readRegistry(this.sessionsDir).find((r) => r.claudeSessionId === claudeSessionId);
|
||||
return entry?.live ?? false;
|
||||
}
|
||||
|
||||
/** Recalcule le cache depuis le disque. Tolérant : ne lève jamais. */
|
||||
async refresh(): Promise<void> {
|
||||
const regBySid = new Map<string, RegistryEntry>();
|
||||
for (const r of readRegistry(this.sessionsDir)) {
|
||||
if (r.claudeSessionId) regBySid.set(r.claudeSessionId, r);
|
||||
}
|
||||
const known = this.ptyManager.knownClaudeSessionIds();
|
||||
|
||||
// Dédoublonnage des JSONL par claudeSessionId (on retient le plus récent).
|
||||
const latest = new Map<string, DiscoveredJsonl>();
|
||||
for (const d of await scanProjects(this.projectsDir)) {
|
||||
const prev = latest.get(d.claudeSessionId);
|
||||
if (!prev || d.mtimeMs > prev.mtimeMs) latest.set(d.claudeSessionId, d);
|
||||
}
|
||||
|
||||
const summaries: SessionSummary[] = [];
|
||||
const byId = new Map<string, DiscoveredJsonl>();
|
||||
for (const d of latest.values()) {
|
||||
if (known.has(d.claudeSessionId)) continue; // déjà gérée comme session managée
|
||||
byId.set(d.claudeSessionId, d);
|
||||
const r = regBySid.get(d.claudeSessionId);
|
||||
const live = r?.live ?? false;
|
||||
summaries.push({
|
||||
id: d.claudeSessionId, // les découvertes s'identifient par leur claudeSessionId
|
||||
cwd: d.cwd,
|
||||
command: 'claude',
|
||||
title: d.title,
|
||||
status: live ? 'running' : 'exited',
|
||||
live,
|
||||
createdAt: new Date(d.mtimeMs).toISOString(),
|
||||
endedAt: null,
|
||||
exitCode: null,
|
||||
clients: 0,
|
||||
source: 'discovered',
|
||||
claudeSessionId: d.claudeSessionId,
|
||||
pid: live ? (r?.pid ?? null) : null,
|
||||
resumable: !live, // morte → --resume direct ; vivante → fork/observe (jamais resume : corruption)
|
||||
attachable: false, // Arboretum ne tient pas le PTY d'une session externe
|
||||
registryStatus: r?.status ?? null,
|
||||
});
|
||||
}
|
||||
|
||||
this.cache = summaries;
|
||||
this.byId = byId;
|
||||
|
||||
// Diff : n'émettre que les sessions nouvelles ou modifiées.
|
||||
const nextJson = new Map<string, string>();
|
||||
for (const s of summaries) {
|
||||
const j = JSON.stringify(s);
|
||||
nextJson.set(s.id, j);
|
||||
if (this.prevJson.get(s.id) !== j) this.emit('discovery_update', s);
|
||||
}
|
||||
this.prevJson = nextJson;
|
||||
}
|
||||
}
|
||||
|
||||
/** Fusionne sessions managées + découvertes en dédoublonnant par claudeSessionId (managé prioritaire). */
|
||||
export function mergeSessions(managed: SessionSummary[], discovered: SessionSummary[]): SessionSummary[] {
|
||||
const managedSids = new Set(managed.map((s) => s.claudeSessionId).filter((x): x is string => x != null));
|
||||
return [...managed, ...discovered.filter((s) => !(s.claudeSessionId && managedSids.has(s.claudeSessionId)))];
|
||||
}
|
||||
160
packages/server/src/core/jsonl-discovery.ts
Normal file
160
packages/server/src/core/jsonl-discovery.ts
Normal file
@@ -0,0 +1,160 @@
|
||||
// Découverte des sessions Claude sur disque (~/.claude/projects/**/*.jsonl).
|
||||
// Transcription TS du spike S4 (scan.mjs) : parseur strictement tolérant, lecture partielle
|
||||
// tête+queue (~1,4 s pour 1300+ fichiers), 99,8 % des transcripts exploitables.
|
||||
// I/O ASYNCHRONES : le scan rend la main à la boucle d'événements entre fichiers, pour ne pas
|
||||
// figer le streaming des terminaux (il tourne périodiquement dans le DiscoveryService).
|
||||
import { readdir, stat, open } from 'node:fs/promises';
|
||||
import { join } from 'node:path';
|
||||
|
||||
const HEAD_BYTES = 256 * 1024;
|
||||
const TAIL_BYTES = 64 * 1024;
|
||||
const TITLE_MAX = 120;
|
||||
|
||||
/** Reproduit le nom de dossier ~/.claude/projects à partir d'un cwd (validé 100 % — spike S4). */
|
||||
export function munge(cwd: string): string {
|
||||
return cwd.replace(/[^A-Za-z0-9]/g, '-');
|
||||
}
|
||||
|
||||
export interface DiscoveredJsonl {
|
||||
claudeSessionId: string;
|
||||
cwd: string;
|
||||
title: string | null;
|
||||
gitBranch: string | null;
|
||||
version: string | null;
|
||||
mtimeMs: number;
|
||||
/** chemin absolu du .jsonl */
|
||||
file: string;
|
||||
}
|
||||
|
||||
async function readChunk(path: string, position: number, length: number): Promise<string> {
|
||||
const fh = await open(path, 'r');
|
||||
try {
|
||||
const buf = Buffer.alloc(length);
|
||||
const { bytesRead } = await fh.read(buf, 0, length, position);
|
||||
return buf.subarray(0, bytesRead).toString('utf8');
|
||||
} finally {
|
||||
await fh.close();
|
||||
}
|
||||
}
|
||||
|
||||
/** Parse tolérant : chaque ligne en try/catch, lignes coupées/malformées ignorées (jamais de throw). */
|
||||
export function parseLines(text: string, partialFirst = false, partialLast = false): Array<Record<string, unknown>> {
|
||||
const lines = text.split('\n');
|
||||
if (partialFirst) lines.shift(); // 1re ligne potentiellement coupée (lecture de queue)
|
||||
if (partialLast) lines.pop(); // dernière ligne potentiellement coupée (lecture de tête)
|
||||
const out: Array<Record<string, unknown>> = [];
|
||||
for (const l of lines) {
|
||||
if (!l.trim()) continue;
|
||||
try {
|
||||
const o = JSON.parse(l) as unknown;
|
||||
if (o && typeof o === 'object') out.push(o as Record<string, unknown>);
|
||||
} catch {
|
||||
/* ligne coupée ou malformée : ignorée */
|
||||
}
|
||||
}
|
||||
return out;
|
||||
}
|
||||
|
||||
interface Meta {
|
||||
sessionId?: string;
|
||||
cwd?: string;
|
||||
gitBranch?: string;
|
||||
version?: string;
|
||||
aiTitle?: string;
|
||||
summary?: string;
|
||||
lastPrompt?: string;
|
||||
firstUserPrompt?: string;
|
||||
}
|
||||
|
||||
function asString(v: unknown): string | undefined {
|
||||
return typeof v === 'string' && v.length > 0 ? v : undefined;
|
||||
}
|
||||
|
||||
function extractMeta(objs: Array<Record<string, unknown>>, meta: Meta): void {
|
||||
const setOnce = (key: 'sessionId' | 'cwd' | 'gitBranch' | 'version', v: unknown): void => {
|
||||
if (meta[key] === undefined) {
|
||||
const s = asString(v);
|
||||
if (s) meta[key] = s;
|
||||
}
|
||||
};
|
||||
for (const o of objs) {
|
||||
setOnce('sessionId', o.sessionId);
|
||||
setOnce('cwd', o.cwd);
|
||||
setOnce('gitBranch', o.gitBranch);
|
||||
setOnce('version', o.version);
|
||||
// Titre : on retient la dernière valeur vue (la plus récente) — head puis tail → la queue gagne.
|
||||
const ai = asString(o.aiTitle);
|
||||
if (ai) meta.aiTitle = ai;
|
||||
const sum = asString(o.summary);
|
||||
if (sum) meta.summary = sum;
|
||||
const lp = asString(o.lastPrompt);
|
||||
if (lp) meta.lastPrompt = lp;
|
||||
// Premier prompt utilisateur (depuis la tête) : repli ultime pour le titre.
|
||||
if (meta.firstUserPrompt === undefined && o.type === 'user') {
|
||||
const msg = o.message;
|
||||
if (msg && typeof msg === 'object') {
|
||||
const text = asString((msg as Record<string, unknown>).content);
|
||||
if (text) meta.firstUserPrompt = text;
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/** Priorité de titre actée S4 : ai-title > summary > last-prompt > 1er prompt user. */
|
||||
function pickTitle(meta: Meta): string | null {
|
||||
const raw = meta.aiTitle ?? meta.summary ?? meta.lastPrompt ?? meta.firstUserPrompt ?? null;
|
||||
if (!raw) return null;
|
||||
const t = raw.replace(/\s+/g, ' ').trim();
|
||||
return t.length > TITLE_MAX ? `${t.slice(0, TITLE_MAX - 1)}…` : t;
|
||||
}
|
||||
|
||||
/**
|
||||
* Scanne `projectsDir` en lecture partielle tête+queue. Tolérant : un dossier/fichier illisible ou
|
||||
* malformé est ignoré (jamais de throw). Un fichier sans sessionId+cwd exploitables est écarté.
|
||||
*/
|
||||
export async function scanProjects(projectsDir: string): Promise<DiscoveredJsonl[]> {
|
||||
let dirs: string[];
|
||||
try {
|
||||
dirs = await readdir(projectsDir);
|
||||
} catch {
|
||||
return []; // dossier absent (Claude jamais lancé) : pas une erreur
|
||||
}
|
||||
const out: DiscoveredJsonl[] = [];
|
||||
for (const dir of dirs) {
|
||||
const dirPath = join(projectsDir, dir);
|
||||
let entries;
|
||||
try {
|
||||
entries = await readdir(dirPath, { withFileTypes: true });
|
||||
} catch {
|
||||
continue;
|
||||
}
|
||||
for (const e of entries) {
|
||||
if (!e.isFile() || !e.name.endsWith('.jsonl')) continue;
|
||||
const fp = join(dirPath, e.name);
|
||||
try {
|
||||
const st = await stat(fp);
|
||||
if (st.size === 0) continue;
|
||||
const meta: Meta = {};
|
||||
const head = await readChunk(fp, 0, Math.min(HEAD_BYTES, st.size));
|
||||
extractMeta(parseLines(head, false, st.size > HEAD_BYTES), meta);
|
||||
if (st.size > HEAD_BYTES + TAIL_BYTES) {
|
||||
const tail = await readChunk(fp, st.size - TAIL_BYTES, TAIL_BYTES);
|
||||
extractMeta(parseLines(tail, true, false), meta);
|
||||
}
|
||||
if (!meta.sessionId || !meta.cwd) continue; // pas une session exploitable
|
||||
out.push({
|
||||
claudeSessionId: meta.sessionId,
|
||||
cwd: meta.cwd,
|
||||
title: pickTitle(meta),
|
||||
gitBranch: meta.gitBranch ?? null,
|
||||
version: meta.version ?? null,
|
||||
mtimeMs: st.mtimeMs,
|
||||
file: fp,
|
||||
});
|
||||
} catch {
|
||||
/* fichier illisible : ignoré */
|
||||
}
|
||||
}
|
||||
}
|
||||
return out;
|
||||
}
|
||||
@@ -1,14 +1,20 @@
|
||||
import { EventEmitter } from 'node:events';
|
||||
import { existsSync, statSync } from 'node:fs';
|
||||
import { randomUUID } from 'node:crypto';
|
||||
import { homedir } from 'node:os';
|
||||
import { join } from 'node:path';
|
||||
import pty from '@homebridge/node-pty-prebuilt-multiarch';
|
||||
import { FLOW, REPLAY_TAIL_BYTES, type SessionSummary } from '@arboretum/shared';
|
||||
import { RingBuffer } from './ring-buffer.js';
|
||||
import { buildSpawnSpec } from './claude-launcher.js';
|
||||
import { findByPid } from './session-registry.js';
|
||||
import type { Db } from '../db/index.js';
|
||||
|
||||
const RING_CAPACITY = 2 * 1024 * 1024;
|
||||
const KILL_GRACE_MS = 5000;
|
||||
/** Capture du claudeSessionId après spawn : poll du registre par pid (waitReady validé S1). */
|
||||
const CLAUDE_ID_POLL_MS = 400;
|
||||
const CLAUDE_ID_TIMEOUT_MS = 60_000;
|
||||
|
||||
/** Lien entre un client WS attaché et une session. La gateway fournit les callbacks d'envoi. */
|
||||
export interface ClientBinding {
|
||||
@@ -36,6 +42,8 @@ interface ManagedSession {
|
||||
paused: boolean;
|
||||
exited: { exitCode: number | null; signal: number | null } | null;
|
||||
killTimer: NodeJS.Timeout | null;
|
||||
/** ID interne du CLI claude, résolu via le registre après spawn (null pour bash / pas encore prêt). */
|
||||
claudeSessionId: string | null;
|
||||
}
|
||||
|
||||
export interface PtyManagerEvents {
|
||||
@@ -46,17 +54,21 @@ export interface PtyManagerEvents {
|
||||
export class PtyManager extends EventEmitter<PtyManagerEvents> {
|
||||
private readonly live = new Map<string, ManagedSession>();
|
||||
|
||||
constructor(private readonly db: Db) {
|
||||
constructor(
|
||||
private readonly db: Db,
|
||||
private readonly sessionsDir: string = join(homedir(), '.claude', 'sessions'),
|
||||
) {
|
||||
super();
|
||||
}
|
||||
|
||||
spawn(opts: { cwd: string; command?: 'claude' | 'bash' }): SessionSummary {
|
||||
spawn(opts: { cwd: string; command?: 'claude' | 'bash'; resume?: { claudeSessionId: string; fork?: boolean } }): SessionSummary {
|
||||
const cwd = opts.cwd;
|
||||
if (!existsSync(cwd) || !statSync(cwd).isDirectory()) {
|
||||
throw Object.assign(new Error(`Not a directory: ${cwd}`), { statusCode: 400 });
|
||||
}
|
||||
const command = opts.command ?? 'claude';
|
||||
const spec = buildSpawnSpec(command);
|
||||
// Un resume/fork est toujours une session claude (le cwd d'origine est garanti par l'appelant — S1).
|
||||
const command = opts.resume ? 'claude' : (opts.command ?? 'claude');
|
||||
const spec = buildSpawnSpec({ command, ...(opts.resume ? { resume: opts.resume } : {}) });
|
||||
const id = randomUUID();
|
||||
const proc = pty.spawn(spec.file, spec.args, {
|
||||
name: 'xterm-256color',
|
||||
@@ -77,26 +89,66 @@ export class PtyManager extends EventEmitter<PtyManagerEvents> {
|
||||
paused: false,
|
||||
exited: null,
|
||||
killTimer: null,
|
||||
claudeSessionId: null,
|
||||
};
|
||||
this.live.set(id, session);
|
||||
this.db
|
||||
.prepare('INSERT INTO sessions (id, cwd, command, created_at) VALUES (?, ?, ?, ?)')
|
||||
.run(id, cwd, command, session.createdAt);
|
||||
.prepare('INSERT INTO sessions (id, cwd, command, created_at, resumed_from) VALUES (?, ?, ?, ?, ?)')
|
||||
.run(id, cwd, command, session.createdAt, opts.resume?.claudeSessionId ?? null);
|
||||
|
||||
proc.onData((data) => this.handleOutput(session, Buffer.from(data, 'utf8')));
|
||||
proc.onExit(({ exitCode, signal }) => this.handleExit(session, exitCode, signal ?? null));
|
||||
if (command === 'claude') this.captureClaudeSessionId(session);
|
||||
|
||||
const summary = this.summarize(session);
|
||||
this.emit('session_update', summary);
|
||||
return summary;
|
||||
}
|
||||
|
||||
/** Résout le claudeSessionId du CLI en pollant le registre par pid, puis le persiste (waitReady — S1). */
|
||||
private captureClaudeSessionId(s: ManagedSession): void {
|
||||
const deadline = Date.now() + CLAUDE_ID_TIMEOUT_MS;
|
||||
const tick = (): void => {
|
||||
if (s.exited || s.claudeSessionId) return;
|
||||
const entry = findByPid(this.sessionsDir, s.proc.pid);
|
||||
if (entry?.claudeSessionId) {
|
||||
s.claudeSessionId = entry.claudeSessionId;
|
||||
this.db.prepare('UPDATE sessions SET claude_session_id = ? WHERE id = ?').run(entry.claudeSessionId, s.id);
|
||||
this.emit('session_update', this.summarize(s));
|
||||
return;
|
||||
}
|
||||
if (Date.now() < deadline) setTimeout(tick, CLAUDE_ID_POLL_MS).unref();
|
||||
};
|
||||
setTimeout(tick, CLAUDE_ID_POLL_MS).unref();
|
||||
}
|
||||
|
||||
/** Session managée VIVANTE portant ce claudeSessionId (garde-fou anti-resume d'une session vivante). */
|
||||
findLiveByClaudeSessionId(claudeSessionId: string): SessionSummary | null {
|
||||
for (const s of this.live.values()) {
|
||||
if (s.claudeSessionId === claudeSessionId) return this.summarize(s);
|
||||
}
|
||||
return null;
|
||||
}
|
||||
|
||||
/**
|
||||
* Ensemble des claudeSessionId déjà connus d'Arboretum (sessions managées vivantes + historique DB).
|
||||
* Sert au DiscoveryService à exclure les doublons (une managée ne doit pas réapparaître en découverte).
|
||||
*/
|
||||
knownClaudeSessionIds(): Set<string> {
|
||||
const rows = this.db
|
||||
.prepare('SELECT claude_session_id FROM sessions WHERE claude_session_id IS NOT NULL')
|
||||
.all() as Array<{ claude_session_id: string }>;
|
||||
const set = new Set(rows.map((r) => r.claude_session_id));
|
||||
for (const s of this.live.values()) if (s.claudeSessionId) set.add(s.claudeSessionId);
|
||||
return set;
|
||||
}
|
||||
|
||||
list(): SessionSummary[] {
|
||||
const liveSummaries = [...this.live.values()].map((s) => this.summarize(s));
|
||||
const liveIds = new Set(this.live.keys());
|
||||
const rows = this.db
|
||||
.prepare('SELECT id, cwd, command, title, created_at, ended_at, exit_code FROM sessions ORDER BY created_at DESC LIMIT 100')
|
||||
.all() as Array<{ id: string; cwd: string; command: string; title: string | null; created_at: string; ended_at: string | null; exit_code: number | null }>;
|
||||
.prepare('SELECT id, cwd, command, title, created_at, ended_at, exit_code, claude_session_id FROM sessions ORDER BY created_at DESC LIMIT 100')
|
||||
.all() as Array<{ id: string; cwd: string; command: string; title: string | null; created_at: string; ended_at: string | null; exit_code: number | null; claude_session_id: string | null }>;
|
||||
const historical: SessionSummary[] = rows
|
||||
.filter((r) => !liveIds.has(r.id))
|
||||
.map((r) => ({
|
||||
@@ -110,6 +162,13 @@ export class PtyManager extends EventEmitter<PtyManagerEvents> {
|
||||
endedAt: r.ended_at,
|
||||
exitCode: r.exit_code,
|
||||
clients: 0,
|
||||
source: 'managed',
|
||||
claudeSessionId: r.claude_session_id,
|
||||
pid: null,
|
||||
// une session claude morte avec un claudeSessionId connu est reprenable (--resume direct).
|
||||
resumable: r.command === 'claude' && r.claude_session_id != null,
|
||||
attachable: false,
|
||||
registryStatus: null,
|
||||
}));
|
||||
return [...liveSummaries, ...historical];
|
||||
}
|
||||
@@ -265,6 +324,13 @@ export class PtyManager extends EventEmitter<PtyManagerEvents> {
|
||||
endedAt: null,
|
||||
exitCode: s.exited?.exitCode ?? null,
|
||||
clients: s.clients.size,
|
||||
source: 'managed',
|
||||
claudeSessionId: s.claudeSessionId,
|
||||
pid: s.exited ? null : s.proc.pid,
|
||||
// une managée vivante ne se resume pas (corruption) ; une managée claude morte oui.
|
||||
resumable: !!s.exited && s.command === 'claude' && s.claudeSessionId != null,
|
||||
attachable: !s.exited,
|
||||
registryStatus: null, // P2 : statut fin des managées via claude-adapter (P3-B)
|
||||
};
|
||||
}
|
||||
}
|
||||
|
||||
92
packages/server/src/core/session-registry.ts
Normal file
92
packages/server/src/core/session-registry.ts
Normal file
@@ -0,0 +1,92 @@
|
||||
// Lecture du registre ~/.claude/sessions (un fichier <pid>.json par session CLI vive).
|
||||
// Module partagé : P2 l'utilise pour la vivacité + le lookup par pid ; P3-B (claude-adapter)
|
||||
// s'en sert comme source primaire d'état fin (busy/idle/waiting + waitingFor).
|
||||
import { readdirSync, readFileSync } from 'node:fs';
|
||||
import { join } from 'node:path';
|
||||
import type { SessionRegistryStatus } from '@arboretum/shared';
|
||||
|
||||
export interface RegistryEntry {
|
||||
claudeSessionId: string | null;
|
||||
cwd: string | null;
|
||||
pid: number;
|
||||
procStart: string | null;
|
||||
status: SessionRegistryStatus | null;
|
||||
waitingFor: string | null;
|
||||
/** vivacité réelle = pid + procStart concordants (jamais la simple présence du fichier). */
|
||||
live: boolean;
|
||||
}
|
||||
|
||||
/** starttime (champ 22 de /proc/<pid>/stat) ou null si le process n'existe pas. Linux-only. */
|
||||
export function readProcStart(pid: number): string | null {
|
||||
try {
|
||||
const stat = readFileSync(`/proc/${pid}/stat`, 'utf8');
|
||||
// Le nom du process (champ 2) peut contenir espaces/parenthèses → parser après le dernier ')'.
|
||||
const after = stat.slice(stat.lastIndexOf(')') + 2).split(' ');
|
||||
return after[19] ?? null; // champ 22 global = index 19 après le champ 3
|
||||
} catch {
|
||||
return null;
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Vivacité par pid + procStart (spike S1/S4) : jamais la présence d'un fichier registre (stale possible).
|
||||
* "Incertain → vivant" : si le process existe mais que le registre n'a pas de procStart, on traite vivant.
|
||||
*/
|
||||
export function isLive(pid: number, procStart: string | number | null | undefined): boolean {
|
||||
const actual = readProcStart(pid);
|
||||
if (actual === null) return false; // process absent → mort
|
||||
if (procStart === null || procStart === undefined) return true; // incertain → vivant
|
||||
return String(procStart) === actual; // procStart discordant → pid recyclé → stale → mort
|
||||
}
|
||||
|
||||
function normalizeStatus(v: unknown): SessionRegistryStatus | null {
|
||||
return v === 'busy' || v === 'idle' || v === 'waiting' ? v : null;
|
||||
}
|
||||
|
||||
function toEntry(r: Record<string, unknown>): RegistryEntry | null {
|
||||
const pid = typeof r.pid === 'number' ? r.pid : Number(r.pid);
|
||||
if (!Number.isInteger(pid)) return null;
|
||||
const procStart = r.procStart === undefined || r.procStart === null ? null : String(r.procStart);
|
||||
return {
|
||||
claudeSessionId: typeof r.sessionId === 'string' ? r.sessionId : null,
|
||||
cwd: typeof r.cwd === 'string' ? r.cwd : null,
|
||||
pid,
|
||||
procStart,
|
||||
status: normalizeStatus(r.status),
|
||||
waitingFor: typeof r.waitingFor === 'string' ? r.waitingFor : null,
|
||||
live: isLive(pid, procStart),
|
||||
};
|
||||
}
|
||||
|
||||
/**
|
||||
* Lit toutes les entrées de `sessionsDir`. Tolérant aux fichiers en cours d'écriture / JSON invalide
|
||||
* (ignorés). Dossier absent → tableau vide (pas une erreur).
|
||||
*/
|
||||
export function readRegistry(sessionsDir: string): RegistryEntry[] {
|
||||
let files: string[];
|
||||
try {
|
||||
files = readdirSync(sessionsDir);
|
||||
} catch {
|
||||
return [];
|
||||
}
|
||||
const out: RegistryEntry[] = [];
|
||||
for (const f of files) {
|
||||
if (!f.endsWith('.json')) continue;
|
||||
try {
|
||||
const r = JSON.parse(readFileSync(join(sessionsDir, f), 'utf8')) as Record<string, unknown>;
|
||||
const entry = toEntry(r);
|
||||
if (entry) out.push(entry);
|
||||
} catch {
|
||||
/* fichier en cours d'écriture / JSON invalide : ignoré */
|
||||
}
|
||||
}
|
||||
return out;
|
||||
}
|
||||
|
||||
/**
|
||||
* Entrée registre correspondant à un pid donné (le CLI nomme son fichier <pid>.json, mais on
|
||||
* matche sur le champ `pid` pour ne pas dépendre du nommage). null si absente/illisible.
|
||||
*/
|
||||
export function findByPid(sessionsDir: string, pid: number): RegistryEntry | null {
|
||||
return readRegistry(sessionsDir).find((e) => e.pid === pid) ?? null;
|
||||
}
|
||||
@@ -27,6 +27,15 @@ const MIGRATIONS: Array<{ id: number; sql: string }> = [
|
||||
);
|
||||
`,
|
||||
},
|
||||
{
|
||||
// P2 — découverte & reprise : corrélation avec les sessions Claude sur disque.
|
||||
id: 2,
|
||||
sql: `
|
||||
ALTER TABLE sessions ADD COLUMN claude_session_id TEXT;
|
||||
ALTER TABLE sessions ADD COLUMN resumed_from TEXT;
|
||||
CREATE INDEX idx_sessions_claude_session_id ON sessions(claude_session_id);
|
||||
`,
|
||||
},
|
||||
];
|
||||
|
||||
export type Db = DatabaseSync;
|
||||
|
||||
@@ -13,10 +13,11 @@ const pkg = JSON.parse(
|
||||
async function main(): Promise<void> {
|
||||
const config = loadConfig();
|
||||
const db = openDb(config.dbPath);
|
||||
const { app, auth, manager } = buildApp(config, db, pkg.version);
|
||||
const { app, auth, manager, discovery } = buildApp(config, db, pkg.version);
|
||||
|
||||
const bootstrapToken = auth.ensureBootstrapToken();
|
||||
await app.listen({ port: config.port, host: config.bind });
|
||||
discovery.start(); // scan initial + rafraîchissement périodique des sessions découvertes
|
||||
|
||||
const url = `http://${config.bind === '0.0.0.0' ? '127.0.0.1' : config.bind}:${config.port}`;
|
||||
app.log.info(`Arboretum v${pkg.version} — ${url}`);
|
||||
@@ -36,6 +37,7 @@ async function main(): Promise<void> {
|
||||
if (shuttingDown) return;
|
||||
shuttingDown = true;
|
||||
app.log.info(`${signal} received — draining sessions then exiting`);
|
||||
discovery.stop();
|
||||
manager.shutdown();
|
||||
setTimeout(() => {
|
||||
void app.close().then(() => process.exit(0));
|
||||
|
||||
@@ -1,10 +1,11 @@
|
||||
import type { FastifyInstance } from 'fastify';
|
||||
import type { CreateSessionRequest, SessionResponse, SessionsListResponse } from '@arboretum/shared';
|
||||
import type { PtyManager } from '../core/pty-manager.js';
|
||||
import { mergeSessions, type DiscoveryService } from '../core/discovery-service.js';
|
||||
|
||||
export function registerSessionRoutes(app: FastifyInstance, manager: PtyManager): void {
|
||||
export function registerSessionRoutes(app: FastifyInstance, manager: PtyManager, discovery: DiscoveryService): void {
|
||||
app.get('/api/v1/sessions', async (): Promise<SessionsListResponse> => {
|
||||
return { sessions: manager.list() };
|
||||
return { sessions: mergeSessions(manager.list(), discovery.list()) };
|
||||
});
|
||||
|
||||
app.post('/api/v1/sessions', async (req, reply) => {
|
||||
@@ -25,6 +26,45 @@ export function registerSessionRoutes(app: FastifyInstance, manager: PtyManager)
|
||||
}
|
||||
});
|
||||
|
||||
// Reprise d'une session morte : nouveau PTY managé `--resume <id>` DANS SON CWD D'ORIGINE (spike S1).
|
||||
// Le cwd n'est jamais fourni par le client : il est lu sur disque via la découverte.
|
||||
app.post('/api/v1/sessions/:id/resume', async (req, reply) => {
|
||||
const { id } = req.params as { id: string };
|
||||
const discovered = discovery.getDiscovered(id);
|
||||
if (!discovered) {
|
||||
return reply.status(404).send({ error: { code: 'NOT_FOUND', message: 'No resumable session with this id' } });
|
||||
}
|
||||
// Garde-fou anti-corruption : jamais de resume direct d'une session vivante (vérif FRAÎCHE).
|
||||
if (discovery.isClaudeSessionLive(id) || manager.findLiveByClaudeSessionId(id)) {
|
||||
return reply.status(409).send({ error: { code: 'SESSION_LIVE', message: 'Session is live — fork it instead' } });
|
||||
}
|
||||
try {
|
||||
const session = manager.spawn({ cwd: discovered.cwd, resume: { claudeSessionId: id } });
|
||||
const res: SessionResponse = { session };
|
||||
return reply.status(201).send(res);
|
||||
} catch (err) {
|
||||
const statusCode = (err as { statusCode?: number }).statusCode ?? 500;
|
||||
return reply.status(statusCode).send({ error: { code: 'SPAWN_FAILED', message: (err as Error).message } });
|
||||
}
|
||||
});
|
||||
|
||||
// Fork : duplique une session (vivante ou morte) sans la corrompre (`--resume <id> --fork-session`).
|
||||
app.post('/api/v1/sessions/:id/fork', async (req, reply) => {
|
||||
const { id } = req.params as { id: string };
|
||||
const discovered = discovery.getDiscovered(id);
|
||||
if (!discovered) {
|
||||
return reply.status(404).send({ error: { code: 'NOT_FOUND', message: 'No session with this id to fork' } });
|
||||
}
|
||||
try {
|
||||
const session = manager.spawn({ cwd: discovered.cwd, resume: { claudeSessionId: id, fork: true } });
|
||||
const res: SessionResponse = { session };
|
||||
return reply.status(201).send(res);
|
||||
} catch (err) {
|
||||
const statusCode = (err as { statusCode?: number }).statusCode ?? 500;
|
||||
return reply.status(statusCode).send({ error: { code: 'SPAWN_FAILED', message: (err as Error).message } });
|
||||
}
|
||||
});
|
||||
|
||||
app.delete('/api/v1/sessions/:id', async (req, reply) => {
|
||||
const { id } = req.params as { id: string };
|
||||
if (!manager.kill(id)) {
|
||||
|
||||
@@ -9,6 +9,7 @@ import {
|
||||
type SessionSummary,
|
||||
} from '@arboretum/shared';
|
||||
import type { ClientBinding, PtyManager } from '../core/pty-manager.js';
|
||||
import type { DiscoveryService } from '../core/discovery-service.js';
|
||||
|
||||
const HEARTBEAT_MS = 30_000;
|
||||
|
||||
@@ -17,7 +18,12 @@ interface ChannelState {
|
||||
binding: ClientBinding;
|
||||
}
|
||||
|
||||
export function registerWsGateway(app: FastifyInstance, manager: PtyManager, serverVersion: string): void {
|
||||
export function registerWsGateway(
|
||||
app: FastifyInstance,
|
||||
manager: PtyManager,
|
||||
discovery: DiscoveryService,
|
||||
serverVersion: string,
|
||||
): void {
|
||||
app.get('/ws', { websocket: true }, (socket: WebSocket, req) => {
|
||||
// L'auth + le check Origin ont eu lieu dans le preValidation global (app.ts).
|
||||
const channels = new Map<number, ChannelState>();
|
||||
@@ -39,8 +45,13 @@ export function registerWsGateway(app: FastifyInstance, manager: PtyManager, ser
|
||||
const onSessionExit = (e: { sessionId: string; exitCode: number | null; signal: number | null }): void => {
|
||||
if (subscribedSessions) send({ type: 'session_exit', ...e });
|
||||
};
|
||||
// Les sessions découvertes (DiscoveryService) sont relayées sur le même flux que les managées.
|
||||
const onDiscoveryUpdate = (session: SessionSummary): void => {
|
||||
if (subscribedSessions) send({ type: 'session_update', session });
|
||||
};
|
||||
manager.on('session_update', onSessionUpdate);
|
||||
manager.on('session_exit', onSessionExit);
|
||||
discovery.on('discovery_update', onDiscoveryUpdate);
|
||||
|
||||
const heartbeat = setInterval(() => {
|
||||
if (!alive) {
|
||||
@@ -148,6 +159,7 @@ export function registerWsGateway(app: FastifyInstance, manager: PtyManager, ser
|
||||
clearInterval(heartbeat);
|
||||
manager.off('session_update', onSessionUpdate);
|
||||
manager.off('session_exit', onSessionExit);
|
||||
discovery.off('discovery_update', onDiscoveryUpdate);
|
||||
for (const [, st] of channels) manager.detach(st.sessionId, st.binding);
|
||||
channels.clear();
|
||||
});
|
||||
|
||||
Reference in New Issue
Block a user