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}