from __future__ import annotations from datetime import datetime, timedelta, timezone class Repository: def __init__(self): self.recovered = [] def running_job_count(self): return 2 def list_stale_jobs(self, before): return [ {"uid": "job-1", "status": "extracting", "updated_at": before - timedelta(seconds=1)}, {"uid": "job-2", "status": "cancelled", "updated_at": before - timedelta(seconds=2)}, ] def recover_job(self, uid, status, safe_error): self.recovered.append((uid, status, safe_error)) def aggregate_metrics(self): return { "job_states": {"failed": 2, "published": 8}, "duration_seconds_p95": 12.5, "failure_stage": {"ocr": 1}, "candidate_acceptance_ratio": 0.8, "review_time_seconds_p95": 55, "ontology_coverage_ratio": 0.75, "projection_lag_seconds": 4, "database_password": "must-not-leak", } def dead_letter_count(self): return 1 def test_limits_metrics_and_stale_job_recovery_are_bounded_and_secret_free(): from app.core.data_research.operations import DataResearchOperations repository = Repository() operations = DataResearchOperations(repository, max_concurrent_jobs=2) assert operations.can_start_job() is False recovered = operations.recover_stale_jobs( now=datetime(2026, 7, 22, tzinfo=timezone.utc), stale_after_seconds=300 ) assert recovered == ["job-1"] assert repository.recovered == [("job-1", "queued", "recovered stale extracting job")] health = operations.health() assert health["dead_letter_count"] == 1 assert health["metrics"]["projection_lag_seconds"] == 4 assert "must-not-leak" not in repr(health) assert "password" not in repr(health).lower() def test_invalid_resource_limits_are_rejected(): import pytest from app.core.data_research.operations import DataResearchOperations with pytest.raises(ValueError, match="max_concurrent_jobs"): DataResearchOperations(Repository(), max_concurrent_jobs=0)