Skip to content

Commit 42fbbca

Browse files
committed
Make in-process loglet readers obey readable tails
1 parent 7667e61 commit 42fbbca

4 files changed

Lines changed: 37 additions & 99 deletions

File tree

crates/bifrost/src/providers/local_loglet/mod.rs

Lines changed: 7 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -138,8 +138,7 @@ impl Loglet for LocalLoglet {
138138
from: LogletOffset,
139139
) -> Result<SendableLogletReadStream, OperationError> {
140140
let readable_tail = OffsetWatch::default();
141-
let read_stream =
142-
LocalLogletReadStream::create(self, filter, from, readable_tail.clone()).await?;
141+
let read_stream = LocalLogletReadStream::create(self, filter, from, readable_tail.clone())?;
143142
Ok(SendableLogletReadStream::new(read_stream, readable_tail))
144143
}
145144

@@ -389,17 +388,19 @@ mod tests {
389388
("record-1", Keys::Single(1)).into(),
390389
("record-2", Keys::Single(2)).into(),
391390
("record-3", Keys::Single(1)).into(),
391+
("record-4", Keys::Single(2)).into(),
392392
]
393393
.into();
394394
let offset = loglet.enqueue_batch(batch).await?.await?;
395395

396396
let key_filter = KeyFilter::Include(1);
397-
let read_stream = loglet
397+
let mut read_stream = loglet
398398
.create_read_stream(key_filter, LogletOffset::OLDEST)
399399
.await?;
400400
read_stream.notify_readable_tail(offset.next());
401401

402402
let records: Vec<_> = read_stream
403+
.by_ref()
403404
.take(2)
404405
.try_collect::<Vec<_>>()
405406
.await?
@@ -419,6 +420,9 @@ mod tests {
419420
eq((LogletOffset::from(3), "record-3".to_owned()))
420421
]
421422
);
423+
let filtered = read_stream.next().await.unwrap()?;
424+
assert_that!(filtered.kind(), eq(crate::RecordKind::Filtered));
425+
assert_that!(read_stream.read_pointer(), eq(offset.next()));
422426

423427
Ok(())
424428
}

crates/bifrost/src/providers/local_loglet/read_stream.rs

Lines changed: 16 additions & 49 deletions
Original file line numberDiff line numberDiff line change
@@ -20,9 +20,9 @@ use tracing::{debug, error, warn};
2020

2121
use restate_core::ShutdownError;
2222
use restate_rocksdb::RocksDbReadPerfGuard;
23-
use restate_types::logs::{KeyFilter, LogletOffset, OffsetWatch, SequenceNumber, TailState};
23+
use restate_types::logs::{KeyFilter, LogletOffset, OffsetWatch, SequenceNumber};
2424

25-
use crate::loglet::{Loglet, LogletReadStream, OperationError};
25+
use crate::loglet::{LogletReadStream, OperationError};
2626
use crate::providers::local_loglet::LogStoreError;
2727
use crate::providers::local_loglet::record_format::decode_and_filter_record;
2828
use crate::{LogEntry, Result};
@@ -37,12 +37,9 @@ pub(crate) struct LocalLogletReadStream {
3737
serde_buffer: BytesMut,
3838
// the next record this stream will attempt to read
3939
read_pointer: LogletOffset,
40-
/// stop when read_pointer is at or beyond this offset
41-
last_known_tail: LogletOffset,
4240
readable_tail_watch: BoxStream<'static, LogletOffset>,
4341
readable_tail: LogletOffset,
4442
iterator: DBRawIteratorWithThreadMode<'static, DB>,
45-
tail_watch: BoxStream<'static, TailState<LogletOffset>>,
4643
terminated: bool,
4744
// IMPORTANT: Do not reorder, this should be dropped last since `iterator` holds a reference
4845
// into the underlying database.
@@ -64,7 +61,7 @@ unsafe fn ignore_iterator_lifetime<'a>(
6461
}
6562

6663
impl LocalLogletReadStream {
67-
pub(crate) async fn create(
64+
pub(crate) fn create(
6865
loglet: Arc<LocalLoglet>,
6966
filter: KeyFilter,
7067
from_offset: LogletOffset,
@@ -95,12 +92,6 @@ impl LocalLogletReadStream {
9592
);
9693

9794
let log_store = &loglet.log_store;
98-
let mut tail_watch = loglet.watch_tail();
99-
let last_known_tail = tail_watch
100-
.next()
101-
.await
102-
.expect("loglet watch returns tail pointer")
103-
.offset();
10495

10596
// ## Safety:
10697
// the iterator is guaranteed to be dropped before the loglet is dropped, we hold to the
@@ -123,8 +114,6 @@ impl LocalLogletReadStream {
123114
read_pointer: from_offset,
124115
iterator: iter,
125116
terminated: false,
126-
tail_watch,
127-
last_known_tail,
128117
readable_tail_watch: Box::pin(readable_tail.to_stream()),
129118
readable_tail: readable_tail.get(),
130119
})
@@ -154,8 +143,15 @@ impl Stream for LocalLogletReadStream {
154143
}
155144

156145
let perf_guard = RocksDbReadPerfGuard::new("local-loglet-next");
146+
let mut filtered_from = None;
157147
loop {
158148
if self.read_pointer >= self.readable_tail {
149+
if let Some(filtered_from) = filtered_from {
150+
return Poll::Ready(Some(Ok(LogEntry::new_filtered_gap(
151+
filtered_from,
152+
self.read_pointer.prev(),
153+
))));
154+
}
159155
let maybe_readable_tail = match self.readable_tail_watch.poll_next_unpin(cx) {
160156
Poll::Ready(tail) => tail,
161157
Poll::Pending => {
@@ -174,36 +170,6 @@ impl Stream for LocalLogletReadStream {
174170
}
175171
}
176172
}
177-
// Are we reading after commit offset?
178-
// We are at tail. We need to wait until new records have been released.
179-
if self.read_pointer >= self.last_known_tail {
180-
let maybe_tail_state = match self.tail_watch.poll_next_unpin(cx) {
181-
Poll::Ready(t) => t,
182-
Poll::Pending => {
183-
perf_guard.forget();
184-
return Poll::Pending;
185-
}
186-
};
187-
188-
match maybe_tail_state {
189-
Some(tail_state) => {
190-
// tail has been updated.
191-
self.last_known_tail = tail_state.offset();
192-
continue;
193-
}
194-
None => {
195-
// system shutdown. Or that the loglet has been unexpectedly shutdown.
196-
self.terminated = true;
197-
return Poll::Ready(Some(Err(OperationError::Shutdown(ShutdownError))));
198-
}
199-
}
200-
}
201-
// tail has been updated.
202-
let last_known_tail = self.last_known_tail;
203-
204-
// assert that we are behind tail
205-
assert!(last_known_tail > self.read_pointer);
206-
207173
// Trim point is the slot **before** the first readable record (if it exists)
208174
// trim point might have been updated since last time.
209175
let trim_point =
@@ -213,10 +179,11 @@ impl Stream for LocalLogletReadStream {
213179
assert!(self.read_pointer > LogletOffset::from(0));
214180

215181
if self.read_pointer < head_offset {
216-
let trim_gap = LogEntry::new_trim_gap(self.read_pointer, trim_point);
182+
let gap_to = trim_point.min(self.readable_tail.prev());
183+
let trim_gap = LogEntry::new_trim_gap(self.read_pointer, gap_to);
217184
// next record should be beyond at the head
218-
self.read_pointer = head_offset;
219-
let key = RecordKey::new(self.loglet_id, trim_point);
185+
self.read_pointer = gap_to.next();
186+
let key = RecordKey::new(self.loglet_id, gap_to);
220187
// park the iterator at the trim point, next iteration will seek it forward.
221188
let key_bytes = key.encode_and_split(&mut self.serde_buffer);
222189
self.iterator.seek(key_bytes);
@@ -251,7 +218,7 @@ impl Stream for LocalLogletReadStream {
251218
loglet_id = self.loglet_id,
252219
read_pointer = %self.read_pointer,
253220
trim_point = %potentially_different_trim_point,
254-
last_known_tail = %self.last_known_tail,
221+
readable_tail = %self.readable_tail,
255222
"poll_next() has moved to a non-existent record, that should not happen!"
256223
);
257224
panic!("poll_next() has moved to a non-existent record, that should not happen!");
@@ -283,7 +250,7 @@ impl Stream for LocalLogletReadStream {
283250
if let Some(record) = maybe_record {
284251
return Poll::Ready(Some(Ok(LogEntry::new_data(key.offset, record))));
285252
}
286-
// Didn't match, loop and read the next record if possible.
253+
filtered_from.get_or_insert(key.offset);
287254
}
288255
}
289256
}

crates/bifrost/src/providers/memory_loglet.rs

Lines changed: 13 additions & 40 deletions
Original file line numberDiff line numberDiff line change
@@ -200,34 +200,22 @@ struct MemoryReadStream {
200200
filter: KeyFilter,
201201
/// The next offset to read from
202202
read_pointer: LogletOffset,
203-
tail_watch: BoxStream<'static, TailState<LogletOffset>>,
204-
/// stop when read_pointer is at or beyond this offset
205-
last_known_tail: LogletOffset,
206203
readable_tail_watch: BoxStream<'static, LogletOffset>,
207204
readable_tail: LogletOffset,
208205
terminated: bool,
209206
}
210207

211208
impl MemoryReadStream {
212-
async fn create(
209+
fn create(
213210
loglet: Arc<MemoryLoglet>,
214211
filter: KeyFilter,
215212
from_offset: LogletOffset,
216213
readable_tail: OffsetWatch,
217214
) -> Self {
218-
let mut tail_watch = loglet.watch_tail();
219-
let last_known_tail = tail_watch
220-
.next()
221-
.await
222-
.expect("loglet watch returns tail pointer")
223-
.offset();
224-
225215
Self {
226216
loglet,
227217
filter,
228218
read_pointer: from_offset,
229-
tail_watch,
230-
last_known_tail,
231219
readable_tail_watch: Box::pin(readable_tail.to_stream()),
232220
readable_tail: readable_tail.get(),
233221
terminated: false,
@@ -257,10 +245,17 @@ impl Stream for MemoryReadStream {
257245
return Poll::Ready(None);
258246
}
259247

248+
let mut filtered_from = None;
260249
loop {
261250
let next_offset = self.read_pointer;
262251

263252
if next_offset >= self.readable_tail {
253+
if let Some(filtered_from) = filtered_from {
254+
return Poll::Ready(Some(Ok(LogEntry::new_filtered_gap(
255+
filtered_from,
256+
next_offset.prev(),
257+
))));
258+
}
264259
match ready!(self.readable_tail_watch.poll_next_unpin(cx)) {
265260
Some(readable_tail) => {
266261
self.readable_tail = readable_tail;
@@ -273,28 +268,6 @@ impl Stream for MemoryReadStream {
273268
}
274269
}
275270

276-
// Are we reading after commit offset?
277-
// We are at tail. We need to wait until new records have been released.
278-
if next_offset >= self.last_known_tail {
279-
match ready!(self.tail_watch.poll_next_unpin(cx)) {
280-
Some(tail_state) => {
281-
self.last_known_tail = tail_state.offset();
282-
continue;
283-
}
284-
None => {
285-
// system shutdown. Or that the loglet has been unexpectedly shutdown.
286-
self.terminated = true;
287-
return Poll::Ready(Some(Err(OperationError::Shutdown(ShutdownError))));
288-
}
289-
}
290-
}
291-
292-
// tail has been updated.
293-
let last_known_tail = self.last_known_tail;
294-
295-
// assert that we are behind tail
296-
assert!(last_known_tail > next_offset);
297-
298271
// Trim point is the the slot **before** the first readable record (if it exists)
299272
// trim point might have been updated since last time.
300273
let trim_point =
@@ -304,9 +277,10 @@ impl Stream for MemoryReadStream {
304277
// Are we reading behind the loglet head? -> TrimGap
305278
assert!(next_offset > LogletOffset::from(0));
306279
if next_offset < head_offset {
307-
let trim_gap = LogEntry::new_trim_gap(next_offset, trim_point);
280+
let gap_to = trim_point.min(self.readable_tail.prev());
281+
let trim_gap = LogEntry::new_trim_gap(next_offset, gap_to);
308282
// next record should be beyond at the head
309-
self.read_pointer = head_offset;
283+
self.read_pointer = gap_to.next();
310284
return Poll::Ready(Some(Ok(trim_gap)));
311285
}
312286

@@ -324,8 +298,7 @@ impl Stream for MemoryReadStream {
324298
if let Some(data_record) = next_record.as_record()
325299
&& !data_record.matches_key_query(&self.filter)
326300
{
327-
// read_pointer is already advanced, just don't return the
328-
// record and fast-forward.
301+
filtered_from.get_or_insert(next_offset);
329302
continue;
330303
}
331304

@@ -350,7 +323,7 @@ impl Loglet for MemoryLoglet {
350323
from: LogletOffset,
351324
) -> Result<SendableLogletReadStream, OperationError> {
352325
let readable_tail = OffsetWatch::default();
353-
let read_stream = MemoryReadStream::create(self, filter, from, readable_tail.clone()).await;
326+
let read_stream = MemoryReadStream::create(self, filter, from, readable_tail.clone());
354327
Ok(SendableLogletReadStream::new(read_stream, readable_tail))
355328
}
356329

crates/bifrost/src/read_stream.rs

Lines changed: 1 addition & 7 deletions
Original file line numberDiff line numberDiff line change
@@ -212,13 +212,7 @@ impl LogReadStream {
212212
/// The read pointer points to the next LSN will be attempted on the next
213213
/// `poll_next()`.
214214
fn calculate_read_pointer(record: &LogEntry) -> Lsn {
215-
// On trim gaps, we fast-forward the read pointer beyond the end of the gap. We do
216-
// this after delivering a TrimGap record. This means that the next read operation
217-
// skips over the boundary of the gap.
218-
record
219-
.trim_gap_to_sequence_number()
220-
.unwrap_or_else(|| record.sequence_number())
221-
.next()
215+
record.next_sequence_number()
222216
}
223217

224218
pub fn safe_known_tail(&self) -> Option<Lsn> {

0 commit comments

Comments
 (0)