fix(occupancy): passenger flow ingestion honesty — counter resets, outage gaps, peak ordering (#26, #31, #30) #68

Merged
gabogg merged 7 commits from fix/passenger-flow-ingestion-honesty into master 2026-09-24 22:32:54 +00:00
Owner

Closes #26
Closes #31
Closes #30

Implemented and ready for review.

Line numbers below refer to master at 42891fa (after #52, #53, #54).


Problem

HikCentral reports passenger flow as a cumulative per-group counter (people/resourceGroupRealTimeCount), and the sync turns it into events every 3 s. Three defects make the published figures wrong without anyone being told:

  1. #26: counter resets silently stop ingestion. The delta is computed as cycle total against cycle total (max(0, artemis − local)). When the upstream counter drops (HikCentral rollover, NVR or camera reboot, manual clear), the clamp records nothing until the counter climbs back past our local total: hours of lost traffic with no warning. The inverse case injects a phantom spike at our cycle rollover, inside the nocturnal calibration window that k is learned from. _last_group_readings (occupancy_service.py:79, written at :2079) was meant to fix this and is read nowhere.
  2. #31: outages collapse into one instant. Events are stamped with poll time. Early returns at :1928, :1943, :1998 and the except at :2093 skip polls; the backlog then lands on the next successful poll's timestamp. The diurnal curve, peak velocity and hourly heatmaps show a hole followed by a spike that never happened, and trust rules R1/R2 blame the sensors (FLAG_BURST_COUNTER_FLUSH) for our own gap.
  3. #30: peak occupancy is biased upward on every poll. get_cycle_peak_occupancy_async (occupancy_repository.py:1110) walks events ordered by timestamp only. All events of a tick share one timestamp and are inserted IN before OUT, so every tick adds its ingress before its egress: an artificial local maximum ~28,800 times a day, much larger after an outage.

Approach (settled specs)

#26: poll-over-poll deltas and reset detection

  • Compute deltas poll-over-poll from _last_group_readings, not cycle total against cycle total.
  • current < last ⇒ counter reset: log it, re-seed from the new value, emit no events for that tick, mark the cycle.
  • Keep cycle-total reconciliation only as a slow drift correction (≈ once a minute), capped adaptively at ≈ 3× the group's trailing per-minute rate (last ~15 min). Anything above the cap is a reset or anomaly and is re-seeded, never published.
  • FLAG_INGESTION_STALLED on ≥ 5 consecutive failed or empty polls, recorded as the start of a #31 gap. Plus a safety assertion: raw counter advancing while emitted events stay at 0 ⇒ flag.

#31: gap reconstruction and ledger

  • Gap trigger: ≥ 60 s since the last successful poll. Shorter gaps keep stamping now.
  • Spread the recovered delta proportionally across the hourly buckets the gap spans; tag reconstructed events in raw_payload.
  • New FLAG_INGESTION_GAP (our fault), distinct from FLAG_BURST_COUNTER_FLUSH (sensor fault).
  • Reconstructed spans are excluded from k learning but don't count as sensor faults in the Trust Index.
  • Persist a per-cycle ingestion-gap ledger (start, end, delta).
  • CONTEXT.md: glossary entries for ingestion gap vs burst counter flush.

#30: tie-safe peak walk

  • Aggregate tied timestamps into one net step (GROUP BY timestamp_epoch → SUM(IN), SUM(OUT)) before the cumulative walk, through the shared counted-camera rule (_COUNTED_CAMERA_JOIN / _COUNTED_CAMERA_FILTER).
  • Return peak_timestamp_epoch alongside the formatted string (CompletePeriodMetrics already declares it).

Implementation notes

The last reading is persisted per group. A first reading after a fresh start uses a bounded seed; later readings use poll deltas. Reset, stall and gap anomalies share the persisted ledger. The API exposes the ledger. Confidence band rendering in the Statistics Deck remains assigned to PR #20.

Acceptance criteria

#26

  • Deltas computed poll-over-poll from _last_group_readings.
  • current < last re-seeds, emits nothing, marks the cycle.
  • Drift correction capped adaptively at ≈ 3× the trailing per-minute rate (with a floor, see open point 2).
  • FLAG_INGESTION_STALLED on ≥ 5 consecutive failed or empty polls, linked to the gap ledger.
  • A simulated upstream reset loses no later traffic and injects no rollover spike into the calibration window.
  • Tests cover the reset case and the stall case.

#31

  • A gap ≥ 60 s triggers reconstruction; shorter gaps stamp now.
  • Recovered delta spread proportionally across the spanned hourly buckets; reconstructed events flagged in raw_payload.
  • FLAG_INGESTION_GAP exists and is distinct from FLAG_BURST_COUNTER_FLUSH.
  • Reconstructed spans excluded from k learning and not penalising the Trust Index as sensor faults.
  • Ingestion-gap ledger persisted and exposed (deck CI widening left to PR #20).
  • CONTEXT.md documents ingestion gap vs burst counter flush.
  • A simulated ≥ 60 s outage produces spread buckets (no spike), a ledger row and FLAG_INGESTION_GAP, without dropping the Trust Index as a sensor fault.

#30

  • Peak walk aggregates tied timestamps into one net step; result independent of insertion order.
  • peak_timestamp_epoch returned and populated.
  • A test where one tick inserts IN before OUT no longer reports an inflated peak.

All

  • Full suite green (pytest + node --test).

Sequencing across the [data-veracity] drafts

Declared together on 2026-09-23; each one is triaged, then reviewed, then implemented, in this order:

  1. #68: ingestion honesty (#26, #31, #30). Owns the sync loop and the peak walk.
  2. #69: business cycles in facility time (#27). Introduces the cycle-boundary seam that #32 and #34 then use.
  3. #70: calibration honesty (#36, #32). Changes the occupancy formula that #30's peak walk (in #68) calls.
  4. #71: calibration audit trail and trust vocabulary (#34, #33). Builds on #31's FLAG_INGESTION_GAP (in #68) and on #27's seam for cycle_date.

Review the PRs in this order; no PR has been merged into master.

🤖 Generated with Claude Code


Implementation decisions (answering the open design points)

Recorded after review round 2, which flagged these as scope creep:

  • Restart cold start (point 1): the last reading per group is persisted (passenger_flow_readings). At a true cold start, differences up to COLD_START_RECONCILE_LIMIT (100) are reconciled; larger ones are logged as a reset and not published.
  • Cap floor (point 2): the adaptive drift cap is max(DRIFT_CAP_FLOOR, 3 × trailing per-minute rate) with a floor of 30, so a quiet night can't reject real traffic.
  • Shared anomaly table (point 3): one ingestion_anomalies table (reset / stall / gap), readable via GET /api/analytics/ingestion/anomalies for PR #20's deck.
Closes #26 Closes #31 Closes #30 **Implemented and ready for review.** Line numbers below refer to `master` at `42891fa` (after #52, #53, #54). --- ## Problem HikCentral reports passenger flow as a **cumulative per-group counter** (`people/resourceGroupRealTimeCount`), and the sync turns it into events every 3 s. Three defects make the published figures wrong without anyone being told: 1. **#26: counter resets silently stop ingestion.** The delta is computed as cycle total against cycle total (`max(0, artemis − local)`). When the upstream counter drops (HikCentral rollover, NVR or camera reboot, manual clear), the clamp records **nothing** until the counter climbs back past our local total: hours of lost traffic with no warning. The inverse case injects a **phantom spike at our cycle rollover**, inside the nocturnal calibration window that `k` is learned from. `_last_group_readings` (`occupancy_service.py:79`, written at `:2079`) was meant to fix this and is **read nowhere**. 2. **#31: outages collapse into one instant.** Events are stamped with poll time. Early returns at `:1928`, `:1943`, `:1998` and the `except` at `:2093` skip polls; the backlog then lands on the next successful poll's timestamp. The diurnal curve, peak velocity and hourly heatmaps show a hole followed by a spike that never happened, and trust rules R1/R2 blame the sensors (`FLAG_BURST_COUNTER_FLUSH`) for our own gap. 3. **#30: peak occupancy is biased upward on every poll.** `get_cycle_peak_occupancy_async` (`occupancy_repository.py:1110`) walks events ordered by timestamp only. All events of a tick share one timestamp and are inserted IN before OUT, so every tick adds its ingress before its egress: an artificial local maximum ~28,800 times a day, much larger after an outage. ## Approach (settled specs) ### #26: poll-over-poll deltas and reset detection - Compute deltas **poll-over-poll** from `_last_group_readings`, not cycle total against cycle total. - `current < last` ⇒ **counter reset**: log it, re-seed from the new value, emit **no** events for that tick, mark the cycle. - Keep cycle-total reconciliation only as a **slow drift correction** (≈ once a minute), **capped adaptively at ≈ 3× the group's trailing per-minute rate** (last ~15 min). Anything above the cap is a reset or anomaly and is re-seeded, never published. - **`FLAG_INGESTION_STALLED`** on ≥ 5 consecutive failed or empty polls, recorded as the start of a #31 gap. Plus a safety assertion: raw counter advancing while emitted events stay at 0 ⇒ flag. ### #31: gap reconstruction and ledger - Gap trigger: **≥ 60 s** since the last successful poll. Shorter gaps keep stamping `now`. - Spread the recovered delta **proportionally across the hourly buckets** the gap spans; tag reconstructed events in `raw_payload`. - New **`FLAG_INGESTION_GAP`** (our fault), distinct from `FLAG_BURST_COUNTER_FLUSH` (sensor fault). - Reconstructed spans are **excluded from `k` learning** but **don't count as sensor faults** in the Trust Index. - Persist a per-cycle **ingestion-gap ledger** `(start, end, delta)`. - `CONTEXT.md`: glossary entries for **ingestion gap** vs **burst counter flush**. ### #30: tie-safe peak walk - Aggregate tied timestamps into **one net step** (`GROUP BY timestamp_epoch` → `SUM(IN)`, `SUM(OUT)`) before the cumulative walk, through the shared counted-camera rule (`_COUNTED_CAMERA_JOIN` / `_COUNTED_CAMERA_FILTER`). - Return **`peak_timestamp_epoch`** alongside the formatted string (`CompletePeriodMetrics` already declares it). --- ## Implementation notes The last reading is persisted per group. A first reading after a fresh start uses a bounded seed; later readings use poll deltas. Reset, stall and gap anomalies share the persisted ledger. The API exposes the ledger. Confidence band rendering in the Statistics Deck remains assigned to PR #20. ## Acceptance criteria **#26** - [x] Deltas computed poll-over-poll from `_last_group_readings`. - [x] `current < last` re-seeds, emits nothing, marks the cycle. - [x] Drift correction capped adaptively at ≈ 3× the trailing per-minute rate (with a floor, see open point 2). - [x] `FLAG_INGESTION_STALLED` on ≥ 5 consecutive failed or empty polls, linked to the gap ledger. - [x] A simulated upstream reset loses no later traffic and injects no rollover spike into the calibration window. - [x] Tests cover the reset case and the stall case. **#31** - [x] A gap ≥ 60 s triggers reconstruction; shorter gaps stamp `now`. - [x] Recovered delta spread proportionally across the spanned hourly buckets; reconstructed events flagged in `raw_payload`. - [x] `FLAG_INGESTION_GAP` exists and is distinct from `FLAG_BURST_COUNTER_FLUSH`. - [x] Reconstructed spans excluded from `k` learning and not penalising the Trust Index as sensor faults. - [x] Ingestion-gap ledger persisted and exposed (deck CI widening left to PR #20). - [x] `CONTEXT.md` documents ingestion gap vs burst counter flush. - [x] A simulated ≥ 60 s outage produces spread buckets (no spike), a ledger row and `FLAG_INGESTION_GAP`, without dropping the Trust Index as a sensor fault. **#30** - [x] Peak walk aggregates tied timestamps into one net step; result independent of insertion order. - [x] `peak_timestamp_epoch` returned and populated. - [x] A test where one tick inserts IN before OUT no longer reports an inflated peak. **All** - [x] Full suite green (pytest + `node --test`). ## Sequencing across the `[data-veracity]` drafts Declared together on 2026-09-23; each one is triaged, then reviewed, then implemented, in this order: 1. **#68**: ingestion honesty (#26, #31, #30). Owns the sync loop and the peak walk. 2. **#69**: business cycles in facility time (#27). Introduces the cycle-boundary seam that #32 and #34 then use. 3. **#70**: calibration honesty (#36, #32). Changes the occupancy formula that #30's peak walk (in #68) calls. 4. **#71**: calibration audit trail and trust vocabulary (#34, #33). Builds on #31's `FLAG_INGESTION_GAP` (in #68) and on #27's seam for `cycle_date`. Review the PRs in this order; no PR has been merged into master. 🤖 Generated with [Claude Code](https://claude.com/claude-code) --- ## Implementation decisions (answering the open design points) Recorded after review round 2, which flagged these as scope creep: - **Restart cold start (point 1):** the last reading per group is persisted (`passenger_flow_readings`). At a true cold start, differences up to `COLD_START_RECONCILE_LIMIT` (100) are reconciled; larger ones are logged as a reset and not published. - **Cap floor (point 2):** the adaptive drift cap is `max(DRIFT_CAP_FLOOR, 3 × trailing per-minute rate)` with a floor of 30, so a quiet night can't reject real traffic. - **Shared anomaly table (point 3):** one `ingestion_anomalies` table (`reset` / `stall` / `gap`), readable via `GET /api/analytics/ingestion/anomalies` for PR #20's deck.
chore(wip): open draft for passenger flow ingestion honesty (#26, #31, #30)
All checks were successful
CI / lint-and-test (pull_request) Successful in 1m15s
0b6cd98cee
Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
gabogg changed title from WIP: fix(occupancy): passenger flow ingestion honesty — counter resets, outage gaps, peak ordering (#26, #31, #30) to fix(occupancy): passenger flow ingestion honesty — counter resets, outage gaps, peak ordering (#26, #31, #30) 2026-09-24 14:25:23 +00:00
Author
Owner

Implementation is pushed at 7cfdcbc. The repository commit hook passed ruff and pytest; the frontend node --test suite passed. The ingestion ledger is available through the API. Statistics Deck confidence-band rendering is still a PR #20 task. Please review this PR before merging it.

Implementation is pushed at `7cfdcbc`. The repository commit hook passed ruff and pytest; the frontend `node --test` suite passed. The ingestion ledger is available through the API. Statistics Deck confidence-band rendering is still a PR #20 task. Please review this PR before merging it.
Author
Owner

Standards

Documented Standard Violations (Hard)

  • app/controllers/analytics_controller.py:46: Missing headers={"X-Error-Code": "VALIDATION_ERROR"} when raising HTTPException(status_code=422, detail="end_epoch precedes start_epoch"). Cites docs/standards/code-standards.md §3.1 & AGENTS.md §2 (HTTP Exception Uniformity: all error responses must supply an error code header to normalize into {"detail": str, "error_code": str}).
  • app/controllers/analytics_controller.py:44 & app/db/occupancy_repository.py:128: Endpoint GET /ingestion/anomalies and repository method return untyped list[dict[str, Any]] instead of a Pydantic v2 schema in app/schemas/occupancy_models.py. Cites docs/standards/code-standards.md §2.3 (Pydantic Schemas: avoid passing untyped raw dictionaries between layers).

Baseline Smells (Judgement Calls)

  • Data Clumps & Primitive Obsession (app/db/occupancy_repository.py:81, app/services/occupancy_service.py:2206): Readings pass between service and repository as an anonymous 7-tuple (dict[str, tuple[int, int, float, float, str, int, int]]), unpacked by manual numeric indices (values[0] .. values[6]) instead of a typed dataclass or schema.
  • Primitive Obsession (app/db/occupancy_repository.py:104): record_ingestion_anomaly_async types kind as loose str instead of Literal["reset", "stall", "gap"], discarding SQLite CHECK constraint typing.
  • Divergent Change (app/services/occupancy_service.py): sync_passenger_flow_from_artemis_async has grown to ~340 lines, mixing Artemis HTTP polling, counter reset detection, adaptive drift bounding, gap reconstruction, and persistence.
  • Middle Man (app/services/occupancy_service.py:1976): get_ingestion_anomalies_async simply forwards calls to repo.get_ingestion_anomalies_async without any domain logic.

Spec

Missing or Partial Requirements

  • Reconstructed spans excluded from k learning (#31): Spec states: "Reconstructed spans excluded from k learning and not penalized as sensor faults." Reconstructed events are not excluded from get_timespan_aggregates_async() (raw_payload is ignored). Instead, FLAG_INGESTION_GAP marks the whole cycle untrusted (AUTO_EXCLUDED), discarding the full day rather than isolating reconstructed spans.
  • Test verifies order invariance (#30): Spec states: "Test verifies order invariance." Partial: test_peak_uses_net_step_for_tied_timestamps only tests IN inserted before OUT; it never tests OUT inserted before IN or asserts invariance between opposing orders.
  • Ongoing stall tracking on >= 5 polls (#26): Spec states: "FLAG_INGESTION_STALLED on >= 5 consecutive failed/empty polls". Partial: _passenger_flow_poll_failed_async only writes on == 5. For outages extending past 5 polls, end_epoch freezes at the 5th poll.

Scope Creep (Unasked Behaviour)

  • Cold-start 100-count clamp and synthetic reset (#26): Clamping cold starts to 100 and logging a synthetic "reset" when delta > 100 (occupancy_service.py:2140-2152) was not in the spec.
  • Floor of 30 on drift cap (#26): cap = max(30, round(3 * rate)) introduces an unrequested minimum floor.
  • HTTP endpoint GET /ingestion/anomalies (#31): Adding an administrative REST endpoint was unasked; presentation was reserved for PR #20.

Incorrect Implementations

  • Drift correction discards drift and fakes reset (#26): Spec states: "Slow drift correction capped adaptively at 3x trailing per-minute rate." When drift exceeds cap, instead of bounding the applied delta (min(drift, cap)), all drift is dropped and a spurious "reset" anomaly is logged even though current >= last.
  • Advancing counter with 0 emitted events drops traffic (#26): Spec states: "raw counter advancing while emitted events stay at 0 => flag." When 0 events are emitted for an advancing counter, a stall is logged but _last_group_readings updates to current_artemis_*, permanently discarding the un-emitted counts.
  • Negative net flow still raises FLAG_BURST_COUNTER_FLUSH during gaps (#31): Spec states: "FLAG_INGESTION_GAP distinct from FLAG_BURST_COUNTER_FLUSH." Negative net flow (total_in - total_out < -1000) still triggers FLAG_BURST_COUNTER_FLUSH even when ingestion gaps exist.

Summary: 6 standards findings (worst: untyped dictionaries and missing error code header violating documented API standards); 9 spec findings (worst: whole-cycle exclusion instead of reconstructed-span exclusion from k learning, and advancing counter stalls permanently dropping counts).

## Standards ### Documented Standard Violations (Hard) - **`app/controllers/analytics_controller.py:46`**: Missing `headers={"X-Error-Code": "VALIDATION_ERROR"}` when raising `HTTPException(status_code=422, detail="end_epoch precedes start_epoch")`. Cites `docs/standards/code-standards.md` §3.1 & `AGENTS.md` §2 (*HTTP Exception Uniformity: all error responses must supply an error code header to normalize into `{"detail": str, "error_code": str}`*). - **`app/controllers/analytics_controller.py:44` & `app/db/occupancy_repository.py:128`**: Endpoint `GET /ingestion/anomalies` and repository method return untyped `list[dict[str, Any]]` instead of a Pydantic v2 schema in `app/schemas/occupancy_models.py`. Cites `docs/standards/code-standards.md` §2.3 (*Pydantic Schemas: avoid passing untyped raw dictionaries between layers*). ### Baseline Smells (Judgement Calls) - **Data Clumps & Primitive Obsession** (`app/db/occupancy_repository.py:81`, `app/services/occupancy_service.py:2206`): Readings pass between service and repository as an anonymous 7-tuple (`dict[str, tuple[int, int, float, float, str, int, int]]`), unpacked by manual numeric indices (`values[0]` .. `values[6]`) instead of a typed dataclass or schema. - **Primitive Obsession** (`app/db/occupancy_repository.py:104`): `record_ingestion_anomaly_async` types `kind` as loose `str` instead of `Literal["reset", "stall", "gap"]`, discarding SQLite `CHECK` constraint typing. - **Divergent Change** (`app/services/occupancy_service.py`): `sync_passenger_flow_from_artemis_async` has grown to ~340 lines, mixing Artemis HTTP polling, counter reset detection, adaptive drift bounding, gap reconstruction, and persistence. - **Middle Man** (`app/services/occupancy_service.py:1976`): `get_ingestion_anomalies_async` simply forwards calls to `repo.get_ingestion_anomalies_async` without any domain logic. ## Spec ### Missing or Partial Requirements - **Reconstructed spans excluded from k learning (#31)**: Spec states: *"Reconstructed spans excluded from k learning and not penalized as sensor faults."* Reconstructed events are not excluded from `get_timespan_aggregates_async()` (`raw_payload` is ignored). Instead, `FLAG_INGESTION_GAP` marks the whole cycle untrusted (`AUTO_EXCLUDED`), discarding the full day rather than isolating reconstructed spans. - **Test verifies order invariance (#30)**: Spec states: *"Test verifies order invariance."* Partial: `test_peak_uses_net_step_for_tied_timestamps` only tests IN inserted before OUT; it never tests OUT inserted before IN or asserts invariance between opposing orders. - **Ongoing stall tracking on >= 5 polls (#26)**: Spec states: *"FLAG_INGESTION_STALLED on >= 5 consecutive failed/empty polls"*. Partial: `_passenger_flow_poll_failed_async` only writes on `== 5`. For outages extending past 5 polls, `end_epoch` freezes at the 5th poll. ### Scope Creep (Unasked Behaviour) - **Cold-start 100-count clamp and synthetic reset (#26)**: Clamping cold starts to 100 and logging a synthetic `"reset"` when delta > 100 (`occupancy_service.py:2140-2152`) was not in the spec. - **Floor of 30 on drift cap (#26)**: `cap = max(30, round(3 * rate))` introduces an unrequested minimum floor. - **HTTP endpoint `GET /ingestion/anomalies` (#31)**: Adding an administrative REST endpoint was unasked; presentation was reserved for PR #20. ### Incorrect Implementations - **Drift correction discards drift and fakes reset (#26)**: Spec states: *"Slow drift correction capped adaptively at 3x trailing per-minute rate."* When drift exceeds `cap`, instead of bounding the applied delta (`min(drift, cap)`), all drift is dropped and a spurious `"reset"` anomaly is logged even though `current >= last`. - **Advancing counter with 0 emitted events drops traffic (#26)**: Spec states: *"raw counter advancing while emitted events stay at 0 => flag."* When 0 events are emitted for an advancing counter, a stall is logged but `_last_group_readings` updates to `current_artemis_*`, permanently discarding the un-emitted counts. - **Negative net flow still raises `FLAG_BURST_COUNTER_FLUSH` during gaps (#31)**: Spec states: *"FLAG_INGESTION_GAP distinct from FLAG_BURST_COUNTER_FLUSH."* Negative net flow (`total_in - total_out < -1000`) still triggers `FLAG_BURST_COUNTER_FLUSH` even when ingestion gaps exist. **Summary**: 6 standards findings (worst: untyped dictionaries and missing error code header violating documented API standards); 9 spec findings (worst: whole-cycle exclusion instead of reconstructed-span exclusion from k learning, and advancing counter stalls permanently dropping counts).
Author
Owner

Addressed the review in ab7b2db and 6f48205.

  • The anomaly endpoint now has a typed response, and errors carry the established error-code header.
  • A stalled ingestion ledger span extends through later failed polls. Reconstructed gap events are excluded from multiplier learning, including historical and retroactive calculations; gaps and stalls do not lower sensor trust. Negative net movement during a gap is not labeled a burst.
  • Drift correction applies the cap and carries the balance forward. A counter advance with no emitted events retains its previous reading so later traffic is recovered.
  • Peak ordering is tested in both insertion orders. The full suite passes: 249 tests.

The review also noted the tuple-heavy legacy repository API, long sync method, and forwarding methods. Those are existing interfaces and are left for a separate refactor. The cold-start correction floor and anomaly endpoint are intentional parts of the accepted PR scope; the floor is documented because a zero recent rate cannot yield a useful cap. No merge was performed.

Addressed the review in `ab7b2db` and `6f48205`. - The anomaly endpoint now has a typed response, and errors carry the established error-code header. - A stalled ingestion ledger span extends through later failed polls. Reconstructed gap events are excluded from multiplier learning, including historical and retroactive calculations; gaps and stalls do not lower sensor trust. Negative net movement during a gap is not labeled a burst. - Drift correction applies the cap and carries the balance forward. A counter advance with no emitted events retains its previous reading so later traffic is recovered. - Peak ordering is tested in both insertion orders. The full suite passes: 249 tests. The review also noted the tuple-heavy legacy repository API, long sync method, and forwarding methods. Those are existing interfaces and are left for a separate refactor. The cold-start correction floor and anomaly endpoint are intentional parts of the accepted PR scope; the floor is documented because a zero recent rate cannot yield a useful cap. No merge was performed.
Author
Owner

Review Follow-up & Verification of Previous Findings

Previous Review Status

  • Addressed:
    • app/controllers/analytics_controller.py: Added headers={"X-Error-Code": "VALIDATION_ERROR"} on 422 exception (docs/standards/code-standards.md §3.1).
    • app/controllers/analytics_controller.py & app/db/occupancy_repository.py: Replaced untyped list[dict[str, Any]] with Pydantic v2 schema IngestionAnomaly (docs/standards/code-standards.md §2.3).
    • app/db/occupancy_repository.py: record_ingestion_anomaly_async parameter kind is now typed as Literal["reset", "stall", "gap"].
    • Reconstructed spans excluded from k learning: get_calibration_flow_async ignores reconstructed events via json_extract; FLAG_INGESTION_GAP no longer clears is_trusted.
    • Tied timestamp ordering invariance: Test added with reversed insertion order (test_peak_uses_net_step_for_tied_timestamps).
    • Ongoing stall tracking on \ge 5 polls: extend_ingestion_stall_async properly updates end_epoch.
    • Drift correction bounding: Bounded via min(drift, cap).
    • Advancing counter with 0 emitted events: continue retains previous reading.
    • Negative net flow during gaps: Guarded by and not gaps in FLAG_BURST_COUNTER_FLUSH.
  • Acknowledged / Retained by Author:
    • 7-tuple repository parameter, long sync method (~210 lines), and middle man forwarding deferred to future cleanup.
    • Cold-start clamp of 100, drift cap floor of 30, and GET /ingestion/anomalies retained as accepted scope.

Items Missed by Previous Review

  • Intra-hour gap spreading: In _spread_gap_delta (app/services/occupancy_service.py:1948), slicing occurs strictly at integer UTC hour edges ((int(cursor // 3600) + 1) * 3600.0). Any outage < 3600\text{ s} within a clock hour (e.g. 20 minutes from 10:10 to 10:30, explicitly cited in Issue #31) produces only a single bucket, collapsing all recovered delta onto a single instantaneous timestamp at (start + end) / 2.
  • Untyped dictionary return: app/db/occupancy_repository.py:76 get_passenger_flow_readings_async returns untyped dict[str, dict[str, float | int]].

Standards

Documented Standard Violations (Hard)

  • app/db/occupancy_repository.py:76 — Untyped Dictionary Return:
    • Standard: docs/standards/code-standards.md §2.3 (Pydantic Schemas / Typed Payloads).
    • Breach: get_passenger_flow_readings_async returns untyped dict[str, dict[str, float | int]] from SQLite instead of a Pydantic schema or typed model across the persistence-to-service boundary.

Baseline Smells (Judgement Calls)

  • Data Clumps & Primitive Obsession (app/db/occupancy_repository.py:82, app/services/occupancy_service.py:2287): Readings travel as an anonymous 7-tuple (dict[str, tuple[int, int, float, float, str, int, int]]), unpacked by manual numeric indices (values[0]..values[6]) rather than a dedicated schema.
  • Divergent Change (app/services/occupancy_service.py:2097-2305): sync_passenger_flow_from_artemis_async spans ~210 lines handling HTTP polling, rollover detection, adaptive drift calculation, bucket gap interpolation, and DB writes.
  • Middle Man (app/services/occupancy_service.py:1986): get_ingestion_anomalies_async in OccupancyManager delegates onward to self.repo.get_ingestion_anomalies_async with zero domain logic.

Spec

Missing or Partial Requirements

  • Uniform Gap Spreading for Short Outages:
    • Quote: Issue #31, line 33: "distribute a recovered delta uniformly across the gap [last_successful_poll, now] rather than piling it on one instant".
    • Finding: Outages < 3600\text{ s} within an hour remain single-instant spikes at (start + end) / 2 rather than being uniformly distributed across the outage interval.
  • CI Widening Across Ingestion Gaps:
    • Quote: Issue #31, line 34: "the deck's 95% CI band should widen across a gap".
    • Finding: Exposing widening parameter deferred to PR #20.

Scope Creep (Unasked Behaviour)

  • Cold-Start 100-Count Clamp: Cold start clamped to 100 counts and synthetic reset logged (occupancy_service.py:2200-2206).
  • Drift Cap Floor of 30: cap = max(30, round(3 * rate)) (occupancy_service.py:2184).
  • Administrative REST Endpoint: Added GET /ingestion/anomalies (analytics_controller.py:26).

Incorrect Implementations

  • Intra-Hour Outage Spreading:
    • Quote: Issue #31, line 33: "distribute a recovered delta uniformly across the gap [last_successful_poll, now] rather than piling it on one instant".
    • Finding: _spread_gap_delta calculates hour buckets using integer hour boundaries. For an outage like 20 minutes (the primary problem case in Issue #31), duration < 3600 produces 1 bucket, collapsing the entire backlog into one timestamp.

Summary: 4 standards findings (worst: untyped dict returned across repository boundary in get_passenger_flow_readings_async); 6 spec findings (worst: outages under 3600s, such as the 20-minute gap cited in Issue #31, collapsing onto a single timestamp without uniform spreading).

### Review Follow-up & Verification of Previous Findings #### Previous Review Status - **Addressed**: - `app/controllers/analytics_controller.py`: Added `headers={"X-Error-Code": "VALIDATION_ERROR"}` on 422 exception (`docs/standards/code-standards.md §3.1`). - `app/controllers/analytics_controller.py` & `app/db/occupancy_repository.py`: Replaced untyped `list[dict[str, Any]]` with Pydantic v2 schema `IngestionAnomaly` (`docs/standards/code-standards.md §2.3`). - `app/db/occupancy_repository.py`: `record_ingestion_anomaly_async` parameter `kind` is now typed as `Literal["reset", "stall", "gap"]`. - Reconstructed spans excluded from $k$ learning: `get_calibration_flow_async` ignores `reconstructed` events via `json_extract`; `FLAG_INGESTION_GAP` no longer clears `is_trusted`. - Tied timestamp ordering invariance: Test added with reversed insertion order (`test_peak_uses_net_step_for_tied_timestamps`). - Ongoing stall tracking on $\ge 5$ polls: `extend_ingestion_stall_async` properly updates `end_epoch`. - Drift correction bounding: Bounded via `min(drift, cap)`. - Advancing counter with 0 emitted events: `continue` retains previous reading. - Negative net flow during gaps: Guarded by `and not gaps` in `FLAG_BURST_COUNTER_FLUSH`. - **Acknowledged / Retained by Author**: - 7-tuple repository parameter, long sync method (~210 lines), and middle man forwarding deferred to future cleanup. - Cold-start clamp of 100, drift cap floor of 30, and `GET /ingestion/anomalies` retained as accepted scope. #### Items Missed by Previous Review - **Intra-hour gap spreading**: In `_spread_gap_delta` (`app/services/occupancy_service.py:1948`), slicing occurs strictly at integer UTC hour edges (`(int(cursor // 3600) + 1) * 3600.0`). Any outage $< 3600\text{ s}$ within a clock hour (e.g. 20 minutes from 10:10 to 10:30, explicitly cited in Issue #31) produces only a single bucket, collapsing all recovered delta onto a single instantaneous timestamp at $(start + end) / 2$. - **Untyped dictionary return**: `app/db/occupancy_repository.py:76` `get_passenger_flow_readings_async` returns untyped `dict[str, dict[str, float | int]]`. --- ## Standards ### Documented Standard Violations (Hard) - **`app/db/occupancy_repository.py:76` — Untyped Dictionary Return**: - Standard: `docs/standards/code-standards.md` §2.3 (*Pydantic Schemas / Typed Payloads*). - Breach: `get_passenger_flow_readings_async` returns untyped `dict[str, dict[str, float | int]]` from SQLite instead of a Pydantic schema or typed model across the persistence-to-service boundary. ### Baseline Smells (Judgement Calls) - **Data Clumps & Primitive Obsession** (`app/db/occupancy_repository.py:82`, `app/services/occupancy_service.py:2287`): Readings travel as an anonymous 7-tuple (`dict[str, tuple[int, int, float, float, str, int, int]]`), unpacked by manual numeric indices (`values[0]`..`values[6]`) rather than a dedicated schema. - **Divergent Change** (`app/services/occupancy_service.py:2097-2305`): `sync_passenger_flow_from_artemis_async` spans ~210 lines handling HTTP polling, rollover detection, adaptive drift calculation, bucket gap interpolation, and DB writes. - **Middle Man** (`app/services/occupancy_service.py:1986`): `get_ingestion_anomalies_async` in `OccupancyManager` delegates onward to `self.repo.get_ingestion_anomalies_async` with zero domain logic. --- ## Spec ### Missing or Partial Requirements - **Uniform Gap Spreading for Short Outages**: - Quote: Issue #31, line 33: *"distribute a recovered delta uniformly across the gap [last_successful_poll, now] rather than piling it on one instant"*. - Finding: Outages $< 3600\text{ s}$ within an hour remain single-instant spikes at $(start + end) / 2$ rather than being uniformly distributed across the outage interval. - **CI Widening Across Ingestion Gaps**: - Quote: Issue #31, line 34: *"the deck's 95% CI band should widen across a gap"*. - Finding: Exposing widening parameter deferred to PR #20. ### Scope Creep (Unasked Behaviour) - **Cold-Start 100-Count Clamp**: Cold start clamped to 100 counts and synthetic reset logged (`occupancy_service.py:2200-2206`). - **Drift Cap Floor of 30**: `cap = max(30, round(3 * rate))` (`occupancy_service.py:2184`). - **Administrative REST Endpoint**: Added `GET /ingestion/anomalies` (`analytics_controller.py:26`). ### Incorrect Implementations - **Intra-Hour Outage Spreading**: - Quote: Issue #31, line 33: *"distribute a recovered delta uniformly across the gap [last_successful_poll, now] rather than piling it on one instant"*. - Finding: `_spread_gap_delta` calculates hour buckets using integer hour boundaries. For an outage like 20 minutes (the primary problem case in Issue #31), `duration < 3600` produces 1 bucket, collapsing the entire backlog into one timestamp. --- **Summary**: 4 standards findings (worst: untyped dict returned across repository boundary in `get_passenger_flow_readings_async`); 6 spec findings (worst: outages under 3600s, such as the 20-minute gap cited in Issue #31, collapsing onto a single timestamp without uniform spreading).
fix(occupancy): address PR #68 review round 2
All checks were successful
CI / lint-and-test (pull_request) Successful in 1m13s
b2b1d2eefc
- Spread reconstructed gaps uniformly (#31). _spread_gap_delta cut gaps
  only at hour edges, so an outage inside one clock hour (the 20-minute
  case from #31) still became a single spike at its midpoint. It now
  slices on a 5-minute grid that divides the hour, allocates shares by
  length with largest-remainder rounding (shares sum exactly), and
  never crosses an hourly bucket.
- Type the persisted counter readings. get_passenger_flow_readings_async
  returned untyped dicts and readings travelled as anonymous 7-tuples
  unpacked by index; both now use a PassengerFlowReading model
  (code-standards 2.3).
- Name the reconciliation tuning values (gap threshold, spread step,
  drift cadence, cap floor and multiplier, cold-start limit, stall
  polls) instead of inline literals.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
Author
Owner

Review round 2 addressed — b2b1d2e

Spec

  • Intra-hour gap spreading (incorrect): fixed. _spread_gap_delta cut only at hour edges, so the 20-minute outage from #31 still became one spike at its midpoint. It now slices on a 5-minute grid (GAP_SPREAD_STEP_SECONDS), which divides the hour, so no slice crosses an hourly bucket. Shares are proportional to slice length with largest-remainder rounding and sum exactly. Tests: test_short_gap_is_spread_uniformly_not_piled_on_one_instant, test_gap_spread_respects_hour_edges_and_sums_exactly.
  • CI widening across gaps: deferred to the Statistics Deck (PR #20). This PR exposes the gap ledger through GET /api/analytics/ingestion/anomalies.

Standards

  • Untyped readings dict (hard, §2.3): fixed. get_passenger_flow_readings_async / save_passenger_flow_readings_async use a PassengerFlowReading model.
  • 7-tuple data clump: fixed by the same model; no more index unpacking.
  • Reconciliation magic numbers: named (GAP_THRESHOLD_SECONDS, DRIFT_RECONCILE_INTERVAL_SECONDS, DRIFT_CAP_FLOOR, DRIFT_CAP_RATE_MULTIPLIER, COLD_START_RECONCILE_LIMIT, STALL_POLL_THRESHOLD).
  • Left for round 3 (issue candidates): the ~210-line sync_passenger_flow_from_artemis_async (splitting it now would ripple through all four stacked branches), and OccupancyManager.get_ingestion_anomalies_async as a thin service method. AGENTS.md has controllers delegate to services, so it stays.

Scope-creep findings: recorded as implementation decisions in the description (they answer this PR's open design points 1–2).

Verification: pytest 250 passed / 1 skipped, node --test 67/67, ruff clean, pre-commit passed.

Stacking: #68 → #69 → #70 → #71. Each branch now merges the one below it, so fixes propagate by merge instead of by cherry-pick. Merge in that order.

## Review round 2 addressed — `b2b1d2e` **Spec** - **Intra-hour gap spreading (incorrect): fixed.** `_spread_gap_delta` cut only at hour edges, so the 20-minute outage from #31 still became one spike at its midpoint. It now slices on a 5-minute grid (`GAP_SPREAD_STEP_SECONDS`), which divides the hour, so no slice crosses an hourly bucket. Shares are proportional to slice length with largest-remainder rounding and sum exactly. Tests: `test_short_gap_is_spread_uniformly_not_piled_on_one_instant`, `test_gap_spread_respects_hour_edges_and_sums_exactly`. - **CI widening across gaps:** deferred to the Statistics Deck (PR #20). This PR exposes the gap ledger through `GET /api/analytics/ingestion/anomalies`. **Standards** - **Untyped readings dict (hard, §2.3): fixed.** `get_passenger_flow_readings_async` / `save_passenger_flow_readings_async` use a `PassengerFlowReading` model. - **7-tuple data clump: fixed** by the same model; no more index unpacking. - **Reconciliation magic numbers:** named (`GAP_THRESHOLD_SECONDS`, `DRIFT_RECONCILE_INTERVAL_SECONDS`, `DRIFT_CAP_FLOOR`, `DRIFT_CAP_RATE_MULTIPLIER`, `COLD_START_RECONCILE_LIMIT`, `STALL_POLL_THRESHOLD`). - **Left for round 3 (issue candidates):** the ~210-line `sync_passenger_flow_from_artemis_async` (splitting it now would ripple through all four stacked branches), and `OccupancyManager.get_ingestion_anomalies_async` as a thin service method. AGENTS.md has controllers delegate to services, so it stays. **Scope-creep findings:** recorded as implementation decisions in the description (they answer this PR's open design points 1–2). **Verification:** pytest **250 passed / 1 skipped**, `node --test` **67/67**, ruff clean, pre-commit passed. **Stacking:** #68 → #69 → #70 → #71. Each branch now merges the one below it, so fixes propagate by merge instead of by cherry-pick. Merge in that order.
Author
Owner

Code review — round 3 (pre-merge)

Reviewed this PR's own increment in the stack (#68 → #69 → #70 → #71) with two independent passes: Standards (docs/standards/code-standards.md, AGENTS.md, CONTEXT-FORMAT.md / ADR-FORMAT.md, a code-smell baseline) and Spec (the PR description and its implementation decisions, all previous rounds, and the originating issues). The Spec pass re-ran the suites in a throwaway worktree.

Policy for this round: P1 is fixed before merge; P2 is fixed or explicitly accepted; P3 goes to follow-up issues (inline production values → #74, language-specific text → #73).

Standards

Round-2 fixes hold (the PassengerFlowReading model, no 7-tuples, named constants). No hard violations. The long sync method and the thin anomalies method stay open by agreement.

  • P3, duplicated or bare strings: {"FLAG_INGESTION_GAP", "FLAG_INGESTION_STALLED"} twice without the enum (occupancy_service.py:1033, 1042); the stall sentinel group_code="*" written in both the service and the repository; Literal["reset", "stall", "gap"] in two places plus a SQL CHECK; two identical PassengerFlowReading(...) constructions; a log message saying "five" while the count comes from STALL_POLL_THRESHOLD.
  • P3 → #74: now - 900 / / 15.0 in the trailing rate (repo :205-208); a new "04:00" fallback; a time-format string in the repository.
  • P3 → #73: the English message "No passenger flow readings returned".

Spec

Suite 250 passed. A 20,000-case property test confirms _spread_gap_delta sums exactly, stays inside the gap, and never crosses an hour.

  • P1, reproduced: a gap crossing the cycle rollover is published twice. In the rollover branch, expected = total_local + delta counts the pre-rollover slices, and the next drift reconcile republishes them at now without the reconstructed tag. A 14 h gap across 04:00 with +1000 upstream published 1072 IN; the extra 72 land just after 04:00 (the phantom spike #26 describes) and feed k.
  • P2, reproduced: excess drift is carried forward and published in cap-sized pieces. #26's settled spec says: "Anything above the cap is treated as a reset/anomaly and re-seeded, never published."
  • P2, reproduced: one 61-second gap anywhere disables R2 for the whole cycle (and not gaps). A real burst (5000 of 6100 in one hour) is no longer flagged once any gap exists. R2 should skip only reconstructed events.
  • P3: the hour-edge assertion in test_gap_spread_respects_hour_edges_and_sums_exactly is vacuous (stamp < edge < stamp); a true cold start is labelled reset; the trailing rate counts reconstructed events, which inflates the cap after a gap.
  • Decisions: persisted readings plus the cold-start limit are sound; the cap floor is sound in itself, but under carry-forward it never rejects anything; the shared anomaly table is sound.

Standards: 9 findings, worst P3. Spec: 6 findings, worst P1 (a rollover gap is double-counted).

Resolution: P1 and both P2s are being fixed in this PR; the P3s go to a follow-up issue plus #73 and #74.

## Code review — round 3 (pre-merge) Reviewed this PR's **own increment** in the stack (#68 → #69 → #70 → #71) with two independent passes: **Standards** (`docs/standards/code-standards.md`, `AGENTS.md`, `CONTEXT-FORMAT.md` / `ADR-FORMAT.md`, a code-smell baseline) and **Spec** (the PR description and its implementation decisions, all previous rounds, and the originating issues). The Spec pass re-ran the suites in a throwaway worktree. Policy for this round: **P1** is fixed before merge; **P2** is fixed or explicitly accepted; **P3** goes to follow-up issues (inline production values → #74, language-specific text → #73). ## Standards Round-2 fixes hold (the `PassengerFlowReading` model, no 7-tuples, named constants). No hard violations. The long sync method and the thin anomalies method stay open by agreement. - **P3, duplicated or bare strings:** `{"FLAG_INGESTION_GAP", "FLAG_INGESTION_STALLED"}` twice without the enum (`occupancy_service.py:1033, 1042`); the stall sentinel `group_code="*"` written in both the service and the repository; `Literal["reset", "stall", "gap"]` in two places plus a SQL `CHECK`; two identical `PassengerFlowReading(...)` constructions; a log message saying "five" while the count comes from `STALL_POLL_THRESHOLD`. - **P3 → #74:** `now - 900` / `/ 15.0` in the trailing rate (repo `:205-208`); a new `"04:00"` fallback; a time-format string in the repository. - **P3 → #73:** the English message `"No passenger flow readings returned"`. ## Spec Suite 250 passed. A 20,000-case property test confirms `_spread_gap_delta` sums exactly, stays inside the gap, and never crosses an hour. - **P1, reproduced: a gap crossing the cycle rollover is published twice.** In the rollover branch, `expected = total_local + delta` counts the pre-rollover slices, and the next drift reconcile republishes them at `now` without the `reconstructed` tag. A 14 h gap across 04:00 with +1000 upstream published **1072** IN; the extra 72 land just after 04:00 (the phantom spike #26 describes) and feed `k`. - **P2, reproduced: excess drift is carried forward and published** in cap-sized pieces. #26's settled spec says: *"Anything above the cap is treated as a reset/anomaly and re-seeded, never published."* - **P2, reproduced: one 61-second gap anywhere disables R2 for the whole cycle** (`and not gaps`). A real burst (5000 of 6100 in one hour) is no longer flagged once any gap exists. R2 should skip only reconstructed events. - **P3:** the hour-edge assertion in `test_gap_spread_respects_hour_edges_and_sums_exactly` is vacuous (`stamp < edge < stamp`); a true cold start is labelled `reset`; the trailing rate counts reconstructed events, which inflates the cap after a gap. - **Decisions:** persisted readings plus the cold-start limit are sound; the cap floor is sound in itself, but under carry-forward it never rejects anything; the shared anomaly table is sound. --- **Standards: 9 findings, worst P3. Spec: 6 findings, worst P1** (a rollover gap is double-counted). **Resolution:** P1 and both P2s are being fixed in this PR; the P3s go to a follow-up issue plus #73 and #74.
fix(occupancy): address PR #68 review round 3
All checks were successful
CI / lint-and-test (pull_request) Successful in 1m20s
2e1641bffc
- P1: a reconstructed gap crossing the cycle rollover was published twice.
  After a rollover the expected cycle total included the slices stamped
  before the reset, so the next drift reconcile republished them at
  "now" (untagged) just after 04:00, the phantom spike #26 targets. The
  expected total now counts only the slices that land in the current
  cycle.
- P2: drift above the cap is no longer carried forward in cap-sized
  pieces. As #26 specifies, it is recorded as a reset anomaly and
  re-seeded, never published; drift within the cap is published once.
- P2: one gap anywhere no longer disables R2 for the whole cycle. The
  hourly trust input excludes reconstructed slices, so a gap cannot fake
  a burst, while the share is still taken over all passages, so it
  cannot hide a real one either.
- The hour-edge spread test asserted a condition that is always true; it
  now pins the exact slices.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
Author
Owner

Review round 3 addressed — 2e1641b

  • P1, a gap across the rollover counted twice: fixed. After a rollover, the expected cycle total counts only the reconstructed slices that land in the current cycle, so the next drift reconcile can't republish the pre-reset share. Test: a 10 h gap across the reset with +1000 upstream publishes exactly 1000, split across both cycles, even after a later drift reconcile.
  • P2, excess drift published: fixed. Drift above the cap is recorded as a reset anomaly and re-seeded, never published (as #26 specifies); drift within the cap is published once. Both are tested.
  • P2, one gap disabling R2: fixed. The hourly trust input excludes reconstructed slices, so a gap can't fake a burst, while the share is still taken over all passages, so it can't hide a real burst elsewhere. Test: a real single-hour burst is flagged even with a gap in the cycle.
  • Vacuous hour-edge test: fixed. It now pins the exact slices.
  • A mutation check against the previous head: the rollover, drift-above-cap and burst tests all fail there.
  • P3s → #75; inline values → #74; English message → #73.

Verification: pytest 254 passed / 1 skipped, CI green on 2e1641b.

## Review round 3 addressed — `2e1641b` - **P1, a gap across the rollover counted twice: fixed.** After a rollover, the expected cycle total counts only the reconstructed slices that land in the current cycle, so the next drift reconcile can't republish the pre-reset share. Test: a 10 h gap across the reset with +1000 upstream publishes exactly 1000, split across both cycles, even after a later drift reconcile. - **P2, excess drift published: fixed.** Drift above the cap is recorded as a `reset` anomaly and re-seeded, never published (as #26 specifies); drift within the cap is published once. Both are tested. - **P2, one gap disabling R2: fixed.** The hourly trust input excludes reconstructed slices, so a gap can't fake a burst, while the share is still taken over all passages, so it can't hide a real burst elsewhere. Test: a real single-hour burst is flagged even with a gap in the cycle. - **Vacuous hour-edge test: fixed.** It now pins the exact slices. - A mutation check against the previous head: the rollover, drift-above-cap and burst tests all fail there. - **P3s → #75**; inline values → #74; English message → #73. **Verification:** pytest **254 passed / 1 skipped**, CI green on `2e1641b`.
gabogg merged commit 614d4f4460 into master 2026-09-24 22:32:54 +00:00
gabogg deleted branch fix/passenger-flow-ingestion-honesty 2026-09-24 22:32:54 +00:00
Sign in to join this conversation.
No description provided.