diff --git a/c/glm.c b/c/glm.c index 4663aba..07fae16 100644 --- a/c/glm.c +++ b/c/glm.c @@ -828,6 +828,26 @@ static int g_gr_on=0; /* grammatica caricata e walker vivo */ static int g_gr_armed=0; /* lazy: parte dal primo byte ammesso dalla radice (salta i preamboli) */ static int g_gr_max=24; static uint64_t g_gr_prop=0, g_gr_acc=0; +static FILE *g_route_fp=NULL; /* ROUTE_TRACE=: dump per-position top-K routing (ids:gates) + * per layer — offline co-activation / coupling analysis. Zero + * effect on computation; measurement only. */ +static int g_route_call=0; +/* COUPLE=<.coli_pairs>: coupling-scored cross-layer prefetch. The routing of layer L + * strongly constrains the routing of L+1/L+2 (measured: median co-activation lift 1.8x + * over independence, p99 40x, and the structure TRANSFERS across workloads — it is a + * property of the model, not the session). An offline table (tools/route_pairs.py, + * built from ROUTE_TRACE dumps) maps (layer, expert) -> top co-activated experts of the + * next layer(s); after FASE A routing we score candidates by summing counts over the + * position's routed set and enqueue the top COUPLE_K non-resident ones into the SAME + * pilot ring (worker, residency re-check, safety invariants unchanged). Unlike PILOT, + * no router matmul is needed — prediction is a table lookup on ids the layer just + * produced. Hints only: a wrong prediction costs bandwidth, never output. */ +#define CP_M 16 +static int g_couple=0, g_couple_k=8, g_couple_d=1; +static int16_t *cp_pred=NULL; /* [(L*2+(dL-1))*E + e]*CP_M + j -> target id (-1 none) */ +static float *cp_cnt=NULL; +static long g_cp_enq=0; +static void couple_prefetch(Model *m, int layer, const int *idx, int Ke); static int g_looka=0; /* LOOKA=1: misura (solo contatori, zero effetti) quanto il routing MoE * e' predicibile IN ANTICIPO — la domanda che decide se un prefetch * pilotato dal router puo' riempire i tempi morti del disco. @@ -1822,8 +1842,16 @@ static void moe(Model *m, Layer *l, int layer, float *x, int S, float *out){ } if(c->norm_topk){ float sm=0; for(int kk=0;kkrouted_scale; + if(g_route_fp){ /* ROUTE_TRACE: one line per (position, layer) */ + fprintf(g_route_fp,"%d %d %d",g_route_call,s,layer); + for(int kk=0;kkn_layers){ int Ke=keff[0]; if(m->enr[layer]>0){ /* [0] vs routing del token precedente */ @@ -2194,6 +2222,77 @@ static void *pilot_worker(void *arg){ } return NULL; } +/* parse .coli_pairs (see tools/route_pairs.py): "COLIPAIRS 1 " then + * "
f:c f:c ..." lines. Needs c->n_experts/n_layers -> called post-init. */ +static void couple_load(Model *m, const char *path){ + Cfg *c=&m->c; int E=c->n_experts, NL=c->n_layers; + FILE *f=fopen(path,"rb"); + if(!f){ fprintf(stderr,"[COUPLE] cannot open %s\n",path); return; } + char magic[16]; int ver=0; long n=0; + if(fscanf(f,"%15s %d %ld",magic,&ver,&n)!=3 || strcmp(magic,"COLIPAIRS") || ver!=1){ + fprintf(stderr,"[COUPLE] %s: bad header\n",path); fclose(f); return; } + size_t cells=(size_t)NL*2*E*CP_M; + cp_pred=malloc(cells*sizeof(int16_t)); cp_cnt=calloc(cells,sizeof(float)); + if(!cp_pred||!cp_cnt){ fprintf(stderr,"[COUPLE] OOM\n"); free(cp_pred); free(cp_cnt); cp_pred=NULL; fclose(f); return; } + for(size_t i=0;i0){ /* line-based: a malformed line cannot eat the next */ + char *p=ln; int L,dL,e; int nc=0; + if(sscanf(p,"%d %d %d%n",&L,&dL,&e,&nc)!=3) continue; + p+=nc; + if(L<0||L>=NL||(dL!=1&&dL!=2)||e<0||e>=E) continue; + size_t base=((size_t)(L*2+(dL-1))*E+e)*CP_M; + int j=0; + while(j=0&&fec; int E=c->n_experts; + if(E>512) return; + if(!pilot_m){ pilot_m=m; pthread_t t; pthread_create(&t,NULL,pilot_worker,NULL); } + for(int dL=1; dL<=g_couple_d; dL++){ + int lt=layer+dL; + if(lt>=c->n_layers || !m->L[lt].sparse) continue; + float sc[512]; memset(sc,0,(size_t)E*sizeof(float)); + for(int kk=0;kk=0;j++) sc[cp_pred[base+j]]+=cp_cnt[base+j]; + } + for(int kk=0;kkbv){bv=sc[e];best=e;} + if(best<0) break; + sc[best]=0; + int found=0; /* residency scan, same locking as pilot */ + pthread_mutex_lock(&g_pilot_mx); + ESlot *P=m->pin[lt]; + for(int z=0;znpin[lt] && !found;z++) if(P[z].eid==best) found=1; + ESlot *Sl=m->ecache[lt]; + for(int z=0;zecn[lt] && !found;z++) if(Sl[z].eid==best) found=1; + pthread_mutex_unlock(&g_pilot_mx); + if(!found){ + unsigned w=__atomic_load_n(&pilot_w,__ATOMIC_RELAXED); + if(w-__atomic_load_n(&pilot_r,__ATOMIC_ACQUIRE)<4096){ + pilot_q[w&4095].l=lt; pilot_q[w&4095].e=best; + __atomic_store_n(&pilot_w,w+1,__ATOMIC_RELEASE); + g_cp_enq++; + } + } + } + } +} static void pilot_prefetch(Model *m, int lnext, const float *x, int S){ Cfg *c=&m->c; Layer *l=&m->L[lnext]; int D=c->hidden, E=c->n_experts; int K = g_pilot_ktopk ? g_pilot_k : c->topk; @@ -2858,6 +2957,7 @@ static void run_text(Model *m, const char *snap, const char *prompt, int ngen){ printf("speculation: %.2f tokens/forward (%llu forwards per %llu tokens) | MTP acceptance %.0f%% (%llu/%llu)\n", m->n_fw?(double)m->n_emit/m->n_fw:1.0, (unsigned long long)m->n_fw, (unsigned long long)m->n_emit, m->mtp_prop?100.0*m->mtp_acc/m->mtp_prop:0.0, (unsigned long long)m->mtp_acc, (unsigned long long)m->mtp_prop); + if(g_cp_enq) printf("couple: %ld cross-layer prefetch hints enqueued\n", g_cp_enq); if(g_gr_prop) printf("grammar: %.0f%% acceptance (%llu/%llu forced drafts)\n", 100.0*g_gr_acc/g_gr_prop, (unsigned long long)g_gr_acc, (unsigned long long)g_gr_prop); #ifdef COLI_CUDA @@ -3906,6 +4006,11 @@ int main(int argc, char **argv){ if(g_pipe_nw<1) g_pipe_nw=1; g_direct = getenv("DIRECT")?atoi(getenv("DIRECT")):0; g_idot = getenv("IDOT")?atoi(getenv("IDOT")):1; /* 0 = kernel f32 esatti (A/B) */ + if(getenv("ROUTE_TRACE")&&*getenv("ROUTE_TRACE")){ + g_route_fp=fopen(getenv("ROUTE_TRACE"),"w"); + if(!g_route_fp) fprintf(stderr,"[ROUTE_TRACE] cannot open %s\n",getenv("ROUTE_TRACE")); + else fprintf(stderr,"[ROUTE_TRACE] logging routing to %s\n",getenv("ROUTE_TRACE")); + } g_repin = getenv("REPIN")?atoi(getenv("REPIN")):0; /* RFC: re-pin ogni n token emessi (0=off) / live re-pin every n emitted tokens (0=off) */ g_absorb = getenv("ABSORB")?atoi(getenv("ABSORB")):-1; /* -1 auto: assorbita per S<=4 */ g_dsa_force = getenv("DSA_FORCE")?atoi(getenv("DSA_FORCE")):0; @@ -3985,6 +4090,13 @@ int main(int argc, char **argv){ /* HOT-STORE: PIN= [PIN_GB=g] -> top expert per frequenza fissi in RAM. * Va PRIMA di cap_for_ram: i pinnati contano nel residente. */ if(getenv("PIN")) pin_load(&m, getenv("PIN"), getenv("PIN_GB")?atof(getenv("PIN_GB")):10.0); + if(getenv("COUPLE")&&*getenv("COUPLE")){ /* coupling-scored cross-layer prefetch */ + g_couple_k=getenv("COUPLE_K")?atoi(getenv("COUPLE_K")):8; + if(g_couple_k<1)g_couple_k=1; if(g_couple_k>32)g_couple_k=32; + g_couple_d=getenv("COUPLE_D")?atoi(getenv("COUPLE_D")):1; + if(g_couple_d<1)g_couple_d=1; if(g_couple_d>2)g_couple_d=2; + couple_load(&m, getenv("COUPLE")); + } /* CACHE CHE IMPARA: l'uso degli expert si accumula in /.coli_usage tra le sessioni; * all'avvio i piu' usati vengono auto-pinnati in RAM (meta' del budget expert: il pin * conosce la TUA storia, la LRU si adatta alla sessione). AUTOPIN=0 disattiva. */ diff --git a/c/tools/route_coupling_report.py b/c/tools/route_coupling_report.py new file mode 100644 index 0000000..c4ec804 --- /dev/null +++ b/c/tools/route_coupling_report.py @@ -0,0 +1,120 @@ +#!/usr/bin/env python3 +"""Copula/Fréchet analysis of MoE routing traces (ROUTE_TRACE output). + +Input lines: " id:gate id:gate ..." — one line per (position, layer). + +Analyses: + 1. Cross-layer dependence screen (consecutive layer pairs L -> L+1): + observed pair co-activation counts vs (a) independence product and + (b) the Fréchet–Hoeffding upper bound min(p_e, p_f). Reports how much + dependence structure exists beyond marginals — the existence proof for + coupling-aware prefetch. + 2. Prefetch simulation at equal budget (held-out split): + predict layer-(L+1) experts from layer-L routing: + marginal : top-B experts of L+1 by marginal frequency (today's heat logic) + coupled : top-B by sum of pair co-activation counts conditioned on the + observed L set (a mixture approximation of the pair-copula) + Metric: recall of the true top-8 at budgets B = 8, 16, 32. + 3. Depth-2 (L -> L+2) repeat of (2) — the horizon where prefetch has a full + disk round-trip to work in (LOOKA kind [2]). +""" +import sys +from collections import defaultdict +import numpy as np + +def load(path): + """group by (forward, position): the call counter increments per moe() CALL + (i.e. per layer), so forwards are reconstructed from layer wrap-arounds""" + rows = defaultdict(dict) # (fwd,pos) -> {layer: [ids]} + fwd, prev_layer = 0, -1 + for line in open(path): + p = line.split() + if len(p) < 4: continue + pos, layer = int(p[1]), int(p[2]) + if layer < prev_layer: fwd += 1 + prev_layer = layer + ids = [int(t.split(":")[0]) for t in p[3:]] + rows[(fwd, pos)][layer] = ids + return rows + +def main(): + path = sys.argv[1] + rows = load(path) + if len(sys.argv) > 2: # transfer mode: train on file1, test on file2 + rows2 = load(sys.argv[2]) + keys = sorted(rows.keys()); keys2 = sorted(rows2.keys()) + base = 1 + max(k[0] for k in keys) + rows.update({(f + base, p): v for (f, p), v in rows2.items()}) + train = keys + test = [(f + base, p) for (f, p) in keys2] + mode = f"TRANSFER train={sys.argv[1].split('/')[-1]} test={sys.argv[2].split('/')[-1]}" + else: + keys = sorted(rows.keys()) + split = int(len(keys) * 0.7) + train, test = keys[:split], keys[split:] + mode = "in-trace 70/30 split" + layers = sorted({L for r in rows.values() for L in r}) + E = 1 + max(i for r in rows.values() for ids in r.values() for i in ids) + print(f"trace: {len(rows)} positions, {len(layers)} routed layers, E={E} " + f"(train {len(train)} / test {len(test)}) [{mode}]") + + # marginals + consecutive-pair co-occurrence from TRAIN + marg = {L: np.zeros(E) for L in layers} + pair = {} # (L, dL) -> sparse dict (e,f)->count + for dL in (1, 2): + for L in layers: + if L + dL in layers: pair[(L, dL)] = defaultdict(int) + for k in train: + r = rows[k] + for L, ids in r.items(): + for e in ids: marg[L][e] += 1 + for (L, dL), d in pair.items(): + if L in r and L + dL in r: + for e in r[L]: + for f in r[L + dL]: d[(e, f)] += 1 + + # 1. dependence screen on L->L+1 + lifts, bound_ratio = [], [] + N = len(train) + for (L, dL), d in pair.items(): + if dL != 1: continue + for (e, f), c in d.items(): + pe, pf = marg[L][e] / N, marg[L + 1][f] / N + if pe * pf <= 0: continue + lifts.append((c / N) / (pe * pf)) + bound_ratio.append((c / N) / min(pe, pf)) + lifts = np.array(lifts); bound_ratio = np.array(bound_ratio) + print(f"\n[1] L->L+1 co-activation, {len(lifts)} observed pairs:") + print(f" lift vs independence: median {np.median(lifts):.2f} " + f"p90 {np.percentile(lifts,90):.2f} p99 {np.percentile(lifts,99):.2f}") + print(f" fraction of pairs at >50% of the Fréchet upper bound: " + f"{(bound_ratio>0.5).mean()*100:.1f}% (>90%: {(bound_ratio>0.9).mean()*100:.1f}%)") + + # 2/3. prefetch recall on TEST + for dL in (1, 2): + print(f"\n[{1+dL}] prefetch L->L+{dL}, recall of true top-8 on held-out positions:") + for B in (8, 16, 32): + hit_m = tot = hit_c = 0 + for k in test: + r = rows[k] + for L in layers: + if L not in r or L + dL not in r or (L, dL) not in pair: continue + true = set(r[L + dL]) + pm = np.argsort(marg[L + dL])[::-1][:B] + hit_m += len(true & set(pm.tolist())) + score = defaultdict(float) + d = pair[(L, dL)] + for (e, f), c in d.items(): + if e in r[L]: score[f] += c + pc = sorted(score, key=score.get, reverse=True)[:B] + if len(pc) < B: # back-fill with marginals + for f in np.argsort(marg[L + dL])[::-1]: + if f not in pc: pc.append(int(f)) + if len(pc) == B: break + hit_c += len(true & set(pc)) + tot += len(true) + print(f" budget {B:3d}/layer: marginal {hit_m/tot*100:5.1f}% " + f"coupled {hit_c/tot*100:5.1f}% (+{(hit_c-hit_m)/tot*100:.1f}pp)") + +if __name__ == "__main__": + main() diff --git a/c/tools/route_pairs.py b/c/tools/route_pairs.py new file mode 100644 index 0000000..6af2937 --- /dev/null +++ b/c/tools/route_pairs.py @@ -0,0 +1,63 @@ +#!/usr/bin/env python3 +"""Build a .coli_pairs cross-layer coupling table from ROUTE_TRACE dumps. + +Input: one or more ROUTE_TRACE files ("call pos layer id:gate ...", one line per +(position, layer)). Output: a text table of the top-M layer-(L+dL) experts most +co-activated with each (layer L, expert e) conditioning event: + + COLIPAIRS 1 +
f1:c1 f2:c2 ... (up to M) + +The engine (COUPLE=) scores layer-(L+dL) candidates by summing counts over +the observed layer-L expert set — the mixture approximation of the pair copula +that measured +3.6..+9.4pp prefetch recall over marginal heat, in- and +cross-domain (see docs in the PR). Counts are raw co-occurrences; the consumer +only needs their ranking, so no normalization is stored. + +Usage: python3 tools/route_pairs.py out.coli_pairs trace1.txt [trace2.txt ...] +""" +import sys +from collections import defaultdict + +M = 16 + +def main(): + out_path, traces = sys.argv[1], sys.argv[2:] + pair = defaultdict(lambda: defaultdict(int)) # (L, dL, e) -> {f: count} + for path in traces: + cur = {} # layer -> ids (within one forward) + prev_layer = -1 + def flush(): + layers = sorted(cur) + for i, L in enumerate(layers): + for dL in (1, 2): + if L + dL in cur: + for e in cur[L]: + d = pair[(L, dL, e)] + for f in cur[L + dL]: d[f] += 1 + # group lines by position within a forward: lines arrive layer-major + # (all positions of layer L, then layer L+1, ...); regroup per position + rows = defaultdict(dict) # (fwd,pos) -> {layer: ids} + fwd = 0 + for line in open(path): + p = line.split() + if len(p) < 4: continue + pos, layer = int(p[1]), int(p[2]) + if layer < prev_layer: fwd += 1 + prev_layer = layer + rows[(fwd, pos)][layer] = [int(t.split(":")[0]) for t in p[3:]] + for r in rows.values(): + cur = r; flush() + print(f"{path}: {len(rows)} positions", file=sys.stderr) + + lines = [] + for (L, dL, e), d in sorted(pair.items()): + top = sorted(d, key=d.get, reverse=True)[:M] + lines.append(f"{L} {dL} {e} " + " ".join(f"{f}:{d[f]}" for f in top)) + with open(out_path, "w") as f: + f.write(f"COLIPAIRS 1 {len(lines)}\n") + f.write("\n".join(lines) + "\n") + print(f"wrote {out_path}: {len(lines)} conditioning entries", file=sys.stderr) + +if __name__ == "__main__": + main()