No reviewers
Labels
No labels
blocked
bug
enhancement
high-priority
low-priority
needs-info
needs-triage
ready-for-agent
ready-for-human
referenced
research
wontfix
No milestone
No project
No assignees
1 participant
Notifications
Due date
No due date set.
Dependencies
No dependencies set
Reference
gabogg/hikcentral!68
Loading…
Reference in a new issue
No description provided.
Delete branch "fix/passenger-flow-ingestion-honesty"
Deleting a branch is permanent. Although the deleted branch may continue to exist for a short time before it actually gets removed, it CANNOT be undone in most cases. Continue?
Closes #26
Closes #31
Closes #30
Implemented and ready for review.
Line numbers below refer to
masterat42891fa(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: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 thatkis learned from._last_group_readings(occupancy_service.py:79, written at:2079) was meant to fix this and is read nowhere.:1928,:1943,:1998and theexceptat:2093skip 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.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
_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.FLAG_INGESTION_STALLEDon ≥ 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
now.raw_payload.FLAG_INGESTION_GAP(our fault), distinct fromFLAG_BURST_COUNTER_FLUSH(sensor fault).klearning but don't count as sensor faults in the Trust Index.(start, end, delta).CONTEXT.md: glossary entries for ingestion gap vs burst counter flush.#30: tie-safe peak walk
GROUP BY timestamp_epoch→SUM(IN),SUM(OUT)) before the cumulative walk, through the shared counted-camera rule (_COUNTED_CAMERA_JOIN/_COUNTED_CAMERA_FILTER).peak_timestamp_epochalongside the formatted string (CompletePeriodMetricsalready 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
_last_group_readings.current < lastre-seeds, emits nothing, marks the cycle.FLAG_INGESTION_STALLEDon ≥ 5 consecutive failed or empty polls, linked to the gap ledger.#31
now.raw_payload.FLAG_INGESTION_GAPexists and is distinct fromFLAG_BURST_COUNTER_FLUSH.klearning and not penalising the Trust Index as sensor faults.CONTEXT.mddocuments ingestion gap vs burst counter flush.FLAG_INGESTION_GAP, without dropping the Trust Index as a sensor fault.#30
peak_timestamp_epochreturned and populated.All
node --test).Sequencing across the
[data-veracity]draftsDeclared together on 2026-09-23; each one is triaged, then reviewed, then implemented, in this order:
FLAG_INGESTION_GAP(in #68) and on #27's seam forcycle_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:
passenger_flow_readings). At a true cold start, differences up toCOLD_START_RECONCILE_LIMIT(100) are reconciled; larger ones are logged as a reset and not published.max(DRIFT_CAP_FLOOR, 3 × trailing per-minute rate)with a floor of 30, so a quiet night can't reject real traffic.ingestion_anomaliestable (reset/stall/gap), readable viaGET /api/analytics/ingestion/anomaliesfor PR #20's deck.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)Implementation is pushed at
7cfdcbc. The repository commit hook passed ruff and pytest; the frontendnode --testsuite 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.Standards
Documented Standard Violations (Hard)
app/controllers/analytics_controller.py:46: Missingheaders={"X-Error-Code": "VALIDATION_ERROR"}when raisingHTTPException(status_code=422, detail="end_epoch precedes start_epoch"). Citesdocs/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: EndpointGET /ingestion/anomaliesand repository method return untypedlist[dict[str, Any]]instead of a Pydantic v2 schema inapp/schemas/occupancy_models.py. Citesdocs/standards/code-standards.md§2.3 (Pydantic Schemas: avoid passing untyped raw dictionaries between layers).Baseline Smells (Judgement Calls)
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.app/db/occupancy_repository.py:104):record_ingestion_anomaly_asynctypeskindas loosestrinstead ofLiteral["reset", "stall", "gap"], discarding SQLiteCHECKconstraint typing.app/services/occupancy_service.py):sync_passenger_flow_from_artemis_asynchas grown to ~340 lines, mixing Artemis HTTP polling, counter reset detection, adaptive drift bounding, gap reconstruction, and persistence.app/services/occupancy_service.py:1976):get_ingestion_anomalies_asyncsimply forwards calls torepo.get_ingestion_anomalies_asyncwithout any domain logic.Spec
Missing or Partial Requirements
get_timespan_aggregates_async()(raw_payloadis ignored). Instead,FLAG_INGESTION_GAPmarks the whole cycle untrusted (AUTO_EXCLUDED), discarding the full day rather than isolating reconstructed spans.test_peak_uses_net_step_for_tied_timestampsonly tests IN inserted before OUT; it never tests OUT inserted before IN or asserts invariance between opposing orders._passenger_flow_poll_failed_asynconly writes on== 5. For outages extending past 5 polls,end_epochfreezes at the 5th poll.Scope Creep (Unasked Behaviour)
"reset"when delta > 100 (occupancy_service.py:2140-2152) was not in the spec.cap = max(30, round(3 * rate))introduces an unrequested minimum floor.GET /ingestion/anomalies(#31): Adding an administrative REST endpoint was unasked; presentation was reserved for PR #20.Incorrect Implementations
cap, instead of bounding the applied delta (min(drift, cap)), all drift is dropped and a spurious"reset"anomaly is logged even thoughcurrent >= last._last_group_readingsupdates tocurrent_artemis_*, permanently discarding the un-emitted counts.FLAG_BURST_COUNTER_FLUSHduring gaps (#31): Spec states: "FLAG_INGESTION_GAP distinct from FLAG_BURST_COUNTER_FLUSH." Negative net flow (total_in - total_out < -1000) still triggersFLAG_BURST_COUNTER_FLUSHeven 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).
Addressed the review in
ab7b2dband6f48205.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.
Review Follow-up & Verification of Previous Findings
Previous Review Status
app/controllers/analytics_controller.py: Addedheaders={"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 untypedlist[dict[str, Any]]with Pydantic v2 schemaIngestionAnomaly(docs/standards/code-standards.md §2.3).app/db/occupancy_repository.py:record_ingestion_anomaly_asyncparameterkindis now typed asLiteral["reset", "stall", "gap"].klearning:get_calibration_flow_asyncignoresreconstructedevents viajson_extract;FLAG_INGESTION_GAPno longer clearsis_trusted.test_peak_uses_net_step_for_tied_timestamps).\ge 5polls:extend_ingestion_stall_asyncproperly updatesend_epoch.min(drift, cap).continueretains previous reading.and not gapsinFLAG_BURST_COUNTER_FLUSH.GET /ingestion/anomaliesretained as accepted scope.Items Missed by Previous Review
_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.app/db/occupancy_repository.py:76get_passenger_flow_readings_asyncreturns untypeddict[str, dict[str, float | int]].Standards
Documented Standard Violations (Hard)
app/db/occupancy_repository.py:76— Untyped Dictionary Return:docs/standards/code-standards.md§2.3 (Pydantic Schemas / Typed Payloads).get_passenger_flow_readings_asyncreturns untypeddict[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)
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.app/services/occupancy_service.py:2097-2305):sync_passenger_flow_from_artemis_asyncspans ~210 lines handling HTTP polling, rollover detection, adaptive drift calculation, bucket gap interpolation, and DB writes.app/services/occupancy_service.py:1986):get_ingestion_anomalies_asyncinOccupancyManagerdelegates onward toself.repo.get_ingestion_anomalies_asyncwith zero domain logic.Spec
Missing or Partial Requirements
< 3600\text{ s}within an hour remain single-instant spikes at(start + end) / 2rather than being uniformly distributed across the outage interval.Scope Creep (Unasked Behaviour)
occupancy_service.py:2200-2206).cap = max(30, round(3 * rate))(occupancy_service.py:2184).GET /ingestion/anomalies(analytics_controller.py:26).Incorrect Implementations
_spread_gap_deltacalculates hour buckets using integer hour boundaries. For an outage like 20 minutes (the primary problem case in Issue #31),duration < 3600produces 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 round 2 addressed —
b2b1d2eSpec
_spread_gap_deltacut 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.GET /api/analytics/ingestion/anomalies.Standards
get_passenger_flow_readings_async/save_passenger_flow_readings_asyncuse aPassengerFlowReadingmodel.GAP_THRESHOLD_SECONDS,DRIFT_RECONCILE_INTERVAL_SECONDS,DRIFT_CAP_FLOOR,DRIFT_CAP_RATE_MULTIPLIER,COLD_START_RECONCILE_LIMIT,STALL_POLL_THRESHOLD).sync_passenger_flow_from_artemis_async(splitting it now would ripple through all four stacked branches), andOccupancyManager.get_ingestion_anomalies_asyncas 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 --test67/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.
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
PassengerFlowReadingmodel, no 7-tuples, named constants). No hard violations. The long sync method and the thin anomalies method stay open by agreement.{"FLAG_INGESTION_GAP", "FLAG_INGESTION_STALLED"}twice without the enum (occupancy_service.py:1033, 1042); the stall sentinelgroup_code="*"written in both the service and the repository;Literal["reset", "stall", "gap"]in two places plus a SQLCHECK; two identicalPassengerFlowReading(...)constructions; a log message saying "five" while the count comes fromSTALL_POLL_THRESHOLD.now - 900// 15.0in the trailing rate (repo:205-208); a new"04:00"fallback; a time-format string in the repository."No passenger flow readings returned".Spec
Suite 250 passed. A 20,000-case property test confirms
_spread_gap_deltasums exactly, stays inside the gap, and never crosses an hour.expected = total_local + deltacounts the pre-rollover slices, and the next drift reconcile republishes them atnowwithout thereconstructedtag. 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 feedk.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.test_gap_spread_respects_hour_edges_and_sums_exactlyis vacuous (stamp < edge < stamp); a true cold start is labelledreset; the trailing rate counts reconstructed events, which inflates the cap after a gap.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.
Review round 3 addressed —
2e1641bresetanomaly and re-seeded, never published (as #26 specifies); drift within the cap is published once. Both are tested.Verification: pytest 254 passed / 1 skipped, CI green on
2e1641b.