Files
govoplan-core/src/govoplan_core/server/dsar.py
T

407 lines
13 KiB
Python

from __future__ import annotations
import re
from datetime import datetime
from typing import Any, Literal
from fastapi import APIRouter, Depends, Header, HTTPException, Query, Request, Response, status
from pydantic import BaseModel, Field, field_validator
from sqlalchemy.orm import Session
from govoplan_core.audit.logging import audit_event
from govoplan_core.auth import ApiPrincipal, get_api_principal, has_scope
from govoplan_core.core.concurrency import (
ConcurrencyError,
MissingPreconditionError,
RevisionConflictError,
assert_revision_precondition,
)
from govoplan_core.core.dsar import DsarSubjectRef
from govoplan_core.db.session import get_session
from govoplan_core.privacy.dsar_workflow import (
DataSubjectRequest,
create_data_subject_request,
data_subject_export,
data_subject_request_dict,
execute_data_subject_erasure,
get_data_subject_request,
list_data_subject_requests,
plan_data_subject_erasure,
search_data_subject_request,
)
READ_SCOPE = "access:privacy:read"
MANAGE_SCOPE = "access:privacy:manage"
EXPORT_SCOPE = "access:privacy:export"
ERASE_SCOPE = "access:privacy:erase"
class DataSubjectSelectorRequest(BaseModel):
account_id: str | None = Field(default=None, max_length=36)
identity_id: str | None = Field(default=None, max_length=36)
membership_id: str | None = Field(default=None, max_length=36)
email: str | None = Field(default=None, max_length=320)
external_references: dict[str, str] = Field(default_factory=dict)
@field_validator("external_references")
@classmethod
def validate_external_references(
cls,
value: dict[str, str],
) -> dict[str, str]:
if len(value) > 50:
raise ValueError("At most 50 external subject references are allowed.")
normalized: dict[str, str] = {}
for raw_key, raw_value in value.items():
key = str(raw_key).strip()
item = str(raw_value).strip()
if not key or not item:
continue
if len(key) > 120 or len(item) > 500:
raise ValueError(
"External subject-reference namespaces are limited to 120 "
"characters and values to 500 characters."
)
normalized[key] = item
return normalized
class DataSubjectRequestCreate(BaseModel):
reference: str = Field(min_length=1, max_length=120)
request_kind: Literal["access", "erasure", "access_and_erasure"]
subject: DataSubjectSelectorRequest
purpose: str = Field(min_length=1, max_length=1000)
legal_basis: str | None = Field(default=None, max_length=1000)
due_at: datetime | None = None
notes: str | None = Field(default=None, max_length=10_000)
class RevisionMutationRequest(BaseModel):
base_revision: int = Field(ge=1)
class DataSubjectExecutionRequest(RevisionMutationRequest):
action_ids: list[str] = Field(min_length=1, max_length=10_000)
confirmation: str = Field(min_length=1, max_length=100)
class DataSubjectRequestResponse(BaseModel):
request: dict[str, Any]
search: dict[str, Any] = Field(default_factory=dict)
erasure_plan: dict[str, Any] = Field(default_factory=dict)
execution: dict[str, Any] = Field(default_factory=dict)
class DataSubjectRequestListResponse(BaseModel):
items: list[dict[str, Any]] = Field(default_factory=list)
router = APIRouter(
prefix="/admin/privacy/data-subject-requests",
tags=["data-subject-requests"],
)
@router.get("", response_model=DataSubjectRequestListResponse)
def list_requests(
limit: int = Query(default=200, ge=1, le=500),
principal: ApiPrincipal = Depends(get_api_principal),
session: Session = Depends(get_session),
) -> DataSubjectRequestListResponse:
_require(principal, READ_SCOPE)
rows = list_data_subject_requests(
session,
tenant_id=principal.tenant_id,
limit=limit,
)
return DataSubjectRequestListResponse(
items=[
data_subject_request_dict(row, include_subject=True)
for row in rows
]
)
@router.post(
"",
response_model=DataSubjectRequestResponse,
status_code=status.HTTP_201_CREATED,
)
def create_request(
payload: DataSubjectRequestCreate,
principal: ApiPrincipal = Depends(get_api_principal),
session: Session = Depends(get_session),
) -> DataSubjectRequestResponse:
_require(principal, MANAGE_SCOPE)
try:
row = create_data_subject_request(
session,
tenant_id=principal.tenant_id,
reference=payload.reference,
request_kind=payload.request_kind,
subject=_subject(payload.subject),
purpose=payload.purpose,
legal_basis=payload.legal_basis,
due_at=payload.due_at,
requested_by_account_id=principal.account_id,
notes=payload.notes,
)
_audit(session, principal, row, "privacy.dsar.created")
session.commit()
session.refresh(row)
return _detail(row)
except ValueError as exc:
session.rollback()
raise HTTPException(status_code=400, detail=str(exc)) from exc
@router.get("/{request_id}", response_model=DataSubjectRequestResponse)
def get_request(
request_id: str,
principal: ApiPrincipal = Depends(get_api_principal),
session: Session = Depends(get_session),
) -> DataSubjectRequestResponse:
_require(principal, READ_SCOPE)
return _detail(_row(session, principal, request_id))
@router.post("/{request_id}/search", response_model=DataSubjectRequestResponse)
def search_request(
request_id: str,
payload: RevisionMutationRequest,
request: Request,
if_match: str | None = Header(default=None, alias="If-Match"),
principal: ApiPrincipal = Depends(get_api_principal),
session: Session = Depends(get_session),
) -> DataSubjectRequestResponse:
_require(principal, MANAGE_SCOPE)
return _mutate(
session,
principal,
request_id,
payload.base_revision,
if_match,
lambda row: search_data_subject_request(
session,
registry=_registry(request),
row=row,
expected_revision=payload.base_revision,
),
"privacy.dsar.searched",
)
@router.post("/{request_id}/erasure-plan", response_model=DataSubjectRequestResponse)
def plan_erasure(
request_id: str,
payload: RevisionMutationRequest,
request: Request,
if_match: str | None = Header(default=None, alias="If-Match"),
principal: ApiPrincipal = Depends(get_api_principal),
session: Session = Depends(get_session),
) -> DataSubjectRequestResponse:
_require(principal, MANAGE_SCOPE)
return _mutate(
session,
principal,
request_id,
payload.base_revision,
if_match,
lambda row: plan_data_subject_erasure(
session,
registry=_registry(request),
row=row,
expected_revision=payload.base_revision,
),
"privacy.dsar.erasure_planned",
)
@router.post("/{request_id}/execute", response_model=DataSubjectRequestResponse)
def execute_erasure(
request_id: str,
payload: DataSubjectExecutionRequest,
request: Request,
if_match: str | None = Header(default=None, alias="If-Match"),
principal: ApiPrincipal = Depends(get_api_principal),
session: Session = Depends(get_session),
) -> DataSubjectRequestResponse:
_require(principal, ERASE_SCOPE)
if payload.confirmation != f"ERASE {request_id}":
raise HTTPException(
status_code=400,
detail=f'Type "ERASE {request_id}" to confirm the selected actions.',
)
return _mutate(
session,
principal,
request_id,
payload.base_revision,
if_match,
lambda row: execute_data_subject_erasure(
session,
registry=_registry(request),
row=row,
expected_revision=payload.base_revision,
action_ids=payload.action_ids,
),
"privacy.dsar.erasure_executed",
details={"action_ids": payload.action_ids},
)
@router.get("/{request_id}/export")
def export_request(
request_id: str,
principal: ApiPrincipal = Depends(get_api_principal),
session: Session = Depends(get_session),
) -> Response:
_require(principal, EXPORT_SCOPE)
row = _row(session, principal, request_id)
content = data_subject_export(row)
_audit(
session,
principal,
row,
"privacy.dsar.exported",
details={"export_bytes": len(content)},
)
session.commit()
safe_reference = re.sub(r"[^A-Za-z0-9._-]+", "-", row.reference).strip("-")
filename = f"dsar-{safe_reference or row.id}.json"
return Response(
content=content,
media_type="application/json",
headers={"Content-Disposition": f'attachment; filename="{filename}"'},
)
def _mutate(
session: Session,
principal: ApiPrincipal,
request_id: str,
base_revision: int,
if_match: str | None,
operation: Any,
audit_action: str,
*,
details: dict[str, Any] | None = None,
) -> DataSubjectRequestResponse:
try:
assert_revision_precondition(
if_match,
resource_type="data_subject_request",
resource_id=request_id,
submitted_base_revision=base_revision,
)
row = get_data_subject_request(
session,
tenant_id=principal.tenant_id,
request_id=request_id,
for_update=True,
)
operation(row)
_audit(session, principal, row, audit_action, details=details)
session.commit()
session.refresh(row)
return _detail(row)
except LookupError as exc:
session.rollback()
raise HTTPException(status_code=404, detail=str(exc)) from exc
except MissingPreconditionError as exc:
session.rollback()
raise HTTPException(status_code=428, detail=exc.as_dict()) from exc
except RevisionConflictError as exc:
session.rollback()
raise HTTPException(status_code=409, detail=exc.as_dict()) from exc
except ConcurrencyError as exc:
session.rollback()
raise HTTPException(status_code=412, detail=str(exc)) from exc
except ValueError as exc:
session.rollback()
raise HTTPException(status_code=400, detail=str(exc)) from exc
def _row(
session: Session,
principal: ApiPrincipal,
request_id: str,
) -> DataSubjectRequest:
try:
return get_data_subject_request(
session,
tenant_id=principal.tenant_id,
request_id=request_id,
)
except LookupError as exc:
raise HTTPException(status_code=404, detail=str(exc)) from exc
def _detail(row: DataSubjectRequest) -> DataSubjectRequestResponse:
return DataSubjectRequestResponse(
request=data_subject_request_dict(row, include_subject=True),
search=dict(row.search_result),
erasure_plan=dict(row.erasure_plan),
execution=dict(row.execution_result),
)
def _subject(payload: DataSubjectSelectorRequest) -> DsarSubjectRef:
return DsarSubjectRef(
account_id=_text(payload.account_id),
identity_id=_text(payload.identity_id),
membership_id=_text(payload.membership_id),
email=_text(payload.email),
external_references={
str(key).strip(): str(value).strip()
for key, value in payload.external_references.items()
if str(key).strip() and str(value).strip()
},
)
def _registry(request: Request) -> object:
registry = getattr(request.app.state, "govoplan_registry", None)
if registry is None:
raise HTTPException(status_code=503, detail="Module registry is unavailable.")
return registry
def _require(principal: ApiPrincipal, scope: str) -> None:
if not has_scope(principal, scope):
raise HTTPException(status_code=403, detail=f"Missing scope: {scope}")
def _audit(
session: Session,
principal: ApiPrincipal,
row: DataSubjectRequest,
action: str,
*,
details: dict[str, Any] | None = None,
) -> None:
audit_event(
session,
tenant_id=principal.tenant_id,
action=action,
user_id=principal.membership_id,
api_key_id=principal.api_key_id,
object_type="data_subject_request",
object_id=row.id,
details={
"reference": row.reference,
"status": row.status,
"resource_revision": row.resource_revision,
**(details or {}),
},
)
def _text(value: str | None) -> str | None:
normalized = (value or "").strip()
return normalized or None
__all__ = ["router"]