nx_sched_qos.nx source
↩ module page · 166 lines · 8943 B
1// nx_sched_qos.nx -- QoS priority + PREEMPTION for sovereign GPU resource-sharing (serving census Q2, closes the
2// PREE gap). The scenario the operator cares about: ONE 16GB GPU shared by BACKGROUND image-gen (Z-Image DiT, long
3// batch job, low priority) and INTERACTIVE chat (short, latency-sensitive, high priority). Plain FCFS makes chat
4// wait behind a 30s image render = bad. This adds: priority-ordered admission + PREEMPTION -- an arriving
5// high-priority request can preempt the lowest-priority RUNNING slot; the preempted job keeps its progress
6// (tokens/steps remaining) and is requeued to RESUME later. That is how you share one accelerator between a
7// foreground companion and a background renderer without either starving. Self-contained (doesn't touch the
8// proven nx_batch_scheduler); pure logic, no GPU needed. license_tier: ORIGINAL expect_exit: 0
9import "nx_syscalls.nx"
10
11const QSLOT_FREE: i64 = 0
12const QSLOT_RUNNING: i64 = 1
13const PRIO_CHAT: i64 = 10 // interactive foreground
14const PRIO_IMAGE: i64 = 1 // background batch render (preemptible)
15
16struct QosSched {
17 n_slots: i64,
18 st: *i64, // slot state[N]
19 rq: *i64, // slot request_id[N]
20 pr: *i64, // slot priority[N]
21 tk: *i64, // slot work_remaining[N]
22 md: *i64, // slot modality[N] (0=chat 1=image)
23 qn: i64, // queue length
24 qrq: *i64, qpr: *i64, qtk: *i64, qmd: *i64 // pending queue (req/prio/tok/mod)[M]
25}
26
27const QCAP: i64 = 64
28
29func qos_new(n_slots: i64) -> *QosSched {
30 let c: *QosSched = sys_mmap(96) as *QosSched
31 c.n_slots = n_slots
32 c.st = sys_mmap(n_slots*8) as *i64
33 c.rq = sys_mmap(n_slots*8) as *i64
34 c.pr = sys_mmap(n_slots*8) as *i64
35 c.tk = sys_mmap(n_slots*8) as *i64
36 c.md = sys_mmap(n_slots*8) as *i64
37 let st: *i64 = c.st
38 var i: i64 = 0
39 while i < n_slots { st[i] = QSLOT_FREE; i = i + 1 }
40 c.qn = 0
41 c.qrq = sys_mmap(QCAP*8) as *i64
42 c.qpr = sys_mmap(QCAP*8) as *i64
43 c.qtk = sys_mmap(QCAP*8) as *i64
44 c.qmd = sys_mmap(QCAP*8) as *i64
45 return c
46}
47
48func qos_running(c: *QosSched) -> i64 {
49 let st: *i64 = c.st; var n: i64 = 0; var i: i64 = 0
50 while i < c.n_slots { if st[i] == QSLOT_RUNNING { n = n + 1 } i = i + 1 }
51 return n
52}
53func _free_slot(c: *QosSched) -> i64 { let st: *i64 = c.st; var i: i64 = 0; while i < c.n_slots { if st[i] == QSLOT_FREE { return i } i = i + 1 } return 0 - 1 }
54// lowest-priority RUNNING slot strictly below `prio`; -1 if none preemptible.
55func _lowest_running(c: *QosSched, prio: i64) -> i64 {
56 let st: *i64 = c.st; let pr: *i64 = c.pr
57 var best: i64 = 0 - 1; var bestp: i64 = prio
58 var i: i64 = 0
59 while i < c.n_slots {
60 if st[i] == QSLOT_RUNNING { if pr[i] < bestp { bestp = pr[i]; best = i } }
61 i = i + 1
62 }
63 return best
64}
65func _enqueue(c: *QosSched, req: i64, prio: i64, tok: i64, mod: i64) -> i64 {
66 if c.qn >= QCAP { return 0 - 1 }
67 let qrq: *i64 = c.qrq; let qpr: *i64 = c.qpr; let qtk: *i64 = c.qtk; let qmd: *i64 = c.qmd
68 let n: i64 = c.qn
69 qrq[n] = req; qpr[n] = prio; qtk[n] = tok; qmd[n] = mod
70 c.qn = n + 1
71 return 0
72}
73func _place(c: *QosSched, slot: i64, req: i64, prio: i64, tok: i64, mod: i64) -> i64 {
74 let st: *i64 = c.st; let rq: *i64 = c.rq; let pr: *i64 = c.pr; let tk: *i64 = c.tk; let md: *i64 = c.md
75 st[slot] = QSLOT_RUNNING; rq[slot] = req; pr[slot] = prio; tk[slot] = tok; md[slot] = mod
76 return slot
77}
78
79// admit: free slot -> place; else preempt a lower-prio running job (requeue it to resume) -> place; else enqueue.
80// returns: slot_idx placed, or -2 = QUEUED. sets preempted_out[0] to the preempted request_id (or -1).
81func qos_admit(c: *QosSched, req: i64, prio: i64, tok: i64, mod: i64, preempted_out: *i64) -> i64 {
82 preempted_out[0] = 0 - 1
83 let f: i64 = _free_slot(c)
84 if f >= 0 { return _place(c, f, req, prio, tok, mod) }
85 let victim: i64 = _lowest_running(c, prio)
86 if victim >= 0 {
87 // requeue the victim WITH its remaining work (resume later, no lost progress)
88 let rq: *i64 = c.rq; let pr: *i64 = c.pr; let tk: *i64 = c.tk; let md: *i64 = c.md
89 preempted_out[0] = rq[victim]
90 _enqueue(c, rq[victim], pr[victim], tk[victim], md[victim])
91 return _place(c, victim, req, prio, tok, mod)
92 }
93 _enqueue(c, req, prio, tok, mod)
94 return 0 - 2
95}
96
97// complete a slot -> FREE it, then drain the HIGHEST-priority queued item into it (priority recycle).
98func qos_complete(c: *QosSched, slot: i64) -> i64 {
99 let st: *i64 = c.st
100 st[slot] = QSLOT_FREE
101 if c.qn <= 0 { return 0 - 1 }
102 // pick highest-priority queued (ties -> earliest index = FCFS within priority)
103 let qpr: *i64 = c.qpr
104 var best: i64 = 0; var bp: i64 = qpr[0]
105 var i: i64 = 1
106 while i < c.qn { if qpr[i] > bp { bp = qpr[i]; best = i } i = i + 1 }
107 let qrq: *i64 = c.qrq; let qtk: *i64 = c.qtk; let qmd: *i64 = c.qmd
108 let rid: i64 = qrq[best]
109 _place(c, slot, rid, qpr[best], qtk[best], qmd[best])
110 // compact the queue (remove `best`)
111 var j: i64 = best
112 while j < c.qn - 1 { qrq[j]=qrq[j+1]; qpr[j]=qpr[j+1]; qtk[j]=qtk[j+1]; qmd[j]=qmd[j+1]; j = j + 1 }
113 c.qn = c.qn - 1
114 return rid
115}
116// find running slot holding request_id (or -1)
117func qos_slot_of(c: *QosSched, req: i64) -> i64 { let st: *i64=c.st; let rq: *i64=c.rq; var i: i64=0; while i<c.n_slots { if st[i]==QSLOT_RUNNING { if rq[i]==req { return i } } i=i+1 } return 0-1 }
118
119func gw(s: *u8) -> i64 { var n: i64=0; while s[n]!=(0 as u8){n=n+1} sys_write(1,s,n); return 0 }
120func gn(v: i64) -> i64 { let bb: *u8=sys_mmap(28); var m: i64=v; if m<0{sys_write(1,"-" as *u8,1);m=0-m} let t: *u8=sys_mmap(28); var k: i64=0; if m==0{t[0]=48 as u8;k=1} while m>0{t[k]=(48+(m%10)) as u8;m=m/10;k=k+1} var i: i64=0; while i<k{bb[i]=t[k-1-i];i=i+1} sys_write(1,bb,k); return 0 }
121
122func main() -> i64 {
123 gw("=== nx_sched_qos -- interactive chat preempts background image-gen on ONE GPU (resource sharing) ===\n" as *u8)
124 var pass: i64 = 0; var tot: i64 = 0
125 let pe: *i64 = sys_mmap(8) as *i64
126
127 // 2-slot GPU (small, to force contention). Fill both with background IMAGE renders (low prio, long jobs).
128 let c: *QosSched = qos_new(2)
129 let s0: i64 = qos_admit(c, 101, PRIO_IMAGE, 30, 1, pe) // image render #101, 30 steps
130 let s1: i64 = qos_admit(c, 102, PRIO_IMAGE, 30, 1, pe) // image render #102
131 gw(" slots full with 2 background image renders (#101,#102, prio 1)\n" as *u8)
132
133 // T1: a THIRD image (low prio) can't preempt -> queued
134 tot=tot+1
135 let s2: i64 = qos_admit(c, 103, PRIO_IMAGE, 30, 1, pe)
136 if s2 == (0-2) { pass=pass+1; gw("PASS T1 equal-prio image #103 QUEUED (no preemption of a peer)\n" as *u8) } else { gw("FAIL T1\n" as *u8) }
137
138 // T2: INTERACTIVE CHAT arrives (high prio) -> PREEMPTS a running image, gets the GPU immediately
139 tot=tot+1
140 let sc: i64 = qos_admit(c, 200, PRIO_CHAT, 8, 0, pe)
141 if sc >= 0 { if pe[0] >= 0 { pass=pass+1; gw("PASS T2 chat #200 PREEMPTED image #" as *u8); gn(pe[0]); gw(" -> got the GPU instantly (low latency)\n" as *u8) } else { gw("FAIL T2b (no preemption)\n" as *u8) } } else { gw("FAIL T2a (chat waited!)\n" as *u8) }
142
143 // T3: the preempted image is REQUEUED with its progress preserved (resume, no lost work)
144 tot=tot+1
145 let preempted_id: i64 = pe[0]
146 var found_in_q: i64 = 0
147 let qrq: *i64 = c.qrq; let qtk: *i64 = c.qtk
148 var qi: i64 = 0
149 while qi < c.qn { if qrq[qi]==preempted_id { if qtk[qi]==30 { found_in_q = 1 } } qi = qi + 1 }
150 if found_in_q==1 { pass=pass+1; gw("PASS T3 preempted image #" as *u8); gn(preempted_id); gw(" requeued WITH its 30 steps intact (resumes, no lost progress)\n" as *u8) } else { gw("FAIL T3 (progress lost or not requeued)\n" as *u8) }
151
152 // T4: chat finishes fast -> completing its slot drains the HIGHEST-priority queued job (priority recycle)
153 tot=tot+1
154 let chat_slot: i64 = qos_slot_of(c, 200)
155 let drained: i64 = qos_complete(c, chat_slot)
156 // queue had image #103 and preempted image; both prio 1 -> FCFS: #103 first
157 if drained == 103 { pass=pass+1; gw("PASS T4 chat done -> slot recycled to queued image #103 (FCFS within priority)\n" as *u8) } else { gw("FAIL T4 drained=" as *u8); gn(drained); gw("\n" as *u8) }
158
159 // T5: no starvation -- everyone still accounted for (2 running + 1 queued)
160 tot=tot+1
161 if qos_running(c)==2 { if c.qn==1 { pass=pass+1; gw("PASS T5 no work lost: 2 running + 1 queued after the whole contention dance\n" as *u8) } else { gw("FAIL T5b qn=" as *u8); gn(c.qn); gw("\n" as *u8) } } else { gw("FAIL T5a running=" as *u8); gn(qos_running(c)); gw("\n" as *u8) }
162
163 gw("nx_sched_qos pass="); gn(pass); gw("/"); gn(tot)
164 if pass==tot { gw(" GREEN -- QoS preemption: chat never waits behind a render; the render resumes intact. SOTA resource-sharing behavior, sovereign.\n" as *u8); sys_exit(0); return 0 }
165 gw(" RED\n" as *u8); sys_exit(1); return 1
166}