#include "wal.hpp" #include "../json_parse.hpp" #include #include #include #include #include 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(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 WalEntry::serialize() const { std::vector 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((sequence >> (i * 8)) & 0xFF)); } // Write timestamp (8 bytes) for (int i = 0; i < 8; ++i) { result.push_back(static_cast((timestamp >> (i * 8)) & 0xFF)); } // Write opType (1 byte) result.push_back(static_cast(opType)); // Write collection (2 bytes length + data) uint16_t collLen = static_cast(collection.size()); result.push_back(static_cast(collLen & 0xFF)); result.push_back(static_cast((collLen >> 8) & 0xFF)); result.insert(result.end(), collection.begin(), collection.end()); // Write documentId (2 bytes length + data) uint16_t docIdLen = static_cast(documentId.size()); result.push_back(static_cast(docIdLen & 0xFF)); result.push_back(static_cast((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(dataStr.size()); for (int i = 0; i < 4; ++i) { result.push_back(static_cast((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(optsStr.size()); for (int i = 0; i < 4; ++i) { result.push_back(static_cast((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(vectorData->size()); for (int i = 0; i < 4; ++i) { result.push_back(static_cast((dim >> (i * 8)) & 0xFF)); } const uint8_t* floatBytes = reinterpret_cast(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(nodeId.size()); result.push_back(static_cast(nodeLen & 0xFF)); result.push_back(static_cast((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((crc >> (i * 8)) & 0xFF)); } // Set length prefix (total size minus 4 bytes for length prefix) uint32_t totalLen = static_cast(result.size() - 4); result[0] = static_cast(totalLen & 0xFF); result[1] = static_cast((totalLen >> 8) & 0xFF); result[2] = static_cast((totalLen >> 16) & 0xFF); result[3] = static_cast((totalLen >> 24) & 0xFF); return result; } std::optional WalEntry::deserialize(const std::vector& 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(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(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(data[offset++]) << (i * 8); } // Read timestamp for (int i = 0; i < 8; ++i) { entry.timestamp |= static_cast(data[offset++]) << (i * 8); } // Read opType entry.opType = static_cast(data[offset++]); // Read collection uint16_t collLen = 0; collLen |= static_cast(data[offset++]); collLen |= static_cast(data[offset++]) << 8; entry.collection = std::string(reinterpret_cast(&data[offset]), collLen); offset += collLen; // Read documentId uint16_t docIdLen = 0; docIdLen |= static_cast(data[offset++]); docIdLen |= static_cast(data[offset++]) << 8; entry.documentId = std::string(reinterpret_cast(&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(data[offset++]) << (i * 8); } std::string dataStr(reinterpret_cast(&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(data[offset++]) << (i * 8); } std::string optsStr(reinterpret_cast(&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(data[offset++]) << (i * 8); } if (offset + dim * sizeof(float) > data.size() - 4) { return std::nullopt; // Truncated vector data } std::vector 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(data[offset]) | (static_cast(data[offset + 1]) << 8); offset += 2; if (offset + nodeLen > data.size() - 4) { return std::nullopt; // Truncated nodeId } entry.nodeId.assign(reinterpret_cast(&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 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(&length), 4); if (!file || length == 0) break; std::vector entryData(length + 4); std::memcpy(entryData.data(), &length, 4); file.read(reinterpret_cast(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 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 lock(writeMutex_); if (!isOpen_.load() || !currentFile_.is_open()) { return 0; } // Assign sequence number and timestamp entry.sequence = ++sequence_; entry.timestamp = static_cast( std::chrono::duration_cast( std::chrono::system_clock::now().time_since_epoch() ).count() ); // Serialize and write auto data = entry.serialize(); currentFile_.write(reinterpret_cast(data.data()), static_cast(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 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 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(&length), 4); if (!scanFile || length == 0) break; std::vector entryData(length + 4); std::memcpy(entryData.data(), &length, 4); scanFile.read(reinterpret_cast(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& 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 WriteAheadLog::getWalFiles() const { std::vector 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::readHeader(std::ifstream& file) { Header header; file.read(reinterpret_cast(&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(&header), sizeof(Header)); currentFile_.flush(); } uint64_t WriteAheadLog::replayFile(const std::filesystem::path& path, uint64_t fromSequence, const std::function& 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(&length), 4); if (!file || length == 0) break; // Read entry data std::vector entryData(length + 4); std::memcpy(entryData.data(), &length, 4); file.read(reinterpret_cast(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