55// http://opensource.org/licenses/MIT>, at your option. You may not use this file except in
66// accordance with one or both of these licenses.
77
8- use std:: collections:: HashMap ;
8+ use std:: collections:: { HashMap , VecDeque } ;
99use std:: ops:: Deref ;
1010use std:: sync:: { Arc , Mutex } ;
1111
12- use lightning:: util:: persist:: KVStore ;
12+ use lightning:: util:: persist:: { KVStore , PageToken } ;
1313use lightning:: util:: ser:: { Readable , Writeable } ;
1414
1515use crate :: logger:: { log_error, LdkLogger } ;
@@ -25,7 +25,7 @@ pub(crate) trait StorableObject: Clone + Readable + Writeable {
2525 fn to_update ( & self ) -> Self :: Update ;
2626}
2727
28- pub ( crate ) trait StorableObjectId : std:: hash:: Hash + PartialEq + Eq {
28+ pub ( crate ) trait StorableObjectId : Clone + std:: hash:: Hash + Eq {
2929 fn encode_to_hex_str ( & self ) -> String ;
3030}
3131
@@ -40,11 +40,16 @@ pub(crate) enum DataStoreUpdateResult {
4040 NotFound ,
4141}
4242
43+ struct InMemoryObjects < SO : StorableObject > {
44+ objects : HashMap < SO :: Id , SO > ,
45+ creation_order : VecDeque < SO :: Id > ,
46+ }
47+
4348pub ( crate ) struct DataStore < SO : StorableObject , L : Deref >
4449where
4550 L :: Target : LdkLogger ,
4651{
47- objects : Mutex < HashMap < SO :: Id , SO > > ,
52+ objects : Mutex < InMemoryObjects < SO > > ,
4853 mutation_lock : tokio:: sync:: Mutex < ( ) > ,
4954 primary_namespace : String ,
5055 secondary_namespace : String ,
6065 objects : Vec < SO > , primary_namespace : String , secondary_namespace : String ,
6166 kv_store : Arc < DynStore > , logger : L ,
6267 ) -> Self {
63- let objects =
64- Mutex :: new ( HashMap :: from_iter ( objects. into_iter ( ) . map ( |obj| ( obj. id ( ) , obj) ) ) ) ;
68+ let mut creation_order = VecDeque :: with_capacity ( objects. len ( ) ) ;
69+ let mut objects_by_id = HashMap :: with_capacity ( objects. len ( ) ) ;
70+ for object in objects {
71+ let id = object. id ( ) ;
72+ creation_order. push_back ( id. clone ( ) ) ;
73+ objects_by_id. insert ( id, object) ;
74+ }
75+ let objects = Mutex :: new ( InMemoryObjects { objects : objects_by_id, creation_order } ) ;
6576 Self {
6677 objects,
6778 mutation_lock : tokio:: sync:: Mutex :: new ( ( ) ) ,
7788
7889 self . persist ( & object) . await ?;
7990 let mut locked_objects = self . objects . lock ( ) . expect ( "lock" ) ;
80- let updated = locked_objects. insert ( object. id ( ) , object) . is_some ( ) ;
91+ let id = object. id ( ) ;
92+ let updated = locked_objects. objects . insert ( id. clone ( ) , object) . is_some ( ) ;
93+ if !updated {
94+ locked_objects. creation_order . push_front ( id) ;
95+ }
8196 Ok ( updated)
8297 }
8398
87102 let id = object. id ( ) ;
88103 let data_to_persist = {
89104 let locked_objects = self . objects . lock ( ) . expect ( "lock" ) ;
90- if let Some ( existing_object) = locked_objects. get ( & id) {
105+ if let Some ( existing_object) = locked_objects. objects . get ( & id) {
91106 let mut updated_object = existing_object. clone ( ) ;
92107 let updated = updated_object. update ( object. to_update ( ) ) ;
93108 if updated {
@@ -104,7 +119,10 @@ where
104119 Some ( updated_object) => {
105120 self . persist ( & updated_object) . await ?;
106121 let mut locked_objects = self . objects . lock ( ) . expect ( "lock" ) ;
107- locked_objects. insert ( id, updated_object) ;
122+ let is_new = locked_objects. objects . insert ( id. clone ( ) , updated_object) . is_none ( ) ;
123+ if is_new {
124+ locked_objects. creation_order . push_front ( id) ;
125+ }
108126 Ok ( true )
109127 } ,
110128 None => Ok ( false ) ,
@@ -113,7 +131,7 @@ where
113131
114132 pub ( crate ) async fn remove ( & self , id : & SO :: Id ) -> Result < ( ) , Error > {
115133 let _guard = self . mutation_lock . lock ( ) . await ;
116- let should_remove = { self . objects . lock ( ) . expect ( "lock" ) . contains_key ( id) } ;
134+ let should_remove = { self . objects . lock ( ) . expect ( "lock" ) . objects . contains_key ( id) } ;
117135 if should_remove {
118136 let store_key = id. encode_to_hex_str ( ) ;
119137 KVStore :: remove (
@@ -135,7 +153,9 @@ where
135153 ) ;
136154 Error :: PersistenceFailed
137155 } ) ?;
138- self . objects . lock ( ) . expect ( "lock" ) . remove ( id) ;
156+ let mut locked_objects = self . objects . lock ( ) . expect ( "lock" ) ;
157+ locked_objects. objects . remove ( id) ;
158+ locked_objects. creation_order . retain ( |object_id| object_id != id) ;
139159 }
140160 Ok ( ( ) )
141161 }
@@ -146,15 +166,15 @@ where
146166 /// Until store reads are async, callers may temporarily see in-memory state that has not yet
147167 /// caught up to a write in progress.
148168 pub ( crate ) fn get ( & self , id : & SO :: Id ) -> Option < SO > {
149- self . objects . lock ( ) . expect ( "lock" ) . get ( id) . cloned ( )
169+ self . objects . lock ( ) . expect ( "lock" ) . objects . get ( id) . cloned ( )
150170 }
151171
152172 pub ( crate ) async fn update ( & self , update : SO :: Update ) -> Result < DataStoreUpdateResult , Error > {
153173 let _guard = self . mutation_lock . lock ( ) . await ;
154174 let id = update. id ( ) ;
155175 let updated_object = {
156176 let locked_objects = self . objects . lock ( ) . expect ( "lock" ) ;
157- let Some ( object) = locked_objects. get ( & id) else {
177+ let Some ( object) = locked_objects. objects . get ( & id) else {
158178 return Ok ( DataStoreUpdateResult :: NotFound ) ;
159179 } ;
160180 let mut updated_object = object. clone ( ) ;
@@ -166,7 +186,7 @@ where
166186
167187 self . persist ( & updated_object) . await ?;
168188 let mut locked_objects = self . objects . lock ( ) . expect ( "lock" ) ;
169- locked_objects. insert ( id, updated_object) ;
189+ locked_objects. objects . insert ( id, updated_object) ;
170190 Ok ( DataStoreUpdateResult :: Updated )
171191 }
172192
@@ -176,7 +196,48 @@ where
176196 /// Until store reads are async, callers may temporarily see in-memory state that has not yet
177197 /// caught up to a write in progress.
178198 pub ( crate ) fn list_filter < F : FnMut ( & & SO ) -> bool > ( & self , f : F ) -> Vec < SO > {
179- self . objects . lock ( ) . expect ( "lock" ) . values ( ) . filter ( f) . cloned ( ) . collect :: < Vec < SO > > ( )
199+ self . objects . lock ( ) . expect ( "lock" ) . objects . values ( ) . filter ( f) . cloned ( ) . collect ( )
200+ }
201+
202+ /// Returns a page of objects, ordered from most recently created to least recently created,
203+ /// together with a token that can be passed to a subsequent call to retrieve the next page.
204+ pub ( crate ) fn list_page (
205+ & self , page_token : Option < PageToken > ,
206+ ) -> Result < ( Vec < SO > , Option < PageToken > ) , Error > {
207+ const PAGE_SIZE : usize = 50 ;
208+
209+ let locked_objects = self . objects . lock ( ) . expect ( "lock" ) ;
210+ let start_index = if let Some ( token) = page_token {
211+ locked_objects
212+ . creation_order
213+ . iter ( )
214+ . position ( |id| id. encode_to_hex_str ( ) == token. as_str ( ) )
215+ . map ( |index| index + 1 )
216+ . ok_or_else ( || {
217+ log_error ! ( self . logger, "Object page token not found: {}" , token) ;
218+ Error :: InvalidPageToken
219+ } ) ?
220+ } else {
221+ 0
222+ } ;
223+
224+ let mut entries = locked_objects
225+ . creation_order
226+ . iter ( )
227+ . skip ( start_index)
228+ . filter_map ( |id| locked_objects. objects . get ( id) . cloned ( ) . map ( |object| ( id, object) ) )
229+ . take ( PAGE_SIZE + 1 )
230+ . collect :: < Vec < _ > > ( ) ;
231+ let has_more = entries. len ( ) > PAGE_SIZE ;
232+ entries. truncate ( PAGE_SIZE ) ;
233+
234+ let next_page_token = if has_more {
235+ entries. last ( ) . map ( |( id, _) | PageToken :: new ( id. encode_to_hex_str ( ) ) )
236+ } else {
237+ None
238+ } ;
239+ let objects = entries. into_iter ( ) . map ( |( _, object) | object) . collect ( ) ;
240+ Ok ( ( objects, next_page_token) )
180241 }
181242
182243 async fn persist ( & self , object : & SO ) -> Result < ( ) , Error > {
@@ -217,7 +278,7 @@ where
217278 /// Until store reads are async, callers may temporarily see in-memory state that has not yet
218279 /// caught up to a write in progress.
219280 pub ( crate ) fn contains_key ( & self , id : & SO :: Id ) -> bool {
220- self . objects . lock ( ) . expect ( "lock" ) . contains_key ( id)
281+ self . objects . lock ( ) . expect ( "lock" ) . objects . contains_key ( id)
221282 }
222283}
223284
@@ -230,6 +291,7 @@ mod tests {
230291 use super :: * ;
231292 use crate :: hex_utils;
232293 use crate :: io:: test_utils:: InMemoryStore ;
294+ use crate :: io:: utils:: read_all_objects;
233295 use crate :: types:: DynStoreWrapper ;
234296
235297 #[ derive( Clone , Copy , Debug , Eq , Hash , PartialEq ) ]
@@ -337,6 +399,116 @@ mod tests {
337399 )
338400 }
339401
402+ #[ tokio:: test]
403+ async fn list_page_paginates_in_reverse_creation_order ( ) {
404+ let store: Arc < DynStore > = Arc :: new ( DynStoreWrapper ( InMemoryStore :: new ( ) ) ) ;
405+ let logger = Arc :: new ( TestLogger :: new ( ) ) ;
406+ let data_store: DataStore < TestObject , Arc < TestLogger > > = DataStore :: new (
407+ Vec :: new ( ) ,
408+ "datastore_test_primary" . to_string ( ) ,
409+ "datastore_test_secondary" . to_string ( ) ,
410+ Arc :: clone ( & store) ,
411+ logger,
412+ ) ;
413+
414+ // Insert more objects than fit in a single page to exercise the pagination loop.
415+ let num_objects = 120u32 ;
416+ for i in 0 ..num_objects {
417+ let id = TestObjectId { id : i. to_be_bytes ( ) } ;
418+ data_store. insert ( TestObject { id, data : [ 7u8 ; 3 ] } ) . await . unwrap ( ) ;
419+ }
420+
421+ let mut listed = Vec :: with_capacity ( num_objects as usize ) ;
422+ let mut page_token = None ;
423+ loop {
424+ let ( page, next_page_token) = data_store. list_page ( page_token) . unwrap ( ) ;
425+ assert ! ( !page. is_empty( ) ) ;
426+ listed. extend ( page) ;
427+ page_token = next_page_token. map ( |token| PageToken :: new ( token. to_string ( ) ) ) ;
428+ if page_token. is_none ( ) {
429+ break ;
430+ }
431+ }
432+
433+ let expected: Vec < TestObject > = ( 0 ..num_objects)
434+ . rev ( )
435+ . map ( |i| TestObject { id : TestObjectId { id : i. to_be_bytes ( ) } , data : [ 7u8 ; 3 ] } )
436+ . collect ( ) ;
437+ assert_eq ! ( listed, expected) ;
438+ }
439+
440+ #[ tokio:: test]
441+ async fn list_page_token_survives_reload_after_unseen_object_is_removed ( ) {
442+ let store: Arc < DynStore > = Arc :: new ( DynStoreWrapper ( InMemoryStore :: new ( ) ) ) ;
443+ let logger = Arc :: new ( TestLogger :: new ( ) ) ;
444+ let primary_namespace = "datastore_test_primary" . to_string ( ) ;
445+ let secondary_namespace = "datastore_test_secondary" . to_string ( ) ;
446+ let data_store: DataStore < TestObject , Arc < TestLogger > > = DataStore :: new (
447+ Vec :: new ( ) ,
448+ primary_namespace. clone ( ) ,
449+ secondary_namespace. clone ( ) ,
450+ Arc :: clone ( & store) ,
451+ Arc :: clone ( & logger) ,
452+ ) ;
453+
454+ for i in 0 ..101u32 {
455+ let id = TestObjectId { id : i. to_be_bytes ( ) } ;
456+ data_store. insert ( TestObject { id, data : [ 7u8 ; 3 ] } ) . await . unwrap ( ) ;
457+ }
458+
459+ let ( first_page, page_token) = data_store. list_page ( None ) . unwrap ( ) ;
460+ assert_eq ! ( first_page. first( ) . unwrap( ) . id. id, 100u32 . to_be_bytes( ) ) ;
461+ assert_eq ! ( first_page. last( ) . unwrap( ) . id. id, 51u32 . to_be_bytes( ) ) ;
462+ let page_token = page_token. unwrap ( ) ;
463+ assert_eq ! ( page_token. as_str( ) , "00000033" ) ;
464+
465+ let oldest_id = TestObjectId { id : 0u32 . to_be_bytes ( ) } ;
466+ data_store. remove ( & oldest_id) . await . unwrap ( ) ;
467+ let reloaded_objects = read_all_objects (
468+ & * store,
469+ & primary_namespace,
470+ & secondary_namespace,
471+ Arc :: clone ( & logger) ,
472+ )
473+ . await
474+ . unwrap ( ) ;
475+ let reloaded_data_store: DataStore < TestObject , Arc < TestLogger > > =
476+ DataStore :: new ( reloaded_objects, primary_namespace, secondary_namespace, store, logger) ;
477+
478+ let ( second_page, next_page_token) =
479+ reloaded_data_store. list_page ( Some ( page_token) ) . unwrap ( ) ;
480+ assert_eq ! ( second_page. first( ) . unwrap( ) . id. id, 50u32 . to_be_bytes( ) ) ;
481+ assert_eq ! ( second_page. last( ) . unwrap( ) . id. id, 1u32 . to_be_bytes( ) ) ;
482+ assert_eq ! ( second_page. len( ) , 50 ) ;
483+ assert ! ( next_page_token. is_none( ) ) ;
484+ }
485+
486+ #[ test]
487+ fn list_page_rejects_invalid_tokens ( ) {
488+ let newest = TestObject { id : TestObjectId { id : 2u32 . to_be_bytes ( ) } , data : [ 2u8 ; 3 ] } ;
489+ let oldest = TestObject { id : TestObjectId { id : 1u32 . to_be_bytes ( ) } , data : [ 1u8 ; 3 ] } ;
490+ let data_store = new_failing_data_store ( vec ! [ newest, oldest] ) ;
491+
492+ let malformed_error =
493+ data_store. list_page ( Some ( PageToken :: new ( "3" . to_string ( ) ) ) ) . unwrap_err ( ) ;
494+ assert_eq ! ( malformed_error, Error :: InvalidPageToken ) ;
495+
496+ let unknown_error =
497+ data_store. list_page ( Some ( PageToken :: new ( "ffffffff" . to_string ( ) ) ) ) . unwrap_err ( ) ;
498+ assert_eq ! ( unknown_error, Error :: InvalidPageToken ) ;
499+ }
500+
501+ #[ test]
502+ fn list_page_only_reads_in_memory ( ) {
503+ let newest = TestObject { id : TestObjectId { id : 2u32 . to_be_bytes ( ) } , data : [ 2u8 ; 3 ] } ;
504+ let oldest = TestObject { id : TestObjectId { id : 1u32 . to_be_bytes ( ) } , data : [ 1u8 ; 3 ] } ;
505+ let data_store = new_failing_data_store ( vec ! [ newest, oldest] ) ;
506+
507+ let ( page, next_page_token) = data_store. list_page ( None ) . unwrap ( ) ;
508+ assert_eq ! ( page, vec![ newest, oldest] ) ;
509+ assert ! ( next_page_token. is_none( ) ) ;
510+ }
511+
340512 #[ tokio:: test]
341513 async fn data_is_persisted ( ) {
342514 let store: Arc < DynStore > = Arc :: new ( DynStoreWrapper ( InMemoryStore :: new ( ) ) ) ;
0 commit comments