nx_mesh_batch.nx source
↩ module page · 196 lines · 11192 B
1// nx_mesh_batch.nx -- WORKER MESH: DYNAMIC BATCHING scheduler + request QUEUE with admission control (the Triton
2// dynamic-batcher analog). A GPU serves far more throughput if concurrent requests are COALESCED into one forward
3// pass instead of run one-at-a-time. This organ is the SCHEDULER that sits in front of a worker:
4// * QUEUE (bounded FIFO) + ADMISSION CONTROL: admit while qlen < cap, else REJECT (backpressure / 503) -- protects
5// the GPU from unbounded pileup under load.
6// * DYNAMIC BATCH COALESCER: form + dispatch a batch when EITHER the queue reaches the preferred batch size
7// (qlen >= max_batch) OR the oldest queued request has waited >= max_delay_ms (bound tail latency) -- exactly
8// Triton's {preferred_batch_size, max_queue_delay}. A dispatched batch of K maps to an n=K worker call.
9// The scheduler is pure integer policy -> mechanically provable. A discrete-event SIM feeds an arrival stream and
10// reports batches formed, avg batch size (the throughput win), rejects, and max queue depth.
11//
12// CLI: (no args) -> self-test GATE (policy table + sim + neg-controls)
13// sim <burst> <max_batch> <max_delay> <cap> -> run the sim for a burst of N arrivals at t=0, print result
14// Closes TWO mesh axes: dynamic-batching AND request-queue-admission. NO fake greens: the sim + neg-controls prove a
15// batch NEVER exceeds max_batch, an empty batch is NEVER dispatched, and admitted NEVER exceeds cap. license_tier: ORIGINAL
16import "nx_syscalls.nx"
17import "nx_runtime.nx"
18
19func mb_w(fd: i64, s: *u8) -> i64 { var n: i64 = 0; while s[n] != (0 as u8) { n = n + 1 } sys_write(fd, s, n); return 0 }
20func mb_wn(fd: i64, v: i64) -> i64 {
21 let bb: *u8 = sys_mmap(28)
22 var m: i64 = v
23 if m < 0 { sys_write(fd, "-" as *u8, 1); m = 0 - m }
24 let tt: *u8 = sys_mmap(28)
25 var k: i64 = 0
26 if m == 0 { tt[0] = 48 as u8; k = 1 }
27 while m > 0 { tt[k] = (48 + (m - (m/10)*10)) as u8; m = m / 10; k = k + 1 }
28 var i: i64 = 0
29 while i < k { bb[i] = tt[k-1-i]; i = i + 1 }
30 sys_write(fd, bb, k)
31 return 0
32}
33func mb_p(s: *u8) -> i64 { return mb_w(1, s) }
34func mb_pn(v: i64) -> i64 { return mb_wn(1, v) }
35func mb_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 }
36func mb_min(a: i64, b: i64) -> i64 { if a < b { return a } return b }
37func mb_streq(a: *u8, b: *u8) -> i64 { var i: i64 = 0; while b[i] != (0 as u8) { if a[i] != b[i] { return 0 } i = i + 1 } if a[i] != (0 as u8) { return 0 } return 1 }
38
39// ---- PURE POLICY ----
40// admit a new request iff the queue has room (else backpressure/reject).
41func mb_admit(qlen: i64, cap: i64) -> i64 { if qlen < cap { return 1 } return 0 }
42// dispatch a batch now? yes if a full preferred batch is queued, OR the oldest has waited past the delay bound.
43func mb_batch_ready(qlen: i64, max_batch: i64, waited_ms: i64, max_delay_ms: i64) -> i64 {
44 if qlen <= 0 { return 0 } // never dispatch an empty batch
45 if qlen >= max_batch { return 1 } // preferred batch full -> go now (throughput)
46 if waited_ms >= max_delay_ms { return 1 } // window elapsed with work waiting -> go (bound latency)
47 return 0
48}
49// how many to take in a batch (never more than the preferred size).
50func mb_batch_take(qlen: i64, max_batch: i64) -> i64 { return mb_min(qlen, max_batch) }
51
52// ---- discrete-event SIM: arrivals[t] requests arrive at ms t; run the scheduler over [0,horizon].
53// outbox[0]=admitted outbox[1]=rejected outbox[2]=batches outbox[3]=total_batched outbox[4]=maxdepth outbox[5]=over_batch(liar) ----
54func mb_sim(arrivals: *i64, horizon: i64, max_batch: i64, max_delay: i64, cap: i64, outbox: *i64) -> i64 {
55 var qlen: i64 = 0
56 var admitted: i64 = 0
57 var rejected: i64 = 0
58 var batches: i64 = 0
59 var total_batched: i64 = 0
60 var maxdepth: i64 = 0
61 var over: i64 = 0
62 var batch_open: i64 = 0
63 var open_valid: i64 = 0
64 var t: i64 = 0
65 while t <= horizon {
66 // arrivals at t
67 let na: i64 = arrivals[t]
68 var j: i64 = 0
69 while j < na {
70 if mb_admit(qlen, cap) == 1 { qlen = qlen + 1; admitted = admitted + 1; if open_valid == 0 { batch_open = t; open_valid = 1 } } else { rejected = rejected + 1 }
71 j = j + 1
72 }
73 if qlen > maxdepth { maxdepth = qlen }
74 // drain any ready batches this tick (a big burst forms several full batches at once)
75 var draining: i64 = 1
76 while draining == 1 {
77 let waited: i64 = t - batch_open
78 if mb_batch_ready(qlen, max_batch, waited, max_delay) == 1 {
79 let take: i64 = mb_batch_take(qlen, max_batch)
80 if take > max_batch { over = over + 1 } // liar tripwire: must never fire
81 if take <= 0 { draining = 0 } else {
82 qlen = qlen - take
83 batches = batches + 1
84 total_batched = total_batched + take
85 if qlen > 0 { batch_open = t; open_valid = 1 } else { open_valid = 0 }
86 }
87 } else { draining = 0 }
88 }
89 t = t + 1
90 }
91 outbox[0] = admitted; outbox[1] = rejected; outbox[2] = batches
92 outbox[3] = total_batched; outbox[4] = maxdepth; outbox[5] = over
93 return 0
94}
95
96func mb_gate() -> i64 {
97 mb_p("=== nx_mesh_batch: dynamic-batching scheduler + queue admission (Triton dynamic-batcher analog) ===\n" as *u8)
98 // policy table
99 var t1: i64 = 0
100 if mb_admit(3, 4) == 1 { if mb_admit(4, 4) == 0 { t1 = 1 } } // room -> admit; full -> reject
101 var t2: i64 = 0
102 if mb_batch_ready(4, 4, 0, 50) == 1 { t2 = 1 } // full batch -> ready even at wait 0
103 var t3: i64 = 0
104 if mb_batch_ready(2, 4, 50, 50) == 1 { t3 = 1 } // window elapsed w/ work -> ready
105 var t4: i64 = 0
106 if mb_batch_ready(2, 4, 10, 50) == 0 { t4 = 1 } // still accumulating -> NOT ready
107 var t5: i64 = 0
108 if mb_batch_ready(0, 4, 999, 50) == 0 { t5 = 1 } // empty -> NEVER ready
109 var t6: i64 = 0
110 if mb_batch_take(10, 4) == 4 { if mb_batch_take(2, 4) == 2 { t6 = 1 } } // take = min(qlen,max_batch)
111 // SIM T7: burst of 10 at t=0, max_batch=4, max_delay=10, cap=100 -> batches 4,4,2 = 3 batches, all 10 served
112 let arr: *i64 = sys_mmap(200 * 8) as *i64
113 var z: i64 = 0
114 while z < 200 { arr[z] = 0; z = z + 1 }
115 arr[0] = 10
116 let ob: *i64 = sys_mmap(64) as *i64
117 mb_sim(arr, 100, 4, 10, 100, ob)
118 var t7: i64 = 0
119 if ob[0]==10 { if ob[1]==0 { if ob[2]==3 { if ob[3]==10 { if ob[5]==0 { t7 = 1 } } } } }
120 // SIM T8 admission/backpressure: burst of 200 at t=0, cap=100 -> admitted 100, rejected 100
121 let arr2: *i64 = sys_mmap(200 * 8) as *i64
122 var z2: i64 = 0
123 while z2 < 200 { arr2[z2] = 0; z2 = z2 + 1 }
124 arr2[0] = 200
125 let ob2: *i64 = sys_mmap(64) as *i64
126 mb_sim(arr2, 100, 8, 10, 100, ob2)
127 var t8: i64 = 0
128 if ob2[0]==100 { if ob2[1]==100 { if ob2[4]<=100 { t8 = 1 } } } // never over cap
129 // neg1 liar: no batch ever exceeded max_batch in either sim
130 var neg1: i64 = 0
131 if ob[5]==0 { if ob2[5]==0 { neg1 = 1 } }
132 // neg2: single arrival waits then flushes as a batch of 1 (no starvation) -- arrival of 1 at t=0, delay 5
133 let arr3: *i64 = sys_mmap(200 * 8) as *i64
134 var z3: i64 = 0
135 while z3 < 200 { arr3[z3] = 0; z3 = z3 + 1 }
136 arr3[0] = 1
137 let ob3: *i64 = sys_mmap(64) as *i64
138 mb_sim(arr3, 100, 4, 5, 100, ob3)
139 var neg2: i64 = 0
140 if ob3[2]==1 { if ob3[3]==1 { neg2 = 1 } } // the lone request IS served (1 batch of 1)
141
142 var avg100: i64 = 0
143 if ob[2] > 0 { avg100 = (ob[3] * 100) / ob[2] }
144 mb_p(" T1 admit/backpressure: " as *u8); mb_pn(t1)
145 mb_p(" | T2 full-batch-ready: " as *u8); mb_pn(t2)
146 mb_p(" | T3 window-flush: " as *u8); mb_pn(t3)
147 mb_p(" | T4 still-accumulating: " as *u8); mb_pn(t4)
148 mb_p(" | T5 empty-never-ready: " as *u8); mb_pn(t5)
149 mb_p(" | T6 take=min: " as *u8); mb_pn(t6)
150 mb_p(" | T7 sim burst10->3 batches: " as *u8); mb_pn(t7)
151 mb_p(" | T8 admission 200@cap100->100 rej: " as *u8); mb_pn(t8)
152 mb_p(" | neg1 never-over-batch: " as *u8); mb_pn(neg1)
153 mb_p(" | neg2 no-starvation: " as *u8); mb_pn(neg2); mb_p("\n" as *u8)
154 mb_p(" sim(burst10,mb4): batches=" as *u8); mb_pn(ob[2]); mb_p(" avg_batch=" as *u8); mb_pn(avg100); mb_p("/100 (=" as *u8); mb_pn(avg100/100); mb_p(".x forward-pass reduction)\n" as *u8)
155
156 let sfd: i64 = sys_openat_wr("knowledge/status/mesh_batch.tsv" as *u8, 0x1a4)
157 if sfd >= 0 {
158 mb_w(sfd, "# nx_mesh_batch -- dynamic batching + queue admission (Triton dynamic-batcher analog)\n" as *u8)
159 mb_w(sfd, "sim_burst10_batches\t" as *u8); mb_wn(sfd, ob[2]); mb_w(sfd, "\n" as *u8)
160 mb_w(sfd, "sim_avg_batch_x100\t" as *u8); mb_wn(sfd, avg100); mb_w(sfd, "\n" as *u8)
161 mb_w(sfd, "admission_reject_over_cap\t" as *u8); mb_wn(sfd, ob2[1]); mb_w(sfd, "\n" as *u8)
162 sys_close(sfd)
163 }
164
165 var pass: i64 = 0
166 if t1==1 { if t2==1 { if t3==1 { if t4==1 { if t5==1 { if t6==1 { if t7==1 { if t8==1 { if neg1==1 { if neg2==1 { pass = 1 } } } } } } } } } }
167 if pass == 1 { mb_p("MESHBATCHGATE verdict=GREEN (queue+admission; dynamic coalesce by size-or-delay; never over-batch; no starvation; liar-killed)\n" as *u8); return 0 }
168 mb_p("MESHBATCHGATE verdict=RED\n" as *u8)
169 return 1
170}
171
172func main(argc: i64, argv: *i64) -> i64 {
173 if argc >= 2 {
174 let cmd: *u8 = argv[1] as *u8
175 if mb_streq(cmd, "sim" as *u8) == 1 {
176 if argc < 6 { mb_w(2, "usage: nx_mesh_batch sim <burst> <max_batch> <max_delay> <cap>\n" as *u8); return 2 }
177 let burst: i64 = mb_atoi(argv[2] as *u8)
178 let mbz: i64 = mb_atoi(argv[3] as *u8)
179 let mdl: i64 = mb_atoi(argv[4] as *u8)
180 let cap: i64 = mb_atoi(argv[5] as *u8)
181 let arr: *i64 = sys_mmap(400 * 8) as *i64
182 var z: i64 = 0
183 while z < 400 { arr[z] = 0; z = z + 1 }
184 arr[0] = burst
185 let ob: *i64 = sys_mmap(64) as *i64
186 mb_sim(arr, 300, mbz, mdl, cap, ob)
187 mb_p("sim burst=" as *u8); mb_pn(burst); mb_p(" max_batch=" as *u8); mb_pn(mbz); mb_p(" max_delay=" as *u8); mb_pn(mdl); mb_p("ms cap=" as *u8); mb_pn(cap); mb_p("\n" as *u8)
188 mb_p(" admitted=" as *u8); mb_pn(ob[0]); mb_p(" rejected=" as *u8); mb_pn(ob[1]); mb_p(" batches=" as *u8); mb_pn(ob[2]); mb_p(" total_batched=" as *u8); mb_pn(ob[3]); mb_p(" maxdepth=" as *u8); mb_pn(ob[4]); mb_p("\n" as *u8)
189 if ob[2] > 0 { mb_p(" avg batch size = " as *u8); mb_pn((ob[3]*100)/ob[2]); mb_p("/100 -> that many requests per GPU forward pass\n" as *u8) }
190 return 0
191 }
192 mb_w(2, "usage: nx_mesh_batch sim <burst> <max_batch> <max_delay> <cap> (no args = gate)\n" as *u8)
193 return 2
194 }
195 return mb_gate()
196}