| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199 |
- """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"
- )
|