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
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
use std::cmp::{max, min};
use std::sync::atomic::Ordering;
use std::thread::JoinHandle;
use crate::common::sort_utils::sort_permutation;
use log::trace;
use parking_lot::{RwLock, RwLockReadGuard};
use crate::segment::common::operation_error::{OperationError, OperationResult};
use crate::segment::entry::StorageSegmentEntry;
use crate::segment::types::SeqNumberType;
use crate::shard::segment_holder::{SegmentHolder, SegmentId};
/// How [`SegmentHolder::flush_all`] performs the flush.
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
pub enum FlushMode {
/// Flush all segments on the current thread and wait for completion.
Sync,
/// Spawn a background flush thread. If one is already running, skip this pass and return the
/// current persisted version.
Background,
}
impl SegmentHolder {
/// Flushes all segments and returns maximum version to persist
///
/// Before flushing, this read-locks all segments. It prevents writes to all not-yet-flushed
/// segments during flushing. All locked segments are flushed and released one by one.
///
/// If there are unsaved changes after flush - detects lowest unsaved change version.
/// If all changes are saved - returns max version.
pub fn flush_all(&self, mode: FlushMode, force: bool) -> OperationResult<SeqNumberType> {
let sync = match mode {
FlushMode::Sync => true,
FlushMode::Background => false,
};
let lock_order: Vec<_> = self.non_appendable_then_appendable_segments_ids().collect();
// Grab and keep to segment RwLock's until the end of this function
let segments = self.read_segment_locks(lock_order.iter().cloned())?;
// We can never have zero segments
// Having zero segments could permanently corrupt the WAL by acknowledging u64::MAX
assert!(
!segments.is_empty(),
"must always have at least one segment",
);
// Read-lock all segments before flushing any, must prevent any writes to any segment
// That is to prevent any copy-on-write operation on two segments from occurring in between
// flushing the two segments. If that would happen, segments could end up in an
// inconsistent state on disk that is not recoverable after a crash.
//
// E.g.: we have a point on an immutable segment. If we use a set-payload operation, we do
// copy-on-write. The point from immutable segment A is deleted, the updated point is
// stored on appendable segment B.
// Because of flush ordering segment B (appendable) is flushed before segment A
// (not-appendable). If the copy-on-write operation happens in between, the point is
// deleted from A, but the new point in B is not persisted. We cannot recover this by
// replaying the WAL in case of a crash because the point in A does not exist anymore,
// making copy-on-write impossible.
// Locking all segments prevents copy-on-write operations from occurring in between
// flushes.
//
// WARNING: Ordering is very important here. Specifically:
// - We MUST lock non-appendable first, then appendable.
// - We MUST flush according to the copy-on-write dependency graph.
// Because of this, two rev(erse) calls are used below here.
//
// Locking must happen in this order because `apply_points_to_appendable` can take two
// write locks, also in this order. If we'd use different ordering we will eventually end
// up with a deadlock.
let mut segment_reads: Vec<_> = segments
.iter()
.map(|segment| Self::aloha_lock_segment_read(segment))
.collect();
let max_applied_version = segment_reads.iter().map(|s| s.version()).max().unwrap_or(0);
if !sync && self.is_background_flushing() {
// There is already a background flush ongoing, return current max persisted version.
// Cap by pending post-flush actions like the main path below, but leave running them
// to the next full pass.
let persisted_version = self.get_max_persisted_version(segment_reads, lock_order);
return Ok(match self.pending_post_flush_ack_cap() {
Some(ack_cap) => min(persisted_version, ack_cap),
None => persisted_version,
});
}
// This lock also prevents multiple parallel sync flushes
// as it is exclusive
let mut background_flush_lock = self.lock_flushing()?;
sort_permutation(&mut segment_reads, &lock_order, |segment_ids| {
self.sort_segment_ids_by_flush_dependency(segment_ids)
});
// Capture all flushers first to improve data consistency
let flushers: Vec<_> = segment_reads
.iter()
.filter_map(|read_segment| read_segment.flusher(force))
.collect();
if sync {
for flusher in flushers {
flusher()?;
}
self.flush_dependency
.lock()
.retain(|_, _, version| *version > max_applied_version);
} else {
let flush_dependency = self.flush_dependency.clone();
*background_flush_lock = Some(
std::thread::Builder::new()
.name("background_flush".to_string())
.spawn(move || {
for flusher in flushers {
flusher()?;
}
flush_dependency
.lock()
.retain(|_, _, version| *version > max_applied_version);
Ok(())
})
.unwrap(),
);
}
// Persisted versions at this point reflect completed flushes only (an ongoing background
// pass was joined by `lock_flushing` above; a freshly spawned one has not updated them
// yet), so the result is a durable waterline: every segment's state up to it is on disk.
// Post-flush actions scheduled at or below it can run now, since the data they clean up is
// durable in its new home by now.
//
// Actions that are not yet ready cap the returned version (and with it the WAL
// acknowledge): the data they have not cleaned up yet contradicts operations past their
// pin, so those operations must stay replayable until the action runs. See
// [`SegmentHolder::register_post_flush_action`].
let persisted_version = self.get_max_persisted_version(segment_reads, lock_order);
let pending_ack_cap = self.run_ready_post_flush_actions(persisted_version)?;
Ok(match pending_ack_cap {
Some(ack_cap) => min(persisted_version, ack_cap),
None => persisted_version,
})
}
fn non_appendable_then_appendable_segments_ids(&self) -> impl Iterator<Item = SegmentId> {
let non_appendable_segments = self.non_appendable_segments_ids();
let appendable_segments = self.appendable_segments_ids();
non_appendable_segments
.into_iter()
.chain(appendable_segments)
}
fn sort_segment_ids_by_flush_dependency(&self, segment_ids: &[SegmentId]) -> Vec<SegmentId> {
let flush_topology = self.flush_dependency.lock().clone();
trace!("{flush_topology:?}");
let mut iter = flush_topology.sort_elements(segment_ids);
let sorted_keys: Vec<_> = iter.by_ref().collect();
// Collect any remaining elements (circular dependencies)
let remaining = iter.into_unordered_vec();
debug_assert!(
remaining.is_empty(),
"circular dependencies detected in flush topology: {remaining:?}",
);
#[cfg(feature = "staging")]
if !remaining.is_empty() {
log::warn!(
"circular dependencies detected in flush topology for segments: {remaining:?}"
);
}
// Build combined order: sorted first, then unordered remainder
let mut final_order: Vec<_> = sorted_keys;
final_order.extend(remaining);
final_order
}
/// Run `f` serialized with the flush pipeline.
///
/// Flusher executions must be serialized end-to-end with their capture: component
/// flushers capture pending state by value (e.g. Gridstore clones its pending update
/// set) and their execution frees storage superseded by that state. Running an extra
/// flush between another flusher's capture and its execution makes the pending sets
/// overlap, and the later execution double-frees blocks that may already have been
/// reused — corrupting the storage. This joins any in-flight background flush (whose
/// flushers were captured earlier) and holds the flush lock while `f` runs, so a
/// flush performed inside `f` cannot interleave with a captured-but-unrun pass.
///
/// Deadlock note: callers may hold segment locks; `flush_all` acquires segment locks
/// before this lock, and flush executions take no segment locks, so the ordering
/// [segment locks → flush lock] is consistent.
pub fn with_flush_serialized<T>(
&self,
f: impl FnOnce() -> OperationResult<T>,
) -> OperationResult<T> {
let _flush_lock = self.lock_flushing()?;
f()
}
// Joins flush thread if exists
// Returns lock to guarantee that there will be no other flush in a different thread
pub(super) fn lock_flushing(
&self,
) -> OperationResult<parking_lot::MutexGuard<'_, Option<JoinHandle<OperationResult<()>>>>> {
let mut lock = self.flush_thread.lock();
let mut join_handle: Option<JoinHandle<OperationResult<()>>> = None;
std::mem::swap(&mut join_handle, &mut lock);
if let Some(join_handle) = join_handle {
// Flush result was reported to segment, so we don't need this value anymore
let flush_result = join_handle
.join()
.map_err(|_err| OperationError::service_error("failed to join flush thread"))?;
flush_result.map_err(|err| {
OperationError::service_error(format!("last background flush failed: {err}"))
})?;
}
Ok(lock)
}
pub(super) fn is_background_flushing(&self) -> bool {
let lock = self.flush_thread.lock();
if let Some(join_handle) = lock.as_ref() {
!join_handle.is_finished()
} else {
false
}
}
/// Calculates the version of the segments that is safe to acknowledge in WAL
///
/// If there are unsaved changes after flush - detects lowest unsaved change version.
/// If all changes are saved - returns max version.
fn get_max_persisted_version(
&self,
segment_reads: Vec<RwLockReadGuard<'_, dyn StorageSegmentEntry>>,
lock_order: Vec<SegmentId>,
) -> SeqNumberType {
// Start with the max_persisted_vesrion at the set overwrite value, which may just be 0
// Any of the segments we flush may increase this if they have a higher persisted version
// The overwrite is required to ensure we acknowledge no-op operations in WAL that didn't hit any segment
//
// Only affects returned version if all changes are saved
let mut max_persisted_version: SeqNumberType = self
.max_persisted_segment_version_overwrite
.load(Ordering::Relaxed);
let mut min_unsaved_version: SeqNumberType = SeqNumberType::MAX;
let mut has_unsaved = false;
for (read_segment, segment_id) in segment_reads.into_iter().zip(lock_order) {
let segment_version = read_segment.version();
let segment_persisted_version = read_segment.persistent_version();
log::trace!(
"Flushed segment {segment_id}:{:?} version: {segment_version} to persisted: {segment_persisted_version}",
read_segment.data_path(),
);
if segment_version > segment_persisted_version {
has_unsaved = true;
min_unsaved_version = min(min_unsaved_version, segment_persisted_version);
}
max_persisted_version = max(max_persisted_version, segment_persisted_version);
drop(read_segment);
}
if has_unsaved {
log::trace!(
"Some segments have unsaved changes, lowest unsaved version: {min_unsaved_version}"
);
min_unsaved_version
} else {
log::trace!(
"All segments flushed successfully, max persisted version: {max_persisted_version}"
);
max_persisted_version
}
}
/// Grab the RwLock's for all the given segment IDs.
fn read_segment_locks(
&self,
segment_ids: impl IntoIterator<Item = SegmentId>,
) -> OperationResult<Vec<&RwLock<dyn StorageSegmentEntry>>> {
segment_ids
.into_iter()
.map(|segment_id| {
self.get(segment_id)
.ok_or_else(|| {
OperationError::service_error(format!("No segment with ID {segment_id}"))
})
.map(|segment| segment.get() as _)
})
.collect()
}
}