code wiki / _hdl_build / nx_pub_recover.nx

nx_pub_recover.nx source

↩ module page · 140 lines · 7809 B

1// nx_pub_recover.nx -- THE PUBLISHER'S DAEMON-RECOVERY rung (ops authority, not ad-hoc workstreams). 2// 3// Operator law: "the nishi publisher should be managing restarts and all this stuff." So recovery is a 4// PUBLISHER responsibility, decided from the MONITOR'S signal -- never an ad-hoc workstream kill. 5// 6// The supervisor (nx_hostctl) already RESPAWNS dead/hung daemons + crash-loop-guards. The GAP it can't see: 7// a WEDGE -- a daemon that still answers but SLOWLY (the :8791 reader at 20s/request passes the binary 8// up/down probe). nx_apm_sweep measures real latency -> apm_metrics.tsv. This organ reads that signal, 9// maps a slow endpoint to its daemon, and emits an IDEMPOTENT + SERIALIZED + LEDGERED recovery REQUEST 10// (the control-plane consumer proc_kill_by_name's it -> the existing keeper respawns a fresh one). 11// 12// By construction: never-brick #26 (a recovery REQUEST is reversible metadata; the action is only a 13// kill+supervised-respawn, never a destructive write); idempotent #10 (a daemon already PENDING is not 14// re-requested); serialized (fl_acquire, one decision at a time); fail-safe #14 (unmapped/healthy -> no-op). 15// pr_recover(metricspath, queuepath, wedge_ms) -> number of recovery requests emitted 16// license_tier: ORIGINAL 17import "nx_syscalls.nx" 18import "nx_arbiter.nx" // fl_acquire / fl_release 19import "nx_framed_append.nx" // fa_appendz (O_APPEND + LOCK_EX, torn-free) 20const PR_MAGIC_100000: i64 = 100000 21 22const PR_WEDGE_MS_DEFAULT: i64 = 5000 23const PR_METRICS: *u8 = "knowledge/status/apm_metrics.tsv" 24const PR_QUEUE: *u8 = "knowledge/publish/recovery_queue.tsv" 25const PR_LOCK: *u8 = "publish_recover" 26const PR_REC_CAP: i64 = 1024 27 28func pr_slen(s: *u8) -> i64 { var n: i64 = 0; while s[n] != (0 as u8) { n = n + 1 } return n } 29func pr_uw(fd: i64, s: *u8) -> i64 { sys_write(fd, s, pr_slen(s)); return 0 } 30func pr_n(buf: *u8, off: i64, v: i64) -> i64 { 31 var o: i64 = off; var m: i64 = v 32 if m < 0 { buf[o] = 45 as u8; o = o + 1; m = 0 - m } 33 let t: *u8 = sys_mmap(24); var k: i64 = 0 34 if m == 0 { t[0] = 48 as u8; k = 1 } 35 while m > 0 { t[k] = (48 + (m - (m/10)*10)) as u8; m = m / 10; k = k + 1 } 36 var i: i64 = 0; while i < k { buf[o + i] = t[k - 1 - i]; i = i + 1 } 37 return o + k 38} 39func pr_s(buf: *u8, off: i64, s: *u8) -> i64 { var o: i64 = off; var i: i64 = 0; while s[i] != (0 as u8) { buf[o] = s[i]; o = o + 1; i = i + 1 } return o } 40func pr_wn(fd: i64, v: i64) -> i64 { let b: *u8 = sys_mmap(24); let e: i64 = pr_n(b, 0, v); sys_write(fd, b, e); return 0 } 41// substring present in hay[0..n)? 42func pr_has(hay: *u8, n: i64, needle: *u8) -> i64 { 43 let nn: i64 = pr_slen(needle); if nn == 0 { return 0 } 44 var i: i64 = 0 45 while i + nn <= n { var j: i64 = 0; var ok: i64 = 1; while j < nn { if hay[i+j] != needle[j] { ok = 0; j = nn } else { j = j + 1 } } if ok == 1 { return 1 } i = i + 1 } 46 return 0 47} 48// extract TAB-field #idx from line[0..linelen) into out (NUL-terminated); returns its length. 49func pr_field(line: *u8, linelen: i64, idx: i64, out: *u8) -> i64 { 50 var f: i64 = 0; var i: i64 = 0; var o: i64 = 0 51 while i < linelen { 52 let c: u8 = line[i] 53 if c == (9 as u8) { if f == idx { out[o] = 0 as u8; return o } f = f + 1 } else { if f == idx { out[o] = c; o = o + 1 } } 54 i = i + 1 55 } 56 out[o] = 0 as u8; return o 57} 58// parse a leading decimal int from a NUL-terminated string. 59func pr_atoi(s: *u8) -> i64 { var v: i64 = 0; var i: i64 = 0; while s[i] != (0 as u8) { let c: i64 = s[i] as i64; if c >= 48 { if c <= 57 { v = v*10 + (c-48) } } i = i + 1 } return v } 60// equal NUL-terminated strings? 61func pr_streq(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 } 62 63// map a probed URL to the daemon binary name nx_hostctl proc_kill_by_name's. returns 1 + name in out, else 0. 64func pr_daemon_for(url: *u8, urllen: i64, out: *u8) -> i64 { 65 if pr_has(url, urllen, "/library" as *u8) == 1 { let o: i64 = pr_s(out, 0, "nx_media_server_auth.elf" as *u8); out[o]=0 as u8; return 1 } 66 if pr_has(url, urllen, "/media" as *u8) == 1 { let o: i64 = pr_s(out, 0, "nx_media_server_auth.elf" as *u8); out[o]=0 as u8; return 1 } 67 if pr_has(url, urllen, "/gallery" as *u8) == 1 { let o: i64 = pr_s(out, 0, "nx_gallery_gateway.elf" as *u8); out[o]=0 as u8; return 1 } 68 if pr_has(url, urllen, "/wiki" as *u8) == 1 { let o: i64 = pr_s(out, 0, "nx_wiki_gw.elf" as *u8); out[o]=0 as u8; return 1 } 69 out[0] = 0 as u8; return 0 // landing/status = sites.elf, owned by hosting -> do NOT auto-request (coordinate) 70} 71 72// is there already a PENDING recovery request for this daemon in the queue? (idempotency #10) 73func pr_already_pending(queuepath: *u8, daemon: *u8) -> i64 { 74 let lp: *i64 = sys_mmap(16) as *i64; lp[0] = 0 75 let b: *u8 = sys_read_file(queuepath, lp) 76 if (b as i64) == 0 { return 0 } 77 let n: i64 = lp[0] 78 let f0: *u8 = sys_mmap(64); let f2: *u8 = sys_mmap(128) 79 var i: i64 = 0 80 while i < n { 81 var e: i64 = i; while e < n { if b[e] == (10 as u8) { break } else { e = e + 1 } } 82 let ll: i64 = e - i 83 if ll > 0 { 84 pr_field(((b as i64)+i) as *u8, ll, 0, f0) 85 pr_field(((b as i64)+i) as *u8, ll, 2, f2) 86 if f0[0] == (80 as u8) { if pr_streq(f2, daemon) == 1 { return 1 } } // 'P' of PENDING + daemon match 87 } 88 i = e + 1 89 } 90 return 0 91} 92 93// THE PUBLISHER RECOVERY DECISION: scan the monitor's metrics, emit idempotent serialized recovery requests. 94func pr_recover(metricspath: *u8, queuepath: *u8, wedge_ms: i64) -> i64 { 95 let lp: *i64 = sys_mmap(16) as *i64; lp[0] = 0 96 let b: *u8 = sys_read_file(metricspath, lp) 97 if (b as i64) == 0 { return 0 } 98 let n: i64 = lp[0] 99 let url: *u8 = sys_mmap(512); let maxf: *u8 = sys_mmap(32); let daemon: *u8 = sys_mmap(64) 100 var emitted: i64 = 0 101 var i: i64 = 0 102 while i < n { 103 var e: i64 = i; while e < n { if b[e] == (10 as u8) { break } else { e = e + 1 } } 104 let ll: i64 = e - i 105 if ll > 0 { 106 pr_field(((b as i64)+i) as *u8, ll, 2, url) // field2 = url 107 pr_field(((b as i64)+i) as *u8, ll, 9, maxf) // field9 = max_ms 108 let mx: i64 = pr_atoi(maxf) 109 if mx >= wedge_ms { 110 let urllen: i64 = pr_slen(url) 111 if pr_daemon_for(url, urllen, daemon) == 1 { 112 let lk: i64 = fl_acquire(PR_LOCK, PR_MAGIC_100000, 1) 113 if pr_already_pending(queuepath, daemon) == 0 { 114 let rec: *u8 = sys_mmap(PR_REC_CAP); var o: i64 = 0 115 o = pr_s(rec, o, "PENDING" as *u8); rec[o]=9 as u8; o=o+1 116 o = pr_n(rec, o, sys_now_realtime_sec()); rec[o]=9 as u8; o=o+1 117 o = pr_s(rec, o, daemon); rec[o]=9 as u8; o=o+1 118 o = pr_s(rec, o, "wedge" as *u8); rec[o]=9 as u8; o=o+1 119 o = pr_n(rec, o, mx); rec[o]=10 as u8; o=o+1 120 fa_appendz(queuepath, rec, PR_REC_CAP) 121 emitted = emitted + 1 122 } 123 fl_release(lk) 124 } 125 } 126 } 127 i = e + 1 128 } 129 return emitted 130} 131 132func main(argc: i64, argv: *i64) -> i64 { 133 var wedge: i64 = PR_WEDGE_MS_DEFAULT 134 if argc >= 2 { let w: i64 = pr_atoi(argv[1] as *u8); if w > 0 { wedge = w } } 135 pr_uw(1, "=== nishi-publisher recovery decision (wedge>=" as *u8); pr_wn(1, wedge); pr_uw(1, "ms) ===\n" as *u8) 136 let c: i64 = pr_recover(PR_METRICS, PR_QUEUE, wedge) 137 pr_uw(1, "recovery requests emitted: " as *u8); pr_wn(1, c) 138 pr_uw(1, " (queue: " as *u8); pr_uw(1, PR_QUEUE); pr_uw(1, ")\n" as *u8) 139 return 0 140}