ソースを参照

fix(recovery): stream the re-mirror pass, and gate unique indexes on drift too

Round 5 of the task-12 review.

- The batched pass resolved every touched id (Documents by value) before the
  first commit, so peak memory was O(total bytes replayed) on an input that
  can be an install's whole history: a slow boot traded for an OOM-killed
  one. Each chunkSize window is now resolved, committed and released before
  the next is resolved. Measured on 1000x1MB documents: peak +1253 MB ->
  +511 MB, flat in the number of ids; the batching win is intact (2000 docs
  on ext4: 85.7s per-document vs 0.53s batched).
- CreateIndex(unique) refused only on !mirrorHealthy(). The re-mirror pass
  bumps drift without flipping health (correct for the read gates, which all
  test both), so the gate was open after a boot that left rows stale and a
  duplicate in a stale row could be accepted. Now tests drift too. Audited
  every other mirror-flag site: no other decision uses health without drift.
- remirrorDocuments() now has an outer catch-all, so 'never throws' covers
  bad_alloc from a document copy, not only the faults guarded per row.
- Deleted the single-id remirrorDocument(): zero call sites and a silently
  changed return contract. Documentation rehomed.
- Recorded honestly: no test induces a UniqueViolation from the pass, and the
  pass re-mirrors documents only, never vector sidecars.
fszontagh 1 ヶ月 前
親
コミット
b18809268a

+ 15 - 3
service/src/database_grpc_impl.cpp

@@ -3358,11 +3358,23 @@ grpc::Status DatabaseGrpcImpl::CreateIndex(
     // could miss duplicates that exist in memory and accept a constraint
     // that is already false. Refuse outright rather than declare a unique
     // constraint the data may not actually satisfy.
-    if (request->unique() && !service_.mirrorHealthy()) {
+    //
+    // ⚠ v2.11.0 T12 round-5 — this gate must check DRIFT as well as health,
+    // and used to check health alone. Every read gate in this file tests
+    // (mirrorHealthy() && mirrorDriftCount() == 0), so drift alone was enough
+    // to route reads to MemoryStore; this one call site was not. The
+    // post-replay re-mirror pass (MemoryStore::remirrorDocuments) creates
+    // exactly the missing state: a row it could not write bumps drift and
+    // deliberately does NOT flip health (flipping it would send every read in
+    // the process to MemoryStore for life over one legacy row - the v2.8.1
+    // fault). Health-only here meant that after such a boot, a duplicate
+    // hiding in a stale LMDB row would be accepted silently.
+    if (request->unique() &&
+        (!service_.mirrorHealthy() || service_.mirrorDriftCount() != 0)) {
         response->set_success(false);
         response->set_error("cannot declare a unique constraint while the LMDB mirror is "
-                            "unhealthy: the duplicate check only consults LMDB, and "
-                            "MemoryStore may be ahead of it right now");
+                            "unhealthy or has drifted: the duplicate check only consults "
+                            "LMDB, and MemoryStore may be ahead of it right now");
         return grpc::Status::OK;
     }
 

+ 97 - 60
service/src/memory_store.cpp

@@ -872,14 +872,6 @@ bool MemoryStore::unloadDocument(const std::string& collection, const std::strin
     return true;
 }
 
-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.
-    // 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) {
@@ -914,87 +906,114 @@ MemoryStore::RemirrorBatchResult MemoryStore::remirrorDocuments(
         ++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)});
-    }
-
-    // ---- 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;
+    // ⚠ CATCH-ALL, deliberately (round 5): every individual step below is
+    // already guarded, but "never throws" must be a property of the
+    // function, not a claim about the steps I remembered. Copying a Document
+    // can throw std::bad_alloc, and getCollection()/shared_lock can throw
+    // std::system_error - neither is a per-row fault this pass can attribute,
+    // and either escaping into the unguarded recover() call reproduces
+    // exactly the boot loop this round's finding 1 was about.
+    try {
+
+    // ---- Phase 1+2, STREAMED: resolve one window, commit it, drop it. ----
+    // v2.11.0 T12 round-5 — this used to resolve EVERY id in `docs` into a
+    // per-project map before committing anything, with `Row` holding its
+    // Document BY VALUE. That traded a slow boot for an OOM-killed one: the
+    // input can be an install's entire history under
+    // --recovery-mode=wal_only, and this repo has measured document shapes
+    // at ~2.5-2.9 MB each (anime_images), so 10k distinct ids would have
+    // been tens of GB of document clones resident at once, on the boot path,
+    // before READY. A slow boot completes; a killed one does not.
+    //
+    // Now at most `chunkSize` documents (plus `deferred`, which is small by
+    // construction - only rows that failed their own transaction) are ever
+    // resident. Each window is resolved, committed, and released before the
+    // next window is resolved.
+    //
+    // Parsing happens per row inside the window, NOT inside the transaction,
+    // precisely so a malformed collection key cannot throw from a
+    // transaction other rows share. A window may span several projects (one
+    // env each, so a transaction cannot span them): it is grouped by project
+    // and committed once per project present in that window. Worst case -
+    // every row in a different project - degrades to a transaction per row,
+    // which is exactly the pre-batching cost and never worse.
+    for (size_t windowStart = 0; windowStart < docs.size(); windowStart += chunkSize) {
+        const size_t windowEnd = std::min(windowStart + chunkSize, docs.size());
+
+        std::map<std::string, std::vector<Row>> byProject;
+        for (size_t i = windowStart; i < windowEnd; ++i) {
+            const std::string& collection = docs[i].first;
+            const std::string& id = docs[i].second;
+            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)});
+        }
 
-        for (size_t start = 0; start < rows.size(); start += chunkSize) {
-            const size_t end = std::min(start + chunkSize, rows.size());
+        for (auto& [project, rows] : byProject) {
+            auto* ds = docStoreResolver_(project);
+            if (!ds) continue;
 
-            bool chunkCommitted = false;
-            if (end - start > 1) {
+            bool committed = false;
+            if (rows.size() > 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});
-                }
+                items.reserve(rows.size());
+                for (auto& r : rows) items.push_back({r.bare, r.id, &r.doc});
                 try {
                     ds->put_batch(items);
-                    chunkCommitted = true;
-                    result.remirrored += end - start;
+                    committed = true;
+                    result.remirrored += rows.size();
                 } 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());
+                                 rows.size(), project, e.what());
                 }
             }
 
-            if (chunkCommitted) continue;
+            if (committed) continue;
 
             // Row by row: either the chunk aborted, or it was a chunk of one.
-            for (size_t i = start; i < end; ++i) {
+            for (auto& r : rows) {
                 try {
-                    ds->put(rows[i].bare, rows[i].id, rows[i].doc);
+                    ds->put(r.bare, r.id, r.doc);
                     ++result.remirrored;
                 } catch (const std::exception&) {
-                    deferred.push_back(Deferred{ds, rows[i]});
+                    deferred.push_back(Deferred{ds, std::move(r)});
                 }
             }
         }
+        // byProject (and every Document in it) is released here, before the
+        // next window is resolved. That release is the whole point.
     }
 
     // ---- Phase 3: one bounded retry for rows that failed alone. ---------
@@ -1002,6 +1021,11 @@ MemoryStore::RemirrorBatchResult MemoryStore::remirrorDocuments(
     // 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.
+    //
+    // ⚠ NOT COVERED BY ANY TEST: no test in the suite induces a
+    // UniqueViolation from this pass, so replacing this retry with a bare
+    // noteFailure() would still pass everything. Recorded in the task-12
+    // report rather than papered over.
     for (const auto& d : deferred) {
         try {
             d.ds->put(d.row.bare, d.row.id, d.row.doc);
@@ -1011,6 +1035,19 @@ MemoryStore::RemirrorBatchResult MemoryStore::remirrorDocuments(
         }
     }
 
+    } catch (const std::exception& e) {
+        spdlog::error("post-replay re-mirror pass aborted early: {} - the rows it had "
+                      "not reached yet stay stale in LMDB; recovery continues", e.what());
+        if (mirrorDriftCount_) mirrorDriftCount_->fetch_add(1, std::memory_order_relaxed);
+        ++result.failed;
+    } catch (...) {
+        spdlog::error("post-replay re-mirror pass aborted early (unknown exception) - "
+                      "the rows it had not reached yet stay stale in LMDB; recovery "
+                      "continues");
+        if (mirrorDriftCount_) mirrorDriftCount_->fetch_add(1, std::memory_order_relaxed);
+        ++result.failed;
+    }
+
     return result;
 }
 

+ 35 - 19
service/src/memory_store.hpp

@@ -328,21 +328,25 @@ public:
      * 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). 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.
+     * original write time and needs no repeating).
      *
-     * Safe to call whether or not this specific id actually needed it:
+     * Safe to call whether or not a given 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), if the mirror isn't wired
-     * (legacy/test bootstrap path, same as mirrorWriteToDocStore's own
-     * no-op condition), or if the row's own re-mirror failed — this NEVER
-     * throws, for the reason spelled out on remirrorDocuments().
+     * the same content.
+     *
+     * ⚠ DOCUMENTS ONLY — never vectors. There is no put_vector() equivalent
+     * here, so a replayed write on a collection with vector_dimension > 0
+     * converges its document and leaves the `_vectors_<collection>` sidecar
+     * stale (SimilaritySearch would score the pre-crash vector). Present
+     * since this pass was introduced; not closed, recorded so "replayed
+     * writes converge" is not read as covering vectors.
+     *
+     * v2.11.0 T12 round-5 — the single-id remirrorDocument() that used to
+     * sit here is GONE: recover() moved to the batched form and it had zero
+     * remaining call sites. Its return semantics had also silently changed
+     * (false, not true, for a system collection or a failed put), so keeping
+     * an untested wrapper with a stale contract was worse than deleting it.
      */
-    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;
@@ -354,10 +358,10 @@ public:
     };
 
     /**
-     * v2.11.0 T12 round-4 — batched form of remirrorDocument(), and the one
-     * PersistenceManager::recover() actually calls.
+     * v2.11.0 T12 round-4 — the batched re-mirror pass, and the only entry
+     * point (see the block above for WHY the pass exists at all).
      *
-     * WHY BATCHED: the single-id form routes to LmdbDocumentStore::put(),
+     * WHY BATCHED: a per-document call 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
@@ -365,9 +369,19 @@ public:
      * --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.
+     * instead, via DocumentStore::put_batch() (whose LMDB override holds the
+     * Task 10 beginWrite/put(WriteTxn&)/commit-then-cache discipline), so N
+     * documents in one project cost ceil(N / chunkSize) fsyncs.
+     *
+     * ⚠ STREAMED, not resolved-then-committed (round 5). Documents are
+     * COPIED out of MemoryStore to be written, and `docs` can name an
+     * install's entire history under --recovery-mode=wal_only. Resolving all
+     * of them before the first commit made peak memory O(total bytes
+     * replayed) — tens of GB on this repo's measured multi-MB document
+     * shapes — which traded a slow boot for an OOM-killed one. Each
+     * `chunkSize` window is resolved, committed and RELEASED before the next
+     * is resolved, so peak memory is O(chunkSize documents). Do not hoist
+     * the resolve step back out of the loop.
      *
      * ⚠ 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
@@ -381,7 +395,9 @@ public:
      * 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.
+     * it). Both are caught per row, and the whole body additionally sits
+     * under a catch-all so an allocation failure while copying a document
+     * cannot escape either.
      *
      * 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 —

+ 1 - 1
service/src/persistence/persistence_manager.cpp

@@ -213,7 +213,7 @@ RecoveryOutcome PersistenceManager::recover(MemoryStore& store) {
         // re-mirroring a doc that replay just deleted from MemoryStore would
         // 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
+        // MemoryStore::remirrorDocuments()'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

+ 5 - 4
service/src/persistence/persistence_manager.hpp

@@ -49,10 +49,11 @@ struct RecoveryOutcome {
     size_t snapshotsAvailable = 0;
 
     // v2.11.0 T12 review (I1) — how many distinct (collection, id) pairs
-    // WAL replay applied as UPDATE/UPSERT and then explicitly re-mirrored
-    // to LMDB afterward (MemoryStore::remirrorDocument(), called once per
-    // id from PersistenceManager::recover() — see that method and
-    // remirrorDocument()'s doc comment for why this step exists: WAL
+    // WAL replay applied as INSERT/UPDATE/UPSERT and then explicitly
+    // re-mirrored to LMDB afterward (one batched
+    // MemoryStore::remirrorDocuments() pass from
+    // PersistenceManager::recover() — see that method and
+    // remirrorDocuments()'s doc comment for why this step exists: WAL
     // replay's INSERT/UPDATE path does not mirror 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