Skip to content

Commit 10b9721

Browse files
committed
fixup! Implement tiered storage
Preserve call-time ordering for ephemeral writes and removes by routing them through the same versioned lock path as primary-backed mutations. Add regression coverage for stale ephemeral writes and removes.
1 parent ba1c36a commit 10b9721

1 file changed

Lines changed: 117 additions & 20 deletions

File tree

src/io/tier_store.rs

Lines changed: 117 additions & 20 deletions
Original file line numberDiff line numberDiff line change
@@ -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

Comments
 (0)