Skip to content

Commit 0481b20

Browse files
committed
fix: isolate drop recreate collection state
1 parent 63cdbb4 commit 0481b20

40 files changed

Lines changed: 1353 additions & 754 deletions

File tree

nodedb-test-support/src/cluster_harness/node/lifecycle/spawn_full.rs

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -96,7 +96,7 @@ impl TestClusterNode {
9696
&data_dir_path.join("test.wal"),
9797
)?);
9898
let wal_records: Arc<[nodedb_wal::WalRecord]> = Arc::from(wal.replay()?.into_boxed_slice());
99-
let replay_tombstones = nodedb_wal::extract_tombstones(&wal_records);
99+
let replay_tombstones = nodedb_wal::extract_tombstones(&wal_records).unwrap();
100100
let (dispatcher, data_sides) = Dispatcher::new(num_cores, 1024);
101101
let (event_producers, event_consumers) = create_event_bus(num_cores);
102102

nodedb-test-support/src/pgwire_harness/restart.rs

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -172,7 +172,7 @@ impl TestServer {
172172
let wal = Arc::new(WalManager::open_for_testing(&wal_path).unwrap());
173173
let wal_records: Arc<[nodedb_wal::WalRecord]> =
174174
Arc::from(wal.replay().unwrap().into_boxed_slice());
175-
let replay_tombstones = nodedb_wal::extract_tombstones(&wal_records);
175+
let replay_tombstones = nodedb_wal::extract_tombstones(&wal_records).unwrap();
176176

177177
let (dispatcher, data_sides) = Dispatcher::new(1, 64);
178178
let (event_producers, event_consumers) = create_event_bus(1);

nodedb-wal/src/replay/filter.rs

Lines changed: 17 additions & 29 deletions
Original file line numberDiff line numberDiff line change
@@ -125,11 +125,9 @@ impl std::ops::Deref for DatabaseTombstones<'_> {
125125
/// Single-pass extraction: walk `records`, decode every
126126
/// `CollectionTombstoned` entry, and return the resulting set.
127127
///
128-
/// Records that fail to decode are logged and skipped — replay callers
129-
/// that need hard failure on corrupt tombstones should validate CRCs
130-
/// upstream; a CRC-passing but structurally corrupt payload is an
131-
/// out-of-band programmer bug, not a run-time recoverable condition.
132-
pub fn extract_tombstones(records: &[WalRecord]) -> TombstoneSet {
128+
/// A malformed tombstone fails extraction. Skipping one would replay writes
129+
/// below an unknown purge boundary and can resurrect a dropped collection.
130+
pub fn extract_tombstones(records: &[WalRecord]) -> crate::Result<TombstoneSet> {
133131
let mut set = TombstoneSet::new();
134132
for record in records {
135133
let Some(kind) = RecordType::from_raw(record.logical_record_type()) else {
@@ -142,22 +140,15 @@ pub fn extract_tombstones(records: &[WalRecord]) -> TombstoneSet {
142140
// they carry only a collection name and an LSN, no secrets.
143141
// If payload-level encryption is later extended to tombstones,
144142
// this call site changes to `decrypt_payload_ring` first.
145-
match CollectionTombstonePayload::from_bytes(&record.payload) {
146-
Ok(payload) => set.insert(
147-
record.header.database_id,
148-
record.header.tenant_id,
149-
payload.collection,
150-
payload.purge_lsn,
151-
),
152-
Err(e) => tracing::warn!(
153-
lsn = record.header.lsn,
154-
tenant_id = record.header.tenant_id,
155-
error = %e,
156-
"skipping malformed CollectionTombstoned record during tombstone extraction",
157-
),
158-
}
143+
let payload = CollectionTombstonePayload::from_bytes(&record.payload)?;
144+
set.insert(
145+
record.header.database_id,
146+
record.header.tenant_id,
147+
payload.collection,
148+
payload.purge_lsn,
149+
);
159150
}
160-
set
151+
Ok(set)
161152
}
162153

163154
#[cfg(test)]
@@ -217,7 +208,7 @@ mod tests {
217208
tombstone_record(7, 1, "orders", 150, 11),
218209
tombstone_record(8, 1, "users", 200, 12),
219210
];
220-
let set = extract_tombstones(&records);
211+
let set = extract_tombstones(&records).unwrap();
221212
assert_eq!(set.len(), 3);
222213
assert_eq!(set.purge_lsn(7, 1, "users"), Some(100));
223214
assert_eq!(set.purge_lsn(7, 1, "orders"), Some(150));
@@ -251,12 +242,12 @@ mod tests {
251242
})
252243
.unwrap(),
253244
];
254-
let set = extract_tombstones(&records);
245+
let set = extract_tombstones(&records).unwrap();
255246
assert_eq!(set.len(), 1);
256247
}
257248

258249
#[test]
259-
fn extract_tolerates_corrupt_payload() {
250+
fn extract_rejects_corrupt_payload() {
260251
// Build a tombstone-typed record whose payload is too short to decode.
261252
let bogus = WalRecord::new(WalRecordArgs {
262253
record_type: RecordType::CollectionTombstoned as u32,
@@ -270,13 +261,10 @@ mod tests {
270261
})
271262
.unwrap();
272263
let records = vec![bogus, tombstone_record(0, 1, "users", 100, 10)];
273-
let set = extract_tombstones(&records);
274-
assert_eq!(
275-
set.len(),
276-
1,
277-
"valid tombstone still captured despite peer corruption"
264+
assert!(
265+
extract_tombstones(&records).is_err(),
266+
"malformed tombstone must fail replay rather than lose a purge boundary"
278267
);
279-
assert_eq!(set.purge_lsn(0, 1, "users"), Some(100));
280268
}
281269

282270
#[test]

nodedb-wal/tests/wal_collection_tombstone.rs

Lines changed: 3 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -95,7 +95,7 @@ fn extract_and_shadow_writes_before_purge_lsn() {
9595
writer.sync().unwrap();
9696

9797
let records = read_all(&path);
98-
let set: TombstoneSet = extract_tombstones(&records);
98+
let set: TombstoneSet = extract_tombstones(&records).unwrap();
9999

100100
assert_eq!(set.len(), 1);
101101
assert_eq!(set.purge_lsn(0, 1, "users"), Some(purge_lsn));
@@ -159,7 +159,7 @@ fn multiple_tombstones_keep_highest_purge_lsn() {
159159
writer.sync().unwrap();
160160

161161
let records = read_all(&path);
162-
let set = extract_tombstones(&records);
162+
let set = extract_tombstones(&records).unwrap();
163163
assert_eq!(set.purge_lsn(0, 1, "users"), Some(500));
164164
assert!(set.is_tombstoned(0, 1, "users", 499));
165165
assert!(!set.is_tombstoned(0, 1, "users", 500));
@@ -179,6 +179,6 @@ fn extract_ignores_unrelated_record_types() {
179179

180180
let records = read_all(&path);
181181
assert_eq!(records.len(), 4);
182-
let set = extract_tombstones(&records);
182+
let set = extract_tombstones(&records).unwrap();
183183
assert!(set.is_empty());
184184
}

nodedb/src/bootstrap/cluster_ready.rs

Lines changed: 11 additions & 16 deletions
Original file line numberDiff line numberDiff line change
@@ -150,25 +150,20 @@ pub async fn await_cluster_ready(
150150
return Err(e);
151151
}
152152

153+
// A pending name-scoped reclaim must complete before the gateway opens.
154+
// Otherwise a same-name CREATE can install a replacement that the delayed
155+
// retry subsequently erases. Fail readiness and let the operator restart
156+
// after the underlying storage fault is resolved.
157+
if let Err(error) = crate::event::collection_gc::pending_reclaim::drain_once(shared).await {
158+
data_groups_gate.fail(format!("pending collection reclaim failed: {error}"));
159+
return Err(anyhow::anyhow!(
160+
"pending collection reclaim failed during startup: {error}"
161+
));
162+
}
163+
153164
data_groups_gate.fire();
154165
transport_gate.fire();
155166

156-
// Boot-repair outstanding engine purges: a node that crashed with a
157-
// dropped collection's catalog row already removed but its redb +
158-
// versioned engine purge incomplete recorded a `_system.pending_reclaim`
159-
// entry (or would have, had it stayed up). Drain that table once at
160-
// boot — re-running the engine purge for each pending entry — so the
161-
// reclaim completes promptly instead of waiting for the first worker
162-
// tick. NOT a readiness gate: a purge hiccup must never wedge boot, so
163-
// this is spawned and any per-entry failure is left for the
164-
// pending-reclaim worker to retry.
165-
{
166-
let drain_shared = Arc::clone(shared);
167-
tokio::spawn(async move {
168-
crate::event::collection_gc::pending_reclaim::drain_once(&drain_shared).await;
169-
});
170-
}
171-
172167
// Warm the QUIC peer cache so the first replicated request
173168
// after boot doesn't pay a cold dial.
174169
if let (Some(transport), Some(topology)) = (

nodedb/src/bootstrap/wal_init.rs

Lines changed: 17 additions & 31 deletions
Original file line numberDiff line numberDiff line change
@@ -68,43 +68,29 @@ pub fn init_wal(
6868
for at-rest catalog encryption"
6969
);
7070

71-
let tombstones = load_tombstones(config, &wal_records);
71+
let tombstones = load_tombstones(config, &wal_records)?;
7272
Ok((wal, wal_records, tombstones))
7373
}
7474

7575
fn load_tombstones(
7676
config: &ServerConfig,
7777
wal_records: &Arc<[nodedb_wal::WalRecord]>,
78-
) -> nodedb_wal::TombstoneSet {
78+
) -> anyhow::Result<nodedb_wal::TombstoneSet> {
7979
let catalog_path = config.catalog_path();
80-
let mut set = nodedb_wal::extract_tombstones(wal_records);
81-
match crate::control::security::catalog::SystemCatalog::open(&catalog_path) {
82-
Ok(catalog) => match catalog.load_wal_tombstones() {
83-
Ok(persisted) => {
84-
if !persisted.is_empty() {
85-
info!(
86-
persisted = persisted.len(),
87-
in_wal = set.len(),
88-
"merging persisted collection tombstones into replay set"
89-
);
90-
}
91-
set.extend(persisted);
92-
}
93-
Err(e) => {
94-
tracing::warn!(
95-
error = %e,
96-
"failed to load _system.wal_tombstones at startup — \
97-
replay will see WAL-extracted tombstones only"
98-
);
99-
}
100-
},
101-
Err(e) => {
102-
tracing::warn!(
103-
error = %e,
104-
"could not open system catalog to load persisted WAL tombstones — \
105-
falling back to WAL-extracted set"
106-
);
107-
}
80+
let mut set = nodedb_wal::extract_tombstones(wal_records)
81+
.map_err(|error| anyhow::anyhow!("extract WAL tombstones: {error}"))?;
82+
let catalog = crate::control::security::catalog::SystemCatalog::open(&catalog_path)
83+
.map_err(|error| anyhow::anyhow!("open catalog for WAL tombstones: {error}"))?;
84+
let persisted = catalog
85+
.load_wal_tombstones()
86+
.map_err(|error| anyhow::anyhow!("load persisted WAL tombstones: {error}"))?;
87+
if !persisted.is_empty() {
88+
info!(
89+
persisted = persisted.len(),
90+
in_wal = set.len(),
91+
"merging persisted collection tombstones into replay set"
92+
);
10893
}
109-
set
94+
set.extend(persisted);
95+
Ok(set)
11096
}

0 commit comments

Comments
 (0)