From a30c12193d93e2e06fa74be135ef3bd31786f8f1 Mon Sep 17 00:00:00 2001 From: archipelago Date: Wed, 7 Oct 2026 21:18:36 -0400 Subject: [PATCH] fix(indeehub): prove idle legacy worker termination before backup --- .../managed-update-recovery-implementation.md | 32 ++++++++ scripts/indeehub-maintenance-controller.py | 55 +++++++++++-- .../test_indeehub_maintenance_controller.py | 82 ++++++++++++++++++- 3 files changed, 158 insertions(+), 11 deletions(-) diff --git a/docs/managed-update-recovery-implementation.md b/docs/managed-update-recovery-implementation.md index 6ba1e263..b9d91366 100644 --- a/docs/managed-update-recovery-implementation.md +++ b/docs/managed-update-recovery-implementation.md @@ -293,3 +293,35 @@ or widening manager write access. The controller uses that fixed command now; 25 pure controller tests pass, including refusal before fence creation on dump failure. A new helper hash requires a matching backend rebuild. Candidate-helper qualification remains separate from actual backend transaction acceptance. + +### Legacy worker termination qualification — 2026-10-07 + +The held disposable-VM operation `4320fe90-8ab5-4496-a4e7-cf11ca0fd376` +exposed the original worker's Node-as-PID1 SIGTERM behavior. A no-network, +no-volume probe from its exact recovery image exited 137 without init and 143 +with init. This is not graceful shutdown or proof of completed jobs. + +The controller now requires an exact six-field nonnegative integer queue +observation; absent `active` can no longer imply idle. Its narrowly scoped legacy +worker path records **forced idle termination**, with `graceful=false` and +`completed_work_claim=false`, only after paused all-zero observations before and +after, closed ingress/frontend, exact operation/original/recovery image and known +command, unchanged saved unit, original stop intent and died event, and proof the +process is dead. Other writers' 137 exits remain refused. All 31 pure controller +regressions passed, including missing/nonzero counts, reopened queue, changed +identity, command, unit and process state. Actual worker-only classification passed +under the hardened service → user scope with the lifecycle lock held. + +The same fixture's API has npm as PID1 and one direct `node dist/main` child. +After fresh empty-business-state proof and durable stop intent, an exact-command, +parent-validated SIGTERM to that child stopped the original container; npm emitted +exit 1. The existing clean-exit gate correctly retained the hold. API termination +classification is still under review; no generic exit-1 allowance was added. +Frontend/worker are confirmed stopped; API is stopped but unconfirmed; storage +members remain running. Backup/fresh-restore, full RPC rollback/success and live +activation are **not passed**. Installed pinned helper remains unchanged; this is +private candidate-helper qualification only. Live Yaya remains unchanged. + +Future worker image source now includes idempotent SIGTERM/SIGINT shutdown in app +commit `29627fc` with four passing Jest tests. Its image has not been built; the +previous frontend/API-only candidate catalog cannot cover that new worker image. diff --git a/scripts/indeehub-maintenance-controller.py b/scripts/indeehub-maintenance-controller.py index a2774abf..ffee9222 100644 --- a/scripts/indeehub-maintenance-controller.py +++ b/scripts/indeehub-maintenance-controller.py @@ -7,6 +7,9 @@ import datetime, hashlib, json, os, pathlib, re, shutil, subprocess, sys, time, NAMES = ('indeedhub','indeedhub-api','indeedhub-ffmpeg','indeedhub-minio','indeedhub-postgres','indeedhub-redis','indeedhub-relay') VOLUMES = ('indeedhub-minio-data','indeedhub-postgres-data','indeedhub-redis-data','indeedhub-relay-data') DATA = pathlib.Path('/var/lib/archipelago') +QUEUE_COUNTS = frozenset(('active','waiting','paused','delayed','failed','completed')) +def valid_queue_counts(counts): + return isinstance(counts,dict) and set(counts)==QUEUE_COUNTS and all(type(v) is int and v>=0 for v in counts.values()) QUEUE_SCRIPT = r'''const {Queue}=require('bullmq'); (async()=>{const q=new Queue('transcode',{connection:{host:process.env.QUEUE_HOST,port:Number(process.env.QUEUE_PORT||6379),password:process.env.QUEUE_PASSWORD,maxRetriesPerRequest:1}}); 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'); @@ -193,8 +196,7 @@ class Controller: self.record['ingress_closed']=True;self.save() def queue(self, action): result=json.loads(self.run(['podman','exec','indeedhub-api','node','-e',QUEUE_SCRIPT,action])) - require(type(result.get('paused')) is bool and isinstance(result.get('counts'),dict),'Invalid queue observation') - for value in result['counts'].values():require(type(value) is int and value>=0,'Invalid job count') + require(type(result.get('paused')) is bool and valid_queue_counts(result.get('counts')),'Incomplete or invalid queue observation') return result def pause_queue(self): if 'queue_was_paused' not in self.record: @@ -207,7 +209,7 @@ class Controller: # completed: this path requires a fresh empty store behind closed ingress. require(self.record.get('stopped',{}).get('indeedhub',{}).get('confirmed'),'Frontend ingress must already be stopped') require(self.record.get('stopped',{}).get('indeedhub-ffmpeg',{}).get('confirmed'),'Transcode worker must already be stopped') - require(self.record.get('queue_pause_confirmed') is True and self.record.get('last_queue_counts',{}).get('active')==0,'Worker queue is not proven idle') + require(self.record.get('queue_pause_confirmed') is True and valid_queue_counts(self.record.get('last_queue_counts')) and self.record['last_queue_counts']['active']==0,'Worker queue is not proven idle') tables=('projects','contents','payments','shareholders','subscriptions','library_items') fields=','.join("'%s',(SELECT count(*) FROM public.%s)"%(name,name) for name in tables) sql="SELECT json_build_object("+fields+",'other_active_transactions',(SELECT count(*) FROM pg_stat_activity WHERE datname=current_database() AND pid<>pg_backend_pid() AND state<>'idle'))" @@ -215,6 +217,41 @@ class Controller: require(set(counts)==set(tables)|{'other_active_transactions'},'Legacy API business-state observation incomplete') require(all(type(value) is int and value==0 for value in counts.values()),'Legacy API has business work or active transactions; completion cannot be inferred') self.record['legacy_api_empty_state']=counts;self.save() + def legacy_idle_worker_termination(self, member, properties): + # Compatibility only for the observed original Node-as-PID1 worker. + # A timeout/SIGKILL is never renamed graceful or completed work. + require(member['name']=='indeedhub-ffmpeg','Forced termination is not allowed for this writer') + self.holds();self.fence_matches() + before=self.record.get('last_queue_counts') + require(self.record.get('ingress_closed') is True and self.record.get('stopped',{}).get('indeedhub',{}).get('confirmed') is True,'Legacy worker ingress is not closed') + require(self.record.get('queue_pause_confirmed') is True and valid_queue_counts(before) and all(v==0 for v in before.values()),'Legacy worker queue was not proven completely empty') + stopped=self.record['stopped'][member['name']] + require(stopped.get('container_id')==member['container_id'] and type(stopped.get('intent_at')) in (int,float),'Legacy worker stop intent changed') + state=dict(line.split('=',1) for line in properties.splitlines() if '=' in line) + require(state.get('ActiveState') in ('inactive','failed') and state.get('SubState') in ('dead','failed') and state.get('ExecMainStatus')=='137','Legacy worker process death is not established') + runtime=json.loads((self.data/'update-transactions'/'supervised'/(self.operation+'.json')).read_text()) + require(runtime.get('id')==self.operation and runtime.get('phase') in ('Editing','Restoring') and runtime.get('target_startup_began') is False,'Legacy worker operation changed') + records=[m for m in runtime['members'] if m['original']['name']==member['name']] + require(len(records)==1,'Ambiguous legacy worker recovery identity') + original=records[0]['original'];recovery=records[0]['recovery_image'] + require(original['container_id']==member['container_id'] and original['image'].removeprefix('sha256:')==member['image_id'].removeprefix('sha256:') and original['config_sha256']==member['config_sha256'] and hashlib.sha256(original['body'].encode()).hexdigest()==member['unit_sha256'],'Legacy worker original identity changed') + require(recovery['source_container_id']==member['container_id'] and recovery['operation_id']==self.operation,'Legacy worker recovery image belongs to another operation') + rows=json.loads(self.run(['podman','image','inspect',recovery['image']])) + require(len(rows)==1 and rows[0]['Id'].removeprefix('sha256:')==recovery['image'].removeprefix('sha256:'),'Legacy worker recovery image changed') + image=rows[0];config=image['Config'] + require(config.get('Cmd')==['node','dist/ffmpeg-worker/worker.js'] and config.get('Entrypoint')==['docker-entrypoint.sh'],'Unrecognized legacy worker command') + created=datetime.datetime.fromisoformat(image['Created'].replace('Z','+00:00')).timestamp() + require(created<=stopped['intent_at'],'Legacy worker recovery image was not captured before stop') + source=pathlib.Path(self.run(['systemctl','--user','show',member['name']+'.service','--property=SourcePath','--value']).decode().strip()) + require(source.is_file() and not source.is_symlink() and source.suffix=='.container' and source.stat().st_uid==os.getuid() and sha(source)==member['unit_sha256'],'Legacy worker saved unit changed') + rows=self.run(['podman','ps','--all','--no-trunc','--filter','id='+member['container_id'],'--format','{{.ID}} {{.State}}']).decode().splitlines() + require(not rows or rows==[member['container_id']+' exited'],'Legacy worker may still be running') + after=self.queue('status') + require(after['paused'] is True and all(v==0 for v in after['counts'].values()),'Legacy worker queue reopened or changed') + return {'classification':'legacy-idle-worker-forced-termination','graceful':False,'completed_work_claim':False, + 'original_container_id':member['container_id'],'original_image_id':member['image_id'], + 'recovery_image_id':recovery['image'],'unit_sha256':member['unit_sha256'], + 'before_counts':before,'after_counts':after['counts'],'process_dead':True} def graceful_stop(self, name): # Save the obligation before systemd can remove an AutoRemove container. stopped=self.record.setdefault('stopped',{}) @@ -225,7 +262,6 @@ class Controller: stopped[name]={'intent_at':time.time(),'container_id':actual['Id']};self.save() self.run(['systemctl','--user','stop',name+'.service'],timeout=180) properties=self.run(['systemctl','--user','show',name+'.service','--property=ActiveState,SubState,Result,ExecMainStatus']).decode() - require('ActiveState=inactive' in properties and 'Result=success' in properties,'Service did not stop successfully') # --rm removes inspect state. Require a persisted Podman died event for # this exact original ID; a forced SIGKILL is never called completed work. events=self.run(['podman','events','--stream=false','--since',str(int(stopped[name]['intent_at'])-1),'--filter','container='+member['container_id'],'--filter','event=died','--format','json']).decode().splitlines() @@ -233,10 +269,13 @@ class Controller: matching=[event for event in matching if event.get('ID',event.get('id'))==member['container_id']] require(matching,'Original process exit evidence unavailable; hold retained') code=matching[-1].get('ContainerExitCode',matching[-1].get('containerExitCode')) - idle_worker = name=='indeedhub-ffmpeg' and self.record.get('queue_pause_confirmed') is True and self.record.get('last_queue_counts',{}).get('active')==0 + idle_worker = name=='indeedhub-ffmpeg' and self.record.get('queue_pause_confirmed') is True and valid_queue_counts(self.record.get('last_queue_counts')) and self.record['last_queue_counts']['active']==0 empty_api = name=='indeedhub-api' and self.record.get('legacy_api_empty_state') is not None - require(str(code)=='0' or (str(code)=='143' and (idle_worker or empty_api)),'Original process did not exit cleanly; active work is not claimed completed') - classification=('idle-worker-terminated-after-queue-drain' if idle_worker else 'empty-business-store-legacy-api-terminated') if str(code)=='143' else 'clean-process-exit' + forced=self.legacy_idle_worker_termination(member,properties) if str(code)=='137' else None + require(forced is not None or ('ActiveState=inactive' in properties and 'Result=success' in properties),'Service did not stop successfully') + require(forced is not None or str(code)=='0' or (str(code)=='143' and (idle_worker or empty_api)),'Original process did not exit cleanly; active work is not claimed completed') + classification=forced['classification'] if forced else (('idle-worker-terminated-after-queue-drain' if idle_worker else 'empty-business-store-legacy-api-terminated') if str(code)=='143' else 'clean-process-exit') + if forced:stopped[name]['legacy_idle_termination']=forced stopped[name].update(confirmed=True,exit_code=int(code),classification=classification,confirmed_at=time.time());self.save() def volume_sources(self): expected=VOLUMES @@ -313,7 +352,7 @@ class Controller: while True: state=self.queue('status');require(state['paused'],'Worker admission reopened') self.record['last_queue_counts']=state['counts'];self.save() - if state['counts'].get('active',0)==0:break + if state['counts']['active']==0:break require(time.monotonic()