code wiki / _hdl_build / nx_research_serve.nx
nx_research_serve.nx source
↩ module page · 216 lines · 11698 B
1// nx_research_serve.nx -- THE RESEARCH REQUEST QUEUE + ADEQUACY FEEDBACK GATE. Lets the autonomous builder and
2// the researcher partner gracefully + loosely (async, via artifacts, integer-only / no floats):
3// - Auto-builder ENQUEUES by appending to knowledge/index/research_queue.txt: <id> PENDING <topic terms...>
4// - This organ DRAINS the queue: each PENDING request is queried across EVERY index (via the search-sources
5// manifest = the same loose federation) and JUDGED for adequacy (distinct docs covering the topic >=
6// SV_MIN_HITS), then routed:
7// ADEQUATE -> status SERVED + research_served.log (builder may now retrieve via federation)
8// INADEQUATE -> status INADEQUATE + research_gaps.txt (THE FEEDBACK FLAG: under-covered topic -> shown
9// to operator + Claude -> improve researcher / fetch deeper)
10// - Queue rewritten with updated statuses (history kept; idempotent: only PENDING processed).
11// Closes the partner loop: builder asks -> researcher serves or flags the gap -> we grow the researcher -> the
12// auto-builder is unblocked -> the team grows toward autonomous S-class. expect_exit: 0 license_tier: ORIGINAL
13import "nx_search_inverted_persist.nx"
14import "nx_itoa_lib.nx" // shared MSB-first emitter (zero-alloc)
15const SV_MAGIC_1024: i64 = 1024
16const SV_MAGIC_2048: i64 = 2048
17
18const SV_MIN_HITS: i64 = 3 // distinct docs to call research ADEQUATE (data-driven; config later)
19const SV_TERMS: i64 = 8
20const SV_ROWIDS: i64 = 4096
21const SV_IDX: i64 = 8192
22const SV_HITS: i64 = 65536
23
24func 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 }
25// MIGRATED to the shared emitter (debt 1785563586). The old body mmapped a scratch buffer
26// per call and never freed it. At PAGE granularity that is 4096B leaked PER CALL -- the
27// defect that took 28.5GB of a 36GB host in nx_ts_lumadiff (2MB input, ~3.66M calls).
28// nxi_* is MSB-first, allocates NOTHING, and emits identical bytes including the sign.
29func sv_num(v: i64) -> i64 { nxi_out(v); return 0 }
30func sv_strlen(s: *u8) -> i64 { var n: i64=0; while s[n]!=(0 as u8){n=n+1} return n }
31func sv_cat(dst: *u8, off: i64, s: *u8) -> i64 { var i: i64=0; while s[i]!=(0 as u8){dst[off+i]=s[i];i=i+1} return off+i }
32func sv_catbuf(dst: *u8, off: i64, src: *u8, n: i64) -> i64 { var i: i64=0; while i<n { dst[off+i]=src[i]; i=i+1 } return off+n }
33func sv_catnum(dst: *u8, off: i64, v: i64) -> i64 {
34 if v == 0 { dst[off]=48 as u8; return off+1 }
35 let tmp: *u8=sys_mmap(28); var m: i64=v; var k: i64=0
36 while m>0 { tmp[k]=(48+(m%10)) as u8; m=m/10; k=k+1 }
37 var i: i64=0; while i<k { dst[off+i]=tmp[k-1-i]; i=i+1 }
38 return off+k
39}
40func sv_scan(buf: *u8, i: i64, end: i64) -> i64 {
41 var p: i64=i; var stop: i64=0
42 while stop==0 { if p>=end { stop=1 } else { if buf[p]==(32 as u8) { stop=1 } else { if buf[p]==(10 as u8) { stop=1 } else { p=p+1 } } } }
43 return p
44}
45func sv_line_end(buf: *u8, i: i64, end: i64) -> i64 {
46 var p: i64=i; var stop: i64=0
47 while stop==0 { if p>=end { stop=1 } else { if buf[p]==(10 as u8) { stop=1 } else { p=p+1 } } }
48 return p
49}
50func sv_load_docmap(path: *u8, paths_out: *i64, cap: i64) -> i64 {
51 let lenbox: *i64 = sys_mmap(16) as *i64
52 let buf: *u8 = sys_read_file(path, lenbox)
53 if buf == 0 as *u8 { return 0 - 1 }
54 let total: i64 = lenbox[0]
55 var nd: i64 = 0; var start: i64 = 0; var i: i64 = 0
56 while i < total {
57 if buf[i] == (10 as u8) { buf[i]=0 as u8; if i>start { if nd<cap { paths_out[nd]=(buf as i64)+start; nd=nd+1 } } start=i+1 }
58 i = i + 1
59 }
60 return nd
61}
62
63// load every index in the manifest (once) into the caller's arrays. Returns count.
64func sv_load_indexes(gidx: *i64, gpaths: *i64, gndocs: *i64) -> i64 {
65 let mbox: *i64 = sys_mmap(16) as *i64
66 let man: *u8 = sys_read_file("knowledge/index/search_sources.txt" as *u8, mbox)
67 if man == 0 as *u8 { return 0 }
68 let mlen: i64 = mbox[0]
69 var n: i64 = 0; var i: i64 = 0
70 while i < mlen {
71 let e1: i64 = sv_scan(man, i, mlen)
72 var j: i64 = e1; if j<mlen { if man[j]==(32 as u8) { man[j]=0 as u8; j=j+1 } }
73 let idx_ptr: i64 = (man as i64) + j
74 let e2: i64 = sv_scan(man, j, mlen)
75 var k: i64 = e2; if k<mlen { if man[k]==(32 as u8) { man[k]=0 as u8; k=k+1 } }
76 let dmap_ptr: i64 = (man as i64) + k
77 let e3: i64 = sv_scan(man, k, mlen); if e3<mlen { if man[e3]==(32 as u8) { man[e3]=0 as u8 } }
78 var e: i64 = sv_line_end(man, e3, mlen); if e<mlen { man[e]=0 as u8; e=e+1 }
79 let idx: *NxInvIndex = nx_inv_load(idx_ptr as *u8)
80 if idx != 0 as *NxInvIndex {
81 let paths: *i64 = sys_mmap(8*SV_ROWIDS) as *i64
82 let nd: i64 = sv_load_docmap(dmap_ptr as *u8, paths, SV_ROWIDS)
83 if nd > 0 { if n < SV_IDX { gidx[n]=idx as i64; gpaths[n]=paths as i64; gndocs[n]=nd; n=n+1 } }
84 }
85 i = e
86 }
87 return n
88}
89
90// adequacy = distinct docs across all indexes matching ALL topic terms (conjunction). ANY-term matching is too
91// loose -- a nonsense topic with one common word ("widget") would falsely pass; requiring every term in the same
92// doc is what tells a real, covered topic from a stray match.
93func sv_hits_for(terms: *i64, nterms: i64, nidx: i64, gidx: *i64, gndocs: *i64, ghits: *u8, gcount: *u8) -> i64 {
94 if nterms <= 0 { return 0 }
95 let res: *NxInvQueryResult = sys_mmap(64) as *NxInvQueryResult
96 let rowids: *i64 = sys_mmap(8*SV_ROWIDS) as *i64
97 var total: i64 = 0
98 var s: i64 = 0
99 while s < nidx {
100 let idx: *NxInvIndex = gidx[s] as *NxInvIndex
101 let nd: i64 = gndocs[s]
102 var z: i64 = 0
103 while z < nd { gcount[z]=0 as u8; z=z+1 } // # distinct terms matching each doc
104 var t: i64 = 0
105 while t < nterms {
106 var z2: i64 = 0
107 while z2 < nd { ghits[z2]=0 as u8; z2=z2+1 } // per-term doc marker (dedup within term)
108 let term: *u8 = terms[t] as *u8
109 nx_inv_query_term(idx, term, sv_strlen(term), rowids, SV_ROWIDS, res)
110 var ri: i64 = 0
111 while ri < res.n_rowids_filled {
112 let rid: i64 = rowids[ri]
113 if rid >= 0 { if rid < nd { ghits[rid]=1 as u8 } }
114 ri = ri + 1
115 }
116 var d2: i64 = 0
117 while d2 < nd { if ghits[d2]==(1 as u8) { gcount[d2]=(gcount[d2] as i64 + 1) as u8 } d2=d2+1 }
118 t = t + 1
119 }
120 var d: i64 = 0
121 while d < nd { if gcount[d] as i64 == nterms { total=total+1 } d=d+1 }
122 s = s + 1
123 }
124 return total
125}
126
127func main() -> i64 {
128 sv_puts("=== nx_research_serve: drain research queue, judge adequacy, flag gaps (no floats) ===\n" as *u8)
129 let gidx: *i64 = sys_mmap(8*SV_IDX) as *i64
130 let gpaths: *i64 = sys_mmap(8*SV_IDX) as *i64
131 let gndocs: *i64 = sys_mmap(8*SV_IDX) as *i64
132 let ghits: *u8 = sys_mmap(SV_HITS)
133 let gcount: *u8 = sys_mmap(SV_HITS)
134 let nidx: i64 = sv_load_indexes(gidx, gpaths, gndocs)
135 if nidx <= 0 { sv_puts(" NO indexes (run nx_search_register first)\n" as *u8); sys_exit(1); return 1 }
136 sv_puts(" loaded indexes="); sv_num(nidx); sv_puts("\n" as *u8)
137
138 let qbox: *i64 = sys_mmap(16) as *i64
139 let q: *u8 = sys_read_file("knowledge/index/research_queue.txt" as *u8, qbox)
140 if q == 0 as *u8 { sv_puts(" queue empty (enqueue: append '<id> PENDING <topic>' to research_queue.txt)\n" as *u8); sys_exit(0); return 0 }
141 let qlen: i64 = qbox[0]
142
143 let nq_fd: i64 = sys_openat_wr("knowledge/index/research_queue.next" as *u8, 0x1a4)
144 let served_fd: i64 = sys_openat_append("knowledge/index/research_served.log" as *u8, 0x1a4)
145 let gaps_fd: i64 = sys_openat_append("knowledge/index/research_gaps.txt" as *u8, 0x1a4)
146
147 var served: i64=0; var inadequate: i64=0; var kept: i64=0
148 var i: i64 = 0
149 while i < qlen {
150 let lstart: i64 = i
151 let le: i64 = sv_line_end(q, i, qlen)
152 if le > lstart {
153 let e1: i64 = sv_scan(q, lstart, le)
154 let id_ptr: i64 = (q as i64) + lstart; let id_len: i64 = e1 - lstart
155 var j: i64 = e1; if j<le { if q[j]==(32 as u8) { j=j+1 } }
156 let e2: i64 = sv_scan(q, j, le)
157 let st_len: i64 = e2 - j
158 var p: i64 = e2; if p<le { if q[p]==(32 as u8) { p=p+1 } }
159 let topic_ptr: i64 = (q as i64) + p; let topic_len: i64 = le - p
160
161 var pending: i64 = 0
162 if st_len == 7 { if q[j]==(80 as u8) { pending=1 } }
163 var ns_ptr: i64 = (q as i64) + j
164 var ns_len: i64 = st_len
165 if pending == 1 {
166 let tbuf: *u8 = sys_mmap(topic_len + 8)
167 var ci: i64 = 0
168 while ci < topic_len { tbuf[ci]=q[p+ci]; ci=ci+1 }
169 tbuf[topic_len]=0 as u8
170 let terms: *i64 = sys_mmap(8*SV_TERMS) as *i64
171 var nterms: i64 = 0
172 var ti: i64 = 0
173 var instr: i64 = 0
174 while ti < topic_len {
175 if tbuf[ti]==(32 as u8) { tbuf[ti]=0 as u8; instr=0 }
176 else { if instr==0 { if nterms<SV_TERMS { terms[nterms]=((tbuf as i64)+ti); nterms=nterms+1 } instr=1 } }
177 ti=ti+1
178 }
179 let hits: i64 = sv_hits_for(terms, nterms, nidx, gidx, gndocs, ghits, gcount)
180 if hits >= SV_MIN_HITS {
181 ns_ptr = "SERVED" as *u8 as i64; ns_len = 6; served=served+1
182 let sl: *u8=sys_mmap(SV_MAGIC_1024); var sp: i64=0
183 sp=sv_catbuf(sl,sp,id_ptr as *u8,id_len); sl[sp]=32 as u8; sp=sp+1
184 sp=sv_cat(sl,sp,"SERVED hits=" as *u8); sp=sv_catnum(sl,sp,hits); sl[sp]=32 as u8; sp=sp+1
185 sp=sv_catbuf(sl,sp,topic_ptr as *u8,topic_len); sl[sp]=10 as u8; sp=sp+1
186 if served_fd>=0 { sys_write(served_fd, sl, sp) }
187 } else {
188 ns_ptr = "INADEQUATE" as *u8 as i64; ns_len = 10; inadequate=inadequate+1
189 let gl: *u8=sys_mmap(SV_MAGIC_1024); var gp: i64=0
190 gp=sv_cat(gl,gp,"GAP topic=" as *u8); gp=sv_catbuf(gl,gp,topic_ptr as *u8,topic_len)
191 gp=sv_cat(gl,gp," hits=" as *u8); gp=sv_catnum(gl,gp,hits)
192 gp=sv_cat(gl,gp," need>=" as *u8); gp=sv_catnum(gl,gp,SV_MIN_HITS)
193 gp=sv_cat(gl,gp," FLAG: improve researcher / fetch deeper for this topic\n" as *u8)
194 if gaps_fd>=0 { sys_write(gaps_fd, gl, gp) }
195 }
196 } else { kept = kept + 1 }
197
198 let nl: *u8=sys_mmap(SV_MAGIC_2048); var np: i64=0
199 np=sv_catbuf(nl,np,id_ptr as *u8,id_len); nl[np]=32 as u8; np=np+1
200 np=sv_catbuf(nl,np,ns_ptr as *u8,ns_len); nl[np]=32 as u8; np=np+1
201 np=sv_catbuf(nl,np,topic_ptr as *u8,topic_len); nl[np]=10 as u8; np=np+1
202 if nq_fd>=0 { sys_write(nq_fd, nl, np) }
203 }
204 i = le; if i < qlen { i = i + 1 }
205 }
206 if nq_fd>=0 { sys_close(nq_fd) }
207 if served_fd>=0 { sys_close(served_fd) }
208 if gaps_fd>=0 { sys_close(gaps_fd) }
209 sys_renameat("knowledge/index/research_queue.next" as *u8, "knowledge/index/research_queue.txt" as *u8)
210
211 sv_puts(" drained: served="); sv_num(served); sv_puts(" inadequate_flagged="); sv_num(inadequate)
212 sv_puts(" already_done="); sv_num(kept); sv_puts("\n" as *u8)
213 sv_puts(" SERVED->research_served.log ; GAPS->research_gaps.txt (feedback flag to operator+Claude)\n" as *u8)
214 sv_puts(" SERVE-OK\n" as *u8)
215 sys_exit(0); return 0
216}