@@ -415,16 +415,19 @@ impl TierStoreInner {
415415
416416 if is_ephemeral_cached_key ( & primary_namespace, & secondary_namespace, & key) {
417417 if let Some ( eph_store) = self . ephemeral_store . as_ref ( ) {
418- let res = KVStore :: write (
419- eph_store. as_ref ( ) ,
420- primary_namespace. as_str ( ) ,
421- secondary_namespace. as_str ( ) ,
422- key. as_str ( ) ,
423- buf,
424- )
425- . await ;
426- self . clean_locks ( & lock_ref, locking_key) ;
427- return res;
418+ let eph_store = Arc :: clone ( eph_store) ;
419+ return self
420+ . execute_locked_write ( lock_ref, locking_key, version, || async move {
421+ KVStore :: write (
422+ eph_store. as_ref ( ) ,
423+ primary_namespace. as_str ( ) ,
424+ secondary_namespace. as_str ( ) ,
425+ key. as_str ( ) ,
426+ buf,
427+ )
428+ . await
429+ } )
430+ . await ;
428431 }
429432 }
430433
@@ -453,16 +456,19 @@ impl TierStoreInner {
453456
454457 if is_ephemeral_cached_key ( & primary_namespace, & secondary_namespace, & key) {
455458 if let Some ( eph_store) = self . ephemeral_store . as_ref ( ) {
456- let res = KVStore :: remove (
457- eph_store. as_ref ( ) ,
458- primary_namespace. as_str ( ) ,
459- secondary_namespace. as_str ( ) ,
460- key. as_str ( ) ,
461- lazy,
462- )
463- . await ;
464- self . clean_locks ( & lock_ref, locking_key) ;
465- return res;
459+ let eph_store = Arc :: clone ( eph_store) ;
460+ return self
461+ . execute_locked_write ( lock_ref, locking_key, version, || async move {
462+ KVStore :: remove (
463+ eph_store. as_ref ( ) ,
464+ primary_namespace. as_str ( ) ,
465+ secondary_namespace. as_str ( ) ,
466+ key. as_str ( ) ,
467+ lazy,
468+ )
469+ . await
470+ } )
471+ . await ;
466472 }
467473 }
468474
@@ -849,6 +855,97 @@ mod tests {
849855 assert_eq ! ( persisted, new_data) ;
850856 }
851857
858+ #[ tokio:: test]
859+ async fn ephemeral_writes_preserve_latest_call_order ( ) {
860+ let base_dir = random_storage_path ( ) ;
861+ let log_path = base_dir. join ( "tier_store_test.log" ) . to_string_lossy ( ) . into_owned ( ) ;
862+ let logger = Arc :: new ( Logger :: new_fs_writer ( log_path, Level :: Trace ) . unwrap ( ) ) ;
863+
864+ let _cleanup = CleanupDir ( base_dir. clone ( ) ) ;
865+
866+ let primary_store: Arc < DynStore > =
867+ Arc :: new ( DynStoreWrapper ( FilesystemStoreV2 :: new ( base_dir. join ( "primary" ) ) . unwrap ( ) ) ) ;
868+ let mut tier = setup_tier_store ( primary_store, logger) ;
869+
870+ let ephemeral_store: Arc < DynStore > =
871+ Arc :: new ( DynStoreWrapper ( FilesystemStoreV2 :: new ( base_dir. join ( "ephemeral" ) ) . unwrap ( ) ) ) ;
872+ tier. set_ephemeral_store ( ephemeral_store) ;
873+
874+ let old_data = vec ! [ 1u8 ; 32 ] ;
875+ let new_data = vec ! [ 2u8 ; 32 ] ;
876+
877+ let old_write = tier. write (
878+ NETWORK_GRAPH_PERSISTENCE_PRIMARY_NAMESPACE ,
879+ NETWORK_GRAPH_PERSISTENCE_SECONDARY_NAMESPACE ,
880+ NETWORK_GRAPH_PERSISTENCE_KEY ,
881+ old_data,
882+ ) ;
883+ let new_write = tier. write (
884+ NETWORK_GRAPH_PERSISTENCE_PRIMARY_NAMESPACE ,
885+ NETWORK_GRAPH_PERSISTENCE_SECONDARY_NAMESPACE ,
886+ NETWORK_GRAPH_PERSISTENCE_KEY ,
887+ new_data. clone ( ) ,
888+ ) ;
889+
890+ new_write. await . unwrap ( ) ;
891+ old_write. await . unwrap ( ) ;
892+
893+ let persisted = tier
894+ . read (
895+ NETWORK_GRAPH_PERSISTENCE_PRIMARY_NAMESPACE ,
896+ NETWORK_GRAPH_PERSISTENCE_SECONDARY_NAMESPACE ,
897+ NETWORK_GRAPH_PERSISTENCE_KEY ,
898+ )
899+ . await
900+ . unwrap ( ) ;
901+ assert_eq ! ( persisted, new_data) ;
902+ }
903+
904+ #[ tokio:: test]
905+ async fn ephemeral_removes_preserve_latest_call_order ( ) {
906+ let base_dir = random_storage_path ( ) ;
907+ let log_path = base_dir. join ( "tier_store_test.log" ) . to_string_lossy ( ) . into_owned ( ) ;
908+ let logger = Arc :: new ( Logger :: new_fs_writer ( log_path, Level :: Trace ) . unwrap ( ) ) ;
909+
910+ let _cleanup = CleanupDir ( base_dir. clone ( ) ) ;
911+
912+ let primary_store: Arc < DynStore > =
913+ Arc :: new ( DynStoreWrapper ( FilesystemStoreV2 :: new ( base_dir. join ( "primary" ) ) . unwrap ( ) ) ) ;
914+ let mut tier = setup_tier_store ( primary_store, logger) ;
915+
916+ let ephemeral_store: Arc < DynStore > =
917+ Arc :: new ( DynStoreWrapper ( FilesystemStoreV2 :: new ( base_dir. join ( "ephemeral" ) ) . unwrap ( ) ) ) ;
918+ tier. set_ephemeral_store ( ephemeral_store) ;
919+
920+ let data = vec ! [ 2u8 ; 32 ] ;
921+
922+ let stale_remove = tier. remove (
923+ NETWORK_GRAPH_PERSISTENCE_PRIMARY_NAMESPACE ,
924+ NETWORK_GRAPH_PERSISTENCE_SECONDARY_NAMESPACE ,
925+ NETWORK_GRAPH_PERSISTENCE_KEY ,
926+ true ,
927+ ) ;
928+ let new_write = tier. write (
929+ NETWORK_GRAPH_PERSISTENCE_PRIMARY_NAMESPACE ,
930+ NETWORK_GRAPH_PERSISTENCE_SECONDARY_NAMESPACE ,
931+ NETWORK_GRAPH_PERSISTENCE_KEY ,
932+ data. clone ( ) ,
933+ ) ;
934+
935+ new_write. await . unwrap ( ) ;
936+ stale_remove. await . unwrap ( ) ;
937+
938+ let persisted = tier
939+ . read (
940+ NETWORK_GRAPH_PERSISTENCE_PRIMARY_NAMESPACE ,
941+ NETWORK_GRAPH_PERSISTENCE_SECONDARY_NAMESPACE ,
942+ NETWORK_GRAPH_PERSISTENCE_KEY ,
943+ )
944+ . await
945+ . unwrap ( ) ;
946+ assert_eq ! ( persisted, data) ;
947+ }
948+
852949 #[ tokio:: test]
853950 async fn backup_write_is_part_of_success_path ( ) {
854951 let base_dir = random_storage_path ( ) ;
0 commit comments