test_postgresql_enterprise_connector.py 5.1 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160
  1. from __future__ import annotations
  2. from contextlib import contextmanager
  3. import json
  4. import re
  5. import pytest
  6. from app.core.connectors.builtin import register_builtin_connectors
  7. from app.core.connectors.builtin.postgresql import PostgreSQLConnector
  8. from app.core.connectors.errors import ConnectorCancelledError
  9. from app.core.connectors.registry import ConnectorRegistry
  10. from app.core.connectors.sdk import OperationRequest
  11. _connections = []
  12. class _Rows:
  13. def __init__(self, rows):
  14. self.rows = rows
  15. def mappings(self):
  16. return self
  17. def all(self):
  18. return self.rows
  19. class _Connection:
  20. def __init__(self, rows):
  21. self.rows = rows
  22. self.statement = None
  23. self.parameters = None
  24. def execute(self, statement, parameters):
  25. self.statement = str(statement)
  26. self.parameters = parameters
  27. return _Rows(self.rows)
  28. @contextmanager
  29. def _provider(source_uid, purpose):
  30. assert source_uid == "source-postgres"
  31. assert purpose == "metadata_collection"
  32. connection = _Connection(
  33. [
  34. {
  35. "schema_name": "public",
  36. "asset_name": "orders",
  37. "asset_type": "TABLE",
  38. "column_name": "id",
  39. "ordinal_position": 1,
  40. "data_type": "bigint",
  41. "is_nullable": "NO",
  42. "column_default": None,
  43. "column_comment": "order key",
  44. },
  45. {
  46. "schema_name": "internal",
  47. "asset_name": "jobs",
  48. "asset_type": "TABLE",
  49. "column_name": "id",
  50. "ordinal_position": 1,
  51. "data_type": "bigint",
  52. "is_nullable": "NO",
  53. "column_default": None,
  54. "column_comment": "job key",
  55. },
  56. ]
  57. )
  58. _connections.append(connection)
  59. yield connection
  60. def _request(operation="discover", **values):
  61. defaults = {
  62. "source_uid": "source-postgres",
  63. "operation": operation,
  64. "config": {"credential_ref": "vault:postgresql/catalog"},
  65. "scope": {"include_schemas": ["public"]},
  66. }
  67. defaults.update(values)
  68. return OperationRequest(**defaults)
  69. def test_postgresql_registry_connector_collects_filtered_read_only_catalog():
  70. _connections.clear()
  71. registry = ConnectorRegistry()
  72. register_builtin_connectors(registry, connection_provider=_provider)
  73. connector = registry.resolve("postgresql", "1.0.0")
  74. result = connector.discover(_request())
  75. assert isinstance(connector, PostgreSQLConnector)
  76. assert connector.manifest.display_name == "PostgreSQL"
  77. assert connector.manifest.capabilities == (
  78. "discover",
  79. "snapshot",
  80. "incremental",
  81. "cancel",
  82. "resume",
  83. "evidence",
  84. )
  85. assert [record["asset_key"] for record in result.records] == [
  86. "source-postgres:public.orders"
  87. ]
  88. assert result.checkpoint["snapshot_summary"]["record_count"] == 1
  89. assert result.evidence["query_kind"] == "read_only_metadata"
  90. assert "postgresql/catalog" not in json.dumps(result.__dict__)
  91. assert len(_connections) == 1
  92. executed_sql = " ".join(_connections[0].statement.lower().split())
  93. assert "information_schema.columns" in executed_sql
  94. assert "pg_catalog." in executed_sql
  95. assert "not in ('pg_catalog', 'information_schema')" in executed_sql
  96. assert not re.search(
  97. r"\b(?:insert|update|delete|alter|drop|truncate|create)\b", executed_sql
  98. )
  99. assert _connections[0].parameters == {}
  100. def test_postgresql_incremental_resume_cancel_and_evidence_are_safe():
  101. connector = PostgreSQLConnector(_provider)
  102. initial = connector.snapshot(_request("snapshot"))
  103. incremental = connector.incremental(
  104. _request("incremental", checkpoint=initial.checkpoint)
  105. )
  106. resumed = connector.resume(_request("resume", checkpoint=initial.checkpoint))
  107. cancelled = connector.cancel(_request("cancel", checkpoint=initial.checkpoint))
  108. evidence = connector.evidence(_request("evidence"))
  109. assert incremental.evidence["diff"]["changed"] is False
  110. assert resumed.status == "succeeded"
  111. assert cancelled.status == "cancelled"
  112. assert evidence.evidence["secret_material"] is False
  113. def test_postgresql_cancel_probe_stops_catalog_collection():
  114. with pytest.raises(ConnectorCancelledError):
  115. PostgreSQLConnector(_provider).discover(
  116. _request(cancel_probe=lambda: True)
  117. )
  118. def test_postgresql_health_and_compatibility_fail_closed_without_driver(monkeypatch):
  119. monkeypatch.setattr(
  120. "importlib.util.find_spec",
  121. lambda name: None if name == "psycopg2" else object(),
  122. )
  123. connector = PostgreSQLConnector(_provider)
  124. health = connector.health({"credential_ref": "vault:postgresql/catalog"})
  125. compatibility = connector.compatibility()
  126. assert health.status == "degraded"
  127. assert health.detail == "optional driver unavailable"
  128. assert compatibility.compatible is False
  129. assert compatibility.connector_version == "1.0.0"
  130. assert compatibility.detail == "optional driver unavailable"