← back to Tk11438 Postgres Migration

verification/ken-rollout/rollout.cjs

113 lines

// TK11438 approved Ken-only transport migration. Never invokes the email wrapper.
const fs=require('fs'),path=require('path'),crypto=require('crypto'),cp=require('child_process'),assert=require('assert/strict'),os=require('os');
const ROOT='/Users/macstudio3/Projects/Ken',DIR=ROOT+'/kalshi-dash',OUT=__dirname,PRIV=OUT+'/private';
const ENV=DIR+'/.env',WRAPPER=DIR+'/follow-the-winners-check.sh',PLIST='/Users/macstudio3/Library/LaunchAgents/com.steve.ken-reconcile-canary.plist',DUMP='/Users/macstudio3/.pm2/dump.pm2';
const manifest=JSON.parse(fs.readFileSync('/Users/macstudio3/Projects/tk11438-postgres-migration/next-batch.json'));
const parse=require('/Users/macstudio3/Projects/patterndesignlab/node_modules/dotenv').parse;
const {Client}=require(DIR+'/node_modules/pg');
const deps='/Users/macstudio3/.npm-global/lib/node_modules/pm2/node_modules/';
const axon=require(deps+'pm2-axon'),rpc=require(deps+'pm2-axon-rpc');
const mode=process.argv[2];assert(['prepare','apply','verify','reload','observe','rollback'].includes(mode));
const sha=b=>crypto.createHash('sha256').update(b).digest('hex');
const json=p=>JSON.parse(fs.readFileSync(p));
const record=(name,x)=>fs.writeFileSync(OUT+'/'+name+'.json',JSON.stringify({at:new Date().toISOString(),...x},null,2)+'\n');
const run=(cmd,args,opts={})=>cp.execFileSync(cmd,args,{encoding:'utf8',timeout:15000,maxBuffer:16*1024*1024,...opts});
const git=(...args)=>run('git',args,{cwd:ROOT}).trim();
const socketURL=s=>{const u=new URL(s);assert(['localhost','127.0.0.1','[::1]',''].includes(u.hostname));u.searchParams.set('host','/tmp');return u.toString();};
const summarize=s=>{const u=new URL(s);return {host:u.hostname,database:u.pathname.slice(1),socket:u.searchParams.get('host'),port:u.port||'5432'};};
const liveRpc=(method,opts={})=>new Promise((resolve,reject)=>{const sock=axon.socket('req'),client=new rpc.Client(sock);const t=setTimeout(()=>{sock.close();reject(Error('RPC timeout: '+method));},30000);sock.on('error',e=>{clearTimeout(t);sock.close();reject(e);});sock.connect('/Users/macstudio3/.pm2/rpc.sock');client.call(method,opts,(err,data)=>{clearTimeout(t);sock.close();err?reject(Error('RPC failed: '+method)):resolve(data);});});
function unique(list){const a=list.filter(x=>x.name==='ken');assert.equal(a.length,1);const x=a[0],e=x.pm2_env||x;assert.equal(e.pm_cwd,DIR);assert.equal(e.pm_exec_path,DIR+'/start.sh');return x;}
function fields(e){return {DATABASE_URL:e.DATABASE_URL,KEN_DATABASE_URL:e.KEN_DATABASE_URL,env:{DATABASE_URL:e.env?.DATABASE_URL,KEN_DATABASE_URL:e.env?.KEN_DATABASE_URL}};}
function setFields(e,v){e.DATABASE_URL=v.DATABASE_URL;e.KEN_DATABASE_URL=v.KEN_DATABASE_URL;e.env.DATABASE_URL=v.env.DATABASE_URL;e.env.KEN_DATABASE_URL=v.env.KEN_DATABASE_URL;}
function atomic(file,bytes,expected){assert.equal(sha(fs.readFileSync(file)),expected,'concurrent change: '+file);const temp=file+'.TK11438-'+process.pid;fs.writeFileSync(temp,bytes,{flag:'wx',mode:fs.statSync(file).mode&0o777});assert.equal(sha(fs.readFileSync(file)),expected,'changed during preparation: '+file);fs.renameSync(temp,file);assert.equal(sha(fs.readFileSync(file)),sha(bytes));}
function patchDump(expected,next){const bytes=fs.readFileSync(DUMP),d=JSON.parse(bytes),e=unique(d);assert.deepEqual(fields(e),expected);setFields(e,next);const reverted=JSON.parse(JSON.stringify(d));setFields(unique(reverted),expected);assert.deepEqual(reverted,JSON.parse(bytes),'out-of-scope dump edit');atomic(DUMP,JSON.stringify(d,null,2),sha(bytes));return {before_sha:sha(bytes),after_sha:sha(fs.readFileSync(DUMP)),only_target_fields:true};}
async function withDb(url,fn){const c=new Client({connectionString:url,application_name:'TK11438-Ken-'+mode,connectionTimeoutMillis:4000,options:'-c default_transaction_read_only=on -c statement_timeout=5000'});try{await c.connect();return await fn(c);}finally{await c.end().catch(()=>{});}}
async function snapshot(label){
 const list=await liveRpc('getMonitorData'),s=unique(list),e=s.pm2_env;assert.equal(e.status,'online');const env=parse(fs.readFileSync(ENV));const identities=[];
 for(const key of ['DATABASE_URL','KEN_DATABASE_URL'])identities.push({key,...await withDb(env[key],async c=>(await c.query("SELECT current_database() database,current_user role,inet_client_addr()::text addr,current_setting('transaction_read_only') readonly")).rows[0])});
 const settings=await withDb(env.DATABASE_URL,async c=>(await c.query("SELECT config->>'safe_mode' safe_mode,config->>'trading_on' trading_on,config->>'kalshi_env' kalshi_env,md5(config::text) config_hash,md5((config-'weather_cache')::text) stable_settings_hash,config->'weather_cache'->>'fetched_at' weather_cache_fetched_at FROM risk_state ORDER BY updated_at DESC LIMIT 1")).rows[0]);
 const kenSettings=await withDb(env.KEN_DATABASE_URL,async c=>(await c.query("SELECT key,md5(value::text) value_hash FROM ken_config WHERE key IN ('live_run','trade_config') ORDER BY key")).rows);
 const canarySource=fs.readFileSync(DIR+'/scripts/reconcile-canary.mjs','utf8');const sql=canarySource.match(/await pool\.query\(`([\s\S]*?)`\)/)[1];
 const sums=await withDb(env.KEN_DATABASE_URL,async c=>(await c.query(sql)).rows[0]);
 const num=x=>Number(x);const reconciliation={returned:Math.abs(num(sums.rollup_returned)-num(sums.ledger_returned))<=200,invested:Math.abs(num(sums.rollup_invested)-num(sums.ledger_invested))<=200,balance:Math.abs(num(sums.balance)-(num(sums.seed)-num(sums.rollup_invested)+num(sums.rollup_returned)))<=200};
 const http=[];const lan=Object.values(os.networkInterfaces()).flat().find(x=>x.family==='IPv4'&&!x.internal&&/^192\.168\./.test(x.address));assert(lan,'LAN interface needed for real non-loopback auth test');
 for(const route of ['/api/setup','/api/trading/history'])for(const auth of [false,true]){
  const r=await fetch('http://'+lan.address+':7810'+route,{headers:auth?{authorization:'Basic '+Buffer.from((env.ADMIN_USER||'admin')+':'+env.ADMIN_PASSWORD).toString('base64')}:{},signal:AbortSignal.timeout(10000)});const text=await r.text();assert.equal(r.status,auth?200:401);if(auth)JSON.parse(text);http.push({route,auth,status:r.status,sha256:sha(text),bytes:Buffer.byteLength(text)});
 }
 let tcp='';try{tcp=run('lsof',['-nP','-a','-p',String(s.pid),'-iTCP']);}catch(err){if(err.status!==1)throw err;tcp=err.stdout||'';}
 const processDbTcp=tcp.split('\n').filter(l=>/:5432\b/.test(l));
 const saved=unique(json(DUMP));const savedMatch=['DATABASE_URL','KEN_DATABASE_URL'].every(k=>e[k]===env[k]&&e.env[k]===env[k]&&saved[k]===env[k]&&saved.env[k]===env[k]);
 const activity=await withDb('postgresql:///postgres?host=/tmp',async c=>(await c.query("SELECT datname,usename,client_addr::text,client_port,pid FROM pg_stat_activity WHERE datname IN ('ken','bertha_betting') AND application_name NOT LIKE 'TK11438%' ORDER BY datname,pid")).rows);
 const proof={label,pid:s.pid,pm_id:e.pm_id,restart_time:e.restart_time,identities,settings,kenSettings,reconciliation,http,process_db_tcp:processDbTcp,saved_match:savedMatch,targets:Object.fromEntries(['DATABASE_URL','KEN_DATABASE_URL'].map(k=>[k,summarize(env[k])])),activity,source_hash:sha(fs.readFileSync(DIR+'/server.js')),wrapper_hash:sha(fs.readFileSync(WRAPPER))};record(label,proof);return proof;
}
(async()=>{
 if(mode==='prepare'){
  assert(!fs.existsSync(PRIV),'Already prepared; do not overwrite baseline');assert.equal(git('status','--porcelain'),'','Ken working tree dirty');
  for(const x of [...manifest.file_changes,...manifest.environment_changes])if(x.sha256)assert.equal(sha(fs.readFileSync(x.path)),x.sha256,'approved source drift: '+x.path);
  const s=unique(await liveRpc('getMonitorData')),saved=unique(json(DUMP)),envBytes=fs.readFileSync(ENV),env=parse(envBytes);const before=fields(s.pm2_env);assert.deepEqual(before,fields(saved),'effective/dump baseline mismatch');for(const k of ['DATABASE_URL','KEN_DATABASE_URL'])assert.equal(before[k],env[k]);
  const desired={DATABASE_URL:socketURL(env.DATABASE_URL),KEN_DATABASE_URL:socketURL(env.KEN_DATABASE_URL)};assert.equal(summarize(desired.DATABASE_URL).database,'bertha_betting');assert.equal(summarize(desired.KEN_DATABASE_URL).database,'ken');
  let nextEnv=envBytes.toString();for(const [k,v]of Object.entries(desired)){assert(!v.includes("'"));const re=new RegExp('^(?:export\\s+)?'+k+'=.*$','gm');assert.equal([...nextEnv.matchAll(re)].length,1);nextEnv=nextEnv.replace(re,k+"='"+v+"'");}assert.deepEqual({...parse(nextEnv),...env},env);for(const k of Object.keys(env))if(!(k in desired))assert.equal(parse(nextEnv)[k],env[k]);
  const w=manifest.file_changes[0],oldWrapper=fs.readFileSync(WRAPPER,'utf8'),newWrapper=oldWrapper.replace(w.old,w.new);assert.notEqual(oldWrapper,newWrapper);assert.equal(newWrapper.replace(w.new,w.old),oldWrapper);
  const oldPlist=fs.readFileSync(PLIST);const plist=JSON.parse(run('/usr/bin/plutil',['-convert','json','-o','-',PLIST]));assert.equal(plist.Label,'com.steve.ken-reconcile-canary');const oldUrl=plist.EnvironmentVariables.KEN_DATABASE_URL;plist.EnvironmentVariables.KEN_DATABASE_URL=socketURL(oldUrl);const newPlist=run('/usr/bin/plutil',['-convert','xml1','-o','-','--','-'],{input:JSON.stringify(plist)});
  fs.mkdirSync(PRIV,{mode:0o700});for(const [f,b]of Object.entries({'env.before':envBytes,'env.after':nextEnv,'wrapper.before':oldWrapper,'wrapper.after':newWrapper,'plist.before':oldPlist,'plist.after':newPlist}))fs.writeFileSync(PRIV+'/'+f,b,{mode:0o600,flag:'wx'});
  const rec={head:git('rev-parse','HEAD'),pid:s.pid,pm_id:s.pm2_env.pm_id,before,desired,source_hash:sha(fs.readFileSync(DIR+'/server.js')),hashes:{env:sha(envBytes),wrapper:sha(oldWrapper),plist:sha(oldPlist)}};fs.writeFileSync(PRIV+'/receipt.json',JSON.stringify(rec,null,2),{mode:0o600,flag:'wx'});
  // Rehearse full file restore and scoped dump round-trip on copies, never live.
  const d=json(DUMP),original=JSON.parse(JSON.stringify(d));setFields(unique(d),{...desired,env:desired});setFields(unique(d),before);assert.deepEqual(d,original);
  for(const name of ['env','wrapper','plist']){const p=PRIV+'/rehearsal-'+name;fs.copyFileSync(PRIV+'/'+name+'.before',p);atomic(p,fs.readFileSync(PRIV+'/'+name+'.after'),rec.hashes[name]);atomic(p,fs.readFileSync(PRIV+'/'+name+'.before'),sha(fs.readFileSync(PRIV+'/'+name+'.after')));assert.equal(sha(fs.readFileSync(p)),rec.hashes[name]);}
  run('/bin/bash',['-n',PRIV+'/wrapper.after']);run('/bin/bash',['-n',PRIV+'/env.after']);run('/usr/bin/plutil',['-lint',PRIV+'/plist.after']);
  const baseline=await snapshot('baseline');assert(Object.values(baseline.reconciliation).every(Boolean),'Existing reconciliation failure: stop before job reload');
  record('preparation',{verdict:'PASS',rollback_rehearsal:'exact file restores plus four-field dump roundtrip PASS',application_mutations:0,email_sends:0,pid:s.pid,pm_id:rec.pm_id});
 }else if(mode==='apply'){
  const rec=json(PRIV+'/receipt.json');assert.equal(json(OUT+'/preparation.json').verdict,'PASS');assert(!fs.existsSync(OUT+'/applied.json'));assert.equal(git('status','--porcelain'),'');assert.equal(git('rev-parse','HEAD'),rec.head);
  const s=unique(await liveRpc('getMonitorData'));assert.equal(s.pid,rec.pid);assert.deepEqual(fields(s.pm2_env),rec.before);assert.equal(sha(fs.readFileSync(DIR+'/server.js')),rec.source_hash);
  for(const [p,k]of [[ENV,'env'],[WRAPPER,'wrapper'],[PLIST,'plist']])assert.equal(sha(fs.readFileSync(p)),rec.hashes[k]);assert.deepEqual(fields(unique(json(DUMP))),rec.before);
  record('mutation-started',{scope:'Ken only',before_pid:s.pid,steps:[]});
  atomic(ENV,fs.readFileSync(PRIV+'/env.after'),rec.hashes.env);atomic(WRAPPER,fs.readFileSync(PRIV+'/wrapper.after'),rec.hashes.wrapper);atomic(PLIST,fs.readFileSync(PRIV+'/plist.after'),rec.hashes.plist);
  run('/bin/bash',['-n',WRAPPER]);run('/usr/bin/plutil',['-lint',PLIST]);
  git('add','--','kalshi-dash/follow-the-winners-check.sh');git('-c','user.name=Steve Abrams','-c','user.email=steve@designerwallcoverings.com','commit','-q','-m','Use Unix socket for Ken hourly signal database');
  const dumpProof=patchDump(rec.before,{...rec.desired,env:rec.desired});
  record('config-applied',{dump:dumpProof,commit:git('rev-parse','HEAD'),application_restart_pending:true});
  await liveRpc('restartProcessId',{id:rec.pm_id,env:rec.desired});
  const after=unique(await liveRpc('getMonitorData'));assert.notEqual(after.pid,rec.pid);record('applied',{verdict:'APPLIED_VERIFY_PENDING',before_pid:rec.pid,after_pid:after.pid,commit:git('rev-parse','HEAD'),email_wrapper_executed:false});
 }else if(mode==='verify'){
  const p=await snapshot('after'),b=json(OUT+'/baseline.json'),rec=json(PRIV+'/receipt.json');assert(p.identities.every(x=>x.addr===null&&x.readonly==='on'));assert.equal(p.process_db_tcp.length,0);assert(p.saved_match);assert.deepEqual(p.identities.map(x=>[x.key,x.database,x.role]),b.identities.map(x=>[x.key,x.database,x.role]));assert.deepEqual(['safe_mode','trading_on','kalshi_env'].map(k=>p.settings[k]),['safe_mode','trading_on','kalshi_env'].map(k=>b.settings[k]));assert.deepEqual(p.kenSettings,b.kenSettings);assert.equal(p.source_hash,b.source_hash);assert(p.activity.some(x=>x.datname==='ken'&&x.client_addr===null));assert(p.activity.some(x=>x.datname==='bertha_betting'&&x.client_addr===null));
  assert.equal(sha(fs.readFileSync(ENV)),sha(fs.readFileSync(PRIV+'/env.after')));assert.equal(sha(fs.readFileSync(PLIST)),sha(fs.readFileSync(PRIV+'/plist.after')));
  const bad=new Client({host:OUT+'/missing-socket',database:'ken',connectionTimeoutMillis:1000});let code;try{await bad.connect();throw Error('unexpected missing socket connection');}catch(err){code=err.code;assert.equal(code,'ENOENT');}finally{await bad.end().catch(()=>{});}
  // Run underlying pure SELECT module only. The wrapper that emails is excluded.
  const wrapperUrl=fs.readFileSync(WRAPPER,'utf8').match(/KEN_DATABASE_URL="([^"]+)"/)?.[1];assert(wrapperUrl);assert.equal(new URL(wrapperUrl).searchParams.get('host'),'/tmp');assert.equal(new URL(wrapperUrl).pathname,'/ken');
  const result=run('/opt/homebrew/bin/node',[DIR+'/follow-the-winners.mjs'],{cwd:DIR,timeout:20000,env:{PATH:process.env.PATH,HOME:process.env.HOME,USER:os.userInfo().username,KEN_DATABASE_URL:wrapperUrl,PGOPTIONS:'-c default_transaction_read_only=on -c statement_timeout=5000'}});
  record('verification',{verdict:'PASS',pid:p.pid,two_pools_socket:true,auth_statuses:p.http.map(x=>x.status),captured_trading_controls_unchanged:true,whole_config_warning:{before:b.settings.config_hash,after:p.settings.config_hash,detail:'AutonomousScan updates weather_cache. Captured switches and ken_config hashes match; uncaptured risk_state keys cannot be proven unchanged retroactively.'},durable_match:true,missing_socket:code,signal_module:{exit:0,configuration_source:'exact KEN_DATABASE_URL from updated hourly wrapper',target:summarize(wrapperUrl),output_sha256:sha(result),bytes:result.length,email_wrapper_executed:false},reconciliation:p.reconciliation,source_unchanged:true});
 }else if(mode==='reload'){
  assert.equal(json(OUT+'/verification.json').verdict,'PASS');assert(!fs.existsSync(OUT+'/reloaded.json'));const target='gui/'+process.getuid()+'/com.steve.ken-reconcile-canary';const before=run('launchctl',['print',target]);assert(!/\n\s+pid = \d+/.test(before),'Canary running: wait until idle');
  record('reload-started',{target,prior_loaded:true,phase:'before-bootout'});
  run('launchctl',['bootout','gui/'+process.getuid(),PLIST]);record('reload-started',{target,prior_loaded:true,phase:'bootout-complete'});
  run('launchctl',['bootstrap','gui/'+process.getuid(),PLIST]);
  const after=run('launchctl',['print',target]);const loadedUrl=after.match(/KEN_DATABASE_URL => (\S+)/)?.[1];assert(loadedUrl);assert.equal(new URL(loadedUrl).searchParams.get('host'),'/tmp');record('reloaded',{verdict:'RELOADED_OBSERVE_PENDING',target,loaded_socket:true,loaded:true});
 }else if(mode==='observe'){
  const proof=await snapshot('monitor'),after=json(OUT+'/after.json');assert.equal(proof.pid,after.pid);assert.equal(proof.restart_time,after.restart_time);assert.equal(proof.process_db_tcp.length,0);assert(proof.saved_match);assert(proof.identities.every(x=>x.addr===null));
  const target='gui/'+process.getuid()+'/com.steve.ken-reconcile-canary',loaded=run('launchctl',['print',target]);
  const loadedUrl=loaded.match(/KEN_DATABASE_URL => (\S+)/)?.[1];assert.equal(new URL(loadedUrl).searchParams.get('host'),'/tmp');assert(/last exit code = 0\b/.test(loaded));assert(!/\n\s+pid = \d+/.test(loaded),'canary still running');assert(/run interval = 1800 seconds/.test(loaded));
  const baseline=json(OUT+'/canary-log-baseline.json'),bytes=fs.readFileSync(baseline.path);assert(bytes.length>=baseline.size,'log rotated; inspect new evidence');const tail=bytes.subarray(baseline.size).toString();assert(tail.includes('[reconcile-canary] OK'));assert(!/DIVERGENCE|\[reconcile-canary\] error:/.test(tail));
  const before=JSON.parse(run('/usr/bin/plutil',['-convert','json','-o','-',PRIV+'/plist.before'])),current=JSON.parse(run('/usr/bin/plutil',['-convert','json','-o','-',PLIST]));current.EnvironmentVariables.KEN_DATABASE_URL=before.EnvironmentVariables.KEN_DATABASE_URL;assert.deepEqual(current,before);
  record('observation',{verdict:'PASS',pid:proof.pid,restart_count_stable:true,two_pools_socket:true,auth_statuses:proof.http.map(x=>x.status),reconciliation_job:{loaded_socket:true,last_exit:0,interval:1800,other_plist_fields_preserved:true,new_log_bytes:tail.length,log_sha256:sha(tail),result:'OK'},email_wrapper_invoked:false});
 }else if(mode==='rollback'){
  const rec=json(PRIV+'/receipt.json');assert(fs.existsSync(OUT+'/mutation-started.json'));assert.equal(sha(fs.readFileSync(DIR+'/server.js')),rec.source_hash);
  for(const [p,k]of [[ENV,'env'],[WRAPPER,'wrapper'],[PLIST,'plist']]){const current=sha(fs.readFileSync(p)),after=sha(fs.readFileSync(PRIV+'/'+k+'.after'));assert([rec.hashes[k],after].includes(current),'peer edit blocks rollback');if(current===after)atomic(p,fs.readFileSync(PRIV+'/'+k+'.before'),after);}
  const f=fields(unique(json(DUMP)));if(JSON.stringify(f)!==JSON.stringify(rec.before))patchDump({...rec.desired,env:rec.desired},rec.before);
  await liveRpc('restartProcessId',{id:rec.pm_id,env:{DATABASE_URL:rec.before.DATABASE_URL,KEN_DATABASE_URL:rec.before.KEN_DATABASE_URL}});
  if(fs.existsSync(OUT+'/reload-started.json')){
   const target='gui/'+process.getuid()+'/com.steve.ken-reconcile-canary';
   const loaded=cp.spawnSync('launchctl',['print',target],{encoding:'utf8',timeout:10000});
   if(loaded.status===0)run('launchctl',['bootout','gui/'+process.getuid(),PLIST]);
   else assert(/could not find service|service not found/i.test(loaded.stderr||''),'unexpected launchd state; inspect before rollback');
   run('launchctl',['bootstrap','gui/'+process.getuid(),PLIST]);
   const restored=run('launchctl',['print',target]),loadedUrl=restored.match(/KEN_DATABASE_URL => (\S+)/)?.[1];
   const originalPlist=JSON.parse(run('/usr/bin/plutil',['-convert','json','-o','-',PRIV+'/plist.before']));
   assert.equal(loadedUrl,originalPlist.EnvironmentVariables.KEN_DATABASE_URL);
   record('reload-rollback',{verdict:'PASS',prior_job_restored:true,target});
  }
  record('rollback',{verdict:'RESTORED_VERIFY_BASELINE_REQUIRED'});
 }
 console.log(JSON.stringify({mode,verdict:'PASS',evidence:OUT}));
})().catch(err=>{record(mode+'-failure',{verdict:'FAIL',code:err.code||err.name,message:String(err.message).replace(/postgres(?:ql)?:\/\/\S+/g,'[REDACTED_URL]')});console.error(JSON.stringify({mode,verdict:'FAIL',code:err.code||err.name,message:String(err.message).replace(/postgres(?:ql)?:\/\/\S+/g,'[REDACTED_URL]')}));process.exitCode=1;});