@@ -40,6 +40,13 @@ pub(crate) enum DataStoreUpdateResult {
4040 NotFound ,
4141}
4242
43+ #[ derive( PartialEq , Eq , Debug , Clone , Copy ) ]
44+ pub ( crate ) enum DataStoreUpdateOrInsertResult {
45+ Inserted ,
46+ Updated ,
47+ Unchanged ,
48+ }
49+
4350pub ( crate ) struct DataStore < SO : StorableObject , L : Deref >
4451where
4552 L :: Target : LdkLogger ,
8188 Ok ( updated)
8289 }
8390
91+ /// Like [`Self::insert`], but when an entry with the object's id already exists, merges the
92+ /// object's full update ([`StorableObject::to_update`]) into it instead of replacing it.
93+ ///
94+ /// Unlike [`Self::update_or_insert`], the caller does not choose what is merged into an
95+ /// existing entry: the full update is always applied.
8496 pub ( crate ) async fn insert_or_update ( & self , object : SO ) -> Result < bool , Error > {
8597 let _guard = self . mutation_lock . lock ( ) . await ;
8698
@@ -170,6 +182,47 @@ where
170182 Ok ( DataStoreUpdateResult :: Updated )
171183 }
172184
185+ /// Applies `update` when an object with its id already exists, or inserts `object` when none
186+ /// does.
187+ ///
188+ /// Like [`Self::update`], but falls back to inserting `object` instead of returning
189+ /// [`DataStoreUpdateResult::NotFound`]. Unlike [`Self::insert_or_update`], the caller chooses
190+ /// exactly what is merged into an existing entry: `update` may carry less than the full
191+ /// object.
192+ ///
193+ /// The existence check and the write share one critical section of the mutation lock, so a
194+ /// concurrent writer cannot land in between and later have its state clobbered by the insert
195+ /// fallback — the check-then-act race that separate [`Self::contains_key`] +
196+ /// [`Self::insert_or_update`] calls reintroduce.
197+ pub ( crate ) async fn update_or_insert (
198+ & self , update : SO :: Update , object : SO ,
199+ ) -> Result < DataStoreUpdateOrInsertResult , Error > {
200+ debug_assert ! ( update. id( ) == object. id( ) , "update and object must share an id" ) ;
201+ let _guard = self . mutation_lock . lock ( ) . await ;
202+
203+ let id = update. id ( ) ;
204+ let ( data_to_persist, result) = {
205+ let locked_objects = self . objects . lock ( ) . expect ( "lock" ) ;
206+ match locked_objects. get ( & id) {
207+ Some ( existing_object) => {
208+ let mut updated_object = existing_object. clone ( ) ;
209+ if updated_object. update ( update) {
210+ ( Some ( updated_object) , DataStoreUpdateOrInsertResult :: Updated )
211+ } else {
212+ ( None , DataStoreUpdateOrInsertResult :: Unchanged )
213+ }
214+ } ,
215+ None => ( Some ( object) , DataStoreUpdateOrInsertResult :: Inserted ) ,
216+ }
217+ } ;
218+
219+ if let Some ( object) = data_to_persist {
220+ self . persist ( & object) . await ?;
221+ self . objects . lock ( ) . expect ( "lock" ) . insert ( id, object) ;
222+ }
223+ Ok ( result)
224+ }
225+
173226 /// Returns in-memory objects matching `f`.
174227 ///
175228 /// The async mutation lock serializes writers, but this synchronous reader cannot wait on it.
@@ -403,6 +456,102 @@ mod tests {
403456 assert_eq ! ( Ok ( true ) , data_store. insert_or_update( new_iou_object) . await ) ;
404457 }
405458
459+ #[ tokio:: test]
460+ async fn update_or_insert_inserts_when_absent ( ) {
461+ let store: Arc < DynStore > = Arc :: new ( DynStoreWrapper ( InMemoryStore :: new ( ) ) ) ;
462+ let logger = Arc :: new ( TestLogger :: new ( ) ) ;
463+ let primary_namespace = "datastore_test_primary" . to_string ( ) ;
464+ let secondary_namespace = "datastore_test_secondary" . to_string ( ) ;
465+ let data_store: DataStore < TestObject , Arc < TestLogger > > = DataStore :: new (
466+ Vec :: new ( ) ,
467+ primary_namespace. clone ( ) ,
468+ secondary_namespace. clone ( ) ,
469+ Arc :: clone ( & store) ,
470+ logger,
471+ ) ;
472+
473+ let id = TestObjectId { id : [ 42u8 ; 4 ] } ;
474+ let object = TestObject { id, data : [ 23u8 ; 3 ] } ;
475+ let update = TestObjectUpdate { id, data : [ 25u8 ; 3 ] } ;
476+ assert_eq ! (
477+ Ok ( DataStoreUpdateOrInsertResult :: Inserted ) ,
478+ data_store. update_or_insert( update, object) . await
479+ ) ;
480+
481+ // The insert path stores the fallback object as-is; the update is not applied to it.
482+ assert_eq ! ( Some ( object) , data_store. get( & id) ) ;
483+ let store_key = id. encode_to_hex_str ( ) ;
484+ assert ! ( KVStore :: read( & * store, & primary_namespace, & secondary_namespace, & store_key)
485+ . await
486+ . is_ok( ) ) ;
487+ }
488+
489+ #[ tokio:: test]
490+ async fn update_or_insert_applies_update_when_present ( ) {
491+ let store: Arc < DynStore > = Arc :: new ( DynStoreWrapper ( InMemoryStore :: new ( ) ) ) ;
492+ let logger = Arc :: new ( TestLogger :: new ( ) ) ;
493+ let id = TestObjectId { id : [ 42u8 ; 4 ] } ;
494+ let existing_object = TestObject { id, data : [ 23u8 ; 3 ] } ;
495+ let data_store: DataStore < TestObject , Arc < TestLogger > > = DataStore :: new (
496+ vec ! [ existing_object] ,
497+ "datastore_test_primary" . to_string ( ) ,
498+ "datastore_test_secondary" . to_string ( ) ,
499+ store,
500+ logger,
501+ ) ;
502+
503+ // When an entry exists, only the update is applied; the fallback object must not replace
504+ // it.
505+ let update = TestObjectUpdate { id, data : [ 24u8 ; 3 ] } ;
506+ let object = TestObject { id, data : [ 99u8 ; 3 ] } ;
507+ assert_eq ! (
508+ Ok ( DataStoreUpdateOrInsertResult :: Updated ) ,
509+ data_store. update_or_insert( update, object) . await
510+ ) ;
511+ assert_eq ! ( data_store. get( & id) . unwrap( ) . data, [ 24u8 ; 3 ] ) ;
512+ }
513+
514+ #[ tokio:: test]
515+ async fn update_or_insert_returns_unchanged_without_persisting ( ) {
516+ let id = TestObjectId { id : [ 42u8 ; 4 ] } ;
517+ let existing_object = TestObject { id, data : [ 23u8 ; 3 ] } ;
518+ let data_store = new_failing_data_store ( vec ! [ existing_object] ) ;
519+
520+ // A no-op update returns `Unchanged` without attempting to persist (the store fails all
521+ // writes) and without falling back to the object.
522+ let update = TestObjectUpdate { id, data : [ 23u8 ; 3 ] } ;
523+ let object = TestObject { id, data : [ 99u8 ; 3 ] } ;
524+ assert_eq ! (
525+ Ok ( DataStoreUpdateOrInsertResult :: Unchanged ) ,
526+ data_store. update_or_insert( update, object) . await
527+ ) ;
528+ assert_eq ! ( Some ( existing_object) , data_store. get( & id) ) ;
529+ }
530+
531+ #[ tokio:: test]
532+ async fn update_or_insert_does_not_mutate_memory_if_persist_fails ( ) {
533+ let existing_id = TestObjectId { id : [ 42u8 ; 4 ] } ;
534+ let existing_object = TestObject { id : existing_id, data : [ 23u8 ; 3 ] } ;
535+ let data_store = new_failing_data_store ( vec ! [ existing_object] ) ;
536+
537+ let update = TestObjectUpdate { id : existing_id, data : [ 24u8 ; 3 ] } ;
538+ let object = TestObject { id : existing_id, data : [ 24u8 ; 3 ] } ;
539+ assert_eq ! (
540+ Err ( Error :: PersistenceFailed ) ,
541+ data_store. update_or_insert( update, object) . await
542+ ) ;
543+ assert_eq ! ( Some ( existing_object) , data_store. get( & existing_id) ) ;
544+
545+ let new_id = TestObjectId { id : [ 55u8 ; 4 ] } ;
546+ let new_object = TestObject { id : new_id, data : [ 34u8 ; 3 ] } ;
547+ let new_update = TestObjectUpdate { id : new_id, data : [ 34u8 ; 3 ] } ;
548+ assert_eq ! (
549+ Err ( Error :: PersistenceFailed ) ,
550+ data_store. update_or_insert( new_update, new_object) . await
551+ ) ;
552+ assert ! ( data_store. get( & new_id) . is_none( ) ) ;
553+ }
554+
406555 #[ tokio:: test]
407556 async fn insert_or_update_does_not_mutate_memory_if_persist_fails ( ) {
408557 let existing_id = TestObjectId { id : [ 42u8 ; 4 ] } ;
0 commit comments