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}