Files
ARR-2.0-0918/tests/test_daily_report_retirement.py
T

345 lines
23 KiB
Python

"""Recoverable daily retirement: real isolated PostgreSQL, synthetic rows only."""
from concurrent.futures import ThreadPoolExecutor
from datetime import date
import hashlib
import json
import os
from pathlib import Path
import threading
import time
import unittest
from arr_ingestion.contracts import ArtifactRef, IngestionError
from arr_ingestion.postgres import DatabaseConfig, PostgresIngestionRepository
from arr_ingestion.repository import JobRegistration
from tests.local_postgres import TemporaryPostgres
PROJECT = Path(__file__).resolve().parents[1]
DAY = date(2026, 9, 16)
ACTOR = "fixture-user"
class RetirementMigrationContractTests(unittest.TestCase):
def test_additive_contract_and_no_fact_or_artifact_deletion_grants(self):
up = (PROJECT / "database/022_daily_report_retirement.sql").read_text()
down = (PROJECT / "database/022_daily_report_retirement.down.sql").read_text()
for text in ("current_database() <> 'booking_test'", "daily_run_retirements_immutable",
"monthly_runs_live_snapshot_unique", "WHERE report_status <> 'withdrawn'",
"'active', 'superseded', 'withdrawn'", "active non-retired daily version"):
self.assertIn(text, up)
self.assertIn("GRANT SELECT, INSERT ON ingestion.daily_run_retirements", up)
self.assertIn("GRANT DELETE ON finance.current_daily_versions", up)
self.assertNotIn("GRANT DELETE ON finance.daily_records", up)
self.assertNotIn("GRANT DELETE ON ingestion.artifacts", up)
self.assertIn("refusing rollback: daily retirement or withdrawn monthly audit exists", down)
@unittest.skipUnless(os.environ.get("ARR_TEST_LOCAL_POSTGRES") == "1", "owned PostgreSQL opt-in required")
class DailyRetirementPostgresTests(unittest.TestCase):
@classmethod
def setUpClass(cls):
cls.database = TemporaryPostgres().__enter__()
cls.addClassCleanup(cls.database.__exit__, None, None, None)
def setUp(self):
self.database.reset_database(schema_version=22)
self.repository = PostgresIngestionRepository(DatabaseConfig("owned-fixture"), connect=self.database.connect)
def sql(self, query, parameters=None):
with self.database.connect() as connection:
cursor = connection.execute(query, parameters)
return cursor.fetchall() if cursor.description else []
def seed(self, job, *, status="accepted", day=DAY, version=None, active=True):
sha = hashlib.sha256(job.encode()).hexdigest()
with self.database.connect() as connection:
source = connection.execute("""INSERT INTO ingestion.artifacts
(artifact_kind, storage_provider, bucket_alias, object_key, original_filename, sha256, byte_size, mime_type)
VALUES ('opera_xml','oss','arr-private',%s,'source.xml',%s,0,'application/xml') RETURNING id""",
(job + "/source.xml", sha)).fetchone()[0]
result = connection.execute("""INSERT INTO ingestion.artifacts
(artifact_kind, storage_provider, bucket_alias, object_key, original_filename, sha256, byte_size)
VALUES ('result_json','local_fixture','fixture',%s,'result.json',%s,0) RETURNING id""",
(job + "/result.json", sha)).fetchone()[0]
run = connection.execute("""INSERT INTO ingestion.processing_runs
(run_key,pipeline_type,source_artifact_id,result_artifact_id,run_status,business_date,
requested_processor_version,requested_rule_set_sha256,delivered_processor_version,
delivered_rule_set_sha256,result_schema_version,delivery_sha256,validated_at,finished_at,failure_code)
VALUES (%s,'opera_daily',%s,%s,%s,%s,'4.4.0',%s,'4.4.0',%s,'4.0',%s,now(),
CASE WHEN %s IN ('accepted','failed','rejected','cancelled') THEN now() END,
CASE WHEN %s IN ('failed','rejected') THEN 'FIXTURE_FAILED' END) RETURNING id""",
(job,source,result,status,day,sha,sha,sha,status,status)).fetchone()[0]
version_id = None
if version is not None:
outputs = []
for kind in ("daily_xlsx", "structured_result_json"):
outputs.append(connection.execute("""INSERT INTO ingestion.artifacts
(artifact_kind,storage_provider,bucket_alias,object_key,original_filename,sha256,byte_size)
VALUES (%s,'local_fixture','fixture',%s,%s,%s,0) RETURNING id""",
(kind,job + "/" + kind,kind,sha)).fetchone()[0])
version_id = connection.execute("""INSERT INTO finance.daily_versions
(business_date,version_no,processing_run_id,source_artifact_id,daily_report_artifact_id,
result_json_artifact_id,structured_result_artifact_id,version_status,processor_version,
rule_set_sha256,result_schema_version,result_sha256,source_rows,retained_rows,
excluded_rate_code_rows,duplicate_rows,validation_failed_rows,price_unmatched_rows,validated_at)
VALUES (%s,%s,%s,%s,%s,%s,%s,'validated','4.4.0',%s,'4.0',%s,1,1,0,0,0,0,now()) RETURNING id""",
(day,version,run,source,outputs[0],result,outputs[1],sha,sha)).fetchone()[0]
connection.execute("""INSERT INTO finance.daily_records
(daily_version_id,source_sequence,source_location,outcome,decision_codes,block_code,
adults,children,company_name,company_key,confirmation_no,disp_room_no,effective_rate_amount,
full_name,no_of_rooms,rate_code,arrival,departure,nights,real_price,total_price,channel_key,
pricing_method,booking_source_match_status)
VALUES (%s,1,'synthetic:1','retained',ARRAY['PRICE_REFERENCE_MATCHED'],'',
2,0,'FIXTURE','FIXTURE',%s,%s,10,'SYNTHETIC',2,'WHO2',%s,%s,2,10,40,'FIT',
'price_reference_exact','missing_group_code')""",
(version_id,job,job,day,day.replace(day=day.day+2)))
connection.execute("""INSERT INTO finance.daily_channel_metrics
(daily_version_id,channel_key,channel_order,row_count) VALUES (%s,'FIT',1,1)""",(version_id,))
if active:
with connection.cursor() as cursor:
self.repository._activate_daily_version(cursor,day,version_id)
if status == "awaiting_review":
delivery = connection.execute("""INSERT INTO ingestion.processing_deliveries
(delivery_key,processing_run_id,envelope_sha256,envelope_json,delivery_status,result_status,
processor_version,rule_set_sha256,result_schema_version,business_date)
VALUES (%s,%s,%s,'{}','recorded_review','review_required','4.4.0',%s,'4.0',%s) RETURNING id""",
(job + "-delivery",run,sha,sha,day)).fetchone()[0]
case = connection.execute("""INSERT INTO ingestion.daily_review_cases
(case_key,processing_run_id,initial_delivery_id,business_date,source_sha256,
processor_version,rule_set_sha256,review_version,case_status)
VALUES (%s,%s,%s,%s,%s,'4.4.0',%s,'1.0','open') RETURNING id""",
("dailyreview-" + sha[:32],run,delivery,day,sha,sha)).fetchone()[0]
connection.execute("""INSERT INTO ingestion.daily_review_items
(review_case_id,company_key,rate_code,effective_rate_amount,affected_records,affected_rooms,affected_room_nights)
VALUES (%s,'FIXTURE','WHO2',1000,1,1,2)""",(case,))
return run,version_id,source
def month(self, version_id, *, snapshot="f" * 64):
with self.database.connect() as connection:
report = connection.execute("""INSERT INTO reporting.monthly_runs
(period_start,as_of_date,version_no,source_snapshot_sha256,report_status,processor_version,
rule_set_sha256,result_schema_version,row_count,channel_count)
VALUES ('2026-09-01','2026-09-16',1,%s,'reserved','fixture',%s,'1.0',0,5) RETURNING id""",
(snapshot,"a" * 64)).fetchone()[0]
connection.execute("""INSERT INTO reporting.monthly_run_daily_versions
(report_id,business_date,daily_version_id) VALUES (%s,%s,%s)""",(report,DAY,version_id))
return report
def test_current_delete_removes_pointer_withdraws_month_and_keeps_original_data(self):
run,version,_ = self.seed("current",version=1)
report = self.month(version)
before = self.sql("SELECT id,object_key,sha256 FROM ingestion.artifacts ORDER BY id")
original_facts = self.sql("SELECT * FROM finance.daily_records ORDER BY id")
self.assertEqual(self.sql("SELECT count(*) FROM finance.v_active_daily_facts"),[(1,)])
receipt = self.repository.retire_daily_job("current",DAY,ACTOR)
self.assertEqual(receipt,{"job_id":"current","business_date":DAY.isoformat(),"is_retired":True,
"current_removed":True,"monthly_refresh_required":True})
self.assertEqual(self.sql("SELECT * FROM finance.current_daily_versions"),[])
self.assertEqual(self.sql("SELECT run_status FROM ingestion.processing_runs WHERE id=%s",(run,)),[("accepted",)])
self.assertEqual(self.sql("SELECT version_status FROM finance.daily_versions WHERE id=%s",(version,)),[("superseded",)])
self.assertEqual(before,self.sql("SELECT id,object_key,sha256 FROM ingestion.artifacts ORDER BY id"))
self.assertEqual(original_facts,self.sql("SELECT * FROM finance.daily_records ORDER BY id"))
self.assertEqual(self.sql("SELECT count(*) FROM finance.v_active_daily_facts"),[(0,)])
self.assertEqual(self.sql("SELECT report_status, withdrawn_at IS NOT NULL FROM reporting.monthly_runs WHERE id=%s",(report,)),[("withdrawn",True)])
event = self.sql("SELECT aggregate_type,event_type,payload FROM ingestion.outbox_events")[0]
self.assertEqual(event[:2],("daily_version","arr.daily_scope_changed"))
self.assertEqual(event[2]["period_start"],"2026-09-01")
self.assertEqual(self.repository.retire_daily_job("current",DAY,"another-user"),receipt)
self.assertEqual(self.sql("SELECT actor_username FROM ingestion.daily_run_retirements"),[(ACTOR,)])
self.assertEqual(self.sql("SELECT count(*) FROM ingestion.outbox_events"),[(1,)])
def test_old_history_delete_keeps_current_and_month(self):
_,old,_ = self.seed("old",version=1)
_,new,_ = self.seed("new",version=2)
report = self.month(new)
receipt = self.repository.retire_daily_job("old",DAY,ACTOR)
self.assertFalse(receipt["current_removed"])
self.assertEqual(self.sql("SELECT daily_version_id FROM finance.current_daily_versions"),[(new,)])
self.assertEqual(self.sql("SELECT report_status FROM reporting.monthly_runs WHERE id=%s",(report,)),[("reserved",)])
self.assertEqual(self.sql("SELECT count(*) FROM ingestion.outbox_events"),[(0,)])
self.assertEqual(self.sql("SELECT count(*) FROM finance.daily_versions WHERE id=%s",(old,)),[(1,)])
def test_pending_price_review_deletes_only_that_task_and_audits_cancellation(self):
_,current,_ = self.seed("current",version=1)
run,_,_ = self.seed("pending",status="awaiting_review")
report = self.month(current)
self.assertFalse(self.repository.retire_daily_job("pending",DAY,ACTOR)["monthly_refresh_required"])
self.assertEqual(self.sql("SELECT run_status FROM ingestion.processing_runs WHERE id=%s",(run,)),[("cancelled",)])
self.assertEqual(self.sql("SELECT case_status,revision FROM ingestion.daily_review_cases"),[("cancelled",1)])
self.assertEqual(self.sql("SELECT event_type,actor_username FROM ingestion.daily_review_events"),[("PRICE_REVIEW_CANCELLED",ACTOR)])
self.assertEqual(self.sql("SELECT daily_version_id FROM finance.current_daily_versions"),[(current,)])
self.assertEqual(self.sql("SELECT report_status FROM reporting.monthly_runs WHERE id=%s",(report,)),[("reserved",)])
self.assertEqual(self.sql("SELECT count(*) FROM ingestion.daily_review_items"),[(1,)])
def test_running_and_generating_tasks_are_blocked_without_changes(self):
for status in ("received","queued","running","validating"):
self.seed(status,status=status)
with self.subTest(status=status), self.assertRaises(IngestionError) as raised:
self.repository.retire_daily_job(status,DAY,ACTOR)
self.assertEqual(raised.exception.code,"DAILY_JOB_BUSY")
self.seed("generating",status="awaiting_review")
self.sql("UPDATE ingestion.daily_review_cases SET case_status='processing'")
with self.assertRaises(IngestionError) as raised:
self.repository.retire_daily_job("generating",DAY,ACTOR)
self.assertEqual(raised.exception.code,"DAILY_JOB_BUSY")
self.assertEqual(self.sql("SELECT count(*) FROM ingestion.daily_run_retirements"),[(0,)])
def test_stale_date_unknown_job_and_invalid_identity_fail_closed(self):
self.seed("current",version=1)
for job,day,code in (("current",date(2026,9,17),"DAILY_DATE_CONFLICT"),
("missing",DAY,"DAILY_JOB_NOT_FOUND"),
("bad/identity",DAY,"DAILY_JOB_NOT_FOUND"),
("current",DAY.isoformat(),"DAILY_DATE_CONFLICT")):
with self.subTest(job=job,day=day),self.assertRaises(IngestionError) as raised:
self.repository.retire_daily_job(job,day,ACTOR)
self.assertEqual(raised.exception.code,code)
self.assertEqual(self.sql("SELECT count(*) FROM ingestion.daily_run_retirements"),[(0,)])
def test_failed_and_cancelled_terminal_records_without_versions_can_be_deleted(self):
for status in ("failed","rejected","cancelled"):
self.seed(status,status=status)
self.assertFalse(self.repository.retire_daily_job(status,DAY,ACTOR)["current_removed"])
self.assertEqual(self.sql("SELECT count(*) FROM ingestion.daily_run_retirements"),[(3,)])
self.assertEqual(self.sql("SELECT count(*) FROM ingestion.outbox_events"),[(0,)])
def test_concurrent_repeated_delete_has_one_audit_and_outbox(self):
self.seed("current",version=1)
with ThreadPoolExecutor(max_workers=2) as pool:
results = list(pool.map(lambda _: self.repository.retire_daily_job("current",DAY,ACTOR),range(2)))
self.assertEqual(results[0],results[1])
self.assertEqual(self.sql("SELECT count(*) FROM ingestion.daily_run_retirements"),[(1,)])
self.assertEqual(self.sql("SELECT count(*) FROM ingestion.outbox_events"),[(1,)])
def test_replacement_racing_old_delete_keeps_new_current_without_resurrection(self):
self.seed("old",version=1)
_,new,_ = self.seed("new",version=2,active=False)
barrier = threading.Barrier(2)
def activate():
barrier.wait()
def operation(cursor):
cursor.execute("SELECT pg_advisory_xact_lock(hashtextextended(%s,0))",("finance-daily:"+DAY.isoformat(),))
self.repository._activate_daily_version(cursor,DAY,new)
self.repository._run_transaction(operation,"fixture activation failed")
def retire():
barrier.wait()
return self.repository.retire_daily_job("old",DAY,ACTOR)
with ThreadPoolExecutor(max_workers=2) as pool:
a,b = pool.submit(activate),pool.submit(retire)
a.result(); b.result()
self.assertEqual(self.sql("SELECT daily_version_id FROM finance.current_daily_versions"),[(new,)])
self.assertEqual(self.sql("SELECT version_status FROM finance.daily_versions WHERE id=%s",(new,)),[("active",)])
def test_outbox_failure_rolls_back_current_removal_month_withdrawal_and_audit(self):
_,version,_ = self.seed("current",version=1)
self.month(version)
self.sql("""CREATE FUNCTION ingestion.fail_retirement_event() RETURNS trigger LANGUAGE plpgsql AS $$
BEGIN RAISE EXCEPTION 'synthetic event failure'; END $$;
CREATE TRIGGER fixture_reject_retirement_event BEFORE INSERT ON ingestion.outbox_events
FOR EACH ROW EXECUTE FUNCTION ingestion.fail_retirement_event()""")
with self.assertRaises(IngestionError):
self.repository.retire_daily_job("current",DAY,ACTOR)
self.assertEqual(self.sql("SELECT daily_version_id FROM finance.current_daily_versions"),[(version,)])
self.assertEqual(self.sql("SELECT version_status FROM finance.daily_versions"),[("active",)])
self.assertEqual(self.sql("SELECT report_status FROM reporting.monthly_runs"),[("reserved",)])
self.assertEqual(self.sql("SELECT count(*) FROM ingestion.daily_run_retirements"),[(0,)])
def test_daily_activation_waits_for_monthly_publication_scope_lock(self):
self.seed("old",version=1)
_,new,_ = self.seed("new",version=2,active=False)
held = self.database.connect()
self.addCleanup(held.close)
held.execute("SELECT pg_advisory_xact_lock(hashtextextended(%s,0))",
("monthly_channel:2026-09-01:2026-09-30:*",))
def activate():
self.repository._run_transaction(
lambda cursor: self.repository._activate_daily_version(cursor,DAY,new),
"fixture activation failed")
with ThreadPoolExecutor(max_workers=1) as pool:
future = pool.submit(activate)
try:
deadline = time.monotonic()+5
while self.sql("SELECT count(*) FROM pg_locks WHERE locktype='advisory' AND NOT granted")[0][0] == 0:
self.assertLess(time.monotonic(),deadline,"activation did not wait for month lock")
time.sleep(.01)
self.assertFalse(future.done())
finally:
held.commit()
future.result(timeout=10)
self.assertEqual(self.sql("SELECT daily_version_id FROM finance.current_daily_versions"),[(new,)])
def test_retirement_audit_and_withdrawn_lineage_are_immutable_and_old_version_cannot_be_current(self):
import psycopg
_,version,_ = self.seed("current",version=1)
report = self.month(version)
self.repository.retire_daily_job("current",DAY,ACTOR)
for query in ("UPDATE ingestion.daily_run_retirements SET actor_username='changed'",
"DELETE FROM ingestion.daily_run_retirements",
"DELETE FROM reporting.monthly_run_daily_versions WHERE report_id="+str(report)):
with self.subTest(query=query),self.assertRaises(psycopg.Error):
self.sql(query)
# Even an accidental lifecycle update cannot repoint a retired version.
self.sql("UPDATE finance.daily_versions SET version_status='active' WHERE id=%s",(version,))
with self.assertRaises(psycopg.Error):
self.sql("INSERT INTO finance.current_daily_versions (business_date,daily_version_id) VALUES (%s,%s)",(DAY,version))
def test_same_retired_job_registration_replay_is_rejected_but_new_upload_identity_is_allowed(self):
_,_,source = self.seed("current",version=1)
sha = hashlib.sha256(b"current").hexdigest()
self.repository.retire_daily_job("current",DAY,ACTOR)
source_ref = ArtifactRef("source_xml","opera_xml","current/source.xml","source.xml",sha,0,"application/xml")
registration = JobRegistration("current",source_ref,"4.4.0",sha,1,"b"*64,None)
with self.assertRaises(IngestionError) as raised:
self.repository.register_job(registration)
self.assertEqual(raised.exception.code,"JOB_RETIRED")
self.repository.register_job(JobRegistration("new-upload",source_ref,"4.4.0",sha,1,"c"*64,None))
self.assertEqual(self.sql("SELECT run_status FROM ingestion.processing_runs WHERE run_key='new-upload'"),[("queued",)])
def test_empty_rollback_reapply_and_rollback_refusal_after_audit(self):
import psycopg
with self.database.connect(autocommit=True) as connection:
connection.execute((PROJECT/"database/022_daily_report_retirement.down.sql").read_text(),prepare=False)
connection.execute((PROJECT/"database/022_daily_report_retirement.sql").read_text(),prepare=False)
self.seed("current",version=1)
self.repository.retire_daily_job("current",DAY,ACTOR)
with self.database.connect(autocommit=True) as connection:
with self.assertRaises(psycopg.Error):
connection.execute((PROJECT/"database/022_daily_report_retirement.down.sql").read_text(),prepare=False)
connection.execute("ROLLBACK")
self.assertEqual(self.sql("SELECT count(*) FROM ingestion.daily_run_retirements"),[(1,)])
def test_upgrade_preserves_existing_history_and_grants_only_pointer_delete(self):
self.database.reset_database(schema_version=21)
run,version,_ = self.seed("existing",version=1)
before = self.sql("SELECT id,business_date,version_status FROM finance.daily_versions")
with self.database.connect(autocommit=True) as connection:
connection.execute("DO $$ BEGIN IF to_regrole('arr_app') IS NULL THEN CREATE ROLE arr_app; END IF; END $$")
connection.execute((PROJECT/"database/022_daily_report_retirement.sql").read_text(),prepare=False)
self.assertEqual(self.sql("SELECT id,business_date,version_status FROM finance.daily_versions"),before)
self.assertEqual(self.sql("SELECT daily_version_id FROM finance.current_daily_versions"),[(version,)])
self.assertEqual(self.sql("""SELECT has_table_privilege('arr_app','ingestion.daily_run_retirements','SELECT'),
has_table_privilege('arr_app','ingestion.daily_run_retirements','INSERT'),
has_table_privilege('arr_app','ingestion.daily_run_retirements','DELETE'),
has_table_privilege('arr_app','finance.current_daily_versions','DELETE'),
has_table_privilege('arr_app','finance.daily_records','DELETE'),
has_table_privilege('arr_app','ingestion.artifacts','DELETE')"""),[(True,True,False,True,False,False)])
def test_daily_lifecycle_gate_rejects_021_then_accepts_complete_022_and_rejects_missing_index(self):
self.database.reset_database(schema_version=21)
with self.assertRaises(IngestionError) as missing:
self.repository.assert_daily_lifecycle_schema()
self.assertEqual(missing.exception.code,"DATABASE_MIGRATION_MISSING")
with self.database.connect(autocommit=True) as connection:
connection.execute((PROJECT/"database/022_daily_report_retirement.sql").read_text(),prepare=False)
self.repository.assert_daily_lifecycle_schema()
self.sql("DROP INDEX reporting.monthly_runs_live_snapshot_unique")
with self.assertRaises(IngestionError) as missing:
self.repository.assert_daily_lifecycle_schema()
self.assertEqual(missing.exception.code,"DATABASE_MIGRATION_MISSING")
def test_daily_lifecycle_gate_rejects_disabled_retirement_audit_guard(self):
self.repository.assert_daily_lifecycle_schema()
self.sql("ALTER TABLE ingestion.daily_run_retirements DISABLE TRIGGER daily_run_retirements_immutable")
with self.assertRaises(IngestionError) as missing:
self.repository.assert_daily_lifecycle_schema()
self.assertEqual(missing.exception.code,"DATABASE_MIGRATION_MISSING")