← back to Tk 10965 Zero Price Analysis

vienna-executor.mjs

167 lines

import fs from 'node:fs';
import path from 'node:path';
import {createHash,randomUUID} from 'node:crypto';
export const DOMAIN='designer-laboratory-sandbox.myshopify.com'; // PRODUCTION despite its name.
export const hash = bytes=>createHash('sha256').update(bytes).digest('hex');
const fail = m=>{throw new Error(m);};
const eq = (a,b)=>JSON.stringify(a)===JSON.stringify(b);
function keys(x, fields, what) {
  if (!x || typeof x!=='object' || Array.isArray(x) || !eq(Object.keys(x).sort(),fields.split(' ').sort())) fail('Malformed '+what);
}
export const key = r=>`${r.inventoryItemId}@${r.locationId}`;
const gid=(s,type)=>typeof s==='string' && new RegExp('^gid://shopify/'+type+'/[1-9][0-9]*$').test(s);
export function readRegular(file) {
  const fd=fs.openSync(file, fs.constants.O_RDONLY|fs.constants.O_NOFOLLOW);
  try {const stat=fs.fstatSync(fd); if(!stat.isFile() || stat.nlink!==1 || stat.size>10_000_000) fail('Unsafe file'); return fs.readFileSync(fd,'utf8');}
  finally {fs.closeSync(fd);}
}
export function validateManifest(bytes, expectedHash, shopId, now=Date.now()) {
  if (hash(bytes)!==expectedHash) fail('Manifest SHA256 mismatch');
  const m=JSON.parse(bytes);
  keys(m,'version scope store createdAt expiresAt records canary','manifest');
  keys(m.store,'domain id','store'); keys(m.scope,'vendor line','scope');
  if(m.version!==1 || m.scope.vendor!=='Phillipe Romano' || m.scope.line!=='Vienna') fail('Foreign manifest scope');
  if(m.store.domain!==DOMAIN || m.store.id!==shopId || !gid(shopId,'Shop')) fail('Store pin mismatch');
  const born=Date.parse(m.createdAt), expires=Date.parse(m.expiresAt);
  if(!Number.isFinite(born) || !Number.isFinite(expires) || new Date(born).toISOString()!==m.createdAt || new Date(expires).toISOString()!==m.expiresAt || born>now || expires<=now || expires<=born || expires-born>86400000 || now-born>86400000) fail('Stale or invalid manifest validity');
  if(!Array.isArray(m.records)||!m.records.length||m.records.length>50000||!Array.isArray(m.canary)||!m.canary.length) fail('Empty or excessive manifest scope');
  const seen=new Set(), variants=new Map(), items=new Map();
  for(const r of m.records) {
    keys(r,'productId variantId inventoryItemId locationId vendor line title variantTitle price inventoryPolicy tracked onHand target','record');
    for(const [f,t] of [['productId','Product'],['variantId','ProductVariant'],['inventoryItemId','InventoryItem'],['locationId','Location']]) if(!gid(r[f],t)) fail('Invalid exact identity '+f);
    if(r.vendor!==m.scope.vendor || r.line!=='Vienna' || typeof r.title!=='string' || !/^Vienna(?:\s|$)/.test(r.title) || typeof r.variantTitle!=='string' || /sample/i.test(r.variantTitle) || r.price!==0 || r.inventoryPolicy!=='DENY' || r.tracked!==true || !Number.isSafeInteger(r.onHand) || r.onHand<=0 || r.target!==0) fail('Foreign or unsafe record scope');
    if(seen.has(key(r))) fail('Duplicate inventory location'); seen.add(key(r));
    const identity=JSON.stringify([r.productId,r.inventoryItemId]);
    if(variants.has(r.variantId)&&variants.get(r.variantId)!==identity) fail('Variant identity collision'); variants.set(r.variantId,identity);
    if(items.has(r.inventoryItemId)&&items.get(r.inventoryItemId)!==r.variantId) fail('Item identity collision'); items.set(r.inventoryItemId,r.variantId);
  }
  if(new Set(m.canary).size!==m.canary.length || m.canary.some(k=>!seen.has(k))) fail('Canary outside frozen scope');
  const selected=new Set(m.canary), products=new Set(m.records.filter(r=>selected.has(key(r))).map(r=>r.productId));
  if(m.records.some(r=>products.has(r.productId)&&!selected.has(key(r)))) fail('Canary must include every frozen location of each selected product');
  return m;
}
function fsyncDirectory(file) { const fd=fs.openSync(path.dirname(file),'r');try{fs.fsyncSync(fd);}finally{fs.closeSync(fd);} }
export function durableWrite(file, data, flag='w') {
  // Atomic fixture replacement preserves the last durable truth on interruption.
  // Refuse aliases before writing, and retain unfinished temp files for inspection.
  let before;
  try { before=fs.lstatSync(file); } catch(e) { if(e.code!=='ENOENT')throw e; }
  if(before && (!before.isFile() || before.nlink!==1)) fail('Unsafe write target');
  const target=flag==='wx'?file:file+'.pending-'+randomUUID();
  const fd=fs.openSync(target,fs.constants.O_WRONLY|fs.constants.O_CREAT|fs.constants.O_EXCL|fs.constants.O_NOFOLLOW,0o600);
  try {if(!fs.fstatSync(fd).isFile()||fs.fstatSync(fd).nlink!==1)fail('Unsafe write descriptor');fs.writeFileSync(fd,data);fs.fsyncSync(fd);} finally {fs.closeSync(fd);}
  if(flag!=='wx') {
    let current; try{current=fs.lstatSync(file);}catch(e){if(e.code!=='ENOENT')throw e;}
    if(Boolean(current)!==Boolean(before) || (current && (!current.isFile()||current.nlink!==1||current.ino!==before.ino||current.dev!==before.dev))) fail('Write target changed');
    fs.renameSync(target,file);
  }
  fsyncDirectory(file);
}
function eventLine(body, prev) {const event={...body,prev};return JSON.stringify({...event,hash:hash(JSON.stringify(event))})+'\n';}
export function readJournal(bytes, manifestHash, m) {
  if(!bytes.endsWith('\n')) fail('Truncated journal');
  const lines=bytes.trimEnd().split('\n'); let prev='0'.repeat(64), header, pending=null;
  const states=new Map(), records=new Map(m.records.map(r=>[key(r),r]));
  for(let i=0;i<lines.length;i++) {
    const e=JSON.parse(lines[i]); const {hash:digest,...body}=e;
    if(digest!==hash(JSON.stringify(body)) || body.prev!==prev || body.seq!==i) fail('Corrupt journal chain'); prev=digest;
    if(i===0) {
      keys(e,'seq type manifestHash store selection prev hash','journal header');
      if(e.type!=='header'||e.manifestHash!==manifestHash||!eq(e.store,m.store)||!['--canary','--all'].includes(e.selection)) fail('Foreign journal');
      header=e; continue;
    }
    keys(e,'seq type key direction before after prev hash','journal event');
    const r=records.get(e.key);
    if(!r || (header.selection==='--canary'&&!m.canary.includes(e.key)) || !['apply','rollback'].includes(e.direction)) fail('Journal scope injection');
    if(e.before!==(e.direction==='apply'?r.onHand:0)||e.after!==(e.direction==='apply'?0:r.onHand)) fail('Journal quantity injection');
    const s=states.get(e.key)||'unattempted';
    if(e.type==='intent') {
      if(pending || (e.direction==='apply'&&!['unattempted','rejected'].includes(s)) || (e.direction==='rollback'&&s!=='applied')) fail('Invalid journal transition');
      pending=e;
    } else if(['success','rejected','verified'].includes(e.type)) {
      if(e.type==='verified') {
        if(pending || s!==(e.direction==='apply'?'applied-unverified':'rolledback-unverified')) fail('Unbound verification');
        states.set(e.key,e.direction==='apply'?'applied':'rolledback');
      } else {
        if(!pending || !eq([pending.key,pending.direction,pending.before,pending.after],[e.key,e.direction,e.before,e.after])) fail('Unbound journal outcome');
        states.set(e.key,e.type==='success'?(e.direction==='apply'?'applied-unverified':'rolledback-unverified'):(e.direction==='apply'?'rejected':'applied'));
        pending=null;
      }
    } else fail('Unknown journal event');
  }
  if(pending) fail('Unresolved mutation intent: reconciliation required; never infer failure');
  return {header,states,seq:lines.length,prev};
}
function assertSnapshot(actual,r,quantity) {
  const expected={...r,onHand:quantity};
  if(!eq(actual,expected)) fail('Identity, quantity or product precondition drift at '+key(r));
}
export async function run(cli) {
  const manifestPath=path.resolve(cli['--manifest']);
  const m=validateManifest(readRegular(manifestPath),cli['--expected-sha256'],cli['--expected-store-id']);
  if(['--plan','--enumerate'].includes(cli.mode)) return {mode:'plan',manifestHash:cli['--expected-sha256'],store:m.store,records:m.records,canary:m.canary,mutations:0};
  const journalPath=path.resolve(cli['--journal']),statePath=path.resolve(cli['--offline-state']);
  if(new Set([manifestPath,journalPath,statePath]).size!==3) fail('Input/output paths must be distinct');
  const exists=fs.existsSync(journalPath);
  let journal;
  if(exists) {
    if(!cli['--journal-sha256']) fail('Existing journal requires external SHA256 checkpoint');
    const bytes=readRegular(journalPath); if(hash(bytes)!==cli['--journal-sha256']) fail('Journal SHA256 mismatch');
    journal=readJournal(bytes,cli['--expected-sha256'],m);
    if(cli.mode!=='--rollback'&&journal.header.selection!==cli.mode) fail('Cannot broaden frozen journal selection');
  } else if(cli.mode==='--rollback'||cli['--journal-sha256']) fail('Missing pinned journal');
  // The only adapter shipped by this local increment is hermetic and explicitly selected.
  // Future live wiring needs separate approval and a reviewed identity/auth contract.
  const {openOfflineAdapter}=await import('./vienna-offline-adapter.mjs');
  const adapter=openOfflineAdapter(statePath,m.store);
  if(!eq(await adapter.shop(),m.store)) fail('Adapter store mismatch');
  const selection=journal?.header.selection||cli.mode;
  const records=m.records.filter(r=>selection==='--all'||m.canary.includes(key(r)));
  if(cli.mode!=='--rollback' && [...(journal?.states.values()||[])].some(s=>s.startsWith('rolledback'))) fail('Rolled back run cannot reapply');
  // Exclusive process lock. A crashed holder leaves the lock for independent reconciliation.
  const stateLock=statePath+'.lock', lock=journalPath+'.lock';
  fs.mkdirSync(stateLock,{mode:0o700});
  try {
    fs.mkdirSync(lock,{mode:0o700});
    try {
    // Recheck checkpoint under lock: concurrent changes cannot pass a stale validation.
    if(exists && hash(readRegular(journalPath))!==cli['--journal-sha256']) fail('Journal changed before lock');
    if(!exists) {
      const line=eventLine({seq:0,type:'header',manifestHash:cli['--expected-sha256'],store:m.store,selection},'0'.repeat(64));
      durableWrite(journalPath,line,'wx'); journal=readJournal(line,cli['--expected-sha256'],m);
    }
    const append=(type,r,direction)=>{
      const line=eventLine({seq:journal.seq,type,key:key(r),direction,before:direction==='apply'?r.onHand:0,after:direction==='apply'?0:r.onHand},journal.prev);
      const fd=fs.openSync(journalPath,fs.constants.O_WRONLY|fs.constants.O_APPEND|fs.constants.O_NOFOLLOW);
      try{const stat=fs.fstatSync(fd);if(!stat.isFile()||stat.nlink!==1)fail('Unsafe journal descriptor');fs.writeFileSync(fd,line);fs.fsyncSync(fd);}finally{fs.closeSync(fd);}
      journal.seq++; journal.prev=JSON.parse(line).hash;
    };
    let writes=0,skipped=0;
    for(const r of records) {
      let s=journal.states.get(key(r))||'unattempted';
      // A durable success whose post-read was interrupted is never written again.
      if(s.endsWith('-unverified')) {
        const direction=s==='applied-unverified'?'apply':'rollback';
        assertSnapshot(await adapter.read(r),r,direction==='apply'?0:r.onHand);
        append('verified',r,direction); s=direction==='apply'?'applied':'rolledback';
      }
      if(cli.mode==='--rollback' && s!=='applied') {skipped++;continue;}
      if(cli.mode!=='--rollback' && s==='applied') {assertSnapshot(await adapter.read(r),r,0);skipped++;continue;}
      const direction=cli.mode==='--rollback'?'rollback':'apply';
      const before=direction==='apply'?r.onHand:0, after=direction==='apply'?0:r.onHand;
      assertSnapshot(await adapter.read(r),r,before);
      append('intent',r,direction); // fsync before crossing the mutation boundary.
      const result=await adapter.set(r,before,after); // exceptions leave unresolved durable intent.
      if(result?.kind==='rejected' && result.confirmedNoWrite===true) {
        append('rejected',r,direction); fail('Confirmed mutation rejection at '+key(r));
      }
      if(result?.kind!=='success'||result.key!==key(r)||result.quantity!==after) fail('Ambiguous mutation result; reconciliation required');
      append('success',r,direction); // only confirmed successes authorize rollback.
      assertSnapshot(await adapter.read(r),r,after);
      append('verified',r,direction); writes++;
    }
    return {mode:cli.mode,store:m.store,writes,skipped,journal:journalPath,journalSha256:hash(readRegular(journalPath)),manifestHash:cli['--expected-sha256'],adapter:'offline-only'};
    } finally {fs.rmdirSync(lock);}
  } finally {fs.rmdirSync(stateLock);}
}