Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
11 changes: 11 additions & 0 deletions src/superlocalmemory/core/transactions/erasure.py
Original file line number Diff line number Diff line change
Expand Up @@ -352,6 +352,17 @@ def finalize(
persisted=persisted,
)

def prove_erased(
self, context: OperationContext, owner: str,
) -> ErasureProofRecord:
"""Read-only re-proof that one owner's projection holds no residue.

Used by the background redrive to complete erase obligations whose
purge already happened. Never writes: callers decide whether and how
to record the proof.
"""
return self._prove_owner(context, owner)

def _prove_owner(
self, context: OperationContext, name: str,
) -> ErasureProofRecord:
Expand Down
155 changes: 148 additions & 7 deletions src/superlocalmemory/server/unified_daemon.py
Original file line number Diff line number Diff line change
Expand Up @@ -6296,6 +6296,7 @@ def _reconcile_pending_projections(
build_transaction_service,
)
from superlocalmemory.core.transactions.obligations import ObligationLedger
from superlocalmemory.core.transactions.owners import ObligationKind

db = getattr(engine, "_db", None)
profile_id = getattr(engine, "_profile_id", None)
Expand All @@ -6317,12 +6318,22 @@ def _reconcile_pending_projections(
done = 0
for operation_id in sorted(op_ids):
try:
context = _context_for_operation(engine, operation_id)
if context is None:
_terminalize_orphan_operation(engine, operation_id)
kinds = _pending_obligation_kinds(db, operation_id)
if not kinds:
# Terminal (or already-driven) obligations surface here via
# operations_missing_manifest, which has no state filter.
# Nothing pending means nothing to do; skipping keeps one
# completed op from occupying a redrive slot every pass.
continue
service.reconcile_operation(db, context)
done += 1
if ObligationKind.ERASE in kinds:
done += _reconcile_erase_operation(engine, db, ledger, operation_id)
if kinds - {ObligationKind.ERASE}:
context = _context_for_operation(engine, operation_id)
if context is None:
_terminalize_orphan_operation(engine, operation_id)
continue
service.reconcile_operation(db, context)
done += 1
except Exception as exc: # noqa: BLE001
logger.warning(
"projection redrive failed for %s: %s", operation_id, exc,
Expand All @@ -6333,9 +6344,134 @@ def _reconcile_pending_projections(
return 0


def _pending_obligation_kinds(db, operation_id: str) -> set[str]:
"""Kinds with non-terminal obligations for one operation.

The redrive fans out by kind: ``apply`` belongs to the ingestion
reconciler, ``erase`` to the erasure reconciler below. Reading kinds
first is what stops an erasure id (no ingestion row by design) from
being misread as an orphaned ingestion.
"""
from superlocalmemory.core.transactions.owners import ObligationState

try:
rows = db.execute(
"SELECT DISTINCT kind FROM projection_obligations "
"WHERE operation_id = ? AND state NOT IN (?, ?)",
(operation_id, str(ObligationState.VERIFIED), str(ObligationState.ERASED)),
)
except Exception as exc: # noqa: BLE001
# Fail open for visibility: an empty set skips this op for one pass
# only (it stays pending and is retried in 30s), but must be loud.
logger.warning(
"erase/apply kind lookup failed for %s: %s", operation_id, exc,
)
return set()
return {str(dict(r).get("kind")) for r in rows}


def _reconcile_erase_operation(engine, db, ledger, operation_id: str) -> int:
"""Drive one erase-kind obligation set toward ERASED without failing it.

Read-only re-proof: an obligation is marked ERASED (no attempt bump)
only when its fact is canonically absent, its tombstone is present, and
its owner proves no residue. Anything else is left pending for the
erasure service — this path never marks FAILED and never bumps attempts,
so a slow purge cannot be mistaken for an orphan.
"""
from superlocalmemory.core.transactions.concrete_owners import (
build_erasure_service,
)
from superlocalmemory.core.transactions.owners import (
ObligationKind,
ObligationState,
)

with db.raw_connection() as conn:
pending = [
o for o in ledger.fetch(conn, operation_id)
if o.kind == ObligationKind.ERASE
and o.state not in (ObligationState.VERIFIED, ObligationState.ERASED)
]
if not pending:
return 0
service = build_erasure_service(engine)
done = 0
# Group by (profile, subject): one erase request normally names one fact,
# but entity/profile erasures fan out. A group that cannot be proven is
# left pending without affecting the groups that can.
groups: dict[tuple[str, str], list] = {}
for ob in pending:
groups.setdefault((ob.profile_id, ob.subject_id), []).append(ob)
for (profile_id, subject_id), obs in groups.items():
done += _reconcile_erase_group(
db, ledger, service, operation_id, profile_id, subject_id, obs,
)
return 1 if done == len(groups) and groups else 0


def _reconcile_erase_group(
db, ledger, service, operation_id: str,
profile_id: str, subject_id: str, obs: list,
) -> int:
"""Prove and close one erase group. Returns 1 iff the group completed."""
from superlocalmemory.core.transactions.erasure import is_tombstoned
from superlocalmemory.core.transactions.owners import (
ObligationKind,
ObligationState,
OperationContext,
)

# NOTE: fact-shaped erasures only. Entity/profile erasures name subjects
# that are not fact_ids and carry multi-fact contexts this redrive cannot
# reconstruct; the tombstone gate below fails those closed (pending).
if db.execute(
"SELECT 1 FROM atomic_facts WHERE fact_id = ? AND profile_id = ? LIMIT 1",
(subject_id, profile_id),
):
return 0
with db.raw_connection() as conn:
if not is_tombstoned(conn, profile_id, subject_id):
return 0
context = OperationContext(
operation_id=operation_id,
profile_id=profile_id,
subject_id=subject_id,
fact_ids=(subject_id,),
)
# Collect proofs before opening the write transaction: proofs are
# read-only, and a mark must never interleave with a later proof read.
proven: list[tuple[str, str | None]] = []
for ob in obs:
try:
proof = service.prove_erased(context, ob.owner)
except Exception: # noqa: BLE001
return 0
if not proof.erased:
return 0
proven.append((ob.owner, proof.checksum or None))
with db.raw_connection() as conn:
for owner, checksum in proven:
ledger.mark(
conn, operation_id, owner, ObligationKind.ERASE,
ObligationState.ERASED,
checksum=checksum,
detail={"phase": "redrive-proven"},
)
# Close the loop for the missing-manifest feed: erase ops never had
# a manifest writer, so without this they would be re-fetched every
# 30s forever and crowd out ingestion redrives.
from superlocalmemory.core.transactions.reconciler import Reconciler

Reconciler(ledger).reconcile(
conn, operation_id, profile_id, canonical_committed=False,
)
return 1


def _terminalize_orphan_operation(engine, operation_id: str) -> None:
from superlocalmemory.core.transactions.obligations import ObligationLedger
from superlocalmemory.core.transactions.owners import ObligationState
from superlocalmemory.core.transactions.owners import ObligationKind, ObligationState
from superlocalmemory.core.transactions.reconciler import Reconciler
from superlocalmemory.core.transactions.service import MAX_APPLY_ATTEMPTS

Expand All @@ -6346,12 +6482,17 @@ def _terminalize_orphan_operation(engine, operation_id: str) -> None:
ledger = ObligationLedger()
with db.raw_connection() as conn:
for obligation in ledger.fetch(conn, operation_id):
if obligation.kind != ObligationKind.APPLY:
# Erase obligations have no ingestion row by design; the
# erase reconciler above owns them. Failing them here as
# orphans is the bug this guard removes.
continue
if obligation.attempts >= MAX_APPLY_ATTEMPTS:
continue
ledger.mark(
conn, operation_id, obligation.owner, obligation.kind,
ObligationState.FAILED,
detail={"phase": "orphan", "error": "canonical record missing"},
detail={"phase": "orphan", "error": "ingestion record missing"},
bump_attempts=True,
)
Reconciler(ledger).reconcile(
Expand Down
102 changes: 102 additions & 0 deletions tests/core/test_projection_spine_integration.py
Original file line number Diff line number Diff line change
Expand Up @@ -250,6 +250,108 @@ def test_redrive_reconciles_orphaned_obligations(stored_engine) -> None:
assert manifest["state"] == ManifestState.DEGRADED.value, manifest["state"]


def test_redrive_does_not_orphan_erase_obligations(stored_engine) -> None:
"""Erase obligations without an ingestion row must never be orphan-FAILED.

Regression for the dashboard-DELETE incident: deleting a fact creates
kind=erase obligations whose operation_id is an erasure id, so there is
no ingestion_operations row. The 30s redrive fed them to the
ingestion-only _context_for_operation and _terminalize_orphan_operation
FAILED them with 'canonical record missing' until exhausted.
"""
from superlocalmemory.server.unified_daemon import (
_reconcile_pending_projections,
)

engine, _operation = stored_engine
op_id, subject = "erase-op-regression-1", "fact-erased-regression-1"
_seed_erase_op(engine, op_id, subject, tombstone=True)

_reconcile_pending_projections(engine, force=True)

by_owner = _erase_states(engine, op_id)
assert set(by_owner) == {"bm25", "temporal", "vector"}
for owner, row in by_owner.items():
assert row["state"] == "erased", f"{owner}: {row}"
assert row["attempts"] == 0, f"{owner}: {row}"
manifest = engine._db.execute(
"SELECT state FROM completion_manifests WHERE operation_id = ?",
(op_id,),
)
assert manifest, "erase close must write a manifest (missing-manifest feed)"


def _seed_erase_op(engine, op_id: str, subject: str, *, tombstone: bool) -> None:
import time as _time

profile_id = engine._profile_id
now = _time.time()
with engine._db.raw_connection() as conn:
if tombstone:
conn.execute(
"INSERT INTO projection_tombstones "
"(profile_id, fact_id, erasure_id, created_at) VALUES (?, ?, ?, ?)",
(profile_id, subject, op_id, now),
)
for owner in ("bm25", "temporal", "vector"):
conn.execute(
"INSERT INTO projection_obligations "
"(operation_id, profile_id, owner, kind, subject_id, state, "
"attempts, created_at, updated_at) "
"VALUES (?, ?, ?, 'erase', ?, 'pending', 0, ?, ?)",
(op_id, profile_id, owner, subject, now, now),
)


def _erase_states(engine, op_id: str) -> dict:
rows = engine._db.execute(
"SELECT owner, state, attempts FROM projection_obligations "
"WHERE operation_id = ?",
(op_id,),
)
return {dict(r)["owner"]: dict(r) for r in rows}


def test_redrive_leaves_erase_pending_when_residue_remains(stored_engine) -> None:
"""A bm25 row for the subject must block ERASED (no vacuous proof)."""
from superlocalmemory.server.unified_daemon import (
_reconcile_pending_projections,
)

engine, _operation = stored_engine
op_id, subject = "erase-op-residue-1", "fact-erased-residue-1"
_seed_erase_op(engine, op_id, subject, tombstone=True)
with engine._db.raw_connection() as conn:
conn.execute(
"INSERT INTO bm25_tokens (fact_id, profile_id, tokens) VALUES (?, ?, '[]')",
(subject, engine._profile_id),
)

_reconcile_pending_projections(engine, force=True)

for owner, row in _erase_states(engine, op_id).items():
assert row["state"] == "pending", f"{owner}: {row}"
assert row["attempts"] == 0, f"{owner}: {row}"


def test_redrive_leaves_erase_pending_without_tombstone(stored_engine) -> None:
"""No tombstone means the delete never ran: never mark ERASED."""
from superlocalmemory.server.unified_daemon import (
_reconcile_pending_projections,
)

engine, _operation = stored_engine
op_id, subject = "erase-op-notomb-1", "fact-erased-notomb-1"
_seed_erase_op(engine, op_id, subject, tombstone=False)

_reconcile_pending_projections(engine, force=True)
_reconcile_pending_projections(engine, force=True)

for owner, row in _erase_states(engine, op_id).items():
assert row["state"] == "pending", f"{owner}: {row}"
assert row["attempts"] == 0, f"{owner}: {row}"


def test_manifest_is_reverifiable(stored_engine) -> None:
from superlocalmemory.core.transactions import Reconciler
from superlocalmemory.server.unified_daemon import (
Expand Down