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}