瀏覽代碼

fix(recovery): guard, batch and generalize the post-replay re-mirror pass

The round-3 re-mirror pass could stop the service from starting and paid one
fsync per document on the boot path.

- Every per-row re-mirror is guarded: a UniqueViolation (reachable because
  the pass runs against a stale index) or a malformed collection key no
  longer escapes recover() and turns a working recovery into a deterministic
  boot loop. Failures log ERROR with collection and id, bump mirror drift
  (not mirror health), and count into RecoveryOutcome::updatesRemirrorFailed.
- Batched: new virtual DocumentStore::put_batch(), overridden by
  LmdbDocumentStore with one WriteTxn / one commit / cache-after. Chunks of
  256, grouped by project. 2000 documents on ext4: 94.0s -> 0.54s. A row
  that throws aborts only its chunk, which is then retried row by row, plus
  one final retry pass that resolves moved-unique-value conflicts.
- WalOpType::INSERT is tracked too, so a delete-then-reinsert inside one
  replay window no longer leaves LMDB with the row deleted.
- Corrected premise: applyDualWriteMirror swallows every LMDB fault except
  UniqueViolation, so a WAL entry never implied its LMDB write succeeded -
  the divergence this pass repairs is not cascade-only.
fszontagh 1 月之前
父節點
當前提交
501202efe5

+ 135 - 14
service/src/memory_store.cpp

@@ -874,23 +874,144 @@ bool MemoryStore::unloadDocument(const std::string& collection, const std::strin
 
 bool MemoryStore::remirrorDocument(const std::string& collection, const std::string& id) {
     // v2.11.0 T12 review (I1) — see the header comment for why this exists.
-    const CollectionData* coll = getCollection(collection);
-    if (!coll) return false;
+    // Delegates to the batched form so there is exactly ONE implementation
+    // of the resolve / skip / catch-per-row rules; a chunk of one is just a
+    // single-document transaction, which is what this used to do anyway.
+    return remirrorDocuments({{collection, id}}, 1).remirrored == 1;
+}
+
+MemoryStore::RemirrorBatchResult MemoryStore::remirrorDocuments(
+    const std::vector<std::pair<std::string, std::string>>& docs,
+    size_t chunkSize) {
+    // v2.11.0 T12 round-4 — see the header comment. This function must never
+    // throw: it runs inside PersistenceManager::recover(), and an escaping
+    // exception there aborts startup permanently.
+    RemirrorBatchResult result;
+    if (docs.empty()) return result;
+    if (chunkSize == 0) chunkSize = 1;
+    if (!docStoreResolver_ || !mirrorHealthy_ || !mirrorDriftCount_) return result;
+
+    // One resolved row: the project decides which env (and therefore which
+    // transaction) it belongs to; `bare` is what the per-project store wants.
+    struct Row {
+        std::string qualified;
+        std::string bare;
+        std::string id;
+        Document doc;
+    };
 
-    std::optional<Document> doc;
-    {
-        std::shared_lock<std::shared_mutex> lock(coll->mutex);
-        auto it = coll->documents.find(id);
-        if (it == coll->documents.end()) return false;
-        doc = it->second;
+    auto noteFailure = [&](const std::string& coll, const std::string& id,
+                           const char* what) {
+        // Drift, not health: a stale row IS drift, and a nonzero drift count
+        // already routes reads to MemoryStore (which is ahead here, and
+        // correct). Deliberately NOT flipping mirrorHealthy_ — one legacy or
+        // conflicting row should not send every read in the process to
+        // MemoryStore for the rest of its life (the v2.8.1 lesson).
+        spdlog::error("post-replay re-mirror failed for collection='{}' id='{}': {} "
+                      "- that row stays stale in LMDB; recovery continues",
+                      coll, id, what);
+        mirrorDriftCount_->fetch_add(1, std::memory_order_relaxed);
+        ++result.failed;
+    };
+
+    // ---- Phase 1: resolve every row, grouped by project. -----------------
+    // Parsing happens HERE, one row at a time, precisely so a malformed
+    // collection key cannot throw from inside a transaction other rows share.
+    std::map<std::string, std::vector<Row>> byProject;
+    for (const auto& [collection, id] : docs) {
+        std::string project;
+        std::string bare;
+        try {
+            auto pc = smartbotic::database::parseProjectCollection(collection);
+            project = std::move(pc.project);
+            bare = std::move(pc.collection);
+        } catch (const std::exception& e) {
+            noteFailure(collection, id, e.what());
+            continue;
+        }
+        // System collections are never mirrored (same rule as
+        // applyDualWriteMirror) — skipped, not a failure.
+        if (bare.empty() || bare[0] == '_') continue;
+
+        const CollectionData* coll = getCollection(collection);
+        if (!coll) continue;
+        std::optional<Document> doc;
+        {
+            std::shared_lock<std::shared_mutex> lock(coll->mutex);
+            auto it = coll->documents.find(id);
+            if (it == coll->documents.end()) continue;  // nothing to mirror
+            doc = it->second;
+        }
+        byProject[project].push_back(Row{collection, std::move(bare), id, std::move(*doc)});
     }
 
-    // No CollectionData mutation here — just reading the current state and
-    // pushing it to LMDB. mirrorWriteToDocStore() is itself a no-op if the
-    // resolver was never wired (legacy/test bootstrap), matching every
-    // other caller's documented behaviour.
-    mirrorWriteToDocStore(collection, id, doc, EventType::UPDATE);
-    return true;
+    // ---- Phase 2: commit in chunks, per project. -------------------------
+    // Rows that fail their own single-row transaction are deferred for one
+    // final retry pass (phase 3).
+    struct Deferred {
+        smartbotic::db::storage::DocumentStore* ds;
+        Row row;
+    };
+    std::vector<Deferred> deferred;
+
+    for (auto& [project, rows] : byProject) {
+        auto* ds = docStoreResolver_(project);
+        if (!ds) continue;
+
+        for (size_t start = 0; start < rows.size(); start += chunkSize) {
+            const size_t end = std::min(start + chunkSize, rows.size());
+
+            bool chunkCommitted = false;
+            if (end - start > 1) {
+                // put_batch() is all-or-nothing (one WriteTxn, one commit,
+                // handles cached only after that commit succeeds - the Task
+                // 10 invariant lives inside it, not here).
+                std::vector<smartbotic::db::storage::DocumentStore::BatchPutItem> items;
+                items.reserve(end - start);
+                for (size_t i = start; i < end; ++i) {
+                    items.push_back({rows[i].bare, rows[i].id, &rows[i].doc});
+                }
+                try {
+                    ds->put_batch(items);
+                    chunkCommitted = true;
+                    result.remirrored += end - start;
+                } catch (const std::exception& e) {
+                    spdlog::warn("post-replay re-mirror: chunk of {} document(s) in "
+                                 "project '{}' aborted ({}); retrying it row by row "
+                                 "so one bad row cannot leave the rest stale",
+                                 end - start, project, e.what());
+                }
+            }
+
+            if (chunkCommitted) continue;
+
+            // Row by row: either the chunk aborted, or it was a chunk of one.
+            for (size_t i = start; i < end; ++i) {
+                try {
+                    ds->put(rows[i].bare, rows[i].id, rows[i].doc);
+                    ++result.remirrored;
+                } catch (const std::exception&) {
+                    deferred.push_back(Deferred{ds, rows[i]});
+                }
+            }
+        }
+    }
+
+    // ---- Phase 3: one bounded retry for rows that failed alone. ---------
+    // This is what resolves the moved-unique-value case: row B holding a
+    // value that used to live on row A conflicts against A's stale posting
+    // until A itself has been re-mirrored. Exactly one extra attempt - no
+    // loop, so a genuinely unresolvable conflict cannot spin the boot path.
+    for (const auto& d : deferred) {
+        try {
+            d.ds->put(d.row.bare, d.row.id, d.row.doc);
+            ++result.remirrored;
+        } catch (const std::exception& e) {
+            noteFailure(d.row.qualified, d.row.id, e.what());
+        }
+    }
+
+    return result;
 }
 
 bool MemoryStore::exists(const std::string& collection, const std::string& id) const {

+ 83 - 11
service/src/memory_store.hpp

@@ -299,11 +299,19 @@ public:
      * entries at boot — do NOT mirror to LMDB (unlike every ordinary write
      * path, where mirrorDocOrUndo()/mirrorWriteToDocStore() always run
      * BEFORE the WAL entry is even logged — see update()/insertWithVector()
-     * above). That asymmetry is harmless for the ordinary write path,
-     * because a WAL entry there can only exist if its LMDB write already
-     * committed successfully first (mirror failure prevents the WAL log
-     * call from ever being reached). It stops being harmless the moment
-     * something logs a WAL entry BEFORE its LMDB write — which
+     * above). Round 4 correction — that asymmetry is NOT harmless for the
+     * ordinary write path either, and the round-3 claim that "a WAL entry
+     * implies the LMDB write already succeeded" is false:
+     * applyDualWriteMirror() SWALLOWS every LMDB fault except
+     * UniqueViolation (log ERROR, bump drift, flip mirror_healthy_, no
+     * rethrow — storage/dual_write_mirror.hpp), so mirrorDocOrUndo()
+     * returns normally and emitPersist() logs the WAL entry anyway. A WAL
+     * entry therefore only implies "the mirror did not throw
+     * UniqueViolation". So the divergence this pass repairs is NOT
+     * cascade-only: it also occurs on any ordinary write whose mirror fault
+     * was swallowed, and re-mirroring at boot repairs those too. The
+     * divergence merely becomes RELIABLY reachable the moment something
+     * logs a WAL entry BEFORE its LMDB write — which
      * relations/relation_cascade.cpp's cascade delete deliberately does
      * (WAL-first, for its own crash-safety reasons: see that file's
      * header). A crash between that WAL write and the LMDB commit leaves a
@@ -313,24 +321,88 @@ public:
      * loadDocumentWithHistory() do not, and LMDB then permanently serves
      * the pre-cascade row until something else happens to rewrite it.
      *
-     * PersistenceManager::recover() calls this once per DISTINCT id that
-     * WAL replay applied as UPDATE/UPSERT (deduplicated, and skipped if a
-     * later DELETE for that id was also replayed — remove() already
+     * PersistenceManager::recover() re-mirrors once per DISTINCT id that
+     * WAL replay applied as INSERT/UPDATE/UPSERT (deduplicated, and skipped
+     * if a later DELETE for that id was also replayed — remove() already
      * mirrored that) — a targeted post-replay pass, not a mirror call on
      * every replayed entry, which would multiply LMDB writes by however
      * many times a hot document was rewritten since the last snapshot for
      * no benefit (the ordinary case's LMDB write already happened at
-     * original write time and needs no repeating).
+     * original write time and needs no repeating). It calls the BATCHED
+     * remirrorDocuments() below, not this single-id form — see that
+     * method's comment for why one-transaction-per-document on the boot
+     * path is not acceptable.
      *
      * Safe to call whether or not this specific id actually needed it:
      * re-mirroring an already-correct row is an idempotent overwrite with
      * the same content. Returns false (no-op) if the id is not currently
-     * in MemoryStore (nothing to mirror) or the mirror isn't wired
+     * in MemoryStore (nothing to mirror), if the mirror isn't wired
      * (legacy/test bootstrap path, same as mirrorWriteToDocStore's own
-     * no-op condition).
+     * no-op condition), or if the row's own re-mirror failed — this NEVER
+     * throws, for the reason spelled out on remirrorDocuments().
      */
     bool remirrorDocument(const std::string& collection, const std::string& id);
 
+    struct RemirrorBatchResult {
+        // Rows whose LMDB write committed as part of this call.
+        uint64_t remirrored = 0;
+        // Rows this call could not re-mirror. Each one is logged at ERROR
+        // naming its collection and id, and bumps the mirror drift counter.
+        // A nonzero value means "those rows stay stale in LMDB", never
+        // "recovery failed" — see below.
+        uint64_t failed = 0;
+    };
+
+    /**
+     * v2.11.0 T12 round-4 — batched form of remirrorDocument(), and the one
+     * PersistenceManager::recover() actually calls.
+     *
+     * WHY BATCHED: the single-id form routes to LmdbDocumentStore::put(),
+     * which opens and commits its OWN WriteTxn. The env is opened without
+     * MDB_NOSYNC, so that is one fsync per document, on the boot path,
+     * before sd_notify(READY=1). WAL size is bounded only by
+     * snapshotIntervalSec (3600) and maxWalSizeMb (100), and
+     * --recovery-mode=wal_only can replay the entire history — tens of
+     * thousands of distinct ids on a busy install, against a 10-minute
+     * systemd start watchdog. This commits in chunks of `chunkSize`
+     * instead, using the Task 10 primitives (beginWrite() /
+     * put(WriteTxn&, …, to_cache) / commitAndCache()), so N documents cost
+     * ceil(N / chunkSize) fsyncs.
+     *
+     * ⚠ NEVER THROWS, and that is the whole point of its error handling. A
+     * re-mirror is a REPAIR pass: failing one row must degrade to "that row
+     * stays stale in LMDB", which is exactly the state the pass exists to
+     * improve on and is strictly better than refusing to boot. An escaping
+     * exception here would propagate out of recover() and turn a working
+     * recovery into a deterministic boot loop on data that booted fine
+     * before. Two throws are genuinely reachable: UniqueViolation (the pass
+     * runs against a STALE index, so a unique value that MOVED between two
+     * rows conflicts with the other row's not-yet-re-mirrored posting) and
+     * std::invalid_argument from parseProjectCollection() on a malformed or
+     * legacy collection key (WAL replay itself never parses collection
+     * keys, so such a key replays fine and only this pass would trip on
+     * it). Both are caught per row.
+     *
+     * ONE BAD ROW MUST NOT POISON ITS CHUNK: a throw from row N aborts that
+     * chunk's transaction (nothing in it committed, and NOTHING is cached —
+     * the Task 10 invariant: to_cache is applied only after a successful
+     * commit), so the whole chunk is then retried ROW BY ROW, each in its
+     * own transaction. The rows that can commit do; only the genuinely bad
+     * one is counted as failed. Rows that fail individually get ONE further
+     * retry pass after every other row has been re-mirrored, which is what
+     * actually resolves the moved-unique-value case: once the row that used
+     * to hold the value has been re-mirrored, its stale posting is gone and
+     * the row that now holds it commits.
+     *
+     * Rows whose collection is missing from MemoryStore, whose id is absent,
+     * or whose collection is `_`-prefixed (system collections are not
+     * mirrored — same rule as applyDualWriteMirror) are silently skipped:
+     * not remirrored, not failed.
+     */
+    RemirrorBatchResult remirrorDocuments(
+        const std::vector<std::pair<std::string, std::string>>& docs,
+        size_t chunkSize = 256);
+
     /**
      * Check if a document exists.
      */

+ 53 - 14
service/src/persistence/persistence_manager.cpp

@@ -206,42 +206,81 @@ RecoveryOutcome PersistenceManager::recover(MemoryStore& store) {
             return false;
         }
         // v2.11.0 T12 review (I1) — track every (collection, id) applied as
-        // UPDATE/UPSERT during this replay, deduplicated, so it can be
+        // INSERT/UPDATE/UPSERT during this replay, deduplicated, so it can be
         // explicitly re-mirrored to LMDB once replay finishes. Erased again
         // on a later DELETE for the same id: remove() (DELETE's own replay
         // path) already mirrors, so there is nothing left to re-mirror, and
         // re-mirroring a doc that replay just deleted from MemoryStore would
-        // be wrong (remirrorDocument() would correctly no-op on a missing
-        // id, but skipping it here avoids the pointless lookup). See
+        // be wrong (remirrorDocuments() would correctly skip a missing id,
+        // but dropping it here avoids the pointless lookup). See
         // MemoryStore::remirrorDocument()'s doc comment for why this step
         // exists at all: WAL replay's INSERT/UPDATE path never mirrors on
         // its own, unlike every ordinary write, and
         // relations/relation_cascade.cpp's WAL-first cascade can leave a
         // replayed UPDATE whose LMDB write never actually happened.
-        std::unordered_map<std::string, std::pair<std::string, std::string>> touchedByUpdate;
+        //
+        // INSERT is tracked too (round 4): loadDocument() does not mirror
+        // either, so a DELETE-then-reinsert of the same id inside one replay
+        // window used to leave LMDB with the row deleted — the DELETE
+        // mirrored, the reinsert did not. Tracking it costs nothing, and it
+        // makes "replayed writes converge" true generally rather than for
+        // two of the three op types.
+        //
+        // Insertion-ORDERED (a vector plus a dedup index, not just a map):
+        // the order rows are re-mirrored in decides which of two rows wins a
+        // transient unique-index conflict, and an unordered_map made that
+        // non-deterministic between boots. WAL order is the order the writes
+        // originally happened in, which is the least surprising choice.
+        std::vector<std::pair<std::string, std::string>> touched;
+        std::unordered_map<std::string, size_t> touchedIndex;
+        std::vector<bool> touchedAlive;
         uint64_t replayed = wal_->replay(fromSequence, [&](const WalEntry& entry) {
             applyWalEntry(store, entry);
             const std::string key = entry.collection + "\x1f" + entry.documentId;
-            if (entry.opType == WalOpType::UPDATE || entry.opType == WalOpType::UPSERT) {
-                touchedByUpdate[key] = {entry.collection, entry.documentId};
+            if (entry.opType == WalOpType::INSERT || entry.opType == WalOpType::UPDATE ||
+                entry.opType == WalOpType::UPSERT) {
+                auto it = touchedIndex.find(key);
+                if (it == touchedIndex.end()) {
+                    touchedIndex.emplace(key, touched.size());
+                    touched.emplace_back(entry.collection, entry.documentId);
+                    touchedAlive.push_back(true);
+                } else {
+                    touchedAlive[it->second] = true;
+                }
             } else if (entry.opType == WalOpType::DELETE) {
-                touchedByUpdate.erase(key);
+                auto it = touchedIndex.find(key);
+                if (it != touchedIndex.end()) touchedAlive[it->second] = false;
             }
         });
         outcome.walEntriesReplayed = replayed;
         spdlog::info("Replayed {} WAL entries from sequence {}", replayed, fromSequence);
 
-        uint64_t remirrored = 0;
-        for (const auto& [key, collAndId] : touchedByUpdate) {
-            (void)key;
-            if (store.remirrorDocument(collAndId.first, collAndId.second)) ++remirrored;
+        std::vector<std::pair<std::string, std::string>> toRemirror;
+        toRemirror.reserve(touched.size());
+        for (size_t i = 0; i < touched.size(); ++i) {
+            if (touchedAlive[i]) toRemirror.push_back(touched[i]);
         }
-        outcome.updatesRemirroredAfterReplay = remirrored;
-        if (remirrored > 0) {
+        // BATCHED, and it never throws — both deliberate, see
+        // MemoryStore::remirrorDocuments(). Per-document transactions here
+        // meant one fsync per replayed id on the boot path before
+        // sd_notify(READY=1); an escaping exception here meant recover()
+        // failed and the service refused to start, permanently, on data that
+        // booted fine before.
+        auto remirror = store.remirrorDocuments(toRemirror);
+        outcome.updatesRemirroredAfterReplay = remirror.remirrored;
+        outcome.updatesRemirrorFailed = remirror.failed;
+        if (remirror.remirrored > 0) {
             spdlog::info("Re-mirrored {} document(s) to LMDB after WAL replay "
                         "(closes the window where a WAL-first writer, e.g. a "
                         "cascade delete, logged an UPDATE whose LMDB write "
-                        "never ran before a crash)", remirrored);
+                        "never ran before a crash)", remirror.remirrored);
+        }
+        if (remirror.failed > 0) {
+            spdlog::error("{} document(s) could not be re-mirrored to LMDB after WAL "
+                          "replay (each logged above with its collection and id). "
+                          "Those rows stay stale in LMDB; recovery is NOT failed by "
+                          "this - a repair pass that cannot repair one row must not "
+                          "stop the service from starting.", remirror.failed);
         }
         // Anchor the WAL's sequence_ counter to the snapshot's walSequence.
         // Without this, the steady-state outcome of `truncateBefore(walSeq)`

+ 9 - 0
service/src/persistence/persistence_manager.hpp

@@ -62,6 +62,15 @@ struct RecoveryOutcome {
     // idempotent) re-mirror.
     uint64_t updatesRemirroredAfterReplay = 0;
 
+    // v2.11.0 T12 round-4 — rows the post-replay re-mirror pass could NOT
+    // write to LMDB (a malformed/legacy collection key, or a unique-index
+    // conflict that survived the retry pass). Each is logged at ERROR with
+    // its collection and id and bumps the mirror drift counter. Nonzero
+    // means "those rows stay stale in LMDB", NOT "recovery failed": a
+    // repair pass that cannot repair one row must never stop the service
+    // from starting, which is what an unguarded throw here used to do.
+    uint64_t updatesRemirrorFailed = 0;
+
     bool isNonTrivial() const {
         return kind != Kind::TrivialSuccess && kind != Kind::FreshInstall;
     }

+ 16 - 8
service/src/relations/relation_cascade.hpp

@@ -51,14 +51,22 @@
 // non-null) — reads are LMDB-first, so this is user-visible, not just an
 // internal inconsistency — plus a live reverse-index posting pointing at a
 // now-deleted parent, invisible until that child is next rewritten through
-// the ordinary write path (which re-mirrors it). This is a real, currently
-// UNCLOSED gap — recorded here rather than fixed because closing it needs
-// either extending loadDocumentWithHistory() to mirror (a MemoryStore
-// change with implications far beyond this module) or a dedicated
-// re-mirror pass keyed off the mutations WAL replay just applied, and both
-// are more than this task's remaining budget covers. Flagged for whoever
-// picks this up next; do not assume LMDB is exempt from the crash window
-// just because MemoryStore is.
+// the ordinary write path (which re-mirrors it).
+//
+// CLOSED (round 3, refined in round 4) by a post-replay re-mirror pass in
+// PersistenceManager::recover(): every DISTINCT (collection, id) that replay
+// applied as INSERT/UPDATE/UPSERT, minus any a later DELETE removed, is
+// pushed to LMDB via MemoryStore::remirrorDocuments() once replay finishes.
+// It is BATCHED (chunked transactions — one fsync per document on the boot
+// path was minutes of startup on a large WAL) and it NEVER THROWS (an
+// escaping exception from a repair pass turned recover() into a refusal to
+// start). It is deliberately not cascade-specific: WAL entries carry no
+// "came from a cascade" marker, and the same divergence occurs whenever
+// applyDualWriteMirror() SWALLOWED an ordinary write's LMDB fault (it
+// swallows everything but UniqueViolation), so a WAL entry never did imply
+// its LMDB write succeeded. Do not assume LMDB is exempt from the crash
+// window just because MemoryStore is — it converges at the NEXT boot, not
+// at the moment of the crash.
 //
 // ⚠ NOT RECURSIVE, AND WHAT THAT MEANS FOR `restrict` (review finding I2):
 // planCascade() only ever looks at relations whose parent is the document

+ 29 - 0
service/src/storage/document_store.hpp

@@ -46,6 +46,35 @@ public:
                      std::string_view id,
                      const smartbotic::database::Document& doc) = 0;
 
+    // v2.11.0 T12 round-4 — put several documents (possibly across several
+    // collections) as ONE atomic unit.
+    //
+    // Contract: all-or-nothing. Either every item is written, or the call
+    // throws and NOTHING it touched was written. Callers that need
+    // per-row failure isolation retry the batch through single-item put()
+    // calls after catching — which is what MemoryStore::remirrorDocuments()
+    // does, so one bad row cannot leave a whole chunk stale.
+    //
+    // Why it exists: the post-WAL-replay re-mirror pass runs on the boot
+    // path before sd_notify(READY=1), and put() commits its own transaction
+    // (one fsync) per document. On a crash boot with a large WAL, or under
+    // --recovery-mode=wal_only, per-document commits are minutes of startup.
+    //
+    // `items` holds NON-OWNING Document pointers - they must outlive the
+    // call. The default implementation is a non-atomic loop over put(), for
+    // backends (and test doubles) with no transaction concept of their own;
+    // LmdbDocumentStore overrides it with a single WriteTxn.
+    struct BatchPutItem {
+        std::string_view collection;
+        std::string_view id;
+        const smartbotic::database::Document* doc;
+    };
+    virtual void put_batch(const std::vector<BatchPutItem>& items) {
+        for (const auto& it : items) {
+            if (it.doc) put(it.collection, it.id, *it.doc);
+        }
+    }
+
     // Read a document. Returns nullopt if absent or collection doesn't exist.
     virtual std::optional<smartbotic::database::Document>
     get(std::string_view collection, std::string_view id) = 0;

+ 17 - 0
service/src/storage/document_store_lmdb.cpp

@@ -530,6 +530,23 @@ void LmdbDocumentStore::put(std::string_view collection,
     for (const auto& [sub, d] : to_cache) cacheCommittedDbi(sub, d);
 }
 
+void LmdbDocumentStore::put_batch(const std::vector<BatchPutItem>& items) {
+    if (items.empty()) return;
+    WriteTxn wtxn(env_);
+    std::vector<std::pair<std::string, unsigned int>> to_cache;
+    for (const auto& item : items) {
+        if (!item.doc) continue;
+        put(wtxn, item.collection, item.id, *item.doc, to_cache);
+    }
+    // Commit BEFORE any handle is cached - the v2.8.0 lesson. If any put()
+    // above threw, wtxn's destructor aborts everything and to_cache is
+    // discarded unapplied, which is exactly right: LMDB closes handles whose
+    // opening transaction aborted, so caching one would poison the sub-db for
+    // the life of the process.
+    wtxn.commit();
+    for (const auto& [sub, d] : to_cache) cacheCommittedDbi(sub, d);
+}
+
 void LmdbDocumentStore::put(WriteTxn& wtxn,
                              std::string_view collection,
                              std::string_view id,

+ 7 - 0
service/src/storage/document_store_lmdb.hpp

@@ -309,6 +309,13 @@ public:
              std::string_view id,
              const smartbotic::database::Document& doc) override;
 
+    // v2.11.0 T12 round-4 — see DocumentStore::put_batch. One WriteTxn for
+    // the whole vector, one commit, then (and only then) the accumulated
+    // handles are cached: the same commit-before-cache discipline as every
+    // other wrapper in this file. Throws without having written anything if
+    // any item fails.
+    void put_batch(const std::vector<BatchPutItem>& items) override;
+
     std::optional<smartbotic::database::Document>
     get(std::string_view collection, std::string_view id) override;
 

+ 139 - 0
tests/test_relation_enforcement.cpp

@@ -1083,6 +1083,143 @@ void test_replayed_cascade_update_remirrors_to_lmdb_after_crash_window() {
     mstore.stop();
 }
 
+// v2.11.0 T12 round-4 — the post-replay re-mirror pass must not be able to
+// stop the service from starting. Two throws are genuinely reachable from it
+// (see MemoryStore::remirrorDocuments): std::invalid_argument from
+// parseProjectCollection() on a malformed/legacy collection key, and a
+// storage fault from the put itself. Unguarded, either escaped recover(),
+// which DatabaseService::initialize() turns into "refuse to start" - a
+// deterministic boot loop on data that booted fine before.
+//
+// This also pins the batching contract: a row that throws inside a chunk
+// aborts that chunk's transaction, and the OTHER rows of that chunk must
+// still converge via the row-by-row retry.
+void test_a_failing_row_does_not_stop_recovery() {
+    TmpEnv t("rel-remirror-fail");
+    LmdbDocumentStore store(t.env);
+    TmpPersistence p("rel-remirror-fail-wal");
+    check(p.pm.start(), "persistence manager started");
+
+    // Row 1: an ordinary, well-formed document. Must converge.
+    Document good; good.id = "w1"; good.collection = "widgets";
+    good.set_data({{"v", 1}});
+    p.pm.logInsert("default:widgets", good);
+    good.set_data({{"v", 2}});
+    p.pm.logUpdate("default:widgets", good);
+
+    // Row 2: same project, so it lands in the SAME chunk as row 1 - but its
+    // id is past LMDB's 511-byte key limit, so the put throws
+    // (MDB_BAD_VALSIZE) from inside that chunk's transaction. This is the
+    // honest way to induce a mid-chunk throw, same technique as
+    // test_aborted_write_does_not_poison_the_collection (v2.8.0).
+    Document oversize; oversize.id = std::string(600, 'k'); oversize.collection = "widgets";
+    oversize.set_data({{"v", 1}});
+    p.pm.logInsert("default:widgets", oversize);
+    p.pm.logUpdate("default:widgets", oversize);
+
+    // Row 3: a malformed collection key. WAL replay itself never parses
+    // collection keys, so this replays fine and only the re-mirror pass
+    // trips on it - previously with an uncaught std::invalid_argument.
+    Document weird; weird.id = "x1"; weird.collection = "legacy";
+    weird.set_data({{"v", 1}});
+    p.pm.logInsert("weird:legacy:key", weird);
+    p.pm.logUpdate("weird:legacy:key", weird);
+    p.pm.stop();
+
+    MemoryStore freshStore(MemoryStore::Config{});
+    freshStore.start();
+    std::atomic<bool> mirrorHealthy{true};
+    std::atomic<uint64_t> mirrorDrift{0};
+    freshStore.setDocumentStoreMirror(
+        [&store](std::string_view) -> smartbotic::db::storage::DocumentStore* { return &store; },
+        &mirrorHealthy, &mirrorDrift);
+
+    PersistenceManager::Config cfg2;
+    cfg2.dataDir = p.path;
+    PersistenceManager pm2(cfg2);
+    auto outcome = pm2.recover(freshStore);
+
+    // THE finding: recovery completes.
+    check(outcome.kind != smartbotic::database::RecoveryOutcome::Kind::Failed,
+          "recovery completed despite two un-mirrorable rows - no boot loop");
+    check(outcome.walEntriesReplayed == 6,
+          "all six WAL entries replayed - the failures did not truncate replay");
+
+    // The failures are counted and reported, not swallowed silently.
+    check(outcome.updatesRemirrorFailed == 2,
+          "both un-mirrorable rows were counted as failed");
+    check(mirrorDrift.load() == 2,
+          "each failed row bumped mirror drift (the operator-visible signal)");
+    check(mirrorHealthy.load(),
+          "mirror health NOT flipped - one legacy row must not send every read "
+          "in the process to MemoryStore for the rest of its life");
+
+    // The other row in the same chunk still converged.
+    check(outcome.updatesRemirroredAfterReplay == 1, "the good row was re-mirrored");
+    auto lmdbGood = store.get("widgets", "w1");
+    check(lmdbGood.has_value(), "the good row reached LMDB despite sharing a chunk with a bad one");
+    if (lmdbGood) {
+        check(lmdbGood->data()["v"] == 2, "and it carries the REPLAYED update, not the insert");
+    }
+
+    // MemoryStore replay itself was unaffected for every row, including the
+    // ones LMDB could not take.
+    check(freshStore.get("default:widgets", "w1").has_value(), "w1 in MemoryStore");
+    check(freshStore.get("weird:legacy:key", "x1").has_value(),
+          "the malformed-key row still recovered into MemoryStore - it is only "
+          "LMDB that cannot hold it");
+
+    freshStore.stop();
+}
+
+// v2.11.0 T12 round-4 (point 3) — INSERT must be tracked by the re-mirror
+// pass too. loadDocument() (INSERT replay) does not mirror, so a
+// delete-then-reinsert of the same id inside one replay window used to leave
+// LMDB with the row DELETED: the DELETE mirrored via remove(), the reinsert
+// did not.
+void test_reinsert_after_delete_converges_in_lmdb() {
+    TmpEnv t("rel-remirror-reinsert");
+    LmdbDocumentStore store(t.env);
+    TmpPersistence p("rel-remirror-reinsert-wal");
+    check(p.pm.start(), "persistence manager started");
+
+    Document d; d.id = "r1"; d.collection = "widgets";
+    d.set_data({{"v", 1}});
+    store.put("widgets", "r1", d);           // LMDB starts in the pre-window state
+    p.pm.logInsert("default:widgets", d);
+    p.pm.logDelete("default:widgets", "r1");
+    d.set_data({{"v", 99}});                 // reinserted with new content
+    p.pm.logInsert("default:widgets", d);
+    p.pm.stop();
+
+    MemoryStore freshStore(MemoryStore::Config{});
+    freshStore.start();
+    std::atomic<bool> mirrorHealthy{true};
+    std::atomic<uint64_t> mirrorDrift{0};
+    freshStore.setDocumentStoreMirror(
+        [&store](std::string_view) -> smartbotic::db::storage::DocumentStore* { return &store; },
+        &mirrorHealthy, &mirrorDrift);
+
+    PersistenceManager::Config cfg2;
+    cfg2.dataDir = p.path;
+    PersistenceManager pm2(cfg2);
+    auto outcome = pm2.recover(freshStore);
+    check(outcome.kind != smartbotic::database::RecoveryOutcome::Kind::Failed, "recovery completed");
+    check(outcome.updatesRemirroredAfterReplay == 1,
+          "the reinserted id was re-mirrored (INSERT is tracked, not just UPDATE/UPSERT)");
+
+    auto mem = freshStore.get("default:widgets", "r1");
+    check(mem.has_value(), "reinserted row present in MemoryStore");
+
+    // LOAD-BEARING: without INSERT tracking the DELETE's own mirror wins and
+    // LMDB has nothing here at all.
+    auto lmdb = store.get("widgets", "r1");
+    check(lmdb.has_value(), "reinserted row present in LMDB - the DELETE's mirror did not win");
+    if (lmdb) check(lmdb->data()["v"] == 99, "and it is the REINSERTED content, not the original");
+
+    freshStore.stop();
+}
+
 }  // namespace
 
 int main() {
@@ -1102,6 +1239,8 @@ int main() {
     test_crash_between_wal_and_commit_recovers();
     test_cascade_refuses_when_grandchild_is_restrict_protected();
     test_replayed_cascade_update_remirrors_to_lmdb_after_crash_window();
+    test_a_failing_row_does_not_stop_recovery();
+    test_reinsert_after_delete_converges_in_lmdb();
 
     std::cout << "passed: " << g_pass << ", failed: " << g_fail << "\n";
     return g_fail == 0 ? 0 : 1;