test_ingestion_service.py 4.2 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137
  1. from __future__ import annotations
  2. from dataclasses import replace
  3. import pytest
  4. class MemoryJobRepository:
  5. def __init__(self):
  6. self.records = {}
  7. def get(self, uid):
  8. return self.records.get(uid)
  9. def get_by_idempotency_key(self, key):
  10. return next(
  11. (item for item in self.records.values() if item.idempotency_key == key),
  12. None,
  13. )
  14. def add(self, record):
  15. self.records[record.uid] = record
  16. return record
  17. def save(self, record):
  18. self.records[record.uid] = record
  19. return record
  20. def payload(**overrides):
  21. value = {
  22. "source_uid": "00000000-0000-0000-0000-000000000001",
  23. "artifact_uid": "00000000-0000-0000-0000-000000000002",
  24. "job_type": "file_extract",
  25. "parser_version": "csv-v1",
  26. "parameters": {"delimiter": ",", "schema": "public"},
  27. }
  28. value.update(overrides)
  29. return value
  30. def test_same_canonical_input_reuses_existing_job():
  31. from app.core.data_research.ingestion import IngestionService
  32. repository = MemoryJobRepository()
  33. service = IngestionService(repository)
  34. first, first_created = service.create_job(payload(), actor_uid="editor-1")
  35. second, second_created = service.create_job(
  36. payload(parameters={"schema": "public", "delimiter": ","}),
  37. actor_uid="editor-2",
  38. )
  39. assert first_created is True
  40. assert second_created is False
  41. assert second.uid == first.uid
  42. assert second.idempotency_key == first.idempotency_key
  43. def test_force_rerun_creates_a_distinct_job_without_changing_canonical_key_rules():
  44. from app.core.data_research.ingestion import IngestionService
  45. repository = MemoryJobRepository()
  46. service = IngestionService(repository, nonce_factory=lambda: "rerun-1")
  47. first, _ = service.create_job(payload(), actor_uid="editor-1")
  48. rerun, created = service.create_job(
  49. payload(force_rerun=True),
  50. actor_uid="editor-1",
  51. )
  52. assert created is True
  53. assert rerun.uid != first.uid
  54. assert rerun.idempotency_key != first.idempotency_key
  55. assert rerun.force_rerun is True
  56. def test_job_transitions_follow_the_declared_state_machine():
  57. from app.core.data_research.errors import InvalidJobTransition
  58. from app.core.data_research.ingestion import IngestionService
  59. repository = MemoryJobRepository()
  60. service = IngestionService(repository)
  61. job, _ = service.create_job(payload(), actor_uid="editor-1")
  62. for status in (
  63. "queued",
  64. "extracting",
  65. "normalizing",
  66. "matching",
  67. "awaiting_review",
  68. "published",
  69. ):
  70. job = service.transition(job.uid, status, statistics={"stage": status})
  71. assert job.status == status
  72. with pytest.raises(InvalidJobTransition, match="published -> queued"):
  73. service.transition(job.uid, "queued")
  74. def test_failed_job_can_retry_and_running_job_can_cancel():
  75. from app.core.data_research.ingestion import IngestionService
  76. repository = MemoryJobRepository()
  77. service = IngestionService(repository)
  78. job, _ = service.create_job(payload(), actor_uid="editor-1")
  79. service.transition(job.uid, "queued")
  80. failed = service.transition(
  81. job.uid,
  82. "failed",
  83. error="password=clear-secret token=abc123 connection refused",
  84. )
  85. assert "clear-secret" not in failed.last_error
  86. assert "abc123" not in failed.last_error
  87. assert "[redacted]" in failed.last_error
  88. retried = service.retry(job.uid)
  89. assert retried.status == "queued"
  90. assert retried.last_error is None
  91. cancelled = service.cancel(job.uid)
  92. assert cancelled.status == "cancelled"
  93. def test_missing_or_invalid_payload_fails_before_repository_write():
  94. from app.core.data_research.errors import IngestionPayloadInvalid
  95. from app.core.data_research.ingestion import IngestionService
  96. repository = MemoryJobRepository()
  97. service = IngestionService(repository)
  98. with pytest.raises(IngestionPayloadInvalid, match="source_uid"):
  99. service.create_job(payload(source_uid=""), actor_uid="editor-1")
  100. with pytest.raises(IngestionPayloadInvalid, match="parameters"):
  101. service.create_job(payload(parameters=[]), actor_uid="editor-1")
  102. assert repository.records == {}