nx_sweep_daemon.nx source
↩ module page · 239 lines · 12980 B
1// nx_sweep_daemon.nx -- STANDING SWEEP daemon, CLI half (tools/ops surface). Runs the data-driven check
2// registry on an interval, NEVER hanging (each check is deadline-killed by nx_sweep_core), serves a
3// PLAIN-ENGLISH /status page, and dedups alerts (one line on the 0->1 red transition, silent while it
4// stays red -- captured to the supervisor log via stdout). Supervised by nx_daemon_supervisor (a
5// daemons.reg row, revive-armed) so it never silently dies. Verbs:
6// once [reg] run every check once, print verdicts, exit 0 iff all PASS (gate/manual/CI)
7// serve <port> [reg] bind :port, sweep every interval, serve /status (the standing daemon)
8// Default registry: knowledge/sweep_checks.reg. license_tier: ORIGINAL
9import "nx_sweep_daemon_lib.nx"
10
11const SV_ARG_VERB: i64 = 1
12const SV_ARG_PORT: i64 = 2 // serve: port
13const SV_ARG_SREG: i64 = 3 // serve: optional registry
14const SV_ARG_IVMS: i64 = 4 // serve: optional interval ms (default SV_INTERVAL_MS)
15const SV_ARG_OREG: i64 = 2 // once: optional registry
16const SV_ARGC_SERVE: i64 = 3 // serve port
17const SV_RC_USAGE: i64 = 2
18const SV_INTERVAL_MS: i64 = 300000 // 5 min between sweeps (a threshold -> named, rule 11)
19// ---- fork-child /status architecture: a forked child ALWAYS serves the last-published snapshot from a
20// seqlock-shared page, so a slow (multi-second) sweep in the PARENT never starves /status. shq[0]=seq
21// (odd=writer active), shq[1]=body len, shm[SV_BODY_OFF..]=body. This is the daemon-supervisor's own pattern.
22const SV_SHM_SZ: i64 = 262144 // shared snapshot page (holds the rendered status HTML)
23const SV_SEQ_I: i64 = 0 // seqlock counter (i64 index into shm)
24const SV_BLEN_I: i64 = 1 // published body length (i64 index)
25const SV_BODY_OFF: i64 = 16 // body bytes begin after the two i64 header slots
26const SV_SEQ_TRIES: i64 = 100000 // seqlock read retry cap (writer-starvation backstop)
27const SV_RESP_EXTRA: i64 = 256 // HTTP header headroom over the body
28const SV_SEQ_MOD: i64 = 2 // seqlock parity: even=committed, odd=writer-active
29const SV_RC_CHILD: i64 = 3 // status-child fatal-bind exit code
30const SV_REQ_CAP: i64 = 8192
31const SV_AF_INET: i64 = 2
32const SV_SOCK_STREAM: i64 = 1
33const SV_SOL_SOCKET: i64 = 1
34const SV_SO_REUSEADDR: i64 = 2
35const SV_PTR1: i64 = 8 // single-i64 scratch
36const SV_OPTLEN: i64 = 4 // setsockopt optlen (int)
37const SV_SADDR_LEN: i64 = 16
38const SV_BACKLOG: i64 = 16
39const SV_PORT_HI_SH: i64 = 8
40const SV_BYTE: i64 = 0xFF
41const SV_ADDR_FAM: i64 = 0 // sockaddr offset: family
42const SV_ADDR_PORT: i64 = 2 // sockaddr offset: port (big-endian)
43const SV_DEC: i64 = 10
44const SV_ASCII_0: i64 = 48
45const SV_ASCII_9: i64 = 57
46
47func sv_puts(s: *u8) -> i64 { var n: i64 = 0; while s[n] != (0 as u8) { n = n + 1 } sys_write(1, s, n); return 0 }
48func sv_seq(a: *u8, b: *u8) -> i64 { var i: i64 = 0; while a[i] != (0 as u8) { if a[i] != b[i] { return 0 } i = i + 1 } if b[i] != (0 as u8) { return 0 } return 1 }
49func sv_atoi(s: *u8) -> i64 {
50 var v: i64 = 0; var i: i64 = 0
51 while s[i] != (0 as u8) { let c: i64 = s[i] as i64; if c < SV_ASCII_0 { return v } if c > SV_ASCII_9 { return v } v = v * SV_DEC + (c - SV_ASCII_0); i = i + 1 }
52 return v
53}
54// print one verdict line for `once`
55func sv_print_verdict(ctx: *i64, i: i64) -> i64 {
56 let a_name: *i64 = sc_arr(ctx, SC_NAME)
57 let a_v: *i64 = sc_arr(ctx, SC_VERDICT)
58 let a_ms: *i64 = sc_arr(ctx, SC_MS)
59 let v: i64 = a_v[i]
60 sv_puts(" " as *u8); sv_puts(a_name[i] as *u8); sv_puts(" -> " as *u8)
61 if v == SW_V_PASS { sv_puts("PASS" as *u8) } else {
62 if v == SW_V_TIMEOUT { sv_puts("TIMEOUT" as *u8) } else {
63 if v == SW_V_ERROR { sv_puts("ERROR" as *u8) } else { sv_puts("FAIL" as *u8) } } }
64 sv_puts(" (" as *u8)
65 var m: i64 = a_ms[i]; if m == 0 { sv_puts("0" as *u8) } else {
66 let d: *u8 = sys_mmap(32); var k: i64 = 0
67 while m > 0 { d[k] = (SV_ASCII_0 + (m % SV_DEC)) as u8; m = m / SV_DEC; k = k + 1 }
68 while k > 0 { sys_write(1, ((d as i64) + k - 1) as *u8, 1); k = k - 1 }
69 }
70 sv_puts("ms)\n" as *u8)
71 return 0
72}
73// alert-dedup: emit ONE line the cycle a check first goes bad (red==1); silent while it stays bad.
74func sv_alerts(ctx: *i64) -> i64 {
75 let a_name: *i64 = sc_arr(ctx, SC_NAME)
76 let a_v: *i64 = sc_arr(ctx, SC_VERDICT)
77 let a_red: *i64 = sc_arr(ctx, SC_RED)
78 var i: i64 = 0
79 while i < ctx[SC_N] {
80 if a_v[i] != SW_V_PASS { if a_red[i] == 1 {
81 sv_puts("ALERT: sweep check '" as *u8); sv_puts(a_name[i] as *u8); sv_puts("' FAILED\n" as *u8)
82 } }
83 i = i + 1
84 }
85 return 0
86}
87// build "HTTP/1.1 200 OK ... <body>" into resp; returns length. op is a caller-owned cell (reused across
88// requests -- no per-request mmap).
89func sv_http_wrap(resp: *u8, cap: i64, body: *u8, blen: i64, op: *i64) -> i64 {
90 op[0] = 0
91 sd_app(resp, op, cap, "HTTP/1.1 200 OK\r\nContent-Type: text/html; charset=utf-8\r\nContent-Length: " as *u8)
92 sd_appn(resp, op, cap, blen)
93 sd_app(resp, op, cap, "\r\nConnection: close\r\n\r\n" as *u8)
94 var i: i64 = 0
95 var o: i64 = op[0]
96 while i < blen { if o < cap - 1 { resp[o] = body[i]; o = o + 1 } i = i + 1 }
97 return o
98}
99
100// seqlock WRITE (parent): publish html[0..hl) into the shared page. Odd seq during copy = tear-signal.
101func sv_shm_write(shm: *u8, html: *u8, hl: i64) -> i64 {
102 let shq: *i64 = shm as *i64
103 var bl: i64 = hl
104 if bl > SV_SHM_SZ - SV_BODY_OFF { bl = SV_SHM_SZ - SV_BODY_OFF }
105 shq[SV_SEQ_I] = shq[SV_SEQ_I] + 1 // -> odd: writer active
106 var i: i64 = 0
107 while i < bl { shm[SV_BODY_OFF + i] = html[i]; i = i + 1 }
108 shq[SV_BLEN_I] = bl
109 shq[SV_SEQ_I] = shq[SV_SEQ_I] + 1 // -> even: snapshot committed
110 return 0
111}
112// build a sockaddr_in for 0.0.0.0:port (big-endian port) into a 16-byte buffer
113func sv_sockaddr(addr: *u8, port: i64) -> i64 {
114 var z: i64 = 0; while z < SV_SADDR_LEN { addr[z] = 0 as u8; z = z + 1 }
115 addr[SV_ADDR_FAM] = SV_AF_INET as u8
116 addr[SV_ADDR_PORT] = ((port >> SV_PORT_HI_SH) & SV_BYTE) as u8
117 addr[SV_ADDR_PORT + 1] = (port & SV_BYTE) as u8
118 return 0
119}
120// the /status child: bind+listen+accept on `port`, seqlock-READ the shared snapshot, serve it as text/html.
121// ALWAYS responsive (the parent's slow sweep runs in a different process). Returns only on a fatal bind error.
122func sv_status_child(shm: *u8, port: i64) -> i64 {
123 let shq: *i64 = shm as *i64
124 let fd: i64 = sys_socket(SV_AF_INET, SV_SOCK_STREAM, 0)
125 if fd < 0 { return 1 }
126 let opt: *i64 = sys_mmap(SV_PTR1) as *i64; opt[0] = 1
127 sys_setsockopt(fd, SV_SOL_SOCKET, SV_SO_REUSEADDR, opt as *u8, SV_OPTLEN)
128 let addr: *u8 = sys_mmap(SV_SADDR_LEN)
129 sv_sockaddr(addr, port)
130 if sys_bind(fd, addr, SV_SADDR_LEN) < 0 { return 1 }
131 if sys_listen(fd, SV_BACKLOG) < 0 { return 1 }
132 let req: *u8 = sys_mmap(SV_REQ_CAP)
133 let body: *u8 = sys_mmap(SV_SHM_SZ)
134 let resp: *u8 = sys_mmap(SV_SHM_SZ + SV_RESP_EXTRA)
135 let wop: *i64 = sys_mmap(SV_PTR1) as *i64
136 var run: i64 = 1
137 while run == 1 {
138 let cfd: i64 = sys_accept(fd)
139 if cfd < 0 { } else {
140 sys_read(cfd, req, SV_REQ_CAP)
141 // seqlock read: copy body only when seq is even AND unchanged across the copy (tear-free)
142 var blen: i64 = 0
143 var got: i64 = 0
144 var tries: i64 = 0
145 while got == 0 {
146 let s1: i64 = shq[SV_SEQ_I]
147 if (s1 % SV_SEQ_MOD) == 0 {
148 var bl: i64 = shq[SV_BLEN_I]
149 if bl > SV_SHM_SZ - SV_BODY_OFF { bl = SV_SHM_SZ - SV_BODY_OFF }
150 if bl < 0 { bl = 0 }
151 var ci: i64 = 0
152 while ci < bl { body[ci] = shm[SV_BODY_OFF + ci]; ci = ci + 1 }
153 if shq[SV_SEQ_I] == s1 { blen = bl; got = 1 }
154 }
155 tries = tries + 1
156 if tries > SV_SEQ_TRIES { got = 1 }
157 }
158 let rl: i64 = sv_http_wrap(resp, SV_SHM_SZ + SV_RESP_EXTRA, body, blen, wop)
159 sys_write(cfd, resp, rl)
160 sys_close(cfd)
161 }
162 }
163 return 0
164}
165
166func main(argc: i64, argv: *i64) -> i64 {
167 if argc < SV_ARGC_SERVE - 1 { sv_puts("usage: nx_sweep_daemon once [reg] | serve <port> [reg]\n" as *u8); return SV_RC_USAGE }
168 let verb: *u8 = argv[SV_ARG_VERB] as *u8
169
170 if sv_seq(verb, "mtime" as *u8) == 1 {
171 // diagnostic: print sd_mtime(path) -- compare to `stat -c %Y` to validate the stat offset on this host
172 if argc <= SV_ARG_OREG { sv_puts("usage: nx_sweep_daemon mtime <path>\n" as *u8); return SV_RC_USAGE }
173 let ctx: *i64 = sd_new()
174 let m: i64 = sd_mtime(ctx, argv[SV_ARG_OREG] as *u8)
175 sv_puts("mtime=" as *u8)
176 let mb: *u8 = sys_mmap(SV_REQ_CAP); let mo: *i64 = sys_mmap(SV_PTR1) as *i64; mo[0] = 0
177 sd_appn(mb, mo, SV_REQ_CAP, m); sys_write(1, mb, mo[0]); sv_puts("\n" as *u8)
178 return 0
179 }
180
181 if sv_seq(verb, "once" as *u8) == 1 {
182 var reg: *u8 = "knowledge/sweep_checks.reg" as *u8
183 if argc > SV_ARG_OREG { reg = argv[SV_ARG_OREG] as *u8 }
184 let ctx: *i64 = sd_new()
185 if sd_load(ctx, reg) < 0 { sv_puts("NX-SWEEP: cannot read registry " as *u8); sv_puts(reg); sv_puts("\n" as *u8); return 1 }
186 sv_puts("=== NX-SWEEP once (" as *u8); sv_puts(reg); sv_puts(") ===\n" as *u8)
187 let bad: i64 = sd_run_all(ctx)
188 var i: i64 = 0
189 while i < ctx[SC_N] { sv_print_verdict(ctx, i); i = i + 1 }
190 sv_alerts(ctx)
191 if bad == 0 { sv_puts("SWEEP verdict=GREEN (all checks pass)\n" as *u8); return 0 }
192 sv_puts("SWEEP verdict=RED\n" as *u8)
193 return 1
194 }
195
196 if sv_seq(verb, "serve" as *u8) == 1 {
197 if argc < SV_ARGC_SERVE { sv_puts("usage: nx_sweep_daemon serve <port> [reg] [interval_ms]\n" as *u8); return SV_RC_USAGE }
198 let port: i64 = sv_atoi(argv[SV_ARG_PORT] as *u8)
199 var reg: *u8 = "knowledge/sweep_checks.reg" as *u8
200 if argc > SV_ARG_SREG { reg = argv[SV_ARG_SREG] as *u8 }
201 var interval: i64 = SV_INTERVAL_MS
202 if argc > SV_ARG_IVMS { let iv: i64 = sv_atoi(argv[SV_ARG_IVMS] as *u8); if iv > 0 { interval = iv } }
203 let ctx: *i64 = sd_new()
204 if sd_load(ctx, reg) < 0 { sv_puts("NX-SWEEP: cannot read registry\n" as *u8); return 1 }
205 ctx[SC_MTIME] = sd_mtime(ctx, reg) // seed the hot-reload gate with the loaded file's mtime
206 // shared snapshot page (inherited by the forked child); publish an INITIAL render so /status is
207 // answerable the instant the child binds -- before the first (blocking) sweep even runs.
208 let shm: *u8 = sys_mmap_shared(SV_SHM_SZ)
209 let shq: *i64 = shm as *i64
210 shq[SV_SEQ_I] = 0
211 shq[SV_BLEN_I] = 0
212 let html: *u8 = sys_mmap(SV_SHM_SZ)
213 sv_shm_write(shm, html, sd_render_status(ctx, html, SV_SHM_SZ))
214 // fork the /status child: it binds the port + serves the shared snapshot, ALWAYS responsive. The
215 // parent never touches the socket -> a multi-second sweep can never starve /status or race the
216 // supervisor health probe (retires the warmup + boot-blocking-accept compromise entirely).
217 let cpid: i64 = sys_fork()
218 if cpid == 0 { sv_status_child(shm, port); sys_exit(SV_RC_CHILD); return SV_RC_CHILD }
219 if cpid < 0 { sv_puts("NX-SWEEP: fork fail\n" as *u8); return 1 }
220 sv_puts("nx_sweep_daemon: /status child forked; parent sweeping every interval, publishing to shm\n" as *u8)
221 // PARENT: the sweep loop. Run checks, alert, render, PUBLISH to shm, sleep. No socket here.
222 let cst: *i64 = sys_mmap(SV_PTR1) as *i64 // child wait-status cell (hoisted; no per-cycle mmap)
223 var go: i64 = 1
224 while go == 1 {
225 sd_reload_if_changed(ctx, reg) // HOT-RELOAD, mtime-gated (proven live 1->2 on append)
226 sd_run_all(ctx) // deadline-killed checks -- never hangs
227 sv_alerts(ctx) // dedup: one line on the 0->1 red transition
228 sd_sample(ctx, "knowledge/sweep_metrics.ring" as *u8, sd_now_sec(ctx), SD_UPWINDOW_S) // producer writes its OWN samples inline (TSDB history + uptime%); leak-free now()
229 sv_shm_write(shm, html, sd_render_status(ctx, html, SV_SHM_SZ))
230 // if the /status child is gone (e.g. its bind lost to another instance already holding the port),
231 // THIS parent is redundant -> exit cleanly rather than orphan a sweeper with no reader.
232 if sys_wait4(cpid, cst, WNOHANG) == cpid { sv_puts("nx_sweep_daemon: /status child gone -> exiting redundant instance\n" as *u8); return 0 }
233 sys_sleep_ms(interval)
234 }
235 return 0
236 }
237 sv_puts("usage: nx_sweep_daemon once [reg] | serve <port> [reg]\n" as *u8)
238 return SV_RC_USAGE
239}