Source code for fhops.evaluation.playback.core

"""Deterministic playback primitives."""

from __future__ import annotations

from collections import defaultdict
from collections.abc import Iterable, Iterator, Sequence
from dataclasses import dataclass, field
from typing import TYPE_CHECKING

from fhops.scenario.contract import Problem

if TYPE_CHECKING:  # pragma: no cover - import for typing only
    import pandas as pd

__all__ = [
    "PlaybackConfig",
    "PlaybackRecord",
    "ShiftSummary",
    "DaySummary",
    "PlaybackResult",
    "run_playback",
    "summarise_shifts",
    "summarise_days",
]


[docs] @dataclass(slots=True) class PlaybackConfig: """Tuning options for deterministic playback execution.""" respect_blackouts: bool = True infer_missing_shifts: bool = True include_idle_records: bool = False # TODO: emit idle rows once availability bridge lands.
[docs] @dataclass(slots=True) class PlaybackRecord: """Atomic shift-level record produced by deterministic playback.""" day: int shift_id: str machine_id: str block_id: str | None hours_worked: float | None = None production_units: float | None = None mobilisation_cost: float | None = None blackout_hit: bool = False landing_id: str | None = None machine_role: str | None = None downtime: bool = False weather_severity: float | None = None metadata: dict[str, object] = field(default_factory=dict) sample_id: int = 0
[docs] @dataclass(slots=True) class ShiftSummary: """Aggregated metrics per machine/shift.""" day: int shift_id: str machine_id: str machine_role: str | None = None available_hours: float = 0.0 total_hours: float = 0.0 production_units: float = 0.0 mobilisation_cost: float = 0.0 idle_hours: float | None = None blackout_conflicts: int = 0 sequencing_violations: int = 0 utilisation_ratio: float | None = None sample_id: int = 0 downtime_hours: float = 0.0 downtime_events: int = 0 weather_severity_total: float = 0.0
[docs] @dataclass(slots=True) class DaySummary: """Aggregated metrics per day across machines.""" day: int available_hours: float = 0.0 total_hours: float = 0.0 production_units: float = 0.0 mobilisation_cost: float = 0.0 completed_blocks: int = 0 idle_hours: float | None = None blackout_conflicts: int = 0 sequencing_violations: int = 0 utilisation_ratio: float | None = None sample_id: int = 0 downtime_hours: float = 0.0 downtime_events: int = 0 weather_severity_total: float = 0.0
[docs] @dataclass(slots=True) class PlaybackResult: """Container grouping playback outputs.""" records: Sequence[PlaybackRecord] shift_summaries: Sequence[ShiftSummary] day_summaries: Sequence[DaySummary] config: PlaybackConfig sequencing_debug: dict[str, object] | None = None sample_id: int = 0 delivered_total: float = 0.0 remaining_work_total: float = 0.0
[docs] def run_playback( problem: Problem, assignments: pd.DataFrame, *, config: PlaybackConfig | None = None, sample_id: int = 0, ) -> PlaybackResult: """Convert solver assignments into playback records and aggregated summaries.""" cfg = config or PlaybackConfig() availability_map = _compute_shift_availability(problem, cfg) from .adapters import assignments_to_records # Local import to avoid circular dependency. record_iter = assignments_to_records(problem, assignments) # Materialise records since downstream summaries iterate multiple times. records: tuple[PlaybackRecord, ...] = tuple(record_iter) sequencing_debug: dict[str, object] | None = None tracker = getattr(record_iter, "sequencing_tracker", None) delivered_total = 0.0 remaining_work_total = 0.0 if tracker is not None: sequencing_debug = tracker.debug_snapshot() delivered_total = float(getattr(tracker, "delivered_total", 0.0) or 0.0) remaining_work_total = float(sum(tracker.remaining_work.values())) shift_summaries = tuple( summarise_shifts( records, availability_map, include_idle=cfg.include_idle_records, sample_id=sample_id, ) ) completed_by_day: dict[int, set[str]] = defaultdict(set) for record in records: if record.block_id and record.metadata.get("block_completed"): completed_by_day[record.day].add(record.block_id) day_summaries = tuple( summarise_days( shift_summaries, availability_map, completed_by_day, sample_id=sample_id, ) ) if cfg.respect_blackouts: # When respecting blackouts, flag entries already emitted; future iterations may # choose to filter or adjust scheduling, but caller still receives full trace. pass return PlaybackResult( records=records, shift_summaries=shift_summaries, day_summaries=day_summaries, config=cfg, sequencing_debug=sequencing_debug, sample_id=sample_id, delivered_total=delivered_total, remaining_work_total=remaining_work_total, )
[docs] def summarise_shifts( records: Iterable[PlaybackRecord], availability_map: dict[tuple[int, str, str], float], *, include_idle: bool = False, sample_id: int = 0, ) -> Iterator[ShiftSummary]: """Aggregate playback records to machine/shift summaries. Parameters ---------- records: Iterator of :class:`PlaybackRecord` instances produced by ``assignments_to_records``. availability_map: Mapping ``(day, shift_id, machine_id) -> available_hours`` computed via ``_compute_shift_availability``; used to seed summaries and compute utilisation ratios. include_idle: When ``True``, emit summaries for every availability entry even if no work occurred (resulting in ``total_hours = 0`` but preserving the availability baseline). sample_id: Identifier propagated through stochastic playback so downstream aggregations can group rows. """ aggregates: dict[tuple[int, str, str], ShiftSummary] = {} seen: set[tuple[int, str, str]] = set() if include_idle: for (day, shift_id, machine_id), available_hours in availability_map.items(): aggregates[(day, shift_id, machine_id)] = ShiftSummary( day=day, shift_id=shift_id, machine_id=machine_id, available_hours=available_hours, sample_id=sample_id, ) for record in records: key = (record.day, record.shift_id, record.machine_id) summary = aggregates.get(key) if summary is None: summary = ShiftSummary( day=record.day, shift_id=record.shift_id, machine_id=record.machine_id, available_hours=availability_map.get(key, 0.0), machine_role=record.machine_role, sample_id=sample_id, ) aggregates[key] = summary seen.add(key) if summary.machine_role is None and record.machine_role is not None: summary.machine_role = record.machine_role if record.hours_worked is not None: summary.total_hours += record.hours_worked if record.production_units is not None: summary.production_units += record.production_units if record.mobilisation_cost is not None: summary.mobilisation_cost += record.mobilisation_cost if record.blackout_hit: summary.blackout_conflicts += 1 if record.metadata.get("sequencing_violation"): summary.sequencing_violations += 1 if record.downtime and record.hours_worked is not None: summary.downtime_hours += record.hours_worked summary.downtime_events += 1 if record.weather_severity: summary.weather_severity_total += float(record.weather_severity) for key, summary in aggregates.items(): if summary.idle_hours is None: summary.idle_hours = max(summary.available_hours - summary.total_hours, 0.0) if summary.available_hours > 0: summary.utilisation_ratio = summary.total_hours / summary.available_hours else: summary.utilisation_ratio = None for key in sorted(aggregates): if include_idle or key in seen: yield aggregates[key]
[docs] def summarise_days( shift_summaries: Iterable[ShiftSummary], availability_map: dict[tuple[int, str, str], float], completed_by_day: dict[int, set[str]], sample_id: int = 0, ) -> Iterator[DaySummary]: """Aggregate playback results to day-level summaries. Parameters ---------- shift_summaries: Iterable produced by :func:`summarise_shifts`. availability_map: Same map passed to ``summarise_shifts``; used to recover day-level availability totals. completed_by_day: Mapping ``day -> set(block_id)`` indicating which blocks completed on each day (consumed for KPI reporting). sample_id: Propagated through stochastic playback. """ aggregates: dict[int, DaySummary] = {} availability_by_day: dict[int, float] = defaultdict(float) for (day, _shift_id, _machine_id), hours in availability_map.items(): availability_by_day[day] += hours for summary in shift_summaries: day_summary = aggregates.get(summary.day) if day_summary is None: day_summary = DaySummary(day=summary.day, sample_id=sample_id) aggregates[summary.day] = day_summary day_summary.available_hours += summary.available_hours day_summary.total_hours += summary.total_hours day_summary.production_units += summary.production_units day_summary.mobilisation_cost += summary.mobilisation_cost day_summary.blackout_conflicts += summary.blackout_conflicts day_summary.sequencing_violations += summary.sequencing_violations day_summary.downtime_hours += summary.downtime_hours day_summary.downtime_events += summary.downtime_events day_summary.weather_severity_total += summary.weather_severity_total for day, available in availability_by_day.items(): day_summary = aggregates.setdefault(day, DaySummary(day=day, sample_id=sample_id)) day_summary.available_hours = max(day_summary.available_hours, available) # Compute per-day idle hours and completed block counts via aggregated data. for day in sorted(aggregates): day_summary = aggregates[day] day_summary.idle_hours = max(day_summary.available_hours - day_summary.total_hours, 0.0) day_summary.completed_blocks = len(completed_by_day.get(day, set())) if day_summary.available_hours > 0: day_summary.utilisation_ratio = day_summary.total_hours / day_summary.available_hours else: day_summary.utilisation_ratio = None yield day_summary
def _compute_shift_availability( problem: Problem, config: PlaybackConfig, ) -> dict[tuple[int, str, str], float]: """Derive available hours per machine/day/shift from scenario data.""" scenario = problem.scenario machines = {machine.id: machine for machine in scenario.machines} day_availability: dict[tuple[str, int], int] = {} for calendar_entry in scenario.calendar: day_availability[(calendar_entry.machine_id, calendar_entry.day)] = calendar_entry.available shift_hours = {} if scenario.timeline and scenario.timeline.shifts: shift_hours = {shift_def.name: shift_def.hours for shift_def in scenario.timeline.shifts} availability: dict[tuple[int, str, str], float] = {} def machine_available(machine_id: str, day: int) -> bool: return day_availability.get((machine_id, day), 1) == 1 if scenario.shift_calendar: for shift_entry in scenario.shift_calendar: if shift_entry.available != 1: continue if shift_entry.machine_id not in machines: continue if not machine_available(shift_entry.machine_id, shift_entry.day): continue hours = shift_hours.get(shift_entry.shift_id) if hours is None and config.infer_missing_shifts: hours = machines[shift_entry.machine_id].daily_hours if hours is None: continue availability[(shift_entry.day, shift_entry.shift_id, shift_entry.machine_id)] = hours elif shift_hours: for day in problem.days: for shift_id, hours in shift_hours.items(): for machine_id in machines: if not machine_available(machine_id, day): continue availability[(day, shift_id, machine_id)] = hours else: for day in problem.days: for machine_id, machine in machines.items(): if not machine_available(machine_id, day): continue availability[(day, "S1", machine_id)] = machine.daily_hours return availability