diff --git a/README.md b/README.md index 99584e0..a76bd45 100644 --- a/README.md +++ b/README.md @@ -99,10 +99,10 @@ Oder manuell: ```bash # .txz von Gitea laden -curl -O https://git.casaderoll.de/michael/MUA-Mikes-Unraid-Agent/raw/branch/main/dist/mua-2026.08.29.r027-x86_64-1.txz +curl -O https://git.casaderoll.de/michael/MUA-Mikes-Unraid-Agent/raw/branch/main/dist/mua-2026.08.31.r028-x86_64-1.txz # Installieren -upgradepkg --install-new mua-2026.08.29.r027-x86_64-1.txz +upgradepkg --install-new mua-2026.08.31.r028-x86_64-1.txz ``` ### 2. Service starten diff --git a/dist/mua-2026.08.31.r028-x86_64-1.txz b/dist/mua-2026.08.31.r028-x86_64-1.txz new file mode 100644 index 0000000..7650463 Binary files /dev/null and b/dist/mua-2026.08.31.r028-x86_64-1.txz differ diff --git a/mcp/helpers.php b/mcp/helpers.php index e26f7c1..5684ffb 100644 --- a/mcp/helpers.php +++ b/mcp/helpers.php @@ -10,7 +10,7 @@ error_reporting(E_ALL); ini_set('display_errors', 0); // Fehler werden als JSON-RPC-Error zurückgegeben, nicht als HTML -const MUA_VERSION = '2026.08.29.r027'; +const MUA_VERSION = '2026.08.31.r028'; const MUA_SERVER_NAME = 'mua'; /** diff --git a/mcp/server.php b/mcp/server.php index 51e88cf..41cf33d 100644 --- a/mcp/server.php +++ b/mcp/server.php @@ -21,7 +21,7 @@ const MUA_PORT = 3002; const MUA_BIND = '0.0.0.0'; const MUA_SESSION_TTL = 3600; // Sekunden const MUA_SERVER_NAME = 'mua'; -const MUA_VERSION = '2026.08.29.r027'; +const MUA_VERSION = '2026.08.31.r028'; // ── Session-Management (in-memory) ─────────────────────────────────────── $sessions = []; // session_id => ['created' => time, 'initialized' => bool] diff --git a/package.json b/package.json index ed3f519..80ba94f 100644 --- a/package.json +++ b/package.json @@ -1,6 +1,6 @@ { "name": "mua", - "version": "2026.08.29.r027", + "version": "2026.08.31.r028", "description": "Mikes Unraid Agent - MCP over HTTP (Streamable HTTP) for Unraid", "type": "module", "main": "src/index.ts", diff --git a/plugin/mua.plg b/plugin/mua.plg index 7211687..52b1adb 100644 --- a/plugin/mua.plg +++ b/plugin/mua.plg @@ -2,13 +2,13 @@ - + - - + + ]> +### 2026.08.31.r028 +- Beendet festhängende Tool-Aufrufe mit einem echten harten Timeout, auch wenn Kindprozesse Ausgabe-Streams offen halten. +- Ersetzt den fehleranfälligen globalen Tool-Zähler durch idempotente, selbstheilende Leases; abgebrochene Aufrufe blockieren MUA nicht mehr dauerhaft mit `Server busy`. +- Der Health-Endpunkt zeigt aktive und maximal erlaubte Tool-Aufrufe sowie das Alter des ältesten Aufrufs an. +- Beendet Prozesse bei erreichtem Ausgabelimit sofort und liefert die begrenzte Ausgabe weiterhin korrekt zurück. + ### 2026.08.29.r027 - Stellt bei aktivierter Toolbox zwei echte MCP-Werkzeuge bereit: `unraid_toolbox_status` und `unraid_toolbox_exec`. - Agenten erkennen und verwenden Medienwerkzeuge direkt im dynamisch gelabelten Container, ohne lokales Docker oder Umwege über die Host-Shell. @@ -167,7 +173,7 @@ Das .txz enthält: install/doinst.sh (läuft nach Installation) =========================================== --> - + &txzURL; &txzSHA256; diff --git a/plugin/plugin.json b/plugin/plugin.json index faa07c0..d62e590 100644 --- a/plugin/plugin.json +++ b/plugin/plugin.json @@ -1,7 +1,7 @@ { "name": "mua", "author": "Michael", - "version": "2026.08.29.r027", + "version": "2026.08.31.r028", "minver": "7.0.0", "pluginDirectory": "/usr/local/emhttp/plugins/mua", "configDirectory": "/boot/config/plugins/mua", diff --git a/src/concurrency.ts b/src/concurrency.ts new file mode 100644 index 0000000..86f8544 --- /dev/null +++ b/src/concurrency.ts @@ -0,0 +1,90 @@ +import { randomUUID } from "node:crypto"; + +export interface ToolCallLease { + id: string; + tool: string; + startedAt: number; + expiresAt: number; +} + +export interface ToolCallGateSnapshot { + active: number; + maximum: number; + oldestAgeMs: number; +} + +/** + * Bounded, idempotent tool-call admission with stale-lease recovery. + * + * A Map is used instead of a numeric counter so that late or duplicate + * releases cannot decrement a newer call or make the counter negative. + */ +export class ToolCallGate { + private readonly leases = new Map(); + + constructor(private readonly maximum: number) { + if (!Number.isInteger(maximum) || maximum < 1) { + throw new Error("maximum concurrent tool calls must be a positive integer"); + } + } + + acquire(tool: string, lifetimeMs: number, now = Date.now()): ToolCallLease | null { + this.pruneExpired(now); + if (this.leases.size >= this.maximum) return null; + const lease: ToolCallLease = { + id: randomUUID(), + tool, + startedAt: now, + expiresAt: now + Math.max(1, lifetimeMs), + }; + this.leases.set(lease.id, lease); + return lease; + } + + release(id: string): void { + this.leases.delete(id); + } + + pruneExpired(now = Date.now()): number { + let removed = 0; + for (const [id, lease] of this.leases) { + if (lease.expiresAt <= now) { + this.leases.delete(id); + removed += 1; + } + } + return removed; + } + + snapshot(now = Date.now()): ToolCallGateSnapshot { + this.pruneExpired(now); + let oldestAgeMs = 0; + for (const lease of this.leases.values()) { + oldestAgeMs = Math.max(oldestAgeMs, now - lease.startedAt); + } + return { active: this.leases.size, maximum: this.maximum, oldestAgeMs }; + } +} + +/** Reject independently of the wrapped promise so a stuck stream cannot hold a slot forever. */ +export async function withHardTimeout( + work: Promise, + timeoutMs: number, + label: string, +): Promise { + let timer: ReturnType | undefined; + const timeout = new Promise((_, reject) => { + timer = setTimeout( + () => reject(new Error(`${label} exceeded hard timeout after ${timeoutMs} ms`)), + timeoutMs, + ); + }); + // The original work may settle after the timeout won the race. Mark a late + // rejection handled so it cannot become an unhandled rejection. + void work.catch(() => {}); + try { + return await Promise.race([work, timeout]); + } finally { + if (timer !== undefined) clearTimeout(timer); + } +} diff --git a/src/helpers.ts b/src/helpers.ts index c1b2c1f..ed4f0b6 100644 --- a/src/helpers.ts +++ b/src/helpers.ts @@ -15,7 +15,7 @@ import { existsSync, mkdirSync, readFileSync, rmSync, statSync, writeFileSync } // ── Konstanten ────────────────────────────────────────────────────────── export const MUA_SERVER_NAME = "mua"; -export const MUA_VERSION = "2026.08.29.r027"; +export const MUA_VERSION = "2026.08.31.r028"; export const MUA_PROTOCOL_VERSION = "2025-03-26"; export const PHP_HELPER = "/usr/local/bin/unraid-docker-mcp-helper.php"; export const STATUS_HELPER = "/usr/local/bin/unraid-mcp-status-helper.php"; @@ -27,6 +27,38 @@ export interface CmdResult { code: number; } +async function awaitProcess( + proc: { kill: () => void; exited: Promise }, + stdout: Promise, + stderr: Promise, + timeoutSec: number, + label: string, +): Promise<[TStdout, TStderr, number]> { + const completion = Promise.all([stdout, stderr, proc.exited]) as Promise< + [TStdout, TStderr, number] + >; + // A killed shell can leave descendants holding stdout/stderr open. Reject + // explicitly instead of waiting forever for those inherited pipes. + let timer: ReturnType | undefined; + const timeout = new Promise((_, reject) => { + timer = setTimeout(() => { + try { + proc.kill(); + } catch { + // The process may already have exited; timeout rejection still frees + // the caller and its concurrency lease. + } + reject(new Error(`${label} timed out after ${timeoutSec} seconds`)); + }, timeoutSec * 1000); + }); + void completion.catch(() => {}); + try { + return await Promise.race([completion, timeout]); + } finally { + if (timer !== undefined) clearTimeout(timer); + } +} + /** * Führe einen lokalen Shell-Befehl aus (via /bin/sh -c). * Timeout in Sekunden. @@ -41,13 +73,13 @@ export async function runLocal( stdout: "pipe", stderr: "pipe", }); - const timeout = setTimeout(() => proc.kill(), timeoutSec * 1000); - const [stdout, stderr, code] = await Promise.all([ + const [stdout, stderr, code] = await awaitProcess( + proc, new Response(proc.stdout).text(), new Response(proc.stderr).text(), - proc.exited, - ]); - clearTimeout(timeout); + timeoutSec, + label, + ); let result = stdout.trim(); if (result === "" && stderr.trim() !== "") result = stderr.trim(); return result; @@ -67,13 +99,13 @@ export async function runShell(cmd: string, timeoutSec = 60): Promise { stdout: "pipe", stderr: "pipe", }); - const timeout = setTimeout(() => proc.kill(), timeoutSec * 1000); - const [stdout, stderr, code] = await Promise.all([ + const [stdout, stderr, code] = await awaitProcess( + proc, new Response(proc.stdout).text(), new Response(proc.stderr).text(), - proc.exited, - ]); - clearTimeout(timeout); + timeoutSec, + "run_shell", + ); return JSON.stringify({ exit_code: code, stdout: sanitizeLogOutput(stdout.trim(), 100000), @@ -174,13 +206,17 @@ async function runArgv( ): Promise { try { const proc = spawn(argv, { stdout: "pipe", stderr: "pipe", cwd: "/" }); - const timeout = setTimeout(() => proc.kill(), timeoutSec * 1000); - const [stdout, stderr, code] = await Promise.all([ - readStreamLimited(proc.stdout, maxStdout), - readStreamLimited(proc.stderr, maxStderr), - proc.exited, - ]); - clearTimeout(timeout); + const stdoutPromise = readStreamLimited(proc.stdout, maxStdout); + const stderrPromise = readStreamLimited(proc.stderr, maxStderr); + stopProcessWhenTruncated(proc, stdoutPromise); + stopProcessWhenTruncated(proc, stderrPromise); + const [stdout, stderr, code] = await awaitProcess( + proc, + stdoutPromise, + stderrPromise, + timeoutSec, + label, + ); if (code !== 0) { throw new Error(`${label} failed (exit ${code}): ${sanitizeLogOutput(stderr.text || stdout.text, maxStderr)}`); } @@ -214,7 +250,17 @@ async function readStreamLimited( if (slice.length > 0) chunks.push(slice); kept += slice.length; } - if (value.length > available) truncated = true; + if (value.length > available) { + truncated = true; + // Do not keep draining an unbounded producer after the configured + // output budget is full. Closing the pipe also helps the producer exit. + try { + await reader.cancel(); + } catch { + // The process may close the stream concurrently. + } + break; + } } const combined = new Uint8Array(chunks.reduce((n, c) => n + c.length, 0)); let offset = 0; @@ -225,6 +271,20 @@ async function readStreamLimited( return { text: new TextDecoder().decode(combined), truncated }; } +function stopProcessWhenTruncated( + proc: { kill: () => void }, + output: Promise<{ text: string; truncated: boolean }>, +): void { + void output.then((result) => { + if (!result.truncated) return; + try { + proc.kill(); + } catch { + // The process may already have exited after its pipe was closed. + } + }).catch(() => {}); +} + /** * Führt ausschließlich freigegebene Leseprogramme direkt als argv aus. * Kein /bin/sh, keine Pipes, Umleitungen, Substitutionen oder Verkettungen. @@ -288,13 +348,17 @@ export async function runReadOnlyCommand( try { const proc = spawn([program, ...args], { stdout: "pipe", stderr: "pipe", cwd: "/" }); - const timeout = setTimeout(() => proc.kill(), timeoutSec * 1000); - const [stdout, stderr, code] = await Promise.all([ - readStreamLimited(proc.stdout, maxOutputChars), - readStreamLimited(proc.stderr, Math.min(8_000, maxOutputChars)), - proc.exited, - ]); - clearTimeout(timeout); + const stdoutPromise = readStreamLimited(proc.stdout, maxOutputChars); + const stderrPromise = readStreamLimited(proc.stderr, Math.min(8_000, maxOutputChars)); + stopProcessWhenTruncated(proc, stdoutPromise); + stopProcessWhenTruncated(proc, stderrPromise); + const [stdout, stderr, code] = await awaitProcess( + proc, + stdoutPromise, + stderrPromise, + timeoutSec, + `read-only ${program}`, + ); return JSON.stringify({ exit_code: code, stdout: sanitizeLogOutput(stdout.text.trim(), maxOutputChars), @@ -427,13 +491,17 @@ export async function toolboxExec( ["/usr/bin/docker", "exec", selected.name, "/bin/sh", "-lc", command], { stdout: "pipe", stderr: "pipe", cwd: "/" }, ); - const timeout = setTimeout(() => proc.kill(), timeoutSec * 1000); - const [stdout, stderr, code] = await Promise.all([ - readStreamLimited(proc.stdout, 100_000), - readStreamLimited(proc.stderr, 20_000), - proc.exited, - ]); - clearTimeout(timeout); + const stdoutPromise = readStreamLimited(proc.stdout, 100_000); + const stderrPromise = readStreamLimited(proc.stderr, 20_000); + stopProcessWhenTruncated(proc, stdoutPromise); + stopProcessWhenTruncated(proc, stderrPromise); + const [stdout, stderr, code] = await awaitProcess( + proc, + stdoutPromise, + stderrPromise, + timeoutSec, + `toolbox ${selected.name}`, + ); return JSON.stringify({ container: selected.name, exit_code: code, diff --git a/src/index.ts b/src/index.ts index 7330ff6..856b08a 100644 --- a/src/index.ts +++ b/src/index.ts @@ -34,12 +34,19 @@ import { import { randomBytes } from "node:crypto"; import { spawn } from "node:child_process"; import { appendFile, chmodSync, statSync, writeFileSync } from "node:fs"; +import { ToolCallGate, withHardTimeout } from "./concurrency"; const PORT = Number(process.env["MUA_PORT"] ?? 3002); const HOST = process.env["MUA_HOST"] ?? "0.0.0.0"; const MAX_REQUEST_BYTES = Number(process.env["MUA_MAX_REQUEST_BYTES"] ?? 1_048_576); const RATE_LIMIT_PER_MINUTE = Number(process.env["MUA_RATE_LIMIT_PER_MINUTE"] ?? 120); const MAX_CONCURRENT_TOOL_CALLS = Number(process.env["MUA_MAX_CONCURRENT_TOOL_CALLS"] ?? 4); +const DEFAULT_TOOL_CALL_TIMEOUT_SECONDS = Number( + process.env["MUA_DEFAULT_TOOL_CALL_TIMEOUT_SECONDS"] ?? 330, +); +const MAX_TOOL_CALL_TIMEOUT_SECONDS = Number( + process.env["MUA_MAX_TOOL_CALL_TIMEOUT_SECONDS"] ?? 1830, +); const CORS_ORIGIN = process.env["MUA_CORS_ORIGIN"] ?? ""; const AUDIT_LOG = process.env["MUA_AUDIT_LOG"] ?? "/var/log/plugins/mua-audit.log"; @@ -90,7 +97,15 @@ function rpcError(id: number | string | null, code: number, message: string) { // ── Session Management ────────────────────────────────────────────────── const sessions = new Map(); const rateLimits = new Map(); -let activeToolCalls = 0; +const toolCallGate = new ToolCallGate(MAX_CONCURRENT_TOOL_CALLS); + +function hardTimeoutMs(args: Record): number { + const requested = Number(args["timeout_seconds"]); + const seconds = Number.isFinite(requested) && requested > 0 + ? requested + 15 + : DEFAULT_TOOL_CALL_TIMEOUT_SECONDS; + return Math.max(1, Math.min(seconds, MAX_TOOL_CALL_TIMEOUT_SECONDS)) * 1000; +} function allowRequest(client: string): boolean { const now = Date.now(); @@ -205,7 +220,9 @@ async function handleMcpRequest( sessionId: sessionId ?? "", }; } - if (activeToolCalls >= MAX_CONCURRENT_TOOL_CALLS) { + const timeoutMs = hardTimeoutMs(args); + const lease = toolCallGate.acquire(toolName, timeoutMs + 5000); + if (!lease) { return { response: rpcResult(id, { content: [{ type: "text", text: "ERROR: Server busy; retry later" }], @@ -214,10 +231,9 @@ async function handleMcpRequest( sessionId: sessionId ?? "", }; } - activeToolCalls += 1; const startedAt = Date.now(); try { - const text = await tool.handler(args); + const text = await withHardTimeout(tool.handler(args), timeoutMs, `MCP tool '${toolName}'`); auditTool(toolName, true, Date.now() - startedAt); const result: Record = { content: [{ type: "text", text }], @@ -243,7 +259,7 @@ async function handleMcpRequest( sessionId: sessionId ?? "", }; } finally { - activeToolCalls -= 1; + toolCallGate.release(lease.id); } } @@ -294,6 +310,7 @@ const server = Bun.serve({ // ── Health-Check (ohne Auth) ──────────────────────────────────────── if (path === "/health") { + const toolCalls = toolCallGate.snapshot(); return Response.json( { status: "ok", @@ -301,6 +318,11 @@ const server = Bun.serve({ version: MUA_VERSION, port: PORT, auth: "required", + tool_calls: { + active: toolCalls.active, + maximum: toolCalls.maximum, + oldest_age_ms: toolCalls.oldestAgeMs, + }, time: new Date().toISOString(), }, { headers: corsHeaders }, diff --git a/src/security.test.ts b/src/security.test.ts index 4b06654..f9dc4bf 100644 --- a/src/security.test.ts +++ b/src/security.test.ts @@ -2,6 +2,7 @@ import { describe, expect, test } from "bun:test"; import { parseConfig, SAFE_DEFAULT_TOOLS } from "./auth"; import { installCommunityApp, runReadOnlyCommand, sanitizeLogOutput } from "./helpers"; import { getToolRisk, toolByName } from "./tools"; +import { ToolCallGate, withHardTimeout } from "./concurrency"; describe("secure tool configuration", () => { test("new and incomplete configs use the read-only baseline", () => { @@ -117,6 +118,37 @@ describe("direct Unraid terminal", () => { }); }); +describe("tool-call concurrency recovery", () => { + test("stale leases heal without corrupting newer leases", () => { + const gate = new ToolCallGate(1); + const first = gate.acquire("first", 10, 100); + expect(first).not.toBeNull(); + expect(gate.acquire("blocked", 10, 105)).toBeNull(); + + const second = gate.acquire("second", 100, 111); + expect(second).not.toBeNull(); + gate.release(first!.id); // late release must not affect the newer lease + expect(gate.snapshot(112).active).toBe(1); + gate.release(second!.id); + expect(gate.snapshot(113).active).toBe(0); + }); + + test("hard timeout rejects work that never settles", async () => { + const never = new Promise(() => {}); + await expect(withHardTimeout(never, 20, "stuck test")).rejects.toThrow( + "exceeded hard timeout", + ); + }); + + test("shell timeout rejects even when a descendant keeps pipes open", async () => { + const tool = toolByName("unraid_system_shell"); + expect(tool).toBeDefined(); + await expect( + tool!.handler({ command: "sleep 2 &", timeout_seconds: 1 }), + ).rejects.toThrow("timed out after 1 seconds"); + }); +}); + describe("read-only shell", () => { test("executes an allowlisted program without a shell", async () => { const result = JSON.parse(await runReadOnlyCommand("ls", ["-ld", "/"], 5));