from __future__ import annotations from dataclasses import replace import pytest class MemoryJobRepository: def __init__(self): self.records = {} def get(self, uid): return self.records.get(uid) def get_by_idempotency_key(self, key): return next( (item for item in self.records.values() if item.idempotency_key == key), None, ) def add(self, record): self.records[record.uid] = record return record def save(self, record): self.records[record.uid] = record return record def payload(**overrides): value = { "source_uid": "00000000-0000-0000-0000-000000000001", "artifact_uid": "00000000-0000-0000-0000-000000000002", "job_type": "file_extract", "parser_version": "csv-v1", "parameters": {"delimiter": ",", "schema": "public"}, } value.update(overrides) return value def test_same_canonical_input_reuses_existing_job(): from app.core.data_research.ingestion import IngestionService repository = MemoryJobRepository() service = IngestionService(repository) first, first_created = service.create_job(payload(), actor_uid="editor-1") second, second_created = service.create_job( payload(parameters={"schema": "public", "delimiter": ","}), actor_uid="editor-2", ) assert first_created is True assert second_created is False assert second.uid == first.uid assert second.idempotency_key == first.idempotency_key def test_force_rerun_creates_a_distinct_job_without_changing_canonical_key_rules(): from app.core.data_research.ingestion import IngestionService repository = MemoryJobRepository() service = IngestionService(repository, nonce_factory=lambda: "rerun-1") first, _ = service.create_job(payload(), actor_uid="editor-1") rerun, created = service.create_job( payload(force_rerun=True), actor_uid="editor-1", ) assert created is True assert rerun.uid != first.uid assert rerun.idempotency_key != first.idempotency_key assert rerun.force_rerun is True def test_job_transitions_follow_the_declared_state_machine(): from app.core.data_research.errors import InvalidJobTransition from app.core.data_research.ingestion import IngestionService repository = MemoryJobRepository() service = IngestionService(repository) job, _ = service.create_job(payload(), actor_uid="editor-1") for status in ( "queued", "extracting", "normalizing", "matching", "awaiting_review", "published", ): job = service.transition(job.uid, status, statistics={"stage": status}) assert job.status == status with pytest.raises(InvalidJobTransition, match="published -> queued"): service.transition(job.uid, "queued") def test_failed_job_can_retry_and_running_job_can_cancel(): from app.core.data_research.ingestion import IngestionService repository = MemoryJobRepository() service = IngestionService(repository) job, _ = service.create_job(payload(), actor_uid="editor-1") service.transition(job.uid, "queued") extracting = service.transition(job.uid, "extracting") assert extracting.attempt_count == 1 failed = service.transition( job.uid, "failed", error="password=clear-secret token=abc123 connection refused", ) assert "clear-secret" not in failed.last_error assert "abc123" not in failed.last_error assert "[redacted]" in failed.last_error assert failed.failure_stage == "extracting" assert failed.attempt_count == 1 retried = service.retry(job.uid) assert retried.status == "queued" assert retried.last_error is None assert retried.failure_stage is None extracting_again = service.transition(job.uid, "extracting") assert extracting_again.attempt_count == 2 cancelled = service.cancel(job.uid) assert cancelled.status == "cancelled" def test_missing_or_invalid_payload_fails_before_repository_write(): from app.core.data_research.errors import IngestionPayloadInvalid from app.core.data_research.ingestion import IngestionService repository = MemoryJobRepository() service = IngestionService(repository) with pytest.raises(IngestionPayloadInvalid, match="source_uid"): service.create_job(payload(source_uid=""), actor_uid="editor-1") with pytest.raises(IngestionPayloadInvalid, match="parameters"): service.create_job(payload(parameters=[]), actor_uid="editor-1") assert repository.records == {}