Files
govoplan-connectors/tests/test_governed_runtime.py

340 lines
12 KiB
Python

from __future__ import annotations
from types import SimpleNamespace
import unittest
from unittest.mock import patch
from sqlalchemy import create_engine
from sqlalchemy.orm import sessionmaker
from govoplan_core.auth import ApiPrincipal
from govoplan_core.core.access import PrincipalRef
from govoplan_core.db.base import Base
from govoplan_connectors.backend.db.models import (
ConnectorConfiguration,
ConnectorDefinition,
ConnectorDefinitionRevision,
ConnectorSimulationRun,
)
from govoplan_connectors.backend.governed_runtime import (
GovernedConnectorError,
create_configuration,
execute_run,
list_configurations,
review_run,
update_configuration,
upsert_definition,
)
from govoplan_connectors.backend.governed_schemas import (
ConnectorConfigurationCreateRequest,
ConnectorConfigurationUpdateRequest,
ConnectorDefinitionUpsertRequest,
ConnectorReviewRequest,
ConnectorRunRequest,
)
def principal(tenant_id: str = "tenant-1") -> ApiPrincipal:
return ApiPrincipal(
principal=PrincipalRef(
account_id="account-1",
membership_id="membership-1",
tenant_id=tenant_id,
scopes=frozenset(
{
"connectors:source:read",
"connectors:source:write",
"connectors:source:admin",
}
),
),
account=SimpleNamespace(id="account-1"),
user=SimpleNamespace(id="user-1"),
)
def definition_payload(
*,
mapping_version: str = "1",
package_ref: str = "municipal-addresses@1",
timeout_seconds: int = 10,
) -> ConnectorDefinitionUpsertRequest:
return ConnectorDefinitionUpsertRequest.model_validate(
{
"definition_key": "municipal.addresses",
"name": "Municipal addresses",
"description": "A package-managed reference connector.",
"origin": "package",
"package_ref": package_ref,
"specification": {
"provider": "municipal-directory",
"protocol": "rest",
"capabilities": ["discover", "read", "dry_run"],
"input_schema": {"type": "object"},
"output_schema": {"type": "object"},
"mapping": {
"version": mapping_version,
"rules": [
{
"source": "external_id",
"target": "address.external_id",
"required": True,
},
{
"source": "street",
"target": "address.street",
"required": True,
},
],
},
"validation_rules": [
{
"kind": "unique",
"field": "address.external_id",
"severity": "error",
"code": "addresses.external_id.ambiguous",
"message": "The external identifier is not unique.",
}
],
"dry_run": {
"supported": True,
"simulation_supported": True,
"max_items": 50,
"redacted_fields": ["address.street"],
"sample_rows": [
{"external_id": "A-1", "street": "Sample street"}
],
},
"audit": {
"event_prefix": "connectors.municipal_addresses",
"expected_events": ["simulation.completed"],
"evidence_fields": ["input_hash", "configuration_hash"],
},
"privacy_classification": "confidential",
"retention_class": "connector-preview-30d",
"operational_limits": {"timeout_seconds": timeout_seconds},
"retry_policy": {"max_attempts": 2},
},
}
)
class GovernedConnectorRuntimeTests(unittest.TestCase):
def setUp(self) -> None:
self.engine = create_engine("sqlite:///:memory:")
self.tables = [
ConnectorDefinition.__table__,
ConnectorDefinitionRevision.__table__,
ConnectorConfiguration.__table__,
ConnectorSimulationRun.__table__,
]
Base.metadata.create_all(self.engine, tables=self.tables)
self.Session = sessionmaker(bind=self.engine)
self.session = self.Session()
self.audit = patch(
"govoplan_connectors.backend.governed_runtime.audit_from_principal"
)
self.audit_mock = self.audit.start()
def tearDown(self) -> None:
self.audit.stop()
self.session.close()
Base.metadata.drop_all(self.engine, tables=reversed(self.tables))
self.engine.dispose()
def _configuration(
self,
*,
ambiguity_policy: str = "manual_review",
):
definition = upsert_definition(
self.session,
principal(),
definition_payload(),
)
return create_configuration(
self.session,
principal(),
ConnectorConfigurationCreateRequest(
definition_id=definition.id,
name=f"Address import {ambiguity_policy}",
endpoint_url="https://directory.example.invalid/v1",
credential_ref="vault://connectors/address-reader",
local_overrides={"retry_policy": {"max_attempts": 5}},
ambiguity_policy=ambiguity_policy,
status="active",
),
)
def test_package_update_is_explicit_and_preserves_local_overrides(self) -> None:
configuration = self._configuration()
updated_definition = upsert_definition(
self.session,
principal(),
definition_payload(
mapping_version="2",
package_ref="municipal-addresses@2",
timeout_seconds=20,
),
)
unchanged = next(
item
for item in list_configurations(self.session, tenant_id="tenant-1")
if item.id == configuration.id
)
self.assertEqual(2, updated_definition.current_revision)
self.assertTrue(unchanged.update_available)
self.assertEqual("1", unchanged.effective_configuration["mapping"]["version"])
self.assertEqual(5, unchanged.effective_configuration["retry_policy"]["max_attempts"])
self.assertEqual(["retry_policy.max_attempts"], unchanged.protected_paths)
adopted = update_configuration(
self.session,
principal(),
configuration_id=configuration.id,
payload=ConnectorConfigurationUpdateRequest(
expected_revision=configuration.resource_revision,
adopt_latest_definition=True,
),
)
self.assertFalse(adopted.update_available)
self.assertEqual("2", adopted.effective_configuration["mapping"]["version"])
self.assertEqual(
20,
adopted.effective_configuration["operational_limits"]["timeout_seconds"],
)
self.assertEqual(5, adopted.effective_configuration["retry_policy"]["max_attempts"])
self.assertEqual(["retry_policy.max_attempts"], adopted.protected_paths)
def test_ambiguous_simulation_requires_review_and_is_idempotent(self) -> None:
configuration = self._configuration()
payload = ConnectorRunRequest(
idempotency_key="simulation-1",
external_revision="directory-etag-22",
input_rows=[
{"external_id": "duplicate", "street": "First"},
{"external_id": "duplicate", "street": "Second"},
],
)
created = execute_run(
self.session,
principal(),
configuration_id=configuration.id,
mode="simulation",
payload=payload,
)
replayed = execute_run(
self.session,
principal(),
configuration_id=configuration.id,
mode="simulation",
payload=payload,
)
self.assertEqual(created.id, replayed.id)
self.assertEqual("manual_review", created.status)
self.assertEqual("pending", created.review_state)
self.assertEqual(2, created.summary["ambiguous"])
self.assertEqual("<redacted>", created.effects[0]["sample"]["address"]["street"])
self.assertEqual("directory-etag-22", created.provenance["external_revision"])
reviewed = review_run(
self.session,
principal(),
run_id=created.id,
payload=ConnectorReviewRequest(
decision="approved",
reason="The duplicate rows represent an approved upstream alias.",
),
)
self.assertEqual("approved", reviewed.review_state)
self.assertEqual("review_approved", reviewed.status)
with self.assertRaisesRegex(GovernedConnectorError, "different run inputs"):
execute_run(
self.session,
principal(),
configuration_id=configuration.id,
mode="simulation",
payload=ConnectorRunRequest(
idempotency_key="simulation-1",
input_rows=[{"external_id": "other", "street": "Other"}],
),
)
def test_ambiguity_policy_can_quarantine_or_reject(self) -> None:
for policy, expected_status, expected_review in (
("quarantine", "quarantined", "quarantined"),
("reject", "rejected", "not_required"),
):
configuration = self._configuration(ambiguity_policy=policy)
result = execute_run(
self.session,
principal(),
configuration_id=configuration.id,
mode="dry_run",
payload=ConnectorRunRequest(
idempotency_key=f"{policy}-1",
input_rows=[
{"external_id": "same", "street": "First"},
{"external_id": "same", "street": "Second"},
],
),
)
self.assertEqual(expected_status, result.status)
self.assertEqual(expected_review, result.review_state)
def test_endpoint_credentials_and_stale_saves_are_rejected(self) -> None:
definition = upsert_definition(
self.session,
principal(),
definition_payload(),
)
with self.assertRaisesRegex(GovernedConnectorError, "must not contain credentials"):
create_configuration(
self.session,
principal(),
ConnectorConfigurationCreateRequest(
definition_id=definition.id,
name="Unsafe",
endpoint_url="https://user:secret@example.invalid/v1",
),
)
configuration = self._configuration()
with self.assertRaisesRegex(GovernedConnectorError, "reload it before saving"):
update_configuration(
self.session,
principal(),
configuration_id=configuration.id,
payload=ConnectorConfigurationUpdateRequest(
expected_revision=configuration.resource_revision + 1,
name="Stale",
),
)
def test_definition_ownership_cannot_change_implicitly(self) -> None:
package_definition = upsert_definition(
self.session,
principal(),
definition_payload(),
)
local_payload = definition_payload().model_copy(
update={"origin": "local", "package_ref": None}
)
with self.assertRaisesRegex(
GovernedConnectorError,
"configuration overrides",
):
upsert_definition(self.session, principal(), local_payload)
self.assertFalse(package_definition.local_definition)
if __name__ == "__main__":
unittest.main()