test_ingestion_models.py 3.4 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103
  1. from __future__ import annotations
  2. from sqlalchemy import CheckConstraint, UniqueConstraint
  3. def test_ingestion_models_expose_stable_control_plane_tables():
  4. from app.models.data_research import (
  5. EvidenceFragment,
  6. ExtractionCandidate,
  7. IngestionJob,
  8. IngestionSource,
  9. SourceArtifact,
  10. )
  11. assert IngestionSource.__tablename__ == "ingestion_sources"
  12. assert IngestionJob.__tablename__ == "ingestion_jobs"
  13. assert SourceArtifact.__tablename__ == "source_artifacts"
  14. assert EvidenceFragment.__tablename__ == "evidence_fragments"
  15. assert ExtractionCandidate.__tablename__ == "extraction_candidates"
  16. assert IngestionJob.__table__.c.source_uid.foreign_keys
  17. assert IngestionJob.__table__.c.artifact_uid.foreign_keys
  18. assert EvidenceFragment.__table__.c.job_uid.foreign_keys
  19. assert ExtractionCandidate.__table__.c.job_uid.foreign_keys
  20. def test_ingestion_job_has_idempotency_and_status_constraints():
  21. from app.models.data_research import IngestionJob
  22. constraints = list(IngestionJob.__table__.constraints)
  23. unique_columns = {
  24. tuple(column.name for column in constraint.columns)
  25. for constraint in constraints
  26. if isinstance(constraint, UniqueConstraint)
  27. }
  28. check_sql = " ".join(
  29. str(constraint.sqltext)
  30. for constraint in constraints
  31. if isinstance(constraint, CheckConstraint)
  32. )
  33. assert ("idempotency_key",) in unique_columns
  34. assert "attempt_count" in IngestionJob.__table__.c
  35. assert "failure_stage" in IngestionJob.__table__.c
  36. assert "attempt_count >= 0" in check_sql
  37. for status in (
  38. "created",
  39. "queued",
  40. "extracting",
  41. "normalizing",
  42. "matching",
  43. "awaiting_review",
  44. "published",
  45. "partial",
  46. "failed",
  47. "cancelled",
  48. ):
  49. assert status in check_sql
  50. def test_catalog_snapshot_has_one_immutable_record_per_job_attempt():
  51. from app.models.data_research import CatalogSnapshot
  52. constraints = list(CatalogSnapshot.__table__.constraints)
  53. unique_columns = {
  54. tuple(column.name for column in constraint.columns)
  55. for constraint in constraints
  56. if isinstance(constraint, UniqueConstraint)
  57. }
  58. assert CatalogSnapshot.__tablename__ == "catalog_snapshots"
  59. assert CatalogSnapshot.__table__.c.job_uid.foreign_keys
  60. assert ("job_uid", "attempt") in unique_columns
  61. assert "content_hash" in CatalogSnapshot.__table__.c
  62. def test_model_serialization_redacts_control_plane_secrets():
  63. from app.models.data_research import IngestionSource, SourceArtifact
  64. source = IngestionSource(
  65. uid="00000000-0000-0000-0000-000000000001",
  66. source_type="database",
  67. name="orders",
  68. config={"secret_ref": "vault://orders", "schema": "public"},
  69. permission_scope={"roles": ["editor"]},
  70. status="active",
  71. )
  72. artifact = SourceArtifact(
  73. uid="00000000-0000-0000-0000-000000000002",
  74. source_uid=source.uid,
  75. filename="orders.csv",
  76. media_type="text/csv",
  77. size_bytes=12,
  78. content_hash="a" * 64,
  79. storage_ref="minio://private/orders.csv?token=secret",
  80. parser_version="csv-v1",
  81. )
  82. assert "config" not in source.to_dict()
  83. assert "vault://orders" not in str(source.to_dict())
  84. assert "storage_ref" not in artifact.to_dict()
  85. assert "token=secret" not in str(artifact.to_dict())
  86. assert artifact.to_dict()["content_hash"] == "a" * 64