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.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_tickets.backend.db.models import Ticket from govoplan_tickets.backend.service import ( ADMIN_SCOPE, ASSIGN_SCOPE, READ_SCOPE, RESOLVE_SCOPE, TRIAGE_SCOPE, can_read_ticket, ) PROVIDER_ID = "tickets.tickets" RESOURCE_TYPE = "ticket" class TicketsSearchSource: def resource_types(self) -> Sequence[SearchResourceType]: return ( SearchResourceType( provider_id=PROVIDER_ID, module_id="tickets", resource_type=RESOURCE_TYPE, label="Tickets", requires_authorization_recheck=True, ), ) def backfill(self, session: object, *, request: SearchBackfillRequest) -> SearchBackfillPage: _assert_source(request.provider_id, request.resource_type) db = _session(session) query = select(Ticket).where( Ticket.tenant_id == request.tenant_id, Ticket.deleted_at.is_(None), ) if request.cursor: query = query.where(Ticket.id > request.cursor) rows = tuple(db.scalars(query.order_by(Ticket.id.asc()).limit(request.limit + 1))) has_more = len(rows) > request.limit selected = rows[: request.limit] high_watermark = db.scalar( select(func.max(Ticket.updated_at)).where( Ticket.tenant_id == request.tenant_id, Ticket.deleted_at.is_(None), ) ) return SearchBackfillPage( documents=tuple(_document(row) 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 else None, ) def authorize( self, session: object, principal: object, *, requests: Sequence[SearchAuthorizationRequest], ) -> Mapping[str, bool]: decisions = {item.reference.key: False for item in requests} db = _session(session) tenant_id = str(getattr(principal, "tenant_id", "") or "") for request in requests: reference = request.reference if ( reference.tenant_id != tenant_id or reference.module_id != "tickets" or reference.resource_type != RESOURCE_TYPE ): continue decisions[reference.key] = can_read_ticket( db, principal, ticket_id=reference.resource_id, ) return decisions def index_changes_for_event( self, session: object, *, event: PlatformEvent, delivery_key: str, ) -> Sequence[SearchIndexChange]: if ( event.module_id != "tickets" 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.scalar( select(Ticket).where( Ticket.tenant_id == event.tenant.id, Ticket.id == event.resource.id, ) ) deleted = row is None or row.deleted_at is not None cursor = event.event_id document = None if deleted else _document(row, change_cursor=cursor) reference = SearchResourceReference( tenant_id=event.tenant.id, module_id="tickets", 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 else cursor, cursor=cursor, document=document, occurred_at=event.occurred_at, ), ) def create_tickets_search_source(_context: ModuleContext) -> TicketsSearchSource: return TicketsSearchSource() def _document(row: Ticket, *, change_cursor: str | None = None) -> SearchDocument: tokens = [ f"scope:{READ_SCOPE}", f"scope:{TRIAGE_SCOPE}", f"scope:{ASSIGN_SCOPE}", f"scope:{RESOLVE_SCOPE}", f"scope:{ADMIN_SCOPE}", ] if row.created_by: tokens.append(f"account:{row.created_by}") for value in (row.assignee, row.reporter, row.requester, *(row.participants or [])): if not isinstance(value, Mapping): continue kind = str(value.get("kind") or "") subject_id = str(value.get("id") or "") prefix = "function" if kind == "function_assignment" else kind if prefix and subject_id and prefix != "external": tokens.append(f"{prefix}:{subject_id}") return SearchDocument( tenant_id=row.tenant_id, module_id="tickets", provider_id=PROVIDER_ID, resource_type=RESOURCE_TYPE, resource_id=row.id, title=row.title, url=f"/tickets?ticketId={quote(row.id, safe='')}", summary=f"{row.ticket_number} · {row.status} · {row.priority}", body=row.search_text[:200_000], keywords=tuple( item[:200] for item in (row.ticket_number, row.ticket_type, row.status, row.priority, row.queue_ref or "") if item ), visibility=row.visibility, acl_tokens=tuple(dict.fromkeys(tokens)) if row.visibility == "restricted" else (), metadata={ "ticket_number": row.ticket_number, "ticket_type": row.ticket_type, "status": row.status, "priority": row.priority, "queue_ref": row.queue_ref, "service_target_at": row.service_target_at.isoformat() if row.service_target_at else None, }, source_revision=str(row.revision), change_cursor=change_cursor, source_updated_at=row.updated_at or row.recorded_at, requires_authorization_recheck=True, ) def _assert_source(provider_id: str, resource_type: str) -> None: if provider_id != PROVIDER_ID or resource_type != RESOURCE_TYPE: raise ValueError("Unsupported Tickets search source.") def _session(value: object) -> Session: if not isinstance(value, Session): raise TypeError("Tickets search requires a SQLAlchemy session.") return value __all__ = [ "PROVIDER_ID", "RESOURCE_TYPE", "TicketsSearchSource", "create_tickets_search_source", ]