Verify original database commitments before releasing restored IndeeHub
This commit is contained in:
@@ -11,6 +11,53 @@ QUEUE_SCRIPT = r'''const {Queue}=require('bullmq');
|
||||
try{const action=process.argv[1];if(action==='pause')await q.pause();else if(action==='resume')await q.resume();else if(action!=='status')throw Error('action');
|
||||
console.log(JSON.stringify({paused:await q.isPaused(),counts:await q.getJobCounts('active','waiting','paused','delayed','failed','completed')}));}
|
||||
finally{await q.close()}})().catch(()=>process.exit(1));'''
|
||||
# The exact three migrations in the privately qualified API candidate. This is
|
||||
# an allowlist of additive schema history, never permission to discard app data.
|
||||
ADDITIVE_MIGRATIONS = {
|
||||
'AddArchipelagoPublicationsAndRentals1791288000000':1791288000000,
|
||||
'AddMediaRegistrationIntents1791374400000':1791374400000,
|
||||
'AddMediaRegistrationRetirements1791374401000':1791374401000,
|
||||
}
|
||||
ADDITIVE_TABLES = {'archipelago_media_registrations','archipelago_publications','archipelago_publication_outbox','archipelago_rental_entitlements','archipelago_registration_intents'}
|
||||
DB_COMMITMENTS_SQL = r'''
|
||||
BEGIN TRANSACTION ISOLATION LEVEL REPEATABLE READ READ ONLY;
|
||||
SELECT format($query$
|
||||
SELECT jsonb_build_object('table',%L,'schema',
|
||||
jsonb_build_object(
|
||||
'columns',(SELECT coalesce(jsonb_agg(jsonb_build_array(a.attnum,a.attname,format_type(a.atttypid,a.atttypmod),a.attnotnull,a.attidentity,a.attgenerated,pg_get_expr(d.adbin,d.adrelid)) ORDER BY a.attnum),'[]'::jsonb) FROM pg_attribute a LEFT JOIN pg_attrdef d ON d.adrelid=a.attrelid AND d.adnum=a.attnum WHERE a.attrelid=%s AND a.attnum>0 AND NOT a.attisdropped),
|
||||
'constraints',(SELECT coalesce(jsonb_agg(jsonb_build_array(conname,pg_get_constraintdef(oid,true)) ORDER BY conname),'[]'::jsonb) FROM pg_constraint WHERE conrelid=%s),
|
||||
'indexes',(SELECT coalesce(jsonb_agg(pg_get_indexdef(indexrelid) ORDER BY indexrelid::regclass::text),'[]'::jsonb) FROM pg_index WHERE indrelid=%s),
|
||||
'triggers',(SELECT coalesce(jsonb_agg(pg_get_triggerdef(oid,true) ORDER BY tgname),'[]'::jsonb) FROM pg_trigger WHERE tgrelid=%s AND NOT tgisinternal),
|
||||
'rls',%L,'policies',(SELECT coalesce(jsonb_agg(to_jsonb(p) ORDER BY policyname),'[]'::jsonb) FROM pg_policies p WHERE schemaname='public' AND tablename=%L)),
|
||||
'rows',count(*),'rows_sha256',encode(sha256(convert_to(coalesce(string_agg(to_jsonb(t)::text,E'\n' ORDER BY to_jsonb(t)::text),''),'UTF8')),'hex')) FROM public.%I t;
|
||||
$query$,c.relname,c.oid,c.oid,c.oid,c.oid,c.relrowsecurity::text||':'||c.relforcerowsecurity::text,c.relname,c.relname)
|
||||
FROM pg_class c JOIN pg_namespace n ON n.oid=c.relnamespace WHERE n.nspname='public' AND c.relkind IN ('r','p') ORDER BY c.relname
|
||||
\gexec
|
||||
SELECT jsonb_build_object('migration_rows',coalesce(jsonb_agg(to_jsonb(m) ORDER BY id),'[]'::jsonb)) FROM public.migrations m;
|
||||
COMMIT;
|
||||
'''
|
||||
def verify_database_compatibility(before, after):
|
||||
require(before.get('operation_id')==after.get('operation_id'),'Data compatibility operation changed')
|
||||
old=before['tables'];new=after['tables'];require(set(old)<=set(new),'Data compatibility lost original tables')
|
||||
extra=set(new)-set(old);require(extra<=ADDITIVE_TABLES,'Data compatibility contains unreviewed tables')
|
||||
for name in old:
|
||||
require(old[name]['schema']==new[name]['schema'],'Data compatibility changed original table schema')
|
||||
if name!='migrations':
|
||||
require(old[name]['rows']==new[name]['rows'] and old[name]['rows_sha256']==new[name]['rows_sha256'],'Data compatibility changed original rows')
|
||||
for name in extra:require(new[name]['rows']==0,'Data compatibility contains new application data')
|
||||
previous=before['migrations'];current=after['migrations']
|
||||
require(current[:len(previous)]==previous,'Data compatibility changed original migration history')
|
||||
added=current[len(previous):];seen=set()
|
||||
for row in added:
|
||||
require(set(row)=={'id','timestamp','name'} and type(row['id']) is int and row['name'] not in seen and ADDITIVE_MIGRATIONS.get(row['name'])==row['timestamp'],'Data compatibility has unreviewed migration history')
|
||||
seen.add(row['name'])
|
||||
ordered=list(ADDITIVE_MIGRATIONS)
|
||||
require([row['name'] for row in added]==ordered[:len(added)],'Data compatibility migration order changed')
|
||||
expected_extra=(ADDITIVE_TABLES-{'archipelago_registration_intents'}) if added else set()
|
||||
if len(added)>=2:expected_extra=ADDITIVE_TABLES
|
||||
require(extra==expected_extra-set(old),'Data compatibility additions do not match reviewed migration evidence')
|
||||
return {'original_tables':len(old),'new_empty_tables':sorted(extra),'reviewed_migrations':[r['name'] for r in added]}
|
||||
|
||||
def require(condition, message):
|
||||
if not condition: raise RuntimeError(message)
|
||||
def atomic(path, value):
|
||||
@@ -60,11 +107,11 @@ class Controller:
|
||||
self.record=json.loads(self.path.read_text()) if self.path.exists() else None
|
||||
if self.record:require(self.record['operation_id']==operation,'Maintenance journal changed')
|
||||
def save(self): atomic(self.path,self.record)
|
||||
def run(self, argv, timeout=30, output=None):
|
||||
def run(self, argv, timeout=30, output=None, input_bytes=None):
|
||||
if self.runner:return self.runner(argv,timeout,output)
|
||||
self.root.mkdir(mode=0o700,parents=True,exist_ok=True)
|
||||
with (self.root/'commands.private.log').open('ab') as errors:
|
||||
result=subprocess.run(argv,stdout=output or subprocess.PIPE,stderr=errors,timeout=timeout,check=True,pass_fds=(self.lock_fd,))
|
||||
result=subprocess.run(argv,stdout=output or subprocess.PIPE,stderr=errors,timeout=timeout,check=True,pass_fds=(self.lock_fd,),input=input_bytes)
|
||||
if output:return b''
|
||||
require(len(result.stdout)<=2*1024*1024,'Command response exceeds bound')
|
||||
return result.stdout
|
||||
@@ -149,8 +196,27 @@ class Controller:
|
||||
rows=json.loads(self.run(['podman','volume','inspect',*expected]))
|
||||
require({row['Name'] for row in rows}==set(expected),'Persistent volume scope changed')
|
||||
return {row['Name']:row['Mountpoint'] for row in rows}
|
||||
def database_commitments(self):
|
||||
raw=self.run(['podman','exec','-i','indeedhub-postgres','psql','-XqAt','--set=ON_ERROR_STOP=1','-U','indeedhub','-d','indeedhub'],timeout=300,input_bytes=DB_COMMITMENTS_SQL.encode())
|
||||
rows=[json.loads(line) for line in raw.decode().splitlines() if line.strip()]
|
||||
tables={};migrations=None
|
||||
for row in rows:
|
||||
if 'migration_rows' in row:
|
||||
require(migrations is None,'Duplicate database migration observation');migrations=row['migration_rows']
|
||||
else:
|
||||
name=row.pop('table');require(name not in tables and re.fullmatch('[a-zA-Z_][a-zA-Z0-9_]*',name),'Invalid database table observation');tables[name]=row
|
||||
require(tables and 'migrations' in tables and isinstance(migrations,list),'Database compatibility observation incomplete')
|
||||
return {'operation_id':self.operation,'tables':tables,'migrations':migrations}
|
||||
def verify_restored_data(self):
|
||||
baseline=self.record.get('database_before')
|
||||
require(baseline and baseline.get('operation_id')==self.operation,'Data compatibility baseline missing')
|
||||
current=self.database_commitments();proof=verify_database_compatibility(baseline,current)
|
||||
self.record['recovery_data_verification']={'operation_id':self.operation,'checked_at':time.time(),'before_sha256':hashlib.sha256(json.dumps(baseline,sort_keys=True).encode()).hexdigest(),'after_sha256':hashlib.sha256(json.dumps(current,sort_keys=True).encode()).hexdigest(),**proof}
|
||||
self.record['recovery_data_verified']=True;self.save()
|
||||
def backup(self):
|
||||
if self.record.get('backup_complete'):return
|
||||
if 'database_before' not in self.record:
|
||||
self.record['database_before']=self.database_commitments();self.save()
|
||||
sources=self.volume_sources();self.record['volume_sources']=sources;self.save()
|
||||
backup=self.root/'backup';backup.mkdir(mode=0o700,exist_ok=True)
|
||||
if 'database.dump' not in self.record.setdefault('artifacts',{}):
|
||||
@@ -241,7 +307,7 @@ class Controller:
|
||||
if outcome!='committed':
|
||||
require(type(runtime.get('target_startup_began')) is bool,'Target-start obligation unavailable')
|
||||
if runtime['target_startup_began']:
|
||||
require(self.record.get('recovery_data_verified') is True,'Data compatibility after rollback needs recorded verification; ingress remains closed')
|
||||
self.verify_restored_data()
|
||||
else:
|
||||
self.record['rollback_data_claim']='No target startup/migration began; only original runtime restored.'
|
||||
|
||||
|
||||
Reference in New Issue
Block a user