1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
//! Bounded, committed deletion of a key range.
//!
//! A range delete whose cardinality follows user data — a reclaimed prefix, a
//! changefeed retention backlog — cannot be a single transaction: the write
//! batch, and so the memory the job holds, would grow with the data. Such a
//! delete is instead paged here, in transactions of at most one page each, every
//! one committed before the next page is scanned. The job's memory is then a
//! function of the page size alone.
//!
//! Pages are scanned inside the transaction that deletes them, so the page and
//! its deletions cannot disagree, and every deleted key is charged to that
//! transaction's write set — which is what makes the write-cardinality guard and
//! the transaction metrics see the work at all. A whole-range `delr` reports
//! neither.
//!
//! A pass carries a key budget so that one large range cannot spend the whole
//! pass and starve the ranges behind it. Running out of budget, or losing the
//! task lease, ends the delete with every page so far committed, which is what
//! lets the next pass continue rather than restart.
use std::ops::Range;
use anyhow::Result;
use tokio_util::sync::CancellationToken;
use crate::catalog::{DatabaseId, NamespaceId};
use crate::key::root::rc::{ReclaimKey, ReclaimState};
use crate::kvs::LockType::Optimistic;
use crate::kvs::TransactionType::{Read, Write};
use crate::kvs::tasklease::LeaseHandler;
use crate::kvs::{Datastore, Key};
/// How a bounded range delete ended.
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
pub(crate) enum PagedOutcome {
/// The range is empty. The page that emptied it reports this, so a caller
/// that holds a claim on the range can retire it without another pass.
Complete,
/// Keys remain because the budget ran out. Every page deleted so far is
/// committed.
Incomplete,
/// Keys remain because the task lease passed to another node. Every page
/// deleted so far is committed, and the caller must stop rather than move on
/// to its next range: a lease check is throttled to one datastore read per
/// maintenance period, so this answer consumed the read the caller's own
/// re-check depends on and that re-check would report the lease held without
/// asking.
LeaseLost,
/// The claim authorising the delete was withdrawn while it ran, so nothing
/// further may be deleted from the range.
Cancelled,
}
/// What rides in each page's transaction besides the page's own deletions.
pub(crate) enum PageCompanion<'a> {
/// Nothing. The committed deletions are the only record of progress, so a
/// later pass resumes by scanning the range from its head and landing on the
/// first key that survives. Correct only where the range's lower bound is
/// fixed — everything below it has already been deleted, so rescanning
/// re-reads no live key.
None,
/// A reclaim queue entry. Its presence is the claim that the range is
/// orphaned, and its cursor advances with every page so a later pass resumes
/// at the last durable deletion instead of rescanning the prefix.
ReclaimEntry {
rc: &'a ReclaimKey<'a>,
/// The exact state this page is allowed to advance. Updating it after
/// every commit turns the cursor write into a compare-and-set, so a
/// concurrent cancellation or page wins instead of being overwritten.
state: ReclaimState,
},
/// A changefeed retention watermark. The window's upper bound is only the
/// keys the database's retention has made stale, and that policy is catalog
/// state a concurrent `DEFINE DATABASE` or `DEFINE TABLE` can extend. The
/// watermark is therefore recomputed inside every page transaction, and the
/// delete stops if the policy now reaches below the bound the window was
/// built from — a deleted changefeed entry cannot be recovered.
ChangefeedRetention {
ns: NamespaceId,
db: DatabaseId,
/// The encoded watermark this window's upper bound was built from.
watermark: &'a [u8],
},
}
/// One bounded range delete.
pub(crate) struct PagedDelete<'a> {
/// The keys to delete. A caller resuming from a durable cursor passes a
/// window already advanced past it.
pub(crate) window: Range<Key>,
pub(crate) companion: PageCompanion<'a>,
/// Whether each key is cleared of every MVCC version rather than
/// soft-deleted.
pub(crate) expunge: bool,
/// Maximum keys deleted per committed page.
pub(crate) page: u32,
}
impl PagedDelete<'_> {
/// Delete the window in committed pages, spending at most `budget` keys.
///
/// `budget` is decremented by what was deleted, so a caller sharing one
/// pass across several ranges can hand the remainder to the next.
pub(crate) async fn run(
mut self,
ds: &Datastore,
lh: &LeaseHandler,
canceller: &CancellationToken,
budget: &mut u64,
) -> Result<PagedOutcome> {
loop {
Datastore::ensure_not_cancelled(canceller)?;
if *budget == 0 {
// Whether the range is finished is a question about the range,
// not about the budget: the page that spent the last of it can
// also be the page that emptied the range, and a full page gives
// no sign either way. One single-key read settles it, so a
// drained range is reported finished now instead of holding a
// claim that names nothing until a later pass looks. Reached at
// most once per range, and only where the last page landed
// exactly on the budget — a short page has already answered it
// below.
let txn = ds.transaction(Read, Optimistic).await?;
let rest = txn.keys(self.window.clone(), 1, 0, None).await;
let _ = txn.cancel().await;
return Ok(match rest?.is_empty() {
true => PagedOutcome::Complete,
false => PagedOutcome::Incomplete,
});
}
// A delete can span many passes, so handing the range to the next
// lease holder costs nothing: every page so far is committed and the
// next holder resumes from what survives.
if !lh.try_maintain_lease().await? {
return Ok(PagedOutcome::LeaseLost);
}
// At least one key however small the page: a zero limit scans
// nothing, and an empty page is how this loop reports a drained
// range, which would retire a claim over data still present.
let limit = (*budget).min(self.page.max(1) as u64) as u32;
let txn = ds.transaction(Write, Optimistic).await?;
if let PageCompanion::ChangefeedRetention {
ns,
db,
watermark,
} = &self.companion
{
// The retention that made this window stale is read inside every
// page transaction, so a policy extended before this read stops
// the delete here rather than deleting entries the new policy
// keeps. Reading the catalog here is also what arms the
// write-conflict check on a conflict-serializing backend, for a
// policy committed between this read and the commit below; on a
// last-writer-wins backend the window is only as fresh as this
// read.
let covered = catch!(
txn,
crate::cf::gc::retention_still_reaches(&txn, *ns, *db, watermark).await
);
if !covered {
let _ = txn.cancel().await;
return Ok(PagedOutcome::Cancelled);
}
}
// One bounded page, scanned inside the transaction that deletes it
// so the page and its deletions cannot disagree. Keys only: a delete
// needs no value.
//
// Ascending order is load-bearing, so this is `keys` and not the
// reverse-order `keysr`: the resume cursor below is the page's last
// key, which only covers everything scanned so far when the page
// holds the window's lowest keys. Under a descending scan the cursor
// would be the page's smallest key and advancing past it would skip
// the rest of the window permanently.
let keys = catch!(txn, txn.keys(self.window.clone(), limit, 0, None).await);
// A page the backend could not fill is the end of the range, so the
// page that empties a range also reports it. Waiting for a following
// empty page instead would spend a whole transaction and scan per
// range to learn nothing, and — where that page is also the one that
// spends the budget — would report a range as unfinished after
// emptying it, keeping a claim naming nothing for another tick.
let exhausted = keys.len() < limit as usize;
let Some(last) = keys.last().cloned() else {
let _ = txn.cancel().await;
return Ok(PagedOutcome::Complete);
};
for key in &keys {
if self.expunge {
catch!(txn, txn.clr(key).await);
} else {
catch!(txn, txn.del(key).await);
}
}
let advanced = if let PageCompanion::ReclaimEntry {
rc,
state,
} = &self.companion
{
let advanced = ReclaimState {
observed_ms: state.observed_ms,
cursor: Some(last.clone()),
};
// The queue entry is the claim authorising these deletes. Advance
// only the exact state this page started from: unlike a blind set,
// this neither resurrects a cancelled claim nor overwrites a cursor
// committed by another node on a last-writer-wins backend.
match txn.putc(*rc, &advanced, Some(state)).await {
Ok(()) => Some(advanced),
Err(e) if crate::kvs::ds::is_conditional_write_conflict(&e) => {
let _ = txn.cancel().await;
return Ok(PagedOutcome::Cancelled);
}
Err(e) => {
let _ = txn.cancel().await;
return Err(e);
}
}
} else {
None
};
match txn.commit().await {
Ok(()) => {}
Err(e) if crate::kvs::ds::is_conditional_write_conflict(&e) => {
let _ = txn.cancel().await;
return Ok(PagedOutcome::Cancelled);
}
Err(e) => {
let _ = txn.cancel().await;
return Err(e);
}
}
if let (
Some(advanced),
PageCompanion::ReclaimEntry {
state,
..
},
) = (advanced, &mut self.companion)
{
*state = advanced;
}
*budget = budget.saturating_sub(keys.len() as u64);
if exhausted {
return Ok(PagedOutcome::Complete);
}
// Resume just after the last key this page deleted. A trailing zero
// byte is the successor of `last` in unsigned byte order, so the
// next page starts at the first key beyond it without re-reading it.
let mut start = last.clone();
start.push(0);
self.window = start..self.window.end.clone();
yield_now!();
}
}
}