Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
16 changes: 15 additions & 1 deletion ARTIFACTS.md
Original file line number Diff line number Diff line change
Expand Up @@ -79,7 +79,12 @@ One JSON object per line per PT ever seen (universe + standalone entities):
"episodes": [{"start": "2025-04-24T07:00", "end": "2025-10-16T23:00",
"ongoing": false, "uncertain": true, "cause_class": "programat",
"cause_raw": "Revizie tehnica CTE Progresu",
"remediere_last": "2025-10-15T23:00"}]}}}
"remediere_last": "2025-10-15T23:00"}],
"episodes_count_deficienta": 2, "est_hours_deficienta": 61.5,
"episodes_deficienta": [{"start": "2025-07-19T08:00", "end": "2025-07-20T08:00",
"ongoing": false, "uncertain": false, "cause_class": "unclassified",
"cause_raw": "Lipsa parametri pentru livrare apa calda de consum",
"remediere_last": null}]}}}
```
`runs` = [start_day_of_year (1-based, local), length_days, cause_class],
episodes clipped to the year; ongoing = ended_before null; uncertain =
Expand All @@ -91,6 +96,14 @@ Run `cause_class` takes one of FOUR values: `avarie`, `programat`,
of non-deficienta runs equals `days` for every PT-year and street-year).
Consumers painting or summing headline days must exclude `deficienta` runs.

Deficienta also carries its own episode array and counters:
`episodes_deficienta`, `episodes_count_deficienta`, `est_hours_deficienta`
(PT only - streets have no episode array at all). These are DISJOINT from the
headline trio: `episodes` / `episodes_count` / `est_hours` never include a
deficienta episode, and the deficienta three never include an oprire one. They
must never be summed into the headline. `days_deficienta > 0` if and only if
`episodes_count_deficienta > 0`.

### strazi/all.ndjson.gz
```json
{"slug": "sos-pantelimon", "name": "Sos Pantelimon", "type": "sos",
Expand All @@ -100,6 +113,7 @@ Consumers painting or summing headline days must exclude `deficienta` runs.
"inferred_pt": null, "inferred_km": null,
"addr": {"64": [3, 0.18], "64a": [3, 0.2], "70": [5, 0.42]},
"years": {"2025": {"days": 181, "days_avarie": 121, "days_programat": 74,
"days_deficienta": 28,
"runs": [[10, 3, "avarie"], "..."]}}}
```
Street universe = streets CMTEB named in outages UNION all named Bucharest
Expand Down
44 changes: 37 additions & 7 deletions pipeline/publish.py
Original file line number Diff line number Diff line change
Expand Up @@ -307,6 +307,12 @@ def build(db_path: str, registry_path: str, harta_html: str,
pt_eps: dict[tuple, list[dict]] = defaultdict(list) # (pt, y) display episodes
pt_epn: Counter = Counter() # (pt, y) oprire eps touching
pt_hours: dict[tuple, float] = defaultdict(float) # (pt, y) split est_hours
# Parallel deficienta structures, deliberately separate from every headline
# map above. Nothing here ever feeds pt_days, city_eps, pt_epn or pt_hours,
# so `days` and everything derived from it stay byte-identical to v1.
pt_eps_defi: dict[tuple, list[dict]] = defaultdict(list) # (pt, y) display eps
pt_epn_defi: Counter = Counter() # (pt, y) eps touching
pt_hours_defi: dict[tuple, float] = defaultdict(float) # (pt, y) split est_hours
pt_blocks: dict[str, int] = {}
sector_votes: dict[str, Counter] = defaultdict(Counter)
city_eps: Counter = Counter() # start-year attributed
Expand All @@ -329,9 +335,23 @@ def build(db_path: str, registry_path: str, harta_html: str,
sector_votes[pt][sector] += 1
if blocks:
pt_blocks[pt] = max(pt_blocks.get(pt, 0), blocks)
# The display record is a pure function of the episode row, so it is
# hoisted above the severity branch and shared. Building it for
# deficienta rows is the new behaviour; the oprire record is unchanged.
epd = {"start": fmt_minute(first), "end": fmt_minute(last),
"ongoing": ended is None, "uncertain": bool(gap),
"cause_class": cls, "cause_raw": craw,
"remediere_last": _rem_iso(rem)}
if sev == "deficienta":
for d in days:
pt_cls[(pt, d.year, "deficienta")].add(d)
# Same proportional year-split rule as the oprire path below, and
# keyed off the same `ycount`, so days_deficienta > 0 and
# episodes_count_deficienta > 0 hold on identical year sets.
for y, c in ycount.items():
pt_eps_defi[(pt, y)].append(epd)
pt_epn_defi[(pt, y)] += 1
pt_hours_defi[(pt, y)] += est_h * (c / len(days))
continue
y0 = local(first).year
city_eps[y0] += 1
Expand All @@ -341,10 +361,6 @@ def build(db_path: str, registry_path: str, harta_html: str,
for d in days:
pt_days[(pt, d.year)].add(d)
pt_cls[(pt, d.year, cls)].add(d)
epd = {"start": fmt_minute(first), "end": fmt_minute(last),
"ongoing": ended is None, "uncertain": bool(gap),
"cause_class": cls, "cause_raw": craw,
"remediere_last": _rem_iso(rem)}
for y, c in ycount.items():
pt_eps[(pt, y)].append(epd)
pt_epn[(pt, y)] += 1
Expand Down Expand Up @@ -718,6 +734,9 @@ def _longest(days_map: dict, ent, y: int) -> int:
write_json(out / "client" / "map" / f"pt-{y}.geojson",
{"type": "FeatureCollection", "features": feats})

def _epkey(e):
return (e["start"], e["end"], e["cause_class"], e["cause_raw"] or "")

pt_lines = []
for pt in sorted(pts_all, key=lambda p: pt_slug[p]):
years_obj = {}
Expand All @@ -726,9 +745,8 @@ def _longest(days_map: dict, ent, y: int) -> int:
defi = pt_cls.get((pt, y, "deficienta"), set())
if not union and not defi:
continue
eps = sorted(pt_eps.get((pt, y), ()),
key=lambda e: (e["start"], e["end"], e["cause_class"],
e["cause_raw"] or ""))
eps = sorted(pt_eps.get((pt, y), ()), key=_epkey)
eps_defi = sorted(pt_eps_defi.get((pt, y), ()), key=_epkey)
years_obj[str(y)] = {
"days": len(union),
"days_avarie": len(pt_cls.get((pt, y, "avarie"), ())),
Expand All @@ -739,6 +757,12 @@ def _longest(days_map: dict, ent, y: int) -> int:
"est_hours": round(pt_hours[(pt, y)], 1),
"runs": _runs(pt_cls, pt, y),
"episodes": eps,
# Additive deficienta counterparts. Disjoint from the three
# above: episodes_count / est_hours / episodes never include a
# deficienta episode, and these never include an oprire one.
"episodes_count_deficienta": pt_epn_defi[(pt, y)],
"est_hours_deficienta": round(pt_hours_defi[(pt, y)], 1),
"episodes_deficienta": eps_defi,
}
c = pt_coord.get(pt)
pt_lines.append({
Expand Down Expand Up @@ -766,6 +790,12 @@ def _longest(days_map: dict, ent, y: int) -> int:
"days": len(union),
"days_avarie": len(st_cls.get((k, y, "avarie"), ())),
"days_programat": len(st_cls.get((k, y, "programat"), ())),
# `defi` was already computed above as the emission gate but was
# never written out, so street-years shipped deficienta RUNS with
# no counter to reconcile them against - and ARTIFACTS.md:91
# claims the union invariant holds "for every PT-year and
# street-year". Emitting it makes publish match the contract.
"days_deficienta": len(defi),
"runs": _runs(st_cls, k, y),
}
blocks = []
Expand Down
174 changes: 174 additions & 0 deletions pipeline/validate.py
Original file line number Diff line number Diff line change
Expand Up @@ -33,6 +33,27 @@
ARTIFACT_KEYS_META = ("generated_at", "data_through", "years", "last_complete_year",
"partial_years", "universe_size", "coverage", "sources_cutover_utc")

# Exact key sets for published rows. These ARE the contract, so `set(row) !=
# want` is the assertion: an unannounced extra key fails as loudly as a missing
# one. tests/test_publish_shapes.py imports these so there is one definition.
ARTIFACT_KEYS_PT_RANK = frozenset({
"slug", "name", "sector", "days", "days_avarie", "days_programat",
"days_deficienta", "episodes", "longest_days", "est_day_eq", "delta_prev"})
ARTIFACT_KEYS_ST_RANK = (ARTIFACT_KEYS_PT_RANK - {"sector"}) | {"sectors", "pt_slugs"}
# Year objects use a SUBSET check instead: they are where additive fields land
# most often, and validate should not block a future additive key.
ARTIFACT_KEYS_PT_YEAR = frozenset({
"days", "days_avarie", "days_programat", "days_deficienta",
"episodes_count", "longest_days", "est_hours", "runs", "episodes",
"episodes_count_deficienta", "est_hours_deficienta", "episodes_deficienta"})
ARTIFACT_KEYS_ST_YEAR = frozenset({
"days", "days_avarie", "days_programat", "days_deficienta", "runs"})

# The three headline cause classes. `deficienta` is the pseudo-class feeding the
# secondary counter, excluded from `days` by contract (ARTIFACTS.md:88-92).
NON_DEFICIENTA = ("avarie", "programat", "unclassified")
MAX_DOY = 366


def _alias(db) -> dict[str, str]:
try:
Expand Down Expand Up @@ -210,6 +231,88 @@ def check_sector_consistency(db, web: Path, harta_html: str,
return out


def _run_day_set(runs: list, classes: tuple[str, ...]) -> set[int]:
"""Day-of-year numbers covered by runs whose cause is in `classes`.

Runs of different classes legitimately OVERLAP - a day can carry both an
avarie and a programat episode, which is exactly why ARTIFACTS.md says the
per-class counts may sum above `days`. So this unions DOY numbers rather
than summing run lengths; summing would over-count and make the contract
look breached on perfectly good data.

Callers must shape-check `runs` first (see _year_object_problems).
"""
out: set[int] = set()
for start, length, cls in runs:
if cls in classes:
out.update(range(start, start + length))
return out


def _year_object_problems(ns: str, yo) -> list[str]:
"""Contract violations inside one entity-year object. `ns` is "pt" or "st".

The load-bearing check is ARTIFACTS.md:90-91, stated there as a hard
guarantee and enforced nowhere until now: the union of non-deficienta run
days equals `days`, and deficienta runs account for exactly
`days_deficienta`.

Returns strings and never raises - one malformed entity must not abort the
whole ndjson scan. Details carry only slugs, years and integers, never
cause_raw, so untrusted CMTEB text cannot forge PASS/FAIL lines in the CI
log (SEC049).
"""
if not isinstance(yo, dict):
return [f"year object is {type(yo).__name__}, not an object"]
probs: list[str] = []
want = ARTIFACT_KEYS_PT_YEAR if ns == "pt" else ARTIFACT_KEYS_ST_YEAR
missing = sorted(want - set(yo))
if missing:
probs.append(f"missing keys {missing}")

runs = yo.get("runs")
if not isinstance(runs, list):
return probs + ["runs missing or not a list"]
for r in runs:
if not (isinstance(r, list) and len(r) == 3
and isinstance(r[0], int) and isinstance(r[1], int)
and isinstance(r[2], str)
and r[0] >= 1 and r[1] >= 1 and r[0] + r[1] - 1 <= MAX_DOY):
return probs + [f"malformed run {r}"]

if isinstance(yo.get("days"), int):
head = len(_run_day_set(runs, NON_DEFICIENTA))
if head != yo["days"]:
probs.append(f"non-deficienta runs cover {head} days != days {yo['days']}")
if isinstance(yo.get("days_deficienta"), int):
defi = len(_run_day_set(runs, ("deficienta",)))
if defi != yo["days_deficienta"]:
probs.append(f"deficienta runs cover {defi} days "
f"!= days_deficienta {yo['days_deficienta']}")
return probs


def _deficienta_problems(yo: dict) -> list[str]:
"""PT-only reconciliation between deficienta counters and their episodes.

publish.py appends to episodes_deficienta and increments
episodes_count_deficienta in one loop, and derives the deficienta day set
from the SAME per-episode year Counter, so both relations below are exact
biconditionals by construction. Any drift is a real bug, never data noise.
"""
probs: list[str] = []
for cnt_key, eps_key in (("episodes_count", "episodes"),
("episodes_count_deficienta", "episodes_deficienta")):
n, eps = yo.get(cnt_key), yo.get(eps_key)
if isinstance(eps, list) and isinstance(n, int) and n != len(eps):
probs.append(f"{cnt_key} {n} != len({eps_key}) {len(eps)}")
n, dd = yo.get("episodes_count_deficienta"), yo.get("days_deficienta")
if isinstance(n, int) and isinstance(dd, int) and (n > 0) != (dd > 0):
probs.append(f"episodes_count_deficienta {n} inconsistent "
f"with days_deficienta {dd}")
return probs


def check_artifacts(web: Path, registry_slugs: set[str] | None = None) -> list[tuple]:
out = []

Expand Down Expand Up @@ -242,6 +345,12 @@ def load(rel: str):
streets_with_addr = 0
addr_numbers = 0
addr_bad: list[str] = []
year_bad: list[str] = [] # runs_days_union
key_bad: list[str] = [] # ndjson_year_keys
recon_bad: list[str] = [] # deficienta_reconciliation
n_year_bad = n_key_bad = n_recon_bad = 0
defi_by_nsyear: Counter = Counter() # (ns, "YYYY") -> sum(days_deficienta)
nsyears_seen: set[tuple] = set()
for ns, rel in (("pt", "pt/all.ndjson.gz"), ("st", "strazi/all.ndjson.gz")):
p = web / rel
if not p.exists():
Expand All @@ -253,6 +362,25 @@ def load(rel: str):
for line in f:
o = json.loads(line)
bag.add(o["slug"])
for ystr, yo in (o.get("years") or {}).items():
nsyears_seen.add((ns, ystr))
for msg in _year_object_problems(ns, yo):
if msg.startswith("missing keys"):
n_key_bad += 1
if len(key_bad) < 5:
key_bad.append(f"{o['slug']}/{ystr}: {msg}")
else:
n_year_bad += 1
if len(year_bad) < 5:
year_bad.append(f"{o['slug']}/{ystr}: {msg}")
if ns == "pt" and isinstance(yo, dict):
for msg in _deficienta_problems(yo):
n_recon_bad += 1
if len(recon_bad) < 5:
recon_bad.append(f"{o['slug']}/{ystr}: {msg}")
if isinstance(yo, dict) and isinstance(
yo.get("days_deficienta"), int):
defi_by_nsyear[(ns, ystr)] += yo["days_deficienta"]
if ns == "st" and o.get("blocks"):
streets_with_blocks += 1
block_pts.update(b["pt"] for b in o["blocks"])
Expand Down Expand Up @@ -288,7 +416,33 @@ def load(rel: str):
f"addr pt_index out of range / bad -1: {addr_bad[:5]}" if addr_bad
else f"addr maps on {streets_with_addr} streets, {addr_numbers} numbers, all resolvable"))

# --- deficienta contract (ARTIFACTS.md:88-92) --------------------------
out.append((FAIL if year_bad else PASS, "runs_days_union",
f"{n_year_bad} entity-years breach the runs/days contract: "
+ "; ".join(year_bad) if year_bad
else f"non-deficienta runs == days and deficienta runs == "
f"days_deficienta across {len(nsyears_seen)} namespace-years"))
out.append((FAIL if key_bad else PASS, "ndjson_year_keys",
f"{n_key_bad} entity-years: " + "; ".join(key_bad) if key_bad
else "year objects carry the contract key set"))
out.append((FAIL if recon_bad else PASS, "deficienta_reconciliation",
f"{n_recon_bad} entity-years: " + "; ".join(recon_bad) if recon_bad
else "episode counts reconcile with episode arrays and day counts"))
# Heuristic, not an identity: a genuinely quiet year is conceivable upstream,
# and FAILing on it would let a CMTEB editorial change block the release.
# Every namespace-year at zero simultaneously is a pipeline regression, so
# only that case FAILs.
vac = sorted(f"{ns}-{y}" for (ns, y) in nsyears_seen
if defi_by_nsyear[(ns, y)] == 0)
lvl = FAIL if vac and len(vac) == len(nsyears_seen) else WARN if vac else PASS
out.append((lvl, "deficienta_non_vacuous",
f"days_deficienta is 0 across every entity in: {vac[:6]}" if vac
else f"days_deficienta populated in all {len(nsyears_seen)} "
f"namespace-years"))

rank_bad = []
row_key_bad: list[str] = []
n_row_key_bad = 0
for y in years:
dist = load(f"city/distribution-{y}.json")
if dist is not None:
Expand All @@ -300,6 +454,23 @@ def load(rel: str):
rows = load(f"rankings/{kind}-{y}.json")
if rows is None:
continue
# Shape pre-pass, deliberately BEFORE the value checks. Two
# reasons: the r["days"]/r["slug"] accesses below would raise
# KeyError on a malformed row and kill validate with a traceback
# instead of a FAIL; and the `break`s below stop at the first bad
# row, so a key regression further down would never be seen.
want = ARTIFACT_KEYS_PT_RANK if kind == "pt" else ARTIFACT_KEYS_ST_RANK
file_bad = False
for i, r in enumerate(rows):
if not isinstance(r, dict) or set(r) != want:
file_bad = True
n_row_key_bad += 1
if len(row_key_bad) < 5:
diff = (sorted(set(r) ^ want) if isinstance(r, dict)
else type(r).__name__)
row_key_bad.append(f"{kind}-{y} row {i}: {diff}")
if file_bad:
continue
days_seq = [r["days"] for r in rows]
if days_seq != sorted(days_seq, reverse=True):
rank_bad.append(f"{kind}-{y} not sorted desc")
Expand All @@ -320,6 +491,9 @@ def load(rel: str):
f["geometry"]["type"] != "Point" or len(f["geometry"]["coordinates"]) != 2
for f in feats):
out.append((FAIL, "map_geojson", f"pt-{y}.geojson malformed"))
out.append((FAIL if row_key_bad else PASS, "rankings_row_keys",
f"{n_row_key_bad} rows: " + "; ".join(row_key_bad) if row_key_bad
else f"ranking rows carry the exact contract keys for {len(years)} years"))
out.append((FAIL if rank_bad else PASS, "rankings_integrity",
"; ".join(rank_bad[:5]) if rank_bad
else f"rankings resolvable + sorted for {len(years)} years"))
Expand Down
10 changes: 7 additions & 3 deletions tests/test_publish_shapes.py
Original file line number Diff line number Diff line change
Expand Up @@ -11,15 +11,19 @@

import pytest

from pipeline import validate

DB = Path("db/termo.db")
HARTA = Path("data/harta.html")

pytestmark = pytest.mark.skipif(not DB.exists() or not HARTA.exists(),
reason="real db/harta not present")

PT_KEYS = {"slug", "name", "sector", "days", "days_avarie", "days_programat",
"days_deficienta", "episodes", "longest_days", "est_day_eq", "delta_prev"}
ST_KEYS = (PT_KEYS - {"sector"}) | {"sectors", "pt_slugs"}
# Imported rather than redeclared: pipeline/validate.py enforces these same key
# sets on CI (where this file is skipped for lack of a db), so two copies would
# be free to drift apart silently.
PT_KEYS = set(validate.ARTIFACT_KEYS_PT_RANK)
ST_KEYS = set(validate.ARTIFACT_KEYS_ST_RANK)


@pytest.fixture(scope="module")
Expand Down
Loading
Loading