@@ -26,8 +26,8 @@ use lightning::routing::scoring::{
2626 ChannelLiquidities , ProbabilisticScorer , ProbabilisticScoringDecayParameters ,
2727} ;
2828use lightning:: util:: persist:: {
29- migrate_kv_store_data_async, KVStore , KVSTORE_NAMESPACE_KEY_ALPHABET ,
30- KVSTORE_NAMESPACE_KEY_MAX_LEN , NETWORK_GRAPH_PERSISTENCE_KEY ,
29+ migrate_kv_store_data_async, KVStore , PaginatedKVStore , PaginatedListResponse ,
30+ KVSTORE_NAMESPACE_KEY_ALPHABET , KVSTORE_NAMESPACE_KEY_MAX_LEN , NETWORK_GRAPH_PERSISTENCE_KEY ,
3131 NETWORK_GRAPH_PERSISTENCE_PRIMARY_NAMESPACE , NETWORK_GRAPH_PERSISTENCE_SECONDARY_NAMESPACE ,
3232 OUTPUT_SWEEPER_PERSISTENCE_KEY , OUTPUT_SWEEPER_PERSISTENCE_PRIMARY_NAMESPACE ,
3333 OUTPUT_SWEEPER_PERSISTENCE_SECONDARY_NAMESPACE , SCORER_PERSISTENCE_KEY ,
@@ -222,7 +222,8 @@ where
222222 } )
223223}
224224
225- /// Read all objects of type `T` from the given namespace, spawning reads in parallel.
225+ /// Read all objects of type `T` from the given namespace, listing keys page-by-page and spawning
226+ /// reads in parallel.
226227pub ( crate ) async fn read_all_objects < T , L > (
227228 kv_store : & DynStore , primary_namespace : & str , secondary_namespace : & str , logger : L ,
228229) -> Result < Vec < T > , std:: io:: Error >
@@ -234,55 +235,72 @@ where
234235 let type_name = std:: any:: type_name :: < T > ( ) ;
235236 let mut res = Vec :: new ( ) ;
236237
237- let mut stored_keys = KVStore :: list ( & * kv_store, primary_namespace, secondary_namespace) . await ?;
238-
239238 const BATCH_SIZE : usize = 50 ;
240239
241240 let mut set = tokio:: task:: JoinSet :: new ( ) ;
242-
243- // Fill JoinSet with tasks if possible
244- while set. len ( ) < BATCH_SIZE && !stored_keys. is_empty ( ) {
245- if let Some ( next_key) = stored_keys. pop ( ) {
246- let fut = KVStore :: read ( kv_store, primary_namespace, secondary_namespace, & next_key) ;
247- set. spawn ( fut) ;
248- debug_assert ! ( set. len( ) <= BATCH_SIZE ) ;
241+ let mut page_token = None ;
242+
243+ loop {
244+ let PaginatedListResponse { keys, next_page_token } = PaginatedKVStore :: list_paginated (
245+ & * kv_store,
246+ primary_namespace,
247+ secondary_namespace,
248+ page_token,
249+ )
250+ . await ?;
251+ let mut stored_keys = keys;
252+
253+ // Fill JoinSet with tasks if possible
254+ while set. len ( ) < BATCH_SIZE && !stored_keys. is_empty ( ) {
255+ if let Some ( next_key) = stored_keys. pop ( ) {
256+ let fut =
257+ KVStore :: read ( kv_store, primary_namespace, secondary_namespace, & next_key) ;
258+ set. spawn ( fut) ;
259+ debug_assert ! ( set. len( ) <= BATCH_SIZE ) ;
260+ }
249261 }
250- }
251262
252- while let Some ( read_res) = set. join_next ( ) . await {
253- // Exit early if we get an IO error.
254- let reader = read_res
255- . map_err ( |e| {
256- log_error ! ( logger, "Failed to read {}: {}" , type_name, e) ;
257- set. abort_all ( ) ;
258- e
259- } ) ?
260- . map_err ( |e| {
261- log_error ! ( logger, "Failed to read {}: {}" , type_name, e) ;
262- set. abort_all ( ) ;
263- e
264- } ) ?;
263+ while let Some ( read_res) = set. join_next ( ) . await {
264+ // Exit early if we get an IO error.
265+ let reader = read_res
266+ . map_err ( |e| {
267+ log_error ! ( logger, "Failed to read {}: {}" , type_name, e) ;
268+ set. abort_all ( ) ;
269+ e
270+ } ) ?
271+ . map_err ( |e| {
272+ log_error ! ( logger, "Failed to read {}: {}" , type_name, e) ;
273+ set. abort_all ( ) ;
274+ e
275+ } ) ?;
276+
277+ // Refill set for every finished future, if we still have something to do.
278+ if let Some ( next_key) = stored_keys. pop ( ) {
279+ let fut =
280+ KVStore :: read ( kv_store, primary_namespace, secondary_namespace, & next_key) ;
281+ set. spawn ( fut) ;
282+ debug_assert ! ( set. len( ) <= BATCH_SIZE ) ;
283+ }
265284
266- // Refill set for every finished future, if we still have something to do.
267- if let Some ( next_key) = stored_keys. pop ( ) {
268- let fut = KVStore :: read ( kv_store, primary_namespace, secondary_namespace, & next_key) ;
269- set. spawn ( fut) ;
270- debug_assert ! ( set. len( ) <= BATCH_SIZE ) ;
285+ // Handle result.
286+ let object = T :: read ( & mut & * reader) . map_err ( |e| {
287+ log_error ! ( logger, "Failed to deserialize {}: {}" , type_name, e) ;
288+ std:: io:: Error :: new (
289+ std:: io:: ErrorKind :: InvalidData ,
290+ format ! ( "Failed to deserialize {}" , type_name) ,
291+ )
292+ } ) ?;
293+ res. push ( object) ;
271294 }
272295
273- // Handle result.
274- let object = T :: read ( & mut & * reader) . map_err ( |e| {
275- log_error ! ( logger, "Failed to deserialize {}: {}" , type_name, e) ;
276- std:: io:: Error :: new (
277- std:: io:: ErrorKind :: InvalidData ,
278- format ! ( "Failed to deserialize {}" , type_name) ,
279- )
280- } ) ?;
281- res. push ( object) ;
282- }
296+ debug_assert ! ( set. is_empty( ) ) ;
297+ debug_assert ! ( stored_keys. is_empty( ) ) ;
283298
284- debug_assert ! ( set. is_empty( ) ) ;
285- debug_assert ! ( stored_keys. is_empty( ) ) ;
299+ page_token = next_page_token;
300+ if page_token. is_none ( ) {
301+ break ;
302+ }
303+ }
286304
287305 Ok ( res)
288306}
@@ -718,19 +736,56 @@ fn recover_incomplete_fs_store_migration(storage_dir_path: &Path) -> Result<(),
718736mod tests {
719737 use std:: fs;
720738 use std:: path:: { Path , PathBuf } ;
739+ use std:: sync:: Arc ;
721740
722741 use lightning:: util:: persist:: { migrate_kv_store_data_async, KVStore } ;
742+ use lightning:: util:: ser:: Writeable ;
743+ use lightning:: util:: test_utils:: TestLogger ;
723744 use lightning_persister:: fs_store:: v1:: FilesystemStore ;
724745 use lightning_persister:: fs_store:: v2:: FilesystemStoreV2 ;
725746
726747 use super :: test_utils:: random_storage_path;
727- use super :: { open_or_migrate_fs_store, read_or_generate_seed_file} ;
748+ use super :: { open_or_migrate_fs_store, read_all_objects, read_or_generate_seed_file} ;
749+ use crate :: io:: test_utils:: InMemoryStore ;
750+ use crate :: types:: { DynStore , DynStoreWrapper } ;
728751
729752 const TEST_PRIMARY_NAMESPACE : & str = "test_primary_namespace" ;
730753 const TEST_SECONDARY_NAMESPACE : & str = "test_secondary_namespace" ;
731754 const TEST_KEY : & str = "test_key" ;
732755 const TEST_VALUE : & [ u8 ] = b"test_value" ;
733756
757+ #[ tokio:: test]
758+ async fn read_all_objects_reads_across_pages ( ) {
759+ let store: Arc < DynStore > = Arc :: new ( DynStoreWrapper ( InMemoryStore :: new ( ) ) ) ;
760+ let logger = Arc :: new ( TestLogger :: new ( ) ) ;
761+
762+ // Write more objects than fit in a single page to exercise the pagination loop.
763+ let num_objects = 120u64 ;
764+ for i in 0 ..num_objects {
765+ let key = format ! ( "key_{:03}" , i) ;
766+ KVStore :: write (
767+ & * store,
768+ TEST_PRIMARY_NAMESPACE ,
769+ TEST_SECONDARY_NAMESPACE ,
770+ & key,
771+ i. encode ( ) ,
772+ )
773+ . await
774+ . unwrap ( ) ;
775+ }
776+
777+ let mut read: Vec < u64 > = read_all_objects (
778+ & * store,
779+ TEST_PRIMARY_NAMESPACE ,
780+ TEST_SECONDARY_NAMESPACE ,
781+ Arc :: clone ( & logger) ,
782+ )
783+ . await
784+ . unwrap ( ) ;
785+ read. sort_unstable ( ) ;
786+ assert_eq ! ( read, ( 0 ..num_objects) . collect:: <Vec <u64 >>( ) ) ;
787+ }
788+
734789 #[ test]
735790 fn generated_seed_is_readable ( ) {
736791 let mut rand_path = random_storage_path ( ) ;
0 commit comments