Skip to main content

ursula_runtime/
cold_index.rs

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    /// Drops every cached generation/page for one stream. Compaction invokes
614    /// this on every replica when the replicated replacement command applies.
615    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
732/// Loads and de-duplicates the chunk references present in a stream's index
733/// pages. Chunks crossing a 64 MiB page boundary intentionally appear in more
734/// than one page.
735pub 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
761/// Selects the oldest contiguous run of undersized raw chunks whose combined
762/// payload reaches the byte target without exceeding the configured maximum.
763pub 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
799/// Atomically at the page-object level replaces a contiguous set of chunk
800/// references with one equivalent object. Every rewritten page always points
801/// at readable old or new bytes, so a retry after a partial S3 failure remains
802/// safe.
803pub 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        // Touch page 0 so page 1 becomes the eviction candidate.
1221        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}