Add operation-owned legacy IndeeHub maintenance controller draft

This commit is contained in:
archipelago
2026-10-07 02:24:50 -04:00
parent 6e7ea8b9d6
commit 6d450caebf
3 changed files with 402 additions and 0 deletions
@@ -0,0 +1,73 @@
# Legacy IndeeHub maintenance controller
Status: isolated source implementation. Ten pure Python fake-runtime regressions
pass; no live invocation or production qualification. The controller is not part
of the already signed private app candidate and needs no new app image/API.
The supervised updater owns the lifecycle flock, seven durable holds, original
Quadlets and private writable-layer recovery images. It records its destructive
obligation before invoking the fixed controller with bounded JSON over stdin:
```
python3 /opt/archipelago/scripts/indeehub-maintenance-controller.py acquire
{ "operation_id": "<uuid>", "original_members": [
{ "name": "indeedhub", "container_id": "<64hex>", "image_id": "<64hex>",
"unit_sha256": "<64hex>", "config_sha256": "<64hex>", "running": true }
// All seven exact members; JSON does not include this illustrative comment.
], "recovery": false }
```
Other actions are `verify` with operation_id, and `release` with operation_id and
outcome committed/restored/aborted. Replies are <=4KiB and report drained, held,
released, or recovering for explicit recovery acquire. Inherited
ARCHY_UPDATE_LOCK_FD stays open and is passed to child commands; the script never
unlocks it. Journal: data/update-transactions/indeehub-maintenance/<uuid>/journal.json.
## Forward sequence
- Validate exact original IDs/source-unit hashes and known port exposure. Only
frontend127.0.0.1:7778 is supported; direct backend/S3 ports refuse before stop.
- Require deployed native AppGate and legacy nginx maintenance guards. Inspect
every known legacy sublocation and any direct7778 proxy; unknown routes refuse.
Save an operation-owned readable sentinel, then verify local ingress returns503.
- Record prior BullMQ transcode pause state, globally pause future job admission,
retain queued/delayed/failed jobs. Gracefully stop frontend ingress; bounded
polling waits for active transcodes to finish before stopping worker and API.
- Require successful systemd shutdown plus an exact original Podman died event
with exit0. Forced exits and missing event evidence retain the hold and are
never labelled completed writes.
- While PostgreSQL remains running, capture a fresh custom dump. Cleanly stop
MinIO/Redis/relay/Postgres, then archive all four complete quiescent volumes
(including SQLite WAL and Redis persistence) with metadata. No volume deletion
or migration rollback. Archive hashes/size and per-step obligations are durable.
- Keep admission closed while the updater renders, starts and verifies targets.
## Interrupted recovery
The node first records phase Restoring with boolean target_startup_began, then
calls acquire with recovery:true. That path preserves the original failure and
fence; it does not retry a killed original into a fictitious successful drain or
claim missing backups exist. The node restores exact saved old runtime under the
same hold. Release before any target startup can state only that original runtime
was restored. If target startup/migration began, recorded data-compatibility
verification is required before restored release; an old image alone does not
prove compatibility with newly changed data. No automatic DB/media restore exists.
## Qualification and remaining integration
`python3 tests/regression/test_indeehub_maintenance_controller.py` passes ten
fake-runtime cases in temporary directories, without services/network/containers.
Source nginx template guard coverage also passes its parser check. Production
adapter compilation, actual Podman event format/systemd clean-exit behavior,
application writer shutdown, interrupted backup and supervised restart still need
isolated lifecycle fixtures and then coordinated node acceptance. A long-lived
WebSocket or active upload can exceed graceful-stop deadlines; the current code
refuses completion and preserves recovery obligations rather than silently
calling interrupted work finished.
The deployment must install the exact qualified controller script and record its
hash alongside the backend artifact. Binary-only deployment does not install it.
The backend must refuse missing/mismatched prerequisites before snapshots/stops.
Native AppGate + nginx guards are separate node source changes owned by the
supervised updater agent. The signed app catalog/private image receipts remain
unchanged. Existing live stop/uninstall intent must not be rewritten as maintenance.
+244
View File
@@ -0,0 +1,244 @@
#!/usr/bin/env python3
"""Operation-owned legacy IndeeHub maintenance. Called only by supervised updater.
No live execution is part of source qualification. Original writable-layer images
must already be durable. Never unlock ARCHY_UPDATE_LOCK_FD or release another hold.
"""
import datetime, hashlib, json, os, pathlib, re, shutil, subprocess, sys, time, uuid
NAMES = ('indeedhub','indeedhub-api','indeedhub-ffmpeg','indeedhub-minio','indeedhub-postgres','indeedhub-redis','indeedhub-relay')
DATA = pathlib.Path('/var/lib/archipelago')
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')}));}
finally{await q.close()}})().catch(()=>process.exit(1));'''
def require(condition, message):
if not condition: raise RuntimeError(message)
def atomic(path, value):
path.parent.mkdir(mode=0o700, parents=True, exist_ok=True)
require(not path.is_symlink(), 'Refuse symbolic journal path')
temporary=path.with_name('.'+path.name+'.'+str(uuid.uuid4()))
with temporary.open('x') as stream:
json.dump(value,stream,separators=(',',':'));stream.flush();os.fsync(stream.fileno())
os.chmod(temporary,0o600);os.replace(temporary,path)
descriptor=os.open(path.parent,os.O_RDONLY);os.fsync(descriptor);os.close(descriptor)
def sha(path):
with path.open('rb') as stream:return hashlib.file_digest(stream,'sha256').hexdigest()
def validate_members(members):
require(isinstance(members,list) and len(members)==7,'Seven exact original members required')
require({m.get('name') for m in members}==set(NAMES),'IndeeHub member scope changed')
for m in members:
require(set(m)=={'name','container_id','image_id','unit_sha256','config_sha256','running'},'Unexpected member fields')
for key in ('container_id','image_id','unit_sha256','config_sha256'):
require(bool(re.fullmatch('[0-9a-f]{64}',m[key])),'Invalid original identity/hash')
require(m['running'] is True,'Legacy barrier currently supports an originally running complete stack only')
return sorted(members,key=lambda m:m['name'])
def validate_nginx_guards(config):
guard='if (-f /var/lib/archipelago/app-maintenance/indeedhub) { return 503; }'
lines=config.splitlines();matched=0;index=0
while index<len(lines):
header=lines[index]
if not re.match(r'^\s*location\b.*\{\s*$',header):index+=1;continue
block=[];depth=0
while index<len(lines):
line=lines[index];block.append(line)
unquoted=re.sub(r"(['\"])(?:\\.|(?!\1).)*\1",'',line).split('#',1)[0]
depth+=unquoted.count('{')-unquoted.count('}');index+=1
if depth==0:break
body='\n'.join(block)
if '/app/indeedhub' in header or re.search(r'proxy_pass\s+https?://127[.]0[.]0[.]1:7778(?:/|;)',body):
require(bool(re.match(r'^\s*location\s+(?:\^~\s+|=\s+)?/app/indeedhub(?:/[^\s{]*)?\s*\{\s*$',header)),'Unrecognized direct IndeeHub proxy exposure')
require(guard in body,'An IndeeHub proxy route is missing its maintenance guard')
matched+=1
require(matched>=3,'Expected complete legacy IndeeHub route guards')
return matched
class Controller:
def __init__(self, data, operation, lock_fd, runner=None):
self.data=pathlib.Path(data);self.operation=operation;self.lock_fd=lock_fd;self.runner=runner
require(str(uuid.UUID(operation))==operation,'Invalid operation UUID')
self.root=self.data/'update-transactions'/'indeehub-maintenance'/operation
self.path=self.root/'journal.json';self.fence=self.data/'app-maintenance'/'indeedhub'
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):
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,))
if output:return b''
require(len(result.stdout)<=2*1024*1024,'Command response exceeds bound')
return result.stdout
def inspect(self, name):
rows=json.loads(self.run(['podman','inspect',name]));require(len(rows)==1,'Unexpected container inspection');return rows[0]
def holds(self):
for name in NAMES:
path=self.data/'update-transactions'/'holds'/name
require(path.is_file() and not path.is_symlink() and path.read_text()==self.operation,'Matching durable lifecycle hold required')
def fence_matches(self):
require(self.fence.is_file() and not self.fence.is_symlink() and self.fence.read_text()==self.operation,'Admission fence ownership changed')
def close_ingress(self):
# The deployed native AppGate and legacy nginx guards consume this exact
# sentinel. This code never edits arbitrary nginx configuration.
config=self.run(['sudo','-n','nginx','-T']).decode()
validate_nginx_guards(config)
self.fence.parent.mkdir(mode=0o755,exist_ok=True)
self.fence.parent.chmod(0o755)
require(not self.fence.parent.is_symlink(),'Admission directory is a symlink')
if self.fence.exists():self.fence_matches()
else:
with self.fence.open('x') as stream:stream.write(self.operation);stream.flush();os.fsync(stream.fileno())
self.fence.chmod(0o644)
fd=os.open(self.fence.parent,os.O_RDONLY);os.fsync(fd);os.close(fd)
# A local legacy probe must be rejected without entering the old app.
import urllib.request,urllib.error
try:
urllib.request.urlopen('http://127.0.0.1/app/indeedhub/__maintenance_probe',timeout=5)
raise RuntimeError('Legacy ingress was not fenced')
except urllib.error.HTTPError as error:
require(error.code==503,'Legacy ingress guard did not return maintenance status')
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')
return result
def pause_queue(self):
if 'queue_was_paused' not in self.record:
original=self.queue('status');self.record['queue_was_paused']=original['paused'];self.record['queue_original_counts']=original['counts'];self.save()
state=self.queue('pause');require(state['paused'],'Worker admission did not close')
self.record['queue_pause_confirmed']=True;self.save()
def graceful_stop(self, name):
# Save the obligation before systemd can remove an AutoRemove container.
stopped=self.record.setdefault('stopped',{})
if stopped.get(name,{}).get('confirmed'):return
member=next(m for m in self.record['original_members'] if m['name']==name)
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()
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()
matching=[json.loads(line) for line in events if line.strip()]
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'))
require(str(code)=='0','Original process did not exit cleanly; active work is not claimed completed')
stopped[name].update(confirmed=True,exit_code=0,confirmed_at=time.time());self.save()
def volume_sources(self):
expected=['indeedhub-minio-data','indeedhub-postgres-data','indeedhub-redis-data','indeedhub-relay-data']
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 backup(self):
if self.record.get('backup_complete'):return
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',{}):
path=backup/'database.dump.partial'
with path.open('wb') as output:self.run(['podman','exec','indeedhub-postgres','pg_dump','-U','indeedhub','-d','indeedhub','--format=custom','--no-owner','--no-acl'],timeout=300,output=output);output.flush();os.fsync(output.fileno())
final=backup/'database.dump';os.replace(path,final);self.record['artifacts']['database.dump']={'bytes':final.stat().st_size,'sha256':sha(final)};self.save()
# Redis stop flushes persisted queue state; its clean process exit is
# checked exactly as every other service. SQLite WAL is archived with DB.
for name in ('indeedhub-minio','indeedhub-redis','indeedhub-relay','indeedhub-postgres'):self.graceful_stop(name)
for volume,source in sources.items():
name=volume+'.tar'
if name in self.record['artifacts']:continue
require(pathlib.Path(source).is_absolute() and source.endswith('/_data'),'Invalid volume mountpoint')
available=shutil.disk_usage(backup).free
measured=int(self.run(['podman','unshare','du','-sb',source]).decode().split()[0])
require(available>measured+512*1024*1024,'Insufficient durable backup space')
partial=backup/(name+'.partial')
with partial.open('wb') as output:self.run(['podman','unshare','tar','--xattrs','--acls','--numeric-owner','-C',source,'-cpf','-','.'],timeout=1800,output=output);output.flush();os.fsync(output.fileno())
final=backup/name;os.replace(partial,final);self.record['artifacts'][name]={'bytes':final.stat().st_size,'sha256':sha(final)};self.save()
self.record['backup_complete']=True;self.record['phase']='Drained';self.save()
def acquire(self, members, recovery=False):
members=validate_members(members);self.holds()
if self.record:require(self.record['original_members']==members,'Original operation terms changed')
else:
self.record={'operation_id':self.operation,'original_members':members,'phase':'Prepared','created_at':time.time()};self.save()
require(self.record['phase']!='Released','Completed maintenance must not be reacquired')
if not recovery and not self.record.get('originals_validated'):
for member in members:
actual=self.inspect(member['name']);require(actual['Id']==member['container_id'] and actual['Image']==member['image_id'],'Original member changed')
bindings=actual['HostConfig'].get('PortBindings') or {}
if member['name']=='indeedhub':require(bindings=={'7777/tcp':[{'HostIp':'127.0.0.1','HostPort':'7778'}]},'Unsupported direct frontend exposure')
else:require(not bindings,'Unsupported direct writer exposure')
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','Original unit source missing')
require(source.stat().st_uid==os.getuid() and sha(source)==member['unit_sha256'],'Original unit source changed')
self.record['originals_validated']=True;self.save()
self.close_ingress()
if recovery:
runtime=json.loads((self.data/'update-transactions'/'supervised'/(self.operation+'.json')).read_text())
require(runtime.get('phase')=='Restoring' and type(runtime.get('target_startup_began')) is bool,'Durable explicit restoring obligation required')
self.record['phase']='Recovering';self.record['target_startup_began']=runtime['target_startup_began'];self.save()
return {'operation_id':self.operation,'state':'recovering'}
if not self.record.get('stopped',{}).get('indeedhub-api',{}).get('confirmed'):
self.pause_queue()
self.graceful_stop('indeedhub')
deadline=time.monotonic()+300
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
require(time.monotonic()<deadline,'Transcodes still active; retained job state, no forced completion')
time.sleep(1)
self.graceful_stop('indeedhub-ffmpeg');self.graceful_stop('indeedhub-api')
self.backup();self.verify();return {'operation_id':self.operation,'state':'drained'}
def verify(self):
self.holds();self.fence_matches()
if self.record and self.record.get('phase')=='Recovering':
runtime=json.loads((self.data/'update-transactions'/'supervised'/(self.operation+'.json')).read_text())
require(runtime.get('phase')=='Restoring','Recovery ownership changed')
return {'operation_id':self.operation,'state':'held'}
require(self.record and self.record.get('backup_complete'),'Drain not complete')
require(self.record['phase'] in ('Drained','Released'),'Invalid maintenance phase')
# Verification remains possible when API/storage endpoints are stopped.
# The native adapter separately validates target/original runtime identity.
for name in NAMES:require(self.record.get('stopped',{}).get(name,{}).get('confirmed'),'Original writer stop evidence missing')
for name,record in self.record['artifacts'].items():
path=self.root/'backup'/name;require(path.is_file() and not path.is_symlink() and path.stat().st_size==record['bytes'],'Backup artifact missing or changed')
return {'operation_id':self.operation,'state':'held'}
def release(self, outcome):
require(outcome in ('committed','restored','aborted'),'Invalid release outcome')
require(self.record is not None,'Unknown maintenance operation')
if self.record['phase']=='Released':
require(self.record.get('outcome')==outcome,'Maintenance outcome changed')
if self.fence.exists():
self.fence_matches();self.fence.unlink()
return {'operation_id':self.operation,'state':'released'}
self.holds();self.fence_matches()
runtime=json.loads((self.data/'update-transactions'/'supervised'/(self.operation+'.json')).read_text());phase=runtime['phase']
require(phase=={'committed':'Committed','restored':'Restored','aborted':'Restored'}[outcome],'Runtime outcome not durably verified')
# Restoring old runtime over a changed database is not sufficient to open
# admission. Node must verify same-schema/additive migration compatibility.
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')
else:
self.record['rollback_data_claim']='No target startup/migration began; only original runtime restored.'
if not self.record.get('queue_was_paused',True):
state=self.queue('resume');require(not state['paused'],'Could not restore queue admission')
self.record['phase']='Released';self.record['outcome']=outcome;self.save()
self.fence.unlink();fd=os.open(self.fence.parent,os.O_RDONLY);os.fsync(fd);os.close(fd)
return {'operation_id':self.operation,'state':'released'}
def main():
os.umask(0o077);require(len(sys.argv)==2 and sys.argv[1] in ('acquire','verify','release'),'Unsupported maintenance action')
raw=sys.stdin.buffer.read(65537);require(len(raw)<=65536,'Maintenance request too large');request=json.loads(raw)
allowed={'operation_id','original_members'} if sys.argv[1]=='acquire' else {'operation_id','outcome'} if sys.argv[1]=='release' else {'operation_id'}
require(set(request) in (allowed, allowed|{'recovery'}) if sys.argv[1]=='acquire' else set(request)==allowed,'Unexpected maintenance fields')
if 'recovery' in request:require(type(request['recovery']) is bool,'Invalid recovery flag')
require(os.getuid()==1000,'Expected node service user');fd=int(os.environ['ARCHY_UPDATE_LOCK_FD']);actual=os.fstat(fd);expected=(DATA/'update-transactions'/'lock').stat();require((actual.st_dev,actual.st_ino)==(expected.st_dev,expected.st_ino),'Inherited lifecycle lock is not the expected file')
controller=Controller(DATA,request['operation_id'],fd)
result=controller.acquire(request['original_members'],request.get('recovery',False)) if sys.argv[1]=='acquire' else controller.verify() if sys.argv[1]=='verify' else controller.release(request['outcome'])
encoded=json.dumps(result);require(len(encoded)<=4096,'Maintenance response exceeds bound');print(encoded)
if __name__=='__main__':
try:main()
except Exception as error:
# Command/env details remain in private journal, never RPC/UI stdout.
print(json.dumps({'error':'Maintenance remains held; inspect its private operation journal','reason':type(error).__name__}),file=sys.stderr);sys.exit(1)
@@ -0,0 +1,85 @@
"""Pure fixture coverage; imports the controller without invoking main/live tools."""
import importlib.util,json,pathlib,tempfile,unittest,uuid
MODULE=pathlib.Path(__file__).resolve().parents[2]/'scripts/indeehub-maintenance-controller.py'
spec=importlib.util.spec_from_file_location('maintenance_controller',MODULE);module=importlib.util.module_from_spec(spec);spec.loader.exec_module(module)
def members():
return [{'name':name,'container_id':f'{n+1:064x}','image_id':'a'*64,'unit_sha256':'b'*64,'config_sha256':'c'*64,'running':True} for n,name in enumerate(module.NAMES)]
class MaintenanceTests(unittest.TestCase):
def setUp(self):
self.tmp=tempfile.TemporaryDirectory();self.addCleanup(self.tmp.cleanup);self.root=pathlib.Path(self.tmp.name);self.operation=str(uuid.uuid4());self.calls=[]
holds=self.root/'update-transactions'/'holds';holds.mkdir(parents=True)
for name in module.NAMES:(holds/name).write_text(self.operation)
self.controller=module.Controller(self.root,self.operation,0,self.command_runner)
def command_runner(self,argv,timeout,output):
self.calls.append(argv)
if argv[:2]==['podman','inspect']:
member=next(m for m in members() if m['name']==argv[2]);return json.dumps([{'Id':member['container_id'],'Image':member['image_id'],'HostConfig':{'PortBindings':{}}}]).encode()
raise AssertionError('Unexpected fixture command '+str(argv))
def test_exact_seven_member_identity_is_required(self):
self.assertEqual(len(module.validate_members(members())),7)
for bad in [members()[:-1],members()+[members()[0]],[dict(members()[0],container_id='invalid'),*members()[1:]]]:
with self.assertRaises(RuntimeError):module.validate_members(bad)
def test_foreign_operation_cannot_release_or_overwrite_hold(self):
path=self.root/'update-transactions'/'holds'/'indeedhub';path.write_text(str(uuid.uuid4()))
with self.assertRaises(RuntimeError):self.controller.holds()
self.assertEqual(self.calls,[])
def test_failed_original_exposure_validation_is_not_skipped_on_retry(self):
for _ in range(2):
with self.assertRaisesRegex(RuntimeError,'frontend exposure'):self.controller.acquire(members())
self.controller=module.Controller(self.root,self.operation,0,self.command_runner)
self.assertFalse(self.controller.record.get('originals_validated',False))
self.assertEqual(len(self.calls),2)
def test_changed_original_terms_cannot_resume_saved_operation(self):
self.controller.record={'operation_id':self.operation,'original_members':module.validate_members(members()),'phase':'Prepared'};self.controller.save()
changed=members();changed[0]['config_sha256']='d'*64
with self.assertRaisesRegex(RuntimeError,'terms changed'):self.controller.acquire(changed)
self.assertEqual(self.calls,[])
def test_release_lost_reply_is_idempotent_but_never_removes_foreign_fence(self):
c=self.controller;c.record={'operation_id':self.operation,'phase':'Released','outcome':'committed'};c.save()
c.fence.parent.mkdir(parents=True);c.fence.write_text(self.operation)
self.assertEqual(c.release('committed')['state'],'released');self.assertFalse(c.fence.exists())
self.assertEqual(c.release('committed')['state'],'released')
c.fence.write_text(str(uuid.uuid4()))
with self.assertRaises(RuntimeError):c.release('committed')
self.assertTrue(c.fence.exists())
def test_backup_file_existence_never_substitutes_for_drain_evidence(self):
c=self.controller;c.record={'operation_id':self.operation,'phase':'Prepared','backup_complete':False};c.save()
c.fence.parent.mkdir(parents=True);c.fence.write_text(self.operation)
(c.root/'backup').mkdir();(c.root/'backup'/'database.dump').write_bytes(b'not evidence')
with self.assertRaisesRegex(RuntimeError,'Drain not complete'):c.verify()
def test_forced_original_exit_never_marks_writer_completed(self):
c=self.controller;c.record={'operation_id':self.operation,'phase':'Prepared','original_members':module.validate_members(members())};c.save()
def command(argv,timeout,output):
if argv[:2]==['podman','inspect']:return self.command_runner(argv,timeout,output)
if argv[:3]==['systemctl','--user','stop']:return b''
if argv[:3]==['systemctl','--user','show']:return b'ActiveState=inactive\nResult=success\n'
if argv[:2]==['podman','events']:return json.dumps({'ID':members()[0]['container_id'],'ContainerExitCode':137}).encode()
raise AssertionError(argv)
c.runner=command
with self.assertRaisesRegex(RuntimeError,'did not exit cleanly'):c.graceful_stop('indeedhub')
self.assertFalse(c.record['stopped']['indeedhub'].get('confirmed',False))
self.assertTrue((self.root/'update-transactions'/'holds'/'indeedhub').exists())
def test_failed_stop_can_restore_without_fabricating_completed_drain(self):
c=self.controller;c.record={'operation_id':self.operation,'original_members':module.validate_members(members()),'phase':'Prepared','stopped':{'indeedhub':{'intent_at':1,'exit_code':137}}};c.save()
runtime=c.data/'update-transactions'/'supervised'/(self.operation+'.json');module.atomic(runtime,{'phase':'Restoring','target_startup_began':False})
c.fence.parent.mkdir(parents=True);c.fence.write_text(self.operation)
c.close_ingress=lambda: c.fence_matches()
self.assertEqual(c.acquire(members(),recovery=True)['state'],'recovering')
self.assertFalse(c.record.get('backup_complete',False))
self.assertEqual(c.verify()['state'],'held')
module.atomic(runtime,{'phase':'Restored','target_startup_began':False})
self.assertEqual(c.release('aborted')['state'],'released')
self.assertIn('No target startup',c.record['rollback_data_claim'])
def test_target_started_rollback_requires_data_compatibility(self):
c=self.controller;c.record={'operation_id':self.operation,'phase':'Recovering'};c.save()
c.fence.parent.mkdir(parents=True);c.fence.write_text(self.operation)
module.atomic(c.data/'update-transactions'/'supervised'/(self.operation+'.json'),{'phase':'Restored','target_startup_began':True})
with self.assertRaisesRegex(RuntimeError,'Data compatibility'):c.release('restored')
self.assertTrue(c.fence.exists())
def test_every_legacy_sublocation_must_be_fenced(self):
guard='if (-f /var/lib/archipelago/app-maintenance/indeedhub) { return 503; }'
blocks=[f'location /app/indeedhub/{suffix} {{\n {guard}\n proxy_pass http://127.0.0.1:7778/;\n}}' for suffix in ('','_next/','ws/')]
self.assertEqual(module.validate_nginx_guards('\n'.join(blocks)),3)
with self.assertRaisesRegex(RuntimeError,'missing its maintenance guard'):module.validate_nginx_guards('\n'.join(blocks).replace(guard,'',1))
with self.assertRaisesRegex(RuntimeError,'Unrecognized direct'):module.validate_nginx_guards('\n'.join(blocks)+'\nlocation /other/ {\n proxy_pass http://127.0.0.1:7778/;\n}')
if __name__=='__main__':unittest.main()