Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
20 changes: 16 additions & 4 deletions src/storage/migration_signal.rs
Original file line number Diff line number Diff line change
Expand Up @@ -332,10 +332,22 @@ pub async fn tally_peers(p2p: &Arc<P2PNode>) -> 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);
Expand Down
68 changes: 62 additions & 6 deletions tests/migration_reclaims_disk.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -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<Vec<String>> {
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
Expand Down Expand Up @@ -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
Expand Down
Loading