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