218 lines
12 KiB
Python
218 lines
12 KiB
Python
"""Owned disposable PostgreSQL acceptance for monthly retirement and races.
|
|
|
|
No existing database, live source, production configuration or hotel data is used.
|
|
"""
|
|
from concurrent.futures import ThreadPoolExecutor
|
|
from datetime import date, timedelta
|
|
import os
|
|
from pathlib import Path
|
|
import tempfile
|
|
import threading
|
|
import unittest
|
|
|
|
from arr_ingestion.postgres import DatabaseConfig as IngestionConfig, PostgresIngestionRepository
|
|
from arr_ingestion.service import IngestionService
|
|
from arr_ingestion.validation import DeliveryValidator
|
|
from arr_processing.local import LocalDailyProcessor
|
|
from arr_processing.policy import load_processor_policy
|
|
from arr_storage.filesystem import FilesystemObjectBackend
|
|
from arr_storage.store import ManagedObjectStore
|
|
from arr_web.programmatic import ProgrammaticUploadCoordinator
|
|
from monthly_reports.contracts import ErrorCode
|
|
from monthly_reports.core import build_monthly_report
|
|
from monthly_reports.publishing import AtomicReportPublisher, OpenpyxlWorkbookBuilder
|
|
from monthly_reports.repository import DatabaseConfig, PostgresReportRepository
|
|
from monthly_reports.service import MonthlyReportService, RunRequest
|
|
from monthly_reports.worker import MonthlyOutboxWorker, PostgresOutboxRepository, WorkerError
|
|
from tests.local_postgres import TemporaryPostgres
|
|
from tests.test_arr_opera_daily_ingest import reservation, xml_document
|
|
|
|
|
|
PROJECT = Path(__file__).resolve().parents[1]
|
|
|
|
|
|
class PausedBuilder:
|
|
def __init__(self):
|
|
self.entered, self.release = threading.Event(), threading.Event()
|
|
self.delegate = OpenpyxlWorkbookBuilder()
|
|
|
|
def build(self, report, work_dir):
|
|
self.entered.set()
|
|
if not self.release.wait(15):
|
|
raise RuntimeError("synthetic monthly builder wait expired")
|
|
return self.delegate.build(report, work_dir)
|
|
|
|
|
|
@unittest.skipUnless(os.environ.get("ARR_TEST_LOCAL_POSTGRES") == "1", "owned disposable PostgreSQL opt-in required")
|
|
class MonthlyLifecyclePostgresTests(unittest.TestCase):
|
|
@classmethod
|
|
def setUpClass(cls):
|
|
cls.database = TemporaryPostgres().__enter__()
|
|
cls.addClassCleanup(cls.database.__exit__, None, None, None)
|
|
cls.policy = load_processor_policy(PROJECT)
|
|
|
|
def setUp(self):
|
|
self.database.reset_database(schema_version=22)
|
|
temporary = tempfile.TemporaryDirectory(prefix="arr-monthly-lifecycle-")
|
|
self.addCleanup(temporary.cleanup)
|
|
self.root = Path(temporary.name)
|
|
self.store = ManagedObjectStore(FilesystemObjectBackend(self.root / "objects", create=True))
|
|
self.ingestion = PostgresIngestionRepository(IngestionConfig("owned-fixture"), connect=self.database.connect)
|
|
service = IngestionService(DeliveryValidator(self.store, self.policy), self.ingestion)
|
|
self.uploads = ProgrammaticUploadCoordinator(self.store, self.ingestion, service, LocalDailyProcessor(self.policy),
|
|
self.policy.processor_version, self.policy.rule_set_sha256)
|
|
config = DatabaseConfig("owned-fixture")
|
|
self.reports = PostgresReportRepository(config, connect=self.database.connect)
|
|
self.outbox = PostgresOutboxRepository(config, connect=self.database.connect)
|
|
self.service = self.monthly_service()
|
|
self.worker = MonthlyOutboxWorker(self.outbox, self.reports, self.service)
|
|
|
|
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 monthly_service(self, builder=None):
|
|
output = self.root / "outputs/monthly"
|
|
return MonthlyReportService(self.reports, builder or OpenpyxlWorkbookBuilder(),
|
|
AtomicReportPublisher(self.root, output, object_store=self.store), output / ".staging")
|
|
|
|
def accepted(self, day):
|
|
xml = xml_document(reservation(1, confirmation=f"SYN-{day.isoformat()}", room=f"SYN-ROOM-{day.day}"))
|
|
xml = (xml.replace("2026-07-27", day.isoformat())
|
|
.replace("2026-07-28", (day + timedelta(days=1)).isoformat())
|
|
.replace("20260727", day.strftime("%Y%m%d"))
|
|
.replace("27-07-26", day.strftime("%d-%m-%y")))
|
|
result = self.uploads.submit(f"synthetic-{day.isoformat()}.XML", xml.encode())
|
|
self.assertEqual(result["status"], "succeeded", result)
|
|
return result
|
|
|
|
def retire(self, uploaded):
|
|
return self.ingestion.retire_daily_job(uploaded["job_id"], date.fromisoformat(uploaded["business_date"]), "synthetic-operator")
|
|
|
|
def drain(self):
|
|
outcomes = []
|
|
for _attempt in range(20):
|
|
outcome = self.worker.process_next()
|
|
if outcome.status == "idle":
|
|
return outcomes
|
|
self.assertIn(outcome.status, {"published", "withdrawn"}, outcome)
|
|
outcomes.append(outcome)
|
|
self.fail("synthetic monthly outbox did not settle")
|
|
|
|
def active(self):
|
|
return self.sql("SELECT id, as_of_date, version_no, row_count FROM reporting.monthly_runs WHERE report_status='active'")
|
|
|
|
def test_last_day_deletion_withdraws_all_old_publications_and_acknowledges_without_empty_report(self):
|
|
uploaded = self.accepted(date(2026, 7, 2))
|
|
self.drain()
|
|
published = self.active()[0]
|
|
count_before = self.sql("SELECT count(*) FROM reporting.monthly_runs")[0][0]
|
|
self.assertTrue(self.retire(uploaded)["current_removed"])
|
|
self.assertEqual(self.active(), [])
|
|
outcomes = self.drain()
|
|
self.assertEqual([outcome.status for outcome in outcomes], ["withdrawn"])
|
|
self.assertEqual(self.sql("SELECT count(*) FROM reporting.monthly_runs")[0][0], count_before)
|
|
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", (published[0],)),
|
|
[("withdrawn", True)])
|
|
self.assertEqual(self.sql("SELECT publish_status FROM ingestion.outbox_events WHERE event_type='arr.daily_scope_changed'"), [("published",)])
|
|
self.assertEqual(self.sql("SELECT count(*) FROM finance.daily_records WHERE outcome='retained'"), [(1,)])
|
|
|
|
def test_deleting_middle_day_rebuilds_only_remaining_facts_at_same_watermark(self):
|
|
first = self.accepted(date(2026, 7, 2))
|
|
middle = self.accepted(date(2026, 7, 3))
|
|
last = self.accepted(date(2026, 7, 4))
|
|
self.drain()
|
|
self.assertEqual(self.active()[0][1:], (date(2026, 7, 4), 1, 3))
|
|
self.retire(middle)
|
|
self.assertEqual(self.active(), [])
|
|
self.drain()
|
|
active = self.active()[0]
|
|
self.assertEqual(active[1], date(2026, 7, 4))
|
|
self.assertEqual(active[3], 2)
|
|
self.assertEqual(self.sql("SELECT daily_version_id FROM reporting.monthly_run_daily_versions WHERE report_id=%s ORDER BY business_date", (active[0],)),
|
|
[(first["daily_version_id"],), (last["daily_version_id"],)])
|
|
|
|
def test_deleting_every_day_withdraws_the_month_and_preserves_published_lineage(self):
|
|
first = self.accepted(date(2026, 7, 2))
|
|
self.drain()
|
|
second = self.accepted(date(2026, 7, 3))
|
|
self.drain()
|
|
published_ids = [row[0] for row in self.sql("SELECT id FROM reporting.monthly_runs ORDER BY id")]
|
|
self.retire(first)
|
|
self.retire(second)
|
|
self.assertEqual([outcome.status for outcome in self.drain()], ["withdrawn", "withdrawn"])
|
|
self.assertEqual(self.active(), [])
|
|
self.assertEqual(self.sql("SELECT count(*) FROM finance.current_daily_versions"), [(0,)])
|
|
self.assertEqual(self.sql("SELECT count(*) FROM finance.v_active_daily_facts"), [(0,)])
|
|
self.assertEqual(self.sql("SELECT count(*) FROM reporting.monthly_runs WHERE report_status='withdrawn'"), [(2,)])
|
|
import psycopg
|
|
with self.assertRaisesRegex(psycopg.errors.RaiseException, "lineage and manifest are immutable"):
|
|
self.sql("DELETE FROM reporting.monthly_run_daily_versions WHERE report_id=%s", (published_ids[0],))
|
|
self.assertEqual(self.sql("SELECT count(*) FROM reporting.monthly_run_daily_versions WHERE report_id=%s", (published_ids[0],)), [(1,)])
|
|
|
|
def test_deleting_latest_day_creates_new_A_snapshot_without_resurrecting_earlier_A_publication(self):
|
|
first = self.accepted(date(2026, 7, 2))
|
|
self.drain()
|
|
original = self.active()[0]
|
|
second = self.accepted(date(2026, 7, 3))
|
|
self.drain()
|
|
self.retire(second)
|
|
self.drain()
|
|
remaining = self.active()[0]
|
|
self.assertEqual(remaining[1:], (date(2026, 7, 2), 3, 1))
|
|
self.assertNotEqual(remaining[0], original[0])
|
|
self.assertEqual(self.sql("SELECT source_snapshot_sha256 FROM reporting.monthly_runs WHERE id IN (%s,%s) ORDER BY id", (original[0], remaining[0]))[0],
|
|
self.sql("SELECT source_snapshot_sha256 FROM reporting.monthly_runs WHERE id=%s", (remaining[0],))[0])
|
|
self.assertEqual(self.sql("SELECT daily_version_id FROM finance.current_daily_versions"), [(first["daily_version_id"],)])
|
|
self.assertEqual(self.sql("SELECT report_status FROM reporting.monthly_runs WHERE id=%s", (original[0],)), [("withdrawn",)])
|
|
|
|
def test_reserved_worker_cannot_publish_a_day_deleted_while_it_was_building(self):
|
|
uploaded = self.accepted(date(2026, 7, 2))
|
|
paused = PausedBuilder()
|
|
service = self.monthly_service(paused)
|
|
with ThreadPoolExecutor(max_workers=1) as pool:
|
|
result = pool.submit(service.run, RunRequest(2026, 7, date(2026, 7, 2)))
|
|
try:
|
|
self.assertTrue(paused.entered.wait(10))
|
|
self.retire(uploaded)
|
|
finally:
|
|
paused.release.set()
|
|
final = result.result(timeout=20)
|
|
self.assertEqual((final.status, final.error_code), ("failed", ErrorCode.SOURCE_SNAPSHOT_STALE))
|
|
self.assertEqual(self.active(), [])
|
|
self.assertEqual(self.sql("SELECT report_status FROM reporting.monthly_runs"), [("withdrawn",)])
|
|
self.assertTrue(all(outcome.status == "withdrawn" for outcome in self.drain()))
|
|
|
|
def test_old_cutoff_worker_cannot_replace_a_later_successful_publication(self):
|
|
self.accepted(date(2026, 7, 2))
|
|
paused = PausedBuilder()
|
|
service = self.monthly_service(paused)
|
|
with ThreadPoolExecutor(max_workers=1) as pool:
|
|
result = pool.submit(service.run, RunRequest(2026, 7, date(2026, 7, 2)))
|
|
try:
|
|
self.assertTrue(paused.entered.wait(10))
|
|
self.accepted(date(2026, 7, 3))
|
|
self.drain()
|
|
later_id = self.active()[0][0]
|
|
finally:
|
|
paused.release.set()
|
|
final = result.result(timeout=20)
|
|
self.assertEqual((final.status, final.error_code), ("failed", ErrorCode.SOURCE_SNAPSHOT_STALE))
|
|
self.assertEqual(self.active()[0][0], later_id)
|
|
self.assertEqual(self.active()[0][1], date(2026, 7, 3))
|
|
|
|
def test_empty_ack_refuses_a_month_that_still_has_active_facts(self):
|
|
self.accepted(date(2026, 7, 2))
|
|
event = self.outbox.claim_next()
|
|
self.assertIsNotNone(event)
|
|
with self.assertRaises(WorkerError) as raised:
|
|
self.outbox.mark_withdrawn(event.event_id, date(2026, 7, 1))
|
|
self.assertEqual(raised.exception.code, "MONTHLY_WORKER_WITHDRAW_ACK_FAILED")
|
|
self.assertEqual(self.sql("SELECT publish_status FROM ingestion.outbox_events WHERE id=%s", (event.event_id,)), [("publishing",)])
|
|
|
|
|
|
if __name__ == "__main__":
|
|
unittest.main()
|