2323//! `BLOCKS_IN_REL_FILE` + intra-segment offset)
2424
2525use std:: collections:: { BTreeMap , BTreeSet , HashSet } ;
26- use std:: io:: { self , Read , Write } ;
26+ use std:: io:: { self , Read } ;
2727use std:: sync:: Arc ;
2828
2929use anyhow:: { Context , Result , anyhow} ;
30+ use roaring:: RoaringBitmap ;
3031use thiserror:: Error ;
3132use tokio_util:: io:: SyncIoBridge ;
3233
@@ -39,7 +40,6 @@ use crate::pg::backup::{BackupSentinelDtoV2, increment, name_from_sentinel_key};
3940use crate :: pg:: wal:: segment:: { SegmentName , wal_segment_size} ;
4041use crate :: pg:: walparser:: {
4142 BlockLocation , ParsePageError , RelFileNode , WalParser , extract_locations_from_wal_file,
42- read_locations_from, write_locations_to,
4343} ;
4444use crate :: storage:: DynStorage ;
4545
@@ -69,11 +69,12 @@ pub enum DeltaError {
6969}
7070
7171/// In-memory delta map: which blocks of which relfiles changed?
72- /// `BTreeSet<u32>` per rel keeps memory bounded for sparse-delta workloads
73- /// (typical delta backup touches < 1% of pages)
72+ /// `RoaringBitmap` per rel run/bitmap-compresses dense rewrites (VACUUM FULL,
73+ /// CREATE INDEX, bulk load) that balloon a `BTreeSet<u32>` to ~13 B/block;
74+ /// sparse OLTP deltas stay comparable. Matches wal-g's `map[RelFileNode]*roaring.Bitmap`
7475#[ derive( Debug , Default , Clone ) ]
7576pub struct PagedFileDeltaMap {
76- by_rel : BTreeMap < RelFileNode , BTreeSet < u32 > > ,
77+ by_rel : BTreeMap < RelFileNode , RoaringBitmap > ,
7778}
7879
7980impl PagedFileDeltaMap {
@@ -86,7 +87,7 @@ impl PagedFileDeltaMap {
8687 }
8788
8889 pub fn len ( & self ) -> usize {
89- self . by_rel . values ( ) . map ( |s| s. len ( ) ) . sum ( )
90+ self . by_rel . values ( ) . map ( |s| s. len ( ) as usize ) . sum ( )
9091 }
9192
9293 pub fn add_location ( & mut self , loc : BlockLocation ) {
@@ -111,7 +112,7 @@ impl PagedFileDeltaMap {
111112 let seg_id = get_rel_file_id_from ( path) ?;
112113 let lo = seg_id as u32 * BLOCKS_IN_REL_FILE ;
113114 let hi = lo. saturating_add ( BLOCKS_IN_REL_FILE ) ;
114- let shifted: BTreeSet < u32 > = blocks. range ( lo..hi) . map ( |& b| b - lo) . collect ( ) ;
115+ let shifted: BTreeSet < u32 > = blocks. range ( lo..hi) . map ( |b| b - lo) . collect ( ) ;
115116 Ok ( Some ( shifted) )
116117 }
117118
@@ -120,7 +121,7 @@ impl PagedFileDeltaMap {
120121 pub fn locations ( & self ) -> Vec < BlockLocation > {
121122 let mut out = Vec :: with_capacity ( self . len ( ) ) ;
122123 for ( rel, blocks) in & self . by_rel {
123- for & b in blocks {
124+ for b in blocks {
124125 out. push ( BlockLocation {
125126 rel : * rel,
126127 block_no : b,
@@ -131,44 +132,11 @@ impl PagedFileDeltaMap {
131132 }
132133}
133134
134- /// On-disk delta file: aggregated block locations + parser state for the
135- /// cross-segment record stitching
136- pub struct DeltaFile {
137- pub locations : Vec < BlockLocation > ,
138- pub wal_parser : WalParser ,
139- }
140-
141- impl DeltaFile {
142- pub fn new ( wal_parser : WalParser ) -> Self {
143- Self {
144- locations : Vec :: new ( ) ,
145- wal_parser,
146- }
147- }
148-
149- /// wal-g `DeltaFile.Save`: write tuples (zero-terminated) then parser state
150- pub fn save < W : Write > ( & self , mut w : W ) -> Result < ( ) , DeltaError > {
151- write_locations_to ( & mut w, & self . locations ) ?;
152- self . wal_parser . save ( & mut w) ?;
153- Ok ( ( ) )
154- }
155-
156- /// Serialized byte length, for `Vec::with_capacity`. Must track `save`:
157- /// 16-byte tuples + 16-byte terminal sentinel + u32 parser len + parser data
158- pub fn serialized_len ( & self ) -> usize {
159- self . locations . len ( ) * 16 + 16 + 4 + self . wal_parser . current_record_data ( ) . len ( )
160- }
161-
162- pub fn load < R : Read > ( mut r : R ) -> Result < Self , DeltaError > {
163- let locations = read_locations_from ( & mut r)
164- . map_err ( |e| DeltaError :: Io ( io:: Error :: other ( e. to_string ( ) ) ) ) ?;
165- let wal_parser = WalParser :: load ( & mut r) ?;
166- Ok ( Self {
167- locations,
168- wal_parser,
169- } )
170- }
171- }
135+ // On-disk sidecar format (wal-g `DeltaFile`): location tuples, all-zero
136+ // terminator, then `WalParser` state (u32 len + bytes). Never materialized as
137+ // a struct — `wal_delta::record_segment` writes it append-only and
138+ // `fold_sidecar_into_map` streams it tuple-by-tuple, so neither side holds the
139+ // whole group's locations in memory
172140
173141// ─── path → RelFileNode parsing ─────────────────────────────────────────────
174142
@@ -566,19 +534,18 @@ async fn build_delta_map_from_sidecars(
566534 let mut g = first_used_delta;
567535 while g < last_complete_group {
568536 let name = delta_group_name ( timeline, g, seg_size) ;
569- let df = get_delta_file ( settings, storage, & name, compression)
537+ delta = fold_sidecar_into_map ( settings, storage, & name, compression, delta )
570538 . await
571- . with_context ( || format ! ( "delta sidecar {name}" ) ) ?;
572- delta . add_locations ( df . locations ) ;
539+ . with_context ( || format ! ( "delta sidecar {name}" ) ) ?
540+ . 0 ;
573541 g += n;
574542 }
575543 // Last complete group: its locations + parser seed for the tail walk
576544 let last_name = delta_group_name ( timeline, last_complete_group, seg_size) ;
577- let last_df = get_delta_file ( settings, storage, & last_name, compression)
545+ let ( d , mut parser ) = fold_sidecar_into_map ( settings, storage, & last_name, compression, delta )
578546 . await
579547 . with_context ( || format ! ( "delta sidecar {last_name}" ) ) ?;
580- delta. add_locations ( last_df. locations ) ;
581- let mut parser = last_df. wal_parser ;
548+ delta = d;
582549
583550 // Trailing partial group: raw WAL from the group start up to end_lsn
584551 let tail_first = first_not_used_delta;
@@ -593,27 +560,58 @@ async fn build_delta_map_from_sidecars(
593560 Ok ( delta)
594561}
595562
596- /// Fetch + decode a `<group>_delta` sidecar from `wal_005/` into a `DeltaFile`
597- async fn get_delta_file (
563+ /// Fetch a `<group>_delta` sidecar from `wal_005/` and fold its location tuples
564+ /// straight into `map`, returning the trailing `WalParser` state. Streams 16-byte
565+ /// tuples through a `SyncIoBridge` so the group's locations are never collected
566+ /// into a `Vec` — they land in the roaring map a tuple at a time
567+ async fn fold_sidecar_into_map (
598568 settings : & crate :: config:: Settings ,
599569 storage : & DynStorage ,
600570 group_name : & str ,
601571 compression : compression:: Method ,
602- ) -> Result < DeltaFile > {
603- use tokio :: io :: AsyncReadExt ;
572+ map : PagedFileDeltaMap ,
573+ ) -> Result < ( PagedFileDeltaMap , WalParser ) > {
604574 let key = delta_storage_key ( group_name, compression) ;
605575 let r = storage
606576 . get ( & key)
607577 . await
608578 . with_context ( || format ! ( "get {key}" ) ) ?;
609579 let decrypted = settings. decrypt ( r) ;
610- let mut decoded = compression:: decode ( compression, decrypted) ;
611- let mut buf = Vec :: new ( ) ;
612- decoded
613- . read_to_end ( & mut buf)
580+ let decoded = compression:: decode ( compression, decrypted) ;
581+ tokio:: task:: spawn_blocking ( move || fold_sidecar_stream ( decoded, map) )
614582 . await
615- . with_context ( || format ! ( "read {key}" ) ) ?;
616- DeltaFile :: load ( buf. as_slice ( ) ) . with_context ( || format ! ( "decode {key}" ) )
583+ . context ( "join sidecar fold" ) ?
584+ . with_context ( || format ! ( "decode {key}" ) )
585+ }
586+
587+ /// Sync side of [`fold_sidecar_into_map`]: read tuples until the all-zero
588+ /// terminator (or EOF), then the parser state
589+ fn fold_sidecar_stream (
590+ decoded : compression:: AsyncReader ,
591+ mut map : PagedFileDeltaMap ,
592+ ) -> Result < ( PagedFileDeltaMap , WalParser ) > {
593+ let mut r = SyncIoBridge :: new ( decoded) ;
594+ let mut buf = [ 0u8 ; 16 ] ;
595+ loop {
596+ match r. read_exact ( & mut buf) {
597+ Ok ( ( ) ) => { }
598+ Err ( e) if e. kind ( ) == io:: ErrorKind :: UnexpectedEof => {
599+ return Ok ( ( map, WalParser :: new ( ) ) ) ;
600+ }
601+ Err ( e) => return Err ( anyhow:: Error :: from ( e) . context ( "read sidecar tuple" ) ) ,
602+ }
603+ let loc = BlockLocation :: new (
604+ u32:: from_le_bytes ( buf[ 0 ..4 ] . try_into ( ) . unwrap ( ) ) ,
605+ u32:: from_le_bytes ( buf[ 4 ..8 ] . try_into ( ) . unwrap ( ) ) ,
606+ u32:: from_le_bytes ( buf[ 8 ..12 ] . try_into ( ) . unwrap ( ) ) ,
607+ u32:: from_le_bytes ( buf[ 12 ..16 ] . try_into ( ) . unwrap ( ) ) ,
608+ ) ;
609+ if loc. is_terminal ( ) {
610+ let parser = WalParser :: load ( & mut r) . context ( "load sidecar parser state" ) ?;
611+ return Ok ( ( map, parser) ) ;
612+ }
613+ map. add_location ( loc) ;
614+ }
617615}
618616
619617/// Full raw-WAL walk of `[start_lsn, end_lsn)`: parse every segment. Fallback
@@ -835,26 +833,27 @@ mod tests {
835833 }
836834
837835 #[ test]
838- fn delta_file_round_trip ( ) {
839- let wp = WalParser :: new ( ) ;
836+ fn sidecar_format_round_trip ( ) {
837+ // Bytes the streaming writer emits: tuples + terminator + parser state
838+ use crate :: pg:: walparser:: { read_locations_from, write_locations_to} ;
840839 let mut buf = Vec :: new ( ) ;
841- wp . save ( & mut buf ) . unwrap ( ) ;
842- let wp_reloaded = WalParser :: load ( buf. as_slice ( ) ) . unwrap ( ) ;
843-
844- let mut df = DeltaFile :: new ( wp_reloaded ) ;
845- df . locations
846- . push ( BlockLocation :: new ( DEFAULT_SPC_NODE , 16384 , 16385 , 7 ) ) ;
847- df . locations
848- . push ( BlockLocation :: new ( DEFAULT_SPC_NODE , 16384 , 16386 , 0 ) ) ;
849-
850- let mut out = Vec :: new ( ) ;
851- df . save ( & mut out ) . unwrap ( ) ;
852-
853- let df2 = DeltaFile :: load ( out . as_slice ( ) ) . unwrap ( ) ;
854- assert_eq ! ( df2 . locations . len ( ) , 2 ) ;
855- assert_eq ! ( df2 . locations [ 0 ] . block_no, 7 ) ;
856- assert_eq ! ( df2 . locations [ 1 ] . block_no , 0 ) ;
857- assert ! ( df2 . wal_parser . current_record_data( ) . is_empty( ) ) ;
840+ write_locations_to (
841+ & mut buf,
842+ & [
843+ BlockLocation :: new ( DEFAULT_SPC_NODE , 16384 , 16385 , 7 ) ,
844+ BlockLocation :: new ( DEFAULT_SPC_NODE , 16384 , 16386 , 0 ) ,
845+ ] ,
846+ )
847+ . unwrap ( ) ;
848+ WalParser :: new ( ) . save ( & mut buf ) . unwrap ( ) ;
849+
850+ let mut cur = buf . as_slice ( ) ;
851+ let locs = read_locations_from ( & mut cur ) . unwrap ( ) ;
852+ assert_eq ! ( locs . len ( ) , 2 ) ;
853+ assert_eq ! ( locs [ 0 ] . block_no , 7 ) ;
854+ assert_eq ! ( locs [ 1 ] . block_no, 0 ) ;
855+ let wp = WalParser :: load ( & mut cur ) . unwrap ( ) ;
856+ assert ! ( wp . current_record_data( ) . is_empty( ) ) ;
858857 }
859858
860859 #[ test]
@@ -922,29 +921,40 @@ mod tests {
922921 }
923922
924923 #[ tokio:: test]
925- async fn get_delta_file_reads_sidecar ( ) {
924+ async fn fold_sidecar_into_map_reads_blocks ( ) {
925+ use crate :: pg:: walparser:: write_locations_to;
926926 use crate :: storage:: fs:: FsStorage ;
927927 let dir = tempfile:: tempdir ( ) . unwrap ( ) ;
928928 let storage: DynStorage = Arc :: new ( FsStorage :: new ( dir. path ( ) ) . unwrap ( ) ) ;
929929 let settings = crate :: config:: Settings :: default ( ) ;
930930 let method = compression:: Method :: None ;
931931
932- let mut df = DeltaFile :: new ( WalParser :: new ( ) ) ;
933- df. locations
934- . push ( BlockLocation :: new ( DEFAULT_SPC_NODE , 16384 , 16385 , 7 ) ) ;
932+ // Sidecar bytes: one tuple + terminator + parser state
935933 let mut raw = Vec :: new ( ) ;
936- df. save ( & mut raw) . unwrap ( ) ;
934+ write_locations_to (
935+ & mut raw,
936+ & [ BlockLocation :: new ( DEFAULT_SPC_NODE , 16384 , 16385 , 7 ) ] ,
937+ )
938+ . unwrap ( ) ;
939+ WalParser :: new ( ) . save ( & mut raw) . unwrap ( ) ;
937940
938941 let group = delta_group_name ( 1 , 0 , DEFAULT_WAL_SEG_SIZE ) ;
939942 let key = delta_storage_key ( & group, method) ;
940943 let len = raw. len ( ) as u64 ;
941944 let r: crate :: compression:: AsyncReader = Box :: pin ( std:: io:: Cursor :: new ( raw) ) ;
942945 storage. put ( & key, r, Some ( len) ) . await . unwrap ( ) ;
943946
944- let got = get_delta_file ( & settings, & storage, & group, method)
945- . await
946- . unwrap ( ) ;
947- assert_eq ! ( got. locations. len( ) , 1 ) ;
948- assert_eq ! ( got. locations[ 0 ] . block_no, 7 ) ;
947+ let ( map, parser) = fold_sidecar_into_map (
948+ & settings,
949+ & storage,
950+ & group,
951+ method,
952+ PagedFileDeltaMap :: new ( ) ,
953+ )
954+ . await
955+ . unwrap ( ) ;
956+ let blocks = map. blocks_for ( "base/16384/16385" ) . unwrap ( ) . unwrap ( ) ;
957+ assert_eq ! ( blocks. into_iter( ) . collect:: <Vec <_>>( ) , vec![ 7u32 ] ) ;
958+ assert ! ( parser. current_record_data( ) . is_empty( ) ) ;
949959 }
950960}
0 commit comments