fix(indeehub): prove legacy API child shutdown before backup
This commit is contained in:
@@ -10,11 +10,43 @@ 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())
|
||||
BUSINESS_COUNTS = frozenset(('projects','contents','payments','shareholders','subscriptions','library_items','other_active_transactions'))
|
||||
def valid_empty_business(counts):
|
||||
return isinstance(counts,dict) and set(counts)==BUSINESS_COUNTS and all(type(v) is int and v==0 for v in counts.values())
|
||||
def valid_empty_queue(observed):
|
||||
return isinstance(observed,dict) and set(observed)=={'paused','counts','extra'} and observed['paused'] is True and valid_queue_counts(observed['counts']) and all(v==0 for v in observed['counts'].values()) and isinstance(observed['extra'],dict) and set(observed['extra'])=={'prioritized','waiting_children'} and all(type(v) is int and v==0 for v in observed['extra'].values())
|
||||
def valid_api_process_proof(proof):
|
||||
if not isinstance(proof,dict) or set(proof)!={'parent','child'}:return False
|
||||
for part in proof.values():
|
||||
if not isinstance(part,dict) or set(part)!={'pid','ppid','starttime','command'}:return False
|
||||
if type(part['pid']) is not int or type(part['ppid']) is not int or not isinstance(part['starttime'],str) or not re.fullmatch('[0-9]+',part['starttime']) or not isinstance(part['command'],str):return False
|
||||
parent=proof['parent'];child=proof['child']
|
||||
return parent['pid']==1 and parent['ppid']==0 and parent['command'].rstrip('\0')=='npm run start:prod' and child['pid']>1 and child['ppid']==1 and child['command']=='node\0dist/main\0'
|
||||
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');
|
||||
console.log(JSON.stringify({paused:await q.isPaused(),counts:await q.getJobCounts('active','waiting','paused','delayed','failed','completed')}));}
|
||||
console.log(JSON.stringify({paused:await q.isPaused(),prefix:q.toKey(''),counts:await q.getJobCounts('active','waiting','paused','delayed','failed','completed')}));}
|
||||
finally{await q.close()}})().catch(()=>process.exit(1));'''
|
||||
|
||||
REDIS_OBSERVE_SCRIPT = "-- Atomic observation only. KEYS[1] must be the previously captured and validated\n-- exact queue prefix including trailing colon (for example bull:transcode:).\n-- No Queue constructor, marker deletion, pause/resume or data write occurs.\nlocal p = KEYS[1]\nif #KEYS ~= 1 or #ARGV ~= 0 or not p or #p == 0 then\n return redis.error_reply('Invalid bound queue prefix')\nend\nlocal counts = {\n active = redis.call('LLEN', p .. 'active'),\n waiting = redis.call('LLEN', p .. 'wait'),\n paused = redis.call('LLEN', p .. 'paused'),\n delayed = redis.call('ZCARD', p .. 'delayed'),\n failed = redis.call('ZCARD', p .. 'failed'),\n completed = redis.call('ZCARD', p .. 'completed')\n}\nreturn cjson.encode({\n paused = redis.call('HEXISTS', p .. 'meta', 'paused') == 1,\n counts = counts,\n extra = {\n prioritized = redis.call('ZCARD', p .. 'prioritized'),\n waiting_children = redis.call('ZCARD', p .. 'waiting-children')\n }\n})\n"
|
||||
API_PROCESS_SCRIPT = r'''const fs=require('fs');
|
||||
const command=p=>fs.readFileSync(`/proc/${p}/cmdline`,'utf8');
|
||||
const parent=p=>Number(fs.readFileSync(`/proc/${p}/status`,'utf8').match(/^PPid:\s+(\d+)$/m)?.[1]);
|
||||
const identity=p=>{const value=fs.readFileSync(`/proc/${p}/stat`,'utf8');return {pid:Number(p),ppid:parent(p),starttime:value.slice(value.lastIndexOf(') ')+2).split(' ')[19],command:command(p)}};
|
||||
function observe(){
|
||||
if(command(1).replace(/\0+$/,'')!=='npm run start:prod')throw Error('Unrecognized API parent');
|
||||
const children=fs.readdirSync('/proc').filter(p=>/^\d+$/.test(p)).filter(p=>{try{return parent(p)===1}catch{return false}});
|
||||
if(children.length!==1)throw Error('Ambiguous API children');
|
||||
const child=identity(children[0]);if(child.command!=='node\0dist/main\0'||child.ppid!==1||!/^\d+$/.test(child.starttime))throw Error('Unrecognized API child');
|
||||
return {parent:identity(1),child};
|
||||
}
|
||||
const proof=observe(),action=process.argv[1];
|
||||
if(action==='signal'){
|
||||
const expected=JSON.parse(process.argv[2]);if(JSON.stringify(proof)!==JSON.stringify(expected)||JSON.stringify(observe())!==JSON.stringify(expected))throw Error('API process identity changed');
|
||||
process.kill(proof.child.pid,'SIGTERM');fs.writeSync(1,JSON.stringify({proof,signal:'SIGTERM',acknowledged:true})+'\n');
|
||||
}else if(action==='probe')fs.writeSync(1,JSON.stringify(proof)+'\n');else throw Error('Unsupported action');'''
|
||||
LEGACY_API_CMD = ['sh','-c',"echo 'Running database migrations...' && npx typeorm migration:run -d dist/database/ormconfig.js && echo 'Migrations complete.' && npm run start:prod"]
|
||||
|
||||
# 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 = {
|
||||
@@ -217,6 +249,97 @@ 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 api_recovery_identity(self, member):
|
||||
self.holds();self.fence_matches()
|
||||
require(member['name']=='indeedhub-api','Legacy API member required')
|
||||
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 API operation changed')
|
||||
rows=[m for m in runtime['members'] if m['original']['name']==member['name']]
|
||||
require(len(rows)==1,'Ambiguous legacy API recovery identity');original=rows[0]['original'];recovery=rows[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 API original identity changed')
|
||||
require(recovery['source_container_id']==member['container_id'] and recovery['operation_id']==self.operation,'Legacy API recovery identity changed')
|
||||
images=json.loads(self.run(['podman','image','inspect',recovery['image']]))
|
||||
require(len(images)==1 and images[0]['Id'].removeprefix('sha256:')==recovery['image'].removeprefix('sha256:'),'Legacy API recovery image changed')
|
||||
image=images[0];require(image['Config'].get('Cmd')==LEGACY_API_CMD and image['Config'].get('Entrypoint')==['docker-entrypoint.sh'],'Unrecognized legacy API command')
|
||||
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 API saved unit changed')
|
||||
return image
|
||||
def api_queue_binding(self, member, api, image):
|
||||
def environment(config):
|
||||
result={}
|
||||
for value in config.get('Env',[]):
|
||||
key,_,text=value.partition('=')
|
||||
if key in ('QUEUE_HOST','QUEUE_PORT','QUEUE_PASSWORD'):
|
||||
require(key not in result,'Ambiguous queue endpoint');result[key]=text
|
||||
return result
|
||||
expected=environment(image['Config']);actual=environment(api['Config'])
|
||||
require(actual==expected and expected.get('QUEUE_PORT','6379')=='6379','Original API queue endpoint changed')
|
||||
redis=next(m for m in self.record['original_members'] if m['name']=='indeedhub-redis')
|
||||
store=self.inspect(redis['name']);require(store['Id']==redis['container_id'] and store['Image']==redis['image_id'],'Original queue store changed')
|
||||
api_networks=api['NetworkSettings']['Networks'];networks=store['NetworkSettings']['Networks']
|
||||
host=expected.get('QUEUE_HOST');matches=[]
|
||||
for name,network in networks.items():
|
||||
if name in api_networks and api_networks[name]['NetworkID']==network['NetworkID'] and host in (network.get('Aliases') or []):matches.append(network['NetworkID'])
|
||||
require(len(matches)==1,'API queue endpoint is not bound to original Redis network alias')
|
||||
return {'container_id':redis['container_id'],'image_id':redis['image_id'],'network_id':matches[0],'host':host,'port':6379}
|
||||
def require_api_stopped(self, member):
|
||||
require(not self.run(['podman','ps','--no-trunc','--filter','name=^'+member['name']+'$','--format','{{.ID}}']).strip(),'An API writer is still running')
|
||||
intent=self.record['stopped'][member['name']]['intent_at']
|
||||
events=self.run(['podman','events','--stream=false','--since',str(int(intent)-1),'--filter','container='+member['container_id'],'--filter','event=oom','--format','json'])
|
||||
require(not events.strip(),'API OOM evidence prevents compatibility termination')
|
||||
def observe_empty_redis_queue(self, member, prefix, binding=None):
|
||||
require(prefix=='bull:transcode:','Unrecognized bound queue prefix')
|
||||
redis=next(m for m in self.record['original_members'] if m['name']=='indeedhub-redis')
|
||||
actual=self.inspect(redis['name']);require(actual['Id']==redis['container_id'] and actual['Image']==redis['image_id'],'Original queue store changed')
|
||||
image=self.api_recovery_identity(member)
|
||||
if binding is not None:
|
||||
networks=actual['NetworkSettings']['Networks']
|
||||
require(binding.get('container_id')==redis['container_id'] and binding.get('image_id')==redis['image_id'] and binding.get('port')==6379 and sum(network.get('NetworkID')==binding.get('network_id') and binding.get('host') in (network.get('Aliases') or []) for network in networks.values())==1,'Original queue endpoint binding changed')
|
||||
values=[v.split('=',1)[1] for v in image['Config'].get('Env',[]) if v.startswith('QUEUE_PASSWORD=')]
|
||||
require(len(values)<=1,'Ambiguous queue credential');password=values[0] if values else ''
|
||||
require('\n' not in password and '\r' not in password,'Unsupported queue credential framing')
|
||||
shell='IFS= read -r secret || exit 1; if [ -n "$secret" ]; then REDISCLI_AUTH="$secret"; export REDISCLI_AUTH; else unset REDISCLI_AUTH; fi; exec redis-cli --no-auth-warning --raw EVAL "$1" 1 "$2"'
|
||||
observed=json.loads(self.run(['podman','exec','-i',redis['container_id'],'sh','-c',shell,'queue-observer',REDIS_OBSERVE_SCRIPT,prefix],input_bytes=(password+'\n').encode()))
|
||||
require(observed.get('paused') is True and valid_queue_counts(observed.get('counts')) and all(v==0 for v in observed['counts'].values()),'Legacy API queue not paused and empty')
|
||||
extra=observed.get('extra');require(isinstance(extra,dict) and set(extra)=={'prioritized','waiting_children'} and all(type(v) is int and v==0 for v in extra.values()),'Legacy API has additional queued work')
|
||||
return observed
|
||||
def signal_legacy_api(self, member):
|
||||
stopped=self.record['stopped'][member['name']]
|
||||
if stopped.get('api_signal'):
|
||||
require(stopped['api_signal'].get('acknowledged') is True,'Unacknowledged legacy API signal; hold retained')
|
||||
return
|
||||
image=self.api_recovery_identity(member)
|
||||
require(datetime.datetime.fromisoformat(image['Created'].replace('Z','+00:00')).timestamp()<=stopped['intent_at'],'API recovery image was not captured before stop')
|
||||
actual=self.inspect(member['name']);require(actual['Id']==member['container_id'] and actual['Image']==member['image_id'] and actual['State']['Running'],'Original API changed before signal')
|
||||
restart=self.run(['systemctl','--user','show',member['name']+'.service','--property=Restart','--value']).decode().strip()
|
||||
require(restart=='no','Legacy API signal requires no automatic restart')
|
||||
binding=self.api_queue_binding(member,actual,image)
|
||||
queue=self.queue('status');prefix=queue.get('prefix');before=self.observe_empty_redis_queue(member,prefix,binding)
|
||||
self.legacy_api_idle()
|
||||
proof=json.loads(self.run(['podman','exec',member['container_id'],'node','-e',API_PROCESS_SCRIPT,'probe']))
|
||||
require(valid_api_process_proof(proof),'Incomplete API process identity proof')
|
||||
intent={'operation_id':self.operation,'container_id':member['container_id'],'proof':proof,'prefix':prefix,'queue_binding':binding,'before_queue':before,'intent_at':time.time(),'acknowledged':False}
|
||||
stopped['api_signal']=intent;self.save()
|
||||
ack=json.loads(self.run(['podman','exec',member['container_id'],'node','-e',API_PROCESS_SCRIPT,'signal',json.dumps(proof,separators=(',',':'))]))
|
||||
require(ack=={'proof':proof,'signal':'SIGTERM','acknowledged':True},'Legacy API signal acknowledgement changed')
|
||||
intent['acknowledged']=True;intent['acknowledged_at']=time.time();self.save()
|
||||
deadline=time.monotonic()+30
|
||||
while self.run(['podman','ps','--no-trunc','--filter','id='+member['container_id'],'--format','{{.ID}}']).strip():
|
||||
require(time.monotonic()<deadline,'Legacy API did not terminate after child signal');time.sleep(0.2)
|
||||
self.require_api_stopped(member)
|
||||
def legacy_api_wrapper_termination(self, member, properties):
|
||||
self.api_recovery_identity(member)
|
||||
stopped=self.record['stopped'][member['name']];signal=stopped.get('api_signal',{})
|
||||
require(signal.get('operation_id')==self.operation and signal.get('container_id')==member['container_id'] and signal.get('acknowledged') is True and signal.get('intent_at',0)>=stopped['intent_at'],'Legacy API signal proof missing')
|
||||
require(valid_empty_business(self.record.get('legacy_api_empty_state')),'Legacy API empty-business proof missing')
|
||||
require(valid_api_process_proof(signal.get('proof')) and valid_empty_queue(signal.get('before_queue')) and type(signal.get('acknowledged_at')) in (int,float) and signal['acknowledged_at']>=signal['intent_at'],'Legacy API durable signal proof incomplete')
|
||||
state=dict(line.split('=',1) for line in properties.splitlines() if '=' in line)
|
||||
require(state.get('ActiveState')=='failed' and state.get('SubState')=='failed' and state.get('ExecMainStatus')=='1' and state.get('Result')=='exit-code','Unrecognized API wrapper exit')
|
||||
require(not self.run(['podman','ps','--no-trunc','--filter','id='+member['container_id'],'--format','{{.ID}}']).strip(),'Legacy API is still running')
|
||||
self.require_api_stopped(member)
|
||||
require(isinstance(signal.get('queue_binding'),dict),'API queue binding proof missing')
|
||||
after=self.observe_empty_redis_queue(member,signal['prefix'],signal['queue_binding'])
|
||||
return {'classification':'empty-business-legacy-api-child-signal-termination','graceful':False,'completed_work_claim':False,'process_dead':True,'signal':signal,'after_queue':after}
|
||||
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.
|
||||
@@ -260,6 +383,7 @@ class Controller:
|
||||
if name not in stopped:
|
||||
actual=self.inspect(name);require(actual['Id']==member['container_id'] and actual['Image']==member['image_id'],'Original container changed before stop')
|
||||
stopped[name]={'intent_at':time.time(),'container_id':actual['Id']};self.save()
|
||||
if name=='indeedhub-api':self.signal_legacy_api(member)
|
||||
self.run(['systemctl','--user','stop',name+'.service'],timeout=180)
|
||||
properties=self.run(['systemctl','--user','show',name+'.service','--property=ActiveState,SubState,Result,ExecMainStatus']).decode()
|
||||
# --rm removes inspect state. Require a persisted Podman died event for
|
||||
@@ -272,6 +396,7 @@ class Controller:
|
||||
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
|
||||
forced=self.legacy_idle_worker_termination(member,properties) if str(code)=='137' else None
|
||||
if name=='indeedhub-api' and str(code)=='1':forced=self.legacy_api_wrapper_termination(member,properties)
|
||||
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')
|
||||
|
||||
Reference in New Issue
Block a user