|
@@ -5,6 +5,7 @@
|
|
|
#include "persistence/history_store.hpp"
|
|
#include "persistence/history_store.hpp"
|
|
|
#include "project_addressing.hpp"
|
|
#include "project_addressing.hpp"
|
|
|
#include "storage/document_store.hpp"
|
|
#include "storage/document_store.hpp"
|
|
|
|
|
+#include "storage/document_store_lmdb.hpp"
|
|
|
#include "storage/dual_write_mirror.hpp"
|
|
#include "storage/dual_write_mirror.hpp"
|
|
|
#include "storage/filter_eval.hpp"
|
|
#include "storage/filter_eval.hpp"
|
|
|
|
|
|
|
@@ -364,7 +365,18 @@ std::string MemoryStore::insert(const std::string& collection, Document doc) {
|
|
|
uint64_t docSize = estimateDocumentSize(doc);
|
|
uint64_t docSize = estimateDocumentSize(doc);
|
|
|
|
|
|
|
|
// v2.0 dual-write under lock — see header comment on mirrorWriteToDocStore.
|
|
// v2.0 dual-write under lock — see header comment on mirrorWriteToDocStore.
|
|
|
- mirrorWriteToDocStore(collection, docId, doc, EventType::INSERT);
|
|
|
|
|
|
|
+ // v2.11.0 T11 — routed through mirrorDocOrUndo: a UniqueViolation must
|
|
|
|
|
+ // undo the map insert, the vector, and the expiration index entry added
|
|
|
|
|
+ // above, so a rejected insert leaves no trace of ever having happened.
|
|
|
|
|
+ mirrorDocOrUndo(collection, docId, doc, EventType::INSERT, [&]() {
|
|
|
|
|
+ if (doc.expiresAt > 0) {
|
|
|
|
|
+ removeFromExpirationIndex(*coll, docId, doc.expiresAt);
|
|
|
|
|
+ }
|
|
|
|
|
+ if (!vec.empty()) {
|
|
|
|
|
+ removeVector(*coll, docId);
|
|
|
|
|
+ }
|
|
|
|
|
+ coll->documents.erase(docId);
|
|
|
|
|
+ });
|
|
|
if (!vec.empty()) {
|
|
if (!vec.empty()) {
|
|
|
mirrorVectorToDocStore(collection, docId, &vec, EventType::INSERT);
|
|
mirrorVectorToDocStore(collection, docId, &vec, EventType::INSERT);
|
|
|
}
|
|
}
|
|
@@ -488,6 +500,15 @@ bool MemoryStore::update(const std::string& collection, const std::string& id, c
|
|
|
// Track memory change (old size)
|
|
// Track memory change (old size)
|
|
|
uint64_t oldSize = estimateDocumentSize(it->second);
|
|
uint64_t oldSize = estimateDocumentSize(it->second);
|
|
|
|
|
|
|
|
|
|
+ // v2.11.0 T11 — snapshot enough of the pre-write state to undo, in case
|
|
|
|
|
+ // the mirror rejects this as a UniqueViolation. Captured before anything
|
|
|
|
|
+ // is mutated: the document itself, and whether/what vector it held.
|
|
|
|
|
+ const Document original = it->second;
|
|
|
|
|
+ std::optional<std::vector<float>> originalVec;
|
|
|
|
|
+ if (auto vit = coll->vectors.find(id); vit != coll->vectors.end()) {
|
|
|
|
|
+ originalVec = vit->second;
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
// Remove from old expiration index
|
|
// Remove from old expiration index
|
|
|
if (it->second.expiresAt > 0) {
|
|
if (it->second.expiresAt > 0) {
|
|
|
removeFromExpirationIndex(*coll, id, it->second.expiresAt);
|
|
removeFromExpirationIndex(*coll, id, it->second.expiresAt);
|
|
@@ -546,7 +567,23 @@ bool MemoryStore::update(const std::string& collection, const std::string& id, c
|
|
|
uint64_t newSize = estimateDocumentSize(updated);
|
|
uint64_t newSize = estimateDocumentSize(updated);
|
|
|
|
|
|
|
|
// v2.0 dual-write under lock
|
|
// v2.0 dual-write under lock
|
|
|
- mirrorWriteToDocStore(collection, id, updated, EventType::UPDATE);
|
|
|
|
|
|
|
+ // v2.11.0 T11 — a UniqueViolation here must put the document, its
|
|
|
|
|
+ // vector, and the expiration index back exactly as `original` had them,
|
|
|
|
|
+ // so a rejected update leaves the row as if the update never happened.
|
|
|
|
|
+ mirrorDocOrUndo(collection, id, updated, EventType::UPDATE, [&]() {
|
|
|
|
|
+ if (updated.expiresAt > 0) {
|
|
|
|
|
+ removeFromExpirationIndex(*coll, id, updated.expiresAt);
|
|
|
|
|
+ }
|
|
|
|
|
+ it->second = original;
|
|
|
|
|
+ if (original.expiresAt > 0) {
|
|
|
|
|
+ addToExpirationIndex(*coll, id, original.expiresAt);
|
|
|
|
|
+ }
|
|
|
|
|
+ if (originalVec.has_value()) {
|
|
|
|
|
+ storeVector(*coll, id, *originalVec);
|
|
|
|
|
+ } else {
|
|
|
|
|
+ removeVector(*coll, id);
|
|
|
|
|
+ }
|
|
|
|
|
+ });
|
|
|
if (!vec.empty()) {
|
|
if (!vec.empty()) {
|
|
|
mirrorVectorToDocStore(collection, id, &vec, EventType::UPDATE);
|
|
mirrorVectorToDocStore(collection, id, &vec, EventType::UPDATE);
|
|
|
} else {
|
|
} else {
|
|
@@ -608,6 +645,19 @@ std::string MemoryStore::upsert(const std::string& collection, Document doc) {
|
|
|
uint64_t oldSize = 0;
|
|
uint64_t oldSize = 0;
|
|
|
|
|
|
|
|
auto it = coll->documents.find(docId);
|
|
auto it = coll->documents.find(docId);
|
|
|
|
|
+
|
|
|
|
|
+ // v2.11.0 T11 — snapshot pre-write state for undo before either branch
|
|
|
|
|
+ // mutates anything. Only meaningful on the update branch (isInsert=false
|
|
|
|
|
+ // undo just erases what it added), but cheap enough to always capture.
|
|
|
|
|
+ const bool hadExisting = it != coll->documents.end();
|
|
|
|
|
+ const Document originalDoc = hadExisting ? it->second : Document{};
|
|
|
|
|
+ std::optional<std::vector<float>> originalVec;
|
|
|
|
|
+ if (hadExisting) {
|
|
|
|
|
+ if (auto vit = coll->vectors.find(docId); vit != coll->vectors.end()) {
|
|
|
|
|
+ originalVec = vit->second;
|
|
|
|
|
+ }
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
if (it == coll->documents.end()) {
|
|
if (it == coll->documents.end()) {
|
|
|
// Insert
|
|
// Insert
|
|
|
isInsert = true;
|
|
isInsert = true;
|
|
@@ -664,8 +714,29 @@ std::string MemoryStore::upsert(const std::string& collection, Document doc) {
|
|
|
uint64_t newSize = estimateDocumentSize(doc);
|
|
uint64_t newSize = estimateDocumentSize(doc);
|
|
|
|
|
|
|
|
// v2.0 dual-write under lock (upsert path — INSERT or UPDATE depending on isInsert)
|
|
// v2.0 dual-write under lock (upsert path — INSERT or UPDATE depending on isInsert)
|
|
|
- mirrorWriteToDocStore(collection, docId, doc,
|
|
|
|
|
- isInsert ? EventType::INSERT : EventType::UPDATE);
|
|
|
|
|
|
|
+ // v2.11.0 T11 — on UniqueViolation, undo whichever branch above ran: an
|
|
|
|
|
+ // insert erases the row it added, an update restores the row (and its
|
|
|
|
|
+ // vector/expiration index) to the pre-upsert snapshot.
|
|
|
|
|
+ mirrorDocOrUndo(collection, docId, doc,
|
|
|
|
|
+ isInsert ? EventType::INSERT : EventType::UPDATE, [&]() {
|
|
|
|
|
+ if (doc.expiresAt > 0) {
|
|
|
|
|
+ removeFromExpirationIndex(*coll, docId, doc.expiresAt);
|
|
|
|
|
+ }
|
|
|
|
|
+ if (isInsert) {
|
|
|
|
|
+ if (!vec.empty()) removeVector(*coll, docId);
|
|
|
|
|
+ coll->documents.erase(docId);
|
|
|
|
|
+ } else {
|
|
|
|
|
+ coll->documents[docId] = originalDoc;
|
|
|
|
|
+ if (originalDoc.expiresAt > 0) {
|
|
|
|
|
+ addToExpirationIndex(*coll, docId, originalDoc.expiresAt);
|
|
|
|
|
+ }
|
|
|
|
|
+ if (originalVec.has_value()) {
|
|
|
|
|
+ storeVector(*coll, docId, *originalVec);
|
|
|
|
|
+ } else {
|
|
|
|
|
+ removeVector(*coll, docId);
|
|
|
|
|
+ }
|
|
|
|
|
+ }
|
|
|
|
|
+ });
|
|
|
if (!vec.empty()) {
|
|
if (!vec.empty()) {
|
|
|
mirrorVectorToDocStore(collection, docId, &vec,
|
|
mirrorVectorToDocStore(collection, docId, &vec,
|
|
|
isInsert ? EventType::INSERT : EventType::UPDATE);
|
|
isInsert ? EventType::INSERT : EventType::UPDATE);
|
|
@@ -792,6 +863,13 @@ bool MemoryStore::updateIfVersion(const std::string& collection, const std::stri
|
|
|
// Track memory change (old size)
|
|
// Track memory change (old size)
|
|
|
uint64_t oldSize = estimateDocumentSize(it->second);
|
|
uint64_t oldSize = estimateDocumentSize(it->second);
|
|
|
|
|
|
|
|
|
|
+ // v2.11.0 T11 — snapshot for undo, same reasoning as update() above.
|
|
|
|
|
+ const Document original = it->second;
|
|
|
|
|
+ std::optional<std::vector<float>> originalVec;
|
|
|
|
|
+ if (auto vit = coll->vectors.find(id); vit != coll->vectors.end()) {
|
|
|
|
|
+ originalVec = vit->second;
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
// Remove from old expiration index
|
|
// Remove from old expiration index
|
|
|
if (it->second.expiresAt > 0) {
|
|
if (it->second.expiresAt > 0) {
|
|
|
removeFromExpirationIndex(*coll, id, it->second.expiresAt);
|
|
removeFromExpirationIndex(*coll, id, it->second.expiresAt);
|
|
@@ -858,7 +936,22 @@ bool MemoryStore::updateIfVersion(const std::string& collection, const std::stri
|
|
|
uint64_t newSize = estimateDocumentSize(updated);
|
|
uint64_t newSize = estimateDocumentSize(updated);
|
|
|
|
|
|
|
|
// v2.0 dual-write under lock — both doc body and vector.
|
|
// v2.0 dual-write under lock — both doc body and vector.
|
|
|
- mirrorWriteToDocStore(collection, id, updated, EventType::UPDATE);
|
|
|
|
|
|
|
+ // v2.11.0 T11 — same undo shape as update(): restore doc/vector/expiry
|
|
|
|
|
+ // to the pre-write snapshot on a rejected write.
|
|
|
|
|
+ mirrorDocOrUndo(collection, id, updated, EventType::UPDATE, [&]() {
|
|
|
|
|
+ if (updated.expiresAt > 0) {
|
|
|
|
|
+ removeFromExpirationIndex(*coll, id, updated.expiresAt);
|
|
|
|
|
+ }
|
|
|
|
|
+ it->second = original;
|
|
|
|
|
+ if (original.expiresAt > 0) {
|
|
|
|
|
+ addToExpirationIndex(*coll, id, original.expiresAt);
|
|
|
|
|
+ }
|
|
|
|
|
+ if (originalVec.has_value()) {
|
|
|
|
|
+ storeVector(*coll, id, *originalVec);
|
|
|
|
|
+ } else {
|
|
|
|
|
+ removeVector(*coll, id);
|
|
|
|
|
+ }
|
|
|
|
|
+ });
|
|
|
if (!vec.empty()) {
|
|
if (!vec.empty()) {
|
|
|
mirrorVectorToDocStore(collection, id, &vec, EventType::UPDATE);
|
|
mirrorVectorToDocStore(collection, id, &vec, EventType::UPDATE);
|
|
|
} else {
|
|
} else {
|
|
@@ -905,6 +998,15 @@ uint64_t MemoryStore::patchDocument(const std::string& collection, const std::st
|
|
|
// Track memory change (old size)
|
|
// Track memory change (old size)
|
|
|
uint64_t oldSize = estimateDocumentSize(it->second);
|
|
uint64_t oldSize = estimateDocumentSize(it->second);
|
|
|
|
|
|
|
|
|
|
+ // v2.11.0 T11 — snapshot for undo. patchDocument mutates it->second in
|
|
|
|
|
+ // place rather than building a separate "updated" object, so the undo
|
|
|
|
|
+ // needs the whole pre-patch document plus whatever vector it held.
|
|
|
|
|
+ const Document original = it->second;
|
|
|
|
|
+ std::optional<std::vector<float>> originalVec;
|
|
|
|
|
+ if (auto vit = coll->vectors.find(id); vit != coll->vectors.end()) {
|
|
|
|
|
+ originalVec = vit->second;
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
// Remove from old expiration index
|
|
// Remove from old expiration index
|
|
|
if (it->second.expiresAt > 0) {
|
|
if (it->second.expiresAt > 0) {
|
|
|
removeFromExpirationIndex(*coll, id, it->second.expiresAt);
|
|
removeFromExpirationIndex(*coll, id, it->second.expiresAt);
|
|
@@ -959,7 +1061,26 @@ uint64_t MemoryStore::patchDocument(const std::string& collection, const std::st
|
|
|
// v2.0 dual-write under lock (patch path — vector only mirrored when
|
|
// v2.0 dual-write under lock (patch path — vector only mirrored when
|
|
|
// the patch actually supplied a new _vector; otherwise the existing
|
|
// the patch actually supplied a new _vector; otherwise the existing
|
|
|
// vector stays in place on both stores)
|
|
// vector stays in place on both stores)
|
|
|
- mirrorWriteToDocStore(collection, id, updatedDoc, EventType::UPDATE);
|
|
|
|
|
|
|
+ // v2.11.0 T11 — undo restores the whole pre-patch document, then
|
|
|
|
|
+ // re-adds the original expiration index entry (removing whatever entry
|
|
|
|
|
+ // the patch attempt left behind first). Vector is only touched back if
|
|
|
|
|
+ // the patch attempt itself touched it.
|
|
|
|
|
+ mirrorDocOrUndo(collection, id, updatedDoc, EventType::UPDATE, [&]() {
|
|
|
|
|
+ if (it->second.expiresAt > 0) {
|
|
|
|
|
+ removeFromExpirationIndex(*coll, id, it->second.expiresAt);
|
|
|
|
|
+ }
|
|
|
|
|
+ it->second = original;
|
|
|
|
|
+ if (original.expiresAt > 0) {
|
|
|
|
|
+ addToExpirationIndex(*coll, id, original.expiresAt);
|
|
|
|
|
+ }
|
|
|
|
|
+ if (!vec.empty()) {
|
|
|
|
|
+ if (originalVec.has_value()) {
|
|
|
|
|
+ storeVector(*coll, id, *originalVec);
|
|
|
|
|
+ } else {
|
|
|
|
|
+ removeVector(*coll, id);
|
|
|
|
|
+ }
|
|
|
|
|
+ }
|
|
|
|
|
+ });
|
|
|
if (!vec.empty()) {
|
|
if (!vec.empty()) {
|
|
|
mirrorVectorToDocStore(collection, id, &vec, EventType::UPDATE);
|
|
mirrorVectorToDocStore(collection, id, &vec, EventType::UPDATE);
|
|
|
}
|
|
}
|
|
@@ -1224,7 +1345,13 @@ bool MemoryStore::setAdd(const std::string& collection, const std::string& setId
|
|
|
coll->updatedAt = doc.updatedAt;
|
|
coll->updatedAt = doc.updatedAt;
|
|
|
|
|
|
|
|
// v2.0 dual-write under lock
|
|
// v2.0 dual-write under lock
|
|
|
- mirrorWriteToDocStore(collection, setId, doc, EventType::INSERT);
|
|
|
|
|
|
|
+ // v2.11.0 T11 — undo the just-created set document on rejection.
|
|
|
|
|
+ mirrorDocOrUndo(collection, setId, doc, EventType::INSERT, [&]() {
|
|
|
|
|
+ if (doc.expiresAt > 0) {
|
|
|
|
|
+ removeFromExpirationIndex(*coll, setId, doc.expiresAt);
|
|
|
|
|
+ }
|
|
|
|
|
+ coll->documents.erase(setId);
|
|
|
|
|
+ });
|
|
|
|
|
|
|
|
lock.unlock();
|
|
lock.unlock();
|
|
|
|
|
|
|
@@ -1255,6 +1382,8 @@ bool MemoryStore::setAdd(const std::string& collection, const std::string& setId
|
|
|
}
|
|
}
|
|
|
}
|
|
}
|
|
|
|
|
|
|
|
|
|
+ const Document original = it->second;
|
|
|
|
|
+
|
|
|
saveToHistory(*coll, it->second);
|
|
saveToHistory(*coll, it->second);
|
|
|
members.push_back(member);
|
|
members.push_back(member);
|
|
|
it->second.set_data(tree);
|
|
it->second.set_data(tree);
|
|
@@ -1266,7 +1395,10 @@ bool MemoryStore::setAdd(const std::string& collection, const std::string& setId
|
|
|
Document updated = it->second;
|
|
Document updated = it->second;
|
|
|
|
|
|
|
|
// v2.0 dual-write under lock (sadd-style: append member, update doc)
|
|
// v2.0 dual-write under lock (sadd-style: append member, update doc)
|
|
|
- mirrorWriteToDocStore(collection, setId, updated, EventType::UPDATE);
|
|
|
|
|
|
|
+ // v2.11.0 T11 — undo restores the document to its pre-add state.
|
|
|
|
|
+ mirrorDocOrUndo(collection, setId, updated, EventType::UPDATE, [&]() {
|
|
|
|
|
+ it->second = original;
|
|
|
|
|
+ });
|
|
|
|
|
|
|
|
lock.unlock();
|
|
lock.unlock();
|
|
|
|
|
|
|
@@ -1977,6 +2109,9 @@ uint64_t MemoryStore::restoreToVersion(const std::string& collection, const std:
|
|
|
uint64_t newVersion = 0;
|
|
uint64_t newVersion = 0;
|
|
|
Document restoredDoc;
|
|
Document restoredDoc;
|
|
|
|
|
|
|
|
|
|
+ // v2.11.0 T11 — snapshot for undo on the "active document" branch.
|
|
|
|
|
+ const Document originalIfActive = wasDeleted ? Document{} : docIt->second;
|
|
|
|
|
+
|
|
|
if (!wasDeleted) {
|
|
if (!wasDeleted) {
|
|
|
// Active document: save current state to history, then overwrite
|
|
// Active document: save current state to history, then overwrite
|
|
|
saveToHistory(*coll, docIt->second);
|
|
saveToHistory(*coll, docIt->second);
|
|
@@ -2027,8 +2162,23 @@ uint64_t MemoryStore::restoreToVersion(const std::string& collection, const std:
|
|
|
|
|
|
|
|
// v2.0 dual-write under lock (restore path — INSERT if doc was previously
|
|
// v2.0 dual-write under lock (restore path — INSERT if doc was previously
|
|
|
// deleted/missing, UPDATE if we replaced an existing doc).
|
|
// deleted/missing, UPDATE if we replaced an existing doc).
|
|
|
- mirrorWriteToDocStore(collection, id, restoredDoc,
|
|
|
|
|
- wasDeleted ? EventType::INSERT : EventType::UPDATE);
|
|
|
|
|
|
|
+ // v2.11.0 T11 — undo depends on which branch ran: a restore-of-deleted
|
|
|
|
|
+ // erases the row it recreated, a restore-over-active puts the document
|
|
|
|
|
+ // back exactly as it was before the restore.
|
|
|
|
|
+ mirrorDocOrUndo(collection, id, restoredDoc,
|
|
|
|
|
+ wasDeleted ? EventType::INSERT : EventType::UPDATE, [&]() {
|
|
|
|
|
+ if (restoredDoc.expiresAt > 0) {
|
|
|
|
|
+ removeFromExpirationIndex(*coll, id, restoredDoc.expiresAt);
|
|
|
|
|
+ }
|
|
|
|
|
+ if (wasDeleted) {
|
|
|
|
|
+ coll->documents.erase(id);
|
|
|
|
|
+ } else {
|
|
|
|
|
+ docIt->second = originalIfActive;
|
|
|
|
|
+ if (originalIfActive.expiresAt > 0) {
|
|
|
|
|
+ addToExpirationIndex(*coll, id, originalIfActive.expiresAt);
|
|
|
|
|
+ }
|
|
|
|
|
+ }
|
|
|
|
|
+ });
|
|
|
|
|
|
|
|
lock.unlock();
|
|
lock.unlock();
|
|
|
|
|
|
|
@@ -2436,6 +2586,17 @@ void MemoryStore::mirrorWriteToDocStore(const std::string& collection, const std
|
|
|
pc.collection, id, doc, eventType);
|
|
pc.collection, id, doc, eventType);
|
|
|
}
|
|
}
|
|
|
|
|
|
|
|
|
|
+void MemoryStore::mirrorDocOrUndo(const std::string& collection, const std::string& id,
|
|
|
|
|
+ const Document& doc, EventType eventType,
|
|
|
|
|
+ const std::function<void()>& undo) {
|
|
|
|
|
+ try {
|
|
|
|
|
+ mirrorWriteToDocStore(collection, id, doc, eventType);
|
|
|
|
|
+ } catch (const smartbotic::db::storage::UniqueViolation&) {
|
|
|
|
|
+ undo();
|
|
|
|
|
+ throw;
|
|
|
|
|
+ }
|
|
|
|
|
+}
|
|
|
|
|
+
|
|
|
void MemoryStore::mirrorVectorToDocStore(const std::string& collection, const std::string& id,
|
|
void MemoryStore::mirrorVectorToDocStore(const std::string& collection, const std::string& id,
|
|
|
const std::vector<float>* vec, EventType eventType) {
|
|
const std::vector<float>* vec, EventType eventType) {
|
|
|
if (!docStoreResolver_ || !mirrorHealthy_ || !mirrorDriftCount_) return;
|
|
if (!docStoreResolver_ || !mirrorHealthy_ || !mirrorDriftCount_) return;
|