Files
govoplan-records/src/govoplan_records/backend/router.py
T

691 lines
24 KiB
Python

from __future__ import annotations
import hashlib
from typing import Any
from fastapi import APIRouter, Depends, HTTPException, Query, status
from sqlalchemy.exc import IntegrityError
from sqlalchemy.orm import Session
from govoplan_core.auth import ApiPrincipal, get_api_principal, has_scope
from govoplan_core.audit.logging import audit_from_principal
from govoplan_core.core.records import RecordFilingRequest, RecordSourceLocator
from govoplan_core.db.session import get_session
from govoplan_records.backend.manifest import ADMIN_SCOPE, READ_SCOPE, WRITE_SCOPE
from govoplan_records.backend.schemas import (
FilePlanNodeWriteRequest,
RecordAppraisalRequest,
RecordArchiveProviderResponse,
RecordCatalogResponse,
RecordClassWriteRequest,
RecordCloseRequest,
RecordCreateRequest,
RecordDetailResponse,
RecordDispositionCreateRequest,
RecordDispositionFinalizeRequest,
RecordDispositionWithdrawRequest,
RecordHoldCreateRequest,
RecordHoldReleaseRequest,
RecordItemCreateRequest,
RecordLifecycleActionRequest,
RecordListResponse,
RecordRecoveryStatusResponse,
RecordSourceProviderResponse,
RecordTransferDispatchRequest,
RecordTransferPackageCreateRequest,
RecordUpdateRequest,
RecordVolumeCreateRequest,
)
from govoplan_records.backend.service import (
RecordConflictError,
RecordNotFoundError,
RecordSourceUnavailableError,
RecordStoreError,
SqlRecordRegistry,
)
from govoplan_records.backend.recovery import (
RecordRecoveryError,
begin_record_atomic_recovery,
bind_record_recovery_operation,
reset_record_recovery_operation,
)
def create_router(registry: object | None = None) -> APIRouter:
router = APIRouter(prefix="/records", tags=["records"])
records = SqlRecordRegistry(registry)
@router.get("/catalog", response_model=RecordCatalogResponse)
def api_catalog(
session: Session = Depends(get_session),
principal: ApiPrincipal = Depends(get_api_principal),
) -> RecordCatalogResponse:
_require(principal, READ_SCOPE)
return RecordCatalogResponse(**records.catalog(session, principal))
@router.post(
"/catalog/file-plan",
response_model=dict[str, Any],
status_code=status.HTTP_201_CREATED,
)
def api_write_file_plan_node(
payload: FilePlanNodeWriteRequest,
session: Session = Depends(get_session),
principal: ApiPrincipal = Depends(get_api_principal),
) -> dict[str, Any]:
_require(principal, ADMIN_SCOPE)
return _write(
session,
lambda: records.write_file_plan_node(
session, principal, payload=payload.model_dump(mode="python")
),
principal=principal,
operation_type="catalog.file_plan.write",
idempotency_key=payload.idempotency_key,
request=payload.model_dump(mode="json"),
resource_type="record_file_plan_node",
resource_id=payload.node_id,
)
@router.post(
"/catalog/classes",
response_model=dict[str, Any],
status_code=status.HTTP_201_CREATED,
)
def api_write_record_class(
payload: RecordClassWriteRequest,
session: Session = Depends(get_session),
principal: ApiPrincipal = Depends(get_api_principal),
) -> dict[str, Any]:
_require(principal, ADMIN_SCOPE)
return _write(
session,
lambda: records.write_record_class(
session, principal, payload=payload.model_dump(mode="python")
),
principal=principal,
operation_type="catalog.class.write",
idempotency_key=payload.idempotency_key,
request=payload.model_dump(mode="json"),
resource_type="record_class",
resource_id=payload.class_id,
)
@router.get("/sources", response_model=RecordSourceProviderResponse)
def api_source_providers(
session: Session = Depends(get_session),
principal: ApiPrincipal = Depends(get_api_principal),
) -> RecordSourceProviderResponse:
_require(principal, WRITE_SCOPE)
return RecordSourceProviderResponse(
providers=records.source_providers(session, principal)
)
@router.get("/archive-providers", response_model=RecordArchiveProviderResponse)
def api_archive_providers(
session: Session = Depends(get_session),
principal: ApiPrincipal = Depends(get_api_principal),
) -> RecordArchiveProviderResponse:
_require(principal, WRITE_SCOPE)
return RecordArchiveProviderResponse(
providers=records.archive_providers(session, principal)
)
@router.get("", response_model=RecordListResponse)
def api_list_records(
query: str | None = Query(default=None, max_length=500),
record_state: str | None = Query(default=None, alias="state", max_length=40),
class_id: str | None = Query(default=None, max_length=255),
file_plan_node_id: str | None = Query(default=None, max_length=255),
offset: int = Query(default=0, ge=0),
limit: int = Query(default=100, ge=1, le=200),
session: Session = Depends(get_session),
principal: ApiPrincipal = Depends(get_api_principal),
) -> RecordListResponse:
_require(principal, READ_SCOPE)
items, total = records.list_records(
session,
principal,
query=query,
state=record_state,
class_id=class_id,
file_plan_node_id=file_plan_node_id,
offset=offset,
limit=limit,
)
return RecordListResponse(
records=items, total=total, offset=offset, limit=limit
)
@router.post("", response_model=dict[str, Any], status_code=status.HTTP_201_CREATED)
def api_create_record(
payload: RecordCreateRequest,
session: Session = Depends(get_session),
principal: ApiPrincipal = Depends(get_api_principal),
) -> dict[str, Any]:
_require(principal, WRITE_SCOPE)
return _write(
session,
lambda: records.create_record(
session, principal, payload=payload.model_dump(mode="python")
),
principal=principal,
operation_type="record.create",
idempotency_key=payload.idempotency_key,
request=payload.model_dump(mode="json"),
resource_type="record",
resource_id=payload.record_id or payload.record_number,
)
@router.get("/{record_id}", response_model=RecordDetailResponse)
def api_get_record(
record_id: str,
revision: int | None = Query(default=None, ge=1),
session: Session = Depends(get_session),
principal: ApiPrincipal = Depends(get_api_principal),
) -> RecordDetailResponse:
_require(principal, READ_SCOPE)
try:
return RecordDetailResponse(
**records.get_record(
session, principal, record_id=record_id, revision=revision
)
)
except RecordStoreError as exc:
raise _http_error(exc) from exc
@router.get("/{record_id}/recovery", response_model=RecordRecoveryStatusResponse)
def api_record_recovery_status(
record_id: str,
session: Session = Depends(get_session),
principal: ApiPrincipal = Depends(get_api_principal),
) -> RecordRecoveryStatusResponse:
_require(principal, ADMIN_SCOPE)
try:
return RecordRecoveryStatusResponse(
**records.recovery_status(session, principal, record_id=record_id)
)
except RecordStoreError as exc:
raise _http_error(exc) from exc
@router.patch("/{record_id}", response_model=dict[str, Any])
def api_update_record(
record_id: str,
payload: RecordUpdateRequest,
session: Session = Depends(get_session),
principal: ApiPrincipal = Depends(get_api_principal),
) -> dict[str, Any]:
_require(principal, WRITE_SCOPE)
return _write(
session,
lambda: records.update_record(
session,
principal,
record_id=record_id,
payload=payload.model_dump(mode="python", exclude_unset=True),
),
principal=principal,
operation_type="record.revise",
idempotency_key=payload.idempotency_key,
request=payload.model_dump(mode="json", exclude_unset=True),
resource_type="record",
resource_id=record_id,
)
@router.post("/{record_id}/close", response_model=dict[str, Any])
def api_close_record(
record_id: str,
payload: RecordCloseRequest,
session: Session = Depends(get_session),
principal: ApiPrincipal = Depends(get_api_principal),
) -> dict[str, Any]:
_require(principal, WRITE_SCOPE)
if payload.restart_retention:
_require(principal, ADMIN_SCOPE)
return _write(
session,
lambda: records.close_record(
session,
principal,
record_id=record_id,
payload=payload.model_dump(mode="python"),
),
principal=principal,
operation_type="record.close",
idempotency_key=payload.idempotency_key,
request=payload.model_dump(mode="json"),
resource_type="record",
resource_id=record_id,
)
@router.post("/{record_id}/reopen", response_model=dict[str, Any])
def api_reopen_record(
record_id: str,
payload: RecordLifecycleActionRequest,
session: Session = Depends(get_session),
principal: ApiPrincipal = Depends(get_api_principal),
) -> dict[str, Any]:
_require(principal, WRITE_SCOPE)
return _write(
session,
lambda: records.reopen_record(
session,
principal,
record_id=record_id,
payload=payload.model_dump(mode="python"),
),
principal=principal,
operation_type="record.reopen",
idempotency_key=payload.idempotency_key,
request=payload.model_dump(mode="json"),
resource_type="record",
resource_id=record_id,
)
@router.post("/{record_id}/appraise", response_model=dict[str, Any])
def api_appraise_record(
record_id: str,
payload: RecordAppraisalRequest,
session: Session = Depends(get_session),
principal: ApiPrincipal = Depends(get_api_principal),
) -> dict[str, Any]:
_require(principal, WRITE_SCOPE)
if payload.override_retention_not_due:
_require(principal, ADMIN_SCOPE)
return _write(
session,
lambda: records.appraise_record(
session,
principal,
record_id=record_id,
payload=payload.model_dump(mode="python"),
),
principal=principal,
operation_type="record.appraise",
idempotency_key=payload.idempotency_key,
request=payload.model_dump(mode="json"),
resource_type="record",
resource_id=record_id,
)
@router.post(
"/{record_id}/holds",
response_model=dict[str, Any],
status_code=status.HTTP_201_CREATED,
)
def api_apply_hold(
record_id: str,
payload: RecordHoldCreateRequest,
session: Session = Depends(get_session),
principal: ApiPrincipal = Depends(get_api_principal),
) -> dict[str, Any]:
_require(principal, WRITE_SCOPE)
return _write(
session,
lambda: records.apply_hold(
session,
principal,
record_id=record_id,
payload=payload.model_dump(mode="python"),
),
principal=principal,
operation_type="hold.apply",
idempotency_key=payload.idempotency_key,
request=payload.model_dump(mode="json"),
resource_type="record",
resource_id=record_id,
)
@router.post(
"/{record_id}/holds/{hold_id}/release",
response_model=dict[str, Any],
)
def api_release_hold(
record_id: str,
hold_id: str,
payload: RecordHoldReleaseRequest,
session: Session = Depends(get_session),
principal: ApiPrincipal = Depends(get_api_principal),
) -> dict[str, Any]:
_require(principal, WRITE_SCOPE)
return _write(
session,
lambda: records.release_hold(
session,
principal,
record_id=record_id,
hold_id=hold_id,
payload=payload.model_dump(mode="python"),
),
principal=principal,
operation_type="hold.release",
idempotency_key=payload.idempotency_key,
request={"hold_id": hold_id, **payload.model_dump(mode="json")},
resource_type="record",
resource_id=record_id,
)
@router.post(
"/{record_id}/dispositions",
response_model=dict[str, Any],
status_code=status.HTTP_201_CREATED,
)
def api_propose_disposition(
record_id: str,
payload: RecordDispositionCreateRequest,
session: Session = Depends(get_session),
principal: ApiPrincipal = Depends(get_api_principal),
) -> dict[str, Any]:
_require(principal, WRITE_SCOPE)
return _write(
session,
lambda: records.propose_disposition(
session,
principal,
record_id=record_id,
payload=payload.model_dump(mode="python"),
),
principal=principal,
operation_type="disposition.propose",
idempotency_key=payload.idempotency_key,
request=payload.model_dump(mode="json"),
resource_type="record",
resource_id=record_id,
)
@router.post(
"/{record_id}/dispositions/{disposition_id}/finalize",
response_model=dict[str, Any],
)
def api_finalize_disposition(
record_id: str,
disposition_id: str,
payload: RecordDispositionFinalizeRequest,
session: Session = Depends(get_session),
principal: ApiPrincipal = Depends(get_api_principal),
) -> dict[str, Any]:
_require(principal, WRITE_SCOPE)
return _write(
session,
lambda: records.finalize_disposition(
session,
principal,
record_id=record_id,
disposition_id=disposition_id,
payload=payload.model_dump(mode="python"),
),
principal=principal,
operation_type="disposition.finalize",
idempotency_key=payload.idempotency_key,
request={
"disposition_id": disposition_id,
**payload.model_dump(mode="json"),
},
resource_type="record",
resource_id=record_id,
)
@router.post(
"/{record_id}/dispositions/{disposition_id}/withdraw",
response_model=dict[str, Any],
)
def api_withdraw_disposition(
record_id: str,
disposition_id: str,
payload: RecordDispositionWithdrawRequest,
session: Session = Depends(get_session),
principal: ApiPrincipal = Depends(get_api_principal),
) -> dict[str, Any]:
_require(principal, WRITE_SCOPE)
return _write(
session,
lambda: records.withdraw_disposition(
session,
principal,
record_id=record_id,
disposition_id=disposition_id,
payload=payload.model_dump(mode="python"),
),
principal=principal,
operation_type="disposition.withdraw",
idempotency_key=payload.idempotency_key,
request={
"disposition_id": disposition_id,
**payload.model_dump(mode="json"),
},
resource_type="record",
resource_id=record_id,
)
@router.post(
"/{record_id}/transfer-packages",
response_model=dict[str, Any],
status_code=status.HTTP_201_CREATED,
)
def api_prepare_transfer_package(
record_id: str,
payload: RecordTransferPackageCreateRequest,
session: Session = Depends(get_session),
principal: ApiPrincipal = Depends(get_api_principal),
) -> dict[str, Any]:
_require(principal, WRITE_SCOPE)
return _write(
session,
lambda: records.prepare_transfer_package(
session,
principal,
record_id=record_id,
payload=payload.model_dump(mode="python"),
),
principal=principal,
operation_type="transfer.prepare",
idempotency_key=payload.idempotency_key,
request=payload.model_dump(mode="json"),
resource_type="record",
resource_id=record_id,
)
@router.post(
"/{record_id}/transfer-packages/{package_id}/dispatch",
response_model=dict[str, Any],
)
def api_dispatch_transfer_package(
record_id: str,
package_id: str,
payload: RecordTransferDispatchRequest,
session: Session = Depends(get_session),
principal: ApiPrincipal = Depends(get_api_principal),
) -> dict[str, Any]:
_require(principal, WRITE_SCOPE)
return _write(
session,
lambda: records.dispatch_transfer_package(
session,
principal,
record_id=record_id,
package_id=package_id,
payload=payload.model_dump(mode="python"),
),
principal=principal,
operation_type="transfer.simulate",
idempotency_key=payload.idempotency_key,
request={"package_id": package_id, **payload.model_dump(mode="json")},
resource_type="record",
resource_id=record_id,
)
@router.post(
"/{record_id}/volumes",
response_model=dict[str, Any],
status_code=status.HTTP_201_CREATED,
)
def api_create_volume(
record_id: str,
payload: RecordVolumeCreateRequest,
session: Session = Depends(get_session),
principal: ApiPrincipal = Depends(get_api_principal),
) -> dict[str, Any]:
_require(principal, WRITE_SCOPE)
return _write(
session,
lambda: records.create_volume(
session,
principal,
record_id=record_id,
payload=payload.model_dump(mode="python"),
),
principal=principal,
operation_type="volume.create",
idempotency_key=payload.idempotency_key,
request=payload.model_dump(mode="json"),
resource_type="record",
resource_id=record_id,
)
@router.post(
"/{record_id}/items",
response_model=dict[str, Any],
status_code=status.HTTP_201_CREATED,
)
def api_file_item(
record_id: str,
payload: RecordItemCreateRequest,
session: Session = Depends(get_session),
principal: ApiPrincipal = Depends(get_api_principal),
) -> dict[str, Any]:
_require(principal, WRITE_SCOPE)
request = RecordFilingRequest(
tenant_id=principal.tenant_id,
record_id=record_id,
source=RecordSourceLocator(
tenant_id=principal.tenant_id,
source_module=payload.source.source_module,
resource_type=payload.source.resource_type,
resource_id=payload.source.resource_id,
source_revision=payload.source.source_revision,
metadata=payload.source.metadata,
),
purpose=payload.purpose,
filing_reason=payload.filing_reason,
idempotency_key=payload.idempotency_key,
volume_id=payload.volume_id,
relationship=payload.relationship,
institutional_context=payload.institutional_context,
metadata=payload.metadata,
)
def operation() -> dict[str, Any]:
result = records.file(session, principal, request=request)
return {
"record_id": result.record_id,
"item_id": result.item_id,
"sequence": result.sequence,
"filed_at": result.filed_at,
"replayed": result.replayed,
"source": {
"source_module": result.source.locator.source_module,
"resource_type": result.source.locator.resource_type,
"resource_id": result.source.locator.resource_id,
"source_revision": result.source.locator.source_revision,
"label": result.source.label,
},
}
return _write(
session,
operation,
principal=principal,
operation_type="item.file",
idempotency_key=payload.idempotency_key,
request=payload.model_dump(mode="json"),
resource_type="record",
resource_id=record_id,
)
return router
def _require(principal: ApiPrincipal, scope: str) -> None:
if not has_scope(principal, scope):
raise HTTPException(status_code=403, detail=f"Missing scope: {scope}")
def _write(
session: Session,
operation,
*,
principal: ApiPrincipal,
operation_type: str,
idempotency_key: str,
request: dict[str, Any],
resource_type: str,
resource_id: str,
):
recovery = None
try:
recovery = begin_record_atomic_recovery(
session,
tenant_id=principal.tenant_id,
operation_type=operation_type,
idempotency_key=idempotency_key,
request=request,
resource_type=resource_type,
resource_id=resource_id,
)
recovery_token = bind_record_recovery_operation(recovery.operation_id)
try:
result = operation()
finally:
reset_record_recovery_operation(recovery_token)
if not recovery.replayed:
audit_from_principal(
session,
principal,
action=f"records.{operation_type}",
object_type=resource_type,
object_id=resource_id,
details={
"recovery_operation_id": recovery.operation_id,
"idempotency_key_sha256": _sha256(idempotency_key),
},
commit=False,
)
recovery.commit_success(session, result=result, resource_id=resource_id)
return result
except (RecordStoreError, IntegrityError) as exc:
session.rollback()
if recovery is not None:
recovery.reject(summary=str(exc), error_type=type(exc).__name__)
if isinstance(exc, IntegrityError):
raise HTTPException(
status_code=409, detail="The record write conflicts with existing data."
) from exc
raise _http_error(exc) from exc
except RecordRecoveryError as exc:
session.rollback()
raise HTTPException(status_code=503, detail=str(exc)) from exc
except Exception as exc:
session.rollback()
if recovery is not None:
recovery.fail(summary=str(exc), error_type=type(exc).__name__)
raise
def _sha256(value: str) -> str:
return hashlib.sha256(value.encode("utf-8")).hexdigest()
def _http_error(exc: RecordStoreError) -> HTTPException:
if isinstance(exc, RecordNotFoundError):
code = 404
elif isinstance(exc, RecordConflictError):
code = 409
elif isinstance(exc, RecordSourceUnavailableError):
code = 503
else:
code = 422
return HTTPException(status_code=code, detail=str(exc))
__all__ = ["create_router"]