Complete permission-aware native search indexing

This commit is contained in:
2026-08-04 03:05:03 +02:00
parent 912db10a8a
commit 1bb338117f
12 changed files with 680 additions and 26 deletions
+68
View File
@@ -6,6 +6,7 @@ from types import SimpleNamespace
from sqlalchemy import create_engine
from sqlalchemy.orm import Session
from govoplan_core.core.events import EventObjectRef, EventTenantRef, PlatformEvent
from govoplan_core.core.search import (
SearchBackfillPage,
SearchDocument,
@@ -45,6 +46,7 @@ class _Registry:
return (
(
SimpleNamespace(
module_id="cases",
registration=SimpleNamespace(
id="cases.records",
order=25,
@@ -86,6 +88,45 @@ class _Source:
}
class _EventSource(_Source):
def index_changes_for_event(self, session, *, event, delivery_key):
del session
if (
event.module_id != "cases"
or event.tenant is None
or event.resource is None
or event.resource.type != "case"
or event.resource.id is None
):
return ()
reference = _source_document(event.resource.id).reference
document = SearchDocument(
tenant_id=event.tenant.id,
module_id="cases",
resource_type="case",
resource_id=event.resource.id,
title=f"Permit {event.resource.id}",
url=f"/cases/{event.resource.id}",
acl_tokens=("account:account-1",),
provider_id="cases.records",
source_revision="event-1",
change_cursor=event.event_id,
requires_authorization_recheck=True,
)
return (
SearchIndexChange(
change_id=f"{delivery_key}:cases.records",
provider_id="cases.records",
kind="upsert",
reference=reference,
source_revision=document.source_revision,
cursor=event.event_id,
document=document,
occurred_at=event.occurred_at,
),
)
class _ResultProvider:
def search(self, session, principal, *, query):
del session, principal
@@ -382,6 +423,33 @@ class SearchServiceTests(unittest.TestCase):
self.session.query(SearchIndexDocument).count(),
)
def test_committed_event_ingestion_is_source_owned_and_idempotent(self) -> None:
service = SearchIndexService(_Registry(_EventSource()))
event = PlatformEvent(
type="cases.case.updated",
module_id="cases",
tenant=EventTenantRef(id="tenant-1"),
resource=EventObjectRef(type="case", id="case-1"),
)
first = service.ingest_event(
self.session,
event=event,
delivery_key="delivery-1",
)
second = service.ingest_event(
self.session,
event=event,
delivery_key="delivery-1",
)
self.assertEqual(1, first["queued"])
self.assertEqual(1, second["duplicates"])
self.assertEqual(1, service.process_changes(self.session)["applied"])
indexed = self.session.query(SearchIndexDocument).one()
self.assertEqual("tenant-1", indexed.tenant_id)
self.assertEqual("case-1", indexed.resource_id)
def test_rebuild_resumes_and_removes_stale_documents(self) -> None:
stale = _source_document("stale-case")
self.service.upsert_document(