import os, sys, json, math, random import numpy as np import torch import torch.nn as nn sys.path.insert(0, '/home/maxwelhelp/all/math2nn') from bench import make_model, sweep_baseline, make_report from bench.protocol import evaluate, DEFAULT_SEEDS META = { 'name': 'markov_stream_regression', 'domain': 'optimizer', 'description': 'Ordered stationary AR(1) covariate stream with nonlinear regression targets; tests gradient estimation under Markov dependence.' } def get_dataset(seed, n_train, n_test): rng = np.random.default_rng(seed) rho = 0.9 def ar(n): x = np.empty((n, 10), dtype=np.float32) x[0] = rng.normal(size=10) q = math.sqrt(1.0 - rho * rho) for i in range(1, n): x[i] = rho * x[i-1] + q * rng.normal(size=10) return x xtr, xte = ar(n_train), ar(n_test) def yfun(x): y = (1.5*x[:, 0] - 1.1*x[:, 1] + .7*x[:, 2]**2 + .4*np.sin(x[:, 3]) + .25*x[:, 4]*x[:, 5] + .15*rng.normal(size=len(x))) return y.astype(np.float32).reshape(-1, 1) return {'xtr': xtr, 'ytr': yfun(xtr), 'xte': xte, 'yte': yfun(xte), 'task': 'regression', 'metric': 'mse', 'out_dim': 1} def seed_all(seed): random.seed(seed); np.random.seed(seed); torch.manual_seed(seed) if torch.cuda.is_available(): torch.cuda.manual_seed_all(seed) def flat_grads(net): return torch.cat([p.grad.detach().reshape(-1) for p in net.parameters() if p.grad is not None]) def per_sample_grad(net, x, y): vals = [] lossf = nn.MSELoss() for j in range(len(x)): net.zero_grad(set_to_none=True) lossf(net(x[j:j+1]), y[j:j+1]).backward() vals.append(flat_grads(net)) return torch.stack(vals) def clip(v, bound): n = torch.linalg.vector_norm(v) return v if float(n) <= bound else v * (bound / (n + 1e-12)) def train_one(ds, seed, lr, method, epochs=8, block=32, levels=2, bound=10.0): seed_all(seed) net = make_model('mlp_tiny', (10,), 1) device = 'cuda' if torch.cuda.is_available() else 'cpu' try: net = net.to(device) x = torch.as_tensor(ds['xtr'], dtype=torch.float32, device=device) y = torch.as_tensor(ds['ytr'], dtype=torch.float32, device=device) xt = torch.as_tensor(ds['xte'], dtype=torch.float32, device=device) yt = torch.as_tensor(ds['yte'], dtype=torch.float32, device=device) opt = torch.optim.Adam(net.parameters(), lr=lr) rng = np.random.default_rng(seed + 991) p = np.array([.5, .3, .2], dtype=float) costs = [] corr_norms, grad_seq, clipped = [], [], 0 cursor = 0 base = block * (2 ** levels) for _ in range(epochs): while cursor + base <= len(x): z = x[cursor:cursor+base]; q = y[cursor:cursor+base]; cursor += base if method == 'baseline': net.zero_grad(set_to_none=True) loss = ((net(z) - q) ** 2).mean() loss.backward(); gh = flat_grads(net); costs.append(base) else: # Shared-prefix coupled multilevel estimator: g0 plus one inverse-probability correction. l = int(rng.choice(3, p=p)) b0 = block net.zero_grad(set_to_none=True) g0 = per_sample_grad(net, z[:b0], q[:b0]).mean(0) if l == 0: gh = g0; costs.append(b0) else: bf = b0 * (2 ** l); bp = bf // 2 gf = per_sample_grad(net, z[:bf], q[:bf]).mean(0) gp = per_sample_grad(net, z[:bp], q[:bp]).mean(0) delta = gf - gp corr_norms.append(float(torch.linalg.vector_norm(delta))) delta = clip(delta, bound) if float(torch.linalg.vector_norm(delta)) >= bound - 1e-6: clipped += 1 gh = g0 + delta / float(p[l]); costs.append(bf) gh = clip(gh, bound) if float(torch.linalg.vector_norm(gh)) >= bound - 1e-6: clipped += 1 net.zero_grad(set_to_none=True) off = 0 for par in net.parameters(): n = par.numel(); par.grad = gh[off:off+n].reshape_as(par).clone(); off += n opt.step(); grad_seq.append(gh.detach().cpu().numpy()) cursor = 0 net.eval() with torch.no_grad(): metric = float(((net(xt)-yt)**2).mean()) a = np.asarray(grad_seq) ac = float(np.corrcoef(a[:-1,0], a[1:,0])[0,1]) if len(a) > 3 else float('nan') return metric, net, {'grad_lag1': ac, 'correction_norm_mean': float(np.mean(corr_norms)) if corr_norms else 0.0, 'clip_fraction': clipped/max(1, len(grad_seq)), 'updates': len(grad_seq), 'mean_grad_norm': float(np.linalg.norm(a, axis=1).mean())} except RuntimeError: if device == 'cuda': torch.cuda.empty_cache() return train_one_cpu(ds, seed, lr, method, epochs, block, levels, bound) raise def train_one_cpu(ds, seed, lr, method, epochs=8, block=32, levels=2, bound=10.0): old = torch.cuda.is_available torch.cuda.is_available = lambda: False try: return train_one(ds, seed, lr, method, epochs, block, levels, bound) finally: torch.cuda.is_available = old def run(): lrs = [0.001, 0.003, 0.006] seeds = tuple(range(8)) def base_fn(cfg): lr = cfg['lr'] return lambda seed: train_one(get_dataset(seed, 768, 256), seed, lr, 'baseline')[0] base = sweep_baseline(base_fn, [{'lr': v} for v in lrs], seeds=seeds[:4]) def idea_fn(cfg): lr = cfg['lr'] return lambda seed: train_one(get_dataset(seed, 768, 256), seed, lr, 'coupled')[0] # Evaluate the baseline-best setting and two nearby settings; all are in the baseline union. idea_trials = [] for cfg in [{'lr': v} for v in lrs]: res = evaluate(idea_fn(cfg), seeds=seeds) idea_trials.append({'cfg': cfg, 'result': res}) idea = min(idea_trials, key=lambda z: z['result']['mean'])['result'] idea['sweep'] = [{'cfg': z['cfg'], 'mean': z['result']['mean'], 'std': z['result']['std']} for z in idea_trials] sig = [] chosen_lr = min(idea_trials, key=lambda z: z['result']['mean'])['cfg']['lr'] for seed in seeds: ds = get_dataset(seed, 768, 256) m, _, s = train_one(ds, seed, chosen_lr, 'coupled') sig.append({'seed': seed, 'metric': m, **s}) signature = { 'prediction': 'shared fine-minus-coarse corrections have smaller norm than raw fine gradients under correlated streams', 'observed_correction_norm_mean': float(np.mean([r['correction_norm_mean'] for r in sig])), 'observed_mean_gradient_norm': float(np.mean([r['mean_grad_norm'] for r in sig])), 'observed_gradient_lag1_mean': float(np.mean([r['grad_lag1'] for r in sig])), 'observed_clip_fraction_mean': float(np.mean([r['clip_fraction'] for r in sig])), 'confirmed': bool(np.mean([r['correction_norm_mean'] for r in sig]) < np.mean([r['mean_grad_norm'] for r in sig])) } rep = make_report('markov_stream_regression', 'mlp_tiny', base, idea, signature) rep['custom_track'] = {'name': META['name'], 'file': 'markov_stream_bench.py', 'domain': META['domain']} os.makedirs('artifacts', exist_ok=True) with open('artifacts/bench_report.json', 'w') as f: json.dump(rep, f, indent=2) print(json.dumps(rep, indent=2)) if __name__ == '__main__': run()