340 lines
12 KiB
Python
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()
|