feat(peer): coordinate outbound transfers with local game mutations
Updating or removing a local game rewrites its on-disk files. Peers that were mid-download of that game would keep streaming bytes from files that are being deleted or replaced, handing them a corrupt or stale copy. There was also no authoritative notion of which game version a peer should serve or accept, so a peer could serve whatever happened to be on disk and downloaders could aggregate files from peers running mismatched versions. This introduces a reader-writer coordination scheme between outbound file transfers (readers) and local mutation operations (writers), and gates both serving and downloading on an authoritative game catalog version. Reader-writer coordination: - Track active outbound transfers per game in a shared `OutboundTransfers` map of (id, CancellationToken), threaded through `Ctx`/`PeerCtx` and registered by a `TransferGuard` in the stream service. The guard is registered *before* the serve-eligibility check to close a TOCTOU window where a writer could miss an in-flight reader. - `stream_file_bytes` now honors a cancellation token at every await point (file read, network send, stream close) via `tokio::select!`, so a transfer aborts promptly instead of hanging on a stalled receiver. - `begin_operation` marks a game active first, then cancels its outbound transfers and waits for the count to reach zero before any Updating/RemovingDownload work touches the filesystem. - Active games are now hidden from library snapshots entirely while an operation is in flight, instead of freezing their last announced state, so peers stop discovering a game that is being mutated. Authoritative version catalog: - Replace the `HashSet<String>` catalog with `GameCatalog`, mapping each game id to its expected version (from the bundled game.db / ETI data). - Serving requires the local `version.ini` to match the catalog version (`local_download_matches_catalog`); peer selection, file aggregation, and majority size validation all filter on the expected version (`peers_with_expected_version`, `aggregated_game_files`, and friends). User-visible changes: - The GUI shows confirmation dialogs before Update and Remove, and surfaces a sharing-status indicator on game cards and the detail modal. - A new `OutboundTransferCountChanged` event lets the UI reflect live outbound transfer activity. Test Plan: - just test - just frontend-test - just clippy
This commit is contained in:
24 files changed
+882
-212
No files matched your search
@@ -7,7 +7,7 @@ use std::{
|
||||
time::{Duration, Instant},
|
||||
};
|
||||
|
||||
use lanspread_db::db::{Availability, Game, GameFileDescription};
|
||||
use lanspread_db::db::{Availability, Game, GameCatalog, GameFileDescription};
|
||||
use lanspread_proto::{GameSummary, LibraryDelta, LibrarySnapshot};
|
||||
|
||||
use crate::library::compute_library_digest;
|
||||
@@ -357,6 +357,54 @@ impl PeerGameDB {
|
||||
games
|
||||
}
|
||||
|
||||
/// Returns catalog games aggregated from peers that advertise the expected catalog version.
|
||||
#[must_use]
|
||||
pub fn get_catalog_games(&self, catalog: &GameCatalog) -> Vec<Game> {
|
||||
let mut aggregated: HashMap<String, Game> = HashMap::new();
|
||||
let mut peer_counts: HashMap<String, u32> = HashMap::new();
|
||||
|
||||
for peer in self.peers.values() {
|
||||
for game in peer.games.values().filter(|game| {
|
||||
catalog.contains(&game.id)
|
||||
&& game_matches_expected_version(game, catalog.expected_version(&game.id))
|
||||
}) {
|
||||
*peer_counts.entry(game.id.clone()).or_insert(0) += 1;
|
||||
}
|
||||
}
|
||||
|
||||
for peer in self.peers.values() {
|
||||
for game in peer.games.values().filter(|game| {
|
||||
catalog.contains(&game.id)
|
||||
&& game_matches_expected_version(game, catalog.expected_version(&game.id))
|
||||
}) {
|
||||
aggregated
|
||||
.entry(game.id.clone())
|
||||
.and_modify(|existing| {
|
||||
existing.peer_count = *peer_counts.get(&game.id).unwrap_or(&0);
|
||||
if game.size > existing.size {
|
||||
existing.size = game.size;
|
||||
}
|
||||
existing.set_downloaded(true);
|
||||
if game.installed {
|
||||
existing.installed = true;
|
||||
}
|
||||
})
|
||||
.or_insert_with(|| {
|
||||
let mut game_clone = summary_to_game(game);
|
||||
if let Some(expected_version) = catalog.expected_version(&game.id) {
|
||||
game_clone.eti_game_version = Some(expected_version.to_string());
|
||||
}
|
||||
game_clone.peer_count = *peer_counts.get(&game.id).unwrap_or(&0);
|
||||
game_clone
|
||||
});
|
||||
}
|
||||
}
|
||||
|
||||
let mut games: Vec<Game> = aggregated.into_values().collect();
|
||||
games.sort_by(|a, b| a.name.cmp(&b.name));
|
||||
games
|
||||
}
|
||||
|
||||
/// Returns the latest version of a game across all peers.
|
||||
#[must_use]
|
||||
pub fn get_latest_version_for_game(&self, game_id: &str) -> Option<String> {
|
||||
@@ -451,6 +499,24 @@ impl PeerGameDB {
|
||||
.collect()
|
||||
}
|
||||
|
||||
/// Returns addresses of peers that have the expected catalog version of a game.
|
||||
#[must_use]
|
||||
pub fn peers_with_expected_version(
|
||||
&self,
|
||||
game_id: &str,
|
||||
expected_version: Option<&str>,
|
||||
) -> Vec<SocketAddr> {
|
||||
self.peers
|
||||
.iter()
|
||||
.filter(|(_, peer)| {
|
||||
peer.games
|
||||
.get(game_id)
|
||||
.is_some_and(|game| game_matches_expected_version(game, expected_version))
|
||||
})
|
||||
.map(|(_, peer)| peer.addr)
|
||||
.collect()
|
||||
}
|
||||
|
||||
/// Returns addresses of peers that have the latest version of a game.
|
||||
#[must_use]
|
||||
pub fn peers_with_latest_version(&self, game_id: &str) -> Vec<SocketAddr> {
|
||||
@@ -514,11 +580,33 @@ impl PeerGameDB {
|
||||
.collect()
|
||||
}
|
||||
|
||||
/// Returns file descriptions from peers that advertise the expected catalog version.
|
||||
#[must_use]
|
||||
pub fn expected_version_game_files_for(
|
||||
&self,
|
||||
game_id: &str,
|
||||
expected_version: Option<&str>,
|
||||
) -> Vec<(SocketAddr, Vec<GameFileDescription>)> {
|
||||
let expected_peers = self.peers_with_expected_version(game_id, expected_version);
|
||||
if expected_peers.is_empty() {
|
||||
return Vec::new();
|
||||
}
|
||||
|
||||
self.game_files_for(game_id)
|
||||
.into_iter()
|
||||
.filter(|(addr, _)| expected_peers.contains(addr))
|
||||
.collect()
|
||||
}
|
||||
|
||||
/// Returns aggregated file descriptions for a game across all peers.
|
||||
#[must_use]
|
||||
pub fn aggregated_game_files(&self, game_id: &str) -> Vec<GameFileDescription> {
|
||||
pub fn aggregated_game_files(
|
||||
&self,
|
||||
game_id: &str,
|
||||
expected_version: Option<&str>,
|
||||
) -> Vec<GameFileDescription> {
|
||||
let mut seen: HashMap<String, GameFileDescription> = HashMap::new();
|
||||
for (_, files) in self.latest_game_files_for(game_id) {
|
||||
for (_, files) in self.expected_version_game_files_for(game_id, expected_version) {
|
||||
for file in files {
|
||||
seen.entry(file.relative_path.clone()).or_insert(file);
|
||||
}
|
||||
@@ -559,8 +647,9 @@ impl PeerGameDB {
|
||||
pub fn validate_file_sizes_majority(
|
||||
&self,
|
||||
game_id: &str,
|
||||
expected_version: Option<&str>,
|
||||
) -> eyre::Result<MajorityValidationResult> {
|
||||
let game_files = self.latest_game_files_for(game_id);
|
||||
let game_files = self.expected_version_game_files_for(game_id, expected_version);
|
||||
if game_files.is_empty() {
|
||||
return Ok((Vec::new(), Vec::new(), HashMap::new()));
|
||||
}
|
||||
@@ -813,6 +902,14 @@ fn game_is_ready(summary: &GameSummary) -> bool {
|
||||
summary.availability == Availability::Ready
|
||||
}
|
||||
|
||||
fn game_matches_expected_version(summary: &GameSummary, expected_version: Option<&str>) -> bool {
|
||||
if !game_is_ready(summary) {
|
||||
return false;
|
||||
}
|
||||
|
||||
expected_version.is_none_or(|expected| summary.eti_version.as_deref() == Some(expected))
|
||||
}
|
||||
|
||||
fn summary_to_game(summary: &GameSummary) -> Game {
|
||||
let eti_game_version = game_is_ready(summary)
|
||||
.then(|| summary.eti_version.clone())
|
||||
@@ -925,6 +1022,41 @@ mod tests {
|
||||
assert!(db.peers_with_latest_version("game").is_empty());
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn catalog_aggregation_counts_only_expected_version_peers() {
|
||||
let old_addr = addr(12003);
|
||||
let expected_addr = addr(12004);
|
||||
let newer_addr = addr(12005);
|
||||
let mut db = PeerGameDB::new();
|
||||
db.upsert_peer("old".to_string(), old_addr);
|
||||
db.upsert_peer("expected".to_string(), expected_addr);
|
||||
db.upsert_peer("newer".to_string(), newer_addr);
|
||||
db.update_peer_games(
|
||||
&"old".to_string(),
|
||||
vec![summary("game", "20240101", Availability::Ready)],
|
||||
);
|
||||
db.update_peer_games(
|
||||
&"expected".to_string(),
|
||||
vec![summary("game", "20250101", Availability::Ready)],
|
||||
);
|
||||
db.update_peer_games(
|
||||
&"newer".to_string(),
|
||||
vec![summary("game", "20260101", Availability::Ready)],
|
||||
);
|
||||
let mut catalog = GameCatalog::empty();
|
||||
catalog.insert("game".to_string(), Some("20250101".to_string()));
|
||||
|
||||
let games = db.get_catalog_games(&catalog);
|
||||
|
||||
assert_eq!(games.len(), 1);
|
||||
assert_eq!(games[0].peer_count, 1);
|
||||
assert_eq!(games[0].eti_game_version.as_deref(), Some("20250101"));
|
||||
assert_eq!(
|
||||
db.peers_with_expected_version("game", Some("20250101")),
|
||||
vec![expected_addr]
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn transport_addr_matches_known_peer_on_ephemeral_port() {
|
||||
let advertised = ip_addr([10, 66, 0, 2], 40000);
|
||||
@@ -979,7 +1111,7 @@ mod tests {
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn validation_uses_latest_version_file_metadata() {
|
||||
fn validation_uses_expected_version_file_metadata() {
|
||||
let old_addr = addr(12003);
|
||||
let new_addr = addr(12004);
|
||||
let mut db = PeerGameDB::new();
|
||||
@@ -1010,21 +1142,21 @@ mod tests {
|
||||
],
|
||||
);
|
||||
|
||||
let aggregated = db.aggregated_game_files("game");
|
||||
let aggregated = db.aggregated_game_files("game", Some("20250101"));
|
||||
let archive = aggregated
|
||||
.iter()
|
||||
.find(|desc| desc.relative_path == "game/archive.eti")
|
||||
.expect("latest archive should be present");
|
||||
.expect("expected-version archive should be present");
|
||||
assert_eq!(archive.size, 20);
|
||||
|
||||
let (validated, peers, file_peer_map) = db
|
||||
.validate_file_sizes_majority("game")
|
||||
.validate_file_sizes_majority("game", Some("20250101"))
|
||||
.expect("old-version file metadata should not create ambiguity");
|
||||
assert_eq!(peers, vec![new_addr]);
|
||||
let archive = validated
|
||||
.iter()
|
||||
.find(|desc| desc.relative_path == "game/archive.eti")
|
||||
.expect("latest archive should validate");
|
||||
.expect("expected-version archive should validate");
|
||||
assert_eq!(archive.size, 20);
|
||||
assert_eq!(file_peer_map.get("game/archive.eti"), Some(&vec![new_addr]));
|
||||
}
|
||||
|
||||
Reference in new issue
Block a user