From 9964b5c1584b43e3368064526eb382b7499c2ed8 Mon Sep 17 00:00:00 2001 From: archipelago Date: Wed, 7 Oct 2026 21:40:53 -0400 Subject: [PATCH] fix(indeehub): prove legacy API child shutdown before backup --- .../managed-update-recovery-implementation.md | 29 ++++ scripts/indeehub-maintenance-controller.py | 127 +++++++++++++++++- .../test_indeehub_maintenance_controller.py | 94 +++++++++++++ 3 files changed, 249 insertions(+), 1 deletion(-) diff --git a/docs/managed-update-recovery-implementation.md b/docs/managed-update-recovery-implementation.md index b9d91366..f0c80bff 100644 --- a/docs/managed-update-recovery-implementation.md +++ b/docs/managed-update-recovery-implementation.md @@ -325,3 +325,32 @@ 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. + +### API child-signal compatibility — 2026-10-07 + +The original API wraps `node dist/main` in npm PID1. Directly signalling that exact +child produces npm exit 1, not a clean exit. The candidate controller now records +this only as **non-graceful empty-business termination**, never as completed work. +It requires the exact original/recovery image and saved unit, automatic restart +policy `no`, fresh complete zero business-table/transaction counts, a paused empty +queue, and durable operation/container/parent/child PID+starttime+command signal +intent followed by syscall acknowledgement. Missing acknowledgement retains the +hold; retries cannot infer one. Extra direct children, reused process identity, +partial proof, OOM, forced exit 137, or replacement API writers are refused. + +Post-stop queue verification now uses atomic, read-only Redis Lua, authenticated +through stdin rather than secret command arguments. The observed original API +queue endpoint must bind to the exact original Redis ID/image and shared network +alias. All six counts, prioritized and waiting-children must remain zero, with +admission paused. No default or absent field can establish empty state. + +All **38 pure controller tests passed**, including actual Node execution of the +process selector. The actual hardened VM Redis observer passed. A uniquely named +recovery-image API process probe, with no persistent mounts or published ports, +passed exact child-starttime/signal-acknowledgement checks and exited 1 without +OOM; bound Redis observations before/after were paused and entirely zero. The +probe was removed. Its receipt is `api-process-probe.receipt.json` in the private +VM fixture directory. This qualifies the primitive, not a full backend update. +The earlier held API diagnostic lacks the new durable proof and remains +unaccepted; it must be restored by the matching rebuilt manager before a fresh +full update/rollback rehearsal. No live Yaya mutation or activation occurred. diff --git a/scripts/indeehub-maintenance-controller.py b/scripts/indeehub-maintenance-controller.py index ffee9222..1e5b0347 100644 --- a/scripts/indeehub-maintenance-controller.py +++ b/scripts/indeehub-maintenance-controller.py @@ -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()=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') diff --git a/tests/regression/test_indeehub_maintenance_controller.py b/tests/regression/test_indeehub_maintenance_controller.py index 48291856..544bc111 100644 --- a/tests/regression/test_indeehub_maintenance_controller.py +++ b/tests/regression/test_indeehub_maintenance_controller.py @@ -259,6 +259,100 @@ class MaintenanceTests(unittest.TestCase): with self.assertRaises(RuntimeError):c.legacy_idle_worker_termination(w,p) c.record['ingress_closed']=True;unit.write_text('changed') with self.assertRaises(RuntimeError):c.legacy_idle_worker_termination(w,p) + def legacy_api_fixture(self): + c,worker,props,r,rp,image,observed,state,unit=self.legacy_forced_fixture() + api=next(m for m in members() if m['name']=='indeedhub-api');api['unit_sha256']=module.sha(unit) + original=r['members'][0]['original'];original.update(name=api['name'],container_id=api['container_id']) + r['members'][0]['recovery_image']['source_container_id']=api['container_id'];module.atomic(rp,r) + image['Config']['Cmd']=module.LEGACY_API_CMD;image['Config']['Env']=[] + observed['extra']={'prioritized':0,'waiting_children':0} + c.record['original_members']=[api if m['name']==api['name'] else m for m in members()] + c.record['legacy_api_empty_state']={k:0 for k in ('projects','contents','payments','shareholders','subscriptions','library_items','other_active_transactions')} + c.record['stopped'][api['name']]={'container_id':api['container_id'],'intent_at':1700000000,'api_signal':{'operation_id':self.operation,'container_id':api['container_id'],'intent_at':1700000001,'acknowledged':True,'prefix':'bull:transcode:','queue_binding':{'container_id':next(x['container_id'] for x in members() if x['name']=='indeedhub-redis'),'image_id':'a'*64,'network_id':'network-id','host':'indeedhub-redis','port':6379}}} + c.record['stopped'][api['name']]['api_signal'].update(acknowledged_at=1700000002,proof={'parent':{'pid':1,'ppid':0,'starttime':'10','command':'npm run start:prod\0'},'child':{'pid':35,'ppid':1,'starttime':'20','command':'node\0dist/main\0'}},before_queue=json.loads(json.dumps(observed))) + previous=c.runner + def runner(argv,timeout,output): + if argv[:2]==['podman','inspect']: + row=json.loads(self.command_runner(argv,timeout,output))[0];row['NetworkSettings']={'Networks':{'app':{'NetworkID':'network-id','Aliases':['indeedhub-redis']}}};return json.dumps([row]).encode() + if argv[:2]==['podman','events']:return b'' + return previous(argv,timeout,output) + c.runner=runner + return c,api,'ActiveState=failed\nSubState=failed\nResult=exit-code\nExecMainStatus=1\n',image,observed,state + def test_api_wrapper_requires_acknowledged_original_signal_and_empty_queue(self): + c,m,p,image,observed,state=self.legacy_api_fixture() + proof=c.legacy_api_wrapper_termination(m,p) + self.assertFalse(proof['graceful']);self.assertFalse(proof['completed_work_claim']);self.assertTrue(proof['process_dead']) + signal=c.record['stopped'][m['name']]['api_signal'] + for field,value in [('acknowledged',False),('container_id','other'),('operation_id','other'),('intent_at',0)]: + original=signal[field];signal[field]=value + with self.assertRaises(RuntimeError):c.legacy_api_wrapper_termination(m,p) + signal[field]=original + for field in ('prioritized','waiting_children'): + observed['extra'][field]=1 + with self.assertRaises(RuntimeError):c.legacy_api_wrapper_termination(m,p) + observed['extra'][field]=0 + for value in ('137','143','0'): + with self.assertRaises(RuntimeError):c.legacy_api_wrapper_termination(m,p.replace('ExecMainStatus=1\n','ExecMainStatus='+value+'\n')) + state['ps']=m['container_id'].encode() + with self.assertRaises(RuntimeError):c.legacy_api_wrapper_termination(m,p) + def test_api_queue_observer_rejects_wrong_prefix_missing_counts_and_ambiguous_credentials(self): + c,m,p,image,observed,state=self.legacy_api_fixture() + with self.assertRaises(RuntimeError):c.observe_empty_redis_queue(m,'wrong:') + observed['counts'].pop('active') + with self.assertRaises(RuntimeError):c.observe_empty_redis_queue(m,'bull:transcode:') + observed['counts']['active']=0 + for env in (['QUEUE_PASSWORD=one','QUEUE_PASSWORD=two'],['QUEUE_PASSWORD=bad\nframe']): + image['Config']['Env']=env + with self.assertRaises(RuntimeError):c.observe_empty_redis_queue(m,'bull:transcode:') + def test_api_partial_durable_proof_never_substitutes_for_signal_acknowledgement(self): + import copy + c,m,p,*_=self.legacy_api_fixture();signal=c.record['stopped'][m['name']]['api_signal'];original=copy.deepcopy(signal) + for key in ('proof','before_queue','acknowledged_at'): + signal.clear();signal.update(copy.deepcopy(original));signal.pop(key) + with self.assertRaises(RuntimeError):c.legacy_api_wrapper_termination(m,p) + signal.clear();signal.update(original) + for counts in ({'projects':0},{'unrelated':0},dict.fromkeys(module.BUSINESS_COUNTS,False)): + c.record['legacy_api_empty_state']=counts + with self.assertRaises(RuntimeError):c.legacy_api_wrapper_termination(m,p) + def test_api_queue_endpoint_must_share_exact_original_redis_alias(self): + c,m,p,image,*_=self.legacy_api_fixture() + image['Config']['Env']=['QUEUE_HOST=indeedhub-redis','QUEUE_PORT=6379'] + api={'Config':{'Env':list(image['Config']['Env'])},'NetworkSettings':{'Networks':{'app':{'NetworkID':'network-id'}}}} + self.assertEqual(c.api_queue_binding(m,api,image)['host'],'indeedhub-redis') + api['NetworkSettings']['Networks']['app']['NetworkID']='different' + with self.assertRaises(RuntimeError):c.api_queue_binding(m,api,image) + api['NetworkSettings']['Networks']['app']['NetworkID']='network-id';api['Config']['Env']=['QUEUE_HOST=foreign'] + with self.assertRaises(RuntimeError):c.api_queue_binding(m,api,image) + def test_api_oom_or_replacement_keeps_termination_unconfirmed(self): + c,m,p,image,observed,state=self.legacy_api_fixture();runner=c.runner + c.runner=lambda argv,timeout,output: b'{"Status":"oom"}' if argv[:2]==['podman','events'] else runner(argv,timeout,output) + with self.assertRaisesRegex(RuntimeError,'OOM'):c.legacy_api_wrapper_termination(m,p) + c.runner=lambda argv,timeout,output: b'replacement-id' if argv[:2]==['podman','ps'] and any('name=^' in x for x in argv) else runner(argv,timeout,output) + with self.assertRaisesRegex(RuntimeError,'writer is still running'):c.legacy_api_wrapper_termination(m,p) + def test_unacknowledged_api_signal_is_never_retried(self): + c,m,p,*_=self.legacy_api_fixture();c.record['stopped'][m['name']]['api_signal']['acknowledged']=False + with self.assertRaisesRegex(RuntimeError,'Unacknowledged'):c.signal_legacy_api(m) + self.assertEqual(self.calls,[]) + def test_actual_node_process_selector_rejects_extra_children_reuse_and_failed_signal(self): + import subprocess + harness=r'''const vm=require('vm'),assert=require('assert'),script=JSON.parse(require('fs').readFileSync(0,'utf8')); +function run(change={},action='probe',expected){ + const rows={1:{ppid:0,cmd:'npm run start:prod\0\0',start:'100'},35:{ppid:1,cmd:'node\0dist/main\0',start:'200'},999:{ppid:0,cmd:'node probe\0',start:'300'}}; + for(const [key,value] of Object.entries(change))rows[key]=value; + let output='',killed=[]; + const fs={readdirSync:()=>Object.keys(rows),readFileSync:p=>{const [,id,file]=p.match(/^\/proc\/(\d+)\/(.*)$/);const r=rows[id];if(!r)throw Error('missing');if(file==='cmdline')return r.cmd;if(file==='status')return `PPid:\t${r.ppid}\n`;return `${id} (node) S `+Array(18).fill('0').join(' ')+' '+r.start+' 0';},writeSync:(_,v)=>{output+=v}}; + vm.runInNewContext(script,{require:n=>{assert.equal(n,'fs');return fs},process:{argv:['node',action,JSON.stringify(expected)],kill:(p,s)=>{if(change.fail)throw Error('signal rejected');killed.push([p,s])}}}); + return {proof:JSON.parse(output),killed}; +} +const proof=run().proof;assert.equal(proof.child.starttime,'200'); +assert.deepEqual(run({},'signal',proof).killed,[[35,'SIGTERM']]); +assert.throws(()=>run({36:{ppid:1,cmd:'other\0',start:'400'}})); +assert.throws(()=>run({35:{ppid:1,cmd:'node other\0',start:'200'}})); +assert.throws(()=>run({35:{ppid:1,cmd:'node\0dist/main\0',start:'201'}},'signal',proof)); +assert.throws(()=>run({1:{ppid:0,cmd:'other\0',start:'100'}})); +console.log('process identity cases passed');''' + result=subprocess.run(['node','-e',harness],input=json.dumps(module.API_PROCESS_SCRIPT),text=True,capture_output=True) + self.assertEqual(result.returncode,0,result.stderr) def test_legacy_worker_sigterm_requires_proven_paused_idle_queue(self): c=self.controller;c.record={'operation_id':self.operation,'phase':'Prepared','original_members':module.validate_members(members())};c.save() worker=next(m for m in members() if m['name']=='indeedhub-ffmpeg')