writeonce/docs/examples/log-watcher/mcp.wo

351 lines
14 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. Connections are
-- owned handles: the drop at each iteration's end IS the close — the
-- Haxe try/close/catch bookkeeping (Mcp.hx:101-105) has no equivalent
-- and needs none.
fn serve(port: Int) -> Int {
let srv = net.listen("127.0.0.1", port);
while true {
if env.stopping() { 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) {}
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) {}
}
}
-- 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 }; }