code wiki / (root) / nx_mesh_batch.nx

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}