20260802_472_enterprise_connectors.py 7.5 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114
  1. """Enterprise connector SDK, machine identity, execution and graph ledger."""
  2. from alembic import op
  3. revision = "20260802_472"
  4. down_revision = "20260802_471"
  5. branch_labels = None
  6. depends_on = None
  7. def upgrade() -> None:
  8. op.execute("""
  9. CREATE TABLE public.connector_manifests (
  10. uid UUID PRIMARY KEY, connector_id VARCHAR(64) NOT NULL, connector_version VARCHAR(40) NOT NULL,
  11. sdk_version VARCHAR(20) NOT NULL, display_name VARCHAR(200) NOT NULL,
  12. capabilities JSONB NOT NULL, config_schema JSONB NOT NULL, status VARCHAR(20) NOT NULL,
  13. created_by UUID NOT NULL, created_at TIMESTAMPTZ NOT NULL DEFAULT CURRENT_TIMESTAMP,
  14. updated_at TIMESTAMPTZ NOT NULL DEFAULT CURRENT_TIMESTAMP,
  15. CONSTRAINT uq_connector_manifest_version UNIQUE (connector_id, connector_version),
  16. CONSTRAINT ck_connector_manifest_status CHECK (status IN ('active','retired')),
  17. CONSTRAINT ck_connector_manifest_capabilities_array CHECK (jsonb_typeof(capabilities)='array'),
  18. CONSTRAINT ck_connector_manifest_schema_object CHECK (jsonb_typeof(config_schema)='object')
  19. );
  20. CREATE TABLE public.connector_principals (
  21. uid UUID PRIMARY KEY, connector_id VARCHAR(64) NOT NULL, source_uid UUID NOT NULL,
  22. business_domain_uid UUID NOT NULL, environment VARCHAR(20) NOT NULL,
  23. allowed_operations TEXT[] NOT NULL, allowed_scopes JSONB NOT NULL,
  24. status VARCHAR(20) NOT NULL, created_by UUID NOT NULL,
  25. created_at TIMESTAMPTZ NOT NULL DEFAULT CURRENT_TIMESTAMP, revoked_at TIMESTAMPTZ,
  26. CONSTRAINT uq_connector_principal_binding UNIQUE (connector_id, source_uid, business_domain_uid, environment),
  27. CONSTRAINT ck_connector_principal_environment CHECK (environment IN ('development','staging','production')),
  28. CONSTRAINT ck_connector_principal_status CHECK (status IN ('active','revoked')),
  29. CONSTRAINT ck_connector_principal_scope CHECK (jsonb_typeof(allowed_scopes)='object')
  30. );
  31. CREATE TABLE public.connector_machine_credentials (
  32. uid UUID PRIMARY KEY, principal_uid UUID NOT NULL REFERENCES public.connector_principals(uid),
  33. token_hash CHAR(64) NOT NULL UNIQUE, status VARCHAR(20) NOT NULL,
  34. issued_by UUID NOT NULL, issued_at TIMESTAMPTZ NOT NULL DEFAULT CURRENT_TIMESTAMP,
  35. expires_at TIMESTAMPTZ NOT NULL, first_used_at TIMESTAMPTZ, use_count INTEGER NOT NULL DEFAULT 0,
  36. revoked_at TIMESTAMPTZ, rotated_from_uid UUID REFERENCES public.connector_machine_credentials(uid),
  37. CONSTRAINT ck_connector_credential_status CHECK (status IN ('active','revoked','rotated','expired','replayed')),
  38. CONSTRAINT ck_connector_credential_ttl CHECK (expires_at <= issued_at + INTERVAL '15 minutes'),
  39. CONSTRAINT ck_connector_credential_use_count CHECK (use_count >= 0)
  40. );
  41. CREATE INDEX ix_connector_credential_principal_status ON public.connector_machine_credentials(principal_uid,status,expires_at);
  42. CREATE TABLE public.connector_runs (
  43. uid UUID PRIMARY KEY, idempotency_key CHAR(64) NOT NULL UNIQUE,
  44. connector_id VARCHAR(64) NOT NULL, connector_version VARCHAR(40) NOT NULL,
  45. source_uid UUID NOT NULL, operation VARCHAR(30) NOT NULL, status VARCHAR(20) NOT NULL,
  46. attempt_count INTEGER NOT NULL DEFAULT 0, checkpoint JSONB NOT NULL DEFAULT '{}'::jsonb,
  47. cursor JSONB NOT NULL DEFAULT '{}'::jsonb, error_category VARCHAR(30), error_code VARCHAR(80),
  48. resumed_from_run_uid UUID REFERENCES public.connector_runs(uid), dry_run BOOLEAN NOT NULL DEFAULT FALSE,
  49. actor_uid UUID NOT NULL, created_at TIMESTAMPTZ NOT NULL DEFAULT CURRENT_TIMESTAMP,
  50. updated_at TIMESTAMPTZ NOT NULL DEFAULT CURRENT_TIMESTAMP,
  51. CONSTRAINT fk_connector_run_manifest FOREIGN KEY (connector_id,connector_version)
  52. REFERENCES public.connector_manifests(connector_id,connector_version),
  53. CONSTRAINT ck_connector_run_operation CHECK (operation IN ('discover','snapshot','incremental','lineage','profile','cancel','resume','evidence')),
  54. CONSTRAINT ck_connector_run_status CHECK (status IN ('running','succeeded','dry_run','failed','cancelled','resumable')),
  55. CONSTRAINT ck_connector_run_attempt CHECK (attempt_count BETWEEN 0 AND 5)
  56. );
  57. CREATE INDEX ix_connector_runs_source_created ON public.connector_runs(source_uid,created_at DESC);
  58. CREATE TABLE public.connector_run_attempts (
  59. uid UUID PRIMARY KEY, run_uid UUID NOT NULL REFERENCES public.connector_runs(uid) ON DELETE CASCADE,
  60. attempt_number INTEGER NOT NULL, status VARCHAR(20) NOT NULL, error_category VARCHAR(30),
  61. started_at TIMESTAMPTZ NOT NULL DEFAULT CURRENT_TIMESTAMP, finished_at TIMESTAMPTZ,
  62. CONSTRAINT uq_connector_run_attempt UNIQUE(run_uid,attempt_number)
  63. );
  64. CREATE TABLE public.connector_checkpoints (
  65. uid UUID PRIMARY KEY, run_uid UUID NOT NULL REFERENCES public.connector_runs(uid) ON DELETE CASCADE,
  66. sequence_number INTEGER NOT NULL, cursor JSONB NOT NULL, checkpoint JSONB NOT NULL,
  67. content_hash CHAR(64) NOT NULL, created_at TIMESTAMPTZ NOT NULL DEFAULT CURRENT_TIMESTAMP,
  68. CONSTRAINT uq_connector_checkpoint_sequence UNIQUE(run_uid,sequence_number)
  69. );
  70. CREATE TABLE public.connector_evidence (
  71. uid UUID PRIMARY KEY, run_uid UUID NOT NULL REFERENCES public.connector_runs(uid) ON DELETE CASCADE,
  72. evidence_type VARCHAR(40) NOT NULL, payload JSONB NOT NULL, content_hash CHAR(64) NOT NULL,
  73. byte_size INTEGER NOT NULL, redacted BOOLEAN NOT NULL,
  74. created_at TIMESTAMPTZ NOT NULL DEFAULT CURRENT_TIMESTAMP,
  75. CONSTRAINT ck_connector_evidence_size CHECK (byte_size BETWEEN 0 AND 32768)
  76. );
  77. CREATE TABLE public.connector_graph_edges (
  78. uid UUID PRIMARY KEY, source_uid UUID NOT NULL, from_type VARCHAR(30) NOT NULL,
  79. from_key VARCHAR(500) NOT NULL, relation_type VARCHAR(40) NOT NULL,
  80. to_type VARCHAR(30) NOT NULL, to_key VARCHAR(500) NOT NULL,
  81. business_domain_uid UUID, run_uid UUID REFERENCES public.connector_runs(uid),
  82. run_status VARCHAR(20), evidence JSONB NOT NULL DEFAULT '{}'::jsonb,
  83. created_at TIMESTAMPTZ NOT NULL DEFAULT CURRENT_TIMESTAMP,
  84. CONSTRAINT uq_connector_graph_edge UNIQUE(source_uid,from_type,from_key,relation_type,to_type,to_key),
  85. CONSTRAINT ck_connector_graph_node_types CHECK (from_type IN ('source','asset','process','business_domain','run') AND to_type IN ('source','asset','process','business_domain','run'))
  86. );
  87. CREATE INDEX ix_connector_graph_source ON public.connector_graph_edges(source_uid,relation_type);
  88. CREATE INDEX ix_connector_graph_to ON public.connector_graph_edges(to_type,to_key);
  89. CREATE TABLE public.connector_audit_events (
  90. uid UUID PRIMARY KEY, principal_uid UUID REFERENCES public.connector_principals(uid),
  91. credential_uid UUID REFERENCES public.connector_machine_credentials(uid), event_type VARCHAR(80) NOT NULL,
  92. actor_uid UUID, success BOOLEAN NOT NULL, safe_detail VARCHAR(500) NOT NULL,
  93. created_at TIMESTAMPTZ NOT NULL DEFAULT CURRENT_TIMESTAMP
  94. );
  95. CREATE INDEX ix_connector_audit_principal_created ON public.connector_audit_events(principal_uid,created_at DESC);
  96. """)
  97. def downgrade() -> None:
  98. op.execute("""
  99. DROP TABLE IF EXISTS public.connector_audit_events;
  100. DROP TABLE IF EXISTS public.connector_graph_edges;
  101. DROP TABLE IF EXISTS public.connector_evidence;
  102. DROP TABLE IF EXISTS public.connector_checkpoints;
  103. DROP TABLE IF EXISTS public.connector_run_attempts;
  104. DROP TABLE IF EXISTS public.connector_runs;
  105. DROP TABLE IF EXISTS public.connector_machine_credentials;
  106. DROP TABLE IF EXISTS public.connector_principals;
  107. DROP TABLE IF EXISTS public.connector_manifests;
  108. """)