@@ -9,7 +9,7 @@ use std::collections::HashMap;
99use std:: ops:: Deref ;
1010use std:: sync:: { Arc , Mutex } ;
1111
12- use lightning:: util:: persist:: KVStore ;
12+ use lightning:: util:: persist:: { KVStore , PageToken , PaginatedKVStore } ;
1313use lightning:: util:: ser:: { Readable , Writeable } ;
1414
1515use crate :: logger:: { log_error, LdkLogger } ;
@@ -44,7 +44,7 @@ pub(crate) struct DataStore<SO: StorableObject, L: Deref>
4444where
4545 L :: Target : LdkLogger ,
4646{
47- objects : Mutex < HashMap < SO :: Id , SO > > ,
47+ objects : Mutex < HashMap < String , SO > > ,
4848 mutation_lock : tokio:: sync:: Mutex < ( ) > ,
4949 primary_namespace : String ,
5050 secondary_namespace : String ,
6060 objects : Vec < SO > , primary_namespace : String , secondary_namespace : String ,
6161 kv_store : Arc < DynStore > , logger : L ,
6262 ) -> Self {
63- let objects =
64- Mutex :: new ( HashMap :: from_iter ( objects. into_iter ( ) . map ( |obj| ( obj. id ( ) , obj) ) ) ) ;
63+ let objects = Mutex :: new ( HashMap :: from_iter (
64+ objects. into_iter ( ) . map ( |obj| ( obj. id ( ) . encode_to_hex_str ( ) , obj) ) ,
65+ ) ) ;
6566 Self {
6667 objects,
6768 mutation_lock : tokio:: sync:: Mutex :: new ( ( ) ) ,
@@ -74,20 +75,22 @@ where
7475
7576 pub ( crate ) async fn insert ( & self , object : SO ) -> Result < bool , Error > {
7677 let _guard = self . mutation_lock . lock ( ) . await ;
78+ let store_key = object. id ( ) . encode_to_hex_str ( ) ;
7779
7880 self . persist ( & object) . await ?;
7981 let mut locked_objects = self . objects . lock ( ) . expect ( "lock" ) ;
80- let updated = locked_objects. insert ( object . id ( ) , object) . is_some ( ) ;
82+ let updated = locked_objects. insert ( store_key , object) . is_some ( ) ;
8183 Ok ( updated)
8284 }
8385
8486 pub ( crate ) async fn insert_or_update ( & self , object : SO ) -> Result < bool , Error > {
8587 let _guard = self . mutation_lock . lock ( ) . await ;
8688
8789 let id = object. id ( ) ;
90+ let store_key = id. encode_to_hex_str ( ) ;
8891 let data_to_persist = {
8992 let locked_objects = self . objects . lock ( ) . expect ( "lock" ) ;
90- if let Some ( existing_object) = locked_objects. get ( & id ) {
93+ if let Some ( existing_object) = locked_objects. get ( & store_key ) {
9194 let mut updated_object = existing_object. clone ( ) ;
9295 let updated = updated_object. update ( object. to_update ( ) ) ;
9396 if updated {
@@ -104,7 +107,7 @@ where
104107 Some ( updated_object) => {
105108 self . persist ( & updated_object) . await ?;
106109 let mut locked_objects = self . objects . lock ( ) . expect ( "lock" ) ;
107- locked_objects. insert ( id , updated_object) ;
110+ locked_objects. insert ( store_key , updated_object) ;
108111 Ok ( true )
109112 } ,
110113 None => Ok ( false ) ,
@@ -113,9 +116,9 @@ where
113116
114117 pub ( crate ) async fn remove ( & self , id : & SO :: Id ) -> Result < ( ) , Error > {
115118 let _guard = self . mutation_lock . lock ( ) . await ;
116- let should_remove = { self . objects . lock ( ) . expect ( "lock" ) . contains_key ( id) } ;
119+ let store_key = id. encode_to_hex_str ( ) ;
120+ let should_remove = { self . objects . lock ( ) . expect ( "lock" ) . contains_key ( & store_key) } ;
117121 if should_remove {
118- let store_key = id. encode_to_hex_str ( ) ;
119122 KVStore :: remove (
120123 & * self . kv_store ,
121124 & self . primary_namespace ,
@@ -135,7 +138,7 @@ where
135138 ) ;
136139 Error :: PersistenceFailed
137140 } ) ?;
138- self . objects . lock ( ) . expect ( "lock" ) . remove ( id ) ;
141+ self . objects . lock ( ) . expect ( "lock" ) . remove ( & store_key ) ;
139142 }
140143 Ok ( ( ) )
141144 }
@@ -146,15 +149,16 @@ where
146149 /// Until store reads are async, callers may temporarily see in-memory state that has not yet
147150 /// caught up to a write in progress.
148151 pub ( crate ) fn get ( & self , id : & SO :: Id ) -> Option < SO > {
149- self . objects . lock ( ) . expect ( "lock" ) . get ( id ) . cloned ( )
152+ self . objects . lock ( ) . expect ( "lock" ) . get ( & id . encode_to_hex_str ( ) ) . cloned ( )
150153 }
151154
152155 pub ( crate ) async fn update ( & self , update : SO :: Update ) -> Result < DataStoreUpdateResult , Error > {
153156 let _guard = self . mutation_lock . lock ( ) . await ;
154157 let id = update. id ( ) ;
158+ let store_key = id. encode_to_hex_str ( ) ;
155159 let updated_object = {
156160 let locked_objects = self . objects . lock ( ) . expect ( "lock" ) ;
157- let Some ( object) = locked_objects. get ( & id ) else {
161+ let Some ( object) = locked_objects. get ( & store_key ) else {
158162 return Ok ( DataStoreUpdateResult :: NotFound ) ;
159163 } ;
160164 let mut updated_object = object. clone ( ) ;
@@ -166,7 +170,7 @@ where
166170
167171 self . persist ( & updated_object) . await ?;
168172 let mut locked_objects = self . objects . lock ( ) . expect ( "lock" ) ;
169- locked_objects. insert ( id , updated_object) ;
173+ locked_objects. insert ( store_key , updated_object) ;
170174 Ok ( DataStoreUpdateResult :: Updated )
171175 }
172176
@@ -179,6 +183,40 @@ where
179183 self . objects . lock ( ) . expect ( "lock" ) . values ( ) . filter ( f) . cloned ( ) . collect :: < Vec < SO > > ( )
180184 }
181185
186+ /// Returns a page of objects, ordered from most recently created to least recently created,
187+ /// together with a token that can be passed to a subsequent call to retrieve the next page.
188+ ///
189+ /// The underlying store is only queried for the page's key order, which the in-memory map
190+ /// doesn't track; the objects themselves are served from memory.
191+ pub ( crate ) async fn list_page (
192+ & self , page_token : Option < PageToken > ,
193+ ) -> Result < ( Vec < SO > , Option < PageToken > ) , Error > {
194+ let _guard = self . mutation_lock . lock ( ) . await ;
195+ let response = PaginatedKVStore :: list_paginated (
196+ & * self . kv_store ,
197+ & self . primary_namespace ,
198+ & self . secondary_namespace ,
199+ page_token,
200+ )
201+ . await
202+ . map_err ( |e| {
203+ log_error ! (
204+ self . logger,
205+ "Listing object data under {}/{} failed due to: {}" ,
206+ & self . primary_namespace,
207+ & self . secondary_namespace,
208+ e
209+ ) ;
210+ Error :: PersistenceFailed
211+ } ) ?;
212+
213+ let locked_objects = self . objects . lock ( ) . expect ( "lock" ) ;
214+ let objects =
215+ response. keys . iter ( ) . filter_map ( |key| locked_objects. get ( key) . cloned ( ) ) . collect ( ) ;
216+
217+ Ok ( ( objects, response. next_page_token ) )
218+ }
219+
182220 async fn persist ( & self , object : & SO ) -> Result < ( ) , Error > {
183221 let ( store_key, data) = Self :: encode_object ( object) ;
184222 self . persist_encoded ( store_key, data) . await
@@ -217,7 +255,7 @@ where
217255 /// Until store reads are async, callers may temporarily see in-memory state that has not yet
218256 /// caught up to a write in progress.
219257 pub ( crate ) fn contains_key ( & self , id : & SO :: Id ) -> bool {
220- self . objects . lock ( ) . expect ( "lock" ) . contains_key ( id )
258+ self . objects . lock ( ) . expect ( "lock" ) . contains_key ( & id . encode_to_hex_str ( ) )
221259 }
222260}
223261
@@ -337,6 +375,74 @@ mod tests {
337375 )
338376 }
339377
378+ #[ tokio:: test]
379+ async fn list_page_paginates_in_reverse_creation_order ( ) {
380+ let store: Arc < DynStore > = Arc :: new ( DynStoreWrapper ( InMemoryStore :: new ( ) ) ) ;
381+ let logger = Arc :: new ( TestLogger :: new ( ) ) ;
382+ let data_store: DataStore < TestObject , Arc < TestLogger > > = DataStore :: new (
383+ Vec :: new ( ) ,
384+ "datastore_test_primary" . to_string ( ) ,
385+ "datastore_test_secondary" . to_string ( ) ,
386+ Arc :: clone ( & store) ,
387+ logger,
388+ ) ;
389+
390+ // Insert more objects than fit in a single page to exercise the pagination loop.
391+ let num_objects = 120u32 ;
392+ for i in 0 ..num_objects {
393+ let id = TestObjectId { id : i. to_be_bytes ( ) } ;
394+ data_store. insert ( TestObject { id, data : [ 7u8 ; 3 ] } ) . await . unwrap ( ) ;
395+ }
396+
397+ let mut listed = Vec :: with_capacity ( num_objects as usize ) ;
398+ let mut page_token = None ;
399+ loop {
400+ let ( page, next_page_token) = data_store. list_page ( page_token) . await . unwrap ( ) ;
401+ assert ! ( !page. is_empty( ) ) ;
402+ listed. extend ( page) ;
403+ page_token = next_page_token;
404+ if page_token. is_none ( ) {
405+ break ;
406+ }
407+ }
408+
409+ let expected: Vec < TestObject > = ( 0 ..num_objects)
410+ . rev ( )
411+ . map ( |i| TestObject { id : TestObjectId { id : i. to_be_bytes ( ) } , data : [ 7u8 ; 3 ] } )
412+ . collect ( ) ;
413+ assert_eq ! ( listed, expected) ;
414+ }
415+
416+ #[ tokio:: test]
417+ async fn list_page_waits_for_mutations ( ) {
418+ let store: Arc < DynStore > = Arc :: new ( DynStoreWrapper ( InMemoryStore :: new ( ) ) ) ;
419+ let logger = Arc :: new ( TestLogger :: new ( ) ) ;
420+ let data_store: Arc < DataStore < TestObject , Arc < TestLogger > > > = Arc :: new ( DataStore :: new (
421+ Vec :: new ( ) ,
422+ "datastore_test_primary" . to_string ( ) ,
423+ "datastore_test_secondary" . to_string ( ) ,
424+ store,
425+ logger,
426+ ) ) ;
427+
428+ let mutation_guard = data_store. mutation_lock . lock ( ) . await ;
429+ let ( started_tx, started_rx) = tokio:: sync:: oneshot:: channel ( ) ;
430+ let task_store = Arc :: clone ( & data_store) ;
431+ let list_task = tokio:: spawn ( async move {
432+ started_tx. send ( ( ) ) . unwrap ( ) ;
433+ task_store. list_page ( None ) . await
434+ } ) ;
435+
436+ started_rx. await . unwrap ( ) ;
437+ tokio:: task:: yield_now ( ) . await ;
438+ assert ! ( !list_task. is_finished( ) ) ;
439+
440+ drop ( mutation_guard) ;
441+ let ( page, next_page_token) = list_task. await . unwrap ( ) . unwrap ( ) ;
442+ assert ! ( page. is_empty( ) ) ;
443+ assert ! ( next_page_token. is_none( ) ) ;
444+ }
445+
340446 #[ tokio:: test]
341447 async fn data_is_persisted ( ) {
342448 let store: Arc < DynStore > = Arc :: new ( DynStoreWrapper ( InMemoryStore :: new ( ) ) ) ;
0 commit comments