fix(indeehub): prove idle legacy worker termination before backup
This commit is contained in:
@@ -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()<deadline,'Transcodes still active; retained job state, no forced completion')
|
||||
time.sleep(1)
|
||||
self.graceful_stop('indeedhub-ffmpeg');self.legacy_api_idle();self.graceful_stop('indeedhub-api')
|
||||
|
||||
Reference in New Issue
Block a user