from datetime import UTC, datetime from production_analytics.enlyze.gateway import EnlyzeDowntime from production_analytics.erp import CurrentWorkplaceStatus from production_analytics.service.downtime import ( DowntimeCategory, DowntimeReconciliationService, ProductionOrderBoundary, ) NOW = datetime(2026, 9, 16, 12, tzinfo=UTC) def source(*, identifier="d", end=NOW, category=None): return EnlyzeDowntime( identifier, "machine", "THRESHOLD", datetime(2026, 9, 16, 10, tzinfo=UTC), end, None, "reason" if category else None, "reason" if category else None, None, None, category, None, ) def status(order="1", *, remaining=3, good=7, target=10, feedback=NOW): return CurrentWorkplaceStatus( "WP", order, None, None, feedback, target, good, remaining, None, None ) class Gateway: def __init__(self, events): self.events = events def get_downtimes(self, machine_id, *, start=None): return self.events class Erp: def __init__(self, current): self.current = current def get_current_workplace_status(self, workplace): return self.current class Repo: def __init__(self): self.boundary = None self.events = {} self.closures = [] def current_boundary(self, machine_id): return self.boundary def save_boundary(self, boundary): self.boundary = boundary def close_order_attribution(self, machine_id, production_order, ended_at): self.closures.append((production_order, ended_at)) for event in self.events.values(): if event.production_order == production_order and event.attributed_end is None: event # database-specific clipping is covered by SQL contract def upsert(self, event): self.events[event.external_id] = event def service(repo, events, current): return DowntimeReconciliationService( Gateway(events), Erp(current), repo, machine_id="machine", workplace="WP", production_order_format="FA-{production_order}", ) def test_categories_open_and_reconciliation_upsert() -> None: repo = Repo() runner = service( repo, [ source(identifier="p", category="PLANNED"), source(identifier="u", category="UNPLANNED"), source(identifier="x", end=None), ], status(), ) assert runner.reconcile_once(NOW) == 3 assert {key: event.category for key, event in repo.events.items()} == { "p": DowntimeCategory.PLANNED, "u": DowntimeCategory.UNPLANNED, "x": DowntimeCategory.UNKNOWN, } assert repo.events["x"].source_end is None # Same UUID is overwritten, never duplicated; delayed classification is accepted. runner.gateway.events = [source(identifier="x", end=NOW, category="PLANNED")] runner.reconcile_once(NOW) assert len(repo.events) == 3 assert repo.events["x"].category is DowntimeCategory.PLANNED assert repo.events["x"].source_end == NOW def test_completion_clips_attribution_and_new_order_closes_previous() -> None: repo = Repo() repo.boundary = ProductionOrderBoundary( "machine", "FA-1", datetime(2026, 9, 16, 9, tzinfo=UTC), None, ) event = source(identifier="open", end=None) first = service(repo, [event], status(remaining=0, feedback=NOW)) first.reconcile_once(NOW) assert repo.boundary is not None and repo.boundary.ended_at == NOW assert repo.events["open"].production_order == "FA-1" assert repo.events["open"].attributed_end == NOW repo.boundary = ProductionOrderBoundary( "machine", "FA-1", datetime(2026, 9, 16, 8, tzinfo=UTC), None, ) service(repo, [], status(order="2", feedback=NOW)).reconcile_once(NOW) assert repo.closures == [("FA-1", NOW), ("FA-1", NOW)] assert repo.boundary.production_order == "FA-2" def test_unknown_is_not_unplanned_and_no_schedule_is_used() -> None: repo = Repo() service(repo, [source(end=datetime(2026, 9, 20, tzinfo=UTC))], status()).reconcile_once(NOW) assert repo.events["d"].category is DowntimeCategory.UNKNOWN