| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160 |
- from __future__ import annotations
- from contextlib import contextmanager
- import json
- import re
- import pytest
- from app.core.connectors.builtin import register_builtin_connectors
- from app.core.connectors.builtin.postgresql import PostgreSQLConnector
- from app.core.connectors.errors import ConnectorCancelledError
- from app.core.connectors.registry import ConnectorRegistry
- from app.core.connectors.sdk import OperationRequest
- _connections = []
- class _Rows:
- def __init__(self, rows):
- self.rows = rows
- def mappings(self):
- return self
- def all(self):
- return self.rows
- class _Connection:
- def __init__(self, rows):
- self.rows = rows
- self.statement = None
- self.parameters = None
- def execute(self, statement, parameters):
- self.statement = str(statement)
- self.parameters = parameters
- return _Rows(self.rows)
- @contextmanager
- def _provider(source_uid, purpose):
- assert source_uid == "source-postgres"
- assert purpose == "metadata_collection"
- connection = _Connection(
- [
- {
- "schema_name": "public",
- "asset_name": "orders",
- "asset_type": "TABLE",
- "column_name": "id",
- "ordinal_position": 1,
- "data_type": "bigint",
- "is_nullable": "NO",
- "column_default": None,
- "column_comment": "order key",
- },
- {
- "schema_name": "internal",
- "asset_name": "jobs",
- "asset_type": "TABLE",
- "column_name": "id",
- "ordinal_position": 1,
- "data_type": "bigint",
- "is_nullable": "NO",
- "column_default": None,
- "column_comment": "job key",
- },
- ]
- )
- _connections.append(connection)
- yield connection
- def _request(operation="discover", **values):
- defaults = {
- "source_uid": "source-postgres",
- "operation": operation,
- "config": {"credential_ref": "vault:postgresql/catalog"},
- "scope": {"include_schemas": ["public"]},
- }
- defaults.update(values)
- return OperationRequest(**defaults)
- def test_postgresql_registry_connector_collects_filtered_read_only_catalog():
- _connections.clear()
- registry = ConnectorRegistry()
- register_builtin_connectors(registry, connection_provider=_provider)
- connector = registry.resolve("postgresql", "1.0.0")
- result = connector.discover(_request())
- assert isinstance(connector, PostgreSQLConnector)
- assert connector.manifest.display_name == "PostgreSQL"
- assert connector.manifest.capabilities == (
- "discover",
- "snapshot",
- "incremental",
- "cancel",
- "resume",
- "evidence",
- )
- assert [record["asset_key"] for record in result.records] == [
- "source-postgres:public.orders"
- ]
- assert result.checkpoint["snapshot_summary"]["record_count"] == 1
- assert result.evidence["query_kind"] == "read_only_metadata"
- assert "postgresql/catalog" not in json.dumps(result.__dict__)
- assert len(_connections) == 1
- executed_sql = " ".join(_connections[0].statement.lower().split())
- assert "information_schema.columns" in executed_sql
- assert "pg_catalog." in executed_sql
- assert "not in ('pg_catalog', 'information_schema')" in executed_sql
- assert not re.search(
- r"\b(?:insert|update|delete|alter|drop|truncate|create)\b", executed_sql
- )
- assert _connections[0].parameters == {}
- def test_postgresql_incremental_resume_cancel_and_evidence_are_safe():
- connector = PostgreSQLConnector(_provider)
- initial = connector.snapshot(_request("snapshot"))
- incremental = connector.incremental(
- _request("incremental", checkpoint=initial.checkpoint)
- )
- resumed = connector.resume(_request("resume", checkpoint=initial.checkpoint))
- cancelled = connector.cancel(_request("cancel", checkpoint=initial.checkpoint))
- evidence = connector.evidence(_request("evidence"))
- assert incremental.evidence["diff"]["changed"] is False
- assert resumed.status == "succeeded"
- assert cancelled.status == "cancelled"
- assert evidence.evidence["secret_material"] is False
- def test_postgresql_cancel_probe_stops_catalog_collection():
- with pytest.raises(ConnectorCancelledError):
- PostgreSQLConnector(_provider).discover(
- _request(cancel_probe=lambda: True)
- )
- def test_postgresql_health_and_compatibility_fail_closed_without_driver(monkeypatch):
- monkeypatch.setattr(
- "importlib.util.find_spec",
- lambda name: None if name == "psycopg2" else object(),
- )
- connector = PostgreSQLConnector(_provider)
- health = connector.health({"credential_ref": "vault:postgresql/catalog"})
- compatibility = connector.compatibility()
- assert health.status == "degraded"
- assert health.detail == "optional driver unavailable"
- assert compatibility.compatible is False
- assert compatibility.connector_version == "1.0.0"
- assert compatibility.detail == "optional driver unavailable"
|