diff --git a/src/storage/migration_signal.rs b/src/storage/migration_signal.rs index e698561b..6ae6bc7e 100644 --- a/src/storage/migration_signal.rs +++ b/src/storage/migration_signal.rs @@ -332,10 +332,22 @@ pub async fn tally_peers(p2p: &Arc) -> PeerTally { let transport = p2p.transport(); let observer = p2p.peer_id().to_hex(); for peer in transport.connected_peers().await { - // No agent recorded is not the same as a peer that reported nothing, but it is - // just as far from evidence of completion, so it lands in the same bucket rather - // than being skipped. - let agent = transport.peer_user_agent(&peer).await; + let mut agent = transport.peer_user_agent(&peer).await; + if agent.is_none() { + // saorsa-core records a peer's agent and drops it together with the peer's + // connection entry, so no agent means the peer was not connected when it was read: + // it disconnected after the list above was taken. A peer that is still gone is not + // a peer this node can see, so it is skipped. Counting it put a departing client in + // the unreported bucket on a testnet. + if !transport.is_peer_connected(&peer).await { + continue; + } + // It came back between the two reads, so the first one is stale and would put a + // peer that may have finished in the unreported bucket. Read what it announced this + // time. If that is still nothing, it stays in the unreported bucket: no agent is not + // evidence of completion. + agent = transport.peer_user_agent(&peer).await; + } let state = agent .as_deref() .map_or(PeerMigrationState::Unreported, peer_state); diff --git a/tests/migration_reclaims_disk.rs b/tests/migration_reclaims_disk.rs index 55cf07f3..0bd76c64 100644 --- a/tests/migration_reclaims_disk.rs +++ b/tests/migration_reclaims_disk.rs @@ -24,7 +24,11 @@ )] use ant_node::storage::migration::{MigrationPhase, MIN_RETIRE_DELAY_HOURS}; -use ant_node::storage::{ChunkStore, ChunkStoreConfig, LmdbStorage, LmdbStorageConfig}; +use ant_node::storage::{ + ChunkStore, ChunkStoreConfig, LmdbStorage, LmdbStorageConfig, LEGACY_ENV_DIR, +}; +#[cfg(target_os = "linux")] +use std::os::fd::AsRawFd; use std::path::Path; use std::sync::Arc; use tempfile::TempDir; @@ -86,10 +90,55 @@ fn walk(path: &Path, size: &dyn Fn(&std::fs::Metadata) -> u64) -> u64 { /// disappearing proves nothing: unlink a file that something still holds open and every /// name is gone while every block is still spoken for, which is a fair description of the /// bug that started all this. +/// +/// On Linux, read only once the filesystem has committed what it has been given. The free +/// space btrfs reports can lag what it has allocated and freed until its transaction +/// commits, every 30 seconds by default, and the fsync after each chunk file does not +/// commit it. Read straight after the copy, the peak has missed up to 28 MB of the file +/// store, and the recovery measured from it then came up short however much space came +/// back. `syncfs` is Linux only, so other targets read straight away, as before. fn free_space(path: &Path) -> u64 { + #[cfg(target_os = "linux")] + { + let committed = commit_filesystem(path); + assert!( + committed.is_ok(), + "could not commit the filesystem before reading it: {committed:?}" + ); + } fs2::available_space(path).expect("the filesystem should report its free space") } +/// `syncfs(2)` on the filesystem holding `path`. It returns once that filesystem has +/// written out its pending changes, which on btrfs means once the running transaction +/// has committed. +#[cfg(target_os = "linux")] +fn commit_filesystem(path: &Path) -> std::io::Result<()> { + let dir = std::fs::File::open(path)?; + // SAFETY: `dir` is open for the whole call, so its descriptor is valid, and `syncfs` + // only reads the descriptor and keeps nothing. + #[allow(clippy::undocumented_unsafe_blocks, unsafe_code)] + let synced = unsafe { libc::syncfs(dir.as_raw_fd()) }; + if synced == 0 { + Ok(()) + } else { + Err(std::io::Error::last_os_error()) + } +} + +/// What is left under `root` of the legacy environment: the environment itself, and any +/// tombstone retirement renamed it to. +fn legacy_entries(root: &Path) -> std::io::Result> { + let mut left = Vec::new(); + for entry in std::fs::read_dir(root)? { + let name = entry?.file_name().to_string_lossy().into_owned(); + if name.starts_with(LEGACY_ENV_DIR) { + left.push(name); + } + } + Ok(left) +} + /// Deterministic content for chunk `n`, filled so it does not compress to nothing. /// /// `n` goes in verbatim at the front rather than being folded into the fill, because a @@ -250,16 +299,23 @@ async fn retiring_the_legacy_environment_returns_its_bytes_to_the_filesystem() { assert_eq!(store.migration_phase(), MigrationPhase::FilesOnly); assert!(freed > 0, "retirement reported no bytes freed"); - // The deletion runs on a detached thread so the node can serve while it happens. - for _ in 0..400 { - if !environment.exists() && allocated_bytes(&root) < environment_blocks + payload { + // The deletion runs on a detached thread so the node can serve while it happens. The + // thread renames the environment to a tombstone first and removes the tombstone last, + // so the old name is gone before anything is deleted and only the tombstone going says + // the deletion is over. Space read before then comes back in pieces: on btrfs, a + // reading taken mid-delete has seen half the environment back. The wait is up to 80 + // seconds, as long as the old wait and the space poll after it gave the deletion + // between them, which leaves room for the thread to retry after 10, 30 and 60 seconds. + for _ in 0..1600 { + if legacy_entries(&root).is_ok_and(|left| left.is_empty()) { break; } tokio::time::sleep(std::time::Duration::from_millis(50)).await; } + let left = legacy_entries(&root); assert!( - !environment.exists(), - "the environment directory is still on disk" + left.as_ref().is_ok_and(Vec::is_empty), + "the environment is still on disk: {left:?}" ); // Dropping the store closes every handle. A file that is unlinked while something