code wiki / _hdl_build / nx_pub_recover_pulse.nx

nx_pub_recover_pulse.nx source

↩ module page · 123 lines · 6592 B

1// nx_pub_recover_pulse.nx -- THE PUBLISHER'S SELF-HEALING PULSE (composes the gated decide + execute rungs). 2// 3// One pass = decide (scan the monitor's metrics -> recovery requests, via a DATA-DRIVEN map url-substring-> 4// daemon, rule #11 not hardcoded) then execute (kill the wedged daemon; the existing nx_hostctl keeper 5// respawns a fresh one). Run on a cadence by a scheduler (one pass per invocation; the daemon-loop wrapper is 6// the next rung). This is the Publisher OWNING ops/restarts end-to-end -- NO direct/ad-hoc workstream action, 7// every step a gated Publisher rung. 8// pulse_pass(metricspath, queuepath, ledgerpath, mappath, wedge_ms) -> recoveries executed 9// license_tier: ORIGINAL 10import "nx_syscalls.nx" 11import "nx_pub_recover_exec.nx" // pe_exec + (transitively) pr_field / pr_has / pr_already_pending / pr_n / pr_s / fl_acquire / fl_release / fa_appendz / proc_* 12const PP_MAGIC_100000: i64 = 100000 13const PP_MAGIC_5000: i64 = 5000 14 15const PP_METRICS: *u8 = "knowledge/status/apm_metrics.tsv" 16const PP_QUEUE: *u8 = "knowledge/publish/recovery_queue.tsv" 17const PP_LEDGER: *u8 = "knowledge/publish/recovery_ledger.tsv" 18const PP_MAP: *u8 = "knowledge/publish/recover_map.tsv" // url-substring<TAB>daemon (data-driven, no hardcode) 19const PP_LOCK: *u8 = "publish_recover" 20const PP_REC_CAP: i64 = 1024 21const PP_COOLDOWN_SEC: i64 = 300 // never kill-storm: at most PP_MAX_RECOVERIES of a daemon per this window 22const PP_MAX_RECOVERIES: i64 = 2 // (#26 never-brick: a daemon that keeps wedging -> BACK OFF, don't re-kill) 23 24func pp_strlen(s: *u8) -> i64 { var n: i64=0; while s[n]!=(0 as u8){n=n+1} return n } 25// crash-loop guard: how many times was `daemon` recovered in [now-window, now] per the recovery ledger 26// (fmt: epoch<TAB>daemon<TAB>killed<TAB>KILLED -> field0=epoch field1=daemon). 27func pp_recent_recoveries(ledgerpath: *u8, daemon: *u8, now: i64, window: i64) -> i64 { 28 let lp: *i64 = sys_mmap(16) as *i64; lp[0] = 0 29 let b: *u8 = sys_read_file(ledgerpath, lp) 30 if (b as i64) == 0 { return 0 } 31 let n: i64 = lp[0] 32 let f0: *u8 = sys_mmap(32); let f1: *u8 = sys_mmap(128) 33 var count: i64 = 0; var i: i64 = 0 34 while i < n { 35 var e: i64 = i; while e < n { if b[e] == (10 as u8) { break } else { e = e + 1 } } 36 let ll: i64 = e - i 37 if ll > 0 { 38 pr_field(((b as i64)+i) as *u8, ll, 0, f0) 39 pr_field(((b as i64)+i) as *u8, ll, 1, f1) 40 var ep: i64=0; var k: i64=0; while f0[k]!=(0 as u8){ let c: i64=f0[k] as i64; if c>=48 { if c<=57 { ep=ep*10+(c-48) } } k=k+1 } 41 if ep >= now - window { if pr_streq(f1, daemon) == 1 { count = count + 1 } } 42 } 43 i = e + 1 44 } 45 return count 46} 47 48// data-driven url->daemon: first map line whose field0 (url-substring) is in the url -> field1 (daemon). 0 if none. 49func pp_map_lookup(url: *u8, urllen: i64, mappath: *u8, out: *u8) -> i64 { 50 let lp: *i64 = sys_mmap(16) as *i64; lp[0] = 0 51 let b: *u8 = sys_read_file(mappath, lp) 52 if (b as i64) == 0 { return 0 } 53 let n: i64 = lp[0] 54 let sub: *u8 = sys_mmap(256); let dae: *u8 = sys_mmap(128) 55 var i: i64 = 0 56 while i < n { 57 var e: i64 = i; while e < n { if b[e] == (10 as u8) { break } else { e = e + 1 } } 58 let ll: i64 = e - i 59 if ll > 0 { 60 if b[i] != (35 as u8) { // skip '#' comment lines 61 pr_field(((b as i64)+i) as *u8, ll, 0, sub) 62 pr_field(((b as i64)+i) as *u8, ll, 1, dae) 63 if pr_has(url, urllen, sub) == 1 { var k: i64=0; while dae[k]!=(0 as u8){out[k]=dae[k];k=k+1} out[k]=0 as u8; return 1 } 64 } 65 } 66 i = e + 1 67 } 68 return 0 69} 70 71// DECIDE (data-driven): scan metrics (field2=url, field9=max_ms); wedged -> map -> idempotent serialized request. 72func pp_decide(metricspath: *u8, queuepath: *u8, ledgerpath: *u8, mappath: *u8, wedge_ms: i64, now: i64) -> i64 { 73 let lp: *i64 = sys_mmap(16) as *i64; lp[0] = 0 74 let b: *u8 = sys_read_file(metricspath, lp) 75 if (b as i64) == 0 { return 0 } 76 let n: i64 = lp[0] 77 let url: *u8 = sys_mmap(512); let maxf: *u8 = sys_mmap(32); let daemon: *u8 = sys_mmap(128) 78 var emitted: i64 = 0; var i: i64 = 0 79 while i < n { 80 var e: i64 = i; while e < n { if b[e] == (10 as u8) { break } else { e = e + 1 } } 81 let ll: i64 = e - i 82 if ll > 0 { 83 pr_field(((b as i64)+i) as *u8, ll, 2, url) 84 pr_field(((b as i64)+i) as *u8, ll, 9, maxf) 85 var mx: i64 = 0; var mi: i64 = 0; while maxf[mi]!=(0 as u8){ let c: i64=maxf[mi] as i64; if c>=48 { if c<=57 { mx=mx*10+(c-48) } } mi=mi+1 } 86 if mx >= wedge_ms { 87 let urllen: i64 = 0 + pp_strlen(url) 88 if pp_map_lookup(url, urllen, mappath, daemon) == 1 { 89 if pp_recent_recoveries(ledgerpath, daemon, now, PP_COOLDOWN_SEC) < PP_MAX_RECOVERIES { 90 let lk: i64 = fl_acquire(PP_LOCK, PP_MAGIC_100000, 1) 91 if pr_already_pending(queuepath, daemon) == 0 { 92 let rec: *u8 = sys_mmap(PP_REC_CAP); var o: i64 = 0 93 o = pr_s(rec, o, "PENDING" as *u8); rec[o]=9 as u8; o=o+1 94 o = pr_n(rec, o, now); rec[o]=9 as u8; o=o+1 95 o = pr_s(rec, o, daemon); rec[o]=9 as u8; o=o+1 96 o = pr_s(rec, o, "wedge" as *u8); rec[o]=9 as u8; o=o+1 97 o = pr_n(rec, o, mx); rec[o]=10 as u8; o=o+1 98 fa_appendz(queuepath, rec, PP_REC_CAP) 99 emitted = emitted + 1 100 } 101 fl_release(lk) 102 } 103 } 104 } 105 } 106 i = e + 1 107 } 108 return emitted 109} 110 111// ONE PULSE PASS = decide then execute. Returns #recoveries executed. 112func pulse_pass(metricspath: *u8, queuepath: *u8, ledgerpath: *u8, mappath: *u8, wedge_ms: i64) -> i64 { 113 pp_decide(metricspath, queuepath, ledgerpath, mappath, wedge_ms, sys_now_realtime_sec()) 114 return pe_exec(queuepath, ledgerpath) 115} 116 117func pp_uw(s: *u8) -> i64 { var n: i64=0; while s[n]!=(0 as u8){n=n+1} sys_write(1,s,n); return 0 } 118func main(argc: i64, argv: *i64) -> i64 { 119 pp_uw("=== nishi-publisher self-healing pulse (decide+execute) ===\n" as *u8) 120 let c: i64 = pulse_pass(PP_METRICS, PP_QUEUE, PP_LEDGER, PP_MAP, PP_MAGIC_5000) 121 pp_uw("recoveries this pass: " as *u8); let nb: *u8=sys_mmap(24); let e: i64=pr_n(nb,0,c); sys_write(1,nb,e); pp_uw("\n" as *u8) 122 return 0 123}