"""Real producer adapters and outbox retries against an actual local receiver.""" import base64 import copy import hashlib import json from pathlib import Path import subprocess import sys import tempfile import threading import unittest from unittest.mock import patch from wsgiref.simple_server import make_server, WSGIRequestHandler sys.path.insert(0,str(Path(__file__).resolve().parents[1]/'scripts')) import native_factory_acceptance as native class Contracts(unittest.TestCase): def test_jobs_do_not_read_operator_credentials_and_refuse_additive_egress(self): result=subprocess.CompletedProcess([],0,b'{"items":[{"metadata":{"name":"existing-allow"}}]}',b'') with patch.object(native,'receiver_check',return_value={}),patch.object(native,'command',return_value=result) as cmd,patch.object(native,'bao',side_effect=AssertionError('no OpenBao call permitted')): with self.assertRaisesRegex(native.LaneError,'existing_egress_policy_requires_review'): native.jobs({},['kubectl'],{},lambda:None) self.assertEqual(cmd.call_count,1) self.assertNotIn('create',cmd.call_args.args[0]) def test_existing_independent_reader_must_be_unambiguous_and_read_only(self): reader={'name':'operator','may_read':True,'may_write':False,'tenants':['*']} self.assertEqual(native.operator_identity([reader]),reader) for rows in [[],[dict(reader,may_write=True)],[dict(reader,tenants=['tenant:other'])], [reader,dict(reader,name='second-operator')]]: with self.assertRaises(native.LaneError):native.operator_identity(rows) def test_packet_cannot_redirect_secret_mount_or_add_network(self): with tempfile.TemporaryDirectory() as tmp: path=Path(tmp)/'packet.json';native.prepare(Path('/home/worsch'),path) sha=lambda:hashlib.sha256(path.read_bytes()).hexdigest() packet=native.validate_packet(path,sha()) job=packet['objects'][2] job['spec']['template']['spec']['volumes'][1]['secret']['secretName']='unrelated-secret' path.write_text(json.dumps(packet)) with self.assertRaisesRegex(native.LaneError,'resource_scope_changed'):native.validate_packet(path,sha()) def test_changed_probe_bytes_or_packet_digest_refuse(self): with tempfile.TemporaryDirectory() as tmp: path=Path(tmp)/'packet.json';meta=native.prepare(Path('/home/worsch'),path) path.write_bytes(path.read_bytes()+b'\n') with self.assertRaisesRegex(native.LaneError,'packet_digest_changed'): native.validate_packet(path,meta['sha256']) class Quiet(WSGIRequestHandler): def log_message(self,*args):pass class LocalReceiver(unittest.TestCase): def test_two_real_outboxes_retry_lost_receipts_without_duplicate_records(self): sys.path.insert(0,'/home/worsch/audit-core') from audit_core.ingestion import IngestionApplication from audit_core.sqlite_backend import SQLiteAuditBackend from audit_core.senders import SenderRegistry with tempfile.TemporaryDirectory() as tmp: root=Path(tmp);packet_path=root/'packet.json';native.prepare(Path('/home/worsch'),packet_path) packet=json.loads(packet_path.read_text());tokens={s:'local-fixture-'+s for s in native.SOURCES} registry=SenderRegistry.from_env({'AUDIT_CORE_SENDERS':json.dumps([ {'name':s,'tokens':[t],'sources':[s],'tenants':['tenant:platform'],'may_write':True, 'may_read':False,'evidence_kind':'load-bearing','secret_policy':'redact'} for s,t in tokens.items()])}) backend=SQLiteAuditBackend(str(root/'receiver.db')) server=make_server('127.0.0.1',0,IngestionApplication(backend,registry),handler_class=Quiet) thread=threading.Thread(target=server.serve_forever,daemon=True);thread.start() try: for cm in packet['objects'][::3]: sender=cm['metadata']['namespace'];d=root/sender;d.mkdir();d.chmod(0o755) for name,body in cm['data'].items():(d/name).write_text(body) (d/'source.zip').write_bytes(base64.b64decode(cm['binaryData']['source.zip'])) config=json.loads((d/'config.json').read_text());config['origin']='http://127.0.0.1:'+str(server.server_port) (d/'config.json').write_text(json.dumps(config));token=d/'token';token.write_text(tokens[sender]);token.chmod(0o444) argv=['docker','run','--rm','--network','host','--read-only','--cap-drop=ALL', '--security-opt=no-new-privileges','--memory=192m', '--tmpfs','/state:rw,nosuid,nodev,size=32m,mode=1777', '--mount','type=bind,src='+str(d)+',dst=/probe,readonly', '--mount','type=bind,src='+str(token)+',dst=/credential/token,readonly', '--entrypoint','python',native.IMAGE,'-I','/probe/producer.py'] result=subprocess.run(argv,capture_output=True,timeout=60) # Only bounded probe JSON may enter test diagnostics. try:receipt=json.loads(result.stdout) except ValueError:receipt={'status':'container_receipt_missing','exit':result.returncode} self.assertEqual(receipt['status'],'passed',receipt) self.assertEqual(receipt['http_statuses'],[202,200]) self.assertEqual(receipt['negatives']['read_routes_denied'],7) self.assertTrue(receipt['process_restart_retry']) record=backend.get(config['event_id']);self.assertEqual(record['source'],sender) print(json.dumps({'local_sender':sender,'receipt':receipt,'stored_record_keys':list(record), 'stored_data_keys':list(record['details']['data'])})) self.assertEqual(backend.verify_chain().events,2) self.assertTrue(backend.verify_chain().intact) finally: server.shutdown();server.server_close();thread.join(timeout=5) if __name__=='__main__':unittest.main()