fix(recovery): support batched SQLite blob writes

This commit is contained in:
2026-08-19 19:20:05 +02:00
parent 21122e058e
commit ad55d47645
5 changed files with 374 additions and 65 deletions
+293 -53
View File
@@ -1,26 +1,39 @@
from __future__ import annotations
from collections.abc import Callable
from dataclasses import dataclass
from datetime import UTC, datetime
import hashlib
from threading import Lock
from typing import Protocol
from sqlalchemy import event
from sqlalchemy import event, select
from sqlalchemy.orm import Session
from govoplan_core.core.recovery import (
RecoveryGuaranteeError,
RecoveryMode,
RecoveryOperation,
RecoveryPlan,
RecoveryStatus,
plan_recovery_operation,
prepare_recovery_operation,
start_recovery_operation,
verify_recovery_evidence_chain,
)
from govoplan_core.core.recovery_runtime import (
DurableRecoveryOperation,
DurableRecoveryStart,
RecoveryOperationBusy,
RecoveryOperationStateConflict,
begin_durable_recovery_operation,
)
from govoplan_core.core.runtime_coordination import process_runtime_identity
from govoplan_core.core.runtime_coordination import (
RuntimeIdentity,
acquire_lease,
process_runtime_identity,
release_lease,
)
from govoplan_core.db.session import get_database
from govoplan_files.backend.db.models import (
FileBlob,
@@ -38,6 +51,8 @@ from govoplan_files.backend.storage.common import FileStorageError
_PENDING_EFFECTS_KEY = "govoplan_files_pending_recovery_effects"
_HOOKS_INSTALLED_KEY = "govoplan_files_recovery_hooks_installed"
_ROLLBACK_ERRORS_KEY = "govoplan_files_recovery_rollback_errors"
_SQLITE_FENCES: set[str] = set()
_SQLITE_FENCES_LOCK = Lock()
class _PendingEffect(Protocol):
@@ -58,6 +73,8 @@ class PendingBlobWrite:
expected_storage_checksum_sha256: str | None = None
expected_storage_size_bytes: int | None = None
expected_envelope_id: str | None = None
sqlite_fence_key: str | None = None
restart_after_rollback: Callable[[], DurableRecoveryOperation] | None = None
def prepare_stored_bytes(
self,
@@ -72,6 +89,19 @@ class PendingBlobWrite:
self.expected_envelope_id = envelope_id
def settle(self, *, committed: bool) -> None:
if not committed and self.restart_after_rollback is not None:
try:
self.operation = self.restart_after_rollback()
except Exception:
# SQLite cannot make the pre-effect ledger row independent of
# the caller transaction. If reconstructing the evidence row
# also fails, still avoid leaving a known new orphan behind.
if self.created_new:
try:
self.backend.delete(self.storage_key)
except (StorageBackendError, OSError):
pass
raise
evidence = _blob_write_evidence(self)
if _blob_write_complete(evidence):
self.operation.succeed(evidence=evidence)
@@ -219,77 +249,145 @@ def begin_blob_write_recovery(
f"files-blob-{disposition}:{blob_id}:{state_token[:48]}"
)
try:
started = begin_durable_recovery_operation(
get_database().SessionLocal,
identity=process_runtime_identity(),
identity = process_runtime_identity()
except RuntimeError as exc:
raise FileStorageError(
"The Files recovery ledger is unavailable; no object was written"
) from exc
lease_resource_key = f"files:blob:{tenant_id}:{blob_id}"
request = {
"tenant_id": tenant_id,
"blob_id": blob_id,
"storage_key": storage_key if created_new else None,
"storage_key_sha256": key_digest,
"semantic_checksum_sha256": semantic_checksum_sha256,
"semantic_size_bytes": semantic_size_bytes,
"protection_discriminator": protection_discriminator,
"disposition": disposition,
}
recovery_plan = RecoveryPlan(
mode=(
RecoveryMode.COMPENSATION
if created_new
else RecoveryMode.FORWARD_RECOVERY
),
preconditions=(
"the caller has Files write authority for the target owner",
"the storage key belongs to the tenant Files namespace",
"the request records digests rather than file contents",
),
compensation_steps=(
"verify that no FileBlob references the newly reserved key",
"delete only that unreferenced key and verify absence",
)
if created_new
else (),
forward_recovery_steps=(
"verify the expected bytes at the existing blob key",
"forward-complete matching integrity metadata or quarantine the blob",
)
if not created_new
else (),
verification_steps=(
"reload FileBlob metadata through an independent session",
"stream and hash the managed object independently",
),
)
precondition_evidence = {
"blob_id": blob_id,
"storage_key_sha256": key_digest,
"semantic_checksum_sha256": semantic_checksum_sha256,
"semantic_size_bytes": semantic_size_bytes,
"created_new": created_new,
}
metadata = {
"resources": ["postgresql", "object-storage"],
"storage_backend": backend.name,
"durability_mode": "independent_transaction",
}
session_factory = get_database().SessionLocal
sqlite_mode = session.get_bind().dialect.name == "sqlite"
sqlite_fence_key: str | None = None
def start_independent(*, reconstructed_after_rollback: bool = False) -> DurableRecoveryStart:
return begin_durable_recovery_operation(
session_factory,
identity=identity,
module_id="files",
operation_type=f"blob-{disposition}",
idempotency_key=idempotency_key,
request={
"tenant_id": tenant_id,
"blob_id": blob_id,
"storage_key": storage_key if created_new else None,
"storage_key_sha256": key_digest,
"semantic_checksum_sha256": semantic_checksum_sha256,
"semantic_size_bytes": semantic_size_bytes,
"protection_discriminator": protection_discriminator,
"disposition": disposition,
},
recovery_plan=RecoveryPlan(
mode=(
RecoveryMode.COMPENSATION
if created_new
else RecoveryMode.FORWARD_RECOVERY
),
preconditions=(
"the caller has Files write authority for the target owner",
"the storage key belongs to the tenant Files namespace",
"the request records digests rather than file contents",
),
compensation_steps=(
"verify that no FileBlob references the newly reserved key",
"delete only that unreferenced key and verify absence",
)
if created_new
else (),
forward_recovery_steps=(
"verify the expected bytes at the existing blob key",
"forward-complete matching integrity metadata or quarantine the blob",
)
if not created_new
else (),
verification_steps=(
"reload FileBlob metadata through an independent session",
"stream and hash the managed object independently",
),
),
request=request,
recovery_plan=recovery_plan,
precondition_evidence={
"blob_id": blob_id,
"storage_key_sha256": key_digest,
"semantic_checksum_sha256": semantic_checksum_sha256,
"semantic_size_bytes": semantic_size_bytes,
"created_new": created_new,
**precondition_evidence,
**(
{
"sqlite_caller_transaction_rolled_back": True,
"effect_may_have_preceded_durable_intent": True,
}
if reconstructed_after_rollback
else {}
),
},
lease_resource_key=(
f"files:blob:{tenant_id}:{blob_id}"
),
lease_resource_key=lease_resource_key,
lease_ttl_seconds=15 * 60,
resource_type="file_blob",
resource_id=blob_id,
metadata={
"resources": ["postgresql", "object-storage"],
"storage_backend": backend.name,
**metadata,
"durability_mode": (
"sqlite_post_rollback_reconstruction"
if reconstructed_after_rollback
else "independent_transaction"
),
},
)
try:
if sqlite_mode:
_reserve_sqlite_fence(lease_resource_key)
sqlite_fence_key = lease_resource_key
started = _begin_caller_transaction_recovery_operation(
session,
session_factory=session_factory,
identity=identity,
module_id="files",
operation_type=f"blob-{disposition}",
idempotency_key=idempotency_key,
request=request,
recovery_plan=recovery_plan,
precondition_evidence={
**precondition_evidence,
"sqlite_caller_transaction": True,
"reduced_crash_durability": True,
},
lease_resource_key=lease_resource_key,
lease_ttl_seconds=15 * 60,
resource_type="file_blob",
resource_id=blob_id,
metadata={
**metadata,
"resources": ["sqlite", "object-storage"],
"durability_mode": "sqlite_caller_transaction",
},
)
else:
started = start_independent()
except (RecoveryOperationBusy, RecoveryOperationStateConflict) as exc:
if sqlite_fence_key is not None:
_release_sqlite_fence(sqlite_fence_key)
raise FileStorageError(
"This managed blob is already owned by another recovery operation"
) from exc
except (RecoveryGuaranteeError, RuntimeError) as exc:
if sqlite_fence_key is not None:
_release_sqlite_fence(sqlite_fence_key)
raise FileStorageError(
"The Files recovery ledger is unavailable; no object was written"
) from exc
if started.replayed or started.operation is None:
if sqlite_fence_key is not None:
_release_sqlite_fence(sqlite_fence_key)
raise FileStorageError(
"The matching Files blob operation was already completed; reload before retrying"
)
@@ -303,6 +401,14 @@ def begin_blob_write_recovery(
semantic_size_bytes=semantic_size_bytes,
protection_discriminator=protection_discriminator,
created_new=created_new,
sqlite_fence_key=sqlite_fence_key,
restart_after_rollback=(
lambda: _require_started_operation(
start_independent(reconstructed_after_rollback=True)
)
if sqlite_mode
else None
),
)
_register_pending_effect(session, pending)
return pending
@@ -398,10 +504,17 @@ def _register_pending_effect(session: Session, effect: _PendingEffect) -> None:
def _after_session_commit(session: Session) -> None:
# SQLAlchemy also emits after_commit for a released SAVEPOINT. A Files
# batch may open one while creating a later recovery row; settle only when
# the outer business transaction has actually become visible.
if session.in_nested_transaction():
return
_settle_pending_effects(session, committed=True)
def _after_session_rollback(session: Session) -> None:
if session.in_nested_transaction():
return
try:
_settle_pending_effects(session, committed=False)
except RecoveryGuaranteeError as exc:
@@ -422,12 +535,139 @@ def _settle_pending_effects(session: Session, *, committed: bool) -> None:
operation.release_unresolved()
except Exception:
pass
finally:
sqlite_fence_key = getattr(effect, "sqlite_fence_key", None)
if sqlite_fence_key:
_release_sqlite_fence(sqlite_fence_key)
if failures:
raise RecoveryGuaranteeError(
f"{len(failures)} Files recovery operation(s) could not be finalized"
) from failures[0]
def _reserve_sqlite_fence(resource_key: str) -> None:
with _SQLITE_FENCES_LOCK:
if resource_key in _SQLITE_FENCES:
raise RecoveryOperationBusy(
f"Another local SQLite transaction owns {resource_key}"
)
_SQLITE_FENCES.add(resource_key)
def _release_sqlite_fence(resource_key: str) -> None:
with _SQLITE_FENCES_LOCK:
_SQLITE_FENCES.discard(resource_key)
def _require_started_operation(started: DurableRecoveryStart) -> DurableRecoveryOperation:
if started.replayed or started.operation is None:
raise RecoveryOperationStateConflict(started.operation_id, started.status)
return started.operation
def _begin_caller_transaction_recovery_operation(
session: Session,
*,
session_factory: Callable[[], Session],
identity: RuntimeIdentity,
module_id: str,
operation_type: str,
idempotency_key: str,
request: dict[str, object],
recovery_plan: RecoveryPlan,
precondition_evidence: dict[str, object],
lease_resource_key: str,
lease_ttl_seconds: int,
resource_type: str,
resource_id: str,
metadata: dict[str, object],
) -> DurableRecoveryStart:
"""Start SQLite recovery evidence inside the caller transaction.
SQLite permits only one writer, so opening the normal independent ledger
transaction after an earlier member of a batch has written will deadlock
until the busy timeout. The caller-transaction mode preserves fencing,
request hashing, checkpoints, commit verification, and rollback
compensation, but cannot make the intent survive a hard process loss
before the caller commits. PostgreSQL never uses this reduced mode.
"""
claim = acquire_lease(
session,
installation_id=identity.installation_id,
resource_key=lease_resource_key,
holder_node_id=identity.node_id,
holder_incarnation=identity.incarnation,
ttl_seconds=lease_ttl_seconds,
metadata={
"module_id": module_id,
"operation_type": operation_type,
"durability_mode": "sqlite_caller_transaction",
},
)
if claim is None:
raise RecoveryOperationBusy(
f"Another runtime owns the recovery fence for {lease_resource_key}"
)
existing = session.execute(
select(RecoveryOperation).where(
RecoveryOperation.installation_id == identity.installation_id,
RecoveryOperation.module_id == module_id,
RecoveryOperation.idempotency_key == idempotency_key,
)
).scalar_one_or_none()
operation = plan_recovery_operation(
session,
installation_id=identity.installation_id,
module_id=module_id,
operation_type=operation_type,
idempotency_key=idempotency_key,
request=request,
recovery_plan=recovery_plan,
resource_type=resource_type,
resource_id=resource_id,
lease_claim=claim,
metadata=metadata,
)
if existing is not None:
if operation.status == RecoveryStatus.SUCCEEDED.value:
release_lease(session, claim)
return DurableRecoveryStart(
operation_id=operation.id,
status=operation.status,
replayed=True,
operation=None,
)
raise RecoveryOperationStateConflict(operation.id, operation.status)
prepare_recovery_operation(
session,
operation,
evidence=precondition_evidence,
lease_claim=claim,
)
start_recovery_operation(
session,
operation,
evidence={"lease_resource_key": lease_resource_key},
lease_claim=claim,
)
if not verify_recovery_evidence_chain(session, operation.id):
raise RecoveryGuaranteeError(
"Recovery checkpoint chain verification failed before side effects"
)
return DurableRecoveryStart(
operation_id=operation.id,
status=operation.status,
replayed=False,
operation=DurableRecoveryOperation(
session_factory=session_factory,
operation_id=operation.id,
lease_claim=claim,
lease_ttl_seconds=lease_ttl_seconds,
),
)
def _blob_write_evidence(effect: PendingBlobWrite) -> dict[str, object]:
database_blob_present: bool | None
database_matches: bool | None