docs: log-watcher program

This commit is contained in:
shoney.arickathil 2026-08-10 09:26:07 +02:00
parent 5a43210448
commit 2171b2d94b
9 changed files with 1434 additions and 0 deletions

View file

@ -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:<port>
```
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.

View file

@ -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;
}
}

View file

@ -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;
}

View file

@ -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 <logfile> [quietPeriod] [pollInterval] continuous service watch");
print_err(" log-watcher run <cron.d-dir> [config.json] supervisor (cron + services)");
print_err(" log-watcher mcp <cron.d-dir> <config.json> 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;
}

View file

@ -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<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 }; }

View file

@ -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;
}
}

View file

@ -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<Text, CronWatch> = {}
pub(read) active: map<Text, Watcher> = {}
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<Text, CronWatch> = {};
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;
}

View file

@ -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;
}
}

View file

@ -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"