Add native permission-aware campaign search source

This commit is contained in:
2026-08-04 03:03:26 +02:00
parent 8f5231147d
commit 91890fdaf5
3 changed files with 406 additions and 0 deletions
+32
View File
@@ -55,12 +55,14 @@ from govoplan_core.core.postbox import (
CAPABILITY_POSTBOX_EVIDENCE, CAPABILITY_POSTBOX_EVIDENCE,
) )
from govoplan_core.core.references import CAPABILITY_ACCESS_REFERENCE_OPTIONS from govoplan_core.core.references import CAPABILITY_ACCESS_REFERENCE_OPTIONS
from govoplan_core.core.search import SearchSourceProviderRegistration
from govoplan_campaign.backend.change_tracking import register_campaign_change_tracking from govoplan_campaign.backend.change_tracking import register_campaign_change_tracking
from govoplan_campaign.backend.db import models as campaign_models # noqa: F401 - populate Campaign ORM metadata from govoplan_campaign.backend.db import models as campaign_models # noqa: F401 - populate Campaign ORM metadata
from govoplan_campaign.backend.documentation import ( from govoplan_campaign.backend.documentation import (
CAMPAIGN_USER_DOCUMENTATION, CAMPAIGN_USER_DOCUMENTATION,
documentation_topics, documentation_topics,
) )
from govoplan_campaign.backend.search_source import create_campaign_search_source
register_campaign_change_tracking() register_campaign_change_tracking()
@@ -378,6 +380,7 @@ manifest = ModuleManifest(
"postbox", "postbox",
"approvals", "approvals",
"reporting", "reporting",
"search",
), ),
provides_interfaces=( provides_interfaces=(
ModuleInterfaceProvider(name="campaigns.access", version="0.1.6"), ModuleInterfaceProvider(name="campaigns.access", version="0.1.6"),
@@ -481,12 +484,24 @@ manifest = ModuleManifest(
version_max_exclusive="0.2.0", version_max_exclusive="0.2.0",
optional=True, optional=True,
), ),
ModuleInterfaceRequirement(
name="search.source",
version_min="1.0.0",
version_max_exclusive="2.0.0",
optional=True,
),
), ),
permissions=PERMISSIONS, permissions=PERMISSIONS,
route_factory=_campaigns_router, route_factory=_campaigns_router,
role_templates=ROLE_TEMPLATES, role_templates=ROLE_TEMPLATES,
tenant_summary_providers=(_tenant_summary,), tenant_summary_providers=(_tenant_summary,),
tenant_summary_batch_providers=(_tenant_summary_batch,), tenant_summary_batch_providers=(_tenant_summary_batch,),
search_sources=(
SearchSourceProviderRegistration(
id="campaigns.campaigns",
factory=create_campaign_search_source,
),
),
nav_items=( nav_items=(
NavItem( NavItem(
path="/campaigns", path="/campaigns",
@@ -592,6 +607,23 @@ manifest = ModuleManifest(
), ),
documentation=( documentation=(
*CAMPAIGN_USER_DOCUMENTATION, *CAMPAIGN_USER_DOCUMENTATION,
DocumentationTopic(
id="campaigns.search.campaigns",
title="Search authorized campaigns",
summary="Expose campaign identity and lifecycle metadata to permission-aware platform Search.",
body=(
"When Search is installed, Campaign contributes current campaign names, external identifiers, "
"descriptions, and lifecycle state. Search rechecks tenant ownership, group ownership, explicit "
"shares, revocation, deletion, and the Campaign read permission before returning a result. "
"Committed Campaign and share changes update the derived index through the durable platform event "
"path; rebuilding Search never changes Campaign evidence."
),
layer="configured",
documentation_types=("admin", "user"),
audience=("campaign_manager", "campaign_operator", "administrator"),
related_modules=("search",),
order=44,
),
DocumentationTopic( DocumentationTopic(
id="campaigns.postbox-delivery", id="campaigns.postbox-delivery",
title="Deliver Campaign messages to Postboxes", title="Deliver Campaign messages to Postboxes",
@@ -0,0 +1,251 @@
from __future__ import annotations
from collections.abc import Mapping, Sequence
from urllib.parse import quote
from sqlalchemy import func, select
from sqlalchemy.orm import Session
from govoplan_core.auth import ApiPrincipal
from govoplan_core.core.events import PlatformEvent
from govoplan_core.core.modules import ModuleContext
from govoplan_core.core.search import (
SearchAuthorizationRequest,
SearchBackfillPage,
SearchBackfillRequest,
SearchDocument,
SearchIndexChange,
SearchResourceReference,
SearchResourceType,
)
from govoplan_campaign.backend.capabilities import CampaignAccessService
from govoplan_campaign.backend.db.models import Campaign, CampaignShare
PROVIDER_ID = "campaigns.campaigns"
RESOURCE_TYPE = "campaign"
READ_SCOPE = "campaigns:campaign:read"
class CampaignSearchSource:
def resource_types(self) -> Sequence[SearchResourceType]:
return (
SearchResourceType(
provider_id=PROVIDER_ID,
module_id="campaigns",
resource_type=RESOURCE_TYPE,
label="Campaigns",
requires_authorization_recheck=True,
),
)
def backfill(
self,
session: object,
*,
request: SearchBackfillRequest,
) -> SearchBackfillPage:
_assert_source(request.provider_id, request.resource_type)
db = _session(session)
statement = select(Campaign).where(
Campaign.tenant_id == request.tenant_id,
Campaign.status != "deleted",
)
if request.cursor:
statement = statement.where(Campaign.id > request.cursor)
rows = list(
db.scalars(
statement.order_by(Campaign.id).limit(request.limit + 1)
)
)
has_more = len(rows) > request.limit
selected = rows[: request.limit]
shares = _shares_by_campaign(db, selected)
high_watermark = db.scalar(
select(func.max(Campaign.updated_at)).where(
Campaign.tenant_id == request.tenant_id,
Campaign.status != "deleted",
)
)
return SearchBackfillPage(
documents=tuple(
_document(row, shares=shares.get(row.id, ()))
for row in selected
),
next_cursor=selected[-1].id if has_more and selected else None,
complete=not has_more,
high_watermark=(
high_watermark.isoformat()
if high_watermark is not None
else None
),
)
def authorize(
self,
session: object,
principal: object,
*,
requests: Sequence[SearchAuthorizationRequest],
) -> Mapping[str, bool]:
decisions = {item.reference.key: False for item in requests}
if not isinstance(principal, ApiPrincipal) or not principal.has(READ_SCOPE):
return decisions
db = _session(session)
access = CampaignAccessService()
user_id = str(getattr(principal.user, "id", "") or principal.membership_id or "")
for request in requests:
reference = request.reference
if (
reference.tenant_id != principal.tenant_id
or reference.module_id != "campaigns"
or reference.resource_type != RESOURCE_TYPE
):
continue
decisions[reference.key] = access.can_read_campaign(
db,
tenant_id=principal.tenant_id,
campaign_id=reference.resource_id,
user_id=user_id,
group_ids=principal.group_ids,
tenant_admin=principal.has("tenant:*"),
)
return decisions
def index_changes_for_event(
self,
session: object,
*,
event: PlatformEvent,
delivery_key: str,
) -> Sequence[SearchIndexChange]:
if (
event.module_id != "campaigns"
or event.tenant is None
or event.resource is None
or event.resource.type != RESOURCE_TYPE
or event.resource.id is None
):
return ()
db = _session(session)
row = db.get(Campaign, event.resource.id)
deleted = row is None or row.tenant_id != event.tenant.id or row.status == "deleted"
cursor = event.event_id
document = None
if not deleted:
document = _document(
row,
shares=tuple(
db.scalars(
select(CampaignShare).where(
CampaignShare.campaign_id == row.id,
CampaignShare.revoked_at.is_(None),
)
)
),
change_cursor=cursor,
)
reference = SearchResourceReference(
tenant_id=event.tenant.id,
module_id="campaigns",
resource_type=RESOURCE_TYPE,
resource_id=event.resource.id,
)
return (
SearchIndexChange(
change_id=f"{delivery_key}:{PROVIDER_ID}",
provider_id=PROVIDER_ID,
kind="delete" if deleted else "upsert",
reference=reference,
source_revision=(
document.source_revision if document is not None else cursor
),
cursor=cursor,
document=document,
occurred_at=event.occurred_at,
),
)
def create_campaign_search_source(_context: ModuleContext) -> CampaignSearchSource:
return CampaignSearchSource()
def _document(
row: Campaign,
*,
shares: Sequence[CampaignShare],
change_cursor: str | None = None,
) -> SearchDocument:
tokens = [f"scope:{READ_SCOPE}"]
if row.owner_user_id:
tokens.append(f"membership:{row.owner_user_id}")
if row.owner_group_id:
tokens.append(f"group:{row.owner_group_id}")
for share in shares:
prefix = "membership" if share.target_type == "user" else share.target_type
if prefix in {"membership", "group"}:
tokens.append(f"{prefix}:{share.target_id}")
updated_at = row.updated_at or row.created_at
return SearchDocument(
tenant_id=row.tenant_id,
module_id="campaigns",
provider_id=PROVIDER_ID,
resource_type=RESOURCE_TYPE,
resource_id=row.id,
title=row.name,
url=f"/campaigns/{quote(row.id, safe='')}",
summary=row.description[:4000] if row.description else None,
body=" ".join(
value for value in (row.external_id, row.description) if value
)[:200_000],
keywords=(row.external_id[:200], row.status[:200]),
visibility="restricted",
acl_tokens=tuple(dict.fromkeys(tokens)),
metadata={
"external_id": row.external_id,
"status": row.status,
"current_version_id": row.current_version_id,
},
source_revision=f"{row.current_version_id or 'none'}:{updated_at.isoformat()}",
change_cursor=change_cursor,
source_updated_at=updated_at,
requires_authorization_recheck=True,
)
def _shares_by_campaign(
session: Session,
rows: Sequence[Campaign],
) -> dict[str, tuple[CampaignShare, ...]]:
ids = [row.id for row in rows]
grouped: dict[str, list[CampaignShare]] = {item: [] for item in ids}
if not ids:
return {}
for share in session.scalars(
select(CampaignShare).where(
CampaignShare.campaign_id.in_(ids),
CampaignShare.revoked_at.is_(None),
)
):
grouped.setdefault(share.campaign_id, []).append(share)
return {key: tuple(value) for key, value in grouped.items()}
def _assert_source(provider_id: str, resource_type: str) -> None:
if provider_id != PROVIDER_ID or resource_type != RESOURCE_TYPE:
raise ValueError("Unsupported Campaign search source.")
def _session(value: object) -> Session:
if not isinstance(value, Session):
raise TypeError("Campaign search requires a SQLAlchemy session.")
return value
__all__ = [
"CampaignSearchSource",
"PROVIDER_ID",
"RESOURCE_TYPE",
"create_campaign_search_source",
]
+123
View File
@@ -0,0 +1,123 @@
from __future__ import annotations
from types import SimpleNamespace
import unittest
from sqlalchemy import create_engine
from sqlalchemy.orm import Session
from govoplan_access.backend.db.models import Account, Group, User
from govoplan_campaign.backend.db.models import Campaign, CampaignShare
from govoplan_campaign.backend.search_source import (
CampaignSearchSource,
PROVIDER_ID,
RESOURCE_TYPE,
)
from govoplan_core.auth import ApiPrincipal
from govoplan_core.core.access import PrincipalRef
from govoplan_core.core.search import (
SearchAuthorizationRequest,
SearchBackfillRequest,
SearchResourceReference,
)
from govoplan_core.db.base import Base
class CampaignSearchSourceTests(unittest.TestCase):
def setUp(self) -> None:
self.engine = create_engine("sqlite://")
Base.metadata.create_all(
self.engine,
tables=(
Account.__table__,
User.__table__,
Group.__table__,
Campaign.__table__,
CampaignShare.__table__,
),
)
self.session = Session(self.engine)
self.session.add_all(
(
Account(
id="account-1",
email="one@example.test",
normalized_email="one@example.test",
),
User(
id="user-1",
tenant_id="tenant-1",
account_id="account-1",
email="one@example.test",
),
Campaign(
id="campaign-1",
tenant_id="tenant-1",
owner_user_id="user-1",
external_id="monthly-letters",
name="Monthly letters",
),
Campaign(
id="campaign-other",
tenant_id="tenant-2",
external_id="other",
name="Other tenant",
),
)
)
self.session.commit()
self.source = CampaignSearchSource()
def tearDown(self) -> None:
self.session.close()
self.engine.dispose()
def test_backfill_and_live_acl_recheck_do_not_cross_tenants(self) -> None:
page = self.source.backfill(
self.session,
request=SearchBackfillRequest(
tenant_id="tenant-1",
provider_id=PROVIDER_ID,
resource_type=RESOURCE_TYPE,
rebuild_id="rebuild-1",
),
)
self.assertEqual(("campaign-1",), tuple(doc.resource_id for doc in page.documents))
reference = SearchResourceReference(
tenant_id="tenant-1",
module_id="campaigns",
resource_type=RESOURCE_TYPE,
resource_id="campaign-1",
)
request = SearchAuthorizationRequest(reference=reference, source_revision="1")
self.assertTrue(
self.source.authorize(
self.session,
_principal({"campaigns:campaign:read"}),
requests=(request,),
)[reference.key]
)
self.assertFalse(
self.source.authorize(
self.session,
_principal(set()),
requests=(request,),
)[reference.key]
)
def _principal(scopes: set[str]) -> ApiPrincipal:
return ApiPrincipal(
principal=PrincipalRef(
account_id="account-1",
membership_id="user-1",
tenant_id="tenant-1",
scopes=frozenset(scopes),
),
account=SimpleNamespace(id="account-1"),
user=SimpleNamespace(id="user-1"),
)
if __name__ == "__main__":
unittest.main()