Skip to content

Commit de24b15

Browse files
fix: restore cleanup invariants
1 parent 0323cfb commit de24b15

3 files changed

Lines changed: 19 additions & 17 deletions

File tree

‎backend/control/measurements.py‎

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -27,9 +27,11 @@ class SourceMeasurement(BaseModel):
2727
duplicate_rate: float
2828

2929
freshness_lag_seconds: int | None = None
30+
cursor_advanced: bool
3031

3132
odp_stream_lag: int | None = None
3233
odp_pending: int | None = None
34+
dlq_count: int = 0
3335

3436
#: {ErrorKind.value: count} — PR-Control-3. Reuses the SAME taxonomy the
3537
#: persisted backend.models.source_measurement.SourceMeasurement.error_kinds

‎tests/integration/test_control_api.py‎

Lines changed: 6 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -3,7 +3,7 @@
33
rows, POST /control/outcomes/evaluate, and GET /control/advisory-report.
44
"""
55

6-
from datetime import datetime, timedelta
6+
from datetime import UTC, datetime, timedelta
77

88
import pytest
99
from sqlalchemy import select
@@ -25,7 +25,7 @@ async def _seed_measurement_row(session, source_id: str, **overrides):
2525
kwargs = dict(
2626
source_id=source_id,
2727
run_id="row-run-1",
28-
measured_at=datetime(2026, 7, 2, tzinfo=datetime.UTC),
28+
measured_at=datetime(2026, 7, 2, tzinfo=UTC),
2929
accepted=0,
3030
duplicates=0,
3131
rejected=1,
@@ -154,7 +154,7 @@ async def test_evaluate_and_advisory_report_recovery_flow(
154154
.all()
155155
)
156156
assert rows
157-
backdated = datetime.now(datetime.UTC) - timedelta(seconds=120)
157+
backdated = datetime.now(UTC) - timedelta(seconds=120)
158158
for row in rows:
159159
row.created_at = backdated
160160
await db_session.commit()
@@ -166,7 +166,7 @@ async def test_evaluate_and_advisory_report_recovery_flow(
166166
db_session,
167167
source_id,
168168
run_id="row-run-post",
169-
measured_at=datetime.now(datetime.UTC),
169+
measured_at=datetime.now(UTC),
170170
accepted=5,
171171
duplicates=0,
172172
rejected=0,
@@ -299,7 +299,7 @@ async def test_kill_switch_engaged_short_circuits_control_cycle(
299299
executed=False,
300300
measurement_before={},
301301
outcome="recovered" if i < 8 else "persisted",
302-
evaluated_at=datetime.now(datetime.UTC),
302+
evaluated_at=datetime.now(UTC),
303303
)
304304
)
305305
await db_session.commit()
@@ -309,7 +309,7 @@ async def test_kill_switch_engaged_short_circuits_control_cycle(
309309

310310
from backend.control.cycle import run_control_cycle_once
311311

312-
result = await run_control_cycle_once(db_session, now=datetime.now(datetime.UTC))
312+
result = await run_control_cycle_once(db_session, now=datetime.now(UTC))
313313
await db_session.commit()
314314

315315
assert result.executions == []

‎tests/integration/test_sources_api.py‎

Lines changed: 11 additions & 11 deletions
Original file line numberDiff line numberDiff line change
@@ -357,7 +357,7 @@ async def test_control_state_with_run_evidence(client, db_session, sample_source
357357
monkeypatch.delenv("ODP_DATABASE_URL", raising=False)
358358
monkeypatch.delenv("ODP_INGEST_URL", raising=False)
359359

360-
from datetime import datetime
360+
from datetime import UTC, datetime
361361

362362
from backend.models.task import CollectionTask, TaskRun, TaskRunEvent
363363

@@ -371,7 +371,7 @@ async def test_control_state_with_run_evidence(client, db_session, sample_source
371371
run = TaskRun(
372372
task_id=task.id,
373373
status="completed",
374-
finished_at=datetime(2026, 7, 2, tzinfo=datetime.UTC),
374+
finished_at=datetime(2026, 7, 2, tzinfo=UTC),
375375
duration_ms=1000,
376376
records_collected=9,
377377
)
@@ -427,7 +427,7 @@ async def test_control_state_degraded_when_errors(
427427
monkeypatch.delenv("ODP_DATABASE_URL", raising=False)
428428
monkeypatch.delenv("ODP_INGEST_URL", raising=False)
429429

430-
from datetime import datetime
430+
from datetime import UTC, datetime
431431

432432
from backend.models.task import CollectionTask, TaskRun, TaskRunEvent
433433

@@ -441,7 +441,7 @@ async def test_control_state_degraded_when_errors(
441441
run = TaskRun(
442442
task_id=task.id,
443443
status="completed",
444-
finished_at=datetime(2026, 7, 2, tzinfo=datetime.UTC),
444+
finished_at=datetime(2026, 7, 2, tzinfo=UTC),
445445
duration_ms=1000,
446446
records_collected=1,
447447
)
@@ -485,7 +485,7 @@ async def test_control_state_trend_fallback_for_pre_measurement_source(
485485
monkeypatch.delenv("ODP_DATABASE_URL", raising=False)
486486
monkeypatch.delenv("ODP_INGEST_URL", raising=False)
487487

488-
from datetime import datetime
488+
from datetime import UTC, datetime
489489

490490
from backend.models.task import CollectionTask, TaskRun, TaskRunEvent
491491

@@ -501,8 +501,8 @@ async def test_control_state_trend_fallback_for_pre_measurement_source(
501501
run = TaskRun(
502502
task_id=task.id,
503503
status="completed",
504-
created_at=datetime(2026, 7, day, tzinfo=datetime.UTC),
505-
finished_at=datetime(2026, 7, day, tzinfo=datetime.UTC),
504+
created_at=datetime(2026, 7, day, tzinfo=UTC),
505+
finished_at=datetime(2026, 7, day, tzinfo=UTC),
506506
duration_ms=1000,
507507
records_collected=stored,
508508
)
@@ -546,14 +546,14 @@ async def test_control_state_trend_fallback_for_pre_measurement_source(
546546

547547

548548
async def _seed_measurement_row(session, source_id: str, **overrides):
549-
from datetime import datetime
549+
from datetime import UTC, datetime
550550

551551
from backend.models.source_measurement import SourceMeasurement as SourceMeasurementRow
552552

553553
kwargs = dict(
554554
source_id=source_id,
555555
run_id="row-run-1",
556-
measured_at=datetime(2026, 7, 2, tzinfo=datetime.UTC),
556+
measured_at=datetime(2026, 7, 2, tzinfo=UTC),
557557
accepted=5,
558558
duplicates=0,
559559
rejected=0,
@@ -930,14 +930,14 @@ async def test_control_state_classification_flips_with_objective_override(
930930
create_resp = await client.post("/api/v1/sources", json=sample_source_data)
931931
source_id = create_resp.json()["data"]["id"]
932932

933-
from datetime import datetime
933+
from datetime import UTC, datetime
934934

935935
from backend.models.source_measurement import SourceMeasurement as SourceMeasurementRow
936936

937937
row = SourceMeasurementRow(
938938
source_id=source_id,
939939
run_id="row-run-1",
940-
measured_at=datetime(2026, 7, 2, tzinfo=datetime.UTC),
940+
measured_at=datetime(2026, 7, 2, tzinfo=UTC),
941941
accepted=97,
942942
duplicates=0,
943943
rejected=3,

0 commit comments

Comments
 (0)