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,17 @@ pub(crate) enum DataStoreUpdateResult {
4040 NotFound ,
4141}
4242
43+ struct InMemoryObjects < SO : StorableObject > {
44+ objects : HashMap < SO :: Id , SO > ,
45+ creation_order : VecDeque < ( u64 , SO :: Id ) > ,
46+ next_creation_order : u64 ,
47+ }
48+
4349pub ( crate ) struct DataStore < SO : StorableObject , L : Deref >
4450where
4551 L :: Target : LdkLogger ,
4652{
47- objects : Mutex < HashMap < SO :: Id , SO > > ,
53+ objects : Mutex < InMemoryObjects < SO > > ,
4854 mutation_lock : tokio:: sync:: Mutex < ( ) > ,
4955 primary_namespace : String ,
5056 secondary_namespace : String ,
6066 objects : Vec < SO > , primary_namespace : String , secondary_namespace : String ,
6167 kv_store : Arc < DynStore > , logger : L ,
6268 ) -> Self {
63- let objects =
64- Mutex :: new ( HashMap :: from_iter ( objects. into_iter ( ) . map ( |obj| ( obj. id ( ) , obj) ) ) ) ;
69+ let next_creation_order = objects. len ( ) as u64 ;
70+ let mut creation_order = VecDeque :: with_capacity ( objects. len ( ) ) ;
71+ let mut objects_by_id = HashMap :: with_capacity ( objects. len ( ) ) ;
72+ for ( index, object) in objects. into_iter ( ) . enumerate ( ) {
73+ let id = object. id ( ) ;
74+ creation_order. push_back ( ( next_creation_order - index as u64 , id. clone ( ) ) ) ;
75+ objects_by_id. insert ( id, object) ;
76+ }
77+ let objects = Mutex :: new ( InMemoryObjects {
78+ objects : objects_by_id,
79+ creation_order,
80+ next_creation_order,
81+ } ) ;
6582 Self {
6683 objects,
6784 mutation_lock : tokio:: sync:: Mutex :: new ( ( ) ) ,
7794
7895 self . persist ( & object) . await ?;
7996 let mut locked_objects = self . objects . lock ( ) . expect ( "lock" ) ;
80- let updated = locked_objects. insert ( object. id ( ) , object) . is_some ( ) ;
97+ let id = object. id ( ) ;
98+ let updated = locked_objects. objects . insert ( id. clone ( ) , object) . is_some ( ) ;
99+ if !updated {
100+ locked_objects. next_creation_order += 1 ;
101+ let creation_order = locked_objects. next_creation_order ;
102+ locked_objects. creation_order . push_front ( ( creation_order, id) ) ;
103+ }
81104 Ok ( updated)
82105 }
83106
87110 let id = object. id ( ) ;
88111 let data_to_persist = {
89112 let locked_objects = self . objects . lock ( ) . expect ( "lock" ) ;
90- if let Some ( existing_object) = locked_objects. get ( & id) {
113+ if let Some ( existing_object) = locked_objects. objects . get ( & id) {
91114 let mut updated_object = existing_object. clone ( ) ;
92115 let updated = updated_object. update ( object. to_update ( ) ) ;
93116 if updated {
@@ -104,7 +127,12 @@ where
104127 Some ( updated_object) => {
105128 self . persist ( & updated_object) . await ?;
106129 let mut locked_objects = self . objects . lock ( ) . expect ( "lock" ) ;
107- locked_objects. insert ( id, updated_object) ;
130+ let is_new = locked_objects. objects . insert ( id. clone ( ) , updated_object) . is_none ( ) ;
131+ if is_new {
132+ locked_objects. next_creation_order += 1 ;
133+ let creation_order = locked_objects. next_creation_order ;
134+ locked_objects. creation_order . push_front ( ( creation_order, id) ) ;
135+ }
108136 Ok ( true )
109137 } ,
110138 None => Ok ( false ) ,
@@ -113,7 +141,7 @@ where
113141
114142 pub ( crate ) async fn remove ( & self , id : & SO :: Id ) -> Result < ( ) , Error > {
115143 let _guard = self . mutation_lock . lock ( ) . await ;
116- let should_remove = { self . objects . lock ( ) . expect ( "lock" ) . contains_key ( id) } ;
144+ let should_remove = { self . objects . lock ( ) . expect ( "lock" ) . objects . contains_key ( id) } ;
117145 if should_remove {
118146 let store_key = id. encode_to_hex_str ( ) ;
119147 KVStore :: remove (
@@ -135,7 +163,9 @@ where
135163 ) ;
136164 Error :: PersistenceFailed
137165 } ) ?;
138- self . objects . lock ( ) . expect ( "lock" ) . remove ( id) ;
166+ let mut locked_objects = self . objects . lock ( ) . expect ( "lock" ) ;
167+ locked_objects. objects . remove ( id) ;
168+ locked_objects. creation_order . retain ( |( _, object_id) | object_id != id) ;
139169 }
140170 Ok ( ( ) )
141171 }
@@ -146,15 +176,15 @@ where
146176 /// Until store reads are async, callers may temporarily see in-memory state that has not yet
147177 /// caught up to a write in progress.
148178 pub ( crate ) fn get ( & self , id : & SO :: Id ) -> Option < SO > {
149- self . objects . lock ( ) . expect ( "lock" ) . get ( id) . cloned ( )
179+ self . objects . lock ( ) . expect ( "lock" ) . objects . get ( id) . cloned ( )
150180 }
151181
152182 pub ( crate ) async fn update ( & self , update : SO :: Update ) -> Result < DataStoreUpdateResult , Error > {
153183 let _guard = self . mutation_lock . lock ( ) . await ;
154184 let id = update. id ( ) ;
155185 let updated_object = {
156186 let locked_objects = self . objects . lock ( ) . expect ( "lock" ) ;
157- let Some ( object) = locked_objects. get ( & id) else {
187+ let Some ( object) = locked_objects. objects . get ( & id) else {
158188 return Ok ( DataStoreUpdateResult :: NotFound ) ;
159189 } ;
160190 let mut updated_object = object. clone ( ) ;
@@ -166,7 +196,7 @@ where
166196
167197 self . persist ( & updated_object) . await ?;
168198 let mut locked_objects = self . objects . lock ( ) . expect ( "lock" ) ;
169- locked_objects. insert ( id, updated_object) ;
199+ locked_objects. objects . insert ( id, updated_object) ;
170200 Ok ( DataStoreUpdateResult :: Updated )
171201 }
172202
@@ -176,7 +206,55 @@ where
176206 /// Until store reads are async, callers may temporarily see in-memory state that has not yet
177207 /// caught up to a write in progress.
178208 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 > > ( )
209+ self . objects . lock ( ) . expect ( "lock" ) . objects . values ( ) . filter ( f) . cloned ( ) . collect ( )
210+ }
211+
212+ /// Returns a page of objects, ordered from most recently created to least recently created,
213+ /// together with a token that can be passed to a subsequent call to retrieve the next page.
214+ pub ( crate ) fn list_page (
215+ & self , page_token : Option < PageToken > ,
216+ ) -> Result < ( Vec < SO > , Option < PageToken > ) , Error > {
217+ const PAGE_SIZE : usize = 50 ;
218+
219+ let locked_objects = self . objects . lock ( ) . expect ( "lock" ) ;
220+ let token_creation_order = page_token
221+ . map ( |token| {
222+ token. as_str ( ) . parse :: < u64 > ( ) . map_err ( |_| {
223+ log_error ! ( self . logger, "Invalid object page token: {}" , token) ;
224+ Error :: PersistenceFailed
225+ } )
226+ } )
227+ . transpose ( ) ?;
228+ if let Some ( order) = token_creation_order {
229+ if order > locked_objects. next_creation_order {
230+ log_error ! (
231+ self . logger,
232+ "Invalid object page token exceeding latest creation order: {}" ,
233+ order
234+ ) ;
235+ return Err ( Error :: PersistenceFailed ) ;
236+ }
237+ }
238+
239+ let mut entries = locked_objects
240+ . creation_order
241+ . iter ( )
242+ . filter ( |( order, _) | token_creation_order. is_none_or ( |token| * order < token) )
243+ . filter_map ( |( order, id) | {
244+ locked_objects. objects . get ( id) . cloned ( ) . map ( |object| ( * order, object) )
245+ } )
246+ . take ( PAGE_SIZE + 1 )
247+ . collect :: < Vec < _ > > ( ) ;
248+ let has_more = entries. len ( ) > PAGE_SIZE ;
249+ entries. truncate ( PAGE_SIZE ) ;
250+
251+ let next_page_token = if has_more {
252+ entries. last ( ) . map ( |( order, _) | PageToken :: new ( order. to_string ( ) ) )
253+ } else {
254+ None
255+ } ;
256+ let objects = entries. into_iter ( ) . map ( |( _, object) | object) . collect ( ) ;
257+ Ok ( ( objects, next_page_token) )
180258 }
181259
182260 async fn persist ( & self , object : & SO ) -> Result < ( ) , Error > {
@@ -217,7 +295,7 @@ where
217295 /// Until store reads are async, callers may temporarily see in-memory state that has not yet
218296 /// caught up to a write in progress.
219297 pub ( crate ) fn contains_key ( & self , id : & SO :: Id ) -> bool {
220- self . objects . lock ( ) . expect ( "lock" ) . contains_key ( id)
298+ self . objects . lock ( ) . expect ( "lock" ) . objects . contains_key ( id)
221299 }
222300}
223301
@@ -337,6 +415,55 @@ mod tests {
337415 )
338416 }
339417
418+ #[ tokio:: test]
419+ async fn list_page_paginates_in_reverse_creation_order ( ) {
420+ let store: Arc < DynStore > = Arc :: new ( DynStoreWrapper ( InMemoryStore :: new ( ) ) ) ;
421+ let logger = Arc :: new ( TestLogger :: new ( ) ) ;
422+ let data_store: DataStore < TestObject , Arc < TestLogger > > = DataStore :: new (
423+ Vec :: new ( ) ,
424+ "datastore_test_primary" . to_string ( ) ,
425+ "datastore_test_secondary" . to_string ( ) ,
426+ Arc :: clone ( & store) ,
427+ logger,
428+ ) ;
429+
430+ // Insert more objects than fit in a single page to exercise the pagination loop.
431+ let num_objects = 120u32 ;
432+ for i in 0 ..num_objects {
433+ let id = TestObjectId { id : i. to_be_bytes ( ) } ;
434+ data_store. insert ( TestObject { id, data : [ 7u8 ; 3 ] } ) . await . unwrap ( ) ;
435+ }
436+
437+ let mut listed = Vec :: with_capacity ( num_objects as usize ) ;
438+ let mut page_token = None ;
439+ loop {
440+ let ( page, next_page_token) = data_store. list_page ( page_token) . unwrap ( ) ;
441+ assert ! ( !page. is_empty( ) ) ;
442+ listed. extend ( page) ;
443+ page_token = next_page_token. map ( |token| PageToken :: new ( token. to_string ( ) ) ) ;
444+ if page_token. is_none ( ) {
445+ break ;
446+ }
447+ }
448+
449+ let expected: Vec < TestObject > = ( 0 ..num_objects)
450+ . rev ( )
451+ . map ( |i| TestObject { id : TestObjectId { id : i. to_be_bytes ( ) } , data : [ 7u8 ; 3 ] } )
452+ . collect ( ) ;
453+ assert_eq ! ( listed, expected) ;
454+ }
455+
456+ #[ test]
457+ fn list_page_only_reads_in_memory ( ) {
458+ let newest = TestObject { id : TestObjectId { id : 2u32 . to_be_bytes ( ) } , data : [ 2u8 ; 3 ] } ;
459+ let oldest = TestObject { id : TestObjectId { id : 1u32 . to_be_bytes ( ) } , data : [ 1u8 ; 3 ] } ;
460+ let data_store = new_failing_data_store ( vec ! [ newest, oldest] ) ;
461+
462+ let ( page, next_page_token) = data_store. list_page ( None ) . unwrap ( ) ;
463+ assert_eq ! ( page, vec![ newest, oldest] ) ;
464+ assert ! ( next_page_token. is_none( ) ) ;
465+ }
466+
340467 #[ tokio:: test]
341468 async fn data_is_persisted ( ) {
342469 let store: Arc < DynStore > = Arc :: new ( DynStoreWrapper ( InMemoryStore :: new ( ) ) ) ;
0 commit comments