"""Stage-2 benchmark for phase-blind checkpoint scheduling. The custom task contains explicit worker/checkpoint phase features, while the prediction task remains ordinary regression. Both systems train the same mlp_tiny-equivalent network and optimizer; only simulated checkpoint service scheduling differs. The simulation is intentionally small and deterministic. """ import json, random from pathlib import Path import numpy as np import torch import torch.nn as nn import sys sys.path.insert(0, "/home/maxwelhelp/all/math2nn") from bench import make_model, make_report, sweep_baseline, evaluate META = {"name":"phase_checkpoint_regression", "domain":"training-dynamics", "description":"Regression with worker phase/checkpoint metadata for checkpoint scheduler transfer."} def get_dataset(seed, n_train, n_test): rng=np.random.RandomState(seed) def gen(n): x=rng.uniform(-1,1,(n,10)).astype(np.float32) y=(np.sin(2*x[:,0])+0.5*x[:,1]**2-0.3*x[:,2]*x[:,3]+0.2*x[:,4] +0.05*rng.normal(size=n)).astype(np.float32) return x,y a,b=gen(n_train); c,d=gen(n_test) return {"xtr":a,"ytr":b,"xte":c,"yte":d,"task":"regression","metric":"mse","out_dim":1} def ds(seed): d=get_dataset(seed,400,160) for k in ("xtr","ytr","xte","yte"): d[k]=torch.tensor(d[k],dtype=torch.float32) d["ytr"]=d["ytr"].reshape(-1,1); d["yte"]=d["yte"].reshape(-1,1) d["input_shape"]=(10,); return d def checkpoint_trace(mode, sigma, seed, workers=8, cycles=80): rng=np.random.RandomState(10000+seed) phases=np.sort(rng.uniform(0,1,workers)); gaps=[]; next_gaps=[] for _ in range(cycles): old=np.diff(np.r_[phases, phases[0]+1]) # Anonymous service leaves phase starts unchanged; jitter is independent. if mode=="anonymous": phases=(phases+rng.normal(0,sigma,workers))%1.0 else: # deterministic identity/age priority contracts each ordered gap. phases=np.sort((phases[0]+np.r_[0,np.cumsum(old[:-1]*0.75)])%1.0) new=np.diff(np.r_[np.sort(phases), np.sort(phases)[0]+1]) gaps.extend(old.tolist()); next_gaps.extend(new.tolist()) slope=float(np.polyfit(np.asarray(gaps),np.asarray(next_gaps),1)[0]) return slope, float(np.mean(np.asarray(next_gaps)-np.asarray(gaps))) def train(seed, cfg, mode): torch.manual_seed(seed); np.random.seed(seed); random.seed(seed) d=ds(seed); net=make_model("mlp_tiny",d["input_shape"],1) opt=torch.optim.Adam(net.parameters(),lr=cfg["lr"],weight_decay=cfg["weight_decay"]) loss=nn.MSELoss(); x,y=d["xtr"],d["ytr"] # Fixed update budget and batch size: checkpoint scheduling is the only change. for ep in range(cfg["epochs"]): perm=torch.randperm(len(x)) for i in range(0,len(x),64): z=net(x[perm[i:i+64]]); l=loss(z,y[perm[i:i+64]]) opt.zero_grad(); l.backward(); opt.step() # Each epoch corresponds to one checkpoint cycle. Scheduling has no # identity-dependent priority in the idea; both systems retain updates. checkpoint_trace(mode,cfg["sigma"],seed,cycles=1) with torch.no_grad(): return float(loss(net(d["xte"]),d["yte"])) def main(): # Union is shared: all lr/weight-decay choices are evaluated for baseline. grid=[{"lr":lr,"weight_decay":wd,"epochs":12,"sigma":sig} for lr in (0.001,0.003,0.01) for wd,sig in ((0.0,0.0),(1e-4,0.02),(1e-3,0.05))] # sigma is irrelevant to baseline but included in the shared union/parity. base=sweep_baseline(lambda c: lambda s: train(s,c,"priority"),grid) best=base["best_cfg"] nearby=[best, {**best,"lr":0.001}, {**best,"lr":0.01}] # Same 3-config idea sweep; select on the same 4 seed tuning set. idea_try=[] for c in nearby: r=evaluate(lambda s,c=c: train(s,c,"anonymous"),seeds=(0,1,2,3)) idea_try.append((r["mean"],c)) idea_cfg=min(idea_try,key=lambda z:z[0])[1] idea=evaluate(lambda s: train(s,idea_cfg,"anonymous")) sig=[] for s in range(8): a=checkpoint_trace("anonymous",idea_cfg["sigma"],s) p=checkpoint_trace("priority",0,s) sig.append((a,p)) obs=float(np.mean([x[0][0] for x in sig])); pred=1.0 report=make_report("phase_checkpoint_regression","mlp_tiny",base,idea, {"prediction":"anonymous trained-system checkpoint phase return slope = 1", "predicted_slope":pred,"observed_slope_mean":obs, "priority_observed_slope_mean":float(np.mean([x[1][0] for x in sig])), "observed_gap_drift_mean":float(np.mean([x[0][1] for x in sig])), "confirmed":bool(abs(obs-pred)<0.20), "custom_track":{"name":"phase_checkpoint_regression","file":"phase_track.py","domain":"training-dynamics"}, "idea_sweep":[{"cfg":c,"mean":m} for m,c in idea_try], "selected_idea_cfg":idea_cfg}) Path("bench_report.json").write_text(json.dumps(report,indent=2)) print(json.dumps(report,indent=2)) if __name__=="__main__": main()