Manage daily reports by date and refresh monthly reports after deletion
This commit is contained in:
1 parent
976a4fa7b2
commit
2a54838d59
41 files changed
+2757
-89
No files matched your search
@@ -0,0 +1,217 @@
|
||||
"""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()
|
||||
Reference in new issue
Block a user