| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691692693694695696697698699700701702703704705706707708709710 |
- #include "wal.hpp"
- #include "../json_parse.hpp"
- #include <algorithm>
- #include <chrono>
- #include <cstring>
- #include <iomanip>
- #include <sstream>
- namespace smartbotic::database {
- // ===== CRC32 Implementation =====
- namespace crc32 {
- // CRC32 lookup table (polynomial 0xEDB88320)
- static const uint32_t table[256] = {
- 0x00000000, 0x77073096, 0xEE0E612C, 0x990951BA, 0x076DC419, 0x706AF48F, 0xE963A535, 0x9E6495A3,
- 0x0EDB8832, 0x79DCB8A4, 0xE0D5E91E, 0x97D2D988, 0x09B64C2B, 0x7EB17CBD, 0xE7B82D07, 0x90BF1D91,
- 0x1DB71064, 0x6AB020F2, 0xF3B97148, 0x84BE41DE, 0x1ADAD47D, 0x6DDDE4EB, 0xF4D4B551, 0x83D385C7,
- 0x136C9856, 0x646BA8C0, 0xFD62F97A, 0x8A65C9EC, 0x14015C4F, 0x63066CD9, 0xFA0F3D63, 0x8D080DF5,
- 0x3B6E20C8, 0x4C69105E, 0xD56041E4, 0xA2677172, 0x3C03E4D1, 0x4B04D447, 0xD20D85FD, 0xA50AB56B,
- 0x35B5A8FA, 0x42B2986C, 0xDBBBC9D6, 0xACBCF940, 0x32D86CE3, 0x45DF5C75, 0xDCD60DCF, 0xABD13D59,
- 0x26D930AC, 0x51DE003A, 0xC8D75180, 0xBFD06116, 0x21B4F4B5, 0x56B3C423, 0xCFBA9599, 0xB8BDA50F,
- 0x2802B89E, 0x5F058808, 0xC60CD9B2, 0xB10BE924, 0x2F6F7C87, 0x58684C11, 0xC1611DAB, 0xB6662D3D,
- 0x76DC4190, 0x01DB7106, 0x98D220BC, 0xEFD5102A, 0x71B18589, 0x06B6B51F, 0x9FBFE4A5, 0xE8B8D433,
- 0x7807C9A2, 0x0F00F934, 0x9609A88E, 0xE10E9818, 0x7F6A0DBB, 0x086D3D2D, 0x91646C97, 0xE6635C01,
- 0x6B6B51F4, 0x1C6C6162, 0x856530D8, 0xF262004E, 0x6C0695ED, 0x1B01A57B, 0x8208F4C1, 0xF50FC457,
- 0x65B0D9C6, 0x12B7E950, 0x8BBEB8EA, 0xFCB9887C, 0x62DD1DDF, 0x15DA2D49, 0x8CD37CF3, 0xFBD44C65,
- 0x4DB26158, 0x3AB551CE, 0xA3BC0074, 0xD4BB30E2, 0x4ADFA541, 0x3DD895D7, 0xA4D1C46D, 0xD3D6F4FB,
- 0x4369E96A, 0x346ED9FC, 0xAD678846, 0xDA60B8D0, 0x44042D73, 0x33031DE5, 0xAA0A4C5F, 0xDD0D7CC9,
- 0x5005713C, 0x270241AA, 0xBE0B1010, 0xC90C2086, 0x5768B525, 0x206F85B3, 0xB966D409, 0xCE61E49F,
- 0x5EDEF90E, 0x29D9C998, 0xB0D09822, 0xC7D7A8B4, 0x59B33D17, 0x2EB40D81, 0xB7BD5C3B, 0xC0BA6CAD,
- 0xEDB88320, 0x9ABFB3B6, 0x03B6E20C, 0x74B1D29A, 0xEAD54739, 0x9DD277AF, 0x04DB2615, 0x73DC1683,
- 0xE3630B12, 0x94643B84, 0x0D6D6A3E, 0x7A6A5AA8, 0xE40ECF0B, 0x9309FF9D, 0x0A00AE27, 0x7D079EB1,
- 0xF00F9344, 0x8708A3D2, 0x1E01F268, 0x6906C2FE, 0xF762575D, 0x806567CB, 0x196C3671, 0x6E6B06E7,
- 0xFED41B76, 0x89D32BE0, 0x10DA7A5A, 0x67DD4ACC, 0xF9B9DF6F, 0x8EBEEFF9, 0x17B7BE43, 0x60B08ED5,
- 0xD6D6A3E8, 0xA1D1937E, 0x38D8C2C4, 0x4FDFF252, 0xD1BB67F1, 0xA6BC5767, 0x3FB506DD, 0x48B2364B,
- 0xD80D2BDA, 0xAF0A1B4C, 0x36034AF6, 0x41047A60, 0xDF60EFC3, 0xA867DF55, 0x316E8EEF, 0x4669BE79,
- 0xCB61B38C, 0xBC66831A, 0x256FD2A0, 0x5268E236, 0xCC0C7795, 0xBB0B4703, 0x220216B9, 0x5505262F,
- 0xC5BA3BBE, 0xB2BD0B28, 0x2BB45A92, 0x5CB36A04, 0xC2D7FFA7, 0xB5D0CF31, 0x2CD99E8B, 0x5BDEAE1D,
- 0x9B64C2B0, 0xEC63F226, 0x756AA39C, 0x026D930A, 0x9C0906A9, 0xEB0E363F, 0x72076785, 0x05005713,
- 0x95BF4A82, 0xE2B87A14, 0x7BB12BAE, 0x0CB61B38, 0x92D28E9B, 0xE5D5BE0D, 0x7CDCEFB7, 0x0BDBDF21,
- 0x86D3D2D4, 0xF1D4E242, 0x68DDB3F8, 0x1FDA836E, 0x81BE16CD, 0xF6B9265B, 0x6FB077E1, 0x18B74777,
- 0x88085AE6, 0xFF0F6A70, 0x66063BCA, 0x11010B5C, 0x8F659EFF, 0xF862AE69, 0x616BFFD3, 0x166CCF45,
- 0xA00AE278, 0xD70DD2EE, 0x4E048354, 0x3903B3C2, 0xA7672661, 0xD06016F7, 0x4969474D, 0x3E6E77DB,
- 0xAED16A4A, 0xD9D65ADC, 0x40DF0B66, 0x37D83BF0, 0xA9BCAE53, 0xDEBB9EC5, 0x47B2CF7F, 0x30B5FFE9,
- 0xBDBDF21C, 0xCABAC28A, 0x53B39330, 0x24B4A3A6, 0xBAD03605, 0xCDD706B3, 0x54DE5729, 0x23D967BF,
- 0xB3667A2E, 0xC4614AB8, 0x5D681B02, 0x2A6F2B94, 0xB40BBE37, 0xC30C8EA1, 0x5A05DF1B, 0x2D02EF8D
- };
- uint32_t calculate(const void* data, size_t length) {
- uint32_t crc = 0xFFFFFFFF;
- const uint8_t* bytes = static_cast<const uint8_t*>(data);
- for (size_t i = 0; i < length; ++i) {
- crc = table[(crc ^ bytes[i]) & 0xFF] ^ (crc >> 8);
- }
- return crc ^ 0xFFFFFFFF;
- }
- uint32_t calculate(const std::string& str) {
- return calculate(str.data(), str.size());
- }
- } // namespace crc32
- // ===== WalEntry Implementation =====
- std::vector<uint8_t> WalEntry::serialize() const {
- std::vector<uint8_t> result;
- // Reserve space for length prefix (will be set at the end)
- result.resize(4);
- // Write sequence (8 bytes)
- for (int i = 0; i < 8; ++i) {
- result.push_back(static_cast<uint8_t>((sequence >> (i * 8)) & 0xFF));
- }
- // Write timestamp (8 bytes)
- for (int i = 0; i < 8; ++i) {
- result.push_back(static_cast<uint8_t>((timestamp >> (i * 8)) & 0xFF));
- }
- // Write opType (1 byte)
- result.push_back(static_cast<uint8_t>(opType));
- // Write collection (2 bytes length + data)
- uint16_t collLen = static_cast<uint16_t>(collection.size());
- result.push_back(static_cast<uint8_t>(collLen & 0xFF));
- result.push_back(static_cast<uint8_t>((collLen >> 8) & 0xFF));
- result.insert(result.end(), collection.begin(), collection.end());
- // Write documentId (2 bytes length + data)
- uint16_t docIdLen = static_cast<uint16_t>(documentId.size());
- result.push_back(static_cast<uint8_t>(docIdLen & 0xFF));
- result.push_back(static_cast<uint8_t>((docIdLen >> 8) & 0xFF));
- result.insert(result.end(), documentId.begin(), documentId.end());
- // Write data flag and data (1 byte flag + 4 bytes length + data)
- if (data) {
- result.push_back(1);
- std::string dataStr = data->dump();
- uint32_t dataLen = static_cast<uint32_t>(dataStr.size());
- for (int i = 0; i < 4; ++i) {
- result.push_back(static_cast<uint8_t>((dataLen >> (i * 8)) & 0xFF));
- }
- result.insert(result.end(), dataStr.begin(), dataStr.end());
- } else {
- result.push_back(0);
- }
- // Write collectionOptions flag and data
- if (collectionOptions) {
- result.push_back(1);
- std::string optsStr = collectionOptions->toJson().dump();
- uint32_t optsLen = static_cast<uint32_t>(optsStr.size());
- for (int i = 0; i < 4; ++i) {
- result.push_back(static_cast<uint8_t>((optsLen >> (i * 8)) & 0xFF));
- }
- result.insert(result.end(), optsStr.begin(), optsStr.end());
- } else {
- result.push_back(0);
- }
- // Write vectorData flag and data (1 byte flag + 4 bytes dim + dim * sizeof(float) bytes)
- if (vectorData) {
- result.push_back(1);
- uint32_t dim = static_cast<uint32_t>(vectorData->size());
- for (int i = 0; i < 4; ++i) {
- result.push_back(static_cast<uint8_t>((dim >> (i * 8)) & 0xFF));
- }
- const uint8_t* floatBytes = reinterpret_cast<const uint8_t*>(vectorData->data());
- result.insert(result.end(), floatBytes, floatBytes + dim * sizeof(float));
- } else {
- result.push_back(0);
- }
- // v1.8.0 — origin nodeId trailer (1 byte flag + 2 byte len + bytes).
- // Always emitted by current writer. Older readers stop after vectorData
- // and never observe this section; new readers handle missing trailer.
- if (!nodeId.empty()) {
- result.push_back(1);
- uint16_t nodeLen = static_cast<uint16_t>(nodeId.size());
- result.push_back(static_cast<uint8_t>(nodeLen & 0xFF));
- result.push_back(static_cast<uint8_t>((nodeLen >> 8) & 0xFF));
- result.insert(result.end(), nodeId.begin(), nodeId.end());
- } else {
- result.push_back(0);
- }
- // Calculate and append checksum (excluding length prefix and checksum itself)
- uint32_t crc = crc32::calculate(result.data() + 4, result.size() - 4);
- for (int i = 0; i < 4; ++i) {
- result.push_back(static_cast<uint8_t>((crc >> (i * 8)) & 0xFF));
- }
- // Set length prefix (total size minus 4 bytes for length prefix)
- uint32_t totalLen = static_cast<uint32_t>(result.size() - 4);
- result[0] = static_cast<uint8_t>(totalLen & 0xFF);
- result[1] = static_cast<uint8_t>((totalLen >> 8) & 0xFF);
- result[2] = static_cast<uint8_t>((totalLen >> 16) & 0xFF);
- result[3] = static_cast<uint8_t>((totalLen >> 24) & 0xFF);
- return result;
- }
- std::optional<WalEntry> WalEntry::deserialize(const std::vector<uint8_t>& data) {
- if (data.size() < 4) {
- return std::nullopt;
- }
- size_t offset = 0;
- // Read length prefix
- uint32_t length = 0;
- for (int i = 0; i < 4; ++i) {
- length |= static_cast<uint32_t>(data[offset++]) << (i * 8);
- }
- if (data.size() < length + 4) {
- return std::nullopt;
- }
- // Verify checksum
- uint32_t storedCrc = 0;
- for (int i = 0; i < 4; ++i) {
- storedCrc |= static_cast<uint32_t>(data[data.size() - 4 + i]) << (i * 8);
- }
- uint32_t calculatedCrc = crc32::calculate(data.data() + 4, length - 4);
- if (storedCrc != calculatedCrc) {
- return std::nullopt; // Checksum mismatch
- }
- WalEntry entry;
- // Read sequence
- for (int i = 0; i < 8; ++i) {
- entry.sequence |= static_cast<uint64_t>(data[offset++]) << (i * 8);
- }
- // Read timestamp
- for (int i = 0; i < 8; ++i) {
- entry.timestamp |= static_cast<uint64_t>(data[offset++]) << (i * 8);
- }
- // Read opType
- entry.opType = static_cast<WalOpType>(data[offset++]);
- // Read collection
- uint16_t collLen = 0;
- collLen |= static_cast<uint16_t>(data[offset++]);
- collLen |= static_cast<uint16_t>(data[offset++]) << 8;
- entry.collection = std::string(reinterpret_cast<const char*>(&data[offset]), collLen);
- offset += collLen;
- // Read documentId
- uint16_t docIdLen = 0;
- docIdLen |= static_cast<uint16_t>(data[offset++]);
- docIdLen |= static_cast<uint16_t>(data[offset++]) << 8;
- entry.documentId = std::string(reinterpret_cast<const char*>(&data[offset]), docIdLen);
- offset += docIdLen;
- // Read data
- uint8_t hasData = data[offset++];
- if (hasData) {
- uint32_t dataLen = 0;
- for (int i = 0; i < 4; ++i) {
- dataLen |= static_cast<uint32_t>(data[offset++]) << (i * 8);
- }
- std::string dataStr(reinterpret_cast<const char*>(&data[offset]), dataLen);
- offset += dataLen;
- try {
- entry.data = smartbotic::db::parse_to_nlohmann(dataStr);
- } catch (const nlohmann::json::exception&) {
- return std::nullopt;
- }
- }
- // Read collectionOptions
- uint8_t hasOpts = data[offset++];
- if (hasOpts) {
- uint32_t optsLen = 0;
- for (int i = 0; i < 4; ++i) {
- optsLen |= static_cast<uint32_t>(data[offset++]) << (i * 8);
- }
- std::string optsStr(reinterpret_cast<const char*>(&data[offset]), optsLen);
- offset += optsLen;
- try {
- entry.collectionOptions = CollectionOptions::fromJson(smartbotic::db::parse_to_nlohmann(optsStr));
- } catch (const nlohmann::json::exception&) {
- return std::nullopt;
- }
- }
- // Read vectorData (present for VEC_PUT, but flag byte always written)
- if (offset < data.size() - 4) { // at least 1 flag byte + 4 checksum bytes remain
- uint8_t hasVec = data[offset++];
- if (hasVec) {
- uint32_t dim = 0;
- for (int i = 0; i < 4; ++i) {
- dim |= static_cast<uint32_t>(data[offset++]) << (i * 8);
- }
- if (offset + dim * sizeof(float) > data.size() - 4) {
- return std::nullopt; // Truncated vector data
- }
- std::vector<float> vec(dim);
- std::memcpy(vec.data(), &data[offset], dim * sizeof(float));
- offset += dim * sizeof(float);
- entry.vectorData = std::move(vec);
- }
- }
- // v1.8.0 — origin nodeId (optional trailer). Pre-v1.8 entries don't have
- // this section; we leave entry.nodeId empty so the persistence layer can
- // treat it as "origin unknown / assume local".
- if (offset < data.size() - 4) {
- uint8_t hasNode = data[offset++];
- if (hasNode) {
- if (offset + 2 > data.size() - 4) {
- return std::nullopt; // Truncated nodeId length
- }
- uint16_t nodeLen = static_cast<uint16_t>(data[offset]) |
- (static_cast<uint16_t>(data[offset + 1]) << 8);
- offset += 2;
- if (offset + nodeLen > data.size() - 4) {
- return std::nullopt; // Truncated nodeId
- }
- entry.nodeId.assign(reinterpret_cast<const char*>(&data[offset]), nodeLen);
- offset += nodeLen;
- }
- }
- entry.checksum = storedCrc;
- return entry;
- }
- uint32_t WalEntry::calculateChecksum() const {
- auto serialized = serialize();
- // Exclude length prefix and checksum
- return crc32::calculate(serialized.data() + 4, serialized.size() - 8);
- }
- // ===== WriteAheadLog Implementation =====
- WriteAheadLog::WriteAheadLog(Config config)
- : config_(std::move(config)) {
- }
- WriteAheadLog::~WriteAheadLog() {
- close();
- }
- bool WriteAheadLog::open() {
- std::lock_guard<std::mutex> lock(writeMutex_);
- if (isOpen_.load()) {
- return true;
- }
- // Create WAL directory if it doesn't exist
- std::error_code ec;
- std::filesystem::create_directories(config_.walDir, ec);
- if (ec) {
- return false;
- }
- // Find existing WAL files and determine highest sequence
- auto walFiles = getWalFiles();
- if (!walFiles.empty()) {
- // Read the last file to get the highest sequence
- for (auto it = walFiles.rbegin(); it != walFiles.rend(); ++it) {
- std::ifstream file(*it, std::ios::binary);
- if (!file) continue;
- auto header = readHeader(file);
- if (!header) continue;
- // Read entries to find highest sequence
- while (file) {
- uint32_t length;
- file.read(reinterpret_cast<char*>(&length), 4);
- if (!file || length == 0) break;
- std::vector<uint8_t> entryData(length + 4);
- std::memcpy(entryData.data(), &length, 4);
- file.read(reinterpret_cast<char*>(entryData.data() + 4), length);
- if (!file) break;
- auto entry = WalEntry::deserialize(entryData);
- if (entry && entry->sequence > sequence_.load()) {
- sequence_.store(entry->sequence);
- }
- }
- }
- }
- // Open current WAL file for appending
- auto currentPath = currentFilePath();
- bool fileExists = std::filesystem::exists(currentPath);
- currentFile_.open(currentPath, std::ios::binary | std::ios::app);
- if (!currentFile_) {
- return false;
- }
- if (!fileExists) {
- // Write header for new file
- currentFileStartSequence_ = sequence_.load() + 1;
- writeHeader();
- } else {
- // Read existing header
- std::ifstream readFile(currentPath, std::ios::binary);
- auto header = readHeader(readFile);
- if (header) {
- currentFileStartSequence_ = header->startSequence;
- }
- }
- currentFileSize_ = std::filesystem::file_size(currentPath);
- isOpen_.store(true);
- return true;
- }
- void WriteAheadLog::close() {
- std::lock_guard<std::mutex> lock(writeMutex_);
- if (!isOpen_.load()) {
- return;
- }
- if (currentFile_.is_open()) {
- currentFile_.flush();
- currentFile_.close();
- }
- isOpen_.store(false);
- }
- void WriteAheadLog::setSequenceFloor(uint64_t seq) {
- // CAS loop — monotonically raise sequence_ to at least `seq`.
- // Concurrent appends (which `++` the same atomic) are tolerated
- // because compare_exchange_weak retries until it observes a
- // current value that's already ≥ seq, in which case we leave it.
- uint64_t current = sequence_.load();
- while (current < seq && !sequence_.compare_exchange_weak(current, seq)) {
- // current was updated by compare_exchange_weak — loop and retest.
- }
- }
- uint64_t WriteAheadLog::append(WalEntry entry) {
- std::lock_guard<std::mutex> lock(writeMutex_);
- if (!isOpen_.load() || !currentFile_.is_open()) {
- return 0;
- }
- // Assign sequence number and timestamp
- entry.sequence = ++sequence_;
- entry.timestamp = static_cast<uint64_t>(
- std::chrono::duration_cast<std::chrono::milliseconds>(
- std::chrono::system_clock::now().time_since_epoch()
- ).count()
- );
- // Serialize and write
- auto data = entry.serialize();
- currentFile_.write(reinterpret_cast<const char*>(data.data()), static_cast<std::streamsize>(data.size()));
- currentFileSize_ += data.size();
- if (config_.syncOnWrite) {
- currentFile_.flush();
- sync();
- }
- // Rotate if file is too large
- if (currentFileSize_ >= config_.maxFileSizeBytes) {
- rotate();
- }
- return entry.sequence;
- }
- void WriteAheadLog::sync() {
- if (currentFile_.is_open()) {
- currentFile_.flush();
- // Note: For true durability, we'd use fsync() here
- // std::filesystem doesn't provide this, would need OS-specific code
- }
- }
- uint64_t WriteAheadLog::replay(uint64_t fromSequence, std::function<void(const WalEntry&)> callback) {
- auto walFiles = getWalFiles();
- uint64_t count = 0;
- for (const auto& file : walFiles) {
- count += replayFile(file, fromSequence, callback);
- }
- return count;
- }
- uint64_t WriteAheadLog::totalSizeBytes() const {
- uint64_t total = 0;
- auto walFiles = getWalFiles();
- for (const auto& file : walFiles) {
- std::error_code ec;
- total += std::filesystem::file_size(file, ec);
- }
- return total;
- }
- void WriteAheadLog::truncateBefore(uint64_t sequence) {
- std::lock_guard<std::mutex> lock(writeMutex_);
- auto walFiles = getWalFiles();
- for (const auto& file : walFiles) {
- // Read header to get start sequence
- std::ifstream readFile(file, std::ios::binary);
- auto header = readHeader(readFile);
- readFile.close();
- if (header && header->startSequence < sequence) {
- // Check if all entries in this file are before the sequence
- bool canDelete = true;
- uint64_t maxSeqInFile = 0;
- std::ifstream scanFile(file, std::ios::binary);
- scanFile.seekg(sizeof(Header));
- while (scanFile) {
- uint32_t length;
- scanFile.read(reinterpret_cast<char*>(&length), 4);
- if (!scanFile || length == 0) break;
- std::vector<uint8_t> entryData(length + 4);
- std::memcpy(entryData.data(), &length, 4);
- scanFile.read(reinterpret_cast<char*>(entryData.data() + 4), length);
- if (!scanFile) break;
- auto entry = WalEntry::deserialize(entryData);
- if (entry) {
- maxSeqInFile = std::max(maxSeqInFile, entry->sequence);
- }
- }
- if (maxSeqInFile < sequence) {
- // Safe to delete this file
- std::error_code ec;
- std::filesystem::remove(file, ec);
- }
- }
- }
- }
- WalEntry WriteAheadLog::makeInsertEntry(const std::string& collection, const Document& doc,
- const std::string& originNodeId) {
- WalEntry entry;
- entry.opType = WalOpType::INSERT;
- entry.collection = collection;
- entry.documentId = doc.id;
- entry.data = doc.toJson();
- entry.nodeId = originNodeId;
- return entry;
- }
- WalEntry WriteAheadLog::makeUpdateEntry(const std::string& collection, const Document& doc,
- const std::string& originNodeId) {
- WalEntry entry;
- entry.opType = WalOpType::UPDATE;
- entry.collection = collection;
- entry.documentId = doc.id;
- entry.data = doc.toJson();
- entry.nodeId = originNodeId;
- return entry;
- }
- WalEntry WriteAheadLog::makeDeleteEntry(const std::string& collection, const std::string& id,
- const std::string& originNodeId) {
- WalEntry entry;
- entry.opType = WalOpType::DELETE;
- entry.collection = collection;
- entry.documentId = id;
- entry.nodeId = originNodeId;
- return entry;
- }
- WalEntry WriteAheadLog::makeUpsertEntry(const std::string& collection, const Document& doc,
- const std::string& originNodeId) {
- WalEntry entry;
- entry.opType = WalOpType::UPSERT;
- entry.collection = collection;
- entry.documentId = doc.id;
- entry.data = doc.toJson();
- entry.nodeId = originNodeId;
- return entry;
- }
- WalEntry WriteAheadLog::makeCreateCollectionEntry(const std::string& collection, const CollectionOptions& options,
- const std::string& originNodeId) {
- WalEntry entry;
- entry.opType = WalOpType::CREATE_COLLECTION;
- entry.collection = collection;
- entry.collectionOptions = options;
- entry.nodeId = originNodeId;
- return entry;
- }
- WalEntry WriteAheadLog::makeDropCollectionEntry(const std::string& collection,
- const std::string& originNodeId) {
- WalEntry entry;
- entry.opType = WalOpType::DROP_COLLECTION;
- entry.collection = collection;
- entry.nodeId = originNodeId;
- return entry;
- }
- WalEntry WriteAheadLog::makeVecPutEntry(const std::string& collection,
- const std::string& docId, const std::vector<float>& vec,
- const std::string& originNodeId) {
- WalEntry entry;
- entry.opType = WalOpType::VEC_PUT;
- entry.collection = collection;
- entry.documentId = docId;
- entry.vectorData = vec;
- entry.nodeId = originNodeId;
- return entry;
- }
- WalEntry WriteAheadLog::makeVecDeleteEntry(const std::string& collection,
- const std::string& docId,
- const std::string& originNodeId) {
- WalEntry entry;
- entry.opType = WalOpType::VEC_DELETE;
- entry.collection = collection;
- entry.documentId = docId;
- entry.nodeId = originNodeId;
- return entry;
- }
- std::filesystem::path WriteAheadLog::currentFilePath() const {
- return config_.walDir / "wal-current.log";
- }
- std::vector<std::filesystem::path> WriteAheadLog::getWalFiles() const {
- std::vector<std::filesystem::path> files;
- std::error_code ec;
- for (const auto& entry : std::filesystem::directory_iterator(config_.walDir, ec)) {
- if (entry.is_regular_file()) {
- const auto& path = entry.path();
- if (path.filename().string().starts_with("wal-") &&
- path.extension() == ".log") {
- files.push_back(path);
- }
- }
- }
- // Sort by filename (which includes sequence number for archived files)
- std::sort(files.begin(), files.end());
- return files;
- }
- void WriteAheadLog::rotate() {
- if (currentFile_.is_open()) {
- currentFile_.flush();
- currentFile_.close();
- }
- // Rename current file with sequence range
- auto currentPath = currentFilePath();
- std::ostringstream oss;
- oss << "wal-" << std::setfill('0') << std::setw(16) << currentFileStartSequence_
- << "-" << std::setw(16) << sequence_.load() << ".log";
- auto archivePath = config_.walDir / oss.str();
- std::error_code ec;
- std::filesystem::rename(currentPath, archivePath, ec);
- // Open new current file
- currentFileStartSequence_ = sequence_.load() + 1;
- currentFile_.open(currentPath, std::ios::binary | std::ios::app);
- if (currentFile_) {
- writeHeader();
- currentFileSize_ = sizeof(Header);
- }
- }
- std::optional<WriteAheadLog::Header> WriteAheadLog::readHeader(std::ifstream& file) {
- Header header;
- file.read(reinterpret_cast<char*>(&header), sizeof(Header));
- if (!file) {
- return std::nullopt;
- }
- // Verify magic
- if (std::memcmp(header.magic, "CALWAL01", 8) != 0) {
- return std::nullopt;
- }
- return header;
- }
- void WriteAheadLog::writeHeader() {
- Header header;
- std::memcpy(header.magic, "CALWAL01", 8);
- header.startSequence = currentFileStartSequence_;
- currentFile_.write(reinterpret_cast<const char*>(&header), sizeof(Header));
- currentFile_.flush();
- }
- uint64_t WriteAheadLog::replayFile(const std::filesystem::path& path, uint64_t fromSequence,
- const std::function<void(const WalEntry&)>& callback) {
- std::ifstream file(path, std::ios::binary);
- if (!file) {
- return 0;
- }
- auto header = readHeader(file);
- if (!header) {
- return 0;
- }
- uint64_t count = 0;
- while (file) {
- // Read entry length
- uint32_t length;
- file.read(reinterpret_cast<char*>(&length), 4);
- if (!file || length == 0) break;
- // Read entry data
- std::vector<uint8_t> entryData(length + 4);
- std::memcpy(entryData.data(), &length, 4);
- file.read(reinterpret_cast<char*>(entryData.data() + 4), length);
- if (!file) break;
- // Deserialize and invoke callback
- auto entry = WalEntry::deserialize(entryData);
- if (entry && entry->sequence > fromSequence) {
- callback(*entry);
- count++;
- }
- }
- return count;
- }
- } // namespace smartbotic::database
|