Files
govoplan-tickets/src/govoplan_tickets/backend/service.py
T

1162 lines
39 KiB
Python

from __future__ import annotations
from collections.abc import Mapping, Sequence
from datetime import UTC, datetime
import hashlib
import json
from typing import Any
from sqlalchemy.orm import Session
from govoplan_core.core.events import (
EventActorRef,
EventObjectRef,
EventTenantRef,
PlatformEvent,
emit_platform_event,
)
from govoplan_core.core.tickets import (
TicketCaseEscalationCommand,
TicketRoutingRequest,
ticket_case_escalation_provider,
ticket_routing_provider,
)
from govoplan_core.security.module_permissions import scopes_grant_compatible
from govoplan_tickets.backend.db.models import (
Ticket,
TicketComment,
TicketEscalation,
TicketHistory,
)
from govoplan_tickets.backend.domain import (
TicketLink,
TicketRecord,
TicketSubjectRef,
validate_transition,
)
CAPABILITY_TICKETS_REGISTRY = "tickets.registry"
READ_SCOPE = "tickets:ticket:read"
REPORT_SCOPE = "tickets:ticket:report"
TRIAGE_SCOPE = "tickets:ticket:triage"
ASSIGN_SCOPE = "tickets:ticket:assign"
RESOLVE_SCOPE = "tickets:ticket:resolve"
ADMIN_SCOPE = "tickets:ticket:admin"
LEGACY_WRITE_SCOPE = "tickets:ticket:write"
class TicketStoreError(ValueError):
pass
class TicketNotFoundError(LookupError):
pass
class TicketConflictError(TicketStoreError):
pass
class TicketIntegrationUnavailableError(TicketStoreError):
pass
def create_ticket(
session: Session,
principal: object,
*,
record: TicketRecord,
idempotency_key: str,
registry: object | None = None,
) -> TicketRecord:
if not _has_any_scope(principal, REPORT_SCOPE, ADMIN_SCOPE, LEGACY_WRITE_SCOPE):
raise PermissionError("Reporting a ticket requires ticket report access.")
tenant_id = _principal_tenant(principal)
if record.tenant_id != tenant_id:
raise TicketStoreError("Tickets cannot cross tenants.")
if record.revision != 1:
raise TicketStoreError("New tickets start at revision 1.")
clean_key = _bounded(idempotency_key, "Ticket idempotency key", 255)
request_sha256 = _request_sha256(record.to_dict())
replay = _history_replay(session, tenant_id, clean_key, request_sha256)
if replay is not None:
return TicketRecord.from_mapping(replay.snapshot)
if session.get(Ticket, record.ticket_id) is not None:
raise TicketConflictError("A ticket with this identifier already exists.")
if (
session.query(Ticket.id)
.filter(
Ticket.tenant_id == tenant_id,
Ticket.ticket_number == record.ticket_number,
)
.first()
is not None
):
raise TicketConflictError("A ticket with this number already exists.")
routed = _apply_routing(session, principal, record=record, registry=registry)
actor_id = _principal_actor(principal)
row = Ticket(id=routed.ticket_id, tenant_id=tenant_id, ticket_number=routed.ticket_number)
_write_row(row, routed)
row.created_by = actor_id
row.updated_by = actor_id
session.add(row)
session.flush()
_append_history(
session,
row=row,
record=routed,
event_type="reported",
actor_id=actor_id,
reason=routed.change_reason,
idempotency_key=clean_key,
request_sha256=request_sha256,
details={"routing": routed.metadata.get("routing")},
)
_emit(session, row, routed, operation="reported", registry=registry)
return routed
def get_ticket(
session: Session,
principal: object,
*,
ticket_id: str,
include_deleted: bool = False,
) -> TicketRecord | None:
row = _ticket_row(session, _principal_tenant(principal), ticket_id)
if row is None or (row.deleted_at is not None and not include_deleted):
return None
if not can_read_ticket(session, principal, ticket_id=ticket_id):
return None
return _record(row)
def list_tickets(
session: Session,
principal: object,
*,
statuses: Sequence[str] = (),
priorities: Sequence[str] = (),
ticket_types: Sequence[str] = (),
queue_ref: str | None = None,
query: str = "",
include_deleted: bool = False,
offset: int = 0,
limit: int = 100,
) -> tuple[tuple[TicketRecord, ...], int]:
tenant_id = _principal_tenant(principal)
if offset < 0 or not 1 <= limit <= 200:
raise TicketStoreError("Ticket list offset must be non-negative and limit between 1 and 200.")
statement = session.query(Ticket).filter(Ticket.tenant_id == tenant_id)
if not include_deleted:
statement = statement.filter(Ticket.deleted_at.is_(None))
if statuses:
statement = statement.filter(Ticket.status.in_(tuple(dict.fromkeys(statuses))))
if priorities:
statement = statement.filter(Ticket.priority.in_(tuple(dict.fromkeys(priorities))))
if ticket_types:
statement = statement.filter(Ticket.ticket_type.in_(tuple(dict.fromkeys(ticket_types))))
if queue_ref:
statement = statement.filter(Ticket.queue_ref == queue_ref)
clean_query = query.strip().casefold()
if clean_query:
statement = statement.filter(Ticket.search_text.contains(clean_query))
candidates = statement.order_by(
Ticket.service_target_at.asc().nullslast(),
Ticket.priority.desc(),
Ticket.updated_at.desc(),
).all()
accessible = tuple(row for row in candidates if _can_read_row(principal, row))
selected = accessible[offset : offset + limit]
return tuple(_record(row) for row in selected), len(accessible)
def triage_ticket(
session: Session,
principal: object,
*,
ticket_id: str,
expected_revision: int,
changes: Mapping[str, object],
recorded_at: datetime,
change_reason: str,
idempotency_key: str,
registry: object | None = None,
) -> TicketRecord:
allowed = {
"ticket_type",
"priority",
"status",
"title",
"description",
"visibility",
"queue_ref",
"service_target_at",
"reporter",
"requester",
"participants",
"metadata",
}
normalized = _only_changes(changes, allowed)
status = normalized.get("status")
if status in {"resolved", "closed"}:
raise TicketStoreError("Resolve or close tickets through the resolution action.")
return _mutate(
session,
principal,
ticket_id=ticket_id,
expected_revision=expected_revision,
changes=normalized,
recorded_at=recorded_at,
change_reason=change_reason,
idempotency_key=idempotency_key,
event_type="triaged",
required_scopes=(TRIAGE_SCOPE,),
registry=registry,
)
def assign_ticket(
session: Session,
principal: object,
*,
ticket_id: str,
expected_revision: int,
assignee: TicketSubjectRef | None,
recorded_at: datetime,
change_reason: str,
idempotency_key: str,
registry: object | None = None,
) -> TicketRecord:
return _mutate(
session,
principal,
ticket_id=ticket_id,
expected_revision=expected_revision,
changes={"assignee": assignee.to_dict() if assignee else None},
recorded_at=recorded_at,
change_reason=change_reason,
idempotency_key=idempotency_key,
event_type="assigned" if assignee else "unassigned",
required_scopes=(ASSIGN_SCOPE,),
registry=registry,
)
def resolve_ticket(
session: Session,
principal: object,
*,
ticket_id: str,
expected_revision: int,
target_status: str,
recorded_at: datetime,
change_reason: str,
idempotency_key: str,
resolution_summary: str | None = None,
registry: object | None = None,
) -> TicketRecord:
if target_status not in {"in_progress", "waiting", "resolved", "closed", "cancelled"}:
raise TicketStoreError("Unsupported ticket resolution action.")
changes: dict[str, object] = {"status": target_status}
if target_status in {"resolved", "closed"}:
changes["resolution_summary"] = _bounded(
resolution_summary,
"Ticket resolution summary",
20_000,
)
changes["resolved_at"] = recorded_at
elif target_status == "in_progress":
changes["resolution_summary"] = None
changes["resolved_at"] = None
return _mutate(
session,
principal,
ticket_id=ticket_id,
expected_revision=expected_revision,
changes=changes,
recorded_at=recorded_at,
change_reason=change_reason,
idempotency_key=idempotency_key,
event_type=("reopened" if target_status == "in_progress" else target_status),
required_scopes=(RESOLVE_SCOPE,),
registry=registry,
)
def replace_participants(
session: Session,
principal: object,
*,
ticket_id: str,
expected_revision: int,
participants: Sequence[TicketSubjectRef],
recorded_at: datetime,
change_reason: str,
idempotency_key: str,
registry: object | None = None,
) -> TicketRecord:
return _mutate(
session,
principal,
ticket_id=ticket_id,
expected_revision=expected_revision,
changes={"participants": [item.to_dict() for item in participants]},
recorded_at=recorded_at,
change_reason=change_reason,
idempotency_key=idempotency_key,
event_type="participants_changed",
required_scopes=(TRIAGE_SCOPE,),
registry=registry,
)
def add_ticket_link(
session: Session,
principal: object,
*,
ticket_id: str,
expected_revision: int,
link: TicketLink,
recorded_at: datetime,
change_reason: str,
idempotency_key: str,
registry: object | None = None,
) -> TicketRecord:
current = _required_ticket(session, principal, ticket_id=ticket_id, lock=True)
record = _record(current)
if any(item.link_id == link.link_id for item in record.links):
raise TicketConflictError("A ticket link with this identifier already exists.")
return _mutate_locked(
session,
principal,
row=current,
expected_revision=expected_revision,
changes={"links": [item.to_dict() for item in (*record.links, link)]},
recorded_at=recorded_at,
change_reason=change_reason,
idempotency_key=idempotency_key,
event_type="link_added",
required_scopes=(TRIAGE_SCOPE,),
details={"link": link.to_dict()},
registry=registry,
)
def remove_ticket_link(
session: Session,
principal: object,
*,
ticket_id: str,
link_id: str,
expected_revision: int,
recorded_at: datetime,
change_reason: str,
idempotency_key: str,
registry: object | None = None,
) -> TicketRecord:
current = _required_ticket(session, principal, ticket_id=ticket_id, lock=True)
record = _record(current)
selected = next((item for item in record.links if item.link_id == link_id), None)
if selected is None:
raise TicketNotFoundError("Ticket link not found.")
return _mutate_locked(
session,
principal,
row=current,
expected_revision=expected_revision,
changes={"links": [item.to_dict() for item in record.links if item.link_id != link_id]},
recorded_at=recorded_at,
change_reason=change_reason,
idempotency_key=idempotency_key,
event_type="link_removed",
required_scopes=(TRIAGE_SCOPE,),
details={"link": selected.to_dict()},
registry=registry,
)
def add_ticket_comment(
session: Session,
principal: object,
*,
ticket_id: str,
expected_revision: int,
comment_id: str,
body: str,
visibility: str,
recorded_at: datetime,
idempotency_key: str,
registry: object | None = None,
) -> tuple[TicketRecord, dict[str, object]]:
if visibility not in {"external", "internal"}:
raise TicketStoreError("Ticket comments are external or internal.")
if visibility == "internal" and not _has_any_scope(principal, TRIAGE_SCOPE, ADMIN_SCOPE, LEGACY_WRITE_SCOPE):
raise PermissionError("Internal ticket comments require triage access.")
if visibility == "external" and not _has_any_scope(
principal,
REPORT_SCOPE,
TRIAGE_SCOPE,
ADMIN_SCOPE,
LEGACY_WRITE_SCOPE,
):
raise PermissionError("External ticket comments require report access.")
row = _required_ticket(session, principal, ticket_id=ticket_id, lock=True)
if not _can_read_row(principal, row):
raise PermissionError("Ticket comment access is denied.")
clean_comment_id = _bounded(comment_id, "Ticket comment identifier", 255)
clean_body = _bounded(body, "Ticket comment", 20_000)
record = _mutate_locked(
session,
principal,
row=row,
expected_revision=expected_revision,
changes={},
recorded_at=recorded_at,
change_reason="Added a ticket comment.",
idempotency_key=idempotency_key,
event_type="comment_added",
required_scopes=(),
details={"comment_id": clean_comment_id, "visibility": visibility},
registry=registry,
)
existing = (
session.query(TicketComment)
.filter(
TicketComment.tenant_id == record.tenant_id,
TicketComment.comment_id == clean_comment_id,
)
.one_or_none()
)
if existing is None:
existing = TicketComment(
tenant_id=record.tenant_id,
ticket_id=record.ticket_id,
comment_id=clean_comment_id,
ticket_revision=record.revision,
visibility=visibility,
body=clean_body,
created_by=_principal_actor(principal),
created_at=recorded_at,
updated_at=recorded_at,
)
session.add(existing)
session.flush()
elif existing.ticket_id != record.ticket_id or existing.body != clean_body or existing.visibility != visibility:
raise TicketConflictError("Ticket comment identifier was reused for another comment.")
return record, _comment_payload(existing)
def list_ticket_comments(
session: Session,
principal: object,
*,
ticket_id: str,
limit: int = 200,
) -> tuple[dict[str, object], ...]:
if not 1 <= limit <= 500:
raise TicketStoreError("Ticket comment limit must be between 1 and 500.")
row = _required_ticket(session, principal, ticket_id=ticket_id)
if not _can_read_row(principal, row):
raise PermissionError("Ticket comment access is denied.")
query = session.query(TicketComment).filter(
TicketComment.tenant_id == row.tenant_id,
TicketComment.ticket_id == row.id,
)
if not _has_any_scope(principal, TRIAGE_SCOPE, ADMIN_SCOPE, LEGACY_WRITE_SCOPE):
query = query.filter(TicketComment.visibility == "external")
rows = query.order_by(TicketComment.created_at.asc(), TicketComment.id.asc()).limit(limit).all()
return tuple(_comment_payload(item) for item in rows)
def ticket_history(
session: Session,
principal: object,
*,
ticket_id: str,
limit: int = 200,
) -> tuple[dict[str, object], ...]:
if not 1 <= limit <= 500:
raise TicketStoreError("Ticket history limit must be between 1 and 500.")
row = _required_ticket(session, principal, ticket_id=ticket_id)
if not _can_read_row(principal, row):
raise PermissionError("Ticket history access is denied.")
items = (
session.query(TicketHistory)
.filter(TicketHistory.tenant_id == row.tenant_id, TicketHistory.ticket_id == row.id)
.order_by(TicketHistory.revision.desc())
.limit(limit)
.all()
)
disclose_details = _has_any_scope(principal, TRIAGE_SCOPE, ADMIN_SCOPE, LEGACY_WRITE_SCOPE)
return tuple(
{
"revision": item.revision,
"event_type": item.event_type,
"occurred_at": _iso(item.occurred_at),
"actor_id": item.actor_id,
"reason": item.reason,
"details": dict(item.details or {}) if disclose_details else {},
}
for item in items
)
def escalate_ticket_to_case(
session: Session,
principal: object,
*,
ticket_id: str,
expected_revision: int,
case_type_key: str,
occurred_at: datetime,
handoff_note: str | None,
idempotency_key: str,
registry: object | None,
) -> tuple[TicketRecord, dict[str, object]]:
if not _has_any_scope(principal, TRIAGE_SCOPE, ADMIN_SCOPE, LEGACY_WRITE_SCOPE):
raise PermissionError("Ticket escalation requires triage access.")
provider = ticket_case_escalation_provider(registry)
if provider is None:
raise TicketIntegrationUnavailableError("Cases escalation is not available in this installation.")
row = _required_ticket(session, principal, ticket_id=ticket_id, lock=True)
clean_key = _bounded(idempotency_key, "Ticket escalation idempotency key", 255)
request = {
"ticket_id": ticket_id,
"expected_revision": expected_revision,
"case_type_key": case_type_key,
"occurred_at": occurred_at,
"handoff_note": handoff_note,
}
request_sha256 = _request_sha256(request)
replay = (
session.query(TicketEscalation)
.filter(
TicketEscalation.tenant_id == row.tenant_id,
TicketEscalation.ticket_id == row.id,
TicketEscalation.idempotency_key == clean_key,
)
.one_or_none()
)
if replay is not None:
if replay.request_sha256 != request_sha256:
raise TicketConflictError("Ticket escalation idempotency key was reused for another request.")
return _record(row), _escalation_payload(replay, replayed=True)
if row.revision != expected_revision:
raise TicketConflictError("Ticket revision conflict: the expected revision is stale.")
current = _record(row)
command = TicketCaseEscalationCommand(
tenant_id=row.tenant_id,
ticket_id=row.id,
ticket_number=row.ticket_number,
title=current.title,
case_type_key=_bounded(case_type_key, "Case type key", 120),
occurred_at=_aware(occurred_at, "Ticket escalation occurred_at"),
idempotency_key=clean_key,
handoff_note=_optional_bounded(handoff_note, "Ticket escalation handoff note", 10_000),
metadata={"ticket_revision": current.revision},
)
try:
result = provider.escalate_ticket(session, principal, command=command)
except PermissionError:
raise
except ValueError as exc:
raise TicketStoreError(str(exc)) from exc
link = TicketLink(
link_id=f"case:{result.provider_id}:{result.case_id}",
kind="case",
owner_module="cases",
resource_type="case",
resource_id=result.case_id,
relation="escalated_to",
label=result.case_number,
url=result.case_url,
metadata={"provider_id": result.provider_id},
)
if any(item.link_id == link.link_id for item in current.links):
links = current.links
else:
links = (*current.links, link)
updated = _mutate_locked(
session,
principal,
row=row,
expected_revision=expected_revision,
changes={"links": [item.to_dict() for item in links]},
recorded_at=occurred_at,
change_reason=f"Escalated ticket to case {result.case_number}.",
idempotency_key=f"ticket-case-link:{clean_key}",
event_type="case_escalated",
required_scopes=(),
details={"case_id": result.case_id, "case_number": result.case_number},
registry=registry,
)
escalation = TicketEscalation(
tenant_id=row.tenant_id,
ticket_id=row.id,
provider_id=result.provider_id,
idempotency_key=clean_key,
request_sha256=request_sha256,
occurred_at=occurred_at,
actor_id=_principal_actor(principal),
case_id=result.case_id,
case_number=result.case_number,
case_url=result.case_url,
handoff_note=command.handoff_note,
outcome={"replayed_by_provider": result.replayed, **dict(result.metadata)},
)
session.add(escalation)
session.flush()
return updated, _escalation_payload(escalation, replayed=False)
def delete_ticket(
session: Session,
principal: object,
*,
ticket_id: str,
expected_revision: int,
occurred_at: datetime,
reason: str,
idempotency_key: str,
registry: object | None = None,
) -> TicketRecord:
return _mutate(
session,
principal,
ticket_id=ticket_id,
expected_revision=expected_revision,
changes={"deleted_at": occurred_at},
recorded_at=occurred_at,
change_reason=reason,
idempotency_key=idempotency_key,
event_type="deleted",
required_scopes=(ADMIN_SCOPE,),
registry=registry,
)
def can_read_ticket(session: Session, principal: object, *, ticket_id: str) -> bool:
row = _ticket_row(session, _principal_tenant(principal), ticket_id)
return row is not None and _can_read_row(principal, row)
def integration_availability(registry: object | None) -> dict[str, object]:
routing = ticket_routing_provider(registry)
escalation = ticket_case_escalation_provider(registry)
active = _active_modules(registry)
return {
"routing": {"available": routing is not None, "provider": type(routing).__name__ if routing else None},
"case_escalation": {"available": escalation is not None, "provider": type(escalation).__name__ if escalation else None},
"modules": {
name: name in active
for name in ("cases", "helpdesk", "projects", "wiki", "files", "search")
},
"consequences": {
"cases": "Escalation is disabled; tickets remain independently resolvable.",
"helpdesk": "Manual queue and target selection remains available; automatic routing is disabled.",
"projects": "Project links remain typed references and are not validated by Tickets.",
"wiki": "Wiki links remain typed references and are not validated by Tickets.",
"files": "Attachments remain external file references; Tickets stores no file bytes.",
"search": "Ticket APIs remain usable; global indexing and discovery are disabled.",
},
}
class SqlTicketRegistry:
def __init__(self, registry: object | None = None) -> None:
self.registry = registry
def create_ticket(self, session: object, principal: object, *, record: TicketRecord, idempotency_key: str) -> TicketRecord:
return create_ticket(_session(session), principal, record=record, idempotency_key=idempotency_key, registry=self.registry)
def get_ticket(self, session: object, principal: object, *, ticket_id: str) -> TicketRecord | None:
return get_ticket(_session(session), principal, ticket_id=ticket_id)
def _mutate(
session: Session,
principal: object,
*,
ticket_id: str,
expected_revision: int,
changes: Mapping[str, object],
recorded_at: datetime,
change_reason: str,
idempotency_key: str,
event_type: str,
required_scopes: Sequence[str],
details: Mapping[str, object] | None = None,
registry: object | None,
) -> TicketRecord:
row = _required_ticket(session, principal, ticket_id=ticket_id, lock=True)
return _mutate_locked(
session,
principal,
row=row,
expected_revision=expected_revision,
changes=changes,
recorded_at=recorded_at,
change_reason=change_reason,
idempotency_key=idempotency_key,
event_type=event_type,
required_scopes=required_scopes,
details=details,
registry=registry,
)
def _mutate_locked(
session: Session,
principal: object,
*,
row: Ticket,
expected_revision: int,
changes: Mapping[str, object],
recorded_at: datetime,
change_reason: str,
idempotency_key: str,
event_type: str,
required_scopes: Sequence[str],
details: Mapping[str, object] | None = None,
registry: object | None,
) -> TicketRecord:
if required_scopes and not _has_any_scope(principal, *required_scopes, ADMIN_SCOPE, LEGACY_WRITE_SCOPE):
raise PermissionError(f"Ticket action requires one of: {', '.join(required_scopes)}.")
if not _can_read_row(principal, row):
raise PermissionError("Ticket access is denied.")
clean_key = _bounded(idempotency_key, "Ticket idempotency key", 255)
clean_reason = _bounded(change_reason, "Ticket change reason", 1_000)
instant = _aware(recorded_at, "Ticket recorded_at")
request = {
"ticket_id": row.id,
"expected_revision": expected_revision,
"changes": changes,
"recorded_at": instant,
"change_reason": clean_reason,
"event_type": event_type,
}
request_sha256 = _request_sha256(request)
replay = _history_replay(session, row.tenant_id, clean_key, request_sha256)
if replay is not None:
return TicketRecord.from_mapping(replay.snapshot)
if row.revision != expected_revision:
raise TicketConflictError("Ticket revision conflict: the expected revision is stale.")
current = _record(row)
payload = current.to_dict()
payload.update(dict(changes))
payload.update(
{
"revision": current.revision + 1,
"recorded_at": instant.isoformat(),
"change_reason": clean_reason,
}
)
updated = TicketRecord.from_mapping(payload)
validate_transition(current.status, updated.status)
_write_row(row, updated)
row.updated_by = _principal_actor(principal)
session.add(row)
session.flush()
_append_history(
session,
row=row,
record=updated,
event_type=event_type,
actor_id=row.updated_by,
reason=clean_reason,
idempotency_key=clean_key,
request_sha256=request_sha256,
details=details or {},
)
_emit(session, row, updated, operation=event_type, registry=registry)
return updated
def _apply_routing(
session: Session,
principal: object,
*,
record: TicketRecord,
registry: object | None,
) -> TicketRecord:
provider = ticket_routing_provider(registry)
if provider is None:
return record
plan = provider.route_ticket(
session,
principal,
request=TicketRoutingRequest(
tenant_id=record.tenant_id,
ticket_id=record.ticket_id,
ticket_type=record.ticket_type,
priority=record.priority,
title=record.title,
received_at=record.received_at,
queue_hint=record.queue_ref,
attributes=dict(record.metadata),
),
)
payload = record.to_dict()
if not record.queue_ref and plan.queue_ref:
payload["queue_ref"] = plan.queue_ref
if record.service_target_at is None and plan.service_target_at is not None:
payload["service_target_at"] = plan.service_target_at.isoformat()
payload["metadata"] = {
**dict(record.metadata),
"routing": {
"provider_id": plan.provider_id,
"explanation": plan.explanation,
**dict(plan.metadata),
},
}
return TicketRecord.from_mapping(payload)
def _append_history(
session: Session,
*,
row: Ticket,
record: TicketRecord,
event_type: str,
actor_id: str | None,
reason: str,
idempotency_key: str,
request_sha256: str,
details: Mapping[str, object],
) -> None:
session.add(
TicketHistory(
tenant_id=row.tenant_id,
ticket_id=row.id,
revision=record.revision,
event_type=event_type,
occurred_at=record.recorded_at,
actor_id=actor_id,
reason=reason,
idempotency_key=idempotency_key,
request_sha256=request_sha256,
snapshot=record.to_dict(),
details=dict(details),
)
)
session.flush()
def _history_replay(session: Session, tenant_id: str, key: str, digest: str) -> TicketHistory | None:
row = (
session.query(TicketHistory)
.filter(TicketHistory.tenant_id == tenant_id, TicketHistory.idempotency_key == key)
.one_or_none()
)
if row is not None and row.request_sha256 != digest:
raise TicketConflictError("Ticket idempotency key was reused for another request.")
return row
def _write_row(row: Ticket, record: TicketRecord) -> None:
row.revision = record.revision
row.ticket_type = record.ticket_type
row.priority = record.priority
row.status = record.status
row.title = record.title
row.description = record.description
row.visibility = record.visibility
row.queue_ref = record.queue_ref
row.assignee = record.assignee.to_dict() if record.assignee else None
row.reporter = record.reporter.to_dict() if record.reporter else None
row.requester = record.requester.to_dict() if record.requester else None
row.participants = [item.to_dict() for item in record.participants]
row.links = [item.to_dict() for item in record.links]
row.metadata_payload = dict(record.metadata)
row.search_text = _search_text(record)
row.received_at = record.received_at
row.recorded_at = record.recorded_at
row.change_reason = record.change_reason
row.service_target_at = record.service_target_at
row.resolved_at = record.resolved_at
row.resolution_summary = record.resolution_summary
row.deleted_at = record.deleted_at
def _record(row: Ticket) -> TicketRecord:
return TicketRecord.from_mapping(
{
"tenant_id": row.tenant_id,
"ticket_id": row.id,
"ticket_number": row.ticket_number,
"revision": row.revision,
"ticket_type": row.ticket_type,
"priority": row.priority,
"status": row.status,
"title": row.title,
"description": row.description,
"visibility": row.visibility,
"queue_ref": row.queue_ref,
"assignee": row.assignee,
"reporter": row.reporter,
"requester": row.requester,
"participants": row.participants or [],
"links": row.links or [],
"metadata": row.metadata_payload or {},
"received_at": _iso(_db_aware(row.received_at)),
"recorded_at": _iso(_db_aware(row.recorded_at)),
"service_target_at": _iso(_db_aware(row.service_target_at)),
"resolved_at": _iso(_db_aware(row.resolved_at)),
"resolution_summary": row.resolution_summary,
"deleted_at": _iso(_db_aware(row.deleted_at)),
"change_reason": row.change_reason,
}
)
def _ticket_row(session: Session, tenant_id: str, ticket_id: str, *, lock: bool = False) -> Ticket | None:
query = session.query(Ticket).filter(Ticket.tenant_id == tenant_id, Ticket.id == ticket_id)
if lock:
query = query.with_for_update()
return query.one_or_none()
def _required_ticket(session: Session, principal: object, *, ticket_id: str, lock: bool = False) -> Ticket:
row = _ticket_row(session, _principal_tenant(principal), ticket_id, lock=lock)
if row is None:
raise TicketNotFoundError("Ticket not found.")
return row
def _can_read_row(principal: object, row: Ticket) -> bool:
if row.tenant_id != _principal_tenant(principal) or not _has_scope(principal, READ_SCOPE):
return False
if _has_any_scope(principal, TRIAGE_SCOPE, ASSIGN_SCOPE, RESOLVE_SCOPE, ADMIN_SCOPE, LEGACY_WRITE_SCOPE):
return True
if row.visibility == "tenant":
return True
subjects = set(_principal_subjects(principal))
values = [row.assignee, row.reporter, row.requester, *(row.participants or [])]
return any(
isinstance(value, Mapping)
and (str(value.get("kind") or ""), str(value.get("id") or "")) in subjects
for value in values
) or str(row.created_by or "") in _principal_actor_ids(principal)
def _emit(session: Session, row: Ticket, record: TicketRecord, *, operation: str, registry: object | None) -> None:
emit_platform_event(
session,
PlatformEvent(
type=f"tickets.ticket.{operation}",
module_id="tickets",
payload={"ticket_id": row.id, "ticket_number": row.ticket_number, "revision": record.revision, "status": record.status},
occurred_at=record.recorded_at,
actor=EventActorRef(type="account", id=row.updated_by or row.created_by),
tenant=EventTenantRef(id=row.tenant_id),
subject=EventObjectRef(type="ticket", id=row.id, label=record.title),
resource=EventObjectRef(type="ticket", id=row.id, label=record.title),
classification="restricted",
),
registry=registry,
)
def _comment_payload(row: TicketComment) -> dict[str, object]:
return {
"comment_id": row.comment_id,
"ticket_revision": row.ticket_revision,
"visibility": row.visibility,
"body": row.body,
"created_by": row.created_by,
"created_at": _iso(row.created_at),
}
def _escalation_payload(row: TicketEscalation, *, replayed: bool) -> dict[str, object]:
return {
"provider_id": row.provider_id,
"case_id": row.case_id,
"case_number": row.case_number,
"case_url": row.case_url,
"occurred_at": _iso(row.occurred_at),
"actor_id": row.actor_id,
"handoff_note": row.handoff_note,
"outcome": dict(row.outcome or {}),
"replayed": replayed,
}
def _search_text(record: TicketRecord) -> str:
values = [
record.ticket_number,
record.ticket_type,
record.priority,
record.status,
record.title,
record.description,
record.queue_ref or "",
record.resolution_summary or "",
*(item.label or item.id for item in record.participants),
*(item.label or item.resource_id for item in record.links),
]
return "\n".join(values).casefold()
def _principal_subjects(principal: object) -> tuple[tuple[str, str], ...]:
values: list[tuple[str, str]] = []
singular = {
"account": getattr(principal, "account_id", None),
"identity": getattr(principal, "identity_id", None),
"function_assignment": getattr(principal, "acting_assignment_id", None),
}
for kind, value in singular.items():
if str(value or "").strip():
values.append((kind, str(value)))
for kind, attribute in (
("group", "group_ids"),
("role", "role_ids"),
("function_assignment", "function_assignment_ids"),
):
values.extend((kind, str(item)) for item in getattr(principal, attribute, ()) or () if str(item or "").strip())
return tuple(dict.fromkeys(values))
def _principal_actor_ids(principal: object) -> tuple[str, ...]:
user = getattr(principal, "user", None)
return tuple(
dict.fromkeys(
str(value)
for value in (
getattr(principal, "account_id", None),
getattr(principal, "identity_id", None),
getattr(principal, "membership_id", None),
getattr(user, "id", None),
)
if str(value or "").strip()
)
)
def _principal_actor(principal: object) -> str | None:
values = _principal_actor_ids(principal)
return values[0] if values else None
def _principal_tenant(principal: object) -> str:
tenant_id = str(getattr(principal, "tenant_id", "") or "").strip()
if not tenant_id:
raise TicketStoreError("Ticket operations require a tenant-bound principal.")
return tenant_id
def _has_scope(principal: object, scope: str) -> bool:
method = getattr(principal, "has", None)
if callable(method):
return bool(method(scope))
return scopes_grant_compatible(frozenset(getattr(principal, "scopes", ()) or ()), scope)
def _has_any_scope(principal: object, *scopes: str) -> bool:
return any(_has_scope(principal, scope) for scope in scopes)
def _only_changes(value: Mapping[str, object], allowed: set[str]) -> dict[str, object]:
unexpected = set(value) - allowed
if unexpected:
raise TicketStoreError("Unsupported ticket changes: " + ", ".join(sorted(unexpected)))
return dict(value)
def _bounded(value: object, label: str, maximum: int) -> str:
clean = str(value or "").strip()
if not clean or len(clean) > maximum:
raise TicketStoreError(f"{label} must contain 1 to {maximum} characters.")
return clean
def _optional_bounded(value: object, label: str, maximum: int) -> str | None:
if value is None:
return None
return _bounded(value, label, maximum)
def _aware(value: datetime, label: str) -> datetime:
if value.tzinfo is None or value.utcoffset() is None:
raise TicketStoreError(f"{label} must include a timezone.")
return value
def _request_sha256(value: object) -> str:
payload = json.dumps(_json_value(value), sort_keys=True, separators=(",", ":"), ensure_ascii=True).encode("utf-8")
return hashlib.sha256(payload).hexdigest()
def _json_value(value: object) -> Any:
if isinstance(value, datetime):
return value.astimezone(UTC).isoformat()
if isinstance(value, Mapping):
return {str(key): _json_value(item) for key, item in value.items()}
if isinstance(value, (list, tuple)):
return [_json_value(item) for item in value]
return value
def _iso(value: datetime | None) -> str | None:
return value.isoformat() if value else None
def _db_aware(value: datetime | None) -> datetime | None:
if value is not None and (value.tzinfo is None or value.utcoffset() is None):
return value.replace(tzinfo=UTC)
return value
def _active_modules(registry: object | None) -> frozenset[str]:
if registry is None:
return frozenset()
method = getattr(registry, "active_module_ids", None)
if callable(method):
return frozenset(str(item) for item in method())
manifests = getattr(registry, "manifests", None)
if callable(manifests):
return frozenset(str(getattr(item, "id", "")) for item in manifests())
return frozenset()
def _session(value: object) -> Session:
if not isinstance(value, Session):
raise TypeError("Tickets require a SQLAlchemy session.")
return value
__all__ = [
"ADMIN_SCOPE",
"ASSIGN_SCOPE",
"CAPABILITY_TICKETS_REGISTRY",
"LEGACY_WRITE_SCOPE",
"READ_SCOPE",
"REPORT_SCOPE",
"RESOLVE_SCOPE",
"SqlTicketRegistry",
"TRIAGE_SCOPE",
"TicketConflictError",
"TicketIntegrationUnavailableError",
"TicketNotFoundError",
"TicketStoreError",
"add_ticket_comment",
"add_ticket_link",
"assign_ticket",
"can_read_ticket",
"create_ticket",
"delete_ticket",
"escalate_ticket_to_case",
"get_ticket",
"integration_availability",
"list_ticket_comments",
"list_tickets",
"remove_ticket_link",
"replace_participants",
"resolve_ticket",
"ticket_history",
"triage_ticket",
]