docs/examples/log-watcher now COMPILES AND RUNS: `wovm lw.wob watch app.log 2 1`
tails a live file, classifies levels and fires its alert
("last entry is error, quiet for 2s"). corpus 71/0, woc 565/0, wovm gates green.
- program mode: the entry is `fn main` taking nothing or one `multi Text`;
runtime/src/main.c builds that list from the program's own arguments (not
the program name, not the image path) and the entry's return value is the
process exit code (low byte); the loader accepts a 0- or 1-arg free-fn entry
- `+` on Text is now WO-E201 pointing at `..`. This was a memory-safety hole,
not a style nit: the emitter lowered it to ADD on two heap pointers, and the
workload's own `out = out + char_of(c)` produced a wild pointer that
segfaulted the VM inside starts_with. Reported off confident types only
- `x == nil` / `x != nil` lower to EQ (a word compare), never EQS: nil is the
zero word and EQS dereferences its operands, so a nil guard would trap
instead of answering
- docs/examples/log-watcher: seven `+`-on-Text sites corrected to `..`
(logtail sanitize, mcp header/body/carry assembly, supervisor detections
line) — the sample was carrying the Haxe habit, and the language reserves
`+` for arithmetic by doctrine
- wovm CLI takes arguments after the image path (`wovm <file.wob> [args...]`)
Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
248 lines
9 KiB
Text
248 lines
9 KiB
Text
-- 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;
|
|
}
|