20260722_100_data_research_ingestion.py 4.5 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110
  1. """Add the data-research ingestion control plane."""
  2. from alembic import op
  3. revision = "20260722_100"
  4. down_revision = "20260719_90"
  5. branch_labels = None
  6. depends_on = None
  7. def upgrade() -> None:
  8. op.execute(
  9. """
  10. CREATE TABLE IF NOT EXISTS public.ingestion_sources (
  11. uid UUID PRIMARY KEY,
  12. source_type VARCHAR(20) NOT NULL
  13. CHECK (source_type IN ('database','file','ddl')),
  14. name VARCHAR(300) NOT NULL,
  15. config JSONB NOT NULL DEFAULT '{}'::jsonb,
  16. permission_scope JSONB NOT NULL DEFAULT '{}'::jsonb,
  17. status VARCHAR(20) NOT NULL DEFAULT 'active'
  18. CHECK (status IN ('active','disabled')),
  19. created_by VARCHAR(100),
  20. created_at TIMESTAMPTZ NOT NULL DEFAULT CURRENT_TIMESTAMP,
  21. updated_at TIMESTAMPTZ NOT NULL DEFAULT CURRENT_TIMESTAMP
  22. );
  23. CREATE TABLE IF NOT EXISTS public.source_artifacts (
  24. uid UUID PRIMARY KEY,
  25. source_uid UUID NOT NULL
  26. REFERENCES public.ingestion_sources(uid),
  27. filename VARCHAR(500) NOT NULL,
  28. media_type VARCHAR(200) NOT NULL,
  29. size_bytes BIGINT NOT NULL CHECK (size_bytes >= 0),
  30. content_hash CHAR(64) NOT NULL,
  31. storage_ref TEXT NOT NULL,
  32. parser_version VARCHAR(100) NOT NULL,
  33. created_at TIMESTAMPTZ NOT NULL DEFAULT CURRENT_TIMESTAMP,
  34. UNIQUE (source_uid, content_hash, parser_version)
  35. );
  36. CREATE TABLE IF NOT EXISTS public.ingestion_jobs (
  37. uid UUID PRIMARY KEY,
  38. source_uid UUID NOT NULL
  39. REFERENCES public.ingestion_sources(uid),
  40. artifact_uid UUID
  41. REFERENCES public.source_artifacts(uid),
  42. job_type VARCHAR(30) NOT NULL,
  43. status VARCHAR(30) NOT NULL DEFAULT 'created'
  44. CHECK (status IN (
  45. 'created','queued','extracting','normalizing','matching',
  46. 'awaiting_review','published','partial','failed','cancelled'
  47. )),
  48. idempotency_key CHAR(64) NOT NULL UNIQUE,
  49. parser_version VARCHAR(100) NOT NULL,
  50. parameters JSONB NOT NULL DEFAULT '{}'::jsonb,
  51. statistics JSONB NOT NULL DEFAULT '{}'::jsonb,
  52. last_error VARCHAR(1000),
  53. actor_uid VARCHAR(100),
  54. force_rerun BOOLEAN NOT NULL DEFAULT FALSE,
  55. created_at TIMESTAMPTZ NOT NULL DEFAULT CURRENT_TIMESTAMP,
  56. updated_at TIMESTAMPTZ NOT NULL DEFAULT CURRENT_TIMESTAMP,
  57. started_at TIMESTAMPTZ,
  58. finished_at TIMESTAMPTZ
  59. );
  60. CREATE TABLE IF NOT EXISTS public.evidence_fragments (
  61. uid UUID PRIMARY KEY,
  62. job_uid UUID NOT NULL
  63. REFERENCES public.ingestion_jobs(uid) ON DELETE CASCADE,
  64. artifact_uid UUID
  65. REFERENCES public.source_artifacts(uid),
  66. locator JSONB NOT NULL,
  67. excerpt TEXT,
  68. confidence DOUBLE PRECISION
  69. CHECK (confidence IS NULL OR (confidence >= 0 AND confidence <= 1)),
  70. created_at TIMESTAMPTZ NOT NULL DEFAULT CURRENT_TIMESTAMP
  71. );
  72. CREATE TABLE IF NOT EXISTS public.extraction_candidates (
  73. uid UUID PRIMARY KEY,
  74. job_uid UUID NOT NULL
  75. REFERENCES public.ingestion_jobs(uid) ON DELETE CASCADE,
  76. candidate_type VARCHAR(40) NOT NULL,
  77. normalized_data JSONB NOT NULL,
  78. evidence_uids JSONB NOT NULL DEFAULT '[]'::jsonb,
  79. confidence DOUBLE PRECISION NOT NULL
  80. CHECK (confidence >= 0 AND confidence <= 1),
  81. status VARCHAR(20) NOT NULL DEFAULT 'candidate'
  82. CHECK (status IN ('candidate','accepted','rejected','ignored')),
  83. created_at TIMESTAMPTZ NOT NULL DEFAULT CURRENT_TIMESTAMP,
  84. updated_at TIMESTAMPTZ NOT NULL DEFAULT CURRENT_TIMESTAMP
  85. );
  86. CREATE INDEX IF NOT EXISTS idx_ingestion_jobs_status_created
  87. ON public.ingestion_jobs(status, created_at);
  88. CREATE INDEX IF NOT EXISTS idx_ingestion_jobs_source
  89. ON public.ingestion_jobs(source_uid, created_at);
  90. CREATE INDEX IF NOT EXISTS idx_evidence_fragments_job
  91. ON public.evidence_fragments(job_uid);
  92. CREATE INDEX IF NOT EXISTS idx_extraction_candidates_job_status
  93. ON public.extraction_candidates(job_uid, status);
  94. """
  95. )
  96. def downgrade() -> None:
  97. # Additive governance data is intentionally preserved on app rollback.
  98. pass