from __future__ import annotations from types import SimpleNamespace import unittest from unittest.mock import patch from sqlalchemy import create_engine, select from sqlalchemy.orm import Session from govoplan_core.auth import ApiPrincipal from govoplan_core.core.access import PrincipalRef from govoplan_core.core.search import ( SearchAuthorizationRequest, SearchBackfillRequest, ) from govoplan_core.core.recovery import RecoveryOperation, RecoveryStatus from govoplan_core.core.runtime_coordination import ( RuntimeIdentity, bind_process_runtime_identity, ) from govoplan_core.db.base import Base from govoplan_connectors.backend.db.models import ( ConnectorConfiguration, ConnectorDefinition, ConnectorDefinitionRevision, ConnectorKnowledgeObject, ConnectorKnowledgeProfile, ConnectorKnowledgeSyncRun, ) from govoplan_connectors.backend.knowledge_connector import ( KNOWLEDGE_PROVIDER_ID, KNOWLEDGE_RESOURCE_TYPE, KnowledgeConnectorError, create_profile, discover_profile, list_objects, migration_dry_run, publish_page, synchronize_profile, update_profile, ) from govoplan_connectors.backend.knowledge_schemas import ( KnowledgeMigrationDryRunRequest, KnowledgeMigrationTargetState, KnowledgeNamespaceMapping, KnowledgeProfileCreateRequest, KnowledgeProfileUpdateRequest, KnowledgePublishRequest, KnowledgeSyncRequest, ) from govoplan_connectors.backend.knowledge_search import ( ExternalKnowledgeSearchSource, ) from govoplan_connectors.backend.mediawiki_transport import ( MediaWikiChangeBatch, MediaWikiPublishResult, MediaWikiTransportError, ) ALL_SCOPES = frozenset( { "connectors:knowledge:read", "connectors:knowledge:admin", "connectors:knowledge:sync", "connectors:knowledge:publish", "connectors:knowledge:migrate", } ) def principal( tenant_id: str = "tenant-1", *, scopes: frozenset[str] = ALL_SCOPES, groups: frozenset[str] = frozenset({"editors"}), ) -> ApiPrincipal: return ApiPrincipal( principal=PrincipalRef( account_id="account-1", membership_id="membership-1", tenant_id=tenant_id, scopes=scopes, group_ids=groups, ), account=SimpleNamespace(id="account-1"), user=SimpleNamespace(id="account-1"), ) class StaticTransport: def __init__(self) -> None: self.batches: list[MediaWikiChangeBatch] = [] self.change_calls = 0 self.publish_calls = 0 self.publish_error: MediaWikiTransportError | None = None def discover(self, *, endpoint_url, credential): del endpoint_url, credential return { "curtimestamp": "2026-08-22T10:00:00Z", "query": { "general": { "generator": "MediaWiki 1.43.1", "phpversion": "8.3.8", }, "extensions": [ {"name": "BlueSpiceFoundation", "version": "4.5.2"}, {"name": "BlueSpicePermissionManager", "version": "4.5.2"}, ], "namespaces": { "0": {"id": 0, "name": "", "content": True}, "4": {"id": 4, "name": "GovWiki", "content": True}, }, "userinfo": { "id": 17, "name": "govoplan", "rights": ["read", "edit"], "groups": ["bot"], }, }, } def changes(self, **kwargs): del kwargs self.change_calls += 1 if not self.batches: raise AssertionError("No deterministic change batch remains") return self.batches.pop(0) def publish(self, **kwargs): self.publish_calls += 1 if self.publish_error is not None: raise self.publish_error return MediaWikiPublishResult( page_id=str(kwargs.get("expected_page_id") or "99"), revision_id="901", title=str(kwargs["title"]), canonical_url="https://wiki.example.invalid/wiki/Published_Guide", evidence={"result": "Success", "fixture": True}, ) def page( *, page_id: str = "42", revision_id: str = "501", title: str = "Citizen Guide", acl_tokens: list[str] | None = None, ) -> dict[str, object]: return { "change_kind": "upsert", "change_cursor": f"rcid:{revision_id}", "pageid": page_id, "ns": 0, "title": title, "fullurl": f"https://wiki.example.invalid/wiki/{title.replace(' ', '_')}", "lastrevid": revision_id, "revisions": [ { "revid": revision_id, "timestamp": "2026-08-22T10:00:00Z", "user": "Ada Admin", "comment": "Reviewed guidance", "contentmodel": "wikitext", "sha1": f"sha-{revision_id}", "slots": { "main": { "content": "Welcome {{UnsupportedBox|important}} [[Services]]" } }, } ], "categories": [{"title": "Category:Citizen service"}], "links": [{"ns": 0, "title": "Services"}], "images": [{"ns": 6, "title": "File:guide.pdf", "pageid": 71}], "discussions": [ { "id": "discussion-1", "author": "Ada Admin", "body": "Please verify this section.", "created_at": "2026-08-20T09:00:00Z", } ], "permissions": { "visibility": "restricted", "acl_tokens": acl_tokens or ["group:editors"], }, } class MediaWikiConnectorTests(unittest.TestCase): def setUp(self) -> None: bind_process_runtime_identity( RuntimeIdentity( installation_id="mediawiki-connector-tests", node_id="node-1", incarnation="incarnation-1", role="worker", software_version="test", composition_hash="a" * 64, ) ) self.engine = create_engine("sqlite+pysqlite:///:memory:") Base.metadata.create_all(self.engine) self.session = Session(self.engine) self.audit = patch( "govoplan_connectors.backend.knowledge_connector.audit_event" ) self.audit.start() self.credential = patch( "govoplan_connectors.backend.knowledge_connector._credential", return_value={"access_token": "fixture-token"}, ) self.credential.start() self.transport = StaticTransport() self.configuration_id = self._seed_configuration() self.profile_id = self._profile() def tearDown(self) -> None: bind_process_runtime_identity(None) self.credential.stop() self.audit.stop() self.session.close() self.engine.dispose() def _seed_configuration(self) -> str: definition = ConnectorDefinition( id="definition-1", tenant_id="tenant-1", definition_key="knowledge.mediawiki", name="MediaWiki", description="Knowledge transport", status="active", current_revision=1, local_definition=True, ) self.session.add(definition) self.session.add( ConnectorDefinitionRevision( id="definition-revision-1", definition_id=definition.id, revision=1, specification={ "provider": "mediawiki", "protocol": "mediawiki_action_api", }, definition_hash="definition-hash", origin="local", created_by="account-1", ) ) configuration = ConnectorConfiguration( id="configuration-1", tenant_id="tenant-1", definition_id=definition.id, name="Institutional knowledge", status="active", endpoint_url="https://wiki.example.invalid", credential_ref="credential-envelope-1", base_definition_revision=1, local_overrides={}, protected_paths=[], effective_configuration={ "provider": "mediawiki", "protocol": "mediawiki_action_api", }, effective_hash="configuration-hash", resource_revision=1, ambiguity_policy="manual_review", updated_by="account-1", ) self.session.add(configuration) self.session.flush() return configuration.id def _profile(self) -> str: created = create_profile( self.session, principal(), KnowledgeProfileCreateRequest( configuration_id=self.configuration_id, desired_maturity="migrate", source_authority_mode="external_mirror", default_visibility="restricted", default_acl_tokens=["group:knowledge-managers"], namespace_mappings=[ KnowledgeNamespaceMapping( source_namespace_id=0, source_name="", target_space_ref="service-guidance", target_path_prefix="imported", ) ], ), ) return created.id def _discover(self): return discover_profile( self.session, principal(), profile_id=self.profile_id, transport=self.transport, ) def _sync(self, raw_pages, *, key="sync-1", force_full=True): self.transport.batches.append( MediaWikiChangeBatch( changes=tuple(raw_pages), next_cursor=None, complete=True, high_watermark="2026-08-22T10:05:00Z", evidence={"fixture": True}, ) ) return synchronize_profile( self.session, principal(), profile_id=self.profile_id, payload=KnowledgeSyncRequest( idempotency_key=key, force_full=force_full, ), transport=self.transport, registry=None, ) def test_discovery_mapping_idempotency_and_migration_loss_diagnostics(self) -> None: discovery = self._discover() self.assertEqual("bluespice", discovery.product) self.assertEqual("4.5.2", discovery.product_version) self.assertEqual("migrate", discovery.maturity) self.assertIn("publish", discovery.capabilities) self.assertIn("permission_metadata", discovery.capabilities) run = self._sync([page()]) self.assertEqual({"create": 1}, run.counts) self.assertIn( "attachments_reference_only", {item.code for item in run.diagnostics} ) objects, _cursor = list_objects( self.session, principal(), profile_id=self.profile_id ) self.assertEqual(1, len(objects)) item = objects[0] self.assertEqual("42", item.external_reference.object_id) self.assertEqual("501", item.external_reference.version) self.assertEqual("imported/Citizen-Guide", item.mapped_data["target_path"]) self.assertEqual("Ada Admin", item.mapped_data["revision_author"]) self.assertEqual( "user", item.mapped_data["revision_author_reference"]["object_type"] ) self.assertEqual("guide.pdf", item.mapped_data["files"][0]["name"]) self.assertEqual("discussion-1", item.mapped_data["discussions"][0]["external_id"]) replay = synchronize_profile( self.session, principal(), profile_id=self.profile_id, payload=KnowledgeSyncRequest( idempotency_key="sync-1", force_full=True, ), transport=self.transport, registry=None, ) self.assertEqual(run.id, replay.id) self.assertEqual(1, self.transport.change_calls) with self.assertRaisesRegex(KnowledgeConnectorError, "different request"): synchronize_profile( self.session, principal(), profile_id=self.profile_id, payload=KnowledgeSyncRequest( idempotency_key="sync-1", force_full=True, limit=25, ), transport=self.transport, registry=None, ) preview = migration_dry_run( self.session, principal(), profile_id=self.profile_id, payload=KnowledgeMigrationDryRunRequest( idempotency_key="migration-1", target_space_ref="service-guidance", supported_macros=["SupportedBox"], existing_targets=[ KnowledgeMigrationTargetState( path="imported/Citizen-Guide", source_external_id="different-page", attachment_names=["guide.pdf"], ) ], ), ) self.assertFalse(preview.can_apply) self.assertEqual("2026-08-22T10:05:00Z", preview.source_revision) self.assertEqual( {"attachment_name_conflict", "target_path_conflict", "unsupported_macro"}, {item.code for item in preview.diagnostics}, ) def test_acl_changes_moves_deletes_and_search_authorization_are_current(self) -> None: self._discover() self._sync([page()]) source = ExternalKnowledgeSearchSource() backfill = source.backfill( self.session, request=SearchBackfillRequest( tenant_id="tenant-1", provider_id=KNOWLEDGE_PROVIDER_ID, resource_type=KNOWLEDGE_RESOURCE_TYPE, rebuild_id="rebuild-1", limit=100, ), ) self.assertEqual(1, len(backfill.documents)) document = backfill.documents[0] request = SearchAuthorizationRequest( reference=document.reference, source_revision=document.source_revision, ) self.assertTrue( source.authorize(self.session, principal(), requests=[request])[ document.reference.key ] ) denied = principal(groups=frozenset({"other"})) self.assertFalse( source.authorize(self.session, denied, requests=[request])[ document.reference.key ] ) self._sync( [ page( revision_id="502", title="Resident Guide", acl_tokens=["group:reviewers"], ) ], key="sync-2", force_full=False, ) self.assertFalse( source.authorize(self.session, principal(), requests=[request])[ document.reference.key ] ) reviewers = principal(groups=frozenset({"reviewers"})) self.assertTrue( source.authorize(self.session, reviewers, requests=[request])[ document.reference.key ] ) stored = self.session.scalar( select(ConnectorKnowledgeObject).where( ConnectorKnowledgeObject.profile_id == self.profile_id, ConnectorKnowledgeObject.external_id == "42", ) ) self.assertEqual("Resident Guide", stored.title) self.assertEqual("imported/Resident-Guide", stored.mapped_data["target_path"]) self._sync( [ { "change_kind": "delete", "change_cursor": "logid:700", "pageid": 0, "ns": 0, "title": "Resident Guide", "timestamp": "2026-08-22T11:00:00Z", "logid": 700, } ], key="sync-3", force_full=False, ) self.assertEqual("deleted", stored.status) self.assertFalse( source.authorize(self.session, reviewers, requests=[request])[ document.reference.key ] ) after_delete = source.backfill( self.session, request=SearchBackfillRequest( tenant_id="tenant-1", provider_id=KNOWLEDGE_PROVIDER_ID, resource_type=KNOWLEDGE_RESOURCE_TYPE, rebuild_id="rebuild-2", limit=100, ), ) self.assertEqual((), after_delete.documents) def test_publication_replay_and_outcome_unknown_are_evidenced(self) -> None: self._discover() payload = KnowledgePublishRequest( idempotency_key="publish-1", title="Published Guide", body="Reviewed body", summary="Publish approved guidance", expected_external_revision="900", ) result = publish_page( self.session, principal(), profile_id=self.profile_id, external_page_id="99", payload=payload, transport=self.transport, registry=None, durable_recovery=True, ) replay = publish_page( self.session, principal(), profile_id=self.profile_id, external_page_id="99", payload=payload, transport=self.transport, registry=None, durable_recovery=False, ) self.assertTrue(result.accepted) self.assertEqual(result.run.id, replay.run.id) self.assertEqual("901", result.external_reference.version) self.assertEqual(1, self.transport.publish_calls) recovery = self.session.scalar( select(RecoveryOperation).where( RecoveryOperation.resource_type == "external_knowledge_page", RecoveryOperation.resource_id == "99", ) ) self.assertEqual(RecoveryStatus.SUCCEEDED.value, recovery.status) self.transport.publish_error = MediaWikiTransportError( "transport_timeout", "Provider response timed out.", retryable=True, outcome_unknown=True, ) with self.assertRaisesRegex(KnowledgeConnectorError, "outcome is unknown"): publish_page( self.session, principal(), profile_id=self.profile_id, external_page_id="100", payload=KnowledgePublishRequest( idempotency_key="publish-unknown", title="Uncertain Guide", body="Body", ), transport=self.transport, registry=None, durable_recovery=False, ) unresolved = self.session.scalar( select(ConnectorKnowledgeSyncRun).where( ConnectorKnowledgeSyncRun.profile_id == self.profile_id, ConnectorKnowledgeSyncRun.idempotency_key == "publish-unknown", ) ) self.assertEqual("outcome_unknown", unresolved.status) def test_fallback_acl_changes_and_profile_pause_fail_closed_immediately(self) -> None: self._discover() fallback_page = page(page_id="43", title="Fallback Guide") fallback_page.pop("permissions") self._sync([fallback_page]) source = ExternalKnowledgeSearchSource() document = source.backfill( self.session, request=SearchBackfillRequest( tenant_id="tenant-1", provider_id=KNOWLEDGE_PROVIDER_ID, resource_type=KNOWLEDGE_RESOURCE_TYPE, rebuild_id="fallback-rebuild", ), ).documents[0] request = SearchAuthorizationRequest( reference=document.reference, source_revision=document.source_revision, ) managers = principal(groups=frozenset({"knowledge-managers"})) self.assertTrue( source.authorize(self.session, managers, requests=[request])[ document.reference.key ] ) profile = self.session.get(ConnectorKnowledgeProfile, self.profile_id) before_hash = self.session.get( ConnectorKnowledgeObject, document.resource_id ).content_hash updated = update_profile( self.session, principal(), profile_id=self.profile_id, payload=KnowledgeProfileUpdateRequest( expected_resource_revision=profile.resource_revision, default_acl_tokens=["group:reviewers"], ), registry=None, ) self.assertFalse( source.authorize(self.session, managers, requests=[request])[ document.reference.key ] ) reviewers = principal(groups=frozenset({"reviewers"})) self.assertTrue( source.authorize(self.session, reviewers, requests=[request])[ document.reference.key ] ) self.assertNotEqual( before_hash, self.session.get(ConnectorKnowledgeObject, document.resource_id).content_hash, ) update_profile( self.session, principal(), profile_id=self.profile_id, payload=KnowledgeProfileUpdateRequest( expected_resource_revision=updated.resource_revision, status="paused", ), registry=None, ) self.assertFalse( source.authorize(self.session, reviewers, requests=[request])[ document.reference.key ] ) def test_profiles_and_objects_are_tenant_isolated(self) -> None: self._discover() self._sync([page()]) with self.assertRaisesRegex(KnowledgeConnectorError, "not found"): list_objects( self.session, principal("tenant-2"), profile_id=self.profile_id, ) profile = self.session.get(ConnectorKnowledgeProfile, self.profile_id) self.assertEqual("tenant-1", profile.tenant_id) if __name__ == "__main__": unittest.main()