"""Add active metadata discovery, field lineage and correction audit.""" from alembic import op revision = "20260730_380" down_revision = "20260730_370" branch_labels = None depends_on = None def upgrade() -> None: op.execute( """ CREATE TABLE public.active_metadata_plans ( uid UUID PRIMARY KEY, source_uid UUID NOT NULL REFERENCES public.ingestion_sources(uid), name VARCHAR(300) NOT NULL, source_kind VARCHAR(20) NOT NULL CHECK (source_kind IN ('database','file','api')), schedule_type VARCHAR(20) NOT NULL CHECK (schedule_type IN ('manual','interval','cron')), schedule_expression VARCHAR(120), discovery_mode VARCHAR(20) NOT NULL CHECK (discovery_mode IN ('snapshot','cursor')), scope JSONB NOT NULL, cursor_state JSONB NOT NULL DEFAULT '{}'::jsonb, owner_uid UUID NOT NULL REFERENCES public.users(id), enabled BOOLEAN NOT NULL DEFAULT TRUE, current_version INTEGER NOT NULL DEFAULT 1 CHECK (current_version > 0), created_by UUID NOT NULL REFERENCES public.users(id), created_at TIMESTAMPTZ NOT NULL DEFAULT CURRENT_TIMESTAMP, updated_at TIMESTAMPTZ NOT NULL DEFAULT CURRENT_TIMESTAMP, CHECK (jsonb_typeof(scope) = 'object'), CHECK (jsonb_typeof(cursor_state) = 'object') ); CREATE TABLE public.active_metadata_runs ( uid UUID PRIMARY KEY, plan_uid UUID NOT NULL REFERENCES public.active_metadata_plans(uid), batch_key VARCHAR(160) NOT NULL, status VARCHAR(20) NOT NULL CHECK (status IN ('completed','failed')), attempt_count INTEGER NOT NULL DEFAULT 1 CHECK (attempt_count > 0), cursor_before JSONB NOT NULL, cursor_after JSONB NOT NULL, snapshot_hash CHAR(64), statistics JSONB NOT NULL, failure_code VARCHAR(80), failure_reason VARCHAR(500), actor_uid UUID NOT NULL REFERENCES public.users(id), started_at TIMESTAMPTZ NOT NULL, finished_at TIMESTAMPTZ NOT NULL, created_at TIMESTAMPTZ NOT NULL DEFAULT CURRENT_TIMESTAMP, UNIQUE (plan_uid, batch_key), CHECK (jsonb_typeof(cursor_before) = 'object'), CHECK (jsonb_typeof(cursor_after) = 'object'), CHECK (jsonb_typeof(statistics) = 'object') ); CREATE TABLE public.active_metadata_assets ( uid UUID PRIMARY KEY, source_uid UUID NOT NULL REFERENCES public.ingestion_sources(uid), asset_key VARCHAR(500) NOT NULL, namespace VARCHAR(200) NOT NULL, name VARCHAR(200) NOT NULL, asset_type VARCHAR(40) NOT NULL, lifecycle_status VARCHAR(30) NOT NULL CHECK (lifecycle_status IN ('active','deletion_candidate','retired')), current_version INTEGER NOT NULL CHECK (current_version > 0), content_hash CHAR(64) NOT NULL, snapshot JSONB NOT NULL, health JSONB NOT NULL DEFAULT '{}'::jsonb, last_run_uid UUID NOT NULL REFERENCES public.active_metadata_runs(uid), created_at TIMESTAMPTZ NOT NULL DEFAULT CURRENT_TIMESTAMP, updated_at TIMESTAMPTZ NOT NULL DEFAULT CURRENT_TIMESTAMP, UNIQUE (source_uid, asset_key), CHECK (jsonb_typeof(snapshot) = 'object'), CHECK (jsonb_typeof(health) = 'object') ); CREATE TABLE public.active_metadata_asset_versions ( uid UUID PRIMARY KEY, asset_uid UUID NOT NULL REFERENCES public.active_metadata_assets(uid), version INTEGER NOT NULL CHECK (version > 0), run_uid UUID NOT NULL REFERENCES public.active_metadata_runs(uid), content_hash CHAR(64) NOT NULL, snapshot JSONB NOT NULL, actor_uid UUID NOT NULL REFERENCES public.users(id), created_at TIMESTAMPTZ NOT NULL DEFAULT CURRENT_TIMESTAMP, UNIQUE (asset_uid, version), CHECK (jsonb_typeof(snapshot) = 'object') ); CREATE TABLE public.active_metadata_changes ( uid UUID PRIMARY KEY, run_uid UUID NOT NULL REFERENCES public.active_metadata_runs(uid), asset_uid UUID NOT NULL REFERENCES public.active_metadata_assets(uid), asset_key VARCHAR(500) NOT NULL, field_name VARCHAR(200), change_type VARCHAR(40) NOT NULL CHECK ( change_type IN ( 'asset_added','field_added','field_changed', 'field_deletion_candidate','deletion_candidate' ) ), before_state JSONB, after_state JSONB, status VARCHAR(20) NOT NULL DEFAULT 'pending' CHECK (status IN ('pending','accepted','rejected')), created_at TIMESTAMPTZ NOT NULL DEFAULT CURRENT_TIMESTAMP ); CREATE TABLE public.active_metadata_lineage ( uid UUID PRIMARY KEY, run_uid UUID NOT NULL REFERENCES public.active_metadata_runs(uid), parse_status VARCHAR(20) NOT NULL CHECK (parse_status IN ('resolved','failed')), source_asset VARCHAR(500), source_field VARCHAR(200), target_asset VARCHAR(500), target_field VARCHAR(200), relation_type VARCHAR(40) NOT NULL, evidence JSONB NOT NULL, failure_reason VARCHAR(500), created_at TIMESTAMPTZ NOT NULL DEFAULT CURRENT_TIMESTAMP, CHECK (jsonb_typeof(evidence) = 'object') ); CREATE TABLE public.active_metadata_health_signals ( uid UUID PRIMARY KEY, asset_uid UUID NOT NULL REFERENCES public.active_metadata_assets(uid), run_uid UUID NOT NULL REFERENCES public.active_metadata_runs(uid), signal_type VARCHAR(30) NOT NULL CHECK ( signal_type IN ('quality','freshness','task_failure','usage') ), value JSONB, status VARCHAR(20) NOT NULL CHECK (status IN ('healthy','warning','critical','unknown')), evidence JSONB NOT NULL, observed_at TIMESTAMPTZ NOT NULL, created_at TIMESTAMPTZ NOT NULL DEFAULT CURRENT_TIMESTAMP, CHECK (jsonb_typeof(evidence) = 'object') ); CREATE TABLE public.active_metadata_corrections ( uid UUID PRIMARY KEY, asset_uid UUID NOT NULL REFERENCES public.active_metadata_assets(uid), field_name VARCHAR(200), proposed_value JSONB NOT NULL, reason VARCHAR(500) NOT NULL, assignee_uid UUID NOT NULL REFERENCES public.users(id), status VARCHAR(20) NOT NULL CHECK (status IN ('pending','resolved','rejected')), resolution JSONB NOT NULL DEFAULT '{}'::jsonb, current_version INTEGER NOT NULL DEFAULT 1 CHECK (current_version > 0), submitted_by UUID NOT NULL REFERENCES public.users(id), resolved_by UUID REFERENCES public.users(id), created_at TIMESTAMPTZ NOT NULL DEFAULT CURRENT_TIMESTAMP, updated_at TIMESTAMPTZ NOT NULL DEFAULT CURRENT_TIMESTAMP, CHECK (jsonb_typeof(proposed_value) = 'object'), CHECK (jsonb_typeof(resolution) = 'object') ); CREATE TABLE public.active_metadata_correction_audits ( uid UUID PRIMARY KEY, correction_uid UUID NOT NULL REFERENCES public.active_metadata_corrections(uid), version INTEGER NOT NULL CHECK (version > 0), action VARCHAR(40) NOT NULL, before_state JSONB, after_state JSONB NOT NULL, actor_uid UUID NOT NULL REFERENCES public.users(id), created_at TIMESTAMPTZ NOT NULL DEFAULT CURRENT_TIMESTAMP, UNIQUE (correction_uid, version), CHECK (jsonb_typeof(after_state) = 'object') ); CREATE INDEX idx_active_metadata_plans_enabled ON public.active_metadata_plans(enabled, source_kind); CREATE INDEX idx_active_metadata_runs_plan_created ON public.active_metadata_runs(plan_uid, created_at DESC); CREATE INDEX idx_active_metadata_assets_source_status ON public.active_metadata_assets(source_uid, lifecycle_status); CREATE INDEX idx_active_metadata_changes_run_status ON public.active_metadata_changes(run_uid, status, change_type); CREATE INDEX idx_active_metadata_lineage_target ON public.active_metadata_lineage(target_asset, target_field); CREATE INDEX idx_active_metadata_health_asset ON public.active_metadata_health_signals(asset_uid, observed_at DESC); CREATE INDEX idx_active_metadata_corrections_assignee ON public.active_metadata_corrections(assignee_uid, status); """ ) def downgrade() -> None: raise RuntimeError( "active metadata history is append-only; " "downgrade requires an approved archival migration" )