20260730_380_active_metadata_lineage.py 9.1 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199
  1. """Add active metadata discovery, field lineage and correction audit."""
  2. from alembic import op
  3. revision = "20260730_380"
  4. down_revision = "20260730_370"
  5. branch_labels = None
  6. depends_on = None
  7. def upgrade() -> None:
  8. op.execute(
  9. """
  10. CREATE TABLE public.active_metadata_plans (
  11. uid UUID PRIMARY KEY,
  12. source_uid UUID NOT NULL REFERENCES public.ingestion_sources(uid),
  13. name VARCHAR(300) NOT NULL,
  14. source_kind VARCHAR(20) NOT NULL
  15. CHECK (source_kind IN ('database','file','api')),
  16. schedule_type VARCHAR(20) NOT NULL
  17. CHECK (schedule_type IN ('manual','interval','cron')),
  18. schedule_expression VARCHAR(120),
  19. discovery_mode VARCHAR(20) NOT NULL
  20. CHECK (discovery_mode IN ('snapshot','cursor')),
  21. scope JSONB NOT NULL,
  22. cursor_state JSONB NOT NULL DEFAULT '{}'::jsonb,
  23. owner_uid UUID NOT NULL REFERENCES public.users(id),
  24. enabled BOOLEAN NOT NULL DEFAULT TRUE,
  25. current_version INTEGER NOT NULL DEFAULT 1 CHECK (current_version > 0),
  26. created_by UUID NOT NULL REFERENCES public.users(id),
  27. created_at TIMESTAMPTZ NOT NULL DEFAULT CURRENT_TIMESTAMP,
  28. updated_at TIMESTAMPTZ NOT NULL DEFAULT CURRENT_TIMESTAMP,
  29. CHECK (jsonb_typeof(scope) = 'object'),
  30. CHECK (jsonb_typeof(cursor_state) = 'object')
  31. );
  32. CREATE TABLE public.active_metadata_runs (
  33. uid UUID PRIMARY KEY,
  34. plan_uid UUID NOT NULL REFERENCES public.active_metadata_plans(uid),
  35. batch_key VARCHAR(160) NOT NULL,
  36. status VARCHAR(20) NOT NULL CHECK (status IN ('completed','failed')),
  37. attempt_count INTEGER NOT NULL DEFAULT 1 CHECK (attempt_count > 0),
  38. cursor_before JSONB NOT NULL,
  39. cursor_after JSONB NOT NULL,
  40. snapshot_hash CHAR(64),
  41. statistics JSONB NOT NULL,
  42. failure_code VARCHAR(80),
  43. failure_reason VARCHAR(500),
  44. actor_uid UUID NOT NULL REFERENCES public.users(id),
  45. started_at TIMESTAMPTZ NOT NULL,
  46. finished_at TIMESTAMPTZ NOT NULL,
  47. created_at TIMESTAMPTZ NOT NULL DEFAULT CURRENT_TIMESTAMP,
  48. UNIQUE (plan_uid, batch_key),
  49. CHECK (jsonb_typeof(cursor_before) = 'object'),
  50. CHECK (jsonb_typeof(cursor_after) = 'object'),
  51. CHECK (jsonb_typeof(statistics) = 'object')
  52. );
  53. CREATE TABLE public.active_metadata_assets (
  54. uid UUID PRIMARY KEY,
  55. source_uid UUID NOT NULL REFERENCES public.ingestion_sources(uid),
  56. asset_key VARCHAR(500) NOT NULL,
  57. namespace VARCHAR(200) NOT NULL,
  58. name VARCHAR(200) NOT NULL,
  59. asset_type VARCHAR(40) NOT NULL,
  60. lifecycle_status VARCHAR(30) NOT NULL
  61. CHECK (lifecycle_status IN ('active','deletion_candidate','retired')),
  62. current_version INTEGER NOT NULL CHECK (current_version > 0),
  63. content_hash CHAR(64) NOT NULL,
  64. snapshot JSONB NOT NULL,
  65. health JSONB NOT NULL DEFAULT '{}'::jsonb,
  66. last_run_uid UUID NOT NULL REFERENCES public.active_metadata_runs(uid),
  67. created_at TIMESTAMPTZ NOT NULL DEFAULT CURRENT_TIMESTAMP,
  68. updated_at TIMESTAMPTZ NOT NULL DEFAULT CURRENT_TIMESTAMP,
  69. UNIQUE (source_uid, asset_key),
  70. CHECK (jsonb_typeof(snapshot) = 'object'),
  71. CHECK (jsonb_typeof(health) = 'object')
  72. );
  73. CREATE TABLE public.active_metadata_asset_versions (
  74. uid UUID PRIMARY KEY,
  75. asset_uid UUID NOT NULL REFERENCES public.active_metadata_assets(uid),
  76. version INTEGER NOT NULL CHECK (version > 0),
  77. run_uid UUID NOT NULL REFERENCES public.active_metadata_runs(uid),
  78. content_hash CHAR(64) NOT NULL,
  79. snapshot JSONB NOT NULL,
  80. actor_uid UUID NOT NULL REFERENCES public.users(id),
  81. created_at TIMESTAMPTZ NOT NULL DEFAULT CURRENT_TIMESTAMP,
  82. UNIQUE (asset_uid, version),
  83. CHECK (jsonb_typeof(snapshot) = 'object')
  84. );
  85. CREATE TABLE public.active_metadata_changes (
  86. uid UUID PRIMARY KEY,
  87. run_uid UUID NOT NULL REFERENCES public.active_metadata_runs(uid),
  88. asset_uid UUID NOT NULL REFERENCES public.active_metadata_assets(uid),
  89. asset_key VARCHAR(500) NOT NULL,
  90. field_name VARCHAR(200),
  91. change_type VARCHAR(40) NOT NULL CHECK (
  92. change_type IN (
  93. 'asset_added','field_added','field_changed',
  94. 'field_deletion_candidate','deletion_candidate'
  95. )
  96. ),
  97. before_state JSONB,
  98. after_state JSONB,
  99. status VARCHAR(20) NOT NULL DEFAULT 'pending'
  100. CHECK (status IN ('pending','accepted','rejected')),
  101. created_at TIMESTAMPTZ NOT NULL DEFAULT CURRENT_TIMESTAMP
  102. );
  103. CREATE TABLE public.active_metadata_lineage (
  104. uid UUID PRIMARY KEY,
  105. run_uid UUID NOT NULL REFERENCES public.active_metadata_runs(uid),
  106. parse_status VARCHAR(20) NOT NULL
  107. CHECK (parse_status IN ('resolved','failed')),
  108. source_asset VARCHAR(500),
  109. source_field VARCHAR(200),
  110. target_asset VARCHAR(500),
  111. target_field VARCHAR(200),
  112. relation_type VARCHAR(40) NOT NULL,
  113. evidence JSONB NOT NULL,
  114. failure_reason VARCHAR(500),
  115. created_at TIMESTAMPTZ NOT NULL DEFAULT CURRENT_TIMESTAMP,
  116. CHECK (jsonb_typeof(evidence) = 'object')
  117. );
  118. CREATE TABLE public.active_metadata_health_signals (
  119. uid UUID PRIMARY KEY,
  120. asset_uid UUID NOT NULL REFERENCES public.active_metadata_assets(uid),
  121. run_uid UUID NOT NULL REFERENCES public.active_metadata_runs(uid),
  122. signal_type VARCHAR(30) NOT NULL CHECK (
  123. signal_type IN ('quality','freshness','task_failure','usage')
  124. ),
  125. value JSONB,
  126. status VARCHAR(20) NOT NULL
  127. CHECK (status IN ('healthy','warning','critical','unknown')),
  128. evidence JSONB NOT NULL,
  129. observed_at TIMESTAMPTZ NOT NULL,
  130. created_at TIMESTAMPTZ NOT NULL DEFAULT CURRENT_TIMESTAMP,
  131. CHECK (jsonb_typeof(evidence) = 'object')
  132. );
  133. CREATE TABLE public.active_metadata_corrections (
  134. uid UUID PRIMARY KEY,
  135. asset_uid UUID NOT NULL REFERENCES public.active_metadata_assets(uid),
  136. field_name VARCHAR(200),
  137. proposed_value JSONB NOT NULL,
  138. reason VARCHAR(500) NOT NULL,
  139. assignee_uid UUID NOT NULL REFERENCES public.users(id),
  140. status VARCHAR(20) NOT NULL
  141. CHECK (status IN ('pending','resolved','rejected')),
  142. resolution JSONB NOT NULL DEFAULT '{}'::jsonb,
  143. current_version INTEGER NOT NULL DEFAULT 1 CHECK (current_version > 0),
  144. submitted_by UUID NOT NULL REFERENCES public.users(id),
  145. resolved_by UUID REFERENCES public.users(id),
  146. created_at TIMESTAMPTZ NOT NULL DEFAULT CURRENT_TIMESTAMP,
  147. updated_at TIMESTAMPTZ NOT NULL DEFAULT CURRENT_TIMESTAMP,
  148. CHECK (jsonb_typeof(proposed_value) = 'object'),
  149. CHECK (jsonb_typeof(resolution) = 'object')
  150. );
  151. CREATE TABLE public.active_metadata_correction_audits (
  152. uid UUID PRIMARY KEY,
  153. correction_uid UUID NOT NULL
  154. REFERENCES public.active_metadata_corrections(uid),
  155. version INTEGER NOT NULL CHECK (version > 0),
  156. action VARCHAR(40) NOT NULL,
  157. before_state JSONB,
  158. after_state JSONB NOT NULL,
  159. actor_uid UUID NOT NULL REFERENCES public.users(id),
  160. created_at TIMESTAMPTZ NOT NULL DEFAULT CURRENT_TIMESTAMP,
  161. UNIQUE (correction_uid, version),
  162. CHECK (jsonb_typeof(after_state) = 'object')
  163. );
  164. CREATE INDEX idx_active_metadata_plans_enabled
  165. ON public.active_metadata_plans(enabled, source_kind);
  166. CREATE INDEX idx_active_metadata_runs_plan_created
  167. ON public.active_metadata_runs(plan_uid, created_at DESC);
  168. CREATE INDEX idx_active_metadata_assets_source_status
  169. ON public.active_metadata_assets(source_uid, lifecycle_status);
  170. CREATE INDEX idx_active_metadata_changes_run_status
  171. ON public.active_metadata_changes(run_uid, status, change_type);
  172. CREATE INDEX idx_active_metadata_lineage_target
  173. ON public.active_metadata_lineage(target_asset, target_field);
  174. CREATE INDEX idx_active_metadata_health_asset
  175. ON public.active_metadata_health_signals(asset_uid, observed_at DESC);
  176. CREATE INDEX idx_active_metadata_corrections_assignee
  177. ON public.active_metadata_corrections(assignee_uid, status);
  178. """
  179. )
  180. def downgrade() -> None:
  181. raise RuntimeError(
  182. "active metadata history is append-only; "
  183. "downgrade requires an approved archival migration"
  184. )