1use std::collections::HashMap;
2use std::collections::HashSet;
3use std::collections::VecDeque;
4use std::future::Future;
5use std::io;
6use std::pin::Pin;
7use std::sync::Arc;
8use std::sync::Mutex;
9
10use ursula_shard::BucketStreamId;
11use ursula_stream::ColdChunkRef;
12use ursula_stream::ExternalPayloadRef;
13use ursula_stream::ObjectPayloadRef;
14use ursula_stream::StreamReadColdIndexSegment;
15
16use crate::cold_store::ColdStoreHandle;
17
18pub type ColdIndexPageStoreFuture<'a, T> = Pin<Box<dyn Future<Output = io::Result<T>> + Send + 'a>>;
19
20const COLD_INDEX_PAGE_MAGIC: &[u8; 8] = b"UCIDX001";
21const COLD_INDEX_PAGE_VERSION: u16 = 2;
22const COLD_INDEX_ENTRY_COLD_CHUNK: u8 = 1;
23const COLD_INDEX_ENTRY_EXTERNAL_SEGMENT: u8 = 2;
24const FNV64_OFFSET_BASIS: u64 = 0xcbf2_9ce4_8422_2325;
25const FNV64_PRIME: u64 = 0x0000_0100_0000_01b3;
26
27#[derive(Debug, Clone, PartialEq, Eq, Hash)]
28pub struct ColdIndexPageKey {
29 pub stream_id: BucketStreamId,
30 pub generation: u64,
31 pub page_id: u64,
32}
33
34impl ColdIndexPageKey {
35 pub fn path(&self) -> String {
36 format!(
37 "{}/cold-index/{:020}/{:020}.idx",
38 self.stream_id, self.generation, self.page_id
39 )
40 }
41}
42
43pub fn cold_index_prefix(stream_id: &BucketStreamId) -> String {
44 format!("{stream_id}/cold-index/")
45}
46
47#[derive(Debug, Clone, PartialEq, Eq)]
48pub struct ColdIndexPage {
49 pub start_offset: u64,
50 pub end_offset: u64,
51 pub cold_chunks: Vec<ColdChunkRef>,
52 pub external_segments: Vec<ObjectPayloadRef>,
53}
54
55impl ColdIndexPage {
56 pub fn covers(&self, offset: u64) -> bool {
57 self.start_offset <= offset && offset < self.end_offset
58 }
59}
60
61#[derive(Debug, Clone)]
62pub struct ColdIndexPageRollback {
63 key: ColdIndexPageKey,
64 previous: Option<ColdIndexPage>,
65 written_chunk: ColdChunkRef,
66}
67
68fn encode_page(key: &ColdIndexPageKey, page: &ColdIndexPage) -> Vec<u8> {
69 let mut body = Vec::new();
70 put_string(&mut body, &key.stream_id.bucket_id);
71 match &key.stream_id.affinity_key {
72 Some(affinity_key) => {
73 put_u8(&mut body, 1);
74 put_string(&mut body, affinity_key);
75 }
76 None => put_u8(&mut body, 0),
77 }
78 put_string(&mut body, &key.stream_id.stream_id);
79 put_u64(&mut body, key.generation);
80 put_u64(&mut body, key.page_id);
81 put_u64(&mut body, page.start_offset);
82 put_u64(&mut body, page.end_offset);
83 put_u32(
84 &mut body,
85 u32::try_from(page.cold_chunks.len()).expect("cold index cold chunk count fits u32"),
86 );
87 for chunk in &page.cold_chunks {
88 put_u8(&mut body, COLD_INDEX_ENTRY_COLD_CHUNK);
89 put_u64(&mut body, chunk.start_offset);
90 put_u64(&mut body, chunk.end_offset);
91 put_u64(&mut body, chunk.object_size);
92 put_string(&mut body, &chunk.s3_path);
93 }
94 put_u32(
95 &mut body,
96 u32::try_from(page.external_segments.len())
97 .expect("cold index external segment count fits u32"),
98 );
99 for object in &page.external_segments {
100 put_u8(&mut body, COLD_INDEX_ENTRY_EXTERNAL_SEGMENT);
101 put_u64(&mut body, object.start_offset);
102 put_u64(&mut body, object.end_offset);
103 put_u64(&mut body, object.object_size);
104 put_string(&mut body, &object.s3_path);
105 }
106
107 let mut bytes = Vec::with_capacity(COLD_INDEX_PAGE_MAGIC.len() + 2 + 4 + body.len() + 8);
108 bytes.extend_from_slice(COLD_INDEX_PAGE_MAGIC);
109 put_u16(&mut bytes, COLD_INDEX_PAGE_VERSION);
110 put_u32(
111 &mut bytes,
112 u32::try_from(body.len()).expect("cold index page body len fits u32"),
113 );
114 bytes.extend_from_slice(&body);
115 put_u64(&mut bytes, checksum64(&body));
116 bytes
117}
118
119fn decode_page(key: &ColdIndexPageKey, bytes: &[u8]) -> io::Result<ColdIndexPage> {
120 let mut cursor = Cursor::new(bytes);
121 let magic = cursor.read_exact(COLD_INDEX_PAGE_MAGIC.len())?;
122 if magic != COLD_INDEX_PAGE_MAGIC {
123 return Err(io::Error::new(
124 io::ErrorKind::InvalidData,
125 "cold index page has invalid magic",
126 ));
127 }
128 let version = cursor.read_u16()?;
129 if !matches!(version, 1 | COLD_INDEX_PAGE_VERSION) {
130 return Err(io::Error::new(
131 io::ErrorKind::InvalidData,
132 format!("unsupported cold index page version {version}"),
133 ));
134 }
135 let body_len = usize::try_from(cursor.read_u32()?).expect("u32 fits usize");
136 let body = cursor.read_exact(body_len)?;
137 let expected_checksum = cursor.read_u64()?;
138 if cursor.remaining() != 0 {
139 return Err(io::Error::new(
140 io::ErrorKind::InvalidData,
141 "cold index page has trailing bytes",
142 ));
143 }
144 let actual_checksum = checksum64(body);
145 if actual_checksum != expected_checksum {
146 return Err(io::Error::new(
147 io::ErrorKind::InvalidData,
148 "cold index page checksum mismatch",
149 ));
150 }
151
152 let mut body = Cursor::new(body);
153 let bucket_id = body.read_string()?;
154 let affinity_key = if version >= 2 {
155 match body.read_u8()? {
156 0 => None,
157 1 => Some(body.read_string()?),
158 _ => {
159 return Err(io::Error::new(
160 io::ErrorKind::InvalidData,
161 "cold index page has invalid affinity marker",
162 ));
163 }
164 }
165 } else {
166 None
167 };
168 let stream_id = body.read_string()?;
169 let generation = body.read_u64()?;
170 let page_id = body.read_u64()?;
171 if bucket_id != key.stream_id.bucket_id
172 || affinity_key != key.stream_id.affinity_key
173 || stream_id != key.stream_id.stream_id
174 || generation != key.generation
175 || page_id != key.page_id
176 {
177 return Err(io::Error::new(
178 io::ErrorKind::InvalidData,
179 "cold index page key metadata mismatch",
180 ));
181 }
182 let start_offset = body.read_u64()?;
183 let end_offset = body.read_u64()?;
184 let cold_chunk_count = body.read_u32()?;
185 let mut cold_chunks =
186 Vec::with_capacity(usize::try_from(cold_chunk_count).expect("u32 fits usize"));
187 for _ in 0..cold_chunk_count {
188 let tag = body.read_u8()?;
189 if tag != COLD_INDEX_ENTRY_COLD_CHUNK {
190 return Err(io::Error::new(
191 io::ErrorKind::InvalidData,
192 "cold index page expected cold chunk entry",
193 ));
194 }
195 cold_chunks.push(ColdChunkRef {
196 start_offset: body.read_u64()?,
197 end_offset: body.read_u64()?,
198 object_size: body.read_u64()?,
199 s3_path: body.read_string()?,
200 object_offset: 0,
201 shared_object: false,
202 payload_digest: String::new(),
203 });
204 }
205 let external_segment_count = body.read_u32()?;
206 let mut external_segments =
207 Vec::with_capacity(usize::try_from(external_segment_count).expect("u32 fits usize"));
208 for _ in 0..external_segment_count {
209 let tag = body.read_u8()?;
210 if tag != COLD_INDEX_ENTRY_EXTERNAL_SEGMENT {
211 return Err(io::Error::new(
212 io::ErrorKind::InvalidData,
213 "cold index page expected external segment entry",
214 ));
215 }
216 external_segments.push(ObjectPayloadRef {
217 start_offset: body.read_u64()?,
218 end_offset: body.read_u64()?,
219 object_size: body.read_u64()?,
220 s3_path: body.read_string()?,
221 object_offset: 0,
222 });
223 }
224 if body.remaining() != 0 {
225 return Err(io::Error::new(
226 io::ErrorKind::InvalidData,
227 "cold index page body has trailing bytes",
228 ));
229 }
230 Ok(ColdIndexPage {
231 start_offset,
232 end_offset,
233 cold_chunks,
234 external_segments,
235 })
236}
237
238fn put_u8(out: &mut Vec<u8>, value: u8) {
239 out.push(value);
240}
241
242fn put_u16(out: &mut Vec<u8>, value: u16) {
243 out.extend_from_slice(&value.to_le_bytes());
244}
245
246fn put_u32(out: &mut Vec<u8>, value: u32) {
247 out.extend_from_slice(&value.to_le_bytes());
248}
249
250fn put_u64(out: &mut Vec<u8>, value: u64) {
251 out.extend_from_slice(&value.to_le_bytes());
252}
253
254fn put_string(out: &mut Vec<u8>, value: &str) {
255 put_u32(
256 out,
257 u32::try_from(value.len()).expect("cold index string len fits u32"),
258 );
259 out.extend_from_slice(value.as_bytes());
260}
261
262fn checksum64(bytes: &[u8]) -> u64 {
263 let mut hash = FNV64_OFFSET_BASIS;
264 for byte in bytes {
265 hash ^= u64::from(*byte);
266 hash = hash.wrapping_mul(FNV64_PRIME);
267 }
268 hash
269}
270
271struct Cursor<'a> {
272 bytes: &'a [u8],
273 offset: usize,
274}
275
276impl<'a> Cursor<'a> {
277 fn new(bytes: &'a [u8]) -> Self {
278 Self { bytes, offset: 0 }
279 }
280
281 fn remaining(&self) -> usize {
282 self.bytes.len().saturating_sub(self.offset)
283 }
284
285 fn read_exact(&mut self, len: usize) -> io::Result<&'a [u8]> {
286 let end = self.offset.checked_add(len).ok_or_else(|| {
287 io::Error::new(
288 io::ErrorKind::InvalidData,
289 "cold index page offset overflow",
290 )
291 })?;
292 if end > self.bytes.len() {
293 return Err(io::Error::new(
294 io::ErrorKind::UnexpectedEof,
295 "cold index page ended early",
296 ));
297 }
298 let slice = &self.bytes[self.offset..end];
299 self.offset = end;
300 Ok(slice)
301 }
302
303 fn read_u8(&mut self) -> io::Result<u8> {
304 Ok(self.read_exact(1)?[0])
305 }
306
307 fn read_u16(&mut self) -> io::Result<u16> {
308 let mut bytes = [0; 2];
309 bytes.copy_from_slice(self.read_exact(2)?);
310 Ok(u16::from_le_bytes(bytes))
311 }
312
313 fn read_u32(&mut self) -> io::Result<u32> {
314 let mut bytes = [0; 4];
315 bytes.copy_from_slice(self.read_exact(4)?);
316 Ok(u32::from_le_bytes(bytes))
317 }
318
319 fn read_u64(&mut self) -> io::Result<u64> {
320 let mut bytes = [0; 8];
321 bytes.copy_from_slice(self.read_exact(8)?);
322 Ok(u64::from_le_bytes(bytes))
323 }
324
325 fn read_string(&mut self) -> io::Result<String> {
326 let len = usize::try_from(self.read_u32()?).expect("u32 fits usize");
327 let bytes = self.read_exact(len)?;
328 String::from_utf8(bytes.to_vec()).map_err(|err| {
329 io::Error::new(
330 io::ErrorKind::InvalidData,
331 format!("cold index page contains invalid UTF-8: {err}"),
332 )
333 })
334 }
335}
336
337pub trait ColdIndexPageStore: Send + Sync {
338 fn put_page<'a>(
339 &'a self,
340 key: &'a ColdIndexPageKey,
341 page: &'a ColdIndexPage,
342 ) -> ColdIndexPageStoreFuture<'a, ()>;
343
344 fn get_page<'a>(
345 &'a self,
346 key: &'a ColdIndexPageKey,
347 ) -> ColdIndexPageStoreFuture<'a, Option<ColdIndexPage>>;
348}
349
350pub async fn write_cold_chunk_index_pages<S: ColdIndexPageStore + ?Sized>(
351 store: &S,
352 stream_id: &BucketStreamId,
353 chunk: &ColdChunkRef,
354) -> io::Result<()> {
355 write_cold_chunk_index_pages_with_rollback(store, stream_id, chunk)
356 .await
357 .map(|_| ())
358}
359
360pub async fn write_cold_chunk_index_pages_with_rollback<S: ColdIndexPageStore + ?Sized>(
361 store: &S,
362 stream_id: &BucketStreamId,
363 chunk: &ColdChunkRef,
364) -> io::Result<Vec<ColdIndexPageRollback>> {
365 if chunk.end_offset <= chunk.start_offset {
366 return Ok(Vec::new());
367 }
368 let first_page_id = chunk.start_offset / ursula_stream::COLD_INDEX_PAGE_SPAN_BYTES;
369 let last_page_id = (chunk.end_offset - 1) / ursula_stream::COLD_INDEX_PAGE_SPAN_BYTES;
370 let mut rollback = Vec::new();
371 for page_id in first_page_id..=last_page_id {
372 let key = ColdIndexPageKey {
373 stream_id: stream_id.clone(),
374 generation: 0,
375 page_id,
376 };
377 let page_start = page_id.saturating_mul(ursula_stream::COLD_INDEX_PAGE_SPAN_BYTES);
378 let page_end = page_start.saturating_add(ursula_stream::COLD_INDEX_PAGE_SPAN_BYTES);
379 let previous = store.get_page(&key).await?;
380 let mut page = previous.clone().unwrap_or_else(|| ColdIndexPage {
381 start_offset: page_start,
382 end_offset: page_end,
383 cold_chunks: Vec::new(),
384 external_segments: Vec::new(),
385 });
386 rollback.push(ColdIndexPageRollback {
387 key: key.clone(),
388 previous,
389 written_chunk: chunk.clone(),
390 });
391 page.cold_chunks.retain(|existing| {
392 existing.start_offset != chunk.start_offset || existing.end_offset != chunk.end_offset
393 });
394 page.cold_chunks.push(chunk.clone());
395 page.cold_chunks.sort_by_key(|chunk| chunk.start_offset);
396 store.put_page(&key, &page).await?;
397 }
398 Ok(rollback)
399}
400
401pub async fn rollback_cold_index_pages<S: ColdIndexPageStore + ?Sized>(
402 store: &S,
403 rollback: Vec<ColdIndexPageRollback>,
404) -> io::Result<()> {
405 for entry in rollback.into_iter().rev() {
406 let Some(current) = store.get_page(&entry.key).await? else {
407 continue;
408 };
409 let current_still_has_written_chunk = current.cold_chunks.iter().any(|chunk| {
410 chunk.start_offset == entry.written_chunk.start_offset
411 && chunk.end_offset == entry.written_chunk.end_offset
412 && chunk.s3_path == entry.written_chunk.s3_path
413 });
414 if !current_still_has_written_chunk {
415 continue;
416 }
417 let page = entry.previous.unwrap_or_else(|| {
418 let page_start = entry
419 .key
420 .page_id
421 .saturating_mul(ursula_stream::COLD_INDEX_PAGE_SPAN_BYTES);
422 let page_end = page_start.saturating_add(ursula_stream::COLD_INDEX_PAGE_SPAN_BYTES);
423 ColdIndexPage {
424 start_offset: page_start,
425 end_offset: page_end,
426 cold_chunks: Vec::new(),
427 external_segments: Vec::new(),
428 }
429 });
430 store.put_page(&entry.key, &page).await?;
431 }
432 Ok(())
433}
434
435pub async fn write_external_segment_index_pages<S: ColdIndexPageStore + ?Sized>(
436 store: &S,
437 stream_id: &BucketStreamId,
438 start_offset: u64,
439 payload: &ExternalPayloadRef,
440) -> io::Result<()> {
441 let object = ObjectPayloadRef {
442 start_offset,
443 end_offset: start_offset.saturating_add(payload.payload_len),
444 s3_path: payload.s3_path.clone(),
445 object_size: payload.object_size,
446 object_offset: 0,
447 };
448 write_object_index_pages(store, stream_id, object).await
449}
450
451async fn write_object_index_pages<S: ColdIndexPageStore + ?Sized>(
452 store: &S,
453 stream_id: &BucketStreamId,
454 object: ObjectPayloadRef,
455) -> io::Result<()> {
456 if object.end_offset <= object.start_offset {
457 return Ok(());
458 }
459 let first_page_id = object.start_offset / ursula_stream::COLD_INDEX_PAGE_SPAN_BYTES;
460 let last_page_id = (object.end_offset - 1) / ursula_stream::COLD_INDEX_PAGE_SPAN_BYTES;
461 for page_id in first_page_id..=last_page_id {
462 let key = ColdIndexPageKey {
463 stream_id: stream_id.clone(),
464 generation: 0,
465 page_id,
466 };
467 let page_start = page_id.saturating_mul(ursula_stream::COLD_INDEX_PAGE_SPAN_BYTES);
468 let page_end = page_start.saturating_add(ursula_stream::COLD_INDEX_PAGE_SPAN_BYTES);
469 let mut page = store
470 .get_page(&key)
471 .await?
472 .unwrap_or_else(|| ColdIndexPage {
473 start_offset: page_start,
474 end_offset: page_end,
475 cold_chunks: Vec::new(),
476 external_segments: Vec::new(),
477 });
478 page.external_segments.retain(|existing| {
479 existing.start_offset != object.start_offset || existing.end_offset != object.end_offset
480 });
481 page.external_segments.push(object.clone());
482 page.external_segments
483 .sort_by_key(|object| object.start_offset);
484 store.put_page(&key, &page).await?;
485 }
486 Ok(())
487}
488
489#[derive(Debug, Default)]
490pub struct InMemoryColdIndexPageStore {
491 pages: Mutex<HashMap<ColdIndexPageKey, Vec<u8>>>,
492}
493
494#[derive(Debug, Clone)]
495pub struct ColdStoreColdIndexPageStore {
496 cold_store: ColdStoreHandle,
497}
498
499impl ColdStoreColdIndexPageStore {
500 pub fn new(cold_store: ColdStoreHandle) -> Self {
501 Self { cold_store }
502 }
503}
504
505impl ColdIndexPageStore for ColdStoreColdIndexPageStore {
506 fn put_page<'a>(
507 &'a self,
508 key: &'a ColdIndexPageKey,
509 page: &'a ColdIndexPage,
510 ) -> ColdIndexPageStoreFuture<'a, ()> {
511 Box::pin(async move {
512 let bytes = encode_page(key, page);
513 self.cold_store
514 .write_cold_index_page(&key.path(), &bytes)
515 .await?;
516 Ok(())
517 })
518 }
519
520 fn get_page<'a>(
521 &'a self,
522 key: &'a ColdIndexPageKey,
523 ) -> ColdIndexPageStoreFuture<'a, Option<ColdIndexPage>> {
524 Box::pin(async move {
525 self.cold_store
526 .read_cold_index_page(&key.path())
527 .await?
528 .map(|bytes| decode_page(key, &bytes))
529 .transpose()
530 })
531 }
532}
533
534impl InMemoryColdIndexPageStore {
535 pub fn new() -> Self {
536 Self::default()
537 }
538}
539
540impl ColdIndexPageStore for InMemoryColdIndexPageStore {
541 fn put_page<'a>(
542 &'a self,
543 key: &'a ColdIndexPageKey,
544 page: &'a ColdIndexPage,
545 ) -> ColdIndexPageStoreFuture<'a, ()> {
546 Box::pin(async move {
547 let bytes = encode_page(key, page);
548 self.pages
549 .lock()
550 .expect("cold index page store mutex poisoned")
551 .insert(key.clone(), bytes);
552 Ok(())
553 })
554 }
555
556 fn get_page<'a>(
557 &'a self,
558 key: &'a ColdIndexPageKey,
559 ) -> ColdIndexPageStoreFuture<'a, Option<ColdIndexPage>> {
560 Box::pin(async move {
561 self.pages
562 .lock()
563 .expect("cold index page store mutex poisoned")
564 .get(key)
565 .map(|bytes| decode_page(key, bytes))
566 .transpose()
567 })
568 }
569}
570
571#[derive(Debug)]
572pub struct ColdIndexPageCache<S: ColdIndexPageStore + ?Sized> {
573 store: Arc<S>,
574 capacity_pages: usize,
575 inner: Mutex<ColdIndexPageCacheInner>,
576}
577
578#[derive(Debug, Default)]
579struct ColdIndexPageCacheInner {
580 next_generation: u64,
581 pages: HashMap<ColdIndexPageKey, ColdIndexPageCacheEntry>,
582 lru: VecDeque<(ColdIndexPageKey, u64)>,
583}
584
585#[derive(Debug)]
586struct ColdIndexPageCacheEntry {
587 page: Arc<ColdIndexPage>,
588 generation: u64,
589}
590
591impl<S: ColdIndexPageStore + ?Sized> ColdIndexPageCache<S> {
592 pub fn new(store: Arc<S>, capacity_pages: usize) -> Self {
593 Self {
594 store,
595 capacity_pages,
596 inner: Mutex::new(ColdIndexPageCacheInner::default()),
597 }
598 }
599
600 pub async fn put_page(&self, key: &ColdIndexPageKey, page: &ColdIndexPage) -> io::Result<()> {
601 self.store.put_page(key, page).await?;
602 self.insert(key.clone(), Arc::new(page.clone()));
603 Ok(())
604 }
605
606 pub async fn get_page(&self, key: &ColdIndexPageKey) -> io::Result<Option<Arc<ColdIndexPage>>> {
607 if let Some(page) = self.get_cached(key) {
608 return Ok(Some(page));
609 }
610 self.reload_page(key).await
611 }
612
613 pub fn invalidate_stream(&self, stream_id: &BucketStreamId) {
616 let mut inner = self.inner.lock().expect("cold index cache mutex poisoned");
617 inner.pages.retain(|key, _| &key.stream_id != stream_id);
618 inner.lru.retain(|(key, _)| &key.stream_id != stream_id);
619 }
620
621 async fn reload_page(&self, key: &ColdIndexPageKey) -> io::Result<Option<Arc<ColdIndexPage>>> {
622 let Some(page) = self.store.get_page(key).await? else {
623 return Ok(None);
624 };
625 let page = Arc::new(page);
626 self.insert(key.clone(), page.clone());
627 Ok(Some(page))
628 }
629
630 pub async fn object_segments_for_read(
631 &self,
632 stream_id: &BucketStreamId,
633 segment: &StreamReadColdIndexSegment,
634 ) -> io::Result<Vec<ObjectPayloadRef>> {
635 let key = ColdIndexPageKey {
636 stream_id: stream_id.clone(),
637 generation: segment.generation,
638 page_id: segment.page_id,
639 };
640 let Some(page) = self.get_page(&key).await? else {
641 return Err(io::Error::new(
642 io::ErrorKind::NotFound,
643 format!("cold index page '{}' does not exist", key.path()),
644 ));
645 };
646 let read_end = segment
647 .read_start_offset
648 .checked_add(u64::try_from(segment.len).expect("cold index read len fits u64"))
649 .ok_or_else(|| {
650 io::Error::new(
651 io::ErrorKind::InvalidInput,
652 "cold index read range overflows",
653 )
654 })?;
655 let mut objects = objects_for_read(&page, segment.read_start_offset, read_end);
656 if !objects_cover_range(&objects, segment.read_start_offset, read_end)
657 && let Some(reloaded) = self.reload_page(&key).await?
658 {
659 objects = objects_for_read(&reloaded, segment.read_start_offset, read_end);
660 }
661 if !objects_cover_range(&objects, segment.read_start_offset, read_end) {
662 return Err(io::Error::new(
663 io::ErrorKind::InvalidData,
664 "cold index page does not cover requested read range",
665 ));
666 }
667 Ok(objects)
668 }
669
670 pub fn cached_page_count(&self) -> usize {
671 self.inner
672 .lock()
673 .expect("cold index page cache mutex poisoned")
674 .pages
675 .len()
676 }
677
678 fn get_cached(&self, key: &ColdIndexPageKey) -> Option<Arc<ColdIndexPage>> {
679 let mut inner = self
680 .inner
681 .lock()
682 .expect("cold index page cache mutex poisoned");
683 let page = inner.pages.get(key)?.page.clone();
684 Self::touch(&mut inner, key.clone());
685 Some(page)
686 }
687
688 fn insert(&self, key: ColdIndexPageKey, page: Arc<ColdIndexPage>) {
689 let mut inner = self
690 .inner
691 .lock()
692 .expect("cold index page cache mutex poisoned");
693 let generation = Self::touch(&mut inner, key.clone());
694 inner
695 .pages
696 .insert(key, ColdIndexPageCacheEntry { page, generation });
697 Self::evict_over_capacity(&mut inner, self.capacity_pages);
698 }
699
700 fn touch(inner: &mut ColdIndexPageCacheInner, key: ColdIndexPageKey) -> u64 {
701 let generation = inner.next_generation;
702 inner.next_generation = inner.next_generation.saturating_add(1);
703 if let Some(entry) = inner.pages.get_mut(&key) {
704 entry.generation = generation;
705 }
706 inner.lru.push_back((key, generation));
707 generation
708 }
709
710 fn evict_over_capacity(inner: &mut ColdIndexPageCacheInner, capacity_pages: usize) {
711 if capacity_pages == 0 {
712 inner.pages.clear();
713 inner.lru.clear();
714 return;
715 }
716 while inner.pages.len() > capacity_pages {
717 let Some((key, generation)) = inner.lru.pop_front() else {
718 break;
719 };
720 let stale = inner
721 .pages
722 .get(&key)
723 .is_none_or(|entry| entry.generation != generation);
724 if stale {
725 continue;
726 }
727 inner.pages.remove(&key);
728 }
729 }
730}
731
732pub async fn load_cold_chunks_from_pages<S: ColdIndexPageStore + ?Sized>(
736 store: &S,
737 keys: &[ColdIndexPageKey],
738) -> io::Result<Vec<ColdChunkRef>> {
739 let mut chunks = Vec::new();
740 let mut seen = HashSet::new();
741 for key in keys {
742 let Some(page) = store.get_page(key).await? else {
743 continue;
744 };
745 for chunk in page.cold_chunks {
746 let identity = (chunk.start_offset, chunk.end_offset, chunk.s3_path.clone());
747 if seen.insert(identity) {
748 chunks.push(chunk);
749 }
750 }
751 }
752 chunks.sort_by(|left, right| {
753 left.start_offset
754 .cmp(&right.start_offset)
755 .then_with(|| left.end_offset.cmp(&right.end_offset))
756 .then_with(|| left.s3_path.cmp(&right.s3_path))
757 });
758 Ok(chunks)
759}
760
761pub fn select_cold_chunk_compaction(
764 chunks: &[ColdChunkRef],
765 target_bytes: u64,
766 max_bytes: u64,
767) -> Option<Vec<ColdChunkRef>> {
768 if target_bytes == 0 || max_bytes < target_bytes {
769 return None;
770 }
771 let mut candidate = Vec::new();
772 let mut bytes = 0_u64;
773 let mut next_offset = None;
774 for chunk in chunks {
775 let logical_bytes = chunk.end_offset.checked_sub(chunk.start_offset)?;
776 let usable = logical_bytes > 0
777 && chunk.object_size == logical_bytes
778 && chunk.object_size < target_bytes;
779 let contiguous = next_offset.is_none_or(|offset| offset == chunk.start_offset);
780 let next_bytes = bytes.checked_add(chunk.object_size);
781 if !usable || !contiguous || next_bytes.is_none_or(|total| total > max_bytes) {
782 candidate.clear();
783 bytes = 0;
784 next_offset = None;
785 if !usable {
786 continue;
787 }
788 }
789 bytes = bytes.checked_add(chunk.object_size)?;
790 next_offset = Some(chunk.end_offset);
791 candidate.push(chunk.clone());
792 if candidate.len() >= 2 && bytes >= target_bytes {
793 return Some(candidate);
794 }
795 }
796 None
797}
798
799pub async fn replace_cold_chunk_index_pages<S: ColdIndexPageStore + ?Sized>(
804 store: &S,
805 stream_id: &BucketStreamId,
806 old_chunks: &[ColdChunkRef],
807 replacement: &ColdChunkRef,
808) -> io::Result<bool> {
809 replace_cold_chunk_index_pages_with_rollback(store, stream_id, old_chunks, replacement)
810 .await
811 .map(|rollback| rollback.is_some())
812}
813
814pub async fn replace_cold_chunk_index_pages_with_rollback<S: ColdIndexPageStore + ?Sized>(
815 store: &S,
816 stream_id: &BucketStreamId,
817 old_chunks: &[ColdChunkRef],
818 replacement: &ColdChunkRef,
819) -> io::Result<Option<Vec<ColdIndexPageRollback>>> {
820 if old_chunks.len() < 2 || replacement.end_offset <= replacement.start_offset {
821 return Ok(None);
822 }
823 let first_page_id = replacement.start_offset / ursula_stream::COLD_INDEX_PAGE_SPAN_BYTES;
824 let last_page_id = (replacement.end_offset - 1) / ursula_stream::COLD_INDEX_PAGE_SPAN_BYTES;
825 let old_identities = old_chunks
826 .iter()
827 .map(|chunk| (chunk.start_offset, chunk.end_offset, chunk.s3_path.as_str()))
828 .collect::<HashSet<_>>();
829 let mut pages = Vec::new();
830 let mut found = HashSet::new();
831 for page_id in first_page_id..=last_page_id {
832 let key = ColdIndexPageKey {
833 stream_id: stream_id.clone(),
834 generation: 0,
835 page_id,
836 };
837 let Some(mut page) = store.get_page(&key).await? else {
838 return Ok(None);
839 };
840 let previous = page.clone();
841 for chunk in &page.cold_chunks {
842 let identity = (chunk.start_offset, chunk.end_offset, chunk.s3_path.as_str());
843 if old_identities.contains(&identity) {
844 found.insert((chunk.start_offset, chunk.end_offset, chunk.s3_path.clone()));
845 }
846 }
847 page.cold_chunks.retain(|chunk| {
848 !old_identities.contains(&(
849 chunk.start_offset,
850 chunk.end_offset,
851 chunk.s3_path.as_str(),
852 ))
853 });
854 page.cold_chunks.retain(|chunk| {
855 chunk.start_offset != replacement.start_offset
856 || chunk.end_offset != replacement.end_offset
857 });
858 page.cold_chunks.push(replacement.clone());
859 page.cold_chunks.sort_by_key(|chunk| chunk.start_offset);
860 pages.push((key, previous, page));
861 }
862 if found.len() != old_identities.len() {
863 return Ok(None);
864 }
865 let mut rollback = Vec::with_capacity(pages.len());
866 for (key, previous, page) in pages {
867 if let Err(err) = store.put_page(&key, &page).await {
868 rollback_cold_index_pages(store, rollback).await?;
869 return Err(err);
870 }
871 rollback.push(ColdIndexPageRollback {
872 key,
873 previous: Some(previous),
874 written_chunk: replacement.clone(),
875 });
876 }
877 Ok(Some(rollback))
878}
879
880fn objects_for_read(page: &ColdIndexPage, read_start: u64, read_end: u64) -> Vec<ObjectPayloadRef> {
881 let mut objects = Vec::new();
882 for chunk in &page.cold_chunks {
883 if let Some(object) = intersect_object(&ObjectPayloadRef::from(chunk), read_start, read_end)
884 {
885 objects.push(object);
886 }
887 }
888 for object in &page.external_segments {
889 if let Some(object) = intersect_object(object, read_start, read_end) {
890 objects.push(object);
891 }
892 }
893 objects.sort_by_key(|object| object.start_offset);
894 objects
895}
896
897fn intersect_object(
898 object: &ObjectPayloadRef,
899 read_start: u64,
900 read_end: u64,
901) -> Option<ObjectPayloadRef> {
902 let start = object.start_offset.max(read_start);
903 let end = object.end_offset.min(read_end);
904 (start < end).then(|| object.clone())
905}
906
907fn objects_cover_range(objects: &[ObjectPayloadRef], start: u64, end: u64) -> bool {
908 let mut expected = start;
909 for object in objects {
910 if object.end_offset <= expected {
911 continue;
912 }
913 if object.start_offset > expected {
914 return false;
915 }
916 expected = object.end_offset;
917 if expected >= end {
918 return true;
919 }
920 }
921 expected == end
922}
923
924#[cfg(test)]
925mod tests {
926 use super::*;
927
928 fn key(page_id: u64) -> ColdIndexPageKey {
929 ColdIndexPageKey {
930 stream_id: BucketStreamId::new("benchcmp", "cold-index"),
931 generation: 7,
932 page_id,
933 }
934 }
935
936 fn page(start_offset: u64, end_offset: u64) -> ColdIndexPage {
937 ColdIndexPage {
938 start_offset,
939 end_offset,
940 cold_chunks: vec![ColdChunkRef {
941 start_offset,
942 end_offset,
943 s3_path: format!("benchcmp/cold-index/chunks/{start_offset:020}.bin"),
944 object_size: end_offset - start_offset,
945 ..Default::default()
946 }],
947 external_segments: Vec::new(),
948 }
949 }
950
951 #[tokio::test]
952 async fn memory_store_round_trips_pages() {
953 let store = InMemoryColdIndexPageStore::new();
954 let key = key(1);
955 let page = page(0, 128);
956
957 assert_eq!(
958 key.path(),
959 "benchcmp/cold-index/cold-index/00000000000000000007/00000000000000000001.idx"
960 );
961 assert_eq!(store.get_page(&key).await.expect("get missing"), None);
962 store.put_page(&key, &page).await.expect("put page");
963 assert_eq!(
964 store.get_page(&key).await.expect("get page"),
965 Some(page.clone())
966 );
967 assert!(page.covers(127));
968 assert!(!page.covers(128));
969 }
970
971 #[tokio::test]
972 async fn rollback_skips_page_updated_by_newer_writer() {
973 let store = InMemoryColdIndexPageStore::new();
974 let stream_id = BucketStreamId::new("benchcmp", "cold-index");
975 let first = ColdChunkRef {
976 start_offset: 0,
977 end_offset: 128,
978 s3_path: "benchcmp/cold-index/chunks/first.bin".to_owned(),
979 object_size: 128,
980 ..Default::default()
981 };
982 let stale = ColdChunkRef {
983 start_offset: 0,
984 end_offset: 128,
985 s3_path: "benchcmp/cold-index/chunks/stale.bin".to_owned(),
986 object_size: 128,
987 ..Default::default()
988 };
989 let newer = ColdChunkRef {
990 start_offset: 0,
991 end_offset: 128,
992 s3_path: "benchcmp/cold-index/chunks/newer.bin".to_owned(),
993 object_size: 128,
994 ..Default::default()
995 };
996 write_cold_chunk_index_pages(&store, &stream_id, &first)
997 .await
998 .expect("write first chunk");
999 let rollback = write_cold_chunk_index_pages_with_rollback(&store, &stream_id, &stale)
1000 .await
1001 .expect("write stale chunk");
1002 write_cold_chunk_index_pages(&store, &stream_id, &newer)
1003 .await
1004 .expect("write newer chunk");
1005
1006 rollback_cold_index_pages(&store, rollback)
1007 .await
1008 .expect("rollback stale chunk");
1009
1010 let page = store
1011 .get_page(&ColdIndexPageKey {
1012 stream_id,
1013 generation: 0,
1014 page_id: 0,
1015 })
1016 .await
1017 .expect("get page")
1018 .expect("page exists");
1019 assert_eq!(page.cold_chunks, vec![newer]);
1020 }
1021
1022 #[tokio::test]
1023 async fn compact_replacement_rollback_restores_input_chunks() {
1024 let store = InMemoryColdIndexPageStore::new();
1025 let stream_id = BucketStreamId::new("benchcmp", "cold-index");
1026 let first = ColdChunkRef {
1027 start_offset: 0,
1028 end_offset: 64,
1029 s3_path: "benchcmp/cold-index/chunks/first.bin".to_owned(),
1030 object_size: 64,
1031 ..Default::default()
1032 };
1033 let second = ColdChunkRef {
1034 start_offset: 64,
1035 end_offset: 128,
1036 s3_path: "benchcmp/cold-index/chunks/second.bin".to_owned(),
1037 object_size: 64,
1038 ..Default::default()
1039 };
1040 let replacement = ColdChunkRef {
1041 start_offset: 0,
1042 end_offset: 128,
1043 s3_path: "benchcmp/cold-index/chunks/compacted.bin".to_owned(),
1044 object_size: 128,
1045 ..Default::default()
1046 };
1047 for chunk in [&first, &second] {
1048 write_cold_chunk_index_pages(&store, &stream_id, chunk)
1049 .await
1050 .expect("write input chunk");
1051 }
1052
1053 let rollback = replace_cold_chunk_index_pages_with_rollback(
1054 &store,
1055 &stream_id,
1056 &[first.clone(), second.clone()],
1057 &replacement,
1058 )
1059 .await
1060 .expect("replace chunks")
1061 .expect("inputs still match");
1062 rollback_cold_index_pages(&store, rollback)
1063 .await
1064 .expect("rollback replacement");
1065
1066 let page = store
1067 .get_page(&ColdIndexPageKey {
1068 stream_id,
1069 generation: 0,
1070 page_id: 0,
1071 })
1072 .await
1073 .expect("get page")
1074 .expect("page exists");
1075 assert_eq!(page.cold_chunks, vec![first, second]);
1076 }
1077
1078 #[tokio::test]
1079 async fn read_reload_repairs_stale_cached_page() {
1080 let store = Arc::new(InMemoryColdIndexPageStore::new());
1081 let stream_id = BucketStreamId::new("benchcmp", "cold-index");
1082 let cache = ColdIndexPageCache::new(store.clone(), 8);
1083 let first = ColdChunkRef {
1084 start_offset: 0,
1085 end_offset: 128,
1086 s3_path: "benchcmp/cold-index/chunks/first.bin".to_owned(),
1087 object_size: 128,
1088 ..Default::default()
1089 };
1090 write_cold_chunk_index_pages(store.as_ref(), &stream_id, &first)
1091 .await
1092 .expect("write first chunk");
1093 assert_eq!(
1094 cache
1095 .object_segments_for_read(&stream_id, &StreamReadColdIndexSegment {
1096 generation: 0,
1097 page_id: 0,
1098 read_start_offset: 0,
1099 len: 1,
1100 },)
1101 .await
1102 .expect("read first byte")
1103 .len(),
1104 1
1105 );
1106
1107 let second = ColdChunkRef {
1108 start_offset: 128,
1109 end_offset: 256,
1110 s3_path: "benchcmp/cold-index/chunks/second.bin".to_owned(),
1111 object_size: 128,
1112 ..Default::default()
1113 };
1114 write_cold_chunk_index_pages(store.as_ref(), &stream_id, &second)
1115 .await
1116 .expect("write second chunk behind cache");
1117
1118 let objects = cache
1119 .object_segments_for_read(&stream_id, &StreamReadColdIndexSegment {
1120 generation: 0,
1121 page_id: 0,
1122 read_start_offset: 128,
1123 len: 1,
1124 })
1125 .await
1126 .expect("reload stale page");
1127 assert_eq!(objects[0].s3_path, second.s3_path);
1128 }
1129
1130 #[test]
1131 fn binary_page_format_round_trips_and_validates() {
1132 let key = key(42);
1133 let mut page = page(128, 256);
1134 page.external_segments.push(ObjectPayloadRef {
1135 start_offset: 256,
1136 end_offset: 300,
1137 s3_path: "benchcmp/cold-index/external/00000000000000000256.bin".to_owned(),
1138 object_size: 44,
1139 ..Default::default()
1140 });
1141 let bytes = encode_page(&key, &page);
1142 assert!(bytes.starts_with(COLD_INDEX_PAGE_MAGIC));
1143
1144 assert_eq!(decode_page(&key, &bytes).expect("decode page"), page);
1145
1146 let mut corrupted = bytes.clone();
1147 let last = corrupted.last_mut().expect("checksum byte");
1148 *last ^= 0xff;
1149 let err = decode_page(&key, &corrupted).expect_err("corrupt checksum");
1150 assert_eq!(err.kind(), io::ErrorKind::InvalidData);
1151
1152 let wrong_key = ColdIndexPageKey {
1153 stream_id: key.stream_id.clone(),
1154 generation: key.generation + 1,
1155 page_id: key.page_id,
1156 };
1157 let err = decode_page(&wrong_key, &bytes).expect_err("key mismatch");
1158 assert_eq!(err.kind(), io::ErrorKind::InvalidData);
1159 }
1160
1161 #[test]
1162 fn affinity_is_part_of_the_page_path_and_binary_identity() {
1163 let key = ColdIndexPageKey {
1164 stream_id: BucketStreamId::with_affinity("benchcmp", "run-42", "journal"),
1165 generation: 7,
1166 page_id: 42,
1167 };
1168 assert_eq!(
1169 key.path(),
1170 "benchcmp/run-42/journal/cold-index/00000000000000000007/00000000000000000042.idx"
1171 );
1172
1173 let page = page(0, 10);
1174 let bytes = encode_page(&key, &page);
1175 assert_eq!(decode_page(&key, &bytes).expect("decode page"), page);
1176
1177 let wrong_key = ColdIndexPageKey {
1178 stream_id: BucketStreamId::with_affinity("benchcmp", "run-43", "journal"),
1179 ..key
1180 };
1181 assert_eq!(
1182 decode_page(&wrong_key, &bytes)
1183 .expect_err("affinity mismatch")
1184 .kind(),
1185 io::ErrorKind::InvalidData
1186 );
1187 }
1188
1189 #[tokio::test]
1190 async fn page_cache_loads_on_miss_and_evicts_lru() {
1191 let store = Arc::new(InMemoryColdIndexPageStore::new());
1192 for page_id in 0..3 {
1193 store
1194 .put_page(&key(page_id), &page(page_id * 100, page_id * 100 + 100))
1195 .await
1196 .expect("put page");
1197 }
1198 let cache = ColdIndexPageCache::new(store, 2);
1199
1200 assert_eq!(
1201 cache
1202 .get_page(&key(0))
1203 .await
1204 .expect("load page")
1205 .expect("page")
1206 .start_offset,
1207 0
1208 );
1209 assert_eq!(
1210 cache
1211 .get_page(&key(1))
1212 .await
1213 .expect("load page")
1214 .expect("page")
1215 .start_offset,
1216 100
1217 );
1218 assert_eq!(cache.cached_page_count(), 2);
1219
1220 assert!(
1222 cache
1223 .get_page(&key(0))
1224 .await
1225 .expect("cached page")
1226 .is_some()
1227 );
1228 assert_eq!(
1229 cache
1230 .get_page(&key(2))
1231 .await
1232 .expect("load page")
1233 .expect("page")
1234 .start_offset,
1235 200
1236 );
1237 assert_eq!(cache.cached_page_count(), 2);
1238 }
1239
1240 #[tokio::test]
1241 async fn zero_capacity_cache_does_not_retain_pages() {
1242 let store = Arc::new(InMemoryColdIndexPageStore::new());
1243 store
1244 .put_page(&key(0), &page(0, 64))
1245 .await
1246 .expect("put page");
1247 let cache = ColdIndexPageCache::new(store, 0);
1248
1249 assert!(cache.get_page(&key(0)).await.expect("load page").is_some());
1250 assert_eq!(cache.cached_page_count(), 0);
1251 }
1252
1253 #[tokio::test]
1254 async fn selects_and_replaces_target_sized_contiguous_chunks() {
1255 let store = InMemoryColdIndexPageStore::new();
1256 let stream_id = BucketStreamId::new("benchcmp", "compact");
1257 let chunks = (0..4)
1258 .map(|index| ColdChunkRef {
1259 start_offset: index * 2,
1260 end_offset: index * 2 + 2,
1261 object_size: 2,
1262 s3_path: format!("old-{index}"),
1263 ..Default::default()
1264 })
1265 .collect::<Vec<_>>();
1266 for chunk in &chunks {
1267 write_cold_chunk_index_pages(&store, &stream_id, chunk)
1268 .await
1269 .expect("write chunk index");
1270 }
1271 let selected =
1272 select_cold_chunk_compaction(&chunks, 8, 16).expect("select compaction candidate");
1273 assert_eq!(selected, chunks);
1274 let replacement = ColdChunkRef {
1275 start_offset: 0,
1276 end_offset: 8,
1277 object_size: 8,
1278 s3_path: "replacement".to_owned(),
1279 ..Default::default()
1280 };
1281 assert!(
1282 replace_cold_chunk_index_pages(&store, &stream_id, &selected, &replacement)
1283 .await
1284 .expect("replace chunks")
1285 );
1286 let loaded = load_cold_chunks_from_pages(&store, &[ColdIndexPageKey {
1287 stream_id,
1288 generation: 0,
1289 page_id: 0,
1290 }])
1291 .await
1292 .expect("load replacement");
1293 assert_eq!(loaded, vec![replacement]);
1294 }
1295}