code wiki / (root) / nx_sweep_daemon.nx

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}