OOB-Dashboard/ppcbudget/scoring.py

322 lines
12 KiB
Python

"""Reconstruct each campaign's budget state timeline from change events.
The export lists only *changes*, so the state between two rows has to be
inferred. Three things make that non-obvious:
1. Rows are written newest-first, so rows sharing a minute are also
newest-first and must be reversed before walking the machine. Sorting on
timestamp alone silently preserves the wrong intra-minute order -- on the
reference file that costs 64 campaign-hours and manufactures 14 phantom
inconsistencies.
2. Delivery state (Paused) overlays budget state. A paused campaign forgoes
nothing to its budget, so paused minutes are excluded from the loss-eligible
total and from the in-budget denominator that sets the spend rate.
3. De-duplication during export can drop an intermediate transition, leaving a
row whose `From` disagrees with the running state. We trust the row over the
inference and place the implied transition at the midpoint of the gap.
"""
from __future__ import annotations
from dataclasses import dataclass, field
from .ingest import BUDGET_STATES, Event, parse_money
IN, OOB, PAUSED, NA = 0, 1, 2, 3
STATE_NAMES = {IN: "in", OOB: "out_of_budget", PAUSED: "paused", NA: "not_eligible"}
MINUTES_PER_DAY = 1440
DEFAULT_MERGE_GAP_MIN = 5
@dataclass(slots=True)
class ChainBreak:
at_min: int
expected_from: str
saw_from: str
ambiguity_min: int
@dataclass(slots=True)
class Episode:
index: int
start_min: int
end_min: int
raw_min: int
active_min: int # raw minus any overlapping pause
@dataclass(slots=True)
class BudgetInfo:
value: float | None = None
source: str = "unknown" # daily_budget_event | budget_rule | perf_report | unknown
time_weighted: float | None = None
changes: int = 0
@dataclass(slots=True)
class CampaignDay:
campaign: str
date_key: str
t0: int
t1: int
in_min: int = 0
oob_min: int = 0 # loss-eligible: pause and out-of-window already removed
paused_min: int = 0
na_min: int = 0
oob_min_raw: int = 0 # what the budget machine alone says
post_reset_oob_min: int = 0
episodes: list[Episode] = field(default_factory=list)
episodes_raw: int = 0
episodes_merged: int = 0
first_oob_min: int | None = None
last_recovery_min: int | None = None
opened_oob: bool = False
closed_oob: bool = False
hourly_oob: list[int] = field(default_factory=lambda: [0] * 24)
hourly_paused: list[int] = field(default_factory=lambda: [0] * 24)
hourly_na: list[int] = field(default_factory=lambda: [0] * 24)
track: list[tuple[int, int, int]] = field(default_factory=list) # (state, start, end)
budget: BudgetInfo = field(default_factory=BudgetInfo)
chain_breaks: list[ChainBreak] = field(default_factory=list)
oob_uncertainty_min: float = 0.0
confidence: str = "clean" # clean | repaired | partial_day
# filled in by metrics.py / perfjoin.py
severity: float = 0.0
diagnosis: str = ""
lost: dict | None = None
perf: object | None = None
@property
def eligible_min(self) -> int:
return self.t1 - self.t0
@property
def active_min(self) -> int:
"""Eligible minutes where the campaign could actually have spent."""
return self.eligible_min - self.paused_min
@property
def oob_share(self) -> float:
return self.oob_min / self.active_min if self.active_min > 0 else 0.0
@property
def in_hours(self) -> float:
return self.in_min / 60
@property
def oob_hours(self) -> float:
return self.oob_min / 60
@property
def paused_hours(self) -> float:
return self.paused_min / 60
def sort_key(e: Event) -> tuple[int, int]:
"""Chronological order. Descending source index reverses the newest-first file."""
return (e.minute, -e.source_index)
def _walk(events: list[Event], t0: int, t1: int, initial: str):
"""Turn a two-state event stream into contiguous spans over [t0, t1)."""
spans: list[tuple[str, int, int]] = []
breaks: list[ChainBreak] = []
state, prev = initial, t0
for e in events:
m = e.minute if e.minute > prev else prev
if e.from_val != state:
# A transition went missing. The row is observation, the running
# state is inference, so trust the row and split the difference.
mid = (prev + m) // 2
if mid > prev:
spans.append((state, prev, mid))
if m > mid:
spans.append((e.from_val, mid, m))
breaks.append(ChainBreak(m, state, e.from_val, m - prev))
elif m > prev:
spans.append((state, prev, m))
state, prev = e.to_val, m
if t1 > prev:
spans.append((state, prev, t1))
return spans, breaks
def _budget_timeline(events: list[Event], rule_events: list[Event],
t0: int, t1: int) -> BudgetInfo:
"""Recover the daily budget, which can change during the day."""
amounts = sorted(events, key=sort_key)
if amounts:
opening = amounts[0].from_num
spans: list[tuple[float | None, int, int]] = []
value, prev = opening, t0
for e in amounts:
m = max(e.minute, prev)
if m > prev:
spans.append((value, prev, m))
value, prev = e.to_num, m
if t1 > prev:
spans.append((value, prev, t1))
weighted = sum(v * (b - a) for v, a, b in spans if v is not None)
covered = sum(b - a for v, a, b in spans if v is not None)
return BudgetInfo(
value=amounts[-1].to_num,
source="daily_budget_event",
time_weighted=(weighted / covered) if covered else None,
changes=len(amounts),
)
for e in rule_events:
amount = parse_money(e.from_val, e.to_val)
if amount is not None:
return BudgetInfo(value=amount, source="budget_rule", time_weighted=amount)
return BudgetInfo()
def score_campaign_day(campaign: str, date_key: str, events: list[Event],
day_end_min: int = MINUTES_PER_DAY,
merge_gap_min: int = DEFAULT_MERGE_GAP_MIN) -> CampaignDay | None:
"""Score one campaign for one day. Returns None if budget state is unknowable."""
budget_events = sorted((e for e in events if e.machine == "budget"), key=sort_key)
if not budget_events:
return None # never impute a state we did not observe
delivery_events = sorted((e for e in events if e.machine == "delivery"), key=sort_key)
created = [e for e in events if e.machine == "created"]
# Eligible window. A campaign created at 07:02 is scored over the remaining
# 16.97h, not a full day -- but only if no budget event precedes creation.
t1 = min(day_end_min, MINUTES_PER_DAY)
t0 = 0
if created:
birth = min(e.minute for e in created)
if birth <= budget_events[0].minute and birth < t1:
t0 = birth
budget_events = [e for e in budget_events if e.minute >= t0]
if not budget_events:
return None
delivery_events = [e for e in delivery_events if e.minute >= t0]
day = CampaignDay(campaign=campaign, date_key=date_key, t0=t0, t1=t1)
budget_spans, breaks = _walk(budget_events, t0, t1, budget_events[0].from_val)
day.chain_breaks = breaks
day.oob_uncertainty_min = sum(b.ambiguity_min for b in breaks) / 2
delivery_initial = delivery_events[0].from_val if delivery_events else "Delivering"
delivery_spans, _ = _walk(delivery_events, t0, t1, delivery_initial)
# Flatten: not-eligible > paused > out of budget > in budget.
flat = bytearray([NA]) * MINUTES_PER_DAY
for state, a, b in budget_spans:
flat[a:b] = bytes([OOB if state == "Out of budget" else IN]) * (b - a)
for state, a, b in delivery_spans:
if state == "Paused":
flat[a:b] = bytes([PAUSED]) * (b - a)
day.in_min = flat.count(IN)
day.oob_min = flat.count(OOB)
day.paused_min = flat.count(PAUSED)
day.na_min = flat.count(NA)
day.post_reset_oob_min = flat[60:t1].count(OOB) if t1 > 60 else 0
day.hourly_oob = [flat[h * 60:(h + 1) * 60].count(OOB) for h in range(24)]
day.hourly_paused = [flat[h * 60:(h + 1) * 60].count(PAUSED) for h in range(24)]
day.hourly_na = [flat[h * 60:(h + 1) * 60].count(NA) for h in range(24)]
oob_spans = [(a, b) for s, a, b in budget_spans if s == "Out of budget"]
day.oob_min_raw = sum(b - a for a, b in oob_spans)
day.episodes_raw = len(oob_spans)
# Amazon can release a sliver of budget that is consumed within the same
# minute, producing zero-length recoveries. Those are pacing noise, not
# genuine outages, so also report a merged count.
merged: list[tuple[int, int]] = []
for a, b in oob_spans:
if merged and a - merged[-1][1] < merge_gap_min:
merged[-1] = (merged[-1][0], b)
else:
merged.append((a, b))
day.episodes_merged = len(merged)
day.episodes = [
Episode(i + 1, a, b, b - a, flat[a:b].count(OOB))
for i, (a, b) in enumerate(merged)
]
# State facts come from the budget machine; a concurrent pause must not
# mask the fact that the budget itself was exhausted. Loss math, above,
# uses the flattened track instead.
day.opened_oob = budget_spans[0][0] == "Out of budget"
day.closed_oob = budget_spans[-1][0] == "Out of budget"
day.first_oob_min = oob_spans[0][0] if oob_spans else None
day.last_recovery_min = None if day.closed_oob or not oob_spans else oob_spans[-1][1]
day.budget = _budget_timeline(
[e for e in events if e.machine == "budget_amount"],
[e for e in events if e.machine == "budget_rule"],
t0, t1,
)
# Run-length encode for the report; spans are contiguous and tile the day.
track: list[tuple[int, int, int]] = []
run_state, run_start = flat[0], 0
for m in range(1, MINUTES_PER_DAY):
if flat[m] != run_state:
track.append((run_state, run_start, m))
run_state, run_start = flat[m], m
track.append((run_state, run_start, MINUTES_PER_DAY))
day.track = track
if breaks:
day.confidence = "repaired"
if day.na_min > 0:
day.confidence = "partial_day"
return day
def score_all(events: list[Event], merge_gap_min: int = DEFAULT_MERGE_GAP_MIN) -> list[CampaignDay]:
"""Score every campaign in every day present in the event stream."""
by_day: dict[str, dict[str, list[Event]]] = {}
for e in events:
by_day.setdefault(e.date_key, {}).setdefault(e.campaign, []).append(e)
results: list[CampaignDay] = []
for date_key in sorted(by_day):
campaigns = by_day[date_key]
# If the export was cut short (a "Today" pull), do not score the
# unobserved remainder of the day as if it were in budget.
last_seen = max(e.minute for evs in campaigns.values() for e in evs)
day_end = MINUTES_PER_DAY if last_seen >= MINUTES_PER_DAY - 1 else last_seen + 1
for campaign, evs in campaigns.items():
day = score_campaign_day(campaign, date_key, evs, day_end, merge_gap_min)
if day is not None:
results.append(day)
return results
def check_invariants(days: list[CampaignDay]) -> list[str]:
"""Structural checks that make chart/table disagreement impossible."""
problems: list[str] = []
for d in days:
total = d.in_min + d.oob_min + d.paused_min + d.na_min
if total != MINUTES_PER_DAY:
problems.append(f"{d.campaign} [{d.date_key}]: minutes sum to {total}, not 1440")
if sum(e.active_min for e in d.episodes) != d.oob_min:
problems.append(f"{d.campaign} [{d.date_key}]: episode minutes != out-of-budget minutes")
if sum(d.hourly_oob) != d.oob_min:
problems.append(f"{d.campaign} [{d.date_key}]: hourly buckets != out-of-budget minutes")
for i in range(1, len(d.track)):
if d.track[i][1] != d.track[i - 1][2]:
problems.append(f"{d.campaign} [{d.date_key}]: timeline has a gap or overlap")
break
return problems