Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
17 commits
Select commit Hold shift + click to select a range
c39b994
fix(wal): tag records with their owning catalog and route replay per …
aikins01 Oct 3, 2026
9e79215
test(transaction): cover graph-owned WAL record replay across owner c…
aikins01 Oct 3, 2026
4ac1384
fix(storage): silence unused tableID parameter in NodeTable::initScan…
aikins01 Oct 4, 2026
49bc83f
fix(transaction): restore table-ID isUnCommitted/getLocalRowIdx overl…
aikins01 Oct 4, 2026
7ad7fc3
fix(wal): translate UPDATE_SEQUENCE through replayed entry IDs
aikins01 Oct 4, 2026
7773a99
fix(storage): make the graph registry shared lock reentrant
aikins01 Oct 5, 2026
a4db7df
fix(wal): tag UPDATE_SEQUENCE records with the sequence name
aikins01 Oct 5, 2026
e54375b
fix(storage): skip undo version records for graphs dropped mid-transa…
aikins01 Oct 5, 2026
5a29617
fix(wal): encode named UPDATE_SEQUENCE records as a distinct record type
aikins01 Oct 5, 2026
aeb4185
fix(storage): key graph registry recursion and liveness checks per ma…
aikins01 Oct 5, 2026
1c0c08c
fix(wal): cover legacy sequence record replay with hand-written WAL b…
aikins01 Oct 5, 2026
d43609f
fix(storage): retire dropped graph catalogs after their records are a…
aikins01 Oct 7, 2026
29e1db7
fix(wal): record and translate ALTER connection endpoints across replay
aikins01 Oct 7, 2026
cd36e9d
fix(wal): translate recovered rel-insert endpoints through replayed e…
aikins01 Oct 7, 2026
6c4fdd0
fix(wal): reject torn owner trailers instead of misrouting them
aikins01 Oct 7, 2026
dbeba1a
fix(storage): reject null bytes in graph names
aikins01 Oct 7, 2026
2797058
test(transaction): pin crafted recovery bytes for endpoints, torn tra…
aikins01 Oct 7, 2026
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
73 changes: 50 additions & 23 deletions docs/checkpoint_recovery.md
Original file line number Diff line number Diff line change
Expand Up @@ -34,19 +34,32 @@ and recovery finishes the checkpoint.

## Checkpoint bundle format

`CheckpointRecord::bundleFormatVersion` is 1 for checkpoints written by builds that include this
recovery format, including nightly builds made from it. Version 0 records come from older builds,
including 0.21.2 and every earlier release, and have no version field. Recovery rejects an
unsupported checkpoint version before applying that checkpoint's shadow pages; opening a database
may already have recovered other graphs before it encounters the unsupported version.

A version 1 checkpoint stamps every shadow file header with
`ShadowFile::CHECKPOINT_BUNDLE_DATABASE_ID` instead of a real database ID. Recovery identifies the
data file through the last database-header page in the shadow, which must match the data file's
database ID before any page is written.
`CheckpointRecord::bundleFormatVersion` is 2 for checkpoints written by current builds, including
nightly builds made from them. Version 1 records come from the first builds of this recovery
format, and version 0 records come from older builds, including 0.21.2 and every earlier release,
and have no version field. Recovery rejects an unsupported checkpoint version before applying that
checkpoint's shadow pages; opening a database may already have recovered other graphs before it
encounters the unsupported version.

Version 1 and version 2 share the bundle layout. A checkpoint of either version stamps every
shadow file header with `ShadowFile::CHECKPOINT_BUNDLE_DATABASE_ID` instead of a real database ID,
and recovery identifies the data file through the last database-header page in the shadow, which
must match the data file's database ID before any page is written.

Version 2 changes the WAL record encoding. Every record a version 2 build writes ends with an
`ownerCatalogName` trailer naming the graph catalog the record replays against; the trailer is
empty for main-database records, and its absence in a version 1 record is detected from the
record's framed length. A sequence update that names its sequence is written as the distinct
`UPDATE_SEQUENCE_NAMED` record type — the name lets replay find the sequence by name, because an
implicit serial sequence has no create record, so its entry ID can shift between logging and
replay — while a version 1 build writes only the plain `UPDATE_SEQUENCE` record. Recovery decodes
both versions' records. An older build cannot: it reads a record type it does not know as an
invalid record type, and a record type it does know decodes with the trailer skipped and replays
against the main catalog, so a WAL written by a version 2 build must be replayed and checkpointed
by a version 2 build before an older build opens the database (see Downgrades).

A graph data file opened on its own (outside the database it was created in) checkpoints through
the same format: its WAL ends in a version 1 `CHECKPOINT` record and its shadow carries the
the same format: its WAL ends in a version 2 `CHECKPOINT` record and its shadow carries the
sentinel. Reopening the parent database recovers such a graph from that record — including a
checkpoint interrupted after the record became durable — after validating the graph's WAL header
against the graph data file and requiring the sentinel in its shadow header.
Expand All @@ -65,18 +78,32 @@ or its committed checkpoint recovered when the parent database is reopened.

## Downgrades

Open a database with a build that writes version 1 if it was last closed while a version 1
checkpoint was pending. Pending means `<db>.wal` or `<db>.wal.checkpoint` ends in a `CHECKPOINT`
record and `<db>.shadow` exists.

Older builds refuse such a database because the shadow header does not match the data file. Their
error message suggests deleting the shadow file. **Do not delete it.** After the commit point, the
shadow files hold the only copy of the committed pages, and deleting them loses committed data.
Reopen the database with a build that writes version 1, let recovery finish, and close it cleanly
before going back to an older build.

A database without a pending checkpoint has the same on-disk format as before and opens with older
builds.
Before opening a database that a version 2 build has written with an older build, run `CHECKPOINT`
on a version 2 build and let it succeed. The checkpoint is the step that empties the WAL; a clean
close does not, because `force_checkpoint_on_close=false` leaves version 2 records in the WAL tail
and an auto-checkpoint fires only once the WAL grows past `checkpoint_threshold`.

A pending version 2 checkpoint means `<db>.wal` or `<db>.wal.checkpoint` ends in a `CHECKPOINT`
record and `<db>.shadow` exists. Older builds refuse such a database because the shadow header does
not match the data file. Their error message suggests deleting the shadow file. **Do not delete it.**
After the commit point, the shadow files hold the only copy of the committed pages, and deleting them
loses committed data. Reopen the database with a build that writes version 2, let recovery finish,
and checkpoint before going back to an older build.

An ordinary version 2 WAL tail — owner trailers on the records but no pending checkpoint — is the
quieter hazard: an older build does not reliably refuse it. A record type the older build knows
decodes with the trailer skipped as unknown trailing bytes and replays against the main catalog,
silently misapplying records meant for a graph. A record type the older build does not know, such
as `UPDATE_SEQUENCE_NAMED`, is no safer: with the default non-throwing replay configuration the
decoding error is treated like a corrupt WAL tail — the older build replays the committed prefix
and truncates the WAL there, silently discarding that transaction and every transaction recorded
after it. Setting `throwOnWalReplayFailure` to true rejects an unknown record before that WAL
is applied, but it still does not make downgrading safe: a WAL containing only recognized
record types passes validation, and its owner trailers are ignored. This is why the successful
`CHECKPOINT`, not a clean close, is the downgrade gate.

A database with neither a pending checkpoint nor any remaining version 2 records has the same
on-disk format as before and opens with older builds.

## Version 0 checkpoints

Expand Down
50 changes: 27 additions & 23 deletions src/catalog/catalog.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -77,6 +77,10 @@ Catalog* Catalog::Get(const main::ClientContext& context) {
return context.getAttachedDatabase()->getCatalog();
}
auto dbManager = main::DatabaseManager::Get(context);
if (auto* replayOwnerCatalog = dbManager->getReplayOwnerCatalog();
replayOwnerCatalog != nullptr) {
return replayOwnerCatalog;
}
if (dbManager->hasDefaultGraph()) {
auto graphCatalog = dbManager->getDefaultGraphCatalog();
if (graphCatalog != nullptr) {
Expand All @@ -87,16 +91,16 @@ Catalog* Catalog::Get(const main::ClientContext& context) {
}

void Catalog::initCatalogSets() {
tables = std::make_unique<CatalogSet>();
sequences = std::make_unique<CatalogSet>();
functions = std::make_unique<CatalogSet>();
types = std::make_unique<CatalogSet>();
indexes = std::make_unique<CatalogSet>();
macros = std::make_unique<CatalogSet>();
internalTables = std::make_unique<CatalogSet>(true /* isInternal */);
internalSequences = std::make_unique<CatalogSet>(true /* isInternal */);
internalFunctions = std::make_unique<CatalogSet>(true /* isInternal */);
graphs = std::make_unique<CatalogSet>();
tables = std::make_unique<CatalogSet>(this);
sequences = std::make_unique<CatalogSet>(this);
functions = std::make_unique<CatalogSet>(this);
types = std::make_unique<CatalogSet>(this);
indexes = std::make_unique<CatalogSet>(this);
macros = std::make_unique<CatalogSet>(this);
internalTables = std::make_unique<CatalogSet>(this, true /* isInternal */);
internalSequences = std::make_unique<CatalogSet>(this, true /* isInternal */);
internalFunctions = std::make_unique<CatalogSet>(this, true /* isInternal */);
graphs = std::make_unique<CatalogSet>(this);
}

bool Catalog::containsTable(const Transaction* transaction, const std::string& tableName,
Expand Down Expand Up @@ -413,10 +417,10 @@ bool Catalog::containsType(const Transaction* transaction, const std::string& ty
return types->containsEntry(transaction, typeName);
}

void Catalog::createIndex(Transaction* transaction, std::unique_ptr<CatalogEntry> indexCatalogEntry,
bool skipLoggingToWAL) {
oid_t Catalog::createIndex(Transaction* transaction,
std::unique_ptr<CatalogEntry> indexCatalogEntry, bool skipLoggingToWAL) {
DASSERT(indexCatalogEntry->getType() == CatalogEntryType::INDEX_ENTRY);
indexes->createEntry(transaction, std::move(indexCatalogEntry), skipLoggingToWAL);
return indexes->createEntry(transaction, std::move(indexCatalogEntry), skipLoggingToWAL);
}

IndexCatalogEntry* Catalog::getIndex(const Transaction* transaction, table_id_t tableID,
Expand Down Expand Up @@ -820,17 +824,17 @@ void Catalog::serializeSnapshot(Serializer& ser, common::transaction_t snapshotT
}

void Catalog::deserialize(Deserializer& deSer) {
tables = CatalogSet::deserialize(deSer);
sequences = CatalogSet::deserialize(deSer);
functions = CatalogSet::deserialize(deSer);
tables = CatalogSet::deserialize(this, deSer);
sequences = CatalogSet::deserialize(this, deSer);
functions = CatalogSet::deserialize(this, deSer);
registerBuiltInFunctions();
types = CatalogSet::deserialize(deSer);
indexes = CatalogSet::deserialize(deSer);
macros = CatalogSet::deserialize(deSer);
internalTables = CatalogSet::deserialize(deSer);
internalSequences = CatalogSet::deserialize(deSer);
internalFunctions = CatalogSet::deserialize(deSer);
graphs = CatalogSet::deserialize(deSer);
types = CatalogSet::deserialize(this, deSer);
indexes = CatalogSet::deserialize(this, deSer);
macros = CatalogSet::deserialize(this, deSer);
internalTables = CatalogSet::deserialize(this, deSer);
internalSequences = CatalogSet::deserialize(this, deSer);
internalFunctions = CatalogSet::deserialize(this, deSer);
graphs = CatalogSet::deserialize(this, deSer);
}

} // namespace catalog
Expand Down
5 changes: 5 additions & 0 deletions src/catalog/catalog_entry/catalog_entry.cpp
Original file line number Diff line number Diff line change
@@ -1,5 +1,6 @@
#include "catalog/catalog_entry/catalog_entry.h"

#include "catalog/catalog.h"
#include "catalog/catalog_entry/graph_catalog_entry.h"
#include "catalog/catalog_entry/index_catalog_entry.h"
#include "catalog/catalog_entry/scalar_macro_catalog_entry.h"
Expand Down Expand Up @@ -78,5 +79,9 @@ void CatalogEntry::copyFrom(const CatalogEntry& other) {
hasParent_ = other.hasParent_;
}

std::string CatalogEntry::getOwningCatalogName() const {
return owningCatalog ? owningCatalog->getCatalogName() : "";
}

} // namespace catalog
} // namespace lbug
25 changes: 24 additions & 1 deletion src/catalog/catalog_set.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -3,7 +3,9 @@
#include <mutex>

#include "binder/ddl/bound_alter_info.h"
#include "catalog/catalog.h"
#include "catalog/catalog_entry/dummy_catalog_entry.h"
#include "catalog/catalog_entry/rel_group_catalog_entry.h"
#include "catalog/catalog_entry/table_catalog_entry.h"
#include "common/assert.h"
#include "common/exception/catalog.h"
Expand All @@ -23,6 +25,12 @@ CatalogSet::CatalogSet(bool isInternal) {
}
}

CatalogSet::CatalogSet(Catalog* catalog, bool isInternal) : catalog{catalog} {
if (isInternal) {
nextOID = INTERNAL_CATALOG_SET_START_OID;
}
}

static bool checkWWConflict(const Transaction* transaction, const CatalogEntry* entry) {
return (entry->getTimestamp() >= Transaction::START_TRANSACTION_ID &&
entry->getTimestamp() != transaction->getID()) ||
Expand Down Expand Up @@ -104,6 +112,7 @@ CatalogEntry* CatalogSet::createEntryNoLock(const Transaction* transaction,
}

void CatalogSet::emplaceNoLock(std::unique_ptr<CatalogEntry> entry) {
entry->setOwningCatalog(catalog);
if (entries.contains(entry->getName())) {
entry->setPrev(std::move(entries.at(entry->getName())));
entries.erase(entry->getName());
Expand Down Expand Up @@ -198,9 +207,14 @@ void CatalogSet::alterTableEntry(Transaction* transaction, const binder::BoundAl
case AlterType::SET_SORTED_BY:
case AlterType::ADD_FROM_TO_CONNECTION:
case AlterType::DROP_FROM_TO_CONNECTION: {
auto addedRelTableOID = common::INVALID_TABLE_ID;
if (alterInfo.alterType == AlterType::ADD_FROM_TO_CONNECTION) {
addedRelTableOID =
newEntry->ptrCast<RelGroupCatalogEntry>()->getRelEntryInfos().back().oid;
}
emplaceNoLock(std::move(newEntry));
if (transaction->shouldAppendToUndoBuffer()) {
transaction->pushAlterCatalogEntry(*this, *entry, alterInfo);
transaction->pushAlterCatalogEntry(*this, *entry, alterInfo, false, addedRelTableOID);
}
} break;
default: {
Expand Down Expand Up @@ -300,8 +314,13 @@ void CatalogSet::serializeSnapshot(Serializer serializer, const Transaction* sna
}

std::unique_ptr<CatalogSet> CatalogSet::deserialize(Deserializer& deserializer) {
return deserialize(nullptr, deserializer);
}

std::unique_ptr<CatalogSet> CatalogSet::deserialize(Catalog* catalog, Deserializer& deserializer) {
std::string debuggingInfo;
auto catalogSet = std::make_unique<CatalogSet>();
catalogSet->catalog = catalog;
deserializer.validateDebuggingInfo(debuggingInfo, "nextOID");
deserializer.deserializeValue<oid_t>(catalogSet->nextOID);
uint64_t numEntries = 0;
Expand All @@ -316,6 +335,10 @@ std::unique_ptr<CatalogSet> CatalogSet::deserialize(Deserializer& deserializer)
return catalogSet;
}

std::string CatalogSet::getOwnerCatalogName() const {
return catalog ? catalog->getCatalogName() : "";
}

// Ideally we should not trigger the following check. Instead, we should throw more informative
// error message at catalog level.
void CatalogSet::validateExistNoLock(const Transaction* transaction,
Expand Down
2 changes: 1 addition & 1 deletion src/extension/extension_manager.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -46,7 +46,7 @@ void ExtensionManager::loadExtension(const std::string& path, main::ClientContex
isOfficial ? ExtensionSource::OFFICIAL : ExtensionSource::USER));
auto transaction = transaction::Transaction::Get(*context);
if (transaction->shouldLogToWAL()) {
transaction->getLocalWAL().logLoadExtension(path);
transaction->getLocalWAL().logLoadExtension("" /* main catalog */, path);
}
}

Expand Down
4 changes: 2 additions & 2 deletions src/graph/on_disk_graph.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -361,8 +361,8 @@ bool OnDiskGraphVertexScanState::next() {
auto endOffset = std::min(endOffsetExclusive,
tableScanState->source == TableScanSource::COMMITTED ?
startOffsetOfNextGroup :
startOffsetOfNextGroup + transaction->getUncommittedOffset(
tableScanState->table->getTableID(), currentOffset));
startOffsetOfNextGroup +
transaction->getUncommittedOffset(*tableScanState->table, currentOffset));
numNodesToScan = std::min(endOffset - currentOffset, DEFAULT_VECTOR_CAPACITY);
auto result = tableScanState->scanNext(transaction, currentOffset, numNodesToScan);
currentOffset += result.numRows;
Expand Down
7 changes: 6 additions & 1 deletion src/include/catalog/catalog.h
Original file line number Diff line number Diff line change
Expand Up @@ -156,7 +156,7 @@ class LBUG_API Catalog {
common::table_id_t tableID) const;

// Create index entry.
void createIndex(transaction::Transaction* transaction,
common::oid_t createIndex(transaction::Transaction* transaction,
std::unique_ptr<CatalogEntry> indexCatalogEntry, bool skipLoggingToWAL = false);
// Drop all index entries within a table.
void dropAllIndexes(transaction::Transaction* transaction, common::table_id_t tableID);
Expand Down Expand Up @@ -216,6 +216,11 @@ class LBUG_API Catalog {
// Get all graph entries.
std::vector<GraphCatalogEntry*> getGraphEntries(
const transaction::Transaction* transaction) const;
// The next graph-entry OID this catalog would assign. Replay compares recorded
// GRAPH_ENTRY IDs against the value from before the replay pass: IDs at or above
// it were assigned after the last checkpoint and shift during recovery, while IDs
// below it belong to persisted entries and stay stable.
common::oid_t peekNextGraphOID() const { return graphs->peekNextOID(); }

// Create graph entry.
void createGraph(transaction::Transaction* transaction, std::string name, bool isAnyGraph);
Expand Down
7 changes: 7 additions & 0 deletions src/include/catalog/catalog_entry/catalog_entry.h
Original file line number Diff line number Diff line change
Expand Up @@ -15,6 +15,8 @@ class ClientContext;

namespace catalog {

class Catalog;

struct LBUG_API ToCypherInfo {
virtual ~ToCypherInfo() = default;

Expand All @@ -40,6 +42,9 @@ class LBUG_API CatalogEntry {
// getter & setter
//===--------------------------------------------------------------------===//
CatalogEntryType getType() const { return type; }
void setOwningCatalog(Catalog* catalog) { owningCatalog = catalog; }
Catalog* getOwningCatalog() const { return owningCatalog; }
std::string getOwningCatalogName() const;
void rename(std::string name_) { this->name = std::move(name_); }
std::string getName() const { return name; }
common::transaction_t getTimestamp() const { return timestamp; }
Expand Down Expand Up @@ -99,6 +104,8 @@ class LBUG_API CatalogEntry {

protected:
CatalogEntryType type;
// Never serialized; re-established from the owning CatalogSet on load.
Catalog* owningCatalog = nullptr;
std::string name;
common::oid_t oid;
common::transaction_t timestamp;
Expand Down
7 changes: 7 additions & 0 deletions src/include/catalog/catalog_entry/index_catalog_entry.h
Original file line number Diff line number Diff line change
Expand Up @@ -73,6 +73,13 @@ class LBUG_API IndexCatalogEntry final : public CatalogEntry {

common::table_id_t getTableID() const { return tableID; }

// The catalog-set key embeds the table ID (see getInternalIndexName), so the
// name must change with it.
void setTableID(common::table_id_t tableID_) {
tableID = tableID_;
rename(getInternalIndexName(tableID_, indexName));
}

std::string getIndexName() const { return indexName; }

std::vector<common::property_id_t> getPropertyIDs() const { return propertyIDs; }
Expand Down
14 changes: 14 additions & 0 deletions src/include/catalog/catalog_set.h
Original file line number Diff line number Diff line change
Expand Up @@ -22,12 +22,15 @@ class Transaction;
using CatalogEntrySet = common::case_insensitive_map_t<catalog::CatalogEntry*>;

namespace catalog {
class Catalog;

class LBUG_API CatalogSet {
friend class storage::UndoBuffer;

public:
CatalogSet() = default;
explicit CatalogSet(bool isInternal);
explicit CatalogSet(Catalog* catalog, bool isInternal = false);
bool containsEntry(const transaction::Transaction* transaction, const std::string& name);
CatalogEntry* getEntry(const transaction::Transaction* transaction, const std::string& name);
common::oid_t createEntry(transaction::Transaction* transaction,
Expand All @@ -45,6 +48,11 @@ class LBUG_API CatalogSet {
void serializeSnapshot(common::Serializer serializer,
const transaction::Transaction* snapshotTxn) const;
static std::unique_ptr<CatalogSet> deserialize(common::Deserializer& deserializer);
static std::unique_ptr<CatalogSet> deserialize(Catalog* catalog,
common::Deserializer& deserializer);

Catalog* getCatalog() const { return catalog; }
std::string getOwnerCatalogName() const;

common::oid_t getNextOID() {
std::unique_lock lck{mtx};
Expand All @@ -53,6 +61,11 @@ class LBUG_API CatalogSet {

common::oid_t getNextOIDNoLock() { return nextOID++; }

common::oid_t peekNextOID() const {
std::shared_lock lck{mtx};
return nextOID;
}

private:
bool containsEntryNoLock(const transaction::Transaction* transaction,
const std::string& name) const;
Expand Down Expand Up @@ -86,6 +99,7 @@ class LBUG_API CatalogSet {

private:
mutable std::shared_mutex mtx;
Catalog* catalog = nullptr;
common::oid_t nextOID = 0;
common::case_insensitive_map_t<std::unique_ptr<CatalogEntry>> entries;
};
Expand Down
Loading
Loading