"""Add the data-research ingestion control plane.""" from alembic import op revision = "20260722_100" down_revision = "20260720_100" branch_labels = None depends_on = None def upgrade() -> None: op.execute( """ CREATE TABLE IF NOT EXISTS public.ingestion_sources ( uid UUID PRIMARY KEY, source_type VARCHAR(20) NOT NULL CHECK (source_type IN ('database','file','ddl')), name VARCHAR(300) NOT NULL, config JSONB NOT NULL DEFAULT '{}'::jsonb, permission_scope JSONB NOT NULL DEFAULT '{}'::jsonb, status VARCHAR(20) NOT NULL DEFAULT 'active' CHECK (status IN ('active','disabled')), created_by VARCHAR(100), created_at TIMESTAMPTZ NOT NULL DEFAULT CURRENT_TIMESTAMP, updated_at TIMESTAMPTZ NOT NULL DEFAULT CURRENT_TIMESTAMP ); CREATE TABLE IF NOT EXISTS public.source_artifacts ( uid UUID PRIMARY KEY, source_uid UUID NOT NULL REFERENCES public.ingestion_sources(uid), filename VARCHAR(500) NOT NULL, media_type VARCHAR(200) NOT NULL, size_bytes BIGINT NOT NULL CHECK (size_bytes >= 0), content_hash CHAR(64) NOT NULL, storage_ref TEXT NOT NULL, parser_version VARCHAR(100) NOT NULL, created_at TIMESTAMPTZ NOT NULL DEFAULT CURRENT_TIMESTAMP, UNIQUE (source_uid, content_hash, parser_version) ); CREATE TABLE IF NOT EXISTS public.ingestion_jobs ( uid UUID PRIMARY KEY, source_uid UUID NOT NULL REFERENCES public.ingestion_sources(uid), artifact_uid UUID REFERENCES public.source_artifacts(uid), job_type VARCHAR(30) NOT NULL, status VARCHAR(30) NOT NULL DEFAULT 'created' CHECK (status IN ( 'created','queued','extracting','normalizing','matching', 'awaiting_review','published','partial','failed','cancelled' )), idempotency_key CHAR(64) NOT NULL UNIQUE, parser_version VARCHAR(100) NOT NULL, parameters JSONB NOT NULL DEFAULT '{}'::jsonb, statistics JSONB NOT NULL DEFAULT '{}'::jsonb, last_error VARCHAR(1000), actor_uid VARCHAR(100), force_rerun BOOLEAN NOT NULL DEFAULT FALSE, created_at TIMESTAMPTZ NOT NULL DEFAULT CURRENT_TIMESTAMP, updated_at TIMESTAMPTZ NOT NULL DEFAULT CURRENT_TIMESTAMP, started_at TIMESTAMPTZ, finished_at TIMESTAMPTZ ); CREATE TABLE IF NOT EXISTS public.evidence_fragments ( uid UUID PRIMARY KEY, job_uid UUID NOT NULL REFERENCES public.ingestion_jobs(uid) ON DELETE CASCADE, artifact_uid UUID REFERENCES public.source_artifacts(uid), locator JSONB NOT NULL, excerpt TEXT, confidence DOUBLE PRECISION CHECK ( confidence IS NULL OR (confidence >= 0 AND confidence <= 1) ), created_at TIMESTAMPTZ NOT NULL DEFAULT CURRENT_TIMESTAMP ); CREATE TABLE IF NOT EXISTS public.extraction_candidates ( uid UUID PRIMARY KEY, job_uid UUID NOT NULL REFERENCES public.ingestion_jobs(uid) ON DELETE CASCADE, candidate_type VARCHAR(40) NOT NULL, normalized_data JSONB NOT NULL, evidence_uids JSONB NOT NULL DEFAULT '[]'::jsonb, confidence DOUBLE PRECISION NOT NULL CHECK (confidence >= 0 AND confidence <= 1), status VARCHAR(20) NOT NULL DEFAULT 'candidate' CHECK ( status IN ('candidate','accepted','rejected','ignored') ), created_at TIMESTAMPTZ NOT NULL DEFAULT CURRENT_TIMESTAMP, updated_at TIMESTAMPTZ NOT NULL DEFAULT CURRENT_TIMESTAMP ); CREATE INDEX IF NOT EXISTS idx_ingestion_jobs_status_created ON public.ingestion_jobs(status, created_at); CREATE INDEX IF NOT EXISTS idx_ingestion_jobs_source ON public.ingestion_jobs(source_uid, created_at); CREATE INDEX IF NOT EXISTS idx_evidence_fragments_job ON public.evidence_fragments(job_uid); CREATE INDEX IF NOT EXISTS idx_extraction_candidates_job_status ON public.extraction_candidates(job_uid, status); """ ) def downgrade() -> None: # Additive governance data is intentionally preserved on app rollback. pass