code wiki / _hdl_build / nx_vroom_daemon.nx
nx_vroom_daemon.nx source
↩ module page · 352 lines · 16535 B
1// nx_vroom_daemon.nx -- sovereign video-room daemon, PLAIN HTTP on loopback :8446,
2// fronted by nginx (the Synology edge terminates the family TLS; operator chose the clean
3// nishifamily.com/video URL). nginx->TLS-backend failed (OpenSSL rejects nx_tls13's
4// ServerHello with "unsolicited extension"), so the daemon speaks plain HTTP to nginx.
5// The media logic stays 100% sovereign: vr_handle (nx_vroom) + nx_room_relay over
6// shared-mmap, server-relay = no WebRTC/STUN/TURN/WebSocket. Built nx_cc->nxasm (no gcc/.sh).
7// Build: _offc/nx_sov_build_run.elf nx_vroom_daemon
8// license_tier: ORIGINAL
9import "nx_syscalls.nx"
10import "nx_http_server.nx"
11import "nx_room_relay.nx"
12import "nx_vroom.nx"
13import "nx_room_stream.nx"
14const VR_MAGIC_8192: i64 = 8192
15const VR_MAGIC_54000: i64 = 54000
16const VR_MAGIC_131072: i64 = 131072
17
18const VR_PORT: i64 = 8446
19const VR_CLIENT: *u8 = "/volume1/homes/elderwesto/nishihost/vroom_client.html" as *u8
20const VR_PLAINCAP: i64 = 131072
21const VR_OUTCAP: i64 = 262144
22const VR_MAXCHILD: i64 = 64
23const VR_MAXREQ: i64 = 512
24const VR_N: i64 = 16
25const VR_SB: i64 = 65536
26const VR_TTL: i64 = 8000
27const VR_BUDGET: i64 = 1000000000
28const CHAT_CAP: i64 = 128
29const CHAT_TEXTMAX: i64 = 256
30const CHAT_SLOT: i64 = 288
31const VR_PAGE_BODY: *u8 = "<!DOCTYPE html><html><head><meta charset=\"utf-8\"><title>Nishi Family Video</title></head><body style=\"font-family:system-ui;background:#0a0a0c;color:#e6e6ec;text-align:center;padding:8vh\"><h1>Nishi Family Video</h1><p>sovereign daemon live (client file missing)</p></body></html>" as *u8
32
33func vd_strlen(s: *u8) -> i64 { var n: i64=0; while s[n]!=(0 as u8){n=n+1} return n }
34func vd_app(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 }
35func vd_apps(dst: *u8, off: i64, s: *u8) -> i64 { return vd_app(dst, off, s, vd_strlen(s)) }
36func vd_wdec(dst: *u8, off: i64, v: i64) -> i64 {
37 if v == 0 { dst[off]=48; return off+1 }
38 var d: i64=0; var y: i64=v
39 while y>0 { d=d+1; y=y/10 }
40 var i: i64=d-1; y=v
41 while i>=0 { dst[off+i]=(48+(y%10)) as u8; y=y/10; i=i-1 }
42 return off+d
43}
44func vd_build_resp(out: *u8, ctype: *u8, body: *u8, body_len: i64) -> i64 {
45 var o: i64 = vd_apps(out, 0, "HTTP/1.1 200 OK\r\nContent-Type: " as *u8)
46 o = vd_apps(out, o, ctype)
47 o = vd_apps(out, o, "\r\nContent-Length: " as *u8)
48 o = vd_wdec(out, o, body_len)
49 o = vd_apps(out, o, "\r\nConnection: keep-alive\r\n\r\n" as *u8)
50 o = vd_app(out, o, body, body_len)
51 return o
52}
53
54// plain-HTTP recv: read until headers + full Content-Length body (frames span many reads)
55func vd_recv(cfd: i64, buf: *u8, cap: i64) -> i64 {
56 var total: i64 = 0
57 var need: i64 = 0
58 var hdr: i64 = 0
59 var loops: i64 = 0
60 while loops < VR_MAGIC_8192 {
61 let n: i64 = sys_read(cfd, (buf as i64 + total) as *u8, cap - total)
62 if n <= 0 { return total }
63 total = total + n
64 if hdr == 0 {
65 if vr_contains(buf, total, "\r\n\r\n" as *u8, 4) == 1 {
66 hdr = 1
67 let cl: i64 = vr_content_length(buf, total)
68 if cl <= 0 { return total }
69 need = vr_body_start(buf, total) + cl
70 if need > cap { need = cap }
71 }
72 }
73 if hdr == 1 { if total >= need { return total } }
74 if total >= cap { return total }
75 loops = loops + 1
76 }
77 return total
78}
79
80func vd_write_all(cfd: i64, buf: *u8, n: i64) -> i64 {
81 var off: i64 = 0
82 while off < n {
83 let w: i64 = sys_write(cfd, (buf as i64 + off) as *u8, n - off)
84 if w <= 0 { return 0 - 1 }
85 off = off + w
86 }
87 return off
88}
89
90// long-lived chunked-push stream: deliver each OTHER live peer's NEW frame as it arrives.
91// HTTP/1.1 chunked transfer-encoding; the browser auto-de-chunks, so each chunk payload is
92// one record [id u64 LE][len u32 LE][frame bytes]. Heartbeat = record id=0 len=0 (client
93// ignores) -> detects disconnect during idle. Per-connection state (lastseq/lastpeer) in
94// caller buffers. Lower latency than the /frames poll; foundation for audio + FEC.
95func vd_stream(cfd: i64, rr_st: *i64, rr_arena: *u8, rm: i64, pr: i64, ttl: i64,
96 lastseq: *i64, lastpeer: *i64, payload: *u8, chunk: *u8) -> i64 {
97 let hdr: *u8 = "HTTP/1.1 200 OK\r\nContent-Type: application/octet-stream\r\nCache-Control: no-cache\r\nConnection: keep-alive\r\nTransfer-Encoding: chunked\r\n\r\n" as *u8
98 if vd_write_all(cfd, hdr, vd_strlen(hdr)) < 0 { return 0 - 1 }
99 let n: i64 = rr_nslots(rr_st)
100 let slotb: i64 = rr_st[1]
101 var i: i64 = 0
102 while i < n { lastseq[i] = 0; lastpeer[i] = 0; i = i + 1 }
103 var ticks: i64 = 0
104 var idle: i64 = 0
105 while ticks < VR_MAGIC_54000 {
106 let now: i64 = sys_now_us() / 1000
107 var sent: i64 = 0
108 var k: i64 = 0
109 while k < n {
110 if rr_room(rr_st, k) == rm {
111 let pk: i64 = rr_peer(rr_st, k)
112 if pk != pr {
113 if (now - rr_seen(rr_st, k)) <= ttl {
114 if pk != lastpeer[k] { lastpeer[k] = pk; lastseq[k] = 0 }
115 let sq: i64 = rr_seq(rr_st, k)
116 if sq > lastseq[k] {
117 let fl: i64 = rr_flen(rr_st, k)
118 var b: i64 = 0
119 while b < 8 { payload[b] = ((pk >> (b*8)) & 0xff) as u8; b = b + 1 }
120 var c: i64 = 0
121 while c < 4 { payload[8+c] = ((fl >> (c*8)) & 0xff) as u8; c = c + 1 }
122 var j: i64 = 0
123 while j < fl { payload[12+j] = rr_arena[k*slotb + j]; j = j + 1 }
124 let cn: i64 = rs_chunk_encode(chunk, 0, payload, 12 + fl)
125 if vd_write_all(cfd, chunk, cn) < 0 { return 0 - 1 }
126 lastseq[k] = sq
127 sent = 1
128 }
129 }
130 }
131 }
132 k = k + 1
133 }
134 if sent == 0 {
135 idle = idle + 1
136 if idle >= 40 {
137 var z: i64 = 0
138 while z < 12 { payload[z] = 0 as u8; z = z + 1 }
139 let hcn: i64 = rs_chunk_encode(chunk, 0, payload, 12)
140 if vd_write_all(cfd, chunk, hcn) < 0 { return 0 - 1 }
141 idle = 0
142 }
143 } else { idle = 0 }
144 sys_sleep_ms(33)
145 ticks = ticks + 1
146 }
147 let term: *u8 = "0\r\n\r\n" as *u8
148 vd_write_all(cfd, term, 5)
149 return 0
150}
151
152// POST an audio chunk into the audio relay (mirrors vr_handle's /post, separate channel).
153func vd_apost(cfd: i64, plain: *u8, pn: i64, rr_st: *i64, rr_arena: *u8, now: i64) -> i64 {
154 let rm: i64 = vr_room_hash(plain, pn)
155 let pr: i64 = vr_id(plain, pn)
156 if rm != 0 { if pr >= 0 {
157 let bs: i64 = vr_body_start(plain, pn)
158 var fl: i64 = pn - bs
159 let cl: i64 = vr_content_length(plain, pn)
160 if cl >= 0 { if cl < fl { fl = cl } }
161 if fl < 0 { fl = 0 }
162 rr_post(rr_st, rr_arena, rm, pr, now, now, (plain as i64 + bs) as *u8, fl)
163 } }
164 let ok: *u8 = "HTTP/1.1 200 OK\r\nContent-Length: 2\r\nConnection: keep-alive\r\n\r\nok" as *u8
165 return vd_write_all(cfd, ok, vd_strlen(ok))
166}
167
168func vd_rd_i64(b: *u8, o: i64) -> i64 { var v: i64=0; var i: i64=0; while i<8 { v = v | ((b[o+i] as i64) << (i*8)); i=i+1 } return v }
169// parse integer value of query param `key` (e.g. "since="); -1 if absent
170func vd_qparam_int(req: *u8, n: i64, key: *u8, keyn: i64) -> i64 {
171 var lim: i64 = 0; var sc: i64 = 1
172 while sc == 1 { if lim < n { if req[lim]==(13 as u8) { sc = 0 } else { lim = lim + 1 } } else { sc = 0 } }
173 var found: i64 = 0 - 1; var i: i64 = 0
174 while i + keyn <= lim {
175 if found < 0 { var j: i64=0; var ok: i64=1; while j<keyn { if req[i+j]!=key[j] { ok=0; j=keyn } else { j=j+1 } } if ok==1 { found=i+keyn } }
176 i = i + 1
177 }
178 if found < 0 { return 0 - 1 }
179 var v: i64 = 0; var got: i64 = 0; var k: i64 = found; var run: i64 = 1
180 while run == 1 { if k >= lim { run = 0 } else { let c: i64 = req[k] as i64; var dig: i64=0; if c>=48 { if c<=57 { dig=1 } } if dig==1 { v=v*10+(c-48); got=1; k=k+1 } else { run=0 } } }
181 if got == 0 { return 0 - 1 }
182 return v
183}
184// POST /chat: append a text message to the room's chat ring (shared across forks)
185func vd_chat_post(cfd: i64, plain: *u8, pn: i64, meta: *i64, ring: *u8) -> i64 {
186 let rm: i64 = vr_room_hash(plain, pn)
187 let from: i64 = vr_id(plain, pn)
188 if rm != 0 {
189 let bs: i64 = vr_body_start(plain, pn)
190 var len: i64 = pn - bs
191 let cl: i64 = vr_content_length(plain, pn)
192 if cl >= 0 { if cl < len { len = cl } }
193 if len < 0 { len = 0 }
194 if len > CHAT_TEXTMAX { len = CHAT_TEXTMAX }
195 let seq: i64 = meta[0] + 1
196 meta[0] = seq
197 let base: i64 = ((seq - 1) % CHAT_CAP) * CHAT_SLOT
198 var b: i64 = 0
199 while b < 8 { ring[base+b] = ((seq >> (b*8)) & 0xff) as u8; b=b+1 }
200 b = 0; while b < 8 { ring[base+8+b] = ((rm >> (b*8)) & 0xff) as u8; b=b+1 }
201 b = 0; while b < 8 { ring[base+16+b] = ((from >> (b*8)) & 0xff) as u8; b=b+1 }
202 b = 0; while b < 8 { ring[base+24+b] = ((len >> (b*8)) & 0xff) as u8; b=b+1 }
203 var t: i64 = 0; while t < len { ring[base+32+t] = plain[bs+t]; t=t+1 }
204 }
205 let ok: *u8 = "HTTP/1.1 200 OK\r\nContent-Length: 2\r\nConnection: keep-alive\r\n\r\nok" as *u8
206 return vd_write_all(cfd, ok, vd_strlen(ok))
207}
208// GET /chatlog?room&since=N: octet body of records [seq u64][from u64][len u32][text] (seq>since, room match)
209func vd_chat_log(cfd: i64, plain: *u8, pn: i64, meta: *i64, ring: *u8, out: *u8) -> i64 {
210 let rm: i64 = vr_room_hash(plain, pn)
211 var since: i64 = vd_qparam_int(plain, pn, "since=" as *u8, 6)
212 if since < 0 { since = 0 }
213 let gseq: i64 = meta[0]
214 var start: i64 = gseq - CHAT_CAP + 1
215 if start < 1 { start = 1 }
216 if start <= since { start = since + 1 }
217 var blen: i64 = 0
218 var s: i64 = start
219 while s <= gseq {
220 let base: i64 = ((s - 1) % CHAT_CAP) * CHAT_SLOT
221 if vd_rd_i64(ring, base) == s { if vd_rd_i64(ring, base+8) == rm { blen = blen + 20 + vd_rd_i64(ring, base+24) } }
222 s = s + 1
223 }
224 var o: i64 = vd_apps(out, 0, "HTTP/1.1 200 OK\r\nContent-Type: application/octet-stream\r\nContent-Length: " as *u8)
225 o = vd_wdec(out, o, blen)
226 o = vd_apps(out, o, "\r\nConnection: keep-alive\r\n\r\n" as *u8)
227 s = start
228 while s <= gseq {
229 let base: i64 = ((s - 1) % CHAT_CAP) * CHAT_SLOT
230 if vd_rd_i64(ring, base) == s {
231 if vd_rd_i64(ring, base+8) == rm {
232 let from: i64 = vd_rd_i64(ring, base+16)
233 let len: i64 = vd_rd_i64(ring, base+24)
234 o = vr_wu64(out, o, s)
235 o = vr_wu64(out, o, from)
236 o = vr_wu32(out, o, len)
237 var t: i64 = 0; while t < len { out[o+t] = ring[base+32+t]; t=t+1 }
238 o = o + len
239 }
240 }
241 s = s + 1
242 }
243 return vd_write_all(cfd, out, o)
244}
245
246func main() -> i64 {
247 let addr_buf: *u8 = sys_mmap(16)
248 nx_http_server_addr_any(addr_buf, VR_PORT)
249 let lv: *i64 = (sys_mmap(8)) as *i64
250 let lfd: i64 = nx_http_server_listen(addr_buf, 16, lv)
251 if lfd < 0 { return 4 }
252 let banner: *u8 = "nishi sovereign video daemon (HTTP, server-relay, behind nginx) on 0.0.0.0:8446\n" as *u8
253 sys_write(1, banner, vd_strlen(banner))
254
255 let rr_st: *i64 = (sys_mmap_shared((2 + 5*VR_N) * 8)) as *i64
256 let rr_arena: *u8 = sys_mmap_shared(VR_N * VR_SB)
257 rr_init(rr_st, VR_N, VR_SB)
258 // second relay for AUDIO chunks (separate channel: /apost stores, /astream pushes)
259 let rr_a_st: *i64 = (sys_mmap_shared((2 + 5*VR_N) * 8)) as *i64
260 let rr_a_arena: *u8 = sys_mmap_shared(VR_N * VR_SB)
261 rr_init(rr_a_st, VR_N, VR_SB)
262 // shared chat ring (text messages, history)
263 let chat_meta: *i64 = (sys_mmap_shared(64)) as *i64
264 chat_meta[0] = 0
265 let chat_ring: *u8 = sys_mmap_shared(CHAT_CAP * CHAT_SLOT)
266
267 let page_resp: *u8 = sys_mmap(VR_MAGIC_131072)
268 var page_n: i64 = 0
269 let chb: *i64 = (sys_mmap(8)) as *i64
270 chb[0] = 0
271 let client_html: *u8 = sys_read_file(VR_CLIENT, chb)
272 if (client_html as i64) != 0 { if chb[0] > 0 { page_n = vd_build_resp(page_resp, "text/html; charset=utf-8" as *u8, client_html, chb[0]) } }
273 if page_n == 0 { page_n = vd_build_resp(page_resp, "text/html; charset=utf-8" as *u8, VR_PAGE_BODY, vd_strlen(VR_PAGE_BODY)) }
274
275 let plain: *u8 = sys_mmap(VR_PLAINCAP)
276 let out: *u8 = sys_mmap(VR_OUTCAP)
277 let ids: *i64 = (sys_mmap(VR_N * 8)) as *i64
278 let st_lastseq: *i64 = (sys_mmap(VR_N * 8)) as *i64
279 let st_lastpeer: *i64 = (sys_mmap(VR_N * 8)) as *i64
280 let st_payload: *u8 = sys_mmap(VR_SB + 64)
281 let st_chunk: *u8 = sys_mmap(VR_SB + 128)
282 let crb: *i64 = (sys_mmap(8)) as *i64
283 let sock_addr: *u8 = sys_mmap(64)
284 let sock_len: *i64 = (sys_mmap(8)) as *i64
285 let reap: *i64 = (sys_mmap(8)) as *i64
286
287 sys_set_socket_timeout(lfd, 5)
288 var served: i64 = 0
289 var live: i64 = 0
290 while served < VR_BUDGET {
291 while sys_wait4(0 - 1, reap, 1) > 0 { live = live - 1 }
292 sock_len[0] = 16
293 let cfd: i64 = sys_accept_with_addr(lfd, sock_addr, sock_len)
294 if cfd < 0 { continue }
295 if live >= VR_MAXCHILD { if sys_wait4(0 - 1, reap, 0) > 0 { live = live - 1 } }
296 let pid: i64 = sys_fork()
297 if pid == 0 {
298 sys_close(lfd)
299 sys_set_socket_timeout(cfd, 30)
300 var keep: i64 = 1
301 var nreq: i64 = 0
302 while keep == 1 {
303 if nreq >= VR_MAXREQ { keep = 0 } else {
304 let pn: i64 = vd_recv(cfd, plain, VR_PLAINCAP)
305 if pn <= 0 { keep = 0 } else {
306 var done: i64 = 0
307 if vr_contains(plain, pn, "/astream" as *u8, 8) == 1 {
308 vd_stream(cfd, rr_a_st, rr_a_arena, vr_room_hash(plain, pn), vr_id(plain, pn), VR_TTL, st_lastseq, st_lastpeer, st_payload, st_chunk)
309 keep = 0; done = 1
310 }
311 if done == 0 { if vr_contains(plain, pn, "/stream" as *u8, 7) == 1 {
312 vd_stream(cfd, rr_st, rr_arena, vr_room_hash(plain, pn), vr_id(plain, pn), VR_TTL, st_lastseq, st_lastpeer, st_payload, st_chunk)
313 keep = 0; done = 1
314 } }
315 if done == 0 { if vr_contains(plain, pn, "/apost" as *u8, 6) == 1 {
316 if vd_apost(cfd, plain, pn, rr_a_st, rr_a_arena, sys_now_us() / 1000) < 0 { keep = 0 }
317 nreq = nreq + 1; done = 1
318 } }
319 if done == 0 { if vr_contains(plain, pn, "/chatlog" as *u8, 8) == 1 {
320 if vd_chat_log(cfd, plain, pn, chat_meta, chat_ring, out) < 0 { keep = 0 }
321 nreq = nreq + 1; done = 1
322 } }
323 if done == 0 { if vr_contains(plain, pn, "/chat" as *u8, 5) == 1 {
324 if vd_chat_post(cfd, plain, pn, chat_meta, chat_ring) < 0 { keep = 0 }
325 nreq = nreq + 1; done = 1
326 } }
327 if done == 0 {
328 let tnow: i64 = sys_now_us() / 1000
329 var rlen: i64 = vr_handle(plain, pn, rr_st, rr_arena, tnow, VR_TTL, out, VR_OUTCAP, ids)
330 var rp: *u8 = out
331 if rlen == 0 {
332 crb[0] = 0
333 let chf: *u8 = sys_read_file(VR_CLIENT, crb)
334 if (chf as i64) != 0 { if crb[0] > 0 { rlen = vd_build_resp(out, "text/html; charset=utf-8" as *u8, chf, crb[0]); rp = out } }
335 if rlen == 0 { rp = page_resp; rlen = page_n }
336 }
337 if vd_write_all(cfd, rp, rlen) < 0 { keep = 0 }
338 nreq = nreq + 1
339 }
340 }
341 }
342 }
343 sys_close(cfd)
344 sys_exit(0)
345 return 0
346 }
347 live = live + 1
348 sys_close(cfd)
349 served = served + 1
350 }
351 return 0
352}