/** * ARIA-patched API Route Handlers * * Erweiterung der npm-Version von claude-max-api-proxy: * - Bei jedem Claude-CLI-`assistant`-Event mit tool_use-Block (Bash, Read, * Edit, Grep, …) wird ein HTTP-POST an die Bridge gefeuert * (ARIA_TOOL_HOOK_URL, default http://aria-bridge:8090/internal/agent-activity). * Bridge spiegelt das als RVS `agent_activity` an App+Diagnostic → * Gedanken-Stream zeigt live was ARIA gerade tool-maessig macht. * - Voller Live-Stream (assistant_text, tool_use mit input, tool_result) * geht an ARIA_STREAM_HOOK_URL → Bridge → RVS `agent_stream` → Diagnostic * "ARIA Live"-View (TeamViewer-mäßiger Mirror der Claude-Code-Session). * - Subprocess-Tracking + POST /v1/cancel-all fuer Not-Aus (Hard-Kill). * - Fire-and-forget, fail-open. Wenn die Bridge nicht antwortet, bricht * der Brain-Call NICHT ab. * * Wird zur Container-Startzeit ueber die npm-Version geschrieben * (siehe docker-compose.yml proxy-Block). */ import { v4 as uuidv4 } from "uuid"; import http from "http"; import fs from "fs"; import { ClaudeSubprocess } from "../subprocess/manager.js"; import { openaiToCli } from "../adapter/openai-to-cli.js"; import { cliResultToOpenai, createDoneChunk, } from "../adapter/cli-to-openai.js"; const TOOL_HOOK_URL = process.env.ARIA_TOOL_HOOK_URL || "http://aria-bridge:8090/internal/agent-activity"; const STREAM_HOOK_URL = process.env.ARIA_STREAM_HOOK_URL || "http://aria-bridge:8090/internal/agent-stream"; const CODE_FILE_HOOK_URL = process.env.ARIA_CODE_FILE_HOOK_URL || "http://aria-bridge:8090/internal/code-file"; // Code-Projekte leben unter /shared/projects// (Volume in proxy + // bridge + brain gemountet). Schreibt/aendert ARIA hier eine Datei, spiegeln // wir den Volltext live in den Code-Editor der App. Nur Dateien unter diesem // Praefix — ARIAs sonstige Datei-Ops (Skills, Configs) bleiben unberuehrt. const PROJECTS_ROOT = "/shared/projects/"; const CODE_FILE_MAX_BYTES = 512 * 1024; /** Zerlegt einen absoluten Pfad unter /shared/projects// → {pid, rel} * oder null wenn er nicht darunter liegt. */ function _parseProjectPath(filePath) { if (typeof filePath !== "string" || !filePath.startsWith(PROJECTS_ROOT)) return null; const rest = filePath.slice(PROJECTS_ROOT.length); const slash = rest.indexOf("/"); if (slash <= 0) return null; return { pid: rest.slice(0, slash), rel: rest.slice(slash + 1) }; } /** Liest die (frisch geschriebene) Datei und pusht sie als code_file an die * Bridge. Fire-and-forget, fail-open. */ function _emitCodeFile(filePath) { try { const parsed = _parseProjectPath(filePath); if (!parsed || !parsed.rel) return; const st = fs.statSync(filePath); if (!st.isFile() || st.size > CODE_FILE_MAX_BYTES) return; const content = fs.readFileSync(filePath, "utf8"); _postJson(CODE_FILE_HOOK_URL, { projectId: parsed.pid, path: parsed.rel, content, version: Date.now(), }); } catch (_) { /* fail-open */ } } // Tool-Output kann sehr lang werden (git log -p, find /). Wir truncaten // hart auf 4 KB pro Event — der User sieht weiterhin den Anfang und einen // "...(N bytes truncated)" Hinweis. Vollstaendiger Output bleibt im Brain // und wird normal verarbeitet, das hier ist NUR fuer den Live-Mirror. const TOOL_RESULT_MAX_CHARS = 4096; const TOOL_INPUT_MAX_CHARS = 2048; // Idle-Timeout: Subprocess wird gekillt wenn ueber IDLE_TIMEOUT_MS keine // Aktivitaet (message/content_delta) ankommt. Loest das alte Hard-Timeout- // Problem fuer lange Agent-Sessions (Pentests etc.) — ARIA darf ewig // arbeiten solange sie regelmaessig was emittiert, aber wenn der Subprocess // hartnaeckig haengt, schlaegt der Watchdog trotzdem zu. // Default 20min Idle. Override via env ARIA_IDLE_TIMEOUT_MS. // 0 = deaktiviert (nicht empfohlen). const IDLE_TIMEOUT_MS = parseInt(process.env.ARIA_IDLE_TIMEOUT_MS || "1200000", 10); /** * Generic Fire-and-forget POST an die Bridge. Keine Awaits, keine Fehler * nach oben. Eingesetzt fuer Tool-Hook + Stream-Hook. */ function _postJson(url, body) { try { const u = new URL(url); const data = JSON.stringify(body); const req = http.request({ method: "POST", hostname: u.hostname, port: u.port || 80, path: u.pathname, headers: { "Content-Type": "application/json", "Content-Length": Buffer.byteLength(data) }, timeout: 2000, }, (res) => { res.resume(); }); req.on("error", () => {}); req.on("timeout", () => req.destroy()); req.write(data); req.end(); } catch (_) { /* niemals weiterwerfen */ } } /** * Pusht einen Tool-Use-Event an die Bridge (alter Gedanken-Stream-Pfad). */ function _emitToolEvent(toolName, projectId) { if (!toolName) return; _postJson(TOOL_HOOK_URL, { tool: String(toolName), projectId: projectId || "" }); } /** * Pusht ein Stream-Event an die Bridge (neuer "ARIA Live"-Pfad). * kind: "start" | "text" | "tool_use" | "tool_result" | "end" */ function _emitStreamEvent(requestId, kind, fields) { _postJson(STREAM_HOOK_URL, { requestId, kind, ts: Date.now(), ...fields }); } function _truncate(str, max) { if (typeof str !== "string") str = String(str ?? ""); if (str.length <= max) return { text: str, truncatedBytes: 0 }; return { text: str.slice(0, max), truncatedBytes: str.length - max }; } // ── Subprocess-Tracking fuer Not-Aus ────────────────────────── // requestId → ClaudeSubprocess. Eintraege werden beim close/result-Event // wieder entfernt. /v1/cancel-all iteriert und ruft .kill() auf jeden. // Wert: { subprocess, projectId }. projectId erlaubt kontext-scoped Cancel // (nur die Subprozesse EINES Projekts killen statt aller). const _activeSubprocesses = new Map(); function _trackSubprocess(requestId, subprocess, projectId) { _activeSubprocesses.set(requestId, { subprocess, projectId: projectId || "" }); const cleanup = () => _activeSubprocesses.delete(requestId); subprocess.on("close", cleanup); subprocess.on("error", cleanup); } /** * Idle-Watchdog: killt den Subprocess wenn ueber IDLE_TIMEOUT_MS hinweg * keine message/content_delta Events ankommen. Wird beim Start gesetzt, * bei jedem Event reset, bei close/error/result gestoppt. * * Stream-Event 'end' wird durch den normalen close-Listener im Handler * gefeuert — wir muessen hier nichts extra emittieren. */ function _attachIdleWatchdog(subprocess, requestId) { if (!IDLE_TIMEOUT_MS || IDLE_TIMEOUT_MS <= 0) return; // disabled let timer = null; let killed = false; function _kill() { if (killed) return; killed = true; const mins = Math.round(IDLE_TIMEOUT_MS / 60000); console.warn(`[aria-idle] killing subprocess ${requestId} after ${mins}min idle`); try { subprocess.kill(); } catch (_) {} _emitStreamEvent(requestId, "end", { reason: "idle_timeout", idleMs: IDLE_TIMEOUT_MS }); } function _reset() { if (killed) return; if (timer) clearTimeout(timer); timer = setTimeout(_kill, IDLE_TIMEOUT_MS); } function _stop() { if (timer) { clearTimeout(timer); timer = null; } } // Initial-Timer setzen _reset(); // Jedes Event vom Subprozess zaehlt als Lebenszeichen subprocess.on("message", _reset); subprocess.on("content_delta", _reset); // Result/close/error → endgueltig stop subprocess.on("result", _stop); subprocess.on("close", _stop); subprocess.on("error", _stop); } /** * Hookt assistant + user Events und pusht beides an Bridge: * - Alt-API: nur Tool-Namen an /internal/agent-activity (Gedanken-Stream) * - Neu-API: voller Stream (text/tool_use/tool_result) an /internal/agent-stream */ function _attachToolHook(subprocess, requestId, projectId) { // tool_use_id → file_path fuer Write/Edit, damit wir beim (erfolgreichen) // tool_result die frisch geschriebene Datei aus /shared lesen koennen. const _pendingFileWrites = new Map(); subprocess.on("assistant", (message) => { try { const blocks = message?.message?.content || []; for (const b of blocks) { if (!b) continue; if (b.type === "tool_use") { if (b.name) _emitToolEvent(b.name, projectId); if ((b.name === "Write" || b.name === "Edit" || b.name === "MultiEdit") && b.id && b.input && typeof b.input.file_path === "string") { _pendingFileWrites.set(b.id, b.input.file_path); } const inputStr = b.input ? JSON.stringify(b.input) : ""; const inp = _truncate(inputStr, TOOL_INPUT_MAX_CHARS); _emitStreamEvent(requestId, "tool_use", { projectId: projectId || "", id: b.id || null, name: b.name || "", input: inp.text, inputTruncatedBytes: inp.truncatedBytes, }); } else if (b.type === "text" && b.text) { _emitStreamEvent(requestId, "text", { projectId: projectId || "", text: b.text }); } else if (b.type === "thinking" && b.thinking) { // Wenn das Modell Extended Thinking emittiert — selten in // Claude Code CLI, aber moeglich. Markieren wir extra. _emitStreamEvent(requestId, "thinking", { text: b.thinking }); } } } catch (_) { /* fail-open */ } }); // tool_result Blocks kommen in user-Messages — die werden vom // subprocess-Manager NICHT als 'user'-Event emittiert (gibt's nicht), // sondern nur ueber das generische 'message'-Event mit type:'user'. // 'message' feuert auch fuer assistant/result — wir filtern auf user // damit wir nicht doppelt rendern (assistant geht ueber den eigenen // assistant-Handler oben). subprocess.on("message", (message) => { try { if (message?.type !== "user") return; const blocks = message?.message?.content || []; for (const b of blocks) { if (b && b.type === "tool_result") { let content = ""; if (typeof b.content === "string") content = b.content; else if (Array.isArray(b.content)) { content = b.content.map(c => (c && c.type === "text" && c.text) ? c.text : "").join(""); } const out = _truncate(content, TOOL_RESULT_MAX_CHARS); _emitStreamEvent(requestId, "tool_result", { id: b.tool_use_id || null, content: out.text, truncatedBytes: out.truncatedBytes, isError: b.is_error === true, }); // Write/Edit erfolgreich → Datei live in den Code-Editor spiegeln. if (b.tool_use_id && b.is_error !== true && _pendingFileWrites.has(b.tool_use_id)) { _emitCodeFile(_pendingFileWrites.get(b.tool_use_id)); _pendingFileWrites.delete(b.tool_use_id); } } } } catch (_) { /* fail-open */ } }); } /** * Handle POST /v1/chat/completions * * Main endpoint for chat requests, supports both streaming and non-streaming */ export async function handleChatCompletions(req, res) { const requestId = uuidv4().replace(/-/g, "").slice(0, 24); const body = req.body; const stream = body.stream === true; try { // Validate request if (!body.messages || !Array.isArray(body.messages) || body.messages.length === 0) { res.status(400).json({ error: { message: "messages is required and must be a non-empty array", type: "invalid_request_error", code: "invalid_messages", }, }); return; } // Convert to CLI input format const cliInput = openaiToCli(body); // ARIA: Projekt-Kontext (vom Brain via aria_project_id). Fuer // kontext-getaggte Activity-/Stream-Events + kontext-scoped Cancel. const ariaProjectId = String(body.aria_project_id || ""); const subprocess = new ClaudeSubprocess(); // ARIA-Patch: Tool-Use-Events + voller Live-Stream an die Bridge. // Plus: Subprocess fuer Not-Aus tracken (Hard-Kill via /v1/cancel-all). // Plus: Idle-Watchdog — Subprocess darf ewig laufen solange Events // kommen, wird aber gekillt nach IDLE_TIMEOUT_MS Inaktivitaet. _attachToolHook(subprocess, requestId, ariaProjectId); _trackSubprocess(requestId, subprocess, ariaProjectId); _attachIdleWatchdog(subprocess, requestId); _emitStreamEvent(requestId, "start", { model: body.model || null, projectId: ariaProjectId }); subprocess.on("result", () => _emitStreamEvent(requestId, "end", { reason: "result" })); subprocess.on("close", (code) => _emitStreamEvent(requestId, "end", { reason: "close", code })); subprocess.on("error", (err) => _emitStreamEvent(requestId, "end", { reason: "error", error: String(err?.message || err) })); if (stream) { await handleStreamingResponse(req, res, subprocess, cliInput, requestId); } else { await handleNonStreamingResponse(res, subprocess, cliInput, requestId); } } catch (error) { const message = error instanceof Error ? error.message : "Unknown error"; console.error("[handleChatCompletions] Error:", message); if (!res.headersSent) { res.status(500).json({ error: { message, type: "server_error", code: null, }, }); } } } /** * Handle streaming response (SSE) * * IMPORTANT: The Express req.on("close") event fires when the request body * is fully received, NOT when the client disconnects. For SSE connections, * we use res.on("close") to detect actual client disconnection. */ async function handleStreamingResponse(req, res, subprocess, cliInput, requestId) { // Set SSE headers res.setHeader("Content-Type", "text/event-stream"); res.setHeader("Cache-Control", "no-cache"); res.setHeader("Connection", "keep-alive"); res.setHeader("X-Request-Id", requestId); // CRITICAL: Flush headers immediately to establish SSE connection // Without this, headers are buffered and client times out waiting res.flushHeaders(); // Send initial comment to confirm connection is alive res.write(":ok\n\n"); return new Promise((resolve, reject) => { let isFirst = true; let lastModel = "claude-sonnet-4"; let isComplete = false; // Handle actual client disconnect (response stream closed) res.on("close", () => { if (!isComplete) { // Client disconnected before response completed - kill subprocess subprocess.kill(); } resolve(); }); // Handle streaming content deltas subprocess.on("content_delta", (event) => { const text = event.event.delta?.text || ""; if (text && !res.writableEnded) { const chunk = { id: `chatcmpl-${requestId}`, object: "chat.completion.chunk", created: Math.floor(Date.now() / 1000), model: lastModel, choices: [{ index: 0, delta: { role: isFirst ? "assistant" : undefined, content: text, }, finish_reason: null, }], }; res.write(`data: ${JSON.stringify(chunk)}\n\n`); isFirst = false; } }); // Handle final assistant message (for model name) subprocess.on("assistant", (message) => { lastModel = message.message.model; }); subprocess.on("result", (_result) => { isComplete = true; if (!res.writableEnded) { // Send final done chunk with finish_reason const doneChunk = createDoneChunk(requestId, lastModel); res.write(`data: ${JSON.stringify(doneChunk)}\n\n`); res.write("data: [DONE]\n\n"); res.end(); } resolve(); }); subprocess.on("error", (error) => { console.error("[Streaming] Error:", error.message); if (!res.writableEnded) { res.write(`data: ${JSON.stringify({ error: { message: error.message, type: "server_error", code: null }, })}\n\n`); res.end(); } resolve(); }); subprocess.on("close", (code) => { // Subprocess exited - ensure response is closed if (!res.writableEnded) { if (code !== 0 && !isComplete) { // Abnormal exit without result - send error res.write(`data: ${JSON.stringify({ error: { message: `Process exited with code ${code}`, type: "server_error", code: null }, })}\n\n`); } res.write("data: [DONE]\n\n"); res.end(); } resolve(); }); // Start the subprocess subprocess.start(cliInput.prompt, { model: cliInput.model, sessionId: cliInput.sessionId, // ARIA: echter System-Prompt-Kanal — manager.js reicht das (sobald // gepatcht) als --system-prompt (VOLLER Replace) an die CLI. Aktuell // ignoriert ein ungepatchter manager diese Extra-Option gefahrlos. systemPrompt: cliInput.systemPrompt, }).catch((err) => { console.error("[Streaming] Subprocess start error:", err); reject(err); }); }); } /** * Handle non-streaming response */ async function handleNonStreamingResponse(res, subprocess, cliInput, requestId) { return new Promise((resolve) => { let finalResult = null; let isComplete = false; // Client-Disconnect-Handler — wenn Brain die HTTP-Verbindung kappt // (z.B. nach Read-Timeout), den noch laufenden Subprocess killen. // Im Streaming-Branch existiert das schon; non-streaming hatte's // bisher nicht → Subprozess lief verwaist weiter, Ressourcen-Leak. res.on("close", () => { if (!isComplete) { console.warn("[NonStreaming] Client disconnected before result — killing subprocess", requestId); try { subprocess.kill(); } catch (_) {} } resolve(); }); subprocess.on("result", (result) => { finalResult = result; }); subprocess.on("error", (error) => { console.error("[NonStreaming] Error:", error.message); isComplete = true; if (!res.headersSent) { res.status(500).json({ error: { message: error.message, type: "server_error", code: null, }, }); } resolve(); }); subprocess.on("close", (code) => { isComplete = true; if (res.writableEnded) { // Client ist eh schon weg — nichts mehr zu senden. resolve(); return; } if (finalResult) { res.json(cliResultToOpenai(finalResult, requestId)); } else if (!res.headersSent) { res.status(500).json({ error: { message: `Claude CLI exited with code ${code} without response`, type: "server_error", code: null, }, }); } resolve(); }); // Start the subprocess subprocess .start(cliInput.prompt, { model: cliInput.model, sessionId: cliInput.sessionId, // ARIA: echter System-Prompt-Kanal (siehe Streaming-Branch). systemPrompt: cliInput.systemPrompt, }) .catch((error) => { res.status(500).json({ error: { message: error.message, type: "server_error", code: null, }, }); resolve(); }); }); } /** * Handle GET /v1/models * * Returns available models */ // Kuratierte Tier-Liste. ARIA laeuft ueber das Claude-Max-Abo via CLI — // waehlbar ist der TIER (opus/sonnet/haiku), nicht eine feste Modellversion; // die CLI loest den Alias aufs aktuelle Modell des Tiers auf. Die id-Strings // muessen von openai-to-cli.js extractModel() erkannt werden (MODEL_MAP). // // Quelle: /shared/config/models.json — damit neue Tier-Namen oder angepasste // Beschreibungen eine reine DATEI-Aenderung sind (kein Code-Edit, kein Neubau, // kein Neustart: handleModels liest pro Request neu; einfach die Datei // bearbeiten und im Diagnostic „Aktualisieren" druecken). Fehlt/kaputt die // Datei, greifen die eingebauten Defaults; die Datei wird dann einmalig mit // diesen Defaults angelegt, damit es was zu editieren gibt. const MODELS_FILE = process.env.ARIA_MODELS_FILE || "/shared/config/models.json"; const DEFAULT_MODELS = [ { id: "claude-sonnet-4", tier: "sonnet", display_name: "Sonnet (aktuell: Sonnet 5)", description: "Schnell & gut — Standard fuer den Alltag." }, { id: "claude-opus-4", tier: "opus", display_name: "Opus (aktuell: Opus 4.8)", description: "Langsamer, aber am schlausten — fuer schwere/lange Aufgaben." }, { id: "claude-haiku-4", tier: "haiku", display_name: "Haiku (aktuell: Haiku 4.5)", description: "Sehr schnell & guenstig, kleinerer Kontext — fuer einfache Tasks." }, ]; function _loadModels() { try { const raw = fs.readFileSync(MODELS_FILE, "utf-8"); const arr = JSON.parse(raw); if (Array.isArray(arr) && arr.length && arr.every(m => m && typeof m.id === "string")) { return arr; } console.error("[aria-models] models.json ungueltig — nutze Defaults"); } catch (_) { // Datei fehlt (oder unlesbar) → Defaults + einmalig seeden zum Editieren try { fs.mkdirSync("/shared/config", { recursive: true }); if (!fs.existsSync(MODELS_FILE)) { fs.writeFileSync(MODELS_FILE, JSON.stringify(DEFAULT_MODELS, null, 2)); console.error("[aria-models] models.json mit Defaults angelegt:", MODELS_FILE); } } catch (e) { console.error("[aria-models] Seeden fehlgeschlagen:", e && e.message); } } return DEFAULT_MODELS; } export function handleModels(_req, res) { const created = Math.floor(Date.now() / 1000); const models = _loadModels(); res.json({ object: "list", data: models.map(m => ({ id: m.id, object: "model", owned_by: "anthropic", created, tier: m.tier || m.id, display_name: m.display_name || m.id, description: m.description || "", })), }); } /** * Handle GET /health * * Health check endpoint */ export function handleHealth(_req, res) { res.json({ status: "ok", provider: "claude-code-cli", timestamp: new Date().toISOString(), }); } // ── Not-Aus Side-Channel ─────────────────────────────────── // // claude-max-api-proxy steuert seine eigene Route-Registrierung — wir // koennen da nicht reinpatchen ohne sed-Operationen am npm-Paket. Saubrer: // ein dedizierter kleiner HTTP-Listener nur fuer den Not-Aus, auf einem // internen Port im aria-net. Bridge ruft den, killt alle aktiven Claude- // Subprocesses. App + Diagnostic sehen den Stream sofort enden. const INTERNAL_PORT = parseInt(process.env.ARIA_PROXY_INTERNAL_PORT || "3457", 10); const INTERNAL_HOST = "0.0.0.0"; // im aria-net erreichbar, nicht nach extern exposed function _cancelAll() { const ids = Array.from(_activeSubprocesses.keys()); let killed = 0; for (const [id, entry] of _activeSubprocesses) { try { entry.subprocess.kill(); killed++; } catch (e) { console.error("[aria-not-aus] kill failed for", id, e?.message); } } _activeSubprocesses.clear(); return { killed, requestIds: ids }; } // Kontext-scoped Cancel: killt NUR die Subprozesse eines Projekts (leer = // Hauptchat). Fuer Barge-In in einem Kontext ohne die parallele Arbeit in // anderen Kontexten abzuwuergen. function _cancelByProject(projectId) { const pid = String(projectId || ""); const ids = []; let killed = 0; for (const [id, entry] of Array.from(_activeSubprocesses)) { if (entry.projectId !== pid) continue; ids.push(id); try { entry.subprocess.kill(); killed++; } catch (e) { console.error("[aria-cancel] kill failed for", id, e?.message); } _activeSubprocesses.delete(id); } return { killed, requestIds: ids, projectId: pid }; } try { const internalServer = http.createServer((req, res) => { if (req.method === "POST" && req.url === "/cancel-all") { const result = _cancelAll(); console.warn("[aria-not-aus] /cancel-all — killed", result.killed, "subprocess(es)"); res.writeHead(200, { "Content-Type": "application/json" }); res.end(JSON.stringify({ ok: true, ...result })); return; } if (req.method === "POST" && req.url === "/cancel") { // Body: {projectId}. Kontext-scoped Barge-In — killt nur die // Subprozesse dieses Kontexts (leer = Hauptchat). let raw = ""; req.on("data", (c) => { raw += c; if (raw.length > 4096) req.destroy(); }); req.on("end", () => { let projectId = ""; try { projectId = String((JSON.parse(raw || "{}")).projectId || ""); } catch (_) {} const result = _cancelByProject(projectId); console.warn("[aria-cancel] /cancel project=%s — killed %d", projectId || "(main)", result.killed); res.writeHead(200, { "Content-Type": "application/json" }); res.end(JSON.stringify({ ok: true, ...result })); }); return; } if (req.method === "GET" && req.url === "/health") { res.writeHead(200, { "Content-Type": "application/json" }); res.end(JSON.stringify({ ok: true, active: _activeSubprocesses.size })); return; } res.writeHead(404).end(); }); internalServer.on("error", (err) => { console.error("[aria-not-aus] internal listener error:", err.message); }); internalServer.listen(INTERNAL_PORT, INTERNAL_HOST, () => { console.log("[aria-not-aus] internal listener on", INTERNAL_HOST + ":" + INTERNAL_PORT); }); } catch (e) { console.error("[aria-not-aus] startup failed:", e?.message); } //# sourceMappingURL=routes.js.map