From 2171b2d94b4b5b9237145cfe2c6e4938e1b38a4e Mon Sep 17 00:00:00 2001 From: "shoney.arickathil" Date: Mon, 10 Aug 2026 09:26:07 +0200 Subject: [PATCH] docs: log-watcher program --- docs/examples/log-watcher/README.md | 138 ++++++++++ docs/examples/log-watcher/cron.wo | 305 ++++++++++++++++++++ docs/examples/log-watcher/logtail.wo | 157 +++++++++++ docs/examples/log-watcher/main.wo | 124 +++++++++ docs/examples/log-watcher/mcp.wo | 351 ++++++++++++++++++++++++ docs/examples/log-watcher/probes.wo | 36 +++ docs/examples/log-watcher/supervisor.wo | 248 +++++++++++++++++ docs/examples/log-watcher/watcher.wo | 64 +++++ docs/examples/log-watcher/wo.toml | 11 + 9 files changed, 1434 insertions(+) create mode 100644 docs/examples/log-watcher/README.md create mode 100644 docs/examples/log-watcher/cron.wo create mode 100644 docs/examples/log-watcher/logtail.wo create mode 100644 docs/examples/log-watcher/main.wo create mode 100644 docs/examples/log-watcher/mcp.wo create mode 100644 docs/examples/log-watcher/probes.wo create mode 100644 docs/examples/log-watcher/supervisor.wo create mode 100644 docs/examples/log-watcher/watcher.wo create mode 100644 docs/examples/log-watcher/wo.toml diff --git a/docs/examples/log-watcher/README.md b/docs/examples/log-watcher/README.md new file mode 100644 index 0000000..f39963b --- /dev/null +++ b/docs/examples/log-watcher/README.md @@ -0,0 +1,138 @@ +# `log-watcher` — the systems-track sample workload + +The Haxe original (`~/projects/log-watcher`, ~1,200 lines compiled to C++) +ported file for file, per the approved +[systems-track design](../../superpowers/specs/2026-08-01-systems-track-design.md) +(Part 4). A single-binary systems daemon: log-tail watcher, cron.d +supervisor, flock/pgrep probes, hand-rolled MCP-over-HTTP server, JSONL +detection sink. Program mode (`fn main`, blocking legal, one shard) plus the +five builtin stdlib modules — `fs`, `proc`, `net`, `time`, `json` — carry +all of it; read each `.wo` next to its `.hx` sibling. + +> **Status: design artifact — the spec's forcing function.** The systems +> track is approved, pre-implementation. Today's `woc` (milestone 1) +> recovers the `class`/`fn` skeletons in these files (`--dump-ast` lists +> every Watcher method) but diagnoses the adopted surface as WO-E101: +> `use`, `typedef`, standalone union aliases (`type CronResult = …`), +> `pub(read)`, `switch`, `try`. This sample exists to force that grammar +> (the blog/ecommerce/pricing precedent) and becomes the track's acceptance +> test: it compiles and detects a real silent death when the track ships. + +## The mapping + +| `.wo` file | `.hx` sibling | carries | could not express | +| --- | --- | --- | --- | +| `main.wo` | Main.hx | subcommand dispatch; config decode into a `typedef` with `?fields` | — | +| `logtail.wo` | LogTail.hx | `TailState` record, bounded tail reads, rotation-by-inode, strict/relaxed classification, `last_lines` | — | +| `watcher.wo` | Watcher.hx | the per-log state machine: ALERT/CLEAR rule, cron `done()`/`result()`, `CronResult` union | — | +| `cron.wo` | Cron.hx | cron.d parse (aliases, redirect target, flock path), next-fire scan with the Vixie dom/dow OR quirk | — | +| `supervisor.wo` | Supervisor.hx | the 250 ms event loop: rescan, collapse, pre-fire lock probe, completions, detections sink | — | +| `probes.wo` | Flock.hx, Pgrep.hx | `proc.run` exit-code probes as `static fn`s, safe-direction fallbacks | — | +| `mcp.wo` | Mcp.hx, Tools.hx | typed request/response records; pure `handle(req) -> resp` kept socket-free; serve loop over `net` | — | + +The spec's Part 4 table folds Supervisor.hx into the other files; +`supervisor.wo` stays separate because the original is a separate file and +the loop is the program. MiniLog.hx and its two tools +(`load_log_db`/`query_log_db`) are deliberately absent: the spec's +out-of-scope list assigns "sqlite-equivalent embedded SQL over RAM" to the +DB engine — a future sample wires MiniLog's idea to `select`. + +## What the port deletes + +| Haxe | Why it's gone | +| --- | --- | +| `import sys.FileSystem / sys.io.File / sys.io.FileSeek` | one `use fs`; `read_at` takes the offset — no seek state, no open/close bookkeeping (RAII: handle drop = close) | +| `LogTail.newState()` | record field defaults; construction is the brace literal | +| `Util.hx` (`say` = println + flush) | `print` flushes on newline (spec Part 2) | +| `loadConfig`'s Dynamic field-poking | one `json.decode(raw) as FileConfig` — typed, `?fields`, nil on mismatch | +| `try … catch (e:Dynamic)` probes | `try … catch (e)` over the trap system; expected absence is `?T`/nil instead | +| `#if portable` linker pragma | `woc build` is a static single binary by doctrine | +| socket `try c.close() catch` dance | connection is an owned value; scope end closes it | +| `while(true)` with no way out | `env.stopping()` — SIGTERM/SIGINT land as a flag, no signal callbacks | + +## What this sample forces (spec follow-ups) + +The spec's five modules cover the I/O; writing real code forced these +additions, in shrinking order of importance: + +1. **`time.local(ms) -> {year, month, day, hour, minute, dow}`** and + **`time.iso(ms) -> Text`** — cron next-fire needs calendar decomposition + in the host timezone; the detections sink needs an ISO stamp. The spec's + `time` sketch (now/mono/sleep) cannot express cron. +2. **`json.Value`** — an opaque, re-encodable JSON value. JSON-RPC echoes + `id` back verbatim (number | string | null); no record type can hold it. + Decode-target position only, like the `as` rule. +3. **`fs.stat` gains `dir: Bool`** — cron.d scanning must skip + subdirectories; `{size, inode, mtime}` cannot. +4. **Text builtins**: `len`, `substr(s, start, len)`, `split`, `split_ws`, + `trim`, `starts_with`, `ends_with`, `index_of`, `last_index_of`, + `to_lower`, `byte_at`, `char_of`, `parse_int`, `"${…}"` interpolation. + No regex module exists — the original's five ERegs are hand-rolled scans + (see `classify_loose`, `is_env_line`, `find_log`, `sanitize`). +5. **Collection builtins**: `push`, `pop` (returns the removed element), + `shift`, `slice`, `reverse`, `sort`, `join`; `map` index (returns `?V`), + `has`, `remove`, `for k, v in m`. + +## Deliberate divergences + +| Original | Port | Why | +| --- | --- | --- | +| `Float` seconds everywhere | `Int` milliseconds | no Float scalar in the language; ms-as-Int matches the runtime (and Money's minor-units precedent) | +| config `services` entries: string **or** `{path}` object | strings only | typed decode; the object form was never used in the deployed config | +| tool schemas built per call | one JSON `const` | static data is static; the wire bytes are identical | +| JSON-RPC envelope via `Json.stringify` of a Dynamic | typed sub-records `json.encode`d into a concatenated envelope | typed json has no heterogeneous-object builder; `id` passes through as `json.Value` | +| `list_logs` mtime: seconds (float) | ms (int) | consistency with every other timestamp in the port | + +## Memory model: why this sample has zero `@gc` and zero `@table` + +The Haxe original transcompiles to C++ (`bin/src/*.cpp`, hxcpp target), +which makes the contrast measurable — there, **everything** is `@gc`: + +1. Every object is GC-heap: `new Watcher_obj` behind `hx::ObjectPtr`, with + `HX_DEFINE_STACK_FRAME` on every method so the collector can scan roots. +2. Every typedef is a `Dynamic` hash object: `TailState` has no C++ struct — + the hot poll path does `this->st->__Field(HX_("lastLevel",…))`, a hashed + string lookup per field access, every 2 s, per watched file + (`LogTail.cpp` has 34 `Dynamic` sites, `Supervisor.cpp` 91, `Mcp.cpp` 123). +3. Allocation per poll: `poll()` returns a fresh anonymous GC object each + call. + +The port's ownership graph is a pure tree, so MVS covers all of it with +owned values and `@gc` earns its keep nowhere: + +| Object | Owner | wo semantics | +| --- | --- | --- | +| `TailState` | its `Watcher` (field) | owned value; the `__Field` hash lookup becomes a fixed-offset load | +| `Watcher` | `Supervisor.services` **or** `.active` — never both | owned in one container; `remove()` IS the destructor — the Haxe kill-self comment, now literal | +| `CronWatch` | `scheduled` map | owned; iteration mutates via borrow | +| `PollResult`, `ParseResult`, every record | callee frame | stack lifetime, `DROP` at scope end, zero heap | +| `Tools`, `Mcp` | `main` → `Mcp.tools` | owned chain | + +hxcpp pays GC on 100% of these objects; MVS pays on 0%. `@gc` would only +enter if a shape changed — one `Watcher` aliased by two live registries at +once, or a shared mutable cache aliased across requests (the OOP spec's +`PriceCache` pattern). log-watcher has neither. + +`@table` is also correctly absent: it configures storage for engine-bound +classes, and program mode has no DB engine (sub-project 3). The C++ shows +exactly where it lands later — `MiniLog.cpp`'s in-memory sqlite, +`entries(id, path, seq, level, body, byteOffset)` + `idx_level`, which is +the future sample the spec names (wiring MiniLog's idea to `select`): +`@table(name: "entries", index: [path, seq])` on a `LogEntry` class turns +`load_log_db`/`query_log_db` into plain `insert`/`select`. Same for the +detections JSONL → a `@table` type with `insert` replacing +open-append-close. Until the engine links, annotating anything here would +claim storage that doesn't exist. + +## Try it (when the track ships) + +```bash +woc build docs/examples/log-watcher # single static binary +log-watcher watch /var/log/myapp.log 10 2 # ALERT/CLEAR on stdout +log-watcher run /etc/cron.d config.json # supervisor +log-watcher mcp /etc/cron.d config.json # MCP on 127.0.0.1: +``` + +Acceptance (spec): compiles; the tail state machine, cron next-fire, and +MCP `handle` pass fixtures ported from the Haxe test suite; `watch` detects +an error-final quiet period against a growing tempfile, end to end. diff --git a/docs/examples/log-watcher/cron.wo b/docs/examples/log-watcher/cron.wo new file mode 100644 index 0000000..9159a4d --- /dev/null +++ b/docs/examples/log-watcher/cron.wo @@ -0,0 +1,305 @@ +-- cron.wo — Cron.hx, file for file. Read-only parser for a cron.d-format +-- directory (plan 03). The directory is always a parameter so tests and +-- demos run on fixtures without root. +-- +-- Haxe leaned on three ERegs (ws split, env-line, redirect) and on Date for +-- next-fire. The regexes are spelled out as scans; Date forces the one +-- time-module addition this sample discovers: time.local(ms) -> +-- { year, month (1-12), day, hour, minute, dow (0-6, Sunday 0) } in the +-- host timezone — cron semantics are local time. + +use fs +use time + +typedef CronEntry = { + schedule: Text -- 5-field expression, aliases expanded + user: Text + command: Text -- verbatim, redirections included + log_path: Text + lock_path: ?Text -- flock lock file, nil for un-flocked commands + src_file: Text + src_line: Int +} + +typedef Skipped = { + src_file: Text + src_line: Int + reason: Text +} + +typedef ParseResult = { + entries: multi CronEntry + skipped: multi Skipped +} + +-- Cron.hx:25-33 (private typedef Fields). +typedef Fields = { + minute: multi Bool + hour: multi Bool + dom: multi Bool + month: multi Bool + dow: multi Bool + dom_r: Bool -- dom field restricted (not "*") + dow_r: Bool +} + +const FILE_CAP = 1048576 -- cron.d files are small; cap read_all anyway + +pub fn parse_dir(dir: Text) -> ParseResult { + let res = ParseResult { entries: [], skipped: [] }; + if fs.exists(dir) == false { return res; } + let names = try fs.list(dir) catch (e) nil; + if names == nil { + skip(res, dir, 0, "unreadable directory"); + return res; + } + sort(names); + for name in names { + if index_of(name, ".") >= 0 { continue; } -- cron ignores dotted names + let st = fs.stat("${dir}/${name}"); + if st == nil { continue; } + if st.dir { continue; } + parse_file("${dir}/${name}", res); + } + return res; +} + +fn parse_file(path: Text, mut res: ParseResult) { + let content = try fs.read_all(path, FILE_CAP) catch (e) nil; + if content == nil { + skip(res, path, 0, "unreadable"); + return; + } + let line_no = 0; + for raw in split(content, "\n") { + line_no = line_no + 1; + let line = trim(raw); + if line == "" or starts_with(line, "#") { continue; } + if is_env_line(line) { continue; } + parse_job(line, path, line_no, res); + } +} + +fn skip(mut res: ParseResult, src_file: Text, src_line: Int, reason: Text) { + push(res.skipped, Skipped { src_file: src_file, src_line: src_line, reason: reason }); +} + +-- Cron.hx:53's ^[A-Za-z_][A-Za-z0-9_]*\s*= — a NAME= assignment line. +fn is_env_line(line: Text) -> Bool { + let c0 = byte_at(line, 0); + if is_word_byte(c0) == false or (c0 >= 48 and c0 <= 57) { return false; } + let i = 1; + while i < len(line) and is_word_byte(byte_at(line, i)) { i = i + 1; } + while i < len(line) and (byte_at(line, i) == 32 or byte_at(line, i) == 9) { i = i + 1; } + return i < len(line) and byte_at(line, i) == 61; -- '=' +} + +fn parse_job(line: Text, src_file: Text, src_line: Int, mut res: ParseResult) { + let tokens = split_ws(line); + let schedule = ""; + let rest: multi Text = []; + if starts_with(tokens[0], "@") { + if tokens[0] == "@reboot" { skip(res, src_file, src_line, "@reboot out of scope"); return; } + let expanded = alias_of(tokens[0]); + if expanded == nil { skip(res, src_file, src_line, "unknown alias ${tokens[0]}"); return; } + schedule = expanded; + rest = slice(tokens, 1, len(tokens)); + } else { + if len(tokens) < 7 { skip(res, src_file, src_line, "malformed: too few fields"); return; } + schedule = join(slice(tokens, 0, 5), " "); + rest = slice(tokens, 5, len(tokens)); + } + if parse_expr(schedule) == nil { skip(res, src_file, src_line, "malformed schedule: ${schedule}"); return; } + if len(rest) < 2 { skip(res, src_file, src_line, "malformed: missing command"); return; } + let command = join(slice(rest, 1, len(rest)), " "); + let log_path = find_log(command); + if log_path == nil { skip(res, src_file, src_line, "not watchable (no plain .log redirection)"); return; } + push(res.entries, CronEntry { + schedule: schedule, user: rest[0], command: command, + log_path: log_path, lock_path: find_lock(command), + src_file: src_file, src_line: src_line, + }); +} + +-- The log path is the first redirection target that is a plain .log path +-- (plan 03 watchability convention: no $VARs, no pipes). Cron.hx:54's +-- redirect EReg becomes a scan: at each '>' take the token that follows — +-- the &>/2>/>> prefixes all end in the same '>'. +fn find_log(command: Text) -> ?Text { + let i = 0; + while i < len(command) { + if byte_at(command, i) != 62 { i = i + 1; continue; } -- '>' + let j = i + 1; + if j < len(command) and byte_at(command, j) == 62 { j = j + 1; } -- '>>' + while j < len(command) and (byte_at(command, j) == 32 or byte_at(command, j) == 9) { j = j + 1; } + let k = j; + while k < len(command) and byte_at(command, k) != 32 and byte_at(command, k) != 9 { k = k + 1; } + let target = substr(command, j, k - j); + if ends_with(target, ".log") and index_of(target, "$") == -1 { return target; } + i = k; + if i == j { i = j + 1; } -- trailing '>' with no target + } + return nil; +} + +-- A flock-wrapped command notes its lock file: the lock path is the first +-- absolute path after `flock`, before any -c (no flock option takes an +-- absolute-path value, so e.g. `-w 600` is skipped naturally). Cron.hx:134-143. +fn find_lock(command: Text) -> ?Text { + let tokens = split_ws(command); + if tokens[0] != "flock" and ends_with(tokens[0], "/flock") == false { return nil; } + let i = 1; + while i < len(tokens) { + if tokens[i] == "-c" { return nil; } + if starts_with(tokens[i], "/") { return tokens[i]; } + i = i + 1; + } + return nil; +} + +-- Next fire time (ms) strictly after `from_ms`, or nil if none within a +-- year (also nil for an invalid expression). Cron.hx:147-166, minute- +-- resolution forward scan. +pub fn next_fire(expr: Text, from_ms: Int) -> ?Int { + let f = parse_expr(expr); + if f == nil { return nil; } + let t = (from_ms / 60000) * 60000 + 60000; + let limit = t + 366 * 24 * 3600 * 1000; + while t < limit { + let d = time.local(t); + if f.month[d.month] == false { t = t + 86400000; continue; } + let day_ok = false; + if f.dom_r and f.dow_r { + day_ok = f.dom[d.day] or f.dow[d.dow]; -- vixie quirk: restricted dom + } else { -- AND dow match on either (OR) + day_ok = f.dom[d.day] and f.dow[d.dow]; + } + if day_ok == false or f.hour[d.hour] == false { + t = (t / 3600000) * 3600000 + 3600000; -- next hour + continue; + } + if f.minute[d.minute] == false { t = t + 60000; continue; } + return t; + } + return nil; +} + +fn parse_expr(expr: Text) -> ?Fields { + let e = trim(expr); + if starts_with(e, "@") { + let a = alias_of(e); + if a == nil { return nil; } + e = a; + } + let p = split_ws(e); + if len(p) != 5 { return nil; } + let minute = parse_field(p[0], 0, 59, NoNames); + let hour = parse_field(p[1], 0, 23, NoNames); + let dom = parse_field(p[2], 1, 31, NoNames); + let month = parse_field(p[3], 1, 12, MonthNames); + let dow = parse_field(p[4], 0, 7, DowNames); + if minute == nil or hour == nil or dom == nil or month == nil or dow == nil { return nil; } + if dow[7] { dow[0] = true; } -- 7 is Sunday too + return Fields { + minute: minute, hour: hour, dom: dom, month: month, dow: dow, + dom_r: p[2] != "*", dow_r: p[4] != "*", + }; +} + +-- Cron.hx:190-222. `names` selected the lookup map in Haxe; a union selects +-- the switch here. +type NameKind = NoNames | MonthNames | DowNames + +fn parse_field(spec: Text, lo: Int, hi: Int, names: NameKind) -> ?multi Bool { + let res: multi Bool = []; + let fill = 0; + while fill <= hi { push(res, false); fill = fill + 1; } + for part in split(spec, ",") { + let step = 1; + let range = part; + let slash = index_of(part, "/"); + if slash >= 0 { + range = substr(part, 0, slash); + let s = parse_int(substr(part, slash + 1, len(part) - slash - 1)); + if s == nil { return nil; } + if s < 1 { return nil; } + step = s; + } + let a: ?Int = nil; + let b: ?Int = nil; + if range == "*" { + a = lo; + b = hi; + } else { + let dash = index_of(range, "-"); + if dash >= 0 { + a = value(substr(range, 0, dash), names); + b = value(substr(range, dash + 1, len(range) - dash - 1), names); + } else { + a = value(range, names); + b = a; + if slash >= 0 { b = hi; } -- "n/step" means "n-hi/step" + } + } + if a == nil or b == nil { return nil; } + if a < lo or b > hi or a > b { return nil; } + let v = a; + while v <= b { res[v] = true; v = v + step; } + } + return res; +} + +fn value(tok: Text, names: NameKind) -> ?Int { + let n = parse_int(tok); + if n != nil { return n; } + let key = substr(to_lower(tok), 0, 3); + switch names { + case NoNames: return nil; + case MonthNames: return month_of(key); + case DowNames: return dow_of(key); + } +} + +-- Cron.hx:38-51 — the three static maps become exhaustive-by-default +-- switches (no map literals needed for fixed tables). +fn alias_of(tok: Text) -> ?Text { + switch tok { + case "@hourly": return "0 * * * *"; + case "@daily", "@midnight": return "0 0 * * *"; + case "@weekly": return "0 0 * * 0"; + case "@monthly": return "0 0 1 * *"; + case "@yearly", "@annually": return "0 0 1 1 *"; + default: return nil; + } +} + +fn month_of(key: Text) -> ?Int { + switch key { + case "jan": return 1; + case "feb": return 2; + case "mar": return 3; + case "apr": return 4; + case "may": return 5; + case "jun": return 6; + case "jul": return 7; + case "aug": return 8; + case "sep": return 9; + case "oct": return 10; + case "nov": return 11; + case "dec": return 12; + default: return nil; + } +} + +fn dow_of(key: Text) -> ?Int { + switch key { + case "sun": return 0; + case "mon": return 1; + case "tue": return 2; + case "wed": return 3; + case "thu": return 4; + case "fri": return 5; + case "sat": return 6; + default: return nil; + } +} diff --git a/docs/examples/log-watcher/logtail.wo b/docs/examples/log-watcher/logtail.wo new file mode 100644 index 0000000..cb5963e --- /dev/null +++ b/docs/examples/log-watcher/logtail.wo @@ -0,0 +1,157 @@ +-- logtail.wo — LogTail.hx, file for file: bounded tail reads over `fs`. +-- One poll = read at most CHUNK new tail bytes, judge only complete lines +-- (plan 02: never scan a file front to back; torn final line held back). +-- +-- The three Haxe imports (sys.FileSystem, sys.io.File, sys.io.FileSeek) +-- collapse into one `use fs`: stat carries the inode, read_at takes the +-- offset as a parameter (no seek state), and the file handle lives and dies +-- inside the builtin — RAII instead of the open/try/close dance. + +use fs + +-- LogTail.hx:6-11. Field defaults replace LogTail.newState(): construction +-- is the brace literal `TailState {}`, the defaults fill in. +typedef TailState = { + offset: Int = 0 -- next read position (past last complete line) + ino: Int = -1 -- inode at last poll, -1 before first sight + last_level: Text = "" -- level of the last complete entry, "" if none + last_newline_at: Int = 0 -- ms clock when complete lines last arrived +} + +typedef PollResult = { + exists: Bool + new_lines: Int + bytes_read: Int +} + +const CHUNK = 65536 + +-- `now` is injected (ms) so tests can drive synthetic time. LogTail.hx:28-75. +pub fn poll(path: Text, mut st: TailState, now: Int) -> PollResult { + let stat = fs.stat(path); + if stat == nil { return PollResult { exists: false, new_lines: 0, bytes_read: 0 }; } + + if st.ino == -1 { + -- first sight: start at most CHUNK before EOF; a partial first line is + -- read as a continuation, which is acceptable + st.ino = stat.inode; + st.offset = 0; + if stat.size > CHUNK { st.offset = stat.size - CHUNK; } + } else if stat.inode != st.ino or stat.size < st.offset { + -- rotation (rename/recreate or truncate): restart from the top + st.ino = stat.inode; + st.offset = 0; + st.last_level = ""; + } + if stat.size - st.offset > CHUNK { st.offset = stat.size - CHUNK; } -- burst: jump to tail + if stat.size <= st.offset { return PollResult { exists: true, new_lines: 0, bytes_read: 0 }; } + + -- Shrunk between stat and read (the Haxe Eof catch): read_at returns what + -- is actually there; the next poll re-syncs. + let chunk = fs.read_at(path, st.offset, stat.size - st.offset); + if len(chunk) == 0 { return PollResult { exists: true, new_lines: 0, bytes_read: 0 }; } + + -- only complete lines count: cut at the last newline, hold the rest + let nl = last_index_of(chunk, "\n"); + if nl == -1 { return PollResult { exists: true, new_lines: 0, bytes_read: len(chunk) }; } + + let lines = split(substr(chunk, 0, nl + 1), "\n"); + pop(lines); -- empty piece after the final newline + st.offset = st.offset + nl + 1; + for line in lines { + -- relaxed rule (service-health): timestamped app logs must classify + -- too, or their errors never trip the alert rule + let lv = classify_loose(sanitize(line)); + if lv != nil { st.last_level = lv; } -- else continuation: inherits + } + if len(lines) > 0 { st.last_newline_at = now; } + return PollResult { exists: true, new_lines: len(lines), bytes_read: len(chunk) }; +} + +-- Strict prefix rule for demo/test logs. LogTail.hx:77-82. +pub fn classify(line: Text) -> ?Text { + if starts_with(line, "info") { return "info"; } + if starts_with(line, "warn") { return "warn"; } + if starts_with(line, "error") { return "error"; } + return nil; +} + +-- Real app logs put a timestamp first ("2026-07-22T15:40:00 - error: …"); +-- the relaxed rule also accepts the level after a leading token. Haxe used +-- an EReg (LogTail.hx:87, ~/^\S+\s+-\s+(info|warn|error)\b/); with no regex +-- in the language the same rule is spelled out: token, lone dash, level +-- with a word boundary. Callers pass a sanitized line — ANSI codes hide +-- the prefix. +pub fn classify_loose(line: Text) -> ?Text { + let lv = classify(line); + if lv != nil { return lv; } + let tokens = split_ws(line); + if len(tokens) < 3 { return nil; } + if tokens[1] != "-" { return nil; } + return level_bounded(tokens[2]); +} + +-- The regex's \b: "error:" and "error," carry the level; "errors" does not. +fn level_bounded(tok: Text) -> ?Text { + for lv in ["info", "warn", "error"] { + if starts_with(tok, lv) { + if len(tok) == len(lv) { return lv; } + if is_word_byte(byte_at(tok, len(lv))) { return nil; } + return lv; + } + } + return nil; +} + +fn is_word_byte(c: Int) -> Bool { + if c >= 48 and c <= 57 { return true; } -- 0-9 + if c >= 65 and c <= 90 { return true; } -- A-Z + if c >= 97 and c <= 122 { return true; } -- a-z + return c == 95; -- _ +} + +-- Tools.hx:24-34, moved here beside its heaviest caller (poll). Log lines +-- can carry ANSI color sequences and stray control bytes; they would break +-- classification and produce invalid JSON at the client. Strip ESC[…letter +-- sequences, drop other C0 bytes (tab stays). The Haxe EReg becomes an +-- explicit byte scan. +pub fn sanitize(line: Text) -> Text { + let out = ""; + let i = 0; + while i < len(line) { + let c = byte_at(line, i); + if c == 27 and i + 1 < len(line) and byte_at(line, i + 1) == 91 { + i = i + 2; -- skip ESC [ + while i < len(line) { + let f = byte_at(line, i); + i = i + 1; + if (f >= 65 and f <= 90) or (f >= 97 and f <= 122) { break; } + } + continue; + } + if c >= 32 or c == 9 { out = out + char_of(c); } + i = i + 1; + } + return out; +} + +-- Last <= n complete lines from the final CHUNK bytes of the file. A +-- cut-off first line is acceptable (same as first-sight poll); an +-- unterminated final line is dropped. nil when the file is missing. +-- LogTail.hx:98-120. +pub fn last_lines(path: Text, n: Int) -> ?multi Text { + let stat = fs.stat(path); + if stat == nil { return nil; } + let start = 0; + if stat.size > CHUNK { start = stat.size - CHUNK; } + if stat.size - start <= 0 { return []; } + let chunk = fs.read_at(path, start, stat.size - start); + if len(chunk) == 0 { return []; } + let nl = last_index_of(chunk, "\n"); + if nl == -1 { return []; } + let lines = split(substr(chunk, 0, nl + 1), "\n"); + pop(lines); -- empty piece after the final newline + if start > 0 and len(lines) > 0 { shift(lines); } -- cut-off first line + if len(lines) > n { return slice(lines, len(lines) - n, len(lines)); } + return lines; +} diff --git a/docs/examples/log-watcher/main.wo b/docs/examples/log-watcher/main.wo new file mode 100644 index 0000000..20d00a9 --- /dev/null +++ b/docs/examples/log-watcher/main.wo @@ -0,0 +1,124 @@ +-- main.wo — Main.hx, file for file: subcommand dispatch, config decode. +-- A free `fn main(args) -> Int` makes this project a PROGRAM (systems-track +-- spec Part 2): `wo run` executes it, the return value is the exit code, +-- blocking builtins are legal on the single shard. +-- +-- Main.hx:5-7's `#if portable` linker pragma has no equivalent and needs +-- none: `woc build` produces a self-contained static binary by doctrine. +-- Util.hx is gone entirely: print flushes on newline (spec Part 2), which +-- is the only thing Util.say existed to do. + +use fs +use env +use time +use json + +-- config.json, decoded typed: Main.hx:50-61's Dynamic field-poking becomes +-- one checked decode — missing optional fields are fine, a shape mismatch +-- yields nil, never a trap. Field names are the wire format (camelCase); +-- values are seconds on the wire, ms internally. +typedef FileConfig = { + ?pollInterval: Int + ?quietPeriod: Int + ?rescanInterval: Int + ?detections: Text + ?services: multi Text + ?logs: multi Text + ?mcp: McpConfig +} +typedef McpConfig = { ?port: Int, ?apiKey: Text } + +fn main(args: multi Text) -> Int { + if len(args) >= 2 and args[0] == "watch" { + let quiet_s = 10; + if len(args) >= 3 { + let q = parse_int(args[2]); + if q != nil { quiet_s = q; } + } + let poll_s = 2; + if len(args) >= 4 { + let p = parse_int(args[3]); + if p != nil { poll_s = p; } + } + let w = Watcher { path: args[1], quiet_ms: quiet_s * 1000, + poll_ms: poll_s * 1000, live: true, activated_at: time.now() }; + print("watching ${args[1]} (quiet ${quiet_s}s, poll ${poll_s}s)"); + while true { + if env.stopping() { return 0; } + w.tick(time.now()); + time.sleep(poll_s * 1000); + } + } + + if len(args) >= 2 and args[0] == "run" { + let cfg = SupConfig {}; -- field defaults = Main.hx:21's literal + if len(args) >= 3 { + if load_config(args[2], cfg) == false { return 1; } + } + print("supervising ${args[1]} (poll ${cfg.poll_ms / 1000}s, quiet ${cfg.quiet_ms / 1000}s, rescan ${cfg.rescan_ms / 1000}s)"); + let sup = Supervisor { cron_dir: args[1], cfg: cfg }; + sup.init(); + return sup.run(); + } + + if len(args) >= 3 and args[0] == "mcp" { + let raw = fs.read_all(args[2], 1048576); + if raw == nil { + print_err("config error: cannot read ${args[2]}"); + return 1; + } + let j = json.decode(raw) as FileConfig; + if j == nil { + print_err("config error: ${args[2]} is not valid JSON"); + return 1; + } + -- the key may live outside the config file (systemd EnvironmentFile / .env) + let api_key = env.get("LOG_WATCHER_API_KEY"); + let port: ?Int = nil; + if j.mcp != nil { + if j.mcp.apiKey != nil { api_key = j.mcp.apiKey; } + port = j.mcp.port; + } + if port == nil or api_key == nil { + print_err("config error: mcp.port and an api key (mcp.apiKey or LOG_WATCHER_API_KEY) are required"); + return 1; + } + let extra: multi Text = []; + if j.logs != nil { + for p in j.logs { push(extra, p); } + } + if j.services != nil { + for s in j.services { push(extra, s); } + } + let mcp = Mcp { tools: Tools { cron_dir: args[1], extra_logs: extra }, api_key: api_key }; + print("mcp server on 127.0.0.1:${port} (${args[1]})"); + return mcp.serve(port); + } + + print_err("usage:"); + print_err(" log-watcher watch [quietPeriod] [pollInterval] continuous service watch"); + print_err(" log-watcher run [config.json] supervisor (cron + services)"); + print_err(" log-watcher mcp MCP server (127.0.0.1, Bearer auth)"); + return 1; +} + +fn load_config(path: Text, mut cfg: SupConfig) -> Bool { + let raw = fs.read_all(path, 1048576); + if raw == nil { + print_err("config error: cannot read ${path}"); + return false; + } + let j = json.decode(raw) as FileConfig; + if j == nil { + print_err("config error: ${path} is not valid JSON"); + return false; + } + if j.pollInterval != nil { cfg.poll_ms = j.pollInterval * 1000; } + if j.quietPeriod != nil { cfg.quiet_ms = j.quietPeriod * 1000; } + if j.rescanInterval != nil { cfg.rescan_ms = j.rescanInterval * 1000; } + if j.detections != nil { cfg.detections = j.detections; } + if j.services != nil { + for s in j.services { push(cfg.services, s); } + } + return true; +} diff --git a/docs/examples/log-watcher/mcp.wo b/docs/examples/log-watcher/mcp.wo new file mode 100644 index 0000000..20ab4c2 --- /dev/null +++ b/docs/examples/log-watcher/mcp.wo @@ -0,0 +1,351 @@ +-- 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, 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 = {}; + 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 = {}; + 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 = {}; + 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 }; } diff --git a/docs/examples/log-watcher/probes.wo b/docs/examples/log-watcher/probes.wo new file mode 100644 index 0000000..5ddfdb8 --- /dev/null +++ b/docs/examples/log-watcher/probes.wo @@ -0,0 +1,36 @@ +-- probes.wo — Flock.hx + Pgrep.hx as static fns over `proc` (the spec's +-- `Flock.held` / `Pgrep.alive` pattern is what the `static` adoption is +-- for). Args-array only: no shell-string form exists, so command injection +-- is unrepresentable (systems-track spec Part 3). + +use fs +use proc + +class Flock { + -- flock(1) probe: is an exclusive lock currently held on `path`? + -- Exit 1 = -n conflict = held; 0 = acquired-and-released = free. Anything + -- else (66 unreadable file, spawn failure = no flock binary) = unknown -> + -- report free, so the supervisor falls back to watching normally (the + -- safe direction). The exists() guard keeps the probe read-only: flock -n + -- would O_CREAT a missing lock file, and a missing file means nobody + -- holds it anyway. Flock.hx:1-12. + static fn held(path: Text) -> Bool { + if fs.exists(path) == false { return false; } + let r = try proc.run("flock", ["-n", path, "true"]) catch (e) nil; + if r == nil { return false; } + return r.code == 1; + } +} + +class Pgrep { + -- pgrep(1) probe: is any process whose command line contains `needle` + -- alive? Exit 0 = yes; 1 = no; spawn failure (no pgrep binary) = unknown + -- -> report not running, the same safe direction as Flock.held. The Haxe + -- `try new Process catch return false` is the same shape with the trap + -- surface instead of Dynamic. Pgrep.hx:1-12. + static fn alive(needle: Text) -> Bool { + let r = try proc.run("pgrep", ["-f", "--", needle]) catch (e) nil; + if r == nil { return false; } + return r.code == 0; + } +} diff --git a/docs/examples/log-watcher/supervisor.wo b/docs/examples/log-watcher/supervisor.wo new file mode 100644 index 0000000..033935e --- /dev/null +++ b/docs/examples/log-watcher/supervisor.wo @@ -0,0 +1,248 @@ +-- supervisor.wo — Supervisor.hx, file for file. (The systems-track spec's +-- Part 4 table has no row for this file; it stays its own module file +-- because Supervisor.hx is its own file — the README records the addition.) +-- +-- Single-threaded event loop (supersedes plan 02's thread-per-watch: this +-- is simpler and the kill-self guarantee holds by construction — a +-- completed cron watch is an owned value removed from the map, and file +-- handles only live inside a single fs builtin call). MVS makes the Haxe +-- comment literal: the remove IS the destructor. + +use fs +use time +use env +use json + +-- Supervisor.hx:4-10. Wire values are seconds (config.json compatibility); +-- internal values are ms — the defaults here are Main.hx:21's literal. +typedef SupConfig = { + poll_ms: Int = 2000 + quiet_ms: Int = 10000 + rescan_ms: Int = 60000 + services: multi Text = [] -- log paths watched continuously + ?detections: Text -- JSONL sink for error events (service-health) +} + +-- Supervisor.hx:12-17. +typedef CronWatch = { + schedules: multi Text -- same-log entries collapse into one watch + next_fire: ?Int -- ms, min over schedules; nil = none within a year + lock_path: ?Text -- flock note; nil unless every entry agrees + lock_probe: ?Bool -- pre-fire probe result; nil = not probed +} + +typedef Completion = { + log_path: Text + result: CronResult +} + +-- The detections JSONL line — field names are the wire format +-- (Supervisor.hx:162-168), so they stay camelCase. +typedef Detection = { + ts: Text + path: Text + source: Text + event: Text + lastError: Text +} + +class Supervisor { + -- The lock is probed shortly BEFORE the fire, never at it: at fire time + -- the job's own flock has usually grabbed the lock already, so an + -- at-activation probe would wrongly skip every window. + static const PROBE_LEAD_MS = 2000 + + pub(read) scheduled: map = {} + pub(read) active: map = {} + pub(read) completions: multi Completion = [] + + services: multi Watcher = [] + cron_dir: Text + cfg: SupConfig + last_rescan: Int = 0 -- 0 = never (Haxe NEGATIVE_INFINITY) + dir_sig: Text = "" + first_scan: Bool = true + + -- The Haxe constructor body (Supervisor.hx:40-45): writeonce has no + -- constructor keyword — brace literal, then explicit init. + fn init() { + for p in self.cfg.services { + push(self.services, Watcher { path: p, quiet_ms: self.cfg.quiet_ms, + poll_ms: self.cfg.poll_ms, live: true, activated_at: 0 }); + } + } + + -- `now` in ms, injected for testability; run() feeds the real clock. + fn tick(now: Int) { + if now - self.last_rescan >= self.cfg.rescan_ms { self.rescan(now); } + + for log_path, cw in self.scheduled { + if has(self.active, log_path) { continue; } + if cw.next_fire == nil { continue; } + if cw.lock_path != nil and cw.lock_probe == nil and now >= cw.next_fire - PROBE_LEAD_MS and now < cw.next_fire { + cw.lock_probe = Flock.held(cw.lock_path); + } + if now < cw.next_fire { continue; } + -- no probe taken (e.g. started past the fire) leaves locked false, + -- so the watch runs — the safe direction + let locked = cw.lock_probe == true; + cw.lock_probe = nil; -- reset for the next window + if locked { + print("SKIP-LOCKED ${log_path}: ${cw.lock_path} still held, window skipped"); + cw.next_fire = compute_next(cw.schedules, now); + } else { + print("WATCH ${log_path}: activated"); + self.active[log_path] = Watcher { path: log_path, quiet_ms: self.cfg.quiet_ms, + poll_ms: self.cfg.poll_ms, live: false, activated_at: now }; + } + } + + for w in self.services { + if w.due(now) { + let was = w.alerted; + w.tick(now); + if w.alerted and was == false { self.detect(w.path, "service", "alert", now); } + } + } + + let finished: multi Text = []; + for log_path, w in self.active { + if w.due(now) { w.tick(now); } + if w.done(now) { push(finished, log_path); } + } + for log_path in finished { + let w = self.active[log_path]; + if w == nil { continue; } -- forced ?-handling; never taken + let res = w.result(); + push(self.completions, Completion { log_path: log_path, result: res }); + let verdict = switch res { + case Ok: "ok"; + case ErrorFinal: "ERROR-final"; + case Miss: "missed (no log activity in the window)"; + }; + print("DONE ${log_path}: ${verdict}"); + if res == ErrorFinal { + print("ALERT ${log_path}: cron job finished with error as the last entry"); + self.detect(log_path, "cron", "error-final", now); + } + remove(self.active, log_path); -- kill-self: nothing of the watch remains + let cw = self.scheduled[log_path]; + if cw != nil { cw.next_fire = compute_next(cw.schedules, now); } + } + } + + fn rescan(now: Int) { + self.last_rescan = now; + let sig = self.signature(); + if sig == self.dir_sig { return; } -- mtime-gated: nothing changed + self.dir_sig = sig; + + let res = parse_dir(self.cron_dir); + for sk in res.skipped { + print("SKIP ${sk.src_file}:${sk.src_line}: ${sk.reason}"); + } + + let fresh: map = {}; + for e in res.entries { + if has(fresh, e.log_path) == false { + fresh[e.log_path] = CronWatch { schedules: [], next_fire: nil, + lock_path: e.lock_path, lock_probe: nil }; + } + let cw = fresh[e.log_path]; + if cw == nil { continue; } -- forced ?-handling; never taken + if cw.lock_path != e.lock_path { cw.lock_path = nil; } -- disagreeing writers: always watch + push(cw.schedules, e.schedule); + } + for log_path, cw in fresh { + cw.next_fire = compute_next(cw.schedules, now); + if has(self.scheduled, log_path) == false { + let scheds = join(cw.schedules, " | "); + let note = ""; + if cw.lock_path != nil { note = " (flock ${cw.lock_path})"; } + print("SCHEDULE ${log_path}: ${scheds}${note}"); + } + } + -- removed entries lose their pending activation here; an already + -- active watch is intentionally left to finish its window + self.scheduled = fresh; -- move: the old map drops here + + if self.first_scan { + self.first_scan = false; + -- restart mid-window: a recently written log is watched now + for log_path, cw in self.scheduled { + if has(self.active, log_path) { continue; } + let st = fs.stat(log_path); + if st == nil { continue; } + if now - st.mtime <= self.cfg.quiet_ms { + print("WATCH ${log_path}: activated (mid-window at startup)"); + self.active[log_path] = Watcher { path: log_path, quiet_ms: self.cfg.quiet_ms, + poll_ms: self.cfg.poll_ms, live: false, activated_at: now }; + } + } + } + } + + -- service-health detections sink: error events only, one JSON line per + -- rule hit, open-append-close (logrotate-safe). The file exists to tell + -- the Sheriff something is wrong and where to look; a sink failure is + -- reported but never takes the supervisor down. Supervisor.hx:143-176. + fn detect(log_path: Text, source: Text, event: Text, now: Int) { + if self.cfg.detections == nil { return; } + let last_error = ""; + let lines = last_lines(log_path, 50); + if lines != nil { + let i = len(lines) - 1; + while i >= 0 { + let clean = sanitize(lines[i]); + if classify_loose(clean) == "error" { + last_error = clean; + break; + } + i = i - 1; + } + } + let line = json.encode(Detection { ts: time.iso(now), path: log_path, + source: source, event: event, lastError: last_error }); + try fs.append(self.cfg.detections, line + "\n") catch (e) { + print("DETECTIONS-SINK-ERROR ${self.cfg.detections}: ${e.msg}"); + } + } + + fn signature() -> Text { + if fs.exists(self.cron_dir) == false { return ""; } + let names = try fs.list(self.cron_dir) catch (e) nil; + if names == nil { return ""; } + sort(names); + let parts: multi Text = []; + for n in names { + if index_of(n, ".") >= 0 { continue; } + let st = fs.stat("${self.cron_dir}/${n}"); + if st == nil { continue; } + if st.dir { continue; } + push(parts, "${n}:${st.mtime}"); + } + return join(parts, "|"); + } + + -- The daemon loop the systems track legalizes: blocking sleep in program + -- mode; SIGTERM/SIGINT set the runtime's flag polled via env.stopping() + -- (spec Part 2 — no signal callbacks). Haxe's bare while(true) gains a + -- clean shutdown for free. + fn run() -> Int { + while true { + if env.stopping() { return 0; } + self.tick(time.now()); + time.sleep(250); + } + } +} + +fn compute_next(schedules: multi Text, now: Int) -> ?Int { + let best: ?Int = nil; + for expr in schedules { + let d = next_fire(expr, now); + if d == nil { continue; } + if best == nil or d < best { best = d; } + } + return best; +} diff --git a/docs/examples/log-watcher/watcher.wo b/docs/examples/log-watcher/watcher.wo new file mode 100644 index 0000000..a7fccb7 --- /dev/null +++ b/docs/examples/log-watcher/watcher.wo @@ -0,0 +1,64 @@ +-- watcher.wo — Watcher.hx, file for file. One watched log. Service mode +-- (live=true) reports ALERT/CLEAR transitions continuously; cron mode is +-- driven by the supervisor via done()/result(). + +-- Watcher.hx:3-7 — the Haxe enum is a tagged union. +type CronResult = Ok | ErrorFinal | Miss + -- Ok: completed, last entry not error + -- ErrorFinal: completed with error as the last entry -> alert + -- Miss: log never produced content during the window + +class Watcher { + -- Haxe `(default, null)` — public read, owner-only write — is `pub(read)`. + pub(read) path: Text + pub(read) alerted: Bool = false + pub(read) saw_content: Bool = false + + -- The constructor body (Watcher.hx:23-30) is gone: `st` self-initializes + -- through TailState's field defaults; construction is the brace literal. + st: TailState = TailState {} + quiet_ms: Int -- quietPeriod, seconds -> ms + poll_ms: Int -- pollInterval, seconds -> ms + live: Bool + activated_at: Int -- ms + last_poll_at: Int = 0 -- 0 = never polled; first due() is true + -- (Haxe used NEGATIVE_INFINITY) + + fn due(now: Int) -> Bool { + return now - self.last_poll_at >= self.poll_ms; + } + + fn tick(now: Int) { + self.last_poll_at = now; + let r = poll(self.path, self.st, now); -- exclusive borrow of st for the call + if r.new_lines > 0 { self.saw_content = true; } + if self.live == false { return; } + -- the alert rule (README): last entry is error and nothing follows + if self.alerted == false and self.saw_content and self.st.last_level == "error" and now - self.st.last_newline_at >= self.quiet_ms { + self.alerted = true; + print("ALERT ${self.path}: last entry is error, quiet for ${self.quiet_ms / 1000}s"); + return; + } + if self.alerted and self.st.last_level != "error" { + self.alerted = false; + print("CLEAR ${self.path}: new ${self.st.last_level} entries arrived"); + } + } + + -- cron mode: the job is judged complete after a quiet period — from the + -- last write if the log produced content, else from activation (a miss) + fn done(now: Int) -> Bool { + if self.saw_content { return now - self.st.last_newline_at >= self.quiet_ms; } + return now - self.activated_at >= self.quiet_ms; + } + + fn result() -> CronResult { + if self.saw_content == false { return Miss; } + if self.st.last_level == "error" { return ErrorFinal; } + return Ok; + } + + fn last_level() -> Text { + return self.st.last_level; + } +} diff --git a/docs/examples/log-watcher/wo.toml b/docs/examples/log-watcher/wo.toml new file mode 100644 index 0000000..7e0caa7 --- /dev/null +++ b/docs/examples/log-watcher/wo.toml @@ -0,0 +1,11 @@ +name = "log-watcher" +version = "0.1.0" +description = "The Haxe log-watcher ported file-for-file: the systems-track sample workload (program mode + fs/proc/net/time/json)" + +# A free `fn main` makes this a PROGRAM (systems-track spec Part 2): one +# shard, blocking builtins legal, exit code = main's return value. No [app] +# listen — the mcp subcommand takes its port from config.json, like the +# original. + +[runtime] +wo = ">= 0.1"