r028: recover stuck tool call slots
This commit is contained in:
@@ -99,10 +99,10 @@ Oder manuell:
|
|||||||
|
|
||||||
```bash
|
```bash
|
||||||
# .txz von Gitea laden
|
# .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
|
# 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
|
### 2. Service starten
|
||||||
|
|||||||
BIN
Binary file not shown.
+1
-1
@@ -10,7 +10,7 @@
|
|||||||
error_reporting(E_ALL);
|
error_reporting(E_ALL);
|
||||||
ini_set('display_errors', 0); // Fehler werden als JSON-RPC-Error zurückgegeben, nicht als HTML
|
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';
|
const MUA_SERVER_NAME = 'mua';
|
||||||
|
|
||||||
/**
|
/**
|
||||||
|
|||||||
+1
-1
@@ -21,7 +21,7 @@ const MUA_PORT = 3002;
|
|||||||
const MUA_BIND = '0.0.0.0';
|
const MUA_BIND = '0.0.0.0';
|
||||||
const MUA_SESSION_TTL = 3600; // Sekunden
|
const MUA_SESSION_TTL = 3600; // Sekunden
|
||||||
const MUA_SERVER_NAME = 'mua';
|
const MUA_SERVER_NAME = 'mua';
|
||||||
const MUA_VERSION = '2026.08.29.r027';
|
const MUA_VERSION = '2026.08.31.r028';
|
||||||
|
|
||||||
// ── Session-Management (in-memory) ───────────────────────────────────────
|
// ── Session-Management (in-memory) ───────────────────────────────────────
|
||||||
$sessions = []; // session_id => ['created' => time, 'initialized' => bool]
|
$sessions = []; // session_id => ['created' => time, 'initialized' => bool]
|
||||||
|
|||||||
+1
-1
@@ -1,6 +1,6 @@
|
|||||||
{
|
{
|
||||||
"name": "mua",
|
"name": "mua",
|
||||||
"version": "2026.08.29.r027",
|
"version": "2026.08.31.r028",
|
||||||
"description": "Mikes Unraid Agent - MCP over HTTP (Streamable HTTP) for Unraid",
|
"description": "Mikes Unraid Agent - MCP over HTTP (Streamable HTTP) for Unraid",
|
||||||
"type": "module",
|
"type": "module",
|
||||||
"main": "src/index.ts",
|
"main": "src/index.ts",
|
||||||
|
|||||||
+10
-4
@@ -2,13 +2,13 @@
|
|||||||
<!DOCTYPE PLUGIN [
|
<!DOCTYPE PLUGIN [
|
||||||
<!ENTITY name "mua">
|
<!ENTITY name "mua">
|
||||||
<!ENTITY author "Michael">
|
<!ENTITY author "Michael">
|
||||||
<!ENTITY version "2026.08.29.r027">
|
<!ENTITY version "2026.08.31.r028">
|
||||||
<!ENTITY launch "Settings/mua">
|
<!ENTITY launch "Settings/mua">
|
||||||
<!ENTITY pluginURL "https://git.casaderoll.de/michael/MUA-Mikes-Unraid-Agent/raw/branch/main/plugin/mua.plg">
|
<!ENTITY pluginURL "https://git.casaderoll.de/michael/MUA-Mikes-Unraid-Agent/raw/branch/main/plugin/mua.plg">
|
||||||
<!ENTITY pluginLOC "/boot/config/plugins/&name;">
|
<!ENTITY pluginLOC "/boot/config/plugins/&name;">
|
||||||
<!ENTITY emhttpLOC "/usr/local/emhttp/plugins/&name;">
|
<!ENTITY emhttpLOC "/usr/local/emhttp/plugins/&name;">
|
||||||
<!ENTITY txzURL "https://git.casaderoll.de/michael/MUA-Mikes-Unraid-Agent/raw/branch/main/dist/mua-2026.08.29.r027-x86_64-1.txz">
|
<!ENTITY txzURL "https://git.casaderoll.de/michael/MUA-Mikes-Unraid-Agent/raw/branch/main/dist/mua-2026.08.31.r028-x86_64-1.txz">
|
||||||
<!ENTITY txzSHA256 "f7fcbc3b3de260f6df0d413ac072d0f3e674bc25425e6a50b8a11b78a4e5007c">
|
<!ENTITY txzSHA256 "ce5c260f55b9ac01a453e79068e192d14c832b368b048ff922931e8995855046">
|
||||||
]>
|
]>
|
||||||
|
|
||||||
<PLUGIN name="&name;"
|
<PLUGIN name="&name;"
|
||||||
@@ -23,6 +23,12 @@
|
|||||||
>
|
>
|
||||||
|
|
||||||
<CHANGES>
|
<CHANGES>
|
||||||
|
### 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
|
### 2026.08.29.r027
|
||||||
- Stellt bei aktivierter Toolbox zwei echte MCP-Werkzeuge bereit: `unraid_toolbox_status` und `unraid_toolbox_exec`.
|
- 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.
|
- 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)
|
install/doinst.sh (läuft nach Installation)
|
||||||
===========================================
|
===========================================
|
||||||
-->
|
-->
|
||||||
<FILE Name="/boot/config/plugins/&name;/mua-2026.08.29.r027-x86_64-1.txz" Run="upgradepkg --install-new" Mode="755" Min="7.0.0">
|
<FILE Name="/boot/config/plugins/&name;/mua-2026.08.31.r028-x86_64-1.txz" Run="upgradepkg --install-new" Mode="755" Min="7.0.0">
|
||||||
<URL>&txzURL;</URL>
|
<URL>&txzURL;</URL>
|
||||||
<SHA256>&txzSHA256;</SHA256>
|
<SHA256>&txzSHA256;</SHA256>
|
||||||
</FILE>
|
</FILE>
|
||||||
|
|||||||
+1
-1
@@ -1,7 +1,7 @@
|
|||||||
{
|
{
|
||||||
"name": "mua",
|
"name": "mua",
|
||||||
"author": "Michael",
|
"author": "Michael",
|
||||||
"version": "2026.08.29.r027",
|
"version": "2026.08.31.r028",
|
||||||
"minver": "7.0.0",
|
"minver": "7.0.0",
|
||||||
"pluginDirectory": "/usr/local/emhttp/plugins/mua",
|
"pluginDirectory": "/usr/local/emhttp/plugins/mua",
|
||||||
"configDirectory": "/boot/config/plugins/mua",
|
"configDirectory": "/boot/config/plugins/mua",
|
||||||
|
|||||||
@@ -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<string, ToolCallLease>();
|
||||||
|
|
||||||
|
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<T>(
|
||||||
|
work: Promise<T>,
|
||||||
|
timeoutMs: number,
|
||||||
|
label: string,
|
||||||
|
): Promise<T> {
|
||||||
|
let timer: ReturnType<typeof setTimeout> | undefined;
|
||||||
|
const timeout = new Promise<never>((_, 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);
|
||||||
|
}
|
||||||
|
}
|
||||||
+101
-33
@@ -15,7 +15,7 @@ import { existsSync, mkdirSync, readFileSync, rmSync, statSync, writeFileSync }
|
|||||||
|
|
||||||
// ── Konstanten ──────────────────────────────────────────────────────────
|
// ── Konstanten ──────────────────────────────────────────────────────────
|
||||||
export const MUA_SERVER_NAME = "mua";
|
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 MUA_PROTOCOL_VERSION = "2025-03-26";
|
||||||
export const PHP_HELPER = "/usr/local/bin/unraid-docker-mcp-helper.php";
|
export const PHP_HELPER = "/usr/local/bin/unraid-docker-mcp-helper.php";
|
||||||
export const STATUS_HELPER = "/usr/local/bin/unraid-mcp-status-helper.php";
|
export const STATUS_HELPER = "/usr/local/bin/unraid-mcp-status-helper.php";
|
||||||
@@ -27,6 +27,38 @@ export interface CmdResult {
|
|||||||
code: number;
|
code: number;
|
||||||
}
|
}
|
||||||
|
|
||||||
|
async function awaitProcess<TStdout, TStderr>(
|
||||||
|
proc: { kill: () => void; exited: Promise<number> },
|
||||||
|
stdout: Promise<TStdout>,
|
||||||
|
stderr: Promise<TStderr>,
|
||||||
|
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<typeof setTimeout> | undefined;
|
||||||
|
const timeout = new Promise<never>((_, 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).
|
* Führe einen lokalen Shell-Befehl aus (via /bin/sh -c).
|
||||||
* Timeout in Sekunden.
|
* Timeout in Sekunden.
|
||||||
@@ -41,13 +73,13 @@ export async function runLocal(
|
|||||||
stdout: "pipe",
|
stdout: "pipe",
|
||||||
stderr: "pipe",
|
stderr: "pipe",
|
||||||
});
|
});
|
||||||
const timeout = setTimeout(() => proc.kill(), timeoutSec * 1000);
|
const [stdout, stderr, code] = await awaitProcess(
|
||||||
const [stdout, stderr, code] = await Promise.all([
|
proc,
|
||||||
new Response(proc.stdout).text(),
|
new Response(proc.stdout).text(),
|
||||||
new Response(proc.stderr).text(),
|
new Response(proc.stderr).text(),
|
||||||
proc.exited,
|
timeoutSec,
|
||||||
]);
|
label,
|
||||||
clearTimeout(timeout);
|
);
|
||||||
let result = stdout.trim();
|
let result = stdout.trim();
|
||||||
if (result === "" && stderr.trim() !== "") result = stderr.trim();
|
if (result === "" && stderr.trim() !== "") result = stderr.trim();
|
||||||
return result;
|
return result;
|
||||||
@@ -67,13 +99,13 @@ export async function runShell(cmd: string, timeoutSec = 60): Promise<string> {
|
|||||||
stdout: "pipe",
|
stdout: "pipe",
|
||||||
stderr: "pipe",
|
stderr: "pipe",
|
||||||
});
|
});
|
||||||
const timeout = setTimeout(() => proc.kill(), timeoutSec * 1000);
|
const [stdout, stderr, code] = await awaitProcess(
|
||||||
const [stdout, stderr, code] = await Promise.all([
|
proc,
|
||||||
new Response(proc.stdout).text(),
|
new Response(proc.stdout).text(),
|
||||||
new Response(proc.stderr).text(),
|
new Response(proc.stderr).text(),
|
||||||
proc.exited,
|
timeoutSec,
|
||||||
]);
|
"run_shell",
|
||||||
clearTimeout(timeout);
|
);
|
||||||
return JSON.stringify({
|
return JSON.stringify({
|
||||||
exit_code: code,
|
exit_code: code,
|
||||||
stdout: sanitizeLogOutput(stdout.trim(), 100000),
|
stdout: sanitizeLogOutput(stdout.trim(), 100000),
|
||||||
@@ -174,13 +206,17 @@ async function runArgv(
|
|||||||
): Promise<string> {
|
): Promise<string> {
|
||||||
try {
|
try {
|
||||||
const proc = spawn(argv, { stdout: "pipe", stderr: "pipe", cwd: "/" });
|
const proc = spawn(argv, { stdout: "pipe", stderr: "pipe", cwd: "/" });
|
||||||
const timeout = setTimeout(() => proc.kill(), timeoutSec * 1000);
|
const stdoutPromise = readStreamLimited(proc.stdout, maxStdout);
|
||||||
const [stdout, stderr, code] = await Promise.all([
|
const stderrPromise = readStreamLimited(proc.stderr, maxStderr);
|
||||||
readStreamLimited(proc.stdout, maxStdout),
|
stopProcessWhenTruncated(proc, stdoutPromise);
|
||||||
readStreamLimited(proc.stderr, maxStderr),
|
stopProcessWhenTruncated(proc, stderrPromise);
|
||||||
proc.exited,
|
const [stdout, stderr, code] = await awaitProcess(
|
||||||
]);
|
proc,
|
||||||
clearTimeout(timeout);
|
stdoutPromise,
|
||||||
|
stderrPromise,
|
||||||
|
timeoutSec,
|
||||||
|
label,
|
||||||
|
);
|
||||||
if (code !== 0) {
|
if (code !== 0) {
|
||||||
throw new Error(`${label} failed (exit ${code}): ${sanitizeLogOutput(stderr.text || stdout.text, maxStderr)}`);
|
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);
|
if (slice.length > 0) chunks.push(slice);
|
||||||
kept += slice.length;
|
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));
|
const combined = new Uint8Array(chunks.reduce((n, c) => n + c.length, 0));
|
||||||
let offset = 0;
|
let offset = 0;
|
||||||
@@ -225,6 +271,20 @@ async function readStreamLimited(
|
|||||||
return { text: new TextDecoder().decode(combined), truncated };
|
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.
|
* Führt ausschließlich freigegebene Leseprogramme direkt als argv aus.
|
||||||
* Kein /bin/sh, keine Pipes, Umleitungen, Substitutionen oder Verkettungen.
|
* Kein /bin/sh, keine Pipes, Umleitungen, Substitutionen oder Verkettungen.
|
||||||
@@ -288,13 +348,17 @@ export async function runReadOnlyCommand(
|
|||||||
|
|
||||||
try {
|
try {
|
||||||
const proc = spawn([program, ...args], { stdout: "pipe", stderr: "pipe", cwd: "/" });
|
const proc = spawn([program, ...args], { stdout: "pipe", stderr: "pipe", cwd: "/" });
|
||||||
const timeout = setTimeout(() => proc.kill(), timeoutSec * 1000);
|
const stdoutPromise = readStreamLimited(proc.stdout, maxOutputChars);
|
||||||
const [stdout, stderr, code] = await Promise.all([
|
const stderrPromise = readStreamLimited(proc.stderr, Math.min(8_000, maxOutputChars));
|
||||||
readStreamLimited(proc.stdout, maxOutputChars),
|
stopProcessWhenTruncated(proc, stdoutPromise);
|
||||||
readStreamLimited(proc.stderr, Math.min(8_000, maxOutputChars)),
|
stopProcessWhenTruncated(proc, stderrPromise);
|
||||||
proc.exited,
|
const [stdout, stderr, code] = await awaitProcess(
|
||||||
]);
|
proc,
|
||||||
clearTimeout(timeout);
|
stdoutPromise,
|
||||||
|
stderrPromise,
|
||||||
|
timeoutSec,
|
||||||
|
`read-only ${program}`,
|
||||||
|
);
|
||||||
return JSON.stringify({
|
return JSON.stringify({
|
||||||
exit_code: code,
|
exit_code: code,
|
||||||
stdout: sanitizeLogOutput(stdout.text.trim(), maxOutputChars),
|
stdout: sanitizeLogOutput(stdout.text.trim(), maxOutputChars),
|
||||||
@@ -427,13 +491,17 @@ export async function toolboxExec(
|
|||||||
["/usr/bin/docker", "exec", selected.name, "/bin/sh", "-lc", command],
|
["/usr/bin/docker", "exec", selected.name, "/bin/sh", "-lc", command],
|
||||||
{ stdout: "pipe", stderr: "pipe", cwd: "/" },
|
{ stdout: "pipe", stderr: "pipe", cwd: "/" },
|
||||||
);
|
);
|
||||||
const timeout = setTimeout(() => proc.kill(), timeoutSec * 1000);
|
const stdoutPromise = readStreamLimited(proc.stdout, 100_000);
|
||||||
const [stdout, stderr, code] = await Promise.all([
|
const stderrPromise = readStreamLimited(proc.stderr, 20_000);
|
||||||
readStreamLimited(proc.stdout, 100_000),
|
stopProcessWhenTruncated(proc, stdoutPromise);
|
||||||
readStreamLimited(proc.stderr, 20_000),
|
stopProcessWhenTruncated(proc, stderrPromise);
|
||||||
proc.exited,
|
const [stdout, stderr, code] = await awaitProcess(
|
||||||
]);
|
proc,
|
||||||
clearTimeout(timeout);
|
stdoutPromise,
|
||||||
|
stderrPromise,
|
||||||
|
timeoutSec,
|
||||||
|
`toolbox ${selected.name}`,
|
||||||
|
);
|
||||||
return JSON.stringify({
|
return JSON.stringify({
|
||||||
container: selected.name,
|
container: selected.name,
|
||||||
exit_code: code,
|
exit_code: code,
|
||||||
|
|||||||
+27
-5
@@ -34,12 +34,19 @@ import {
|
|||||||
import { randomBytes } from "node:crypto";
|
import { randomBytes } from "node:crypto";
|
||||||
import { spawn } from "node:child_process";
|
import { spawn } from "node:child_process";
|
||||||
import { appendFile, chmodSync, statSync, writeFileSync } from "node:fs";
|
import { appendFile, chmodSync, statSync, writeFileSync } from "node:fs";
|
||||||
|
import { ToolCallGate, withHardTimeout } from "./concurrency";
|
||||||
|
|
||||||
const PORT = Number(process.env["MUA_PORT"] ?? 3002);
|
const PORT = Number(process.env["MUA_PORT"] ?? 3002);
|
||||||
const HOST = process.env["MUA_HOST"] ?? "0.0.0.0";
|
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 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 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 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 CORS_ORIGIN = process.env["MUA_CORS_ORIGIN"] ?? "";
|
||||||
const AUDIT_LOG = process.env["MUA_AUDIT_LOG"] ?? "/var/log/plugins/mua-audit.log";
|
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 ──────────────────────────────────────────────────
|
// ── Session Management ──────────────────────────────────────────────────
|
||||||
const sessions = new Map<string, { createdAt: number; lastActivity: number }>();
|
const sessions = new Map<string, { createdAt: number; lastActivity: number }>();
|
||||||
const rateLimits = new Map<string, { startedAt: number; count: number }>();
|
const rateLimits = new Map<string, { startedAt: number; count: number }>();
|
||||||
let activeToolCalls = 0;
|
const toolCallGate = new ToolCallGate(MAX_CONCURRENT_TOOL_CALLS);
|
||||||
|
|
||||||
|
function hardTimeoutMs(args: Record<string, unknown>): 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 {
|
function allowRequest(client: string): boolean {
|
||||||
const now = Date.now();
|
const now = Date.now();
|
||||||
@@ -205,7 +220,9 @@ async function handleMcpRequest(
|
|||||||
sessionId: sessionId ?? "",
|
sessionId: sessionId ?? "",
|
||||||
};
|
};
|
||||||
}
|
}
|
||||||
if (activeToolCalls >= MAX_CONCURRENT_TOOL_CALLS) {
|
const timeoutMs = hardTimeoutMs(args);
|
||||||
|
const lease = toolCallGate.acquire(toolName, timeoutMs + 5000);
|
||||||
|
if (!lease) {
|
||||||
return {
|
return {
|
||||||
response: rpcResult(id, {
|
response: rpcResult(id, {
|
||||||
content: [{ type: "text", text: "ERROR: Server busy; retry later" }],
|
content: [{ type: "text", text: "ERROR: Server busy; retry later" }],
|
||||||
@@ -214,10 +231,9 @@ async function handleMcpRequest(
|
|||||||
sessionId: sessionId ?? "",
|
sessionId: sessionId ?? "",
|
||||||
};
|
};
|
||||||
}
|
}
|
||||||
activeToolCalls += 1;
|
|
||||||
const startedAt = Date.now();
|
const startedAt = Date.now();
|
||||||
try {
|
try {
|
||||||
const text = await tool.handler(args);
|
const text = await withHardTimeout(tool.handler(args), timeoutMs, `MCP tool '${toolName}'`);
|
||||||
auditTool(toolName, true, Date.now() - startedAt);
|
auditTool(toolName, true, Date.now() - startedAt);
|
||||||
const result: Record<string, unknown> = {
|
const result: Record<string, unknown> = {
|
||||||
content: [{ type: "text", text }],
|
content: [{ type: "text", text }],
|
||||||
@@ -243,7 +259,7 @@ async function handleMcpRequest(
|
|||||||
sessionId: sessionId ?? "",
|
sessionId: sessionId ?? "",
|
||||||
};
|
};
|
||||||
} finally {
|
} finally {
|
||||||
activeToolCalls -= 1;
|
toolCallGate.release(lease.id);
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -294,6 +310,7 @@ const server = Bun.serve({
|
|||||||
|
|
||||||
// ── Health-Check (ohne Auth) ────────────────────────────────────────
|
// ── Health-Check (ohne Auth) ────────────────────────────────────────
|
||||||
if (path === "/health") {
|
if (path === "/health") {
|
||||||
|
const toolCalls = toolCallGate.snapshot();
|
||||||
return Response.json(
|
return Response.json(
|
||||||
{
|
{
|
||||||
status: "ok",
|
status: "ok",
|
||||||
@@ -301,6 +318,11 @@ const server = Bun.serve({
|
|||||||
version: MUA_VERSION,
|
version: MUA_VERSION,
|
||||||
port: PORT,
|
port: PORT,
|
||||||
auth: "required",
|
auth: "required",
|
||||||
|
tool_calls: {
|
||||||
|
active: toolCalls.active,
|
||||||
|
maximum: toolCalls.maximum,
|
||||||
|
oldest_age_ms: toolCalls.oldestAgeMs,
|
||||||
|
},
|
||||||
time: new Date().toISOString(),
|
time: new Date().toISOString(),
|
||||||
},
|
},
|
||||||
{ headers: corsHeaders },
|
{ headers: corsHeaders },
|
||||||
|
|||||||
@@ -2,6 +2,7 @@ import { describe, expect, test } from "bun:test";
|
|||||||
import { parseConfig, SAFE_DEFAULT_TOOLS } from "./auth";
|
import { parseConfig, SAFE_DEFAULT_TOOLS } from "./auth";
|
||||||
import { installCommunityApp, runReadOnlyCommand, sanitizeLogOutput } from "./helpers";
|
import { installCommunityApp, runReadOnlyCommand, sanitizeLogOutput } from "./helpers";
|
||||||
import { getToolRisk, toolByName } from "./tools";
|
import { getToolRisk, toolByName } from "./tools";
|
||||||
|
import { ToolCallGate, withHardTimeout } from "./concurrency";
|
||||||
|
|
||||||
describe("secure tool configuration", () => {
|
describe("secure tool configuration", () => {
|
||||||
test("new and incomplete configs use the read-only baseline", () => {
|
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<string>(() => {});
|
||||||
|
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", () => {
|
describe("read-only shell", () => {
|
||||||
test("executes an allowlisted program without a shell", async () => {
|
test("executes an allowlisted program without a shell", async () => {
|
||||||
const result = JSON.parse(await runReadOnlyCommand("ls", ["-ld", "/"], 5));
|
const result = JSON.parse(await runReadOnlyCommand("ls", ["-ld", "/"], 5));
|
||||||
|
|||||||
Reference in New Issue
Block a user