- Task 5: net.close on every path out of a serve iteration (400
included) and the listener on stop; measured 4 -> 54 fds over 50
requests before, 4 -> 4 over 200 after. The loop's comment claimed
the iteration-end drop IS the close -- wrong twice (net.Conn is a
scalar, and a drop would not close an fd); it now says what is true
- Task 6: LW_SOAK=<seconds> in the acceptance script -- each mode under
load, resident+descriptor deltas against a WARMED baseline (warm-up
includes load: cold-to-high-water is not growth), 256 KiB / zero
tolerance; LW_ACCEPT_WOVM soaks another build
- the soak caught ~1.6 MiB/min of in-arena leaks ASan cannot see (the
arena is one allocation to LeakSanitizer); an arena size-class
census + pointer trace attributed five bugs:
- jparse_string sized every decoded string at "rest of the input"
and relabeled len after -- blocks filed on free lists their next
allocation never reads (fs.read_all's mis-size, again); copy out
exact, free at the taken size
- `!=` never dropped fresh operands (headers["authorization"] !=
"Bearer ${key}" leaked both sides per request); Ne now reaps as Eq
- an Int interpolation segment is a fresh int_to_text, not a borrow;
is_borrowed_value_t asks the segment's type
- json.encode(Ctor{...}) had no owner -- record + both field copies
leaked per tool call; its bespoke lowering now drops the argument
- a discarded expression statement owns its result: `pop(lines);`
leaked the popped element; reader builtins excluded
- after: arena live bytes flat per request on every handler; release
soak 30 s per mode watch 0 / run 0 / mcp +20 KiB, descriptors flat;
ASan build flat at 14600 KiB across 601686 requests in 90 s past its
~1200-request quarantine warm-up
- gates: oop-accept ALL CRITERIA MET, oop-e2e 71/0, woc-test 565/0,
wovm-test green, log-watcher 7/0 (soak opt-in, fast path <1 min)
Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
355 lines
15 KiB
Text
355 lines
15 KiB
Text
-- mcp.wo — Mcp.hx + Tools.hx, file for file. MCP over Streamable HTTP,
|
|
-- hand-rolled subset: stateless, tools-only, no SSE, no sessions. handle()
|
|
-- is a pure function so the whole protocol is testable without sockets —
|
|
-- the original's best design decision, preserved (spec Part 4 requires it).
|
|
--
|
|
-- MiniLog.hx and its two tools (load_log_db / query_log_db) are NOT here:
|
|
-- the systems-track spec names "sqlite-equivalent embedded SQL over RAM"
|
|
-- out of scope — that is the DB engine's job, and a future sample wires
|
|
-- MiniLog's idea to `select`.
|
|
|
|
use fs
|
|
use net
|
|
use time
|
|
use env
|
|
use json
|
|
|
|
typedef HttpReq = { method: Text, path: Text, headers: map<Text, Text>, body: Text }
|
|
typedef HttpResp = { status: Int, body: Text }
|
|
typedef ToolOut = { is_error: Bool, text: Text }
|
|
|
|
-- The JSON-RPC request, decoded typed (Haxe poked a Dynamic). `id` is the
|
|
-- one field typing cannot pin: number | string | null, echoed back verbatim
|
|
-- — the sample forces json.Value, an opaque re-encodable value.
|
|
typedef RpcReq = { ?id: json.Value, ?method: Text, ?params: RpcParams }
|
|
typedef RpcParams = { ?name: Text, ?arguments: RpcArgs }
|
|
typedef RpcArgs = { ?path: Text, ?lines: Int, ?pattern: Text, ?maxMatches: Int }
|
|
|
|
typedef ErrBody = { error: Text }
|
|
typedef RpcErr = { code: Int, message: Text }
|
|
typedef ToolText = { type: Text, text: Text }
|
|
|
|
-- Wire-format rows (camelCase = the JSON the original emitted).
|
|
typedef CronRow = { schedule: Text, command: Text, logPath: Text, nextFire: ?Text, running: Bool }
|
|
typedef LogRow = { path: Text, sizeBytes: ?Int, mtime: ?Int, source: Text }
|
|
typedef Match = { byteOffset: Int, line: Text }
|
|
typedef SearchResult = { matches: multi Match, searchedBytes: Int, note: Text }
|
|
|
|
class Mcp {
|
|
static const PROTOCOL = "2025-03-26"
|
|
static const BODY_MAX = 65536
|
|
|
|
tools: Tools
|
|
api_key: Text
|
|
|
|
-- Mcp.hx:21-46.
|
|
fn handle(req: HttpReq) -> HttpResp {
|
|
if req.headers["authorization"] != "Bearer ${self.api_key}" {
|
|
return HttpResp { status: 401, body: json.encode(ErrBody { error: "unauthorized" }) };
|
|
}
|
|
if req.method != "POST" {
|
|
return HttpResp { status: 405, body: json.encode(ErrBody { error: "POST only" }) };
|
|
}
|
|
if req.path != "/mcp" {
|
|
return HttpResp { status: 404, body: json.encode(ErrBody { error: "not found" }) };
|
|
}
|
|
if len(req.body) > BODY_MAX {
|
|
return HttpResp { status: 413, body: json.encode(ErrBody { error: "body too large" }) };
|
|
}
|
|
|
|
let j = json.decode(req.body) as RpcReq; -- checked decode: ?RpcReq, never a trap
|
|
if j == nil { return rpc_error(nil, -32700, "parse error"); }
|
|
if j.method == nil { return rpc_error(j.id, -32600, "invalid request: no method"); }
|
|
if j.id == nil { return HttpResp { status: 202, body: "" }; } -- notification
|
|
|
|
switch j.method {
|
|
case "initialize":
|
|
return rpc_result(j.id, "{\"protocolVersion\":\"${PROTOCOL}\",\"capabilities\":{\"tools\":{}},\"serverInfo\":{\"name\":\"log-watcher\",\"version\":\"0.1\"}}");
|
|
case "ping":
|
|
return rpc_result(j.id, "{}");
|
|
case "tools/list":
|
|
return rpc_result(j.id, "{\"tools\":${TOOL_SCHEMAS}}");
|
|
case "tools/call":
|
|
return self.call_tool(j.id, j.params);
|
|
default:
|
|
return rpc_error(j.id, -32601, "unknown method: ${j.method}");
|
|
}
|
|
}
|
|
|
|
-- Mcp.hx:54-84. Tool-layer failures return isError, never protocol
|
|
-- errors; a trap inside a tool is caught at this boundary.
|
|
fn call_tool(id: json.Value, params: ?RpcParams) -> HttpResp {
|
|
if params == nil { return rpc_error(id, -32600, "tools/call: missing params.name"); }
|
|
if params.name == nil { return rpc_error(id, -32600, "tools/call: missing params.name"); }
|
|
let o = try self.dispatch(params.name, params.arguments)
|
|
catch (e) err("tool failed: ${e.msg}");
|
|
let content = json.encode(ToolText { type: "text", text: o.text });
|
|
let flag = "false";
|
|
if o.is_error { flag = "true"; }
|
|
return rpc_result(id, "{\"content\":[${content}],\"isError\":${flag}}");
|
|
}
|
|
|
|
fn dispatch(name: Text, args: ?RpcArgs) -> ToolOut {
|
|
switch name {
|
|
case "get_running_crons":
|
|
return self.tools.get_running_crons(time.now());
|
|
case "list_logs":
|
|
return self.tools.list_logs();
|
|
case "tail_log":
|
|
if args == nil { return err("tail_log: path is required"); }
|
|
if args.path == nil { return err("tail_log: path is required"); }
|
|
return self.tools.tail_log(args.path, args.lines);
|
|
case "search_log":
|
|
if args == nil { return err("search_log: path and pattern are required"); }
|
|
if args.path == nil or args.pattern == nil { return err("search_log: path and pattern are required"); }
|
|
return self.tools.search_log(args.path, args.pattern, args.maxMatches);
|
|
default:
|
|
return err("unknown tool: ${name}");
|
|
}
|
|
}
|
|
|
|
-- Blocking accept loop: one HTTP request per connection, respond, close.
|
|
-- Localhost only. A bad request never kills the loop. A connection is a
|
|
-- DESCRIPTOR, not an owned heap value: `net.Conn` is a scalar to the
|
|
-- ownership pass (nothing to drop, and dropping an fd number would be
|
|
-- nonsense), so the close is the source's own job on every path out of an
|
|
-- iteration — the 400 included. Measured before this was here: exactly one
|
|
-- descriptor leaked per request, so the server died at the process limit.
|
|
fn serve(port: Int) -> Int {
|
|
let srv = net.listen("127.0.0.1", port);
|
|
while true {
|
|
if env.stopping() { net.close(srv); return 0; }
|
|
let c = net.accept(srv);
|
|
let req = try self.read_request(c) catch (e) nil;
|
|
if req == nil {
|
|
try net.write(c, "HTTP/1.1 400 Bad Request\r\nContent-Length: 0\r\nConnection: close\r\n\r\n") catch (e) {}
|
|
net.close(c);
|
|
continue;
|
|
}
|
|
let resp = self.handle(req);
|
|
let head = "HTTP/1.1 ${resp.status} ${status_text(resp.status)}\r\n";
|
|
head = head .. "Content-Type: application/json\r\nConnection: close\r\n";
|
|
head = head .. "Content-Length: ${len(resp.body)}\r\n\r\n";
|
|
try net.write(c, head .. resp.body) catch (e) {}
|
|
net.close(c);
|
|
}
|
|
}
|
|
|
|
-- Mcp.hx:109-125 read line-by-line off the socket; `net` has only bounded
|
|
-- reads, so: buffer to the header terminator, then the body to
|
|
-- Content-Length (capped just past BODY_MAX so the 413 still fires).
|
|
fn read_request(c: net.Conn) -> ?HttpReq {
|
|
let buf = "";
|
|
let header_end = -1;
|
|
while header_end == -1 {
|
|
let got = net.read(c, 8192);
|
|
if len(got) == 0 { return nil; } -- peer closed mid-headers
|
|
buf = buf .. got;
|
|
header_end = index_of(buf, "\r\n\r\n");
|
|
if len(buf) > BODY_MAX * 2 { return nil; } -- runaway header block
|
|
}
|
|
let lines = split(substr(buf, 0, header_end), "\r\n");
|
|
let req_line = split_ws(trim(lines[0]));
|
|
let headers: map<Text, Text> = {};
|
|
let i = 1;
|
|
while i < len(lines) {
|
|
let line = trim(lines[i]);
|
|
i = i + 1;
|
|
if line == "" { continue; }
|
|
let colon = index_of(line, ":");
|
|
if colon > 0 {
|
|
headers[to_lower(substr(line, 0, colon))] = trim(substr(line, colon + 1, len(line) - colon - 1));
|
|
}
|
|
}
|
|
let want = 0;
|
|
let cl = headers["content-length"];
|
|
if cl != nil {
|
|
let n = parse_int(cl);
|
|
if n != nil { want = n; }
|
|
if want < 0 { want = 0; }
|
|
}
|
|
if want > BODY_MAX + 1 { want = BODY_MAX + 1; } -- read enough to trigger 413, no more
|
|
let body = substr(buf, header_end + 4, len(buf) - header_end - 4);
|
|
while len(body) < want {
|
|
let got = net.read(c, want - len(body));
|
|
if len(got) == 0 { break; }
|
|
body = body .. got;
|
|
}
|
|
let method = "";
|
|
if len(req_line) > 0 { method = req_line[0]; }
|
|
let path = "/";
|
|
if len(req_line) > 1 { path = req_line[1]; }
|
|
return HttpReq { method: method, path: path, headers: headers, body: body };
|
|
}
|
|
}
|
|
|
|
fn rpc_result(id: json.Value, result_json: Text) -> HttpResp {
|
|
let idj = json.encode(id);
|
|
return HttpResp { status: 200, body: "{\"jsonrpc\":\"2.0\",\"id\":${idj},\"result\":${result_json}}" };
|
|
}
|
|
|
|
fn rpc_error(id: json.Value, code: Int, message: Text) -> HttpResp {
|
|
let idj = json.encode(id);
|
|
let e = json.encode(RpcErr { code: code, message: message });
|
|
return HttpResp { status: 200, body: "{\"jsonrpc\":\"2.0\",\"id\":${idj},\"error\":${e}}" };
|
|
}
|
|
|
|
fn status_text(code: Int) -> Text {
|
|
switch code {
|
|
case 200: return "OK";
|
|
case 202: return "Accepted";
|
|
case 400: return "Bad Request";
|
|
case 401: return "Unauthorized";
|
|
case 404: return "Not Found";
|
|
case 405: return "Method Not Allowed";
|
|
case 413: return "Payload Too Large";
|
|
default: return "Error";
|
|
}
|
|
}
|
|
|
|
-- Static data stays static: the Haxe original rebuilt these objects per
|
|
-- tools/list call (Mcp.hx:133-165) for no benefit. Four tools — the two
|
|
-- minilog tools live with the DB engine, per the spec's out-of-scope list.
|
|
const TOOL_SCHEMAS = "[{\"name\":\"get_running_crons\",\"description\":\"List cron.d entries with schedule, command, log path, next fire time, and whether each is running right now.\",\"inputSchema\":{\"type\":\"object\",\"properties\":{}}},{\"name\":\"list_logs\",\"description\":\"List the known log files (path, size, mtime, source).\",\"inputSchema\":{\"type\":\"object\",\"properties\":{}}},{\"name\":\"tail_log\",\"description\":\"Last N complete lines of a known log file.\",\"inputSchema\":{\"type\":\"object\",\"properties\":{\"path\":{\"type\":\"string\",\"description\":\"log path from list_logs\"},\"lines\":{\"type\":\"integer\",\"description\":\"default 20, max 200\"}},\"required\":[\"path\"]}},{\"name\":\"search_log\",\"description\":\"Case-insensitive substring search over the recent tail (last 4 MiB) of a known log file, newest matches first.\",\"inputSchema\":{\"type\":\"object\",\"properties\":{\"path\":{\"type\":\"string\",\"description\":\"log path from list_logs\"},\"pattern\":{\"type\":\"string\"},\"maxMatches\":{\"type\":\"integer\",\"description\":\"default 20, max 100\"}},\"required\":[\"path\",\"pattern\"]}}]"
|
|
|
|
-- Tools.hx — the MCP tool layer: plain functions over parse_dir + the log
|
|
-- allowlist. Read-only; every probe degrades to "not running" on ambiguity
|
|
-- (same safe direction as Flock.held).
|
|
class Tools {
|
|
cron_dir: Text
|
|
extra_logs: multi Text
|
|
|
|
static const TAIL_DEFAULT = 20
|
|
static const TAIL_MAX = 200
|
|
static const SEARCH_CAP = 4194304 -- last 4 MiB only
|
|
static const SEARCH_MAX_DEFAULT = 20
|
|
static const SEARCH_MAX_CAP = 100
|
|
|
|
fn allowed_paths() -> multi Text {
|
|
let seen: map<Text, Bool> = {};
|
|
let res: multi Text = [];
|
|
for e in parse_dir(self.cron_dir).entries {
|
|
if has(seen, e.log_path) == false {
|
|
seen[e.log_path] = true;
|
|
push(res, e.log_path);
|
|
}
|
|
}
|
|
for p in self.extra_logs {
|
|
if has(seen, p) == false {
|
|
seen[p] = true;
|
|
push(res, p);
|
|
}
|
|
}
|
|
return res;
|
|
}
|
|
|
|
-- needle for the un-flocked liveness probe: the command's first
|
|
-- /-starting token (the script path), else its first token.
|
|
static fn needle(command: Text) -> Text {
|
|
let tokens = split_ws(command);
|
|
for t in tokens {
|
|
if starts_with(t, "/") { return t; }
|
|
}
|
|
return tokens[0];
|
|
}
|
|
|
|
fn get_running_crons(now: Int) -> ToolOut {
|
|
let rows: multi CronRow = [];
|
|
for e in parse_dir(self.cron_dir).entries {
|
|
let nf = next_fire(e.schedule, now);
|
|
let nf_text: ?Text = nil;
|
|
if nf != nil { nf_text = time.iso(nf); }
|
|
let running = false;
|
|
if e.lock_path != nil {
|
|
running = Flock.held(e.lock_path);
|
|
} else {
|
|
running = Pgrep.alive(Tools.needle(e.command));
|
|
}
|
|
push(rows, CronRow { schedule: e.schedule, command: e.command,
|
|
logPath: e.log_path, nextFire: nf_text, running: running });
|
|
}
|
|
return out(json.encode(rows));
|
|
}
|
|
|
|
fn list_logs() -> ToolOut {
|
|
let cron_set: map<Text, Bool> = {};
|
|
for e in parse_dir(self.cron_dir).entries { cron_set[e.log_path] = true; }
|
|
let rows: multi LogRow = [];
|
|
for p in self.allowed_paths() {
|
|
let st = fs.stat(p);
|
|
let size: ?Int = nil;
|
|
let mtime: ?Int = nil;
|
|
if st != nil {
|
|
size = st.size;
|
|
mtime = st.mtime;
|
|
}
|
|
let source = "config";
|
|
if has(cron_set, p) { source = "cron"; }
|
|
push(rows, LogRow { path: p, sizeBytes: size, mtime: mtime, source: source });
|
|
}
|
|
return out(json.encode(rows));
|
|
}
|
|
|
|
fn check_path(path: Text) -> ?ToolOut {
|
|
for p in self.allowed_paths() {
|
|
if p == path { return nil; }
|
|
}
|
|
let known = join(self.allowed_paths(), ", ");
|
|
return err("unknown log path: ${path} — known logs: ${known}");
|
|
}
|
|
|
|
fn tail_log(path: Text, lines: ?Int) -> ToolOut {
|
|
let bad = self.check_path(path);
|
|
if bad != nil { return bad; }
|
|
let n = TAIL_DEFAULT;
|
|
if lines != nil { n = lines; }
|
|
if n > TAIL_MAX { n = TAIL_MAX; }
|
|
if n < 1 { n = 1; }
|
|
let got = last_lines(path, n);
|
|
if got == nil { return err("log file missing: ${path}"); }
|
|
let clean: multi Text = [];
|
|
for l in got { push(clean, sanitize(l)); }
|
|
return out(json.encode(clean));
|
|
}
|
|
|
|
fn search_log(path: Text, pattern: Text, max_matches: ?Int) -> ToolOut {
|
|
let bad = self.check_path(path);
|
|
if bad != nil { return bad; }
|
|
if pattern == "" { return err("empty pattern"); }
|
|
let st = fs.stat(path);
|
|
if st == nil { return err("log file missing: ${path}"); }
|
|
let cap = SEARCH_MAX_DEFAULT;
|
|
if max_matches != nil { cap = max_matches; }
|
|
if cap > SEARCH_MAX_CAP { cap = SEARCH_MAX_CAP; }
|
|
if cap < 1 { cap = 1; }
|
|
let start = 0;
|
|
if st.size > SEARCH_CAP { start = st.size - SEARCH_CAP; }
|
|
let lc = to_lower(pattern);
|
|
let matches: multi Match = []; -- oldest->newest, reversed at the end
|
|
let offset = start;
|
|
let carry = ""; -- partial line across chunk boundary
|
|
while offset < st.size {
|
|
let want = st.size - offset;
|
|
if want > CHUNK { want = CHUNK; }
|
|
let chunk = fs.read_at(path, offset, want);
|
|
if len(chunk) == 0 { break; }
|
|
let lines = split(carry .. chunk, "\n");
|
|
carry = pop(lines); -- last piece has no newline yet
|
|
for line in lines {
|
|
if index_of(to_lower(line), lc) >= 0 {
|
|
push(matches, Match { byteOffset: offset, line: sanitize(line) });
|
|
}
|
|
}
|
|
offset = offset + len(chunk);
|
|
}
|
|
reverse(matches); -- newest first
|
|
if len(matches) > cap { matches = slice(matches, 0, cap); }
|
|
let note = "searched whole file";
|
|
if start > 0 { note = "searched last 4 MiB only"; }
|
|
return out(json.encode(SearchResult { matches: matches,
|
|
searchedBytes: st.size - start, note: note }));
|
|
}
|
|
}
|
|
|
|
fn out(text: Text) -> ToolOut { return ToolOut { is_error: false, text: text }; }
|
|
fn err(text: Text) -> ToolOut { return ToolOut { is_error: true, text: text }; }
|