| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691692693694695696697698699700701702703704705706707708709710711712713714715716717718719720721722723724725726727728729730731732733734735736737738739740741742743744745746747748749750751752753754755756757758759760761762763764765766767768769770771772773774775776777778779780781782783784785786787788789790791792793794795796797798799800801802803804805806807808809810811812813814815816817818819820821822823824825826827828829830831832833834835836837838839840841842843844845846847848849850851852853854855856857858859860861862863864865866867868869870871872873874875876877878879880881882883884885886887888889890891892893894895896897898899900901902903904905906907908909910911912913914915916917918919920921922923924925926927928929930931932933934935936937938939940941942943944945946947948949950951952953954955956957958959960961962963964965966967968969970971972973974975976977978979980981982983984985986987988989990991992993994995996997998999100010011002100310041005100610071008100910101011101210131014101510161017101810191020102110221023102410251026102710281029103010311032103310341035103610371038103910401041104210431044104510461047104810491050105110521053105410551056105710581059106010611062106310641065106610671068106910701071107210731074107510761077107810791080108110821083108410851086108710881089109010911092109310941095109610971098109911001101110211031104110511061107110811091110111111121113111411151116111711181119112011211122112311241125112611271128112911301131113211331134113511361137113811391140114111421143114411451146114711481149115011511152115311541155115611571158115911601161116211631164116511661167116811691170117111721173117411751176117711781179118011811182118311841185118611871188118911901191119211931194119511961197119811991200120112021203120412051206120712081209121012111212121312141215121612171218121912201221122212231224122512261227122812291230123112321233123412351236123712381239124012411242124312441245124612471248124912501251125212531254125512561257125812591260126112621263126412651266126712681269127012711272127312741275127612771278127912801281128212831284128512861287128812891290129112921293129412951296129712981299130013011302130313041305130613071308130913101311131213131314131513161317131813191320132113221323132413251326132713281329133013311332133313341335133613371338133913401341134213431344134513461347134813491350135113521353135413551356135713581359136013611362136313641365136613671368136913701371137213731374137513761377137813791380138113821383138413851386138713881389139013911392139313941395139613971398139914001401140214031404140514061407140814091410141114121413141414151416141714181419142014211422142314241425142614271428142914301431143214331434143514361437143814391440144114421443144414451446144714481449145014511452145314541455145614571458145914601461146214631464146514661467146814691470147114721473147414751476147714781479148014811482148314841485148614871488148914901491149214931494149514961497149814991500150115021503150415051506150715081509151015111512151315141515151615171518151915201521152215231524152515261527152815291530153115321533153415351536153715381539154015411542154315441545154615471548154915501551155215531554155515561557155815591560156115621563156415651566156715681569157015711572157315741575157615771578157915801581158215831584158515861587158815891590159115921593159415951596159715981599160016011602160316041605160616071608160916101611161216131614161516161617161816191620162116221623162416251626162716281629163016311632163316341635163616371638163916401641164216431644164516461647164816491650165116521653165416551656165716581659166016611662166316641665166616671668166916701671167216731674 |
- from __future__ import annotations
- import hashlib
- import os
- import subprocess
- import uuid
- from concurrent.futures import ThreadPoolExecutor
- from contextlib import contextmanager
- from pathlib import Path
- from threading import Barrier, Event, Lock
- import pytest
- from sqlalchemy import create_engine, inspect, text
- from sqlalchemy.engine import make_url
- from sqlalchemy.exc import DBAPIError, IntegrityError
- from sqlalchemy.orm import Session
- from app.core.connectors.builtin.oracle import OracleConnector
- from app.core.connectors.builtin.rest_catalog import RestCatalogConnector
- from app.core.connectors.builtin.sqlserver import SqlServerConnector
- from app.core.connectors.errors import (
- ConnectorAuthenticationError,
- ConnectorConfigurationError,
- ConnectorRateLimitError,
- ConnectorUpstreamError,
- )
- from app.core.connectors.identity import ConnectorIdentityRepository
- from app.core.connectors.registry import ConnectorRegistry
- from app.core.connectors.repository import ConnectorRepository
- from app.core.connectors.runtime import (
- ConnectorRuntime,
- deterministic_idempotency_key,
- )
- from app.core.connectors.sdk import (
- SDK_VERSION,
- SECRET_REF_PATTERN,
- Connector,
- ConnectorManifest,
- OperationRequest,
- OperationResult,
- )
- ROOT = Path(__file__).resolve().parents[2]
- pytestmark = pytest.mark.integration
- def _uid():
- return str(uuid.uuid4())
- def _alembic(url, command, target):
- subprocess.run(
- [
- str(ROOT / ".venv/bin/alembic"),
- "-c",
- str(ROOT / "alembic.ini"),
- command,
- target,
- ],
- cwd=ROOT,
- env={**os.environ, "SQLALCHEMY_DATABASE_URI": url},
- check=True,
- capture_output=True,
- text=True,
- )
- def _clear_connector_test_data(engine):
- tables = set(inspect(engine).get_table_names(schema="public"))
- with engine.begin() as connection:
- for table in (
- "connector_audit_events",
- "connector_graph_edges",
- "connector_evidence",
- "connector_checkpoints",
- "connector_run_attempts",
- "connector_runs",
- "connector_machine_credentials",
- "connector_principals",
- "connector_source_bindings",
- "connector_manifests",
- "connector_rate_limits",
- ):
- if table in tables:
- connection.execute(text(f"DELETE FROM public.{table}"))
- def _create_wp03_rollback_database(base_url: str) -> tuple[str, str]:
- """Create a uniquely named DB so historical WP03 rollback never crosses WP11 fences."""
- database = f"wp03_rollback_{uuid.uuid4().hex}"
- url = make_url(base_url)
- control_url = url.set(database="postgres").render_as_string(hide_password=False)
- control = create_engine(control_url, isolation_level="AUTOCOMMIT", pool_pre_ping=True)
- try:
- quoted = control.dialect.identifier_preparer.quote(database)
- with control.connect() as connection:
- connection.execute(text(f"CREATE DATABASE {quoted}"))
- finally:
- control.dispose()
- return database, url.set(database=database).render_as_string(hide_password=False)
- def _drop_wp03_rollback_database(base_url: str, database: str) -> None:
- url = make_url(base_url)
- control_url = url.set(database="postgres").render_as_string(hide_password=False)
- control = create_engine(control_url, isolation_level="AUTOCOMMIT", pool_pre_ping=True)
- try:
- quoted = control.dialect.identifier_preparer.quote(database)
- with control.connect() as connection:
- connection.execute(
- text(
- "SELECT pg_terminate_backend(pid) FROM pg_stat_activity "
- "WHERE datname=:database AND pid<>pg_backend_pid()"
- ),
- {"database": database},
- )
- connection.execute(text(f"DROP DATABASE IF EXISTS {quoted}"))
- finally:
- control.dispose()
- def test_connector_flask_postgres_security_bindings_and_machine_actions(
- monkeypatch,
- ):
- url = os.environ.get("TEST_DATABASE_URL")
- if not url:
- pytest.skip("TEST_DATABASE_URL is required")
- engine = create_engine(url, pool_pre_ping=True)
- _clear_connector_test_data(engine)
- engine.dispose()
- _alembic(url, "upgrade", "head")
- monkeypatch.setenv("DATABASE_URL", url)
- from app import create_app
- actor = _uid()
- monkeypatch.setattr(
- "app.core.system.permissions.authenticate_request",
- lambda: {"id": actor, "sub": actor, "roles": ["admin"]},
- )
- class SpyTransport:
- def __init__(self):
- self.calls = []
- def get_json(self, **values):
- self.calls.append(values)
- return {
- "assets": [
- {
- "key": "approved-asset",
- "name": "Approved",
- "namespace": "security",
- "type": "table",
- }
- ]
- }
- spy = SpyTransport()
- registry = ConnectorRegistry()
- registry.register(RestCatalogConnector(spy))
- app = create_app()
- app.config.update(TESTING=True)
- app.extensions["connector_registry"] = registry
- client = app.test_client()
- source, domain = _uid(), _uid()
- binding_response = client.post(
- "/api/datasource/connectors/source-bindings",
- json={
- "connector_id": "rest-catalog",
- "version": "1.0.0",
- "source_uid": source,
- "business_domain_uid": domain,
- "environment": "staging",
- "approved_config": {
- "base_url": "https://approved.example.test",
- "allowed_host": "approved.example.test",
- "credential_ref": "env:DATAOPS_CONNECTOR_APPROVED",
- },
- },
- )
- assert binding_response.status_code == 201
- binding = binding_response.get_json()["data"]
- principal_response = client.post(
- "/api/datasource/connectors/principals",
- json={
- "connector_id": "rest-catalog",
- "version": "1.0.0",
- "source_uid": source,
- "business_domain_uid": domain,
- "environment": "staging",
- "operations": ["discover", "cancel", "resume"],
- "scopes": {},
- "source_binding_uid": binding["binding_uid"],
- "source_binding_version": binding["binding_version"],
- },
- )
- assert principal_response.status_code == 201
- principal = principal_response.get_json()["data"]["principal_uid"]
- rebound_response = client.post(
- "/api/datasource/connectors/source-bindings",
- json={
- "binding_uid": binding["binding_uid"],
- "connector_id": "rest-catalog",
- "version": "1.0.0",
- "source_uid": source,
- "business_domain_uid": domain,
- "environment": "staging",
- "approved_config": {
- "base_url": "https://approved.example.test",
- "allowed_host": "approved.example.test",
- "credential_ref": "env:DATAOPS_CONNECTOR_APPROVED_V2",
- },
- },
- )
- assert rebound_response.status_code == 201
- rebound = rebound_response.get_json()["data"]
- assert rebound["binding_version"] == 2
- assert rebound["rebound_principals"] == 1
- with create_engine(url).connect() as connection:
- versions = connection.execute(
- text("""
- SELECT p.source_binding_version,
- ARRAY_AGG(b.status ORDER BY b.binding_version)
- FROM public.connector_principals p
- JOIN public.connector_source_bindings b
- ON b.uid=p.source_binding_uid
- WHERE p.uid=CAST(:principal AS uuid)
- GROUP BY p.source_binding_version
- """),
- {"principal": principal},
- ).one()
- assert versions == (2, ["revoked", "approved"])
- def issue(principal_uid):
- response = client.post(
- f"/api/datasource/connectors/principals/{principal_uid}/credentials",
- json={"ttl_seconds": 300},
- )
- assert response.status_code == 201
- return response.get_json()["data"]["credential"]
- hint = uuid.uuid4().hex + uuid.uuid4().hex
- token = issue(principal)
- empty_config = client.post(
- "/api/datasource/connectors/machine/runs",
- headers={"X-Connector-Credential": token},
- json={
- "connector_id": "rest-catalog",
- "version": "1.0.0",
- "source_uid": source,
- "business_domain_uid": domain,
- "environment": "staging",
- "process_key": "security-probe",
- "operation": "discover",
- "config": {},
- "scope": {},
- "idempotency_key": hint,
- },
- )
- assert empty_config.status_code == 400 and spy.calls == []
- attack = client.post(
- "/api/datasource/connectors/machine/runs",
- headers={"X-Connector-Credential": token},
- json={
- "connector_id": "rest-catalog",
- "version": "1.0.0",
- "source_uid": source,
- "business_domain_uid": domain,
- "environment": "staging",
- "process_key": "security-probe",
- "operation": "discover",
- "config": {
- "base_url": "https://attacker.example.test",
- "allowed_host": "attacker.example.test",
- "credential_ref": "env:DATAOPS_CONNECTOR_APPROVED",
- },
- "scope": {},
- "idempotency_key": hint,
- },
- )
- assert attack.status_code == 400 and spy.calls == []
- created = client.post(
- "/api/datasource/connectors/machine/runs",
- headers={"X-Connector-Credential": token},
- json={
- "connector_id": "rest-catalog",
- "version": "1.0.0",
- "source_uid": source,
- "business_domain_uid": domain,
- "environment": "staging",
- "process_key": "security-probe",
- "operation": "discover",
- "scope": {},
- "idempotency_key": hint,
- },
- )
- assert created.status_code == 201
- assert len(spy.calls) == 1
- assert spy.calls[0]["url"].startswith("https://approved.example.test/")
- assert spy.calls[0]["allowed_host"] == "approved.example.test"
- assert spy.calls[0]["credential_ref"] == "env:DATAOPS_CONNECTOR_APPROVED_V2"
- assert client.post(
- f"/api/datasource/connectors/runs/{hint}/cancel"
- ).status_code == 400
- assert client.post(
- f"/api/datasource/connectors/runs/{hint}/resume"
- ).status_code == 400
- second_source, second_domain = _uid(), _uid()
- second_binding = client.post(
- "/api/datasource/connectors/source-bindings",
- json={
- "connector_id": "rest-catalog",
- "version": "1.0.0",
- "source_uid": second_source,
- "business_domain_uid": second_domain,
- "environment": "production",
- "approved_config": {
- "base_url": "https://second.example.test",
- "allowed_host": "second.example.test",
- "credential_ref": "env:DATAOPS_CONNECTOR_SECOND",
- },
- },
- ).get_json()["data"]
- second_principal = client.post(
- "/api/datasource/connectors/principals",
- json={
- "connector_id": "rest-catalog",
- "version": "1.0.0",
- "source_uid": second_source,
- "business_domain_uid": second_domain,
- "environment": "production",
- "operations": ["discover", "cancel", "resume"],
- "scopes": {},
- "source_binding_uid": second_binding["binding_uid"],
- "source_binding_version": second_binding["binding_version"],
- },
- ).get_json()["data"]["principal_uid"]
- cross = client.post(
- f"/api/datasource/connectors/machine/runs/{hint}/cancel",
- headers={"X-Connector-Credential": issue(second_principal)},
- )
- assert cross.status_code == 403
- replay = client.post(
- f"/api/datasource/connectors/machine/runs/{hint}/cancel",
- headers={"X-Connector-Credential": token},
- )
- assert replay.status_code == 401
- with create_engine(url).begin() as connection:
- connection.execute(
- text("""
- UPDATE public.connector_runs
- SET status='running',cancel_requested=FALSE
- WHERE client_hint_hash=:hint
- """),
- {"hint": hashlib.sha256(hint.encode()).hexdigest()},
- )
- cancelled = client.post(
- f"/api/datasource/connectors/machine/runs/{hint}/cancel",
- headers={"X-Connector-Credential": issue(principal)},
- )
- assert cancelled.status_code == 200
- resumed = client.post(
- f"/api/datasource/connectors/machine/runs/{hint}/resume",
- headers={"X-Connector-Credential": issue(principal)},
- )
- assert resumed.status_code == 200
- conflict = client.post(
- "/api/datasource/connectors/machine/runs",
- headers={"X-Connector-Credential": issue(second_principal)},
- json={
- "connector_id": "rest-catalog",
- "version": "1.0.0",
- "source_uid": second_source,
- "business_domain_uid": second_domain,
- "environment": "production",
- "process_key": "different-process",
- "operation": "discover",
- "scope": {},
- "idempotency_key": hint,
- },
- )
- assert conflict.status_code == 409
- listed = client.get("/api/datasource/connectors/source-bindings")
- assert listed.status_code == 200
- serialized = str(listed.get_json())
- assert "credential_ref" not in serialized
- assert "DATAOPS_CONNECTOR_APPROVED" not in serialized
- with create_engine(url).connect() as connection:
- persisted = connection.execute(
- text("""
- SELECT COALESCE(string_agg(payload::text,''),'')
- FROM public.connector_evidence
- WHERE run_uid IN (
- SELECT uid FROM public.connector_runs WHERE principal_uid=CAST(:principal AS uuid)
- )
- """),
- {"principal": principal},
- ).scalar_one()
- assert "DATAOPS_CONNECTOR_APPROVED" not in persisted
- revoked = client.post(
- f"/api/datasource/connectors/source-bindings/{binding['binding_uid']}/revoke"
- )
- assert revoked.status_code == 200
- assert revoked.get_json()["data"] == {
- "revoked": True,
- "principals_deactivated": 1,
- }
- with create_engine(url).connect() as connection:
- statuses = connection.execute(
- text("""
- SELECT p.status,COALESCE(MAX(c.status),'none')
- FROM public.connector_principals p
- LEFT JOIN public.connector_machine_credentials c
- ON c.principal_uid=p.uid
- WHERE p.uid=CAST(:principal AS uuid)
- GROUP BY p.status
- """),
- {"principal": principal},
- ).one()
- assert statuses[0] == "revoked"
- def test_oracle_sqlserver_bindings_authentication_and_trusted_environment(monkeypatch):
- url = os.environ.get("TEST_DATABASE_URL")
- if not url:
- pytest.skip("TEST_DATABASE_URL is required")
- engine = create_engine(url, pool_pre_ping=True)
- _clear_connector_test_data(engine)
- _alembic(url, "upgrade", "head")
- monkeypatch.setenv("DATABASE_URL", url)
- from app import create_app
- actor = _uid()
- monkeypatch.setattr(
- "app.core.system.permissions.authenticate_request",
- lambda: {"id": actor, "sub": actor, "roles": ["admin"]},
- )
- provider_calls = []
- class Rows:
- def mappings(self):
- return self
- def all(self):
- return [
- {
- "schema_name": "APP",
- "asset_name": "ORDERS",
- "asset_type": "TABLE",
- "column_name": "ID",
- "ordinal_position": 1,
- "data_type": "NUMBER",
- "is_nullable": "NO",
- "column_default": None,
- "column_comment": None,
- }
- ]
- class Connection:
- def execute(self, _statement, _parameters):
- return Rows()
- @contextmanager
- def provider(
- source_uid,
- purpose,
- *,
- environment="production",
- allow_insecure_development=False,
- ):
- provider_calls.append(
- {
- "source_uid": source_uid,
- "purpose": purpose,
- "environment": environment,
- "allow_insecure_development": allow_insecure_development,
- }
- )
- yield Connection()
- class TrackingOracle(OracleConnector):
- seen_configs = []
- def discover(self, request):
- self.seen_configs.append(dict(request.config))
- return super().discover(request)
- class TrackingSqlServer(SqlServerConnector):
- seen_configs = []
- def discover(self, request):
- self.seen_configs.append(dict(request.config))
- return super().discover(request)
- registry = ConnectorRegistry()
- oracle = TrackingOracle(provider)
- sqlserver = TrackingSqlServer(provider)
- registry.register(oracle)
- registry.register(sqlserver)
- app = create_app()
- app.config.update(TESTING=True)
- app.extensions["connector_registry"] = registry
- client = app.test_client()
- rejected = client.post(
- "/api/datasource/connectors/principals",
- json={
- "connector_id": "oracle",
- "version": "1.0.0",
- "source_uid": _uid(),
- "business_domain_uid": _uid(),
- "environment": "staging",
- "operations": ["discover"],
- "scopes": {},
- },
- )
- assert rejected.status_code == 400
- for connector_id, environment, approved_config in (
- (
- "oracle",
- "staging",
- {"credential_ref": "env:DATAOPS_CONNECTOR_ORACLE"},
- ),
- (
- "sqlserver",
- "production",
- {
- "credential_ref": "env:DATAOPS_CONNECTOR_SQLSERVER",
- "allow_insecure_development": True,
- },
- ),
- ):
- source, domain = _uid(), _uid()
- binding_response = client.post(
- "/api/datasource/connectors/source-bindings",
- json={
- "connector_id": connector_id,
- "version": "1.0.0",
- "source_uid": source,
- "business_domain_uid": domain,
- "environment": environment,
- "approved_config": approved_config,
- },
- )
- assert binding_response.status_code == 201
- binding = binding_response.get_json()["data"]
- principal_response = client.post(
- "/api/datasource/connectors/principals",
- json={
- "connector_id": connector_id,
- "version": "1.0.0",
- "source_uid": source,
- "business_domain_uid": domain,
- "environment": environment,
- "operations": ["discover"],
- "scopes": {},
- "source_binding_uid": binding["binding_uid"],
- "source_binding_version": binding["binding_version"],
- },
- )
- assert principal_response.status_code == 201
- principal = principal_response.get_json()["data"]["principal_uid"]
- credential_response = client.post(
- f"/api/datasource/connectors/principals/{principal}/credentials"
- )
- assert credential_response.status_code == 201
- token = credential_response.get_json()["data"]["credential"]
- payload = {
- "connector_id": connector_id,
- "version": "1.0.0",
- "source_uid": source,
- "business_domain_uid": domain,
- "environment": environment,
- "process_key": f"{connector_id}-catalog",
- "operation": "discover",
- "scope": {},
- }
- if connector_id == "sqlserver":
- forged = client.post(
- "/api/datasource/connectors/machine/runs",
- headers={"X-Connector-Credential": token},
- json={**payload, "environment": "development"},
- )
- assert forged.status_code == 403
- created = client.post(
- "/api/datasource/connectors/machine/runs",
- headers={"X-Connector-Credential": token},
- json=payload,
- )
- assert created.status_code == 201
- assert created.get_json()["data"]["records"][0]["name"] == "ORDERS"
- assert oracle.seen_configs == [
- {"credential_ref": "env:DATAOPS_CONNECTOR_ORACLE"}
- ]
- assert sqlserver.seen_configs == [
- {
- "credential_ref": "env:DATAOPS_CONNECTOR_SQLSERVER",
- "allow_insecure_development": True,
- }
- ]
- assert provider_calls[0]["environment"] == "production"
- assert provider_calls[0]["allow_insecure_development"] is False
- assert provider_calls[1]["environment"] == "production"
- assert provider_calls[1]["allow_insecure_development"] is True
- with engine.connect() as connection:
- configs = connection.execute(
- text("""
- SELECT connector_id,safe_config
- FROM public.connector_runs
- WHERE connector_id IN ('oracle','sqlserver')
- ORDER BY connector_id
- """)
- ).all()
- assert configs == [
- ("oracle", {"credential_ref": "env:DATAOPS_CONNECTOR_ORACLE"}),
- (
- "sqlserver",
- {
- "credential_ref": "env:DATAOPS_CONNECTOR_SQLSERVER",
- "allow_insecure_development": True,
- },
- ),
- ]
- engine.dispose()
- def test_connector_postgres_attempt_lease_rejects_late_terminal_writes():
- url = os.environ.get("TEST_DATABASE_URL")
- if not url:
- pytest.skip("TEST_DATABASE_URL is required")
- engine = create_engine(url, pool_pre_ping=True)
- _clear_connector_test_data(engine)
- _alembic(url, "upgrade", "head")
- actor, source = _uid(), _uid()
- connector_id = f"lease-test-{uuid.uuid4().hex[:8]}"
- with engine.begin() as connection:
- connection.execute(
- text("""
- INSERT INTO public.connector_manifests
- (uid,connector_id,connector_version,sdk_version,display_name,
- capabilities,config_schema,status,created_by)
- VALUES(CAST(:uid AS uuid),:connector,'1.0.0','1.0','lease-test',
- '["discover","cancel","resume"]'::jsonb,
- CAST(:schema AS jsonb),
- 'active',CAST(:actor AS uuid))
- """),
- {
- "uid": _uid(),
- "connector": connector_id,
- "actor": actor,
- "schema": '{"type":"object","additionalProperties":false,"properties":{}}',
- },
- )
- key = uuid.uuid4().hex + uuid.uuid4().hex
- first = Session(engine)
- second = Session(engine)
- try:
- repository = ConnectorRepository(first)
- repository.claim(
- key,
- {
- "idempotency_key": key,
- "request_hash": key,
- "client_hint_hash": None,
- "connector_id": connector_id,
- "connector_version": "1.0.0",
- "source_uid": source,
- "operation": "discover",
- "config": {},
- "scope": {},
- "checkpoint": {},
- "cursor": {},
- "dry_run": True,
- "actor_uid": actor,
- },
- )
- lease_one, lease_two = _uid(), _uid()
- assert repository.update(
- key, status="running", attempt_count=1, lease_token=lease_one
- ).acquired
- assert repository.cancel(key)["status"] == "cancelled"
- with Session(engine) as observer:
- observed = ConnectorRepository(observer).get(key)
- assert observed["status"] == "cancelled"
- assert observed["cancel_requested"] is True
- assert ConnectorRepository(observer).is_cancel_requested(key) is True
- resumed = ConnectorRepository(second).update(
- key, status="resumable", cancel_requested=False
- )
- assert resumed.acquired
- assert ConnectorRepository(second).update(
- key, status="running", attempt_count=2, lease_token=lease_two
- ).acquired
- stale_success = ConnectorRepository(first).update(
- key,
- status="succeeded",
- expected_attempt=1,
- lease_token=lease_one,
- )
- stale_failure = ConnectorRepository(first).update(
- key,
- status="failed",
- error_category="upstream",
- error_code="late",
- expected_attempt=1,
- lease_token=lease_one,
- )
- assert stale_success.acquired is False
- assert stale_failure.acquired is False
- winner = ConnectorRepository(second).update(
- key,
- status="succeeded",
- expected_attempt=2,
- lease_token=lease_two,
- )
- assert winner.acquired
- final = ConnectorRepository(second).get(key)
- assert final["status"] == "succeeded"
- assert final["attempt_count"] == 2
- attempts = second.execute(
- text("""
- SELECT attempt_number,status,lease_token::text
- FROM public.connector_run_attempts
- WHERE run_uid=CAST(:run AS uuid)
- ORDER BY attempt_number
- """),
- {"run": final["uid"]},
- ).all()
- assert attempts == [
- (1, "cancelled", lease_one),
- (2, "succeeded", lease_two),
- ]
- finally:
- first.close()
- second.close()
- engine.dispose()
- def test_connector_postgres_repository_identity_graph_cas_and_473_rollback():
- shared_url = os.environ.get("TEST_DATABASE_URL")
- if not shared_url:
- pytest.skip("TEST_DATABASE_URL is required")
- rollback_database, url = _create_wp03_rollback_database(shared_url)
- engine = create_engine(url, pool_pre_ping=True)
- actor, source, domain = _uid(), _uid(), _uid()
- connector_id, version = f"pg-test-{uuid.uuid4().hex[:8]}", "1.0.0"
- _alembic(url, "upgrade", "20260802_476")
- try:
- tables = inspect(engine).get_table_names(schema="public")
- assert {"connector_runs", "connector_rate_limits"}.issubset(tables)
- binding_uid = _uid()
- with engine.begin() as connection:
- connection.execute(
- text("""
- INSERT INTO public.connector_manifests
- (uid,connector_id,connector_version,sdk_version,display_name,capabilities,config_schema,status,created_by)
- VALUES(CAST(:uid AS uuid),:connector,:version,'1.0','test','["discover","cancel","resume"]'::jsonb,
- '{"type":"object"}'::jsonb,'active',CAST(:actor AS uuid))
- """),
- {
- "uid": _uid(),
- "connector": connector_id,
- "version": version,
- "actor": actor,
- },
- )
- connection.execute(
- text("""
- INSERT INTO public.connector_source_bindings
- (uid,binding_version,connector_id,connector_version,source_uid,
- business_domain_uid,environment,approved_config,status,approved_by)
- VALUES(CAST(:uid AS uuid),1,:connector,:version,CAST(:source AS uuid),
- CAST(:domain AS uuid),'staging',
- '{"credential_ref":"env:DATAOPS_CONNECTOR_TEST"}'::jsonb,
- 'approved',CAST(:actor AS uuid))
- """),
- {
- "uid": binding_uid,
- "connector": connector_id,
- "version": version,
- "source": source,
- "domain": domain,
- "actor": actor,
- },
- )
- idem = uuid.uuid4().hex + uuid.uuid4().hex
- def claim(_index):
- with engine.begin() as connection:
- return connection.execute(
- text("""
- INSERT INTO public.connector_runs
- (uid,idempotency_key,connector_id,connector_version,source_uid,operation,status,attempt_count,
- checkpoint,cursor,dry_run,actor_uid,safe_config,scope,request_hash)
- VALUES(CAST(:uid AS uuid),:key,:connector,:version,CAST(:source AS uuid),'discover','running',0,
- '{}'::jsonb,'{}'::jsonb,TRUE,CAST(:actor AS uuid),'{}'::jsonb,'{}'::jsonb,:key)
- ON CONFLICT(idempotency_key) DO NOTHING RETURNING uid
- """),
- {
- "uid": _uid(),
- "key": idem,
- "connector": connector_id,
- "version": version,
- "source": source,
- "actor": actor,
- },
- ).scalar_one_or_none()
- with ThreadPoolExecutor(max_workers=4) as executor:
- assert sum(item is not None for item in executor.map(claim, range(4))) == 1
- principal = _uid()
- with engine.begin() as connection:
- connection.execute(
- text("""
- INSERT INTO public.connector_principals
- (uid,connector_id,connector_version,source_uid,business_domain_uid,environment,
- allowed_operations,allowed_scopes,status,created_by,
- source_binding_uid,source_binding_version)
- VALUES(CAST(:uid AS uuid),:connector,:version,CAST(:source AS uuid),CAST(:domain AS uuid),
- 'staging',ARRAY['discover'],'{}'::jsonb,'active',CAST(:actor AS uuid),
- CAST(:binding AS uuid),1)
- """),
- {
- "uid": principal,
- "connector": connector_id,
- "version": version,
- "source": source,
- "domain": domain,
- "actor": actor,
- "binding": binding_uid,
- },
- )
- with pytest.raises(IntegrityError), engine.begin() as connection:
- connection.execute(
- text("""
- INSERT INTO public.connector_machine_credentials(uid,principal_uid,token_hash,status,issued_by,expires_at)
- VALUES(CAST(:uid AS uuid),CAST(:principal AS uuid),:hash,'active',CAST(:actor AS uuid),
- CURRENT_TIMESTAMP+INTERVAL '901 seconds')
- """),
- {
- "uid": _uid(),
- "principal": principal,
- "hash": "a" * 64,
- "actor": actor,
- },
- )
- session = Session(engine)
- identity = ConnectorIdentityRepository(session)
- issued = identity.issue(principal, ttl_seconds=300, actor_uid=actor)
- authenticated = identity.authenticate(
- issued["credential"],
- connector_id=connector_id,
- connector_version=version,
- source_uid=source,
- business_domain_uid=domain,
- environment="staging",
- operation="discover",
- scope={},
- )
- assert authenticated["principal_uid"] == principal
- with pytest.raises(ConnectorAuthenticationError):
- identity.authenticate(
- issued["credential"],
- connector_id=connector_id,
- connector_version=version,
- source_uid=source,
- business_domain_uid=domain,
- environment="staging",
- operation="discover",
- scope={},
- )
- second = identity.issue(principal, ttl_seconds=300, actor_uid=actor)
- rotated = identity.rotate(
- second["credential_uid"], ttl_seconds=300, actor_uid=actor
- )
- assert identity.revoke(rotated["credential_uid"], actor) is True
- limiter_key = f"{connector_id}:{source}"
- for _ in range(30):
- with Session(engine) as limiter_session:
- ConnectorRepository(limiter_session).acquire_rate_limit(limiter_key)
- with Session(engine) as limiter_session, pytest.raises(ConnectorRateLimitError):
- ConnectorRepository(limiter_session).acquire_rate_limit(limiter_key)
- run_key = uuid.uuid4().hex + uuid.uuid4().hex
- repository = ConnectorRepository(session, principal)
- repository.claim(
- run_key,
- {
- "connector_id": connector_id,
- "connector_version": version,
- "source_uid": source,
- "operation": "discover",
- "checkpoint": {},
- "cursor": {},
- "dry_run": False,
- "principal_uid": principal,
- "business_domain_uid": domain,
- "environment": "staging",
- "process_key": "catalog-sync",
- "config": {"credential_ref": "secret:test/key"},
- "scope": {},
- },
- )
- running_graph = repository.graph(run_uid=repository.get(run_key)["uid"])
- assert running_graph["summary"]["edge_count"] == 3
- assert {node["type"] for node in running_graph["nodes"]} == {
- "source",
- "process",
- "business_domain",
- "run",
- }
- assert {edge["run_status"] for edge in running_graph["edges"]} == {"running"}
- run_lease = _uid()
- repository.update(
- run_key, status="running", attempt_count=1, lease_token=run_lease
- )
- repository.update(
- run_key,
- status="succeeded",
- result=OperationResult(
- records=({"asset_key": f"{source}:APP.ORDERS"},),
- cursor={"page": 2},
- checkpoint={"cursor": {"page": 2}, "snapshot": [{"key": "orders"}]},
- evidence={"safe": True},
- ),
- checkpoint={"cursor": {"page": 2}},
- cursor={"page": 2},
- error_category=None,
- expected_attempt=1,
- lease_token=run_lease,
- )
- graph = repository.graph(source_uid=source, business_domain_uid=domain)
- assert {node["type"] for node in graph["nodes"]} == {
- "source",
- "asset",
- "process",
- "business_domain",
- "run",
- }
- assert graph["summary"]["edge_count"] == 5
- process_graph = repository.graph(process_key="catalog-sync")
- assert process_graph["summary"]["edge_count"] >= 1
- assert {"process", "run"}.issubset(
- {node["type"] for node in process_graph["nodes"]}
- )
- failed_key = uuid.uuid4().hex + uuid.uuid4().hex
- repository.claim(
- failed_key,
- {
- "connector_id": connector_id,
- "connector_version": version,
- "source_uid": source,
- "operation": "discover",
- "checkpoint": {"page": 4},
- "cursor": {"page": 4},
- "dry_run": False,
- "principal_uid": principal,
- "business_domain_uid": domain,
- "environment": "staging",
- "process_key": "failed-catalog-sync",
- "config": {"credential_ref": "env:DATAOPS_CONNECTOR_TEST"},
- "scope": {},
- },
- )
- failed_lease = _uid()
- repository.update(
- failed_key,
- status="running",
- attempt_count=1,
- lease_token=failed_lease,
- )
- failed = repository.update(
- failed_key,
- status="failed",
- error_category="upstream",
- error_code="ConnectorUpstreamError",
- expected_attempt=1,
- lease_token=failed_lease,
- ).record
- assert failed["checkpoint"] == {"page": 4}
- failed_graph = repository.graph(run_uid=failed["uid"])
- assert failed_graph["summary"]["edge_count"] == 3
- assert {edge["run_status"] for edge in failed_graph["edges"]} == {"failed"}
- cancelled_key = uuid.uuid4().hex + uuid.uuid4().hex
- repository.claim(
- cancelled_key,
- {
- "connector_id": connector_id,
- "connector_version": version,
- "source_uid": source,
- "operation": "discover",
- "checkpoint": {"page": 1},
- "cursor": {"page": 1},
- "dry_run": False,
- "principal_uid": principal,
- "business_domain_uid": domain,
- "environment": "staging",
- "process_key": "catalog-sync",
- "config": {"credential_ref": "secret:test/key"},
- "scope": {},
- },
- )
- cancelled = Event()
- def cancel_worker():
- with Session(engine) as worker_session:
- result = ConnectorRepository(worker_session, principal).cancel(
- cancelled_key
- )
- cancelled.set()
- return result
- def completion_worker():
- assert cancelled.wait(timeout=5)
- with Session(engine) as worker_session:
- return ConnectorRepository(worker_session, principal).update(
- cancelled_key,
- status="succeeded",
- result=OperationResult(),
- checkpoint={},
- cursor={},
- error_category=None,
- )
- with ThreadPoolExecutor(max_workers=2) as executor:
- cancel_future = executor.submit(cancel_worker)
- complete_future = executor.submit(completion_worker)
- assert cancel_future.result()["status"] == "cancelled"
- completion_outcome = complete_future.result()
- assert completion_outcome.acquired is False
- after = completion_outcome.record
- assert after["status"] == "cancelled"
- cancelled_graph = repository.graph(run_uid=after["uid"])
- assert cancelled_graph["summary"]["edge_count"] == 3
- assert {edge["run_status"] for edge in cancelled_graph["edges"]} == {
- "cancelled"
- }
- with pytest.raises(IntegrityError), engine.begin() as connection:
- connection.execute(
- text("""
- INSERT INTO public.connector_runs(uid,idempotency_key,connector_id,connector_version,source_uid,
- operation,status,attempt_count,checkpoint,cursor,dry_run,actor_uid,safe_config,scope)
- VALUES(CAST(:uid AS uuid),:key,:connector,:version,CAST(:source AS uuid),'discover','running',0,
- '{}'::jsonb,'{}'::jsonb,FALSE,CAST(:actor AS uuid),'{}'::jsonb,'{}'::jsonb)
- """),
- {
- "uid": _uid(),
- "key": uuid.uuid4().hex + uuid.uuid4().hex,
- "connector": connector_id,
- "version": version,
- "source": source,
- "actor": actor,
- },
- )
- class RuntimeProbeConnector(Connector):
- manifest = ConnectorManifest(
- connector_id=connector_id,
- version=version,
- sdk_version=SDK_VERSION,
- display_name="PostgreSQL Runtime Probe",
- capabilities=(
- "discover",
- "snapshot",
- "incremental",
- "lineage",
- "profile",
- "cancel",
- "resume",
- "evidence",
- ),
- config_schema={
- "type": "object",
- "additionalProperties": False,
- "required": ["credential_ref"],
- "properties": {
- "credential_ref": {
- "type": "string",
- "pattern": SECRET_REF_PATTERN.pattern,
- }
- },
- },
- )
- def __init__(self):
- self.discover_calls = 0
- self.resume_calls = 0
- self.discover_succeeds = False
- self.resume_failures = 2
- self.call_lock = Lock()
- def discover(self, request):
- with self.call_lock:
- self.discover_calls += 1
- if self.discover_succeeds:
- return OperationResult(evidence={"concurrent": True})
- raise ConnectorUpstreamError()
- snapshot = incremental = lineage = profile = evidence = discover
- def cancel(self, request):
- return OperationResult(status="cancelled")
- def resume(self, request):
- with self.call_lock:
- self.resume_calls += 1
- call = self.resume_calls
- if call <= self.resume_failures:
- raise ConnectorUpstreamError()
- return OperationResult(
- checkpoint={"page": 8},
- cursor={"page": 9},
- evidence={"resumed": True},
- )
- probe = RuntimeProbeConnector()
- registry = ConnectorRegistry()
- registry.register(probe)
- capped_source = _uid()
- capped_request = OperationRequest(
- source_uid=capped_source,
- operation="discover",
- config={"credential_ref": "env:DATAOPS_CONNECTOR_TEST"},
- principal_uid=principal,
- business_domain_uid=domain,
- environment="staging",
- process_key="attempt-cap",
- )
- capped_runtime = ConnectorRuntime(
- registry,
- store=ConnectorRepository(session, principal),
- sleeper=lambda _seconds: None,
- )
- with pytest.raises(ConnectorUpstreamError):
- capped_runtime.execute(connector_id, version, capped_request)
- with pytest.raises(ConnectorUpstreamError):
- capped_runtime.execute(connector_id, version, capped_request)
- with pytest.raises(ConnectorConfigurationError):
- capped_runtime.execute(connector_id, version, capped_request)
- capped_key = capped_request.idempotency_key or deterministic_idempotency_key(
- connector_id, version, capped_request
- )
- capped = ConnectorRepository(session).get(capped_key)
- assert capped["attempt_count"] == 5 and probe.discover_calls == 5
- attempts = session.execute(
- text("""
- SELECT COUNT(*),MAX(attempt_number)
- FROM public.connector_run_attempts
- WHERE run_uid=CAST(:run AS uuid)
- """),
- {"run": capped["uid"]},
- ).one()
- assert attempts == (5, 5)
- rejected_source = _uid()
- rejected_key = f"{connector_id}:{rejected_source}"
- for _ in range(30):
- ConnectorRepository(session).acquire_rate_limit(rejected_key)
- rejected_request = OperationRequest(
- source_uid=rejected_source,
- operation="discover",
- config={"credential_ref": "env:DATAOPS_CONNECTOR_TEST"},
- dry_run=True,
- )
- rejected_idempotency = deterministic_idempotency_key(
- connector_id, version, rejected_request
- )
- with pytest.raises(ConnectorRateLimitError):
- ConnectorRuntime(
- registry, store=ConnectorRepository(session), sleeper=lambda _: None
- ).execute(connector_id, version, rejected_request)
- assert ConnectorRepository(session).get(rejected_idempotency) is None
- resume_source = _uid()
- resume_key = uuid.uuid4().hex + uuid.uuid4().hex
- resume_repository = ConnectorRepository(session, principal)
- resume_repository.claim(
- resume_key,
- {
- "connector_id": connector_id,
- "connector_version": version,
- "source_uid": resume_source,
- "operation": "discover",
- "checkpoint": {"page": 7},
- "cursor": {"page": 7},
- "dry_run": False,
- "principal_uid": principal,
- "business_domain_uid": domain,
- "environment": "staging",
- "process_key": "resume-state-machine",
- "config": {"credential_ref": "env:DATAOPS_CONNECTOR_TEST"},
- "scope": {},
- },
- )
- resume_repository.update(
- resume_key, status="running", attempt_count=1, lease_token=_uid()
- )
- resume_repository.cancel(resume_key)
- resume_runtime = ConnectorRuntime(
- registry,
- store=resume_repository,
- max_attempts=2,
- sleeper=lambda _seconds: None,
- )
- with pytest.raises(ConnectorUpstreamError):
- resume_runtime.resume(resume_key)
- resume_failed = resume_repository.get(resume_key)
- assert resume_failed["status"] == "failed"
- assert resume_failed["attempt_count"] == 3
- assert resume_failed["checkpoint"] == {"page": 7}
- resumed_result = resume_runtime.resume(resume_key)
- assert resumed_result.checkpoint == {"page": 8}
- resumed = resume_repository.get(resume_key)
- assert resumed["status"] == "succeeded" and resumed["attempt_count"] == 4
- persisted = session.execute(
- text("""
- SELECT
- (SELECT COUNT(*) FROM public.connector_run_attempts WHERE run_uid=CAST(:run AS uuid)),
- (SELECT COUNT(*) FROM public.connector_evidence WHERE run_uid=CAST(:run AS uuid)),
- (SELECT COUNT(*) FROM public.connector_checkpoints WHERE run_uid=CAST(:run AS uuid))
- """),
- {"run": resumed["uid"]},
- ).one()
- assert persisted == (4, 1, 1)
- probe.discover_succeeds = True
- probe.discover_calls = 0
- retry_hint = uuid.uuid4().hex + uuid.uuid4().hex
- retry_source = _uid()
- retry_request = OperationRequest(
- source_uid=retry_source,
- operation="discover",
- config={"credential_ref": "env:DATAOPS_CONNECTOR_TEST"},
- idempotency_key=retry_hint,
- principal_uid=principal,
- business_domain_uid=domain,
- environment="staging",
- process_key="concurrent-retry",
- )
- retry_key = deterministic_idempotency_key(
- connector_id, version, retry_request
- )
- retry_seed = ConnectorRepository(session, principal)
- retry_seed.claim(
- retry_key,
- {
- "connector_id": connector_id,
- "connector_version": version,
- "source_uid": retry_source,
- "operation": "discover",
- "checkpoint": {"page": 1},
- "cursor": {"page": 1},
- "dry_run": False,
- "principal_uid": principal,
- "business_domain_uid": domain,
- "environment": "staging",
- "process_key": "concurrent-retry",
- "config": {"credential_ref": "env:DATAOPS_CONNECTOR_TEST"},
- "scope": {},
- },
- )
- retry_lease = _uid()
- retry_seed.update(
- retry_key, status="running", attempt_count=1, lease_token=retry_lease
- )
- retry_seed.update(
- retry_key,
- status="failed",
- error_category="upstream",
- error_code="ConnectorUpstreamError",
- expected_attempt=1,
- lease_token=retry_lease,
- )
- retry_gate = Barrier(2)
- class RetryBarrierRepository(ConnectorRepository):
- def claim(self, key, record):
- claimed, created = super().claim(key, record)
- if not created:
- retry_gate.wait(timeout=5)
- return claimed, created
- def competing_retry():
- with Session(engine) as worker_session:
- worker_session.execute(text("SET statement_timeout='5s'"))
- worker_session.commit()
- try:
- result = ConnectorRuntime(
- registry,
- store=RetryBarrierRepository(worker_session, principal),
- max_attempts=1,
- sleeper=lambda _seconds: None,
- ).execute(connector_id, version, retry_request)
- return "success", result.status
- except Exception as error:
- return "error", type(error).__name__
- with ThreadPoolExecutor(max_workers=2) as executor:
- retry_futures = [executor.submit(competing_retry) for _index in range(2)]
- retry_results = [future.result(timeout=10) for future in retry_futures]
- assert sorted(item[0] for item in retry_results) == ["error", "success"]
- assert [item[1] for item in retry_results if item[0] == "error"] == [
- "ConnectorConfigurationError"
- ]
- assert probe.discover_calls == 1
- retried = ConnectorRepository(session).get(retry_key)
- assert retried["status"] == "succeeded" and retried["attempt_count"] == 2
- retry_attempts = session.execute(
- text("""
- SELECT ARRAY_AGG(attempt_number ORDER BY attempt_number)
- FROM public.connector_run_attempts
- WHERE run_uid=CAST(:run AS uuid)
- """),
- {"run": retried["uid"]},
- ).scalar_one()
- assert retry_attempts == [1, 2]
- probe.resume_failures = 0
- probe.resume_calls = 0
- competing_resume_key = uuid.uuid4().hex + uuid.uuid4().hex
- competing_resume_source = _uid()
- resume_seed = ConnectorRepository(session, principal)
- resume_seed.claim(
- competing_resume_key,
- {
- "connector_id": connector_id,
- "connector_version": version,
- "source_uid": competing_resume_source,
- "operation": "discover",
- "checkpoint": {"page": 5},
- "cursor": {"page": 5},
- "dry_run": False,
- "principal_uid": principal,
- "business_domain_uid": domain,
- "environment": "staging",
- "process_key": "concurrent-resume",
- "config": {"credential_ref": "env:DATAOPS_CONNECTOR_TEST"},
- "scope": {},
- },
- )
- resume_seed.update(
- competing_resume_key,
- status="running",
- attempt_count=1,
- lease_token=_uid(),
- )
- resume_seed.cancel(competing_resume_key)
- resume_gate = Barrier(2)
- class ResumeBarrierRepository(ConnectorRepository):
- def __init__(self, worker_session, actor_uid):
- super().__init__(worker_session, actor_uid)
- self.waited = False
- def get(self, key):
- record = super().get(key)
- if not self.waited and record and record.get("status") == "cancelled":
- self.waited = True
- resume_gate.wait(timeout=5)
- return record
- def competing_resume():
- with Session(engine) as worker_session:
- worker_session.execute(text("SET statement_timeout='5s'"))
- worker_session.commit()
- try:
- result = ConnectorRuntime(
- registry,
- store=ResumeBarrierRepository(worker_session, principal),
- max_attempts=1,
- sleeper=lambda _seconds: None,
- ).resume(competing_resume_key)
- return "success", result.status
- except Exception as error:
- return "error", type(error).__name__
- with ThreadPoolExecutor(max_workers=2) as executor:
- resume_futures = [executor.submit(competing_resume) for _index in range(2)]
- resume_results = [future.result(timeout=10) for future in resume_futures]
- assert sorted(item[0] for item in resume_results) == ["error", "success"]
- assert [item[1] for item in resume_results if item[0] == "error"] == [
- "ConnectorConfigurationError"
- ]
- assert probe.resume_calls == 1
- concurrently_resumed = ConnectorRepository(session).get(competing_resume_key)
- assert concurrently_resumed["status"] == "succeeded"
- assert concurrently_resumed["attempt_count"] == 2
- resume_attempts = session.execute(
- text("""
- SELECT ARRAY_AGG(attempt_number ORDER BY attempt_number)
- FROM public.connector_run_attempts
- WHERE run_uid=CAST(:run AS uuid)
- """),
- {"run": concurrently_resumed["uid"]},
- ).scalar_one()
- assert resume_attempts == [1, 2]
- session.close()
- with engine.begin() as connection:
- for table in (
- "connector_audit_events",
- "connector_graph_edges",
- "connector_evidence",
- "connector_checkpoints",
- "connector_run_attempts",
- "connector_runs",
- "connector_machine_credentials",
- "connector_principals",
- "connector_source_bindings",
- "connector_manifests",
- "connector_rate_limits",
- ):
- connection.execute(text(f"DELETE FROM public.{table}"))
- _alembic(url, "downgrade", "20260802_472")
- ambiguous_connector = f"ambiguous-{uuid.uuid4().hex[:8]}"
- ambiguous_source, ambiguous_domain = _uid(), _uid()
- with engine.begin() as connection:
- for ambiguous_version in ("1.0.0", "2.0.0"):
- connection.execute(
- text("""
- INSERT INTO public.connector_manifests
- (uid,connector_id,connector_version,sdk_version,display_name,
- capabilities,config_schema,status,created_by)
- VALUES(CAST(:uid AS uuid),:connector,:version,'1.0','ambiguous',
- '["discover"]'::jsonb,'{"type":"object"}'::jsonb,
- 'active',CAST(:actor AS uuid))
- """),
- {
- "uid": _uid(),
- "connector": ambiguous_connector,
- "version": ambiguous_version,
- "actor": actor,
- },
- )
- connection.execute(
- text("""
- INSERT INTO public.connector_principals
- (uid,connector_id,source_uid,business_domain_uid,environment,
- allowed_operations,allowed_scopes,status,created_by)
- VALUES(CAST(:uid AS uuid),:connector,CAST(:source AS uuid),
- CAST(:domain AS uuid),'staging',ARRAY['discover'],
- '{}'::jsonb,'active',CAST(:actor AS uuid))
- """),
- {
- "uid": _uid(),
- "connector": ambiguous_connector,
- "source": ambiguous_source,
- "domain": ambiguous_domain,
- "actor": actor,
- },
- )
- with pytest.raises(subprocess.CalledProcessError) as ambiguous_upgrade:
- _alembic(url, "upgrade", "20260802_473")
- assert "rejected before backfill" in ambiguous_upgrade.value.stderr
- with engine.connect() as connection:
- assert connection.execute(
- text("SELECT version_num FROM alembic_version")
- ).scalar_one() == "20260802_472"
- assert "connector_version" not in {
- item["name"]
- for item in inspect(engine).get_columns(
- "connector_principals", schema="public"
- )
- }
- single_connector = f"single-{uuid.uuid4().hex[:8]}"
- single_principal, single_source, single_domain = _uid(), _uid(), _uid()
- with engine.begin() as connection:
- connection.execute(
- text("DELETE FROM public.connector_principals")
- )
- connection.execute(text("DELETE FROM public.connector_manifests"))
- connection.execute(
- text("""
- INSERT INTO public.connector_manifests
- (uid,connector_id,connector_version,sdk_version,display_name,
- capabilities,config_schema,status,created_by)
- VALUES(CAST(:uid AS uuid),:connector,'1.0.0','1.0','single',
- '["discover"]'::jsonb,'{"type":"object"}'::jsonb,
- 'active',CAST(:actor AS uuid))
- """),
- {"uid": _uid(), "connector": single_connector, "actor": actor},
- )
- connection.execute(
- text("""
- INSERT INTO public.connector_principals
- (uid,connector_id,source_uid,business_domain_uid,environment,
- allowed_operations,allowed_scopes,status,created_by)
- VALUES(CAST(:uid AS uuid),:connector,CAST(:source AS uuid),
- CAST(:domain AS uuid),'staging',ARRAY['discover'],
- '{}'::jsonb,'active',CAST(:actor AS uuid))
- """),
- {
- "uid": single_principal,
- "connector": single_connector,
- "source": single_source,
- "domain": single_domain,
- "actor": actor,
- },
- )
- _alembic(url, "upgrade", "20260802_475")
- with pytest.raises(subprocess.CalledProcessError) as unbound_upgrade:
- _alembic(url, "upgrade", "20260802_476")
- assert "bind every active enterprise principal first" in (
- unbound_upgrade.value.stderr
- )
- binding_uid = _uid()
- revoked_principal = _uid()
- revoked_source, revoked_domain, revoked_binding = _uid(), _uid(), _uid()
- with engine.begin() as connection:
- connection.execute(
- text("""
- INSERT INTO public.connector_source_bindings
- (uid,binding_version,connector_id,connector_version,source_uid,
- business_domain_uid,environment,approved_config,status,approved_by)
- VALUES(CAST(:uid AS uuid),1,:connector,'1.0.0',CAST(:source AS uuid),
- CAST(:domain AS uuid),'staging','{}'::jsonb,'approved',CAST(:actor AS uuid))
- """),
- {
- "uid": binding_uid,
- "connector": single_connector,
- "source": single_source,
- "domain": single_domain,
- "actor": actor,
- },
- )
- connection.execute(
- text("""
- UPDATE public.connector_principals
- SET source_binding_uid=CAST(:binding AS uuid),source_binding_version=1
- WHERE uid=CAST(:principal AS uuid)
- """),
- {"binding": binding_uid, "principal": single_principal},
- )
- connection.execute(
- text("""
- INSERT INTO public.connector_principals
- (uid,connector_id,connector_version,source_uid,business_domain_uid,
- environment,allowed_operations,allowed_scopes,status,created_by)
- VALUES(CAST(:uid AS uuid),:connector,'1.0.0',CAST(:source AS uuid),
- CAST(:domain AS uuid),'staging',ARRAY['discover'],
- '{}'::jsonb,'revoked',CAST(:actor AS uuid))
- """),
- {
- "uid": revoked_principal,
- "connector": single_connector,
- "source": revoked_source,
- "domain": revoked_domain,
- "actor": actor,
- },
- )
- _alembic(url, "upgrade", "20260802_476")
- with pytest.raises(DBAPIError), engine.begin() as connection:
- connection.execute(
- text("""
- UPDATE public.connector_principals SET status='active'
- WHERE uid=CAST(:principal AS uuid)
- """),
- {"principal": revoked_principal},
- )
- with engine.begin() as connection:
- connection.execute(
- text("""
- INSERT INTO public.connector_source_bindings
- (uid,binding_version,connector_id,connector_version,source_uid,
- business_domain_uid,environment,approved_config,status,approved_by)
- VALUES(CAST(:uid AS uuid),1,:connector,'1.0.0',CAST(:source AS uuid),
- CAST(:domain AS uuid),'staging','{}'::jsonb,'approved',CAST(:actor AS uuid))
- """),
- {
- "uid": revoked_binding,
- "connector": single_connector,
- "source": revoked_source,
- "domain": revoked_domain,
- "actor": actor,
- },
- )
- connection.execute(
- text("""
- UPDATE public.connector_principals
- SET source_binding_uid=CAST(:binding AS uuid),source_binding_version=1
- WHERE uid=CAST(:principal AS uuid)
- """),
- {"binding": revoked_binding, "principal": revoked_principal},
- )
- connection.execute(
- text("""
- UPDATE public.connector_principals SET status='active'
- WHERE uid=CAST(:principal AS uuid)
- """),
- {"principal": revoked_principal},
- )
- assert connection.execute(
- text("""
- SELECT status FROM public.connector_principals
- WHERE uid=CAST(:principal AS uuid)
- """),
- {"principal": revoked_principal},
- ).scalar_one() == "active"
- with pytest.raises(DBAPIError), engine.begin() as connection:
- connection.execute(
- text("""
- INSERT INTO public.connector_manifests
- (uid,connector_id,connector_version,sdk_version,display_name,
- capabilities,config_schema,status,created_by)
- VALUES(CAST(:uid AS uuid),:connector,'2.0.0','1.0','ambiguous',
- '["discover"]'::jsonb,'{"type":"object"}'::jsonb,
- 'active',CAST(:actor AS uuid))
- """),
- {"uid": _uid(), "connector": single_connector, "actor": actor},
- )
- with pytest.raises(subprocess.CalledProcessError) as guarded_downgrade:
- _alembic(url, "downgrade", "20260802_474")
- assert "explicitly remove source bindings" in guarded_downgrade.value.stderr
- _alembic(url, "downgrade", "20260802_475")
- with engine.begin() as connection:
- connection.execute(
- text("""
- UPDATE public.connector_principals
- SET source_binding_uid=NULL,source_binding_version=NULL
- WHERE uid IN (CAST(:principal AS uuid),CAST(:revoked AS uuid))
- """),
- {"principal": single_principal, "revoked": revoked_principal},
- )
- connection.execute(
- text(
- "DELETE FROM public.connector_source_bindings WHERE uid=CAST(:uid AS uuid)"
- ),
- {"uid": binding_uid},
- )
- connection.execute(
- text(
- "DELETE FROM public.connector_source_bindings WHERE uid=CAST(:uid AS uuid)"
- ),
- {"uid": revoked_binding},
- )
- _alembic(url, "downgrade", "20260802_474")
- assert "connector_source_bindings" not in inspect(engine).get_table_names(
- schema="public"
- )
- _clear_connector_test_data(engine)
- _alembic(url, "upgrade", "20260802_476")
- finally:
- try:
- _clear_connector_test_data(engine)
- finally:
- engine.dispose()
- _drop_wp03_rollback_database(shared_url, rollback_database)
- with create_engine(shared_url, pool_pre_ping=True).connect() as connection:
- assert connection.execute(text("SELECT version_num FROM alembic_version")).scalar_one() == "20260818_545"
|