//! Command handlers for peer commands. use std::{ collections::{HashMap, HashSet}, fmt, future::Future, net::IpAddr, path::{Path, PathBuf}, sync::Arc, time::Duration, }; use lanspread_db::{ content_manifest::{CatalogContentManifest, ContentId}, db::GameDB, }; use lanspread_proto::PeerEndpoint; use tokio_util::sync::CancellationToken; #[cfg(test)] use crate::local_games::{rescan_local_game, scan_local_library}; use crate::{ DownloadAttemptKey, DownloadFailureReason, DownloadVerificationActivity, InstallOperation, PeerEvent, PeerEventSender, StreamInstallSettings, apply_launch_settings_to_verified_tree, content_quarantine::ContentQuarantine, context::{Ctx, OperationGuard, OperationKind}, download::{ DownloadCompletion, DownloadGameRequest, DownloadOwnershipReadiness, ValidatedDownloadManifest, download_game_files, download_ownership_matches_content, download_ownership_readiness, remove_downloaded_payload, }, events, install, local_games::{ LocalLibraryScan, game_from_summary, local_dir_is_directory, local_download_matches_catalog, rescan_local_game_with_recovery_failures, scan_local_library_with_recovery_failures, version_ini_is_regular_file, }, mark_launch_settings_applied, peer_db::PeerGameDB, quic_runtime::QuicConnector, scoped_blocking::scoped_blocking, services::{HandshakeCtx, ReservedCandidateHandshake}, stream_install::{ ReceiveStreamedInstallRequest, StreamInstallReceiveError, StreamInstallReceiveErrorKind, StreamInstallTotalDeadline, receive_streamed_install, }, transfer_status::{DownloadAttemptReporter, DownloadAttemptStatus}, }; // ============================================================================= // Command handlers // ============================================================================= const OUTBOUND_TRANSFER_DRAIN_POLL_INTERVAL: Duration = Duration::from_millis(10); const OUTBOUND_TRANSFER_DRAIN_TIMEOUT: Duration = Duration::from_secs(5); async fn retain_guard_until_operation_completes(guard: G, operation: F) -> F::Output where F: Future, { let output = operation.await; drop(guard); output } async fn register_download_attempt( ctx: &Ctx, tx_notify_ui: &PeerEventSender, key: DownloadAttemptKey, cancellation: CancellationToken, ) -> DownloadAttemptStatus { let status = DownloadAttemptStatus::new(key, cancellation, tx_notify_ui.clone()); ctx.active_downloads .write() .await .insert(status.key().id.clone(), status.signal()); status } /// Immutable filesystem target selected when a command enters the peer core. /// /// The configured games directory may change while a spawned operation waits to /// run. Admission compares this captured target with the configured directory, /// and every later filesystem step uses these paths rather than rereading /// `Ctx::game_dir`. #[derive(Clone, Debug, Eq, PartialEq)] struct OperationTarget { games_folder: PathBuf, game_id: String, } impl OperationTarget { fn new(games_folder: PathBuf, game_id: String) -> Self { Self { games_folder, game_id, } } fn game_id(&self) -> &str { &self.game_id } fn game_root(&self) -> PathBuf { self.games_folder.join(&self.game_id) } } fn select_content_sources( peer_game_db: &PeerGameDB, game_id: &str, content_id: ContentId, quarantine: &ContentQuarantine, ) -> Vec { let mut endpoints = peer_game_db.peer_endpoints_with_content(game_id, content_id); endpoints.sort_unstable_by_key(|endpoint| (endpoint.peer_id, endpoint.addr)); endpoints.dedup(); endpoints .into_iter() .filter(|endpoint| !quarantine.is_quarantined(endpoint, content_id)) .collect() } fn load_catalog_manifest(ctx: &Ctx, game_id: &str) -> eyre::Result> { let catalog = Arc::clone(&ctx.catalog); let game_id = game_id.to_owned(); scoped_blocking(move || catalog.manifest(&game_id)) } #[derive(Clone, Copy, Debug, Eq, PartialEq)] enum StreamInstallFailureDisposition { RetryAndQuarantine, Retry, Stop, } const fn stream_install_failure_disposition( kind: StreamInstallReceiveErrorKind, ) -> StreamInstallFailureDisposition { match kind { StreamInstallReceiveErrorKind::Integrity => { StreamInstallFailureDisposition::RetryAndQuarantine } StreamInstallReceiveErrorKind::Transport => StreamInstallFailureDisposition::Retry, StreamInstallReceiveErrorKind::Cancelled | StreamInstallReceiveErrorKind::Setup => { StreamInstallFailureDisposition::Stop } } } const MAX_STREAM_INSTALL_SOURCE_IP_ATTEMPTS: usize = 4; #[derive(Clone, Copy, Debug, Eq, PartialEq)] enum StreamInstallSourceAdmission { Admitted, AlreadyAttempted, Exhausted, } #[derive(Default)] struct StreamInstallSourceAttempts { source_ips: HashSet, } impl StreamInstallSourceAttempts { fn admit(&mut self, source_ip: IpAddr) -> StreamInstallSourceAdmission { if self.source_ips.contains(&source_ip) { return StreamInstallSourceAdmission::AlreadyAttempted; } if self.source_ips.len() >= MAX_STREAM_INSTALL_SOURCE_IP_ATTEMPTS { return StreamInstallSourceAdmission::Exhausted; } self.source_ips.insert(source_ip); StreamInstallSourceAdmission::Admitted } } fn admit_stream_install_source( attempts: &mut StreamInstallSourceAttempts, source: &PeerEndpoint, game_id: &str, ) -> StreamInstallSourceAdmission { let admission = attempts.admit(source.addr.ip()); match admission { StreamInstallSourceAdmission::Admitted => {} StreamInstallSourceAdmission::AlreadyAttempted => log::debug!( "Skipping streamed-install source {} at {} for {game_id}: endpoint IP was already attempted", source.peer_id, source.addr ), StreamInstallSourceAdmission::Exhausted => log::warn!( "Streamed install for {game_id} exhausted its {MAX_STREAM_INSTALL_SOURCE_IP_ATTEMPTS}-source-IP attempt budget" ), } admission } #[derive(Debug)] struct StreamDownloadError { reason: Option, error: eyre::Report, } impl StreamDownloadError { fn cancelled(error: impl Into) -> Self { Self { reason: None, error: error.into(), } } fn sources_exhausted(error: impl Into) -> Self { Self { reason: Some(DownloadFailureReason::VerifiedCatalogSourcesExhausted), error: error.into(), } } fn operation_failed(error: impl Into) -> Self { Self { reason: Some(DownloadFailureReason::OperationFailed), error: error.into(), } } } impl fmt::Display for StreamDownloadError { fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { self.error.fmt(formatter) } } /// Handles the `ListGames` command. pub async fn handle_list_games_command(ctx: &Ctx, tx_notify_ui: &PeerEventSender) { log::info!("ListGames command received"); events::emit_peer_game_list(&ctx.peer_game_db, tx_notify_ui).await; } /// Handles the `DownloadGameFiles` command. #[allow(clippy::too_many_lines)] pub async fn handle_download_game_files_command( ctx: &Ctx, tx_notify_ui: &PeerEventSender, id: String, install_after_download: bool, ) { log::info!("Got PeerCommand::DownloadGameFiles"); let attempt = DownloadAttemptKey::next(id.clone()); if !catalog_contains(ctx, &id) { log::warn!("Ignoring download command for non-catalog game {id}"); send_download_failed( tx_notify_ui, &attempt, DownloadFailureReason::OperationFailed, ); return; } let games_folder = { ctx.game_dir.read().await.clone() }; let target = OperationTarget::new(games_folder.clone(), id.clone()); let catalog_manifest = match load_catalog_manifest(ctx, &id) { Ok(manifest) => manifest, Err(error) => { log::error!("Failed to load catalog content manifest for {id}: {error}"); send_download_failed( tx_notify_ui, &attempt, DownloadFailureReason::OperationFailed, ); return; } }; let content_id = catalog_manifest.content_id(); let manifest_games_folder = games_folder.clone(); let manifest = scoped_blocking(move || { ValidatedDownloadManifest::from_catalog(&manifest_games_folder, catalog_manifest) }); let manifest = match manifest { Ok(manifest) => manifest, Err(error) => { log::error!("Rejected catalog download manifest for {id}: {error}"); send_download_failed( tx_notify_ui, &attempt, DownloadFailureReason::OperationFailed, ); return; } }; let local_dl_available = { if ctx.recovery_quarantine.is_blocked(&games_folder, &id) { false } else { let active_operations = ctx.active_operations.read().await; local_download_matches_catalog( &games_folder, ctx.state_dir.as_ref(), &id, &active_operations, ctx.catalog.catalog(), ) .await && download_ownership_matches_content( &games_folder, ctx.state_dir.as_ref(), &id, content_id, ) .await } }; let network_permit = match ctx.network.try_acquire() { Ok(permit) => permit, Err(error) if local_dl_available => { log::info!( "Using locally downloaded files for game {id} while Local network sharing is unavailable: {error}" ); finish_cached_download(ctx, tx_notify_ui, &attempt, install_after_download, target); return; } Err(error) => { log::warn!("Cannot download {id} while Local network sharing is unavailable: {error}"); send_download_failed( tx_notify_ui, &attempt, DownloadFailureReason::OperationFailed, ); return; } }; let sources = { let peer_game_db = ctx.peer_game_db.read().await; select_content_sources(&peer_game_db, &id, content_id, &ctx.content_quarantine) }; if sources.is_empty() { if local_dl_available { finish_cached_download(ctx, tx_notify_ui, &attempt, install_after_download, target); } else { log::error!("No eligible exact catalog content source available for game {id}"); send_download_failed( tx_notify_ui, &attempt, DownloadFailureReason::VerifiedCatalogSourcesExhausted, ); } return; } match begin_operation(ctx, tx_notify_ui, &target, OperationKind::Downloading).await { BeginOperationResult::Started => {} BeginOperationResult::AlreadyActive => { log::warn!("Operation for {id} already in progress; ignoring new download request"); return; } BeginOperationResult::DrainTimedOut => { log::error!("Timed out waiting for outbound transfers before downloading {id}"); send_download_failed( tx_notify_ui, &attempt, DownloadFailureReason::OperationFailed, ); return; } BeginOperationResult::PublicationFailed => { log::error!("Failed to withdraw published availability before downloading {id}"); send_download_failed( tx_notify_ui, &attempt, DownloadFailureReason::OperationFailed, ); return; } BeginOperationResult::RecoveryBlocked => { log::warn!("Ignoring download for recovery-blocked game {id}"); send_download_failed( tx_notify_ui, &attempt, DownloadFailureReason::OperationFailed, ); return; } BeginOperationResult::GameDirChanged => { log::warn!( "Game directory changed while preparing download for {id}; retry the request against the current library" ); send_download_failed( tx_notify_ui, &attempt, DownloadFailureReason::OperationFailed, ); return; } } let tx_notify_ui_clone = tx_notify_ui.clone(); let download_id = id.clone(); let cancel_token = network_permit.child_token(); let quic = network_permit.connector().clone(); let ctx_clone = ctx.clone(); let operation_target = target.clone(); let download_attempt = register_download_attempt(ctx, tx_notify_ui, attempt, cancel_token.clone()).await; ctx.task_tracker .spawn(retain_guard_until_operation_completes( network_permit, async move { let download_state_guard = OperationGuard::download(download_id.clone(), cancel_token.clone()); let result = download_game_files(DownloadGameRequest { attempt: &download_attempt, manifest, state_dir: ctx_clone.state_dir.as_ref(), sources: &sources, content_id, quarantine: &ctx_clone.content_quarantine, tx_notify_ui: tx_notify_ui_clone.clone(), cancel_token: cancel_token.clone(), quic, }) .await; download_attempt.close_source_admission(); download_attempt.clear_activity(); let result = match result { Ok(completion) => { settle_download_completion(&ctx_clone, &operation_target, completion) .await .map_err(|error| (error, Some(DownloadFailureReason::OperationFailed))) } Err(error) => { let reason = error.reason(); Err((error.into_report(), reason)) } }; match result { Ok(()) => { let Some(prepared) = prepare_install_operation( &ctx_clone, &tx_notify_ui_clone, &operation_target, ) .await else { finish_successful_download_after_refresh( &ctx_clone, &tx_notify_ui_clone, &operation_target, download_state_guard, &download_attempt, ) .await; return; }; if install_after_download { if transition_download_to_install( &ctx_clone, &tx_notify_ui_clone, &download_id, prepared.operation_kind, ) .await { clear_active_download(&ctx_clone, download_attempt.key()).await; download_attempt.emit_finished(); run_started_install_operation( &ctx_clone, &tx_notify_ui_clone, operation_target.clone(), prepared, download_state_guard, cancel_token.clone(), ) .await; } else { finish_successful_download_after_refresh( &ctx_clone, &tx_notify_ui_clone, &operation_target, download_state_guard, &download_attempt, ) .await; } } else { finish_successful_download_after_refresh( &ctx_clone, &tx_notify_ui_clone, &operation_target, download_state_guard, &download_attempt, ) .await; } } Err((e, reason)) => { let download_was_cancelled = cancel_token.is_cancelled(); if download_was_cancelled { log::info!("Download cancelled for {download_id}: {e}"); } else { log::error!("Download failed for {download_id}: {e}"); } finish_failed_download_after_refresh( &ctx_clone, &tx_notify_ui_clone, &operation_target, download_state_guard, &download_attempt, reason, ) .await; } } }, )); } fn finish_cached_download( ctx: &Ctx, tx_notify_ui: &PeerEventSender, attempt: &DownloadAttemptKey, install_after_download: bool, target: OperationTarget, ) { let id = &attempt.id; log::info!("Using locally downloaded files for game {id}; skipping peer transfer"); let status = DownloadAttemptStatus::new( attempt.clone(), CancellationToken::new(), tx_notify_ui.clone(), ); status.emit_begin(); status.close_source_admission(); status.emit_finished(); if install_after_download { spawn_install_operation(ctx, tx_notify_ui, target); } } async fn settle_download_completion( ctx: &Ctx, target: &OperationTarget, completion: DownloadCompletion, ) -> eyre::Result<()> { let game_id = target.game_id(); let uncertain_commit = match completion { DownloadCompletion::Durable => None, DownloadCompletion::RecoveryRequired(cause) => { log::warn!( "Download committed for {game_id}, but readiness is quarantined pending recovery: {cause}" ); Some(cause) } }; install::recover_game_root(&target.game_root(), ctx.state_dir.as_ref(), game_id) .await .map_err(|recovery_error| match uncertain_commit { Some(cause) => recovery_error.wrap_err(format!( "download recovery did not settle for {game_id}; original completion error: {cause}" )), None => recovery_error.wrap_err(format!( "download completion recovery did not settle for {game_id}" )), })?; log::info!("Settled committed download state for {game_id} before publication"); Ok(()) } /// Handles the `InstallGame` command. pub async fn handle_install_game_command(ctx: &Ctx, tx_notify_ui: &PeerEventSender, id: String) { let games_folder = ctx.game_dir.read().await.clone(); spawn_install_operation(ctx, tx_notify_ui, OperationTarget::new(games_folder, id)); } async fn stream_install_target_is_ready(ctx: &Ctx, target: &OperationTarget) -> bool { let id = target.game_id(); if ctx.recovery_quarantine.is_blocked(&target.games_folder, id) { log::warn!("Ignoring streamed install command for recovery-blocked game {id}"); return false; } if download_ownership_readiness(&target.games_folder, ctx.state_dir.as_ref(), id).await == DownloadOwnershipReadiness::RecoveryRequired { log::warn!("Ignoring streamed install command for recovery-quarantined game {id}"); return false; } if local_dir_is_directory(&target.game_root()).await { log::warn!("Ignoring streamed install command for already-installed game {id}"); return false; } true } async fn begin_stream_install_operation( ctx: &Ctx, tx_notify_ui: &PeerEventSender, target: &OperationTarget, attempt: &DownloadAttemptKey, ) -> bool { let id = target.game_id(); match begin_operation(ctx, tx_notify_ui, target, OperationKind::Downloading).await { BeginOperationResult::Started => true, BeginOperationResult::AlreadyActive => { log::warn!("Operation for {id} already in progress; ignoring streamed install request"); false } BeginOperationResult::DrainTimedOut => { log::error!("Timed out waiting for outbound transfers before streamed install of {id}"); send_download_failed( tx_notify_ui, attempt, DownloadFailureReason::OperationFailed, ); false } BeginOperationResult::PublicationFailed => { log::error!( "Failed to withdraw published availability before streamed install of {id}" ); send_download_failed( tx_notify_ui, attempt, DownloadFailureReason::OperationFailed, ); false } BeginOperationResult::RecoveryBlocked => { log::warn!("Ignoring streamed install for recovery-blocked game {id}"); send_download_failed( tx_notify_ui, attempt, DownloadFailureReason::OperationFailed, ); false } BeginOperationResult::GameDirChanged => { log::warn!( "Game directory changed while preparing streamed install for {id}; retry the request against the current library" ); send_download_failed( tx_notify_ui, attempt, DownloadFailureReason::OperationFailed, ); false } } } pub async fn handle_stream_install_game_command( ctx: &Ctx, tx_notify_ui: &PeerEventSender, id: String, settings: StreamInstallSettings, ) { let attempt = DownloadAttemptKey::next(id.clone()); if !catalog_contains(ctx, &id) { log::warn!("Ignoring streamed install command for non-catalog game {id}"); send_download_failed( tx_notify_ui, &attempt, DownloadFailureReason::OperationFailed, ); return; } let manifest = match load_catalog_manifest(ctx, &id) { Ok(manifest) => manifest, Err(error) => { log::error!("Failed to load catalog content manifest for {id}: {error}"); send_download_failed( tx_notify_ui, &attempt, DownloadFailureReason::OperationFailed, ); return; } }; if !manifest.supports_streamed_install() { log::warn!("Ignoring streamed install for {id}: catalog has no extracted-file manifest"); send_download_failed( tx_notify_ui, &attempt, DownloadFailureReason::OperationFailed, ); return; } let games_folder = { ctx.game_dir.read().await.clone() }; let target = OperationTarget::new(games_folder, id.clone()); if !stream_install_target_is_ready(ctx, &target).await { send_download_failed( tx_notify_ui, &attempt, DownloadFailureReason::OperationFailed, ); return; } let network_permit = match ctx.network.try_acquire() { Ok(permit) => permit, Err(error) => { log::warn!( "Cannot stream install {id} while Local network sharing is unavailable: {error}" ); send_download_failed( tx_notify_ui, &attempt, DownloadFailureReason::OperationFailed, ); return; } }; if !begin_stream_install_operation(ctx, tx_notify_ui, &target, &attempt).await { return; } let cancel_token = network_permit.child_token(); let quic = network_permit.connector().clone(); let download_attempt = register_download_attempt(ctx, tx_notify_ui, attempt, cancel_token.clone()).await; let download_guard = OperationGuard::download(id.clone(), cancel_token.clone()); // Root-dependent preflight is repeated after admission. A preceding // same-game operation may have changed ownership or install state after the // optimistic check above but before this operation acquired the gate. if !stream_install_target_is_ready(ctx, &target).await { finish_failed_stream_download( ctx, tx_notify_ui, &target, download_guard, &download_attempt, Some(DownloadFailureReason::OperationFailed), ) .await; return; } let ctx_clone = ctx.clone(); let tx_notify_ui = tx_notify_ui.clone(); let operation = StreamInstallOperation { ctx: ctx_clone, tx_notify_ui, target, manifest, settings, quic, cancel_token, download_guard, download_attempt, }; ctx.task_tracker .spawn(retain_guard_until_operation_completes( network_permit, run_stream_install_operation(operation), )); } struct StreamInstallOperation { ctx: Ctx, tx_notify_ui: PeerEventSender, target: OperationTarget, manifest: Arc, settings: StreamInstallSettings, quic: QuicConnector, cancel_token: CancellationToken, download_guard: OperationGuard, download_attempt: DownloadAttemptStatus, } async fn stream_install_sources( ctx: &Ctx, target: &OperationTarget, manifest: &CatalogContentManifest, ) -> Vec { let peer_game_db = ctx.peer_game_db.read().await; select_content_sources( &peer_game_db, target.game_id(), manifest.content_id(), &ctx.content_quarantine, ) } async fn select_stream_install_sources_or_finish( ctx: &Ctx, tx_notify_ui: &PeerEventSender, target: &OperationTarget, manifest: &CatalogContentManifest, download_guard: OperationGuard, download_attempt: &DownloadAttemptStatus, ) -> Option<(Vec, OperationGuard)> { let sources = stream_install_sources(ctx, target, manifest).await; if !sources.is_empty() { return Some((sources, download_guard)); } log::error!( "No eligible exact catalog content source available for streamed install of {}", target.game_id() ); finish_failed_stream_download( ctx, tx_notify_ui, target, download_guard, download_attempt, Some(DownloadFailureReason::VerifiedCatalogSourcesExhausted), ) .await; None } /// Handles the `UninstallGame` command. pub async fn handle_uninstall_game_command(ctx: &Ctx, tx_notify_ui: &PeerEventSender, id: String) { let games_folder = ctx.game_dir.read().await.clone(); let target = OperationTarget::new(games_folder, id); let ctx = ctx.clone(); let tx_notify_ui = tx_notify_ui.clone(); ctx.task_tracker.clone().spawn(async move { run_uninstall_operation(&ctx, &tx_notify_ui, target).await; }); } pub async fn handle_remove_downloaded_game_command( ctx: &Ctx, tx_notify_ui: &PeerEventSender, id: String, ) { let games_folder = ctx.game_dir.read().await.clone(); let target = OperationTarget::new(games_folder, id); let ctx = ctx.clone(); let tx_notify_ui = tx_notify_ui.clone(); ctx.task_tracker.clone().spawn(async move { run_remove_downloaded_operation(&ctx, &tx_notify_ui, target).await; }); } pub async fn handle_cancel_download_command( ctx: &Ctx, _tx_notify_ui: &PeerEventSender, id: String, ) { let signal = ctx.active_downloads.read().await.get(&id).cloned(); let Some(signal) = signal else { log::warn!("Ignoring cancel request for inactive download {id}"); return; }; log::info!("Cancelling download for game {id}"); signal.cancel_silently(); } async fn run_stream_install_operation(operation: StreamInstallOperation) { let StreamInstallOperation { ctx, tx_notify_ui, target, manifest, settings, quic, cancel_token, download_guard, download_attempt, } = operation; let id = target.game_id.clone(); download_attempt.emit_begin(); let Some((sources, download_guard)) = select_stream_install_sources_or_finish( &ctx, &tx_notify_ui, &target, &manifest, download_guard, &download_attempt, ) .await else { return; }; let receive_result = receive_streamed_install_from_peers(StreamInstallReceiveRequest { ctx: &ctx, tx_notify_ui: &tx_notify_ui, target: &target, manifest: &manifest, sources: &sources, quic: &quic, cancel_token: &cancel_token, attempt: download_attempt.reporter(), }) .await; download_attempt.close_source_admission(); download_attempt.clear_activity(); match receive_result { Ok(transaction) => { let promotion = StreamInstallPromotionPreparation { ctx: &ctx, tx_notify_ui: &tx_notify_ui, target: &target, transaction, settings: &settings, cancel_token: &cancel_token, download_guard, download_attempt: &download_attempt, }; let Some((transaction, download_guard)) = prepare_streamed_install_for_promotion(promotion).await else { return; }; if transition_download_to_install(&ctx, &tx_notify_ui, &id, OperationKind::Installing) .await { commit_streamed_install(StreamInstallCommit { ctx, tx_notify_ui, target, transaction, manifest, cancel_token, operation_guard: download_guard, download_attempt, }) .await; return; } if let Err(err) = transaction.rollback() { log::error!("Failed to roll back streamed install for {id}: {err}"); } finish_failed_stream_download( &ctx, &tx_notify_ui, &target, download_guard, &download_attempt, Some(DownloadFailureReason::OperationFailed), ) .await; } Err(err) => { finish_stream_receive_error( &ctx, &tx_notify_ui, &target, download_guard, &download_attempt, err, ) .await; } } } async fn finish_stream_receive_error( ctx: &Ctx, tx_notify_ui: &PeerEventSender, target: &OperationTarget, download_guard: OperationGuard, download_attempt: &DownloadAttemptStatus, error: StreamDownloadError, ) { if error.reason.is_none() { log::info!( "Streamed install download cancelled for {}: {error}", target.game_id() ); } else { log::error!( "Streamed install download failed for {}: {error}", target.game_id() ); } finish_failed_stream_download( ctx, tx_notify_ui, target, download_guard, download_attempt, error.reason, ) .await; } struct StreamInstallPromotionPreparation<'a> { ctx: &'a Ctx, tx_notify_ui: &'a PeerEventSender, target: &'a OperationTarget, transaction: install::StreamedInstallTransaction, settings: &'a StreamInstallSettings, cancel_token: &'a CancellationToken, download_guard: OperationGuard, download_attempt: &'a DownloadAttemptStatus, } async fn prepare_streamed_install_for_promotion( preparation: StreamInstallPromotionPreparation<'_>, ) -> Option<(install::StreamedInstallTransaction, OperationGuard)> { let StreamInstallPromotionPreparation { ctx, tx_notify_ui, target, transaction, settings, cancel_token, download_guard, download_attempt, } = preparation; let id = target.game_id(); if let Err(error) = apply_launch_settings_to_verified_tree( transaction.staging_dir(), Some(settings.account_name()), Some(settings.language()), Some(settings.persona_name()), ) { log::error!( "Failed to apply launch settings to verified streamed install for {id}: {error}" ); if let Err(rollback_error) = transaction.rollback() { log::error!( "Failed to roll back streamed install for {id} after launch-settings failure: {rollback_error}" ); } finish_failed_stream_download( ctx, tx_notify_ui, target, download_guard, download_attempt, Some(DownloadFailureReason::OperationFailed), ) .await; return None; } // Launch-settings mutation is intentionally finite and non-detachable. // Honor cancellation after it reaches a stable tree boundary, before the // verified staging directory can be promoted. if cancel_token.is_cancelled() { log::info!("Streamed install for {id} was cancelled after applying launch settings"); let reason = match transaction.rollback() { Ok(()) => None, Err(rollback_error) => { log::error!( "Failed to roll back cancelled streamed install for {id}: {rollback_error}" ); Some(DownloadFailureReason::OperationFailed) } }; finish_failed_stream_download( ctx, tx_notify_ui, target, download_guard, download_attempt, reason, ) .await; return None; } Some((transaction, download_guard)) } struct StreamInstallReceiveRequest<'a> { ctx: &'a Ctx, tx_notify_ui: &'a PeerEventSender, target: &'a OperationTarget, manifest: &'a Arc, sources: &'a [PeerEndpoint], quic: &'a QuicConnector, cancel_token: &'a CancellationToken, attempt: DownloadAttemptReporter, } fn settle_failed_stream_receive( ctx: &Ctx, source: &PeerEndpoint, content_id: ContentId, id: &str, transaction: install::StreamedInstallTransaction, error: StreamInstallReceiveError, ) -> Result<(StreamInstallFailureDisposition, StreamInstallReceiveError), StreamDownloadError> { let disposition = stream_install_failure_disposition(error.kind()); if disposition == StreamInstallFailureDisposition::RetryAndQuarantine { ctx.content_quarantine .record_integrity_failure(source, content_id); } transaction.rollback().map_err(|rollback_error| { StreamDownloadError::operation_failed(eyre::eyre!( "streamed install attempt from {} at {} failed for {id}: {error}; rollback also failed: {rollback_error}", source.peer_id, source.addr )) })?; Ok((disposition, error)) } enum StreamInstallFailureAction { Retry { error: StreamInstallReceiveError, invalid_source: bool, }, Exhausted(StreamInstallReceiveError), Stop(StreamDownloadError), } fn stream_install_failure_action( disposition: StreamInstallFailureDisposition, error_kind: StreamInstallReceiveErrorKind, error: StreamInstallReceiveError, total_deadline: StreamInstallTotalDeadline, source: &PeerEndpoint, game_id: &str, ) -> StreamInstallFailureAction { if disposition != StreamInstallFailureDisposition::Stop && total_deadline.is_elapsed() { log::warn!("Streamed install for {game_id} exhausted its shared total receive deadline"); return StreamInstallFailureAction::Exhausted(error); } match disposition { StreamInstallFailureDisposition::RetryAndQuarantine | StreamInstallFailureDisposition::Retry => { log::warn!( "Streamed install attempt from {} at {} failed for {game_id}; trying another peer if available: {error}", source.peer_id, source.addr ); StreamInstallFailureAction::Retry { error, invalid_source: disposition == StreamInstallFailureDisposition::RetryAndQuarantine, } } StreamInstallFailureDisposition::Stop => { let error = eyre::Report::new(error); StreamInstallFailureAction::Stop(match error_kind { StreamInstallReceiveErrorKind::Cancelled => StreamDownloadError::cancelled(error), StreamInstallReceiveErrorKind::Setup => { StreamDownloadError::operation_failed(error) } StreamInstallReceiveErrorKind::Integrity | StreamInstallReceiveErrorKind::Transport => { unreachable!("retryable stream failures cannot stop immediately") } }) } } } fn exhausted_stream_receive_error( id: &str, last_receive_error: Option, ) -> StreamDownloadError { match last_receive_error { Some(error) => StreamDownloadError::sources_exhausted(eyre::Report::new(error)), None => StreamDownloadError::sources_exhausted(eyre::eyre!( "streamed install download failed for {id}: no peer attempts were made" )), } } fn begin_stream_receive_attempt( game_root: &Path, state_dir: &Path, id: &str, retry_invalid_source: &mut bool, attempt: &DownloadAttemptReporter, ) -> Result { let transaction = install::begin_streamed_install(game_root, state_dir, id) .map_err(StreamDownloadError::operation_failed)?; if std::mem::take(retry_invalid_source) { attempt.set_activity(DownloadVerificationActivity::RetryingInvalidSource); } else { attempt.set_activity(DownloadVerificationActivity::VerifyingDownloadedChunks); } Ok(transaction) } async fn receive_streamed_install_from_peers( request: StreamInstallReceiveRequest<'_>, ) -> Result { let StreamInstallReceiveRequest { ctx, tx_notify_ui, target, manifest, sources, quic, cancel_token, attempt, } = request; let id = target.game_id(); let game_root = target.game_root(); let mut last_receive_error = None; let mut retry_invalid_source = false; let content_id = manifest.content_id(); let total_deadline = StreamInstallTotalDeadline::for_manifest(manifest) .map_err(StreamDownloadError::operation_failed)?; let mut source_attempts = StreamInstallSourceAttempts::default(); for source in sources { if let Err(error) = total_deadline.ensure_active(id, cancel_token) { if error.kind() == StreamInstallReceiveErrorKind::Cancelled { return Err(StreamDownloadError::cancelled(eyre::Report::new(error))); } last_receive_error = Some(error); break; } if ctx.content_quarantine.is_quarantined(source, content_id) { log::debug!( "Skipping quarantined streamed-install source {} at {} for {id}", source.peer_id, source.addr ); continue; } match admit_stream_install_source(&mut source_attempts, source, id) { StreamInstallSourceAdmission::Admitted => {} StreamInstallSourceAdmission::AlreadyAttempted => continue, StreamInstallSourceAdmission::Exhausted => break, } let transaction = begin_stream_receive_attempt( &game_root, ctx.state_dir.as_ref(), id, &mut retry_invalid_source, &attempt, )?; let receive_result = receive_streamed_install(ReceiveStreamedInstallRequest { endpoint: *source, game_id: id, manifest: Arc::clone(manifest), staging_dir: transaction.staging_dir(), attempt: attempt.clone(), tx_notify_ui: tx_notify_ui.clone(), quic, cancel_token: cancel_token.clone(), total_deadline, }) .await; match receive_result { Ok(()) => return Ok(transaction), Err(err) => { let error_kind = err.kind(); let (disposition, err) = settle_failed_stream_receive(ctx, source, content_id, id, transaction, err)?; match stream_install_failure_action( disposition, error_kind, err, total_deadline, source, id, ) { StreamInstallFailureAction::Retry { error, invalid_source, } => { last_receive_error = Some(error); retry_invalid_source = invalid_source; } StreamInstallFailureAction::Exhausted(error) => { last_receive_error = Some(error); break; } StreamInstallFailureAction::Stop(error) => return Err(error), } } } } Err(exhausted_stream_receive_error(id, last_receive_error)) } async fn finish_failed_stream_download( ctx: &Ctx, tx_notify_ui: &PeerEventSender, target: &OperationTarget, guard: OperationGuard, status: &DownloadAttemptStatus, direct_reason: Option, ) { status.close_source_admission(); status.clear_activity(); if settle_and_end_download( ctx, tx_notify_ui, target, guard, status.key(), "streamed install failure", ) .await && let Some(reason) = status.resolve_failure(direct_reason) { status.emit_failed(reason); } } struct StreamInstallCommit { ctx: Ctx, tx_notify_ui: PeerEventSender, target: OperationTarget, transaction: install::StreamedInstallTransaction, manifest: Arc, cancel_token: CancellationToken, operation_guard: OperationGuard, download_attempt: DownloadAttemptStatus, } async fn commit_streamed_install(commit: StreamInstallCommit) { let StreamInstallCommit { ctx, tx_notify_ui, target, transaction, manifest, cancel_token, operation_guard, download_attempt, } = commit; let id = target.game_id().to_owned(); match promote_streamed_install( ctx.state_dir.as_ref(), &id, transaction, &manifest, &cancel_token, ) { Ok(()) => { // Promotion is now past its durable rename boundary. Only publish // download success after cancellation can no longer veto it. clear_active_download(&ctx, download_attempt.key()).await; download_attempt.emit_finished(); if settle_and_end_operation( &ctx, &tx_notify_ui, &target, operation_guard, "streamed install completion", ) .await { events::send( &tx_notify_ui, PeerEvent::InstallGameFinished { id: id.clone() }, ); } else { events::send( &tx_notify_ui, PeerEvent::InstallGameFailed { id: id.clone() }, ); } } Err(error) => { let direct_reason = match error { install::StreamedInstallCommitError::Cancelled(error) => { log::info!("Streamed install was cancelled before promotion for {id}: {error}"); None } install::StreamedInstallCommitError::OperationFailed(error) => { log::error!("Streamed install commit failed for {id}: {error}"); Some(DownloadFailureReason::OperationFailed) } }; let settled = settle_and_end_operation( &ctx, &tx_notify_ui, &target, operation_guard, "streamed install commit failure", ) .await; clear_active_download(&ctx, download_attempt.key()).await; if settled && let Some(reason) = download_attempt.resolve_failure(direct_reason) { download_attempt.emit_failed(reason); events::send( &tx_notify_ui, PeerEvent::InstallGameFailed { id: id.clone() }, ); } } } } fn promote_streamed_install( state_dir: &Path, id: &str, transaction: install::StreamedInstallTransaction, manifest: &CatalogContentManifest, cancel_token: &CancellationToken, ) -> Result<(), install::StreamedInstallCommitError> { transaction.commit_classified(manifest, cancel_token)?; if let Err(error) = mark_launch_settings_applied(state_dir, id) { log::warn!( "Streamed install for {id} was promoted, but its launch-settings marker could not be written; first play will retry the rewrite: {error}" ); } Ok(()) } fn spawn_install_operation(ctx: &Ctx, tx_notify_ui: &PeerEventSender, target: OperationTarget) { let ctx = ctx.clone(); let tx_notify_ui = tx_notify_ui.clone(); ctx.task_tracker.clone().spawn(async move { run_install_operation(&ctx, &tx_notify_ui, target).await; }); } async fn run_install_operation(ctx: &Ctx, tx_notify_ui: &PeerEventSender, target: OperationTarget) { let id = target.game_id().to_owned(); let Some(prepared) = prepare_install_operation(ctx, tx_notify_ui, &target).await else { return; }; match begin_operation(ctx, tx_notify_ui, &target, prepared.operation_kind).await { BeginOperationResult::Started => {} BeginOperationResult::AlreadyActive => { log::warn!("Operation for {id} already in progress; ignoring install command"); return; } BeginOperationResult::DrainTimedOut => { log::error!("Timed out waiting for outbound transfers before install/update of {id}"); events::send( tx_notify_ui, PeerEvent::InstallGameFailed { id: id.clone() }, ); return; } BeginOperationResult::PublicationFailed => { log::error!("Failed to withdraw published availability before install/update of {id}"); events::send( tx_notify_ui, PeerEvent::InstallGameFailed { id: id.clone() }, ); return; } BeginOperationResult::RecoveryBlocked => { log::warn!("Ignoring install for recovery-blocked game {id}"); events::send( tx_notify_ui, PeerEvent::InstallGameFailed { id: id.clone() }, ); return; } BeginOperationResult::GameDirChanged => { log::warn!( "Game directory changed while preparing install for {id}; retry the request against the current library" ); events::send( tx_notify_ui, PeerEvent::InstallGameFailed { id: id.clone() }, ); return; } } let cancel_token = ctx.shutdown.child_token(); let operation_guard = OperationGuard::cancellable(id.clone(), cancel_token.clone()); let Some(revalidated) = revalidate_install_operation(ctx, tx_notify_ui, &target, prepared.operation_kind).await else { let _ = settle_and_end_operation( ctx, tx_notify_ui, &target, operation_guard, "install preflight rejection", ) .await; return; }; run_started_install_operation( ctx, tx_notify_ui, target, revalidated, operation_guard, cancel_token, ) .await; } struct PreparedInstallOperation { operation: InstallOperation, operation_kind: OperationKind, } async fn prepare_install_operation( ctx: &Ctx, tx_notify_ui: &PeerEventSender, target: &OperationTarget, ) -> Option { let id = target.game_id(); if !catalog_contains(ctx, id) { log::warn!("Ignoring install command for non-catalog game {id}"); return None; } if ctx.recovery_quarantine.is_blocked(&target.games_folder, id) { log::warn!("Ignoring install command for recovery-blocked game {id}"); events::send( tx_notify_ui, PeerEvent::InstallGameFailed { id: id.to_string() }, ); return None; } if download_ownership_readiness(&target.games_folder, ctx.state_dir.as_ref(), id).await == DownloadOwnershipReadiness::RecoveryRequired { log::warn!("Ignoring install command for recovery-quarantined game {id}"); events::send( tx_notify_ui, PeerEvent::InstallGameFailed { id: id.to_string() }, ); return None; } let game_root = target.game_root(); if !version_ini_is_regular_file(&game_root).await { log::warn!("Ignoring install command for {id}: version.ini sentinel is absent"); events::send( tx_notify_ui, PeerEvent::InstallGameFailed { id: id.to_string() }, ); return None; } let local_present = local_dir_is_directory(&game_root).await; let operation = if local_present { InstallOperation::Updating } else { InstallOperation::Installing }; let operation_kind = match operation { InstallOperation::Installing => OperationKind::Installing, InstallOperation::Updating => OperationKind::Updating, }; Some(PreparedInstallOperation { operation, operation_kind, }) } async fn revalidate_install_operation( ctx: &Ctx, tx_notify_ui: &PeerEventSender, target: &OperationTarget, expected_kind: OperationKind, ) -> Option { let revalidated = prepare_install_operation(ctx, tx_notify_ui, target).await?; if revalidated.operation_kind == expected_kind { return Some(revalidated); } let id = target.game_id(); log::warn!( "Install state changed while admitting {id}; refusing stale {expected_kind:?} operation as {:?}", revalidated.operation_kind ); events::send( tx_notify_ui, PeerEvent::InstallGameFailed { id: id.to_string() }, ); None } async fn run_started_install_operation( ctx: &Ctx, tx_notify_ui: &PeerEventSender, target: OperationTarget, prepared: PreparedInstallOperation, operation_guard: OperationGuard, cancel_token: CancellationToken, ) { let id = target.game_id().to_owned(); let game_root = target.game_root(); let operation = prepared.operation; let result = { let state_dir = ctx.state_dir.as_ref(); match operation { InstallOperation::Installing => { install::install( &game_root, state_dir, &id, ctx.unpacker.clone(), cancel_token.clone(), ) .await } InstallOperation::Updating => { install::update( &game_root, state_dir, &id, ctx.unpacker.clone(), cancel_token.clone(), ) .await } } }; match result { Ok(()) => { if settle_and_end_operation( ctx, tx_notify_ui, &target, operation_guard, "install completion", ) .await { events::send( tx_notify_ui, PeerEvent::InstallGameFinished { id: id.clone() }, ); } else { events::send( tx_notify_ui, PeerEvent::InstallGameFailed { id: id.clone() }, ); } } Err(err) => { log::error!("Install operation failed for {id}: {err}"); let _ = settle_and_end_operation( ctx, tx_notify_ui, &target, operation_guard, "install failure", ) .await; events::send( tx_notify_ui, PeerEvent::InstallGameFailed { id: id.clone() }, ); } } } async fn run_uninstall_operation( ctx: &Ctx, tx_notify_ui: &PeerEventSender, target: OperationTarget, ) { let id = target.game_id().to_owned(); if !catalog_contains(ctx, &id) { log::warn!("Ignoring uninstall command for non-catalog game {id}"); return; } match begin_operation(ctx, tx_notify_ui, &target, OperationKind::Uninstalling).await { BeginOperationResult::Started => {} BeginOperationResult::AlreadyActive => { log::warn!("Operation for {id} already in progress; ignoring uninstall command"); return; } BeginOperationResult::DrainTimedOut => { log::error!("Timed out waiting for outbound transfers before uninstall of {id}"); events::send( tx_notify_ui, PeerEvent::UninstallGameFailed { id: id.clone() }, ); return; } BeginOperationResult::PublicationFailed => { log::error!("Failed to withdraw published availability before uninstall of {id}"); events::send( tx_notify_ui, PeerEvent::UninstallGameFailed { id: id.clone() }, ); return; } BeginOperationResult::RecoveryBlocked => { log::warn!("Ignoring uninstall for recovery-blocked game {id}"); events::send( tx_notify_ui, PeerEvent::UninstallGameFailed { id: id.clone() }, ); return; } BeginOperationResult::GameDirChanged => { log::warn!( "Game directory changed before uninstall admission for {id}; refusing to retarget the command" ); events::send( tx_notify_ui, PeerEvent::UninstallGameFailed { id: id.clone() }, ); return; } } let game_root = target.game_root(); let operation_guard = OperationGuard::new(id.clone()); let result = install::uninstall(&game_root, ctx.state_dir.as_ref(), &id); match result { Ok(()) => { if settle_and_end_operation( ctx, tx_notify_ui, &target, operation_guard, "uninstall completion", ) .await { events::send( tx_notify_ui, PeerEvent::UninstallGameFinished { id: id.clone() }, ); } else { events::send( tx_notify_ui, PeerEvent::UninstallGameFailed { id: id.clone() }, ); } } Err(err) => { log::error!("Uninstall operation failed for {id}: {err}"); let _ = settle_and_end_operation( ctx, tx_notify_ui, &target, operation_guard, "uninstall failure", ) .await; events::send( tx_notify_ui, PeerEvent::UninstallGameFailed { id: id.clone() }, ); } } } async fn run_remove_downloaded_operation( ctx: &Ctx, tx_notify_ui: &PeerEventSender, target: OperationTarget, ) { let id = target.game_id().to_owned(); if !catalog_contains(ctx, &id) { log::warn!("Ignoring downloaded-file removal for non-catalog game {id}"); return; } match begin_operation(ctx, tx_notify_ui, &target, OperationKind::RemovingDownload).await { BeginOperationResult::Started => {} BeginOperationResult::AlreadyActive => { log::warn!("Operation for {id} already in progress; ignoring downloaded-file removal"); return; } BeginOperationResult::DrainTimedOut => { log::error!("Timed out waiting for outbound transfers before removal of {id}"); events::send( tx_notify_ui, PeerEvent::RemoveDownloadedGameFailed { id: id.clone() }, ); return; } BeginOperationResult::PublicationFailed => { log::error!( "Failed to withdraw published availability before downloaded-file removal of {id}" ); events::send( tx_notify_ui, PeerEvent::RemoveDownloadedGameFailed { id: id.clone() }, ); return; } BeginOperationResult::RecoveryBlocked => { log::warn!("Ignoring downloaded-file removal for recovery-blocked game {id}"); events::send( tx_notify_ui, PeerEvent::RemoveDownloadedGameFailed { id: id.clone() }, ); return; } BeginOperationResult::GameDirChanged => { log::warn!( "Game directory changed before downloaded-file removal admission for {id}; refusing to retarget the command" ); events::send( tx_notify_ui, PeerEvent::RemoveDownloadedGameFailed { id: id.clone() }, ); return; } } let operation_guard = OperationGuard::new(id.clone()); let result = remove_downloaded_payload(&target.games_folder, ctx.state_dir.as_ref(), &id).await; match result { Ok(()) => { if settle_and_end_operation( ctx, tx_notify_ui, &target, operation_guard, "downloaded-file removal completion", ) .await { events::send( tx_notify_ui, PeerEvent::RemoveDownloadedGameFinished { id: id.clone() }, ); } else { events::send( tx_notify_ui, PeerEvent::RemoveDownloadedGameFailed { id: id.clone() }, ); } } Err(err) => { log::error!("Downloaded-file removal failed for {id}: {err}"); let _ = settle_and_end_operation( ctx, tx_notify_ui, &target, operation_guard, "downloaded-file removal failure", ) .await; events::send( tx_notify_ui, PeerEvent::RemoveDownloadedGameFailed { id: id.clone() }, ); } } } #[derive(Clone, Copy, Debug, PartialEq, Eq)] enum BeginOperationResult { Started, AlreadyActive, GameDirChanged, DrainTimedOut, RecoveryBlocked, PublicationFailed, } async fn begin_operation( ctx: &Ctx, tx_notify_ui: &PeerEventSender, target: &OperationTarget, operation: OperationKind, ) -> BeginOperationResult { begin_operation_with_drain_timeout( ctx, tx_notify_ui, target, operation, OUTBOUND_TRANSFER_DRAIN_TIMEOUT, ) .await } async fn begin_operation_with_drain_timeout( ctx: &Ctx, tx_notify_ui: &PeerEventSender, target: &OperationTarget, operation: OperationKind, drain_timeout: Duration, ) -> BeginOperationResult { let admission = ctx.operation_admission.lock().await; let game_dir = ctx.game_dir.read().await; if *game_dir != target.games_folder { return BeginOperationResult::GameDirChanged; } if ctx .recovery_quarantine .is_blocked(&target.games_folder, target.game_id()) { return BeginOperationResult::RecoveryBlocked; } drop(game_dir); let started = { let mut active_operations = ctx.active_operations.write().await; if active_operations.contains_key(target.game_id()) { false } else { // Acquire both guards before making the operation visible. A fresh // Hello can therefore observe either the pre-operation projection // or the withdrawn one, never an active operation that is still // advertised as available. let mut library = ctx.local_library.write().await; let withdrawn_revision = match library.withdraw_for_operation(target.game_id()) { Ok(revision) => revision, Err(error) => { log::error!( "Cannot admit {operation:?} for {}: {error}", target.game_id() ); return BeginOperationResult::PublicationFailed; } }; active_operations.insert(target.game_id.clone(), operation); // The network projection is already withdrawn under both guards. // Tell the UI that the game is busy before its local snapshot loses // availability, so it can keep the operation visible in the list. events::send_active_operations_snapshot(tx_notify_ui, &active_operations); if let Some(revision) = withdrawn_revision { let game_db = GameDB::from( library .games .values() .map(game_from_summary) .collect::>(), ); *ctx.local_game_db.write().await = Some(game_db.clone()); events::send( tx_notify_ui, PeerEvent::LocalLibraryChanged { games: game_db.all_games().into_iter().cloned().collect(), }, ); ctx.state_sync.publish_library_revision(revision); } true } }; if !started { return BeginOperationResult::AlreadyActive; } // Once admitted, a directory change observes the active operation and is // rejected. Release the admission barrier before draining transfers. drop(admission); if operation_requires_outbound_drain(operation) && !cancel_and_wait_for_outbound_transfers( ctx, OutboundTransferScope::Game(target.game_id()), drain_timeout, ) .await { if let Err(error) = refresh_local_game_for_ending_operation(ctx, tx_notify_ui, target).await { // The withdrawn projection is safe to retain. Do not restore stale // availability when the authoritative rescan cannot complete. log::error!( "Failed to restore local availability after outbound drain timeout for {}: {error}", target.game_id() ); } end_operation(ctx, tx_notify_ui, target.game_id()).await; return BeginOperationResult::DrainTimedOut; } BeginOperationResult::Started } fn operation_requires_outbound_drain(operation: OperationKind) -> bool { matches!( operation, OperationKind::Downloading | OperationKind::Updating | OperationKind::RemovingDownload ) } #[derive(Clone, Copy)] enum OutboundTransferScope<'a> { Game(&'a str), All, } impl OutboundTransferScope<'_> { fn matches(self, game_id: &str) -> bool { match self { Self::Game(target) => target == game_id, Self::All => true, } } fn description(self) -> String { match self { Self::Game(game_id) => format!("for {game_id}"), Self::All => "across all games".to_string(), } } } fn outbound_transfer_tokens( active: &HashMap>, scope: OutboundTransferScope<'_>, ) -> Vec { active .iter() .filter(|(game_id, _)| scope.matches(game_id)) .flat_map(|(_, transfers)| transfers.iter().map(|(_, token)| token.clone())) .collect() } fn outbound_transfer_count( active: &HashMap>, scope: OutboundTransferScope<'_>, ) -> usize { active .iter() .filter(|(game_id, _)| scope.matches(game_id)) .map(|(_, transfers)| transfers.len()) .sum() } async fn cancel_and_wait_for_outbound_transfers( ctx: &Ctx, scope: OutboundTransferScope<'_>, drain_timeout: Duration, ) -> bool { let tokens_to_cancel = { let active = ctx.active_outbound_transfers.read().await; outbound_transfer_tokens(&active, scope) }; for token in tokens_to_cancel { token.cancel(); } let drained = tokio::time::timeout(drain_timeout, async { loop { let count = { let active = ctx.active_outbound_transfers.read().await; outbound_transfer_count(&active, scope) }; if count == 0 { break; } tokio::time::sleep(OUTBOUND_TRANSFER_DRAIN_POLL_INTERVAL).await; } }) .await .is_ok(); if !drained { let count = { let active = ctx.active_outbound_transfers.read().await; outbound_transfer_count(&active, scope) }; let scope = scope.description(); log::error!( "Timed out after {drain_timeout:?} waiting for {count} outbound transfer(s) to drain {scope}" ); } drained } async fn transition_download_to_install( ctx: &Ctx, tx_notify_ui: &PeerEventSender, id: &str, operation: OperationKind, ) -> bool { let transitioned = { let mut active_operations = ctx.active_operations.write().await; match active_operations.get_mut(id) { Some(current) if *current == OperationKind::Downloading => { *current = operation; true } Some(current) => { log::warn!( "Cannot transition {id} from download to install; current operation is {current:?}" ); false } None => { log::warn!( "Cannot transition {id} from download to install; operation is not active" ); false } } }; if transitioned { events::emit_active_operations(&ctx.active_operations, tx_notify_ui).await; } transitioned } async fn end_operation(ctx: &Ctx, tx_notify_ui: &PeerEventSender, id: &str) { if ctx.active_operations.write().await.remove(id).is_some() { events::emit_active_operations(&ctx.active_operations, tx_notify_ui).await; } } async fn clear_active_download(ctx: &Ctx, attempt: &DownloadAttemptKey) { let mut active_downloads = ctx.active_downloads.write().await; if active_downloads .get(&attempt.id) .is_some_and(|active| active.key() == attempt) { active_downloads.remove(&attempt.id); } } fn send_download_failed( tx_notify_ui: &PeerEventSender, attempt: &DownloadAttemptKey, reason: DownloadFailureReason, ) { let status = DownloadAttemptStatus::new( attempt.clone(), CancellationToken::new(), tx_notify_ui.clone(), ); status.close_source_admission(); status.emit_failed(reason); } async fn end_download_operation( ctx: &Ctx, tx_notify_ui: &PeerEventSender, attempt: &DownloadAttemptKey, ) { clear_active_download(ctx, attempt).await; end_operation(ctx, tx_notify_ui, &attempt.id).await; } async fn settle_target_state( ctx: &Ctx, tx_notify_ui: &PeerEventSender, target: &OperationTarget, ) -> eyre::Result<()> { install::recover_game_root( &target.game_root(), ctx.state_dir.as_ref(), target.game_id(), ) .await?; refresh_local_game_for_ending_operation(ctx, tx_notify_ui, target).await } async fn settle_and_end_operation( ctx: &Ctx, tx_notify_ui: &PeerEventSender, target: &OperationTarget, guard: OperationGuard, label: &str, ) -> bool { if let Err(error) = settle_target_state(ctx, tx_notify_ui, target).await { log::error!( "Failed to settle {label} for {}: {error}; retaining the operation gate until process restart", target.game_id() ); return false; } // Settlement is complete before the guard is disarmed. If this task is // cancelled during the async map cleanup, the still-present operation entry // remains fail-closed; if removal already landed, recovery/publication is // known complete. guard.disarm(); end_operation(ctx, tx_notify_ui, target.game_id()).await; true } async fn settle_and_end_download( ctx: &Ctx, tx_notify_ui: &PeerEventSender, target: &OperationTarget, guard: OperationGuard, attempt: &DownloadAttemptKey, label: &str, ) -> bool { if let Err(error) = settle_target_state(ctx, tx_notify_ui, target).await { log::error!( "Failed to settle {label} for {}: {error}; retaining the operation gate until process restart", target.game_id() ); return false; } guard.disarm(); end_download_operation(ctx, tx_notify_ui, attempt).await; true } async fn finish_successful_download_after_refresh( ctx: &Ctx, tx_notify_ui: &PeerEventSender, target: &OperationTarget, guard: OperationGuard, status: &DownloadAttemptStatus, ) { status.clear_activity(); if settle_and_end_download( ctx, tx_notify_ui, target, guard, status.key(), "download completion", ) .await { status.emit_finished(); } } async fn finish_failed_download_after_refresh( ctx: &Ctx, tx_notify_ui: &PeerEventSender, target: &OperationTarget, guard: OperationGuard, status: &DownloadAttemptStatus, direct_reason: Option, ) { status.close_source_admission(); status.clear_activity(); if settle_and_end_download( ctx, tx_notify_ui, target, guard, status.key(), "download failure", ) .await && let Some(reason) = status.resolve_failure(direct_reason) { status.emit_failed(reason); } } fn catalog_contains(ctx: &Ctx, id: &str) -> bool { ctx.catalog.catalog().contains(id) } async fn begin_local_recovery( ctx: &Ctx, tx_notify_ui: &PeerEventSender, game_dir: &Path, force_empty_snapshot: bool, ) -> eyre::Result<()> { let (had_cached_games, revision) = { let mut library = ctx.local_library.write().await; let had_cached_games = !library.games.is_empty(); let revision = library.clear_for_recovery(game_dir)?; (had_cached_games, revision) }; ctx.recovery_quarantine.begin(game_dir.to_path_buf()); *ctx.local_game_db.write().await = None; if force_empty_snapshot || had_cached_games { events::send( tx_notify_ui, PeerEvent::LocalLibraryChanged { games: Vec::new() }, ); } if let Some(revision) = revision { ctx.state_sync.publish_library_revision(revision); } Ok(()) } /// Handles the `SetGameDir` command. /// /// Recovery, scanning, quarantine settlement, and library-delta delivery /// attempts all complete before this function returns. `Ok` means the returned /// canonical root is bound, even when a recovery failure left one or more games /// quarantined. `Err` is returned only before changing the configured root. pub async fn handle_set_game_dir_command( ctx: &Ctx, tx_notify_ui: &PeerEventSender, game_dir: PathBuf, ) -> Result { handle_set_game_dir_command_with_drain_timeout( ctx, tx_notify_ui, game_dir, OUTBOUND_TRANSFER_DRAIN_TIMEOUT, ) .await } async fn handle_set_game_dir_command_with_drain_timeout( ctx: &Ctx, tx_notify_ui: &PeerEventSender, requested_game_dir: PathBuf, drain_timeout: Duration, ) -> Result { let _admission = ctx.operation_admission.lock().await; let game_dir = match scoped_blocking(|| crate::canonicalize_game_dir(&requested_game_dir)) { Ok(game_dir) => game_dir, Err(error) => { let error = error.to_string(); log::warn!( "Rejecting invalid game directory {}: {error}", requested_game_dir.display() ); return Err(error); } }; if let Err(error) = install::intent::scan_active_install_intents(ctx.state_dir.as_ref(), &game_dir) { let error = format!( "cannot set game directory to {} while install intent state is unresolved: {error}", game_dir.display() ); log::warn!("{error}"); return Err(error); } let current_game_dir = ctx.game_dir.read().await.clone(); let active_ids = active_operation_ids(ctx).await; if !active_ids.is_empty() { let mut active_ids = active_ids.into_iter().collect::>(); active_ids.sort(); let error = format!( "cannot set game directory while operations are active for: {}", active_ids.join(", ") ); log::warn!( "Rejecting game directory refresh/change to {} while operations are active for: {}", game_dir.display(), active_ids.join(", ") ); return Err(error); } if !cancel_and_wait_for_outbound_transfers(ctx, OutboundTransferScope::All, drain_timeout).await { log::error!( "Keeping game directory {} unchanged because outbound transfers did not drain", current_game_dir.display() ); return Err(format!( "outbound transfers did not drain; game directory remains {}", current_game_dir.display() )); } if current_game_dir == game_dir { log::info!( "Game directory {} unchanged; retrying inactive recovery before refresh", game_dir.display() ); if let Err(error) = begin_local_recovery(ctx, tx_notify_ui, &game_dir, true).await { log::error!("Failed to begin local recovery: {error}"); return Err(format!("failed to begin local recovery: {error}")); } match load_local_library(ctx, tx_notify_ui).await { Ok(()) => log::info!("Local game database refreshed successfully"), Err(error) => log::error!( "Game directory {} remains accepted but its library refresh failed: {error}", game_dir.display() ), } return Ok(game_dir); } // Begin the target recovery epoch before changing the configured root. // Requests that captured the old root now observe a root mismatch, while // requests that later capture the new root observe Recovering. if let Err(error) = begin_local_recovery(ctx, tx_notify_ui, &game_dir, true).await { log::error!( "Failed to begin recovery for {}: {error}", game_dir.display() ); return Err(format!( "failed to begin recovery for {}: {error}", game_dir.display() )); } *ctx.game_dir.write().await = game_dir.clone(); log::info!("Game directory set to: {}", game_dir.display()); match load_local_library_with_policy(ctx, tx_notify_ui, LocalLibraryEventPolicy::ForceSnapshot) .await { Ok(()) => log::info!("Local game database loaded successfully"), Err(error) => log::error!( "Game directory {} remains accepted but its library load failed: {error}", game_dir.display() ), } Ok(game_dir) } /// Loads the configured local library and announces the result. pub async fn load_local_library(ctx: &Ctx, tx_notify_ui: &PeerEventSender) -> eyre::Result<()> { load_local_library_with_policy(ctx, tx_notify_ui, LocalLibraryEventPolicy::OnChange).await } async fn load_local_library_with_policy( ctx: &Ctx, tx_notify_ui: &PeerEventSender, event_policy: LocalLibraryEventPolicy, ) -> eyre::Result<()> { let game_dir = { ctx.game_dir.read().await.clone() }; begin_local_recovery(ctx, tx_notify_ui, &game_dir, false).await?; let active_operations = ctx.active_operations.read().await; let active_ids = active_operations.keys().cloned().collect(); let recovery_report = install::recover_on_startup(&game_dir, ctx.state_dir.as_ref(), &active_ids).await?; drop(active_operations); for (id, error) in recovery_report.failures() { log::error!("Keeping game {id} quarantined after recovery failure: {error}"); } let recovery_error = recovery_report.summary_error(); let failed_ids = recovery_report.failed_ids(); scan_and_announce_local_library(ctx, tx_notify_ui, &game_dir, event_policy, &failed_ids) .await?; if !ctx.recovery_quarantine.settle(&game_dir, failed_ids) { eyre::bail!( "local recovery result for {} was superseded before settlement", game_dir.display() ); } match recovery_error { Some(error) => Err(error), None => Ok(()), } } async fn scan_and_announce_local_library( ctx: &Ctx, tx_notify_ui: &PeerEventSender, game_dir: &Path, event_policy: LocalLibraryEventPolicy, recovery_failed_ids: &HashSet, ) -> eyre::Result<()> { let catalog = ctx.catalog.catalog(); let scan = scan_local_library_with_recovery_failures( game_dir, ctx.state_dir.as_ref(), catalog, recovery_failed_ids, ) .await?; match update_and_announce_games_with_policy(ctx, tx_notify_ui, scan, event_policy, None).await { LocalLibraryPublication::Published => Ok(()), LocalLibraryPublication::Rejected => { eyre::bail!("local library scan became obsolete before publication") } } } /// Refreshes the game whose operation has completed before clearing its /// active-operation snapshot, while preserving freeze behavior for other games. async fn refresh_local_game_for_ending_operation( ctx: &Ctx, tx_notify_ui: &PeerEventSender, target: &OperationTarget, ) -> eyre::Result<()> { let catalog = ctx.catalog.catalog(); let failed_ids = ctx.recovery_quarantine.failed_ids(&target.games_folder); let scan = rescan_local_game_with_recovery_failures( &target.games_folder, ctx.state_dir.as_ref(), catalog, target.game_id(), &failed_ids, ) .await?; match update_and_announce_games_with_policy( ctx, tx_notify_ui, scan, LocalLibraryEventPolicy::OnChange, Some(target.game_id()), ) .await { LocalLibraryPublication::Published => Ok(()), LocalLibraryPublication::Rejected => { eyre::bail!("local game refresh became obsolete before publication") } } } #[derive(Clone, Copy, Debug, PartialEq, Eq)] enum LocalLibraryEventPolicy { OnChange, ForceSnapshot, } #[derive(Clone, Copy, Debug, PartialEq, Eq)] enum LocalLibraryPublication { Published, Rejected, } async fn active_operation_ids(ctx: &Ctx) -> HashSet { ctx.active_operations.read().await.keys().cloned().collect() } /// Handles the `GetPeerCount` command. pub async fn handle_get_peer_count_command(ctx: &Ctx, tx_notify_ui: &PeerEventSender) { log::info!("GetPeerCount command received"); events::emit_peer_count(&ctx.peer_game_db, tx_notify_ui).await; } /// Connects to a peer directly, bypassing mDNS discovery. pub async fn handle_connect_peer_command( ctx: &Ctx, tx_notify_ui: &PeerEventSender, endpoint: PeerEndpoint, ) { log::info!("Direct connect command received for {}", endpoint.addr); let network_permit = match ctx.network.try_acquire() { Ok(permit) => permit, Err(error) => { log::warn!( "Cannot directly connect to {} while Local network sharing is unavailable: {error}", endpoint.addr ); return; } }; let network_ctx = network_permit.service_context(ctx.clone()); let handshake_ctx = HandshakeCtx::from_network(&network_ctx, tx_notify_ui); let handshake = match ReservedCandidateHandshake::reserve(handshake_ctx, endpoint).await { Ok(handshake) => handshake, Err(err) => { log::warn!( "Failed to reserve direct connect to {}: {err}", endpoint.addr ); return; } }; ctx.task_tracker .spawn(retain_guard_until_operation_completes( network_permit, async move { if let Err(err) = handshake.run().await { log::warn!("Failed direct connect to {}: {err}", endpoint.addr); } }, )); } // ============================================================================= // Game announcement helpers // ============================================================================= /// Updates the local game database and announces changes to peers. pub async fn update_and_announce_games( ctx: &Ctx, tx_notify_ui: &PeerEventSender, scan: LocalLibraryScan, ) { let _ = update_and_announce_games_with_policy( ctx, tx_notify_ui, scan, LocalLibraryEventPolicy::OnChange, None, ) .await; } async fn update_and_announce_games_with_policy( ctx: &Ctx, tx_notify_ui: &PeerEventSender, scan: LocalLibraryScan, event_policy: LocalLibraryEventPolicy, ending_operation_id: Option<&str>, ) -> LocalLibraryPublication { let LocalLibraryScan { source_game_dir, mut game_db, mut summaries, revision, } = scan; // Hold the configured-root read guard through publication. A directory // switch either happens first (and this scan is rejected) or waits until // both cached library representations have been updated. let current_game_dir = ctx.game_dir.read().await; if *current_game_dir != source_game_dir { log::debug!( "Discarding local library scan from {} because the configured directory is now {}", source_game_dir.display(), current_game_dir.display() ); return LocalLibraryPublication::Rejected; } let mut active_operation_ids = active_operation_ids(ctx).await; if let Some(id) = ending_operation_id { active_operation_ids.remove(id); } if !active_operation_ids.is_empty() { for id in &active_operation_ids { summaries.remove(id); } game_db = GameDB::from(summaries.values().map(game_from_summary).collect()); } // Resolve every manifest needed by the prospective wire publication before // making its revision visible. Hello/Pong responders are cache-only, so a // malformed or unreadable artifact must reject the whole scan without // partially mutating either local cache or emitting UI/state-sync updates. let eligible = crate::library::catalog_eligible_game_ids(&summaries, ctx.catalog.catalog()); let catalog = Arc::clone(&ctx.catalog); let summaries = match scoped_blocking(move || { crate::library::prime_library_manifests(&eligible, &catalog)?; Ok::<_, eyre::Report>(summaries) }) { Ok(summaries) => summaries, Err(error) => { log::error!("Rejecting local library publication: {error}"); return LocalLibraryPublication::Rejected; } }; let mut library_guard = ctx.local_library.write().await; if let Some(published_revision) = library_guard.source_revision_for(&source_game_dir) && revision < published_revision { log::debug!( "Discarding local library scan at source revision {revision} because source revision {published_revision} is already published" ); return LocalLibraryPublication::Rejected; } let published_revision = match library_guard.update_from_scan(&source_game_dir, summaries, revision) { Ok(published_revision) => published_revision, Err(error) => { log::error!("Rejecting local library publication: {error}"); return LocalLibraryPublication::Rejected; } }; { let mut db_guard = ctx.local_game_db.write().await; *db_guard = Some(game_db.clone()); } drop(library_guard); drop(current_game_dir); let all_games = game_db.all_games().into_iter().cloned().collect::>(); if published_revision.is_some() || event_policy == LocalLibraryEventPolicy::ForceSnapshot { events::send( tx_notify_ui, PeerEvent::LocalLibraryChanged { games: all_games.clone(), }, ); } else { log::debug!("Skipping unchanged local library event"); } let Some(published_revision) = published_revision else { return LocalLibraryPublication::Published; }; ctx.state_sync.publish_library_revision(published_revision); LocalLibraryPublication::Published } #[cfg(test)] mod tests { use std::{ collections::{BTreeMap, HashMap}, net::SocketAddr, path::{Path, PathBuf}, sync::{ Arc, atomic::{AtomicBool, Ordering}, }, time::{Duration, SystemTime, UNIX_EPOCH}, }; use lanspread_db::{ content_manifest::{ Blake3Digest, CATALOG_CONTENT_INDEX_NAME, CatalogContentIdentity, CatalogContentIndex, CatalogContentIndexEntry, CatalogContentManifestBody, CatalogExtractedEntry, CatalogFileEntry, ContentId, write_canonical_content_index_atomic, }, db::Availability, }; use lanspread_proto::{ GameAvailability, LibrarySnapshot, PeerEndpoint, PeerId, RuntimeSessionId, }; use tokio::sync::{RwLock, oneshot}; use tokio_util::{sync::CancellationToken, task::TaskTracker}; use super::*; use crate::{ ActiveOperation, ActiveOperationKind, CallToPlayLocalAction, CallToPlayLocalIntent, UnpackFuture, Unpacker, download::seed_pending_download_ownership_for_test, identity::PeerIdentity, install::intent::{InstallIntent, InstallIntentState, intent_path, write_intent}, network_generation::NetworkControl, test_support::{TempDir, catalog_bundle}, }; struct FakeUnpacker; struct DropFlag(Arc); impl Drop for DropFlag { fn drop(&mut self) { self.0.store(true, Ordering::SeqCst); } } impl Unpacker for FakeUnpacker { fn unpack<'a>( &'a self, _archive: &'a Path, dest: &'a Path, _cancel_token: CancellationToken, ) -> UnpackFuture<'a> { Box::pin(async move { tokio::fs::write(dest.join("payload.txt"), b"installed").await?; Ok(()) }) } } fn write_file(path: &Path, bytes: &[u8]) { if let Some(parent) = path.parent() { std::fs::create_dir_all(parent).expect("parent dir should be created"); } std::fs::write(path, bytes).expect("file should be written"); } fn operation_target(game_dir: &Path) -> OperationTarget { OperationTarget::new(game_dir.to_path_buf(), "game".to_string()) } fn streamed_manifest(entries: Vec) -> CatalogContentManifest { let version = b"20250101"; let version_digest = Blake3Digest::hash(version); CatalogContentManifest::seal( CatalogContentManifestBody::new( "game", "20250101", vec![ CatalogFileEntry::file( "version.ini", u64::try_from(version.len()).expect("test version length should fit u64"), version_digest, vec![version_digest], ) .expect("test version entry should validate"), ], entries, ) .expect("test streamed manifest body should validate"), ) .expect("test streamed manifest should seal") } fn downloadable_manifest() -> CatalogContentManifest { let version = b"20250101"; let archive = b"archive"; let version_digest = Blake3Digest::hash(version); let archive_digest = Blake3Digest::hash(archive); CatalogContentManifest::seal( CatalogContentManifestBody::new( "game", "20250101", vec![ CatalogFileEntry::file( "game.eti", u64::try_from(archive.len()).expect("test archive length should fit u64"), archive_digest, vec![archive_digest], ) .expect("test archive entry should validate"), CatalogFileEntry::file( "version.ini", u64::try_from(version.len()).expect("test version length should fit u64"), version_digest, vec![version_digest], ) .expect("test version entry should validate"), ], Vec::new(), ) .expect("test download manifest body should validate"), ) .expect("test download manifest should seal") } async fn seed_exact_download_ownership(ctx: &Ctx, games_folder: &Path, content_id: ContentId) { crate::download::seed_download_ownership_for_test( ctx.state_dir.as_ref(), games_folder, "game", &["game.eti"], ) .await; let games_folder_key = crate::state_paths::games_folder_key(games_folder); let record_path = crate::state_paths::download_ownership_path( ctx.state_dir.as_ref(), "game", &games_folder_key, ); let mut record: serde_json::Value = serde_json::from_slice( &std::fs::read(&record_path).expect("seeded ownership should be readable"), ) .expect("seeded ownership should be valid JSON"); record["committed_content_id"] = serde_json::Value::String(content_id.to_string()); std::fs::write( record_path, serde_json::to_vec(&record).expect("updated ownership should serialize"), ) .expect("updated ownership should be written"); } fn test_ctx(game_dir: PathBuf) -> Ctx { test_ctx_with_catalog(game_dir, catalog_bundle([("game", "20250101")])) } fn test_ctx_with_catalog( game_dir: PathBuf, catalog: Arc, ) -> Ctx { test_ctx_with_catalog_and_network(game_dir, catalog, NetworkControl::enabled_for_test()) } fn test_ctx_with_catalog_and_network( game_dir: PathBuf, catalog: Arc, network: NetworkControl, ) -> Ctx { let state_dir = game_dir.join(".test-state"); let recovery_root = game_dir.clone(); let ctx = Ctx::new( Arc::new(RwLock::new(PeerGameDB::new())), Arc::new(PeerIdentity::generate().expect("test identity should generate")), game_dir, state_dir, Arc::new(FakeUnpacker), CancellationToken::new(), TaskTracker::new(), catalog, Arc::new(RwLock::new(HashMap::new())), Arc::new(crate::NoopStreamInstallProvider), network, ) .expect("test context should initialize"); assert!( ctx.recovery_quarantine .settle(&recovery_root, HashSet::new()) ); ctx } async fn register_test_download( ctx: &Ctx, tx_notify_ui: &PeerEventSender, ) -> (DownloadAttemptStatus, CancellationToken) { let cancellation = CancellationToken::new(); let status = register_download_attempt( ctx, tx_notify_ui, DownloadAttemptKey::next("game".to_owned()), cancellation.clone(), ) .await; (status, cancellation) } fn create_call_to_play_intent() -> CallToPlayLocalIntent { let now = i64::try_from( SystemTime::now() .duration_since(UNIX_EPOCH) .expect("system time should follow the Unix epoch") .as_millis(), ) .expect("current Unix time should fit i64"); CallToPlayLocalIntent { call_id: None, action: CallToPlayLocalAction::Create { game_id: "game".to_owned(), max_players: 4, scheduled_for: None, deadline: now + 60_000, }, } } #[test] fn cancelled_download_owner_does_not_emit_failed_event() { let (tx, mut rx) = crate::peer_event_channel(); let status = DownloadAttemptStatus::new( DownloadAttemptKey::next("game".to_owned()), CancellationToken::new(), tx, ); assert!(status.signal().cancel_silently()); status.close_source_admission(); status.clear_activity(); assert_eq!(status.resolve_failure(None), None); assert!(rx.try_recv().is_err()); } #[test] fn uncancelled_download_error_emits_failed_event() { let (tx, mut rx) = crate::peer_event_channel(); let attempt = DownloadAttemptKey::next("game".to_owned()); send_download_failed(&tx, &attempt, DownloadFailureReason::OperationFailed); assert!(matches!( rx.try_recv(), Ok(PeerEvent::DownloadGameFilesFailed { attempt: emitted_attempt, reason: DownloadFailureReason::OperationFailed, }) if emitted_attempt == attempt )); } #[tokio::test] async fn retained_guard_outlives_cancelled_operation_cleanup() { let dropped = Arc::new(AtomicBool::new(false)); let cancel_token = CancellationToken::new(); let task_cancel_token = cancel_token.clone(); let (started_tx, started_rx) = oneshot::channel(); let (cleanup_tx, cleanup_rx) = oneshot::channel(); let (release_tx, release_rx) = oneshot::channel(); let task = tokio::spawn(retain_guard_until_operation_completes( DropFlag(Arc::clone(&dropped)), async move { started_tx .send(()) .expect("operation-start signal should be observed"); task_cancel_token.cancelled().await; cleanup_tx .send(()) .expect("cleanup-start signal should be observed"); release_rx .await .expect("operation cleanup should be released"); }, )); started_rx .await .expect("retained operation should report startup"); cancel_token.cancel(); cleanup_rx .await .expect("retained operation should enter cleanup"); assert!(!dropped.load(Ordering::SeqCst)); assert!(!task.is_finished()); release_tx .send(()) .expect("cleanup release should reach retained operation"); task.await.expect("retained operation should not panic"); assert!(dropped.load(Ordering::SeqCst)); } #[tokio::test] async fn disabled_network_uses_exact_cached_download_without_peer_db_access() { let games = TempDir::new("lanspread-handler-disabled-cached-download"); write_file(&games.game_root().join("version.ini"), b"20250101"); write_file(&games.game_root().join("game.eti"), b"archive"); let manifest = downloadable_manifest(); let content_id = manifest.content_id(); let catalog = Arc::new( lanspread_db::content_manifest::CatalogBundle::from_manifests([manifest]) .expect("test catalog should validate"), ); let ctx = test_ctx_with_catalog_and_network( games.path().to_path_buf(), catalog, NetworkControl::disabled_for_test(), ); seed_exact_download_ownership(&ctx, games.path(), content_id).await; let peer_db = ctx.peer_game_db.write().await; let (tx, mut rx) = crate::peer_event_channel(); tokio::time::timeout( Duration::from_secs(1), handle_download_game_files_command(&ctx, &tx, "game".to_owned(), false), ) .await .expect("cached download must not wait for PeerDB while sharing is disabled"); drop(peer_db); let begin_attempt = match recv_event(&mut rx).await { PeerEvent::DownloadGameFilesBegin { attempt } => attempt, event => panic!("expected download begin, got {event:?}"), }; let finished_attempt = match recv_event(&mut rx).await { PeerEvent::DownloadGameFilesFinished { attempt } => attempt, event => panic!("expected download finish, got {event:?}"), }; assert_eq!(begin_attempt, finished_attempt); assert_eq!(begin_attempt.id, "game"); assert!(ctx.active_operations.read().await.is_empty()); assert!(ctx.active_downloads.read().await.is_empty()); assert_no_event(&mut rx).await; } #[tokio::test] async fn disabled_network_rejects_nonlocal_download_before_peer_db_or_mutation() { let games = TempDir::new("lanspread-handler-disabled-download"); let ctx = test_ctx_with_catalog_and_network( games.path().to_path_buf(), catalog_bundle([("game", "20250101")]), NetworkControl::disabled_for_test(), ); let peer_db = ctx.peer_game_db.write().await; let (tx, mut rx) = crate::peer_event_channel(); tokio::time::timeout( Duration::from_secs(1), handle_download_game_files_command(&ctx, &tx, "game".to_owned(), false), ) .await .expect("disabled download must not wait for PeerDB"); drop(peer_db); assert!(matches!( recv_event(&mut rx).await, PeerEvent::DownloadGameFilesFailed { attempt, reason: DownloadFailureReason::OperationFailed, } if attempt.id == "game" )); assert!(!games.game_root().exists()); assert!(ctx.active_operations.read().await.is_empty()); assert!(ctx.active_downloads.read().await.is_empty()); assert_no_event(&mut rx).await; } #[tokio::test] async fn disabled_network_rejects_stream_install_before_peer_db_or_staging() { let games = TempDir::new("lanspread-handler-disabled-stream-install"); let manifest = streamed_manifest(vec![ CatalogExtractedEntry::file("payload.txt", 7, Blake3Digest::hash(b"payload")) .expect("test streamed entry should validate"), ]); let catalog = Arc::new( lanspread_db::content_manifest::CatalogBundle::from_manifests([manifest]) .expect("test catalog should validate"), ); let ctx = test_ctx_with_catalog_and_network( games.path().to_path_buf(), catalog, NetworkControl::disabled_for_test(), ); let peer_db = ctx.peer_game_db.write().await; let (tx, mut rx) = crate::peer_event_channel(); tokio::time::timeout( Duration::from_secs(1), handle_stream_install_game_command( &ctx, &tx, "game".to_owned(), StreamInstallSettings::default(), ), ) .await .expect("disabled streamed install must not wait for PeerDB"); drop(peer_db); assert!(matches!( recv_event(&mut rx).await, PeerEvent::DownloadGameFilesFailed { attempt, reason: DownloadFailureReason::OperationFailed, } if attempt.id == "game" )); assert!(!games.game_root().exists()); assert!(ctx.active_operations.read().await.is_empty()); assert!(ctx.active_downloads.read().await.is_empty()); assert_no_event(&mut rx).await; } #[tokio::test] async fn disabled_network_rejects_direct_connect_before_ticket_reservation() { let games = TempDir::new("lanspread-handler-disabled-direct-connect"); let ctx = test_ctx_with_catalog_and_network( games.path().to_path_buf(), catalog_bundle([("game", "20250101")]), NetworkControl::disabled_for_test(), ); let peer_db = ctx.peer_game_db.write().await; let (tx, mut rx) = crate::peer_event_channel(); tokio::time::timeout( Duration::from_secs(1), handle_connect_peer_command(&ctx, &tx, endpoint(7, 12_007)), ) .await .expect("disabled direct connect must not wait for PeerDB"); drop(peer_db); assert_eq!( ctx.peer_game_db.read().await.negotiation_claim_counts(), (0, 0) ); assert_no_event(&mut rx).await; } #[tokio::test] async fn disabled_network_rejects_call_to_play_before_store_mutation() { let games = TempDir::new("lanspread-handler-disabled-call-to-play"); let ctx = test_ctx_with_catalog_and_network( games.path().to_path_buf(), catalog_bundle([("game", "20250101")]), NetworkControl::disabled_for_test(), ); let mut store = ctx.call_to_play.write().await; let before = store .current_publication() .expect("initial Call-to-Play projection should load") .view; let (tx, mut rx) = crate::peer_event_channel(); let (reply_tx, reply_rx) = oneshot::channel(); tokio::time::timeout( Duration::from_secs(1), crate::handle_apply_call_to_play_intent( &ctx, &tx, create_call_to_play_intent(), "Player".to_owned(), reply_tx, ), ) .await .expect("closed admission must reject without waiting for the held store lock"); assert!( reply_rx .await .expect("handler should reply") .expect_err("disabled admission must reject") .contains("disabled or changing state") ); let after = store .current_publication() .expect("unchanged Call-to-Play projection should load") .view; assert_eq!(after, before); drop(store); assert_no_event(&mut rx).await; } #[tokio::test] async fn enabled_call_to_play_admission_waits_then_mutates_once() { let games = TempDir::new("lanspread-handler-enabled-call-to-play"); let ctx = test_ctx(games.path().to_path_buf()); let store = ctx.call_to_play.write().await; let (tx, mut rx) = crate::peer_event_channel(); let (reply_tx, mut reply_rx) = oneshot::channel(); let task = tokio::spawn({ let ctx = ctx.clone(); let tx = tx.clone(); async move { crate::handle_apply_call_to_play_intent( &ctx, &tx, create_call_to_play_intent(), "Player".to_owned(), reply_tx, ) .await; } }); tokio::task::yield_now().await; assert!(matches!( reply_rx.try_recv(), Err(tokio::sync::oneshot::error::TryRecvError::Empty) )); drop(store); task.await .expect("admitted Call-to-Play task should finish"); let receipt = reply_rx .await .expect("handler should reply") .expect("enabled admission should accept the intent"); assert_eq!(receipt.revision, 1); assert!(matches!( recv_event(&mut rx).await, PeerEvent::CallToPlayView(view) if view.events.len() == 1 )); assert_no_event(&mut rx).await; } #[tokio::test] async fn streamed_install_without_extracted_catalog_manifest_fails_before_admission() { let games = TempDir::new("lanspread-handler-stream-capability"); let ctx = test_ctx(games.path().to_path_buf()); let (tx, mut rx) = crate::peer_event_channel(); handle_stream_install_game_command( &ctx, &tx, "game".to_string(), StreamInstallSettings::default(), ) .await; assert!(matches!( recv_event(&mut rx).await, PeerEvent::DownloadGameFilesFailed { attempt, reason: DownloadFailureReason::OperationFailed, } if attempt.id == "game" )); assert!(ctx.active_operations.read().await.is_empty()); assert!(ctx.active_downloads.read().await.is_empty()); assert_no_event(&mut rx).await; } #[tokio::test] async fn version_only_local_download_never_shortcuts_exact_catalog_ownership() { let games = TempDir::new("lanspread-handler-local-content-ownership"); write_file(&games.game_root().join("version.ini"), b"20250101"); write_file(&games.game_root().join("game.eti"), b"archive"); let ctx = test_ctx(games.path().to_path_buf()); crate::download::seed_download_ownership_for_test( ctx.state_dir.as_ref(), games.path(), "game", &["game.eti"], ) .await; let catalog_content_id = ctx .catalog .manifest("game") .expect("test manifest should load") .content_id(); assert!( !download_ownership_matches_content( games.path(), ctx.state_dir.as_ref(), "game", catalog_content_id, ) .await ); let (tx, mut rx) = crate::peer_event_channel(); handle_download_game_files_command(&ctx, &tx, "game".to_string(), false).await; assert!(matches!( recv_event(&mut rx).await, PeerEvent::DownloadGameFilesFailed { attempt, reason: DownloadFailureReason::VerifiedCatalogSourcesExhausted, } if attempt.id == "game" )); assert!(ctx.active_operations.read().await.is_empty()); assert!(ctx.active_downloads.read().await.is_empty()); assert_no_event(&mut rx).await; } #[tokio::test] async fn already_active_download_emits_no_status_for_the_new_attempt() { let games = TempDir::new("lanspread-handler-download-already-active"); let manifest = downloadable_manifest(); let content_id = manifest.content_id(); let catalog = Arc::new( lanspread_db::content_manifest::CatalogBundle::from_manifests([manifest]) .expect("test catalog should validate"), ); let ctx = test_ctx_with_catalog(games.path().to_path_buf(), catalog); { let mut peer_game_db = ctx.peer_game_db.write().await; upsert( &mut peer_game_db, endpoint(1, 12_001), Some(("game", content_id)), ); } ctx.active_operations .write() .await .insert("game".to_owned(), OperationKind::Downloading); let (tx, mut rx) = crate::peer_event_channel(); handle_download_game_files_command(&ctx, &tx, "game".to_owned(), false).await; assert_eq!( ctx.active_operations.read().await.get("game"), Some(&OperationKind::Downloading) ); assert!(ctx.active_downloads.read().await.is_empty()); assert_no_event(&mut rx).await; } #[test] fn streamed_install_marks_settings_only_after_successful_promotion() { let games = TempDir::new("lanspread-handler-stream-promotion"); let state = TempDir::new("lanspread-handler-stream-promotion-state"); let marker = crate::state_paths::launch_settings_applied_path(state.path(), "game"); let transaction = install::begin_streamed_install(&games.game_root(), state.path(), "game") .expect("streamed install transaction should begin"); write_file(&transaction.staging_dir().join("payload.txt"), b"installed"); assert!(!marker.exists()); let manifest = streamed_manifest(vec![ CatalogExtractedEntry::file("payload.txt", 9, Blake3Digest::hash(b"installed")) .expect("test streamed entry should validate"), ]); promote_streamed_install( state.path(), "game", transaction, &manifest, &CancellationToken::new(), ) .expect("verified staging should be promoted"); assert_eq!( std::fs::read(games.game_root().join("local/payload.txt")) .expect("promoted payload should be readable"), b"installed" ); assert!(marker.is_file()); } #[tokio::test] async fn recovery_required_completion_settles_before_success() { let games = TempDir::new("lanspread-handler-download-recovery"); let root = games.game_root(); write_file(&root.join("version.ini"), b"20250101"); write_file(&root.join("archive.eti"), b"archive"); let ctx = test_ctx(games.path().to_path_buf()); seed_pending_download_ownership_for_test( ctx.state_dir.as_ref(), games.path(), "game", &["archive.eti"], &["archive.eti"], ) .await; settle_download_completion( &ctx, &operation_target(games.path()), DownloadCompletion::RecoveryRequired(eyre::eyre!("injected uncertainty")), ) .await .expect("committed download recovery should settle"); assert_eq!( download_ownership_readiness(games.path(), ctx.state_dir.as_ref(), "game").await, DownloadOwnershipReadiness::Settled ); } #[tokio::test] async fn failed_completion_recovery_remains_quarantined() { let games = TempDir::new("lanspread-handler-download-recovery-failure"); let root = games.game_root(); write_file(&root.join("version.ini"), b"20250101"); write_file(&root.join("archive.eti"), b"archive"); std::fs::create_dir_all(root.join(".version.ini.tmp")) .expect("invalid scratch directory should be created"); let ctx = test_ctx(games.path().to_path_buf()); seed_pending_download_ownership_for_test( ctx.state_dir.as_ref(), games.path(), "game", &["archive.eti"], &["archive.eti"], ) .await; let error = settle_download_completion( &ctx, &operation_target(games.path()), DownloadCompletion::RecoveryRequired(eyre::eyre!("injected uncertainty")), ) .await .expect_err("invalid scratch shape should keep recovery quarantined"); assert!( error .to_string() .contains("download recovery did not settle") ); assert_eq!( download_ownership_readiness(games.path(), ctx.state_dir.as_ref(), "game").await, DownloadOwnershipReadiness::RecoveryRequired ); let scan = scan_local_library(games.path(), ctx.state_dir.as_ref(), ctx.catalog.catalog()) .await .expect("quarantined library should still scan"); let game = scan .summaries .get("game") .expect("quarantined game should remain visible"); assert!(!game.downloaded); assert_eq!(game.availability, Availability::LocalOnly); } #[tokio::test] async fn failed_root_stays_quarantined_while_healthy_root_loads_and_retry_clears() { let games = TempDir::new("lanspread-handler-runtime-recovery"); let broken = games.path().join("broken"); let healthy = games.path().join("healthy"); for root in [&broken, &healthy] { write_file(&root.join("version.ini"), b"20250101"); write_file(&root.join("archive.eti"), b"archive"); write_file(&root.join("local/payload.txt"), b"installed"); } std::fs::create_dir_all(broken.join(".version.ini.tmp")) .expect("invalid scratch directory should be created"); write_file(&healthy.join(".version.ini.tmp"), b"stale scratch"); let ctx = test_ctx_with_catalog( games.path().to_path_buf(), catalog_bundle([("broken", "20250101"), ("healthy", "20250101")]), ); let (tx, _rx) = crate::peer_event_channel(); let error = load_local_library(&ctx, &tx) .await .expect_err("per-game recovery failure should remain visible to the caller"); assert!(error.to_string().contains("broken")); assert!(ctx.recovery_quarantine.is_blocked(games.path(), "broken")); assert!(!ctx.recovery_quarantine.is_blocked(games.path(), "healthy")); let library = ctx.local_library.read().await; assert!(!library.games["broken"].downloaded); assert!(!library.games["broken"].installed); assert!(library.games["healthy"].downloaded); assert!(library.games["healthy"].installed); drop(library); assert!(!healthy.join(".version.ini.tmp").exists()); std::fs::remove_dir_all(broken.join(".version.ini.tmp")) .expect("invalid scratch directory should be removable"); load_local_library(&ctx, &tx) .await .expect("corrected recovery should settle"); assert!(!ctx.recovery_quarantine.is_blocked(games.path(), "broken")); let library = ctx.local_library.read().await; assert!(library.games["broken"].downloaded); assert!(library.games["broken"].installed); } #[cfg(unix)] #[tokio::test] async fn startup_recovery_quarantines_symlink_game_root_without_scanning_outside() { use std::os::unix::fs::symlink; let games = TempDir::new("lanspread-handler-symlink-recovery-games"); let outside = TempDir::new("lanspread-handler-symlink-recovery-outside"); write_file(&outside.path().join("version.ini"), b"20250101"); write_file(&outside.path().join("archive.eti"), b"archive"); write_file(&outside.path().join("local/payload.txt"), b"installed"); write_file(&outside.path().join("canary.txt"), b"outside"); symlink(outside.path(), games.path().join("game")) .expect("game-root symlink should be created"); let ctx = test_ctx(games.path().to_path_buf()); let (tx, _rx) = crate::peer_event_channel(); let error = load_local_library(&ctx, &tx) .await .expect_err("unsafe game root should remain a visible recovery failure"); assert!(error.to_string().contains("game")); assert!(ctx.recovery_quarantine.is_blocked(games.path(), "game")); assert!(!ctx.local_library.read().await.games.contains_key("game")); assert_eq!( std::fs::read(outside.path().join("canary.txt")) .expect("canary should remain readable"), b"outside" ); assert_eq!( std::fs::read(outside.path().join("local/payload.txt")) .expect("outside install should remain readable"), b"installed" ); } #[tokio::test] async fn operation_admission_rejects_recovering_and_failed_games() { let games = TempDir::new("lanspread-handler-recovery-admission"); let ctx = test_ctx(games.path().to_path_buf()); let (tx, _rx) = crate::peer_event_channel(); ctx.recovery_quarantine.begin(games.path().to_path_buf()); let target = operation_target(games.path()); assert_eq!( begin_operation(&ctx, &tx, &target, OperationKind::Downloading).await, BeginOperationResult::RecoveryBlocked ); assert!(ctx.active_operations.read().await.is_empty()); assert!( ctx.recovery_quarantine .settle(games.path(), HashSet::from(["game".to_string()])) ); assert_eq!( begin_operation(&ctx, &tx, &target, OperationKind::Downloading).await, BeginOperationResult::RecoveryBlocked ); assert!(ctx.active_operations.read().await.is_empty()); } #[tokio::test] async fn final_download_refresh_failure_keeps_fail_closed_operation_gate() { let games = TempDir::new("lanspread-handler-download-refresh-failure"); write_file(&games.game_root(), b"not a directory"); let ctx = test_ctx(games.path().to_path_buf()); let (tx, mut rx) = crate::peer_event_channel(); let (status, cancel) = register_test_download(&ctx, &tx).await; ctx.active_operations .write() .await .insert("game".to_string(), OperationKind::Downloading); let guard = OperationGuard::download("game".to_string(), cancel.clone()); finish_successful_download_after_refresh( &ctx, &tx, &operation_target(games.path()), guard, &status, ) .await; assert!(ctx.active_operations.read().await.contains_key("game")); assert!(ctx.active_downloads.read().await.contains_key("game")); assert!(cancel.is_cancelled()); assert_no_event(&mut rx).await; } async fn recv_event(rx: &mut crate::PeerEventReceiver) -> PeerEvent { tokio::time::timeout(Duration::from_secs(1), rx.recv()) .await .expect("event should arrive") .expect("event channel should remain open") } async fn assert_no_event(rx: &mut crate::PeerEventReceiver) { assert!( tokio::time::timeout(Duration::from_millis(50), rx.recv()) .await .is_err(), "event channel should stay quiet" ); } fn drain_events(rx: &mut crate::PeerEventReceiver) -> Vec { let mut events = Vec::new(); while let Ok(event) = rx.try_recv() { events.push(event); } events } fn addr(port: u16) -> SocketAddr { SocketAddr::from(([127, 0, 0, 1], port)) } fn peer_id(seed: u8) -> PeerId { PeerId::from_bytes([seed; 32]) } fn endpoint(seed: u8, port: u16) -> PeerEndpoint { PeerEndpoint::new(peer_id(seed), addr(port)) } fn endpoint_at(seed: u8, ip: [u8; 4], port: u16) -> PeerEndpoint { PeerEndpoint::new(peer_id(seed), SocketAddr::from((ip, port))) } fn upsert(db: &mut PeerGameDB, endpoint: PeerEndpoint, game: Option<(&str, ContentId)>) { let ticket = db .begin_candidate_negotiation(endpoint) .expect("authenticated test peer should reserve"); db.commit_authenticated_snapshot( endpoint, ticket, RuntimeSessionId::from_bytes([1; 16]), Some(LibrarySnapshot { revision: 1, games: game .into_iter() .map(|(game_id, content_id)| GameAvailability { game_id: game_id.to_owned(), content_id, }) .collect(), }), ) .expect("authenticated snapshot should commit") .expect("candidate ticket should be current"); } fn summary( id: &str, version: &str, availability: Availability, ) -> crate::library::LocalGameSummary { crate::library::LocalGameSummary { id: id.to_string(), name: id.to_string(), size: 42, downloaded: availability == Availability::Ready, installed: true, eti_version: Some(version.to_string()), availability, } } fn assert_local_update(event: PeerEvent, installed: bool, downloaded: bool) { let _ = local_update_game(event, installed, downloaded); } fn local_update_game( event: PeerEvent, installed: bool, downloaded: bool, ) -> lanspread_db::db::Game { let PeerEvent::LocalLibraryChanged { games } = event else { panic!("expected LocalLibraryChanged"); }; let game = games .into_iter() .find(|game| game.id == "game") .expect("game should be announced"); assert_eq!(game.installed, installed); assert_eq!(game.downloaded, downloaded); game } fn assert_active_update(event: PeerEvent, expected: &[ActiveOperation]) { let PeerEvent::ActiveOperationsChanged { active_operations } = event else { panic!("expected ActiveOperationsChanged"); }; assert_eq!(active_operations, expected); } fn active_update(id: &str, operation: ActiveOperationKind) -> [ActiveOperation; 1] { [ActiveOperation { id: id.to_string(), operation, }] } #[test] fn content_sources_require_exact_content_and_authenticated_identity() { let matching_a = endpoint(20, 12_100); let matching_b = endpoint(21, 12_101); let wrong_content = endpoint(22, 12_102); let no_game = endpoint(23, 12_103); let content_id = ContentId::from_bytes([1; 32]); let other_content_id = ContentId::from_bytes([2; 32]); let mut db = PeerGameDB::new(); upsert(&mut db, matching_a, Some(("game", content_id))); upsert(&mut db, matching_b, Some(("game", content_id))); upsert(&mut db, wrong_content, Some(("game", other_content_id))); upsert(&mut db, no_game, None); let quarantine = ContentQuarantine::default(); quarantine.record_integrity_failure(&matching_a, content_id); assert_eq!( select_content_sources(&db, "game", content_id, &quarantine), vec![matching_b] ); assert_eq!( select_content_sources(&db, "game", other_content_id, &quarantine), vec![wrong_content] ); } #[test] fn streamed_install_same_ip_identities_and_ports_share_one_attempt() { let mut attempts = StreamInstallSourceAttempts::default(); let first = endpoint_at(1, [192, 168, 1, 10], 12_000); let same_ip_sybil = endpoint_at(2, [192, 168, 1, 10], 13_000); assert_eq!( attempts.admit(first.addr.ip()), StreamInstallSourceAdmission::Admitted ); assert_eq!( attempts.admit(same_ip_sybil.addr.ip()), StreamInstallSourceAdmission::AlreadyAttempted ); assert_eq!(attempts.source_ips.len(), 1); } #[test] fn streamed_install_stops_after_four_distinct_source_ip_attempts() { let mut attempts = StreamInstallSourceAttempts::default(); for seed in 1..=u8::try_from(MAX_STREAM_INSTALL_SOURCE_IP_ATTEMPTS) .expect("attempt limit should fit u8") { let source = endpoint_at(seed, [10, 0, 0, seed], 12_000 + u16::from(seed)); assert_eq!( attempts.admit(source.addr.ip()), StreamInstallSourceAdmission::Admitted ); } let excess = endpoint_at(100, [10, 0, 0, 100], 13_000); assert_eq!( attempts.admit(excess.addr.ip()), StreamInstallSourceAdmission::Exhausted ); assert_eq!( attempts.source_ips.len(), MAX_STREAM_INSTALL_SOURCE_IP_ATTEMPTS ); } #[test] fn streamed_install_retries_and_quarantines_only_typed_integrity_failures() { assert_eq!( stream_install_failure_disposition(StreamInstallReceiveErrorKind::Integrity), StreamInstallFailureDisposition::RetryAndQuarantine ); assert_eq!( stream_install_failure_disposition(StreamInstallReceiveErrorKind::Transport), StreamInstallFailureDisposition::Retry ); assert_eq!( stream_install_failure_disposition(StreamInstallReceiveErrorKind::Cancelled), StreamInstallFailureDisposition::Stop ); assert_eq!( stream_install_failure_disposition(StreamInstallReceiveErrorKind::Setup), StreamInstallFailureDisposition::Stop ); } #[test] fn alternate_stream_setup_failure_does_not_emit_retry_activity() { let games = TempDir::new("lanspread-handler-stream-retry-setup-failure"); let state = TempDir::new("lanspread-handler-stream-retry-setup-state"); std::fs::create_dir_all(games.game_root().join("local")) .expect("installed tree should be created"); let (tx, mut rx) = crate::peer_event_channel(); let status = DownloadAttemptStatus::new( DownloadAttemptKey::next("game".to_owned()), CancellationToken::new(), tx, ); let mut retry_invalid_source = true; let error = match begin_stream_receive_attempt( &games.game_root(), state.path(), "game", &mut retry_invalid_source, &status.reporter(), ) { Ok(transaction) => { transaction .rollback() .expect("unexpected transaction should roll back"); panic!("already-installed staging setup must fail"); } Err(error) => error, }; assert_eq!(error.reason, Some(DownloadFailureReason::OperationFailed)); assert!(retry_invalid_source); assert!(rx.try_recv().is_err()); } #[tokio::test] async fn active_game_projection_uses_peer_acceptable_revisions() { let temp = TempDir::new("lanspread-handler-active-hide"); let root = temp.game_root(); write_file(&root.join("version.ini"), b"20250101"); write_file(&root.join("game.eti"), b"archive"); let ctx = test_ctx(temp.path().to_path_buf()); let (tx, mut rx) = crate::peer_event_channel(); let catalog = ctx.catalog.catalog(); // 1. Initial scan: the game is ready and announced let scan = scan_local_library(temp.path(), ctx.state_dir.as_ref(), catalog) .await .expect("scan should succeed"); let source_revision = scan.revision; update_and_announce_games(&ctx, &tx, scan).await; let PeerEvent::LocalLibraryChanged { games } = recv_event(&mut rx).await else { panic!("expected LocalLibraryChanged"); }; assert_eq!(games.len(), 1); assert_eq!(games[0].id, "game"); let initial_snapshot = { let library = ctx.local_library.read().await; crate::library::build_library_snapshot( library.publication(ctx.catalog.catalog()), &ctx.catalog, ) .expect("catalog snapshot should build") }; let initial_revision = initial_snapshot.revision; assert_eq!(initial_snapshot.games.len(), 1); // 2. Set the game as active/in-progress and scan again ctx.active_operations .write() .await .insert("game".to_string(), OperationKind::Installing); let scan = scan_local_library(temp.path(), ctx.state_dir.as_ref(), catalog) .await .expect("scan should succeed"); assert_eq!(scan.revision, source_revision); update_and_announce_games(&ctx, &tx, scan).await; let PeerEvent::LocalLibraryChanged { games } = recv_event(&mut rx).await else { panic!("expected LocalLibraryChanged"); }; assert!( games.is_empty(), "active game should be hidden/unannounced during operations" ); let hidden_snapshot = { let library = ctx.local_library.read().await; crate::library::build_library_snapshot( library.publication(ctx.catalog.catalog()), &ctx.catalog, ) .expect("hidden catalog snapshot should build") }; assert!(hidden_snapshot.revision > initial_revision); assert!(hidden_snapshot.games.is_empty()); let hidden_revision = hidden_snapshot.revision; // 3. Operation completion republishes the game from the same unchanged // disk revision. The operation remains registered until publication, // so the ending ID is explicitly excluded from the hidden projection. let scan = scan_local_library(temp.path(), ctx.state_dir.as_ref(), catalog) .await .expect("ending-operation scan should succeed"); assert_eq!(scan.revision, source_revision); assert_eq!( update_and_announce_games_with_policy( &ctx, &tx, scan, LocalLibraryEventPolicy::OnChange, Some("game"), ) .await, LocalLibraryPublication::Published ); let PeerEvent::LocalLibraryChanged { games } = recv_event(&mut rx).await else { panic!("expected LocalLibraryChanged"); }; assert_eq!(games.len(), 1); assert_eq!(games[0].id, "game"); let restored_snapshot = { let library = ctx.local_library.read().await; crate::library::build_library_snapshot( library.publication(ctx.catalog.catalog()), &ctx.catalog, ) .expect("restored catalog snapshot should build") }; assert!(restored_snapshot.revision > hidden_revision); assert_eq!(restored_snapshot.games.len(), 1); assert_eq!(restored_snapshot.games[0].game_id, "game"); } #[tokio::test] async fn invalid_prospective_manifest_rejects_publication_without_visible_mutation() { let games = TempDir::new("lanspread-handler-invalid-publication-game-root"); let root = games.game_root(); write_file(&root.join("version.ini"), b"20250101"); write_file(&root.join("game.eti"), b"archive"); let manifests = TempDir::new("lanspread-handler-invalid-publication-manifests"); write_file(&manifests.path().join("game.json"), b"not JSON"); let content_index = CatalogContentIndex::from_entries([CatalogContentIndexEntry { game_id: "game".to_owned(), game_version: "20250101".to_owned(), identity: CatalogContentIdentity { content_id: ContentId::from_bytes([1; 32]), supports_streamed_install: false, }, }]) .expect("test content index should validate"); write_canonical_content_index_atomic( &manifests.path().join(CATALOG_CONTENT_INDEX_NAME), &content_index, ) .expect("test content index should publish"); let catalog = Arc::new( lanspread_db::content_manifest::CatalogBundle::new( manifests.path(), BTreeMap::from([("game".to_owned(), "20250101".to_owned())]), ) .expect("catalog construction should defer manifest body parsing"), ); let ctx = test_ctx_with_catalog(games.path().to_path_buf(), catalog); let (tx, mut rx) = crate::peer_event_channel(); let scan = scan_local_library(games.path(), ctx.state_dir.as_ref(), ctx.catalog.catalog()) .await .expect("ready game should scan before manifest priming"); assert_eq!( update_and_announce_games_with_policy( &ctx, &tx, scan, LocalLibraryEventPolicy::OnChange, None, ) .await, LocalLibraryPublication::Rejected ); let library = ctx.local_library.read().await; assert_eq!(library.revision, 0); assert!(library.games.is_empty()); drop(library); assert!(ctx.local_game_db.read().await.is_none()); assert_no_event(&mut rx).await; } #[tokio::test] async fn local_library_rejects_scan_from_previous_game_directory() { let current = TempDir::new("lanspread-handler-current-scan-root"); let previous = TempDir::new("lanspread-handler-previous-scan-root"); write_file(¤t.game_root().join("version.ini"), b"20250101"); write_file(¤t.game_root().join("game.eti"), b"archive"); write_file(&previous.game_root().join("version.ini"), b"20250101"); write_file(&previous.game_root().join("game.eti"), b"archive"); write_file( &previous.game_root().join("local/payload.txt"), b"installed", ); let ctx = test_ctx(current.path().to_path_buf()); let (tx, mut rx) = crate::peer_event_channel(); let catalog = ctx.catalog.catalog(); let current_scan = scan_local_library(current.path(), ctx.state_dir.as_ref(), catalog) .await .expect("current root should scan"); update_and_announce_games(&ctx, &tx, current_scan).await; assert_local_update(recv_event(&mut rx).await, false, true); let obsolete_scan = scan_local_library(previous.path(), ctx.state_dir.as_ref(), catalog) .await .expect("previous root should scan"); update_and_announce_games(&ctx, &tx, obsolete_scan).await; assert_no_event(&mut rx).await; let library = ctx.local_library.read().await; assert!(!library.games["game"].installed); drop(library); let game_db = ctx.local_game_db.read().await; assert!( !game_db .as_ref() .expect("current database should remain published") .get_game_by_id("game") .expect("current game should remain published") .installed ); } #[tokio::test] async fn local_library_rejects_scan_older_than_published_revision() { let temp = TempDir::new("lanspread-handler-stale-scan"); let root = temp.game_root(); write_file(&root.join("version.ini"), b"20250101"); write_file(&root.join("game.eti"), b"archive"); let ctx = test_ctx(temp.path().to_path_buf()); let (tx, mut rx) = crate::peer_event_channel(); let catalog = ctx.catalog.catalog(); let older_scan = scan_local_library(temp.path(), ctx.state_dir.as_ref(), catalog) .await .expect("older snapshot should scan"); seed_pending_download_ownership_for_test( ctx.state_dir.as_ref(), temp.path(), "game", &["game.eti"], &["game.eti"], ) .await; let newer_scan = rescan_local_game(temp.path(), ctx.state_dir.as_ref(), catalog, "game") .await .expect("quarantined snapshot should scan"); assert!(newer_scan.revision > older_scan.revision); let newer_revision = newer_scan.revision; update_and_announce_games(&ctx, &tx, newer_scan).await; assert_local_update(recv_event(&mut rx).await, false, false); update_and_announce_games(&ctx, &tx, older_scan).await; assert_no_event(&mut rx).await; let library = ctx.local_library.read().await; assert_eq!( library.source_revision_for(temp.path()), Some(newer_revision) ); assert!(!library.games["game"].downloaded); assert_eq!(library.games["game"].availability, Availability::LocalOnly); drop(library); let game_db = ctx.local_game_db.read().await; assert!( !game_db .as_ref() .expect("newer database should remain published") .get_game_by_id("game") .expect("newer game should remain published") .downloaded ); } #[tokio::test] async fn begin_operation_reports_authoritative_active_operation_snapshot() { let temp = TempDir::new("lanspread-handler-active-begin"); let root = temp.game_root(); write_file(&root.join("version.ini"), b"20250101"); write_file(&root.join("game.eti"), b"archive"); let ctx = test_ctx(temp.path().to_path_buf()); let (tx, mut rx) = crate::peer_event_channel(); let target = operation_target(temp.path()); assert_eq!( begin_operation(&ctx, &tx, &target, OperationKind::Updating).await, BeginOperationResult::Started ); assert_active_update( recv_event(&mut rx).await, &[ActiveOperation { id: "game".to_string(), operation: ActiveOperationKind::Updating, }], ); } #[tokio::test] async fn begin_operation_announces_busy_before_withdrawn_ui_snapshot() { let temp = TempDir::new("lanspread-handler-active-withdrawal"); let root = temp.game_root(); write_file(&root.join("version.ini"), b"20250101"); write_file(&root.join("game.eti"), b"archive"); let ctx = test_ctx(temp.path().to_path_buf()); let (tx, mut rx) = crate::peer_event_channel(); let scan = scan_local_library(temp.path(), ctx.state_dir.as_ref(), ctx.catalog.catalog()) .await .expect("initial scan should succeed"); update_and_announce_games(&ctx, &tx, scan).await; let PeerEvent::LocalLibraryChanged { games } = recv_event(&mut rx).await else { panic!("expected initial LocalLibraryChanged"); }; assert_eq!(games.len(), 1); let initial_revision = ctx.local_library.read().await.revision; assert_eq!( begin_operation( &ctx, &tx, &operation_target(temp.path()), OperationKind::Installing, ) .await, BeginOperationResult::Started ); assert_active_update( recv_event(&mut rx).await, &active_update("game", ActiveOperationKind::Installing), ); let PeerEvent::LocalLibraryChanged { games } = recv_event(&mut rx).await else { panic!("withdrawn UI snapshot must follow the active-operation event"); }; assert!(games.is_empty()); let snapshot = { let library = ctx.local_library.read().await; assert!(library.revision > initial_revision); crate::library::build_library_snapshot( library.publication(ctx.catalog.catalog()), &ctx.catalog, ) .expect("withdrawn snapshot should build") }; assert!(snapshot.games.is_empty()); assert!( ctx.local_game_db .read() .await .as_ref() .expect("local database should remain published") .all_games() .is_empty() ); } #[tokio::test] async fn begin_operation_timeout_clears_active_operation_snapshot() { let temp = TempDir::new("lanspread-handler-active-drain-timeout"); let root = temp.game_root(); write_file(&root.join("version.ini"), b"20250101"); write_file(&root.join("game.eti"), b"archive"); let ctx = test_ctx(temp.path().to_path_buf()); let (tx, mut rx) = crate::peer_event_channel(); let target = operation_target(temp.path()); let scan = scan_local_library(temp.path(), ctx.state_dir.as_ref(), ctx.catalog.catalog()) .await .expect("initial scan should succeed"); update_and_announce_games(&ctx, &tx, scan).await; assert_local_update(recv_event(&mut rx).await, false, true); let initial_revision = ctx.local_library.read().await.revision; let token = CancellationToken::new(); ctx.active_outbound_transfers .write() .await .insert("game".to_string(), vec![(1, token.clone())]); assert_eq!( begin_operation_with_drain_timeout( &ctx, &tx, &target, OperationKind::Updating, Duration::from_millis(1), ) .await, BeginOperationResult::DrainTimedOut ); assert!(token.is_cancelled()); assert_active_update( recv_event(&mut rx).await, &[ActiveOperation { id: "game".to_string(), operation: ActiveOperationKind::Updating, }], ); let PeerEvent::LocalLibraryChanged { games } = recv_event(&mut rx).await else { panic!("withdrawn UI snapshot must follow the active-operation event"); }; assert!(games.is_empty()); assert_local_update(recv_event(&mut rx).await, false, true); assert_active_update(recv_event(&mut rx).await, &[]); assert!( !ctx.active_operations.read().await.contains_key("game"), "timed-out drain should not leave the operation stuck active" ); let snapshot = { let library = ctx.local_library.read().await; assert_eq!(library.revision, initial_revision + 2); crate::library::build_library_snapshot( library.publication(ctx.catalog.catalog()), &ctx.catalog, ) .expect("restored snapshot should build") }; assert_eq!(snapshot.games.len(), 1); assert_eq!(snapshot.games[0].game_id, "game"); } #[tokio::test] async fn begin_operation_revision_exhaustion_is_fail_closed() { let temp = TempDir::new("lanspread-handler-active-revision-exhaustion"); let root = temp.game_root(); write_file(&root.join("version.ini"), b"20250101"); write_file(&root.join("game.eti"), b"archive"); let ctx = test_ctx(temp.path().to_path_buf()); let (tx, mut rx) = crate::peer_event_channel(); let scan = scan_local_library(temp.path(), ctx.state_dir.as_ref(), ctx.catalog.catalog()) .await .expect("initial scan should succeed"); update_and_announce_games(&ctx, &tx, scan).await; assert_local_update(recv_event(&mut rx).await, false, true); let before_games = { let mut library = ctx.local_library.write().await; library.revision = u64::MAX; library.games.clone() }; let before_db = ctx .local_game_db .read() .await .as_ref() .expect("initial database should be published") .clone(); assert_eq!( begin_operation( &ctx, &tx, &operation_target(temp.path()), OperationKind::Installing, ) .await, BeginOperationResult::PublicationFailed ); assert!(ctx.active_operations.read().await.is_empty()); let library = ctx.local_library.read().await; assert_eq!(library.revision, u64::MAX); assert_eq!(library.games, before_games); drop(library); assert_eq!( ctx.local_game_db .read() .await .as_ref() .expect("database should remain published") .all_games(), before_db.all_games(), ); assert_no_event(&mut rx).await; } #[test] fn payload_replacing_operations_require_outbound_transfer_drain() { assert!(operation_requires_outbound_drain( OperationKind::Downloading )); assert!(operation_requires_outbound_drain(OperationKind::Updating)); assert!(operation_requires_outbound_drain( OperationKind::RemovingDownload )); assert!(!operation_requires_outbound_drain( OperationKind::Installing )); assert!(!operation_requires_outbound_drain( OperationKind::Uninstalling )); } #[tokio::test] async fn unchanged_settled_scan_is_not_reemitted() { let temp = TempDir::new("lanspread-handler-settled-unchanged"); let root = temp.game_root(); write_file(&root.join("version.ini"), b"20250101"); write_file(&root.join("game.eti"), b"archive"); let ctx = test_ctx(temp.path().to_path_buf()); let (tx, mut rx) = crate::peer_event_channel(); let catalog = ctx.catalog.catalog(); let scan = scan_local_library(temp.path(), ctx.state_dir.as_ref(), catalog) .await .expect("first scan should succeed"); update_and_announce_games(&ctx, &tx, scan).await; assert_local_update(recv_event(&mut rx).await, false, true); let scan = scan_local_library(temp.path(), ctx.state_dir.as_ref(), catalog) .await .expect("second scan should succeed"); update_and_announce_games(&ctx, &tx, scan).await; assert_no_event(&mut rx).await; } #[tokio::test] async fn unchanged_operation_refresh_still_reports_settled_snapshot() { let temp = TempDir::new("lanspread-handler-operation-unchanged"); let root = temp.game_root(); write_file(&root.join("version.ini"), b"20250101"); write_file(&root.join("game.eti"), b"archive"); write_file(&root.join("local").join("old.txt"), b"old"); let ctx = test_ctx(temp.path().to_path_buf()); let (tx, mut rx) = crate::peer_event_channel(); let catalog = ctx.catalog.catalog(); let scan = scan_local_library(temp.path(), ctx.state_dir.as_ref(), catalog) .await .expect("initial scan should succeed"); update_and_announce_games(&ctx, &tx, scan).await; assert_local_update(recv_event(&mut rx).await, true, true); run_install_operation(&ctx, &tx, operation_target(temp.path())).await; assert_active_update( recv_event(&mut rx).await, &active_update("game", ActiveOperationKind::Updating), ); let PeerEvent::LocalLibraryChanged { games } = recv_event(&mut rx).await else { panic!("operation admission should withdraw the ready game"); }; assert!(games.is_empty()); assert_local_update(recv_event(&mut rx).await, true, true); assert_active_update(recv_event(&mut rx).await, &[]); assert!(matches!( recv_event(&mut rx).await, PeerEvent::InstallGameFinished { id } if id == "game" )); assert_no_event(&mut rx).await; } #[tokio::test] async fn install_refreshes_settled_state_before_operation_clear() { let temp = TempDir::new("lanspread-handler-install"); let root = temp.game_root(); write_file(&root.join("version.ini"), b"20250101"); write_file(&root.join("game.eti"), b"archive"); let ctx = test_ctx(temp.path().to_path_buf()); let (tx, mut rx) = crate::peer_event_channel(); run_install_operation(&ctx, &tx, operation_target(temp.path())).await; assert_active_update( recv_event(&mut rx).await, &active_update("game", ActiveOperationKind::Installing), ); assert_local_update(recv_event(&mut rx).await, true, true); assert_active_update(recv_event(&mut rx).await, &[]); assert!(matches!( recv_event(&mut rx).await, PeerEvent::InstallGameFinished { id } if id == "game" )); assert!(ctx.active_operations.read().await.is_empty()); } #[tokio::test] async fn install_rechecks_new_ownership_quarantine_after_admission() { let temp = TempDir::new("lanspread-handler-install-stale-preflight"); let root = temp.game_root(); write_file(&root.join("version.ini"), b"20250101"); write_file(&root.join("game.eti"), b"archive"); let ctx = test_ctx(temp.path().to_path_buf()); let target = operation_target(temp.path()); let (tx, mut rx) = crate::peer_event_channel(); let prepared = prepare_install_operation(&ctx, &tx, &target) .await .expect("initial preflight should accept the settled download"); assert_eq!( begin_operation(&ctx, &tx, &target, prepared.operation_kind).await, BeginOperationResult::Started ); let operation_guard = OperationGuard::new("game".to_string()); seed_pending_download_ownership_for_test( ctx.state_dir.as_ref(), temp.path(), "game", &["game.eti"], &["game.eti"], ) .await; assert!( revalidate_install_operation(&ctx, &tx, &target, prepared.operation_kind) .await .is_none(), "post-admission ownership quarantine must invalidate the stale preflight" ); assert!( settle_and_end_operation( &ctx, &tx, &target, operation_guard, "test stale install preflight", ) .await ); assert!(!root.join("local").exists()); assert!(ctx.active_operations.read().await.is_empty()); assert!( drain_events(&mut rx) .iter() .any(|event| matches!(event, PeerEvent::InstallGameFailed { id } if id == "game")) ); } #[tokio::test] async fn streamed_install_rechecks_new_ownership_quarantine_after_admission() { let temp = TempDir::new("lanspread-handler-stream-stale-preflight"); let root = temp.game_root(); write_file(&root.join("version.ini"), b"20250101"); write_file(&root.join("game.eti"), b"archive"); let ctx = test_ctx(temp.path().to_path_buf()); let target = operation_target(temp.path()); let (tx, mut rx) = crate::peer_event_channel(); assert!(stream_install_target_is_ready(&ctx, &target).await); assert_eq!( begin_operation(&ctx, &tx, &target, OperationKind::Downloading).await, BeginOperationResult::Started ); let (status, cancel) = register_test_download(&ctx, &tx).await; let guard = OperationGuard::download("game".to_string(), cancel); seed_pending_download_ownership_for_test( ctx.state_dir.as_ref(), temp.path(), "game", &["game.eti"], &["game.eti"], ) .await; assert!(!stream_install_target_is_ready(&ctx, &target).await); finish_failed_stream_download( &ctx, &tx, &target, guard, &status, Some(DownloadFailureReason::OperationFailed), ) .await; assert!(!root.join("local").exists()); assert!(ctx.active_operations.read().await.is_empty()); assert!(ctx.active_downloads.read().await.is_empty()); assert!(drain_events(&mut rx).iter().any(|event| matches!( event, PeerEvent::DownloadGameFilesFailed { attempt, reason: DownloadFailureReason::OperationFailed, } if attempt.id == "game" ))); } #[tokio::test] async fn download_handoff_waits_for_readers_and_auto_installs() { let temp = TempDir::new("lanspread-handler-download-handoff"); let root = temp.game_root(); write_file(&root.join("version.ini"), b"20250101"); write_file(&root.join("game.eti"), b"archive"); let ctx = test_ctx(temp.path().to_path_buf()); let (prepare_tx, _prepare_rx) = crate::peer_event_channel(); let (download_status, _download_cancel) = register_test_download(&ctx, &prepare_tx).await; ctx.active_operations .write() .await .insert("game".to_string(), OperationKind::Downloading); let target = operation_target(temp.path()); let prepared = prepare_install_operation(&ctx, &prepare_tx, &target) .await .expect("downloaded game should be installable"); let read_guard = ctx.active_operations.read().await; let (tx, mut rx) = crate::peer_event_channel(); let install_task = tokio::spawn({ let ctx = ctx.clone(); let tx = tx.clone(); let target = target.clone(); let attempt = download_status.key().clone(); async move { assert!( transition_download_to_install(&ctx, &tx, "game", prepared.operation_kind) .await ); clear_active_download(&ctx, &attempt).await; let cancel_token = CancellationToken::new(); let operation_guard = OperationGuard::cancellable("game".to_string(), cancel_token.clone()); run_started_install_operation( &ctx, &tx, target, prepared, operation_guard, cancel_token, ) .await; } }); tokio::task::yield_now().await; assert_eq!(read_guard.get("game"), Some(&OperationKind::Downloading)); drop(read_guard); install_task.await.expect("handoff task should finish"); assert_active_update( recv_event(&mut rx).await, &active_update("game", ActiveOperationKind::Installing), ); assert_local_update(recv_event(&mut rx).await, true, true); assert_active_update(recv_event(&mut rx).await, &[]); assert!(matches!( recv_event(&mut rx).await, PeerEvent::InstallGameFinished { id } if id == "game" )); assert!(ctx.active_operations.read().await.is_empty()); assert!(ctx.active_downloads.read().await.is_empty()); } #[tokio::test] async fn cancel_download_command_only_cancels_active_token() { let temp = TempDir::new("lanspread-handler-cancel-download"); let ctx = test_ctx(temp.path().to_path_buf()); let (tx, mut rx) = crate::peer_event_channel(); let (_status, cancel) = register_test_download(&ctx, &tx).await; ctx.active_operations .write() .await .insert("game".to_string(), OperationKind::Downloading); handle_cancel_download_command(&ctx, &tx, "game".to_string()).await; assert!(cancel.is_cancelled()); assert_eq!( ctx.active_operations.read().await.get("game"), Some(&OperationKind::Downloading), "the running transfer owns operation cleanup after cancellation" ); assert_no_event(&mut rx).await; } #[tokio::test] async fn update_refreshes_settled_state_before_operation_clear() { let temp = TempDir::new("lanspread-handler-update"); let root = temp.game_root(); write_file(&root.join("version.ini"), b"20250101"); write_file(&root.join("game.eti"), b"archive"); write_file(&root.join("local").join("old.txt"), b"old"); let ctx = test_ctx(temp.path().to_path_buf()); let (tx, mut rx) = crate::peer_event_channel(); run_install_operation(&ctx, &tx, operation_target(temp.path())).await; assert_active_update( recv_event(&mut rx).await, &active_update("game", ActiveOperationKind::Updating), ); assert_local_update(recv_event(&mut rx).await, true, true); assert_active_update(recv_event(&mut rx).await, &[]); assert!(matches!( recv_event(&mut rx).await, PeerEvent::InstallGameFinished { id } if id == "game" )); assert!(ctx.active_operations.read().await.is_empty()); } #[tokio::test] async fn install_update_uninstall_sequence_reports_new_version_and_settled_state() { let temp = TempDir::new("lanspread-handler-sequence"); let root = temp.game_root(); write_file(&root.join("version.ini"), b"20240101"); write_file(&root.join("game.eti"), b"old archive"); let ctx = test_ctx(temp.path().to_path_buf()); let (tx, mut rx) = crate::peer_event_channel(); run_install_operation(&ctx, &tx, operation_target(temp.path())).await; assert_active_update( recv_event(&mut rx).await, &active_update("game", ActiveOperationKind::Installing), ); let game = local_update_game(recv_event(&mut rx).await, true, true); assert_eq!(game.local_version.as_deref(), Some("20240101")); assert_active_update(recv_event(&mut rx).await, &[]); assert!(matches!( recv_event(&mut rx).await, PeerEvent::InstallGameFinished { id } if id == "game" )); write_file(&root.join("version.ini"), b"20250101"); write_file(&root.join("game.eti"), b"new archive"); run_install_operation(&ctx, &tx, operation_target(temp.path())).await; assert_active_update( recv_event(&mut rx).await, &active_update("game", ActiveOperationKind::Updating), ); let PeerEvent::LocalLibraryChanged { games } = recv_event(&mut rx).await else { panic!("update admission should withdraw the ready game"); }; assert!(games.is_empty()); let game = local_update_game(recv_event(&mut rx).await, true, true); assert_eq!(game.local_version.as_deref(), Some("20250101")); assert_active_update(recv_event(&mut rx).await, &[]); assert!(matches!( recv_event(&mut rx).await, PeerEvent::InstallGameFinished { id } if id == "game" )); run_uninstall_operation(&ctx, &tx, operation_target(temp.path())).await; assert_active_update( recv_event(&mut rx).await, &active_update("game", ActiveOperationKind::Uninstalling), ); let PeerEvent::LocalLibraryChanged { games } = recv_event(&mut rx).await else { panic!("uninstall admission should withdraw the ready game"); }; assert!(games.is_empty()); let game = local_update_game(recv_event(&mut rx).await, false, true); assert_eq!(game.local_version.as_deref(), Some("20250101")); assert_active_update(recv_event(&mut rx).await, &[]); assert!(matches!( recv_event(&mut rx).await, PeerEvent::UninstallGameFinished { id } if id == "game" )); assert!(ctx.active_operations.read().await.is_empty()); } #[tokio::test] async fn uninstall_refreshes_settled_state_before_operation_clear() { let temp = TempDir::new("lanspread-handler-uninstall"); let root = temp.game_root(); write_file(&root.join("version.ini"), b"20250101"); write_file(&root.join("game.eti"), b"archive"); write_file(&root.join("local").join("old.txt"), b"old"); let ctx = test_ctx(temp.path().to_path_buf()); let (tx, mut rx) = crate::peer_event_channel(); run_uninstall_operation(&ctx, &tx, operation_target(temp.path())).await; assert_active_update( recv_event(&mut rx).await, &active_update("game", ActiveOperationKind::Uninstalling), ); assert_local_update(recv_event(&mut rx).await, false, true); assert_active_update(recv_event(&mut rx).await, &[]); assert!(matches!( recv_event(&mut rx).await, PeerEvent::UninstallGameFinished { id } if id == "game" )); assert!(ctx.active_operations.read().await.is_empty()); } #[tokio::test] async fn remove_downloaded_refreshes_settled_state_before_operation_clear() { let temp = TempDir::new("lanspread-handler-remove-downloaded"); let root = temp.game_root(); write_file(&root.join("version.ini"), b"20250101"); write_file(&root.join("game.eti"), b"archive"); let ctx = test_ctx(temp.path().to_path_buf()); crate::download::seed_download_ownership_for_test( ctx.state_dir.as_ref(), temp.path(), "game", &["game.eti"], ) .await; let (tx, mut rx) = crate::peer_event_channel(); let catalog = ctx.catalog.catalog(); let scan = scan_local_library(temp.path(), ctx.state_dir.as_ref(), catalog) .await .expect("initial scan should succeed"); update_and_announce_games(&ctx, &tx, scan).await; assert_local_update(recv_event(&mut rx).await, false, true); run_remove_downloaded_operation(&ctx, &tx, operation_target(temp.path())).await; assert_active_update( recv_event(&mut rx).await, &active_update("game", ActiveOperationKind::RemovingDownload), ); let PeerEvent::LocalLibraryChanged { games } = recv_event(&mut rx).await else { panic!("download removal admission should withdraw the ready game"); }; assert!(games.is_empty()); assert_active_update(recv_event(&mut rx).await, &[]); assert!(matches!( recv_event(&mut rx).await, PeerEvent::RemoveDownloadedGameFinished { id } if id == "game" )); assert!(root.is_dir()); assert!(!root.join("version.ini").exists()); assert!(!root.join("game.eti").exists()); assert!(ctx.active_operations.read().await.is_empty()); } #[tokio::test] async fn install_command_does_not_retarget_after_set_game_dir() { let current = TempDir::new("lanspread-handler-install-captured-root"); let next = TempDir::new("lanspread-handler-install-next-root"); for root in [current.game_root(), next.game_root()] { write_file(&root.join("version.ini"), b"20250101"); write_file(&root.join("game.eti"), b"archive"); } let ctx = test_ctx(current.path().to_path_buf()); let (tx, mut rx) = crate::peer_event_channel(); // On the current-thread test runtime, this ready lock acquisition does // not yield to the newly spawned command. Move the configured root // after capture so its operation must reject the stale target. handle_install_game_command(&ctx, &tx, "game".to_string()).await; *ctx.game_dir.write().await = next.path().to_path_buf(); assert_eq!(*ctx.game_dir.read().await, next.path()); assert!(!current.game_root().join("local").exists()); assert!(!next.game_root().join("local").exists()); let terminal_event = tokio::time::timeout(Duration::from_secs(1), async { loop { let event = recv_event(&mut rx).await; if matches!(event, PeerEvent::InstallGameFailed { ref id } if id == "game") { break event; } assert!(!matches!( event, PeerEvent::InstallGameFinished { ref id } if id == "game" )); } }) .await .expect("captured-root install should report failure"); assert!(matches!( terminal_event, PeerEvent::InstallGameFailed { id } if id == "game" )); } #[tokio::test] async fn uninstall_command_does_not_retarget_after_set_game_dir() { let current = TempDir::new("lanspread-handler-uninstall-captured-root"); let next = TempDir::new("lanspread-handler-uninstall-next-root"); for root in [current.game_root(), next.game_root()] { write_file(&root.join("version.ini"), b"20250101"); write_file(&root.join("game.eti"), b"archive"); write_file(&root.join("local/payload.txt"), b"installed"); } let ctx = test_ctx(current.path().to_path_buf()); let (tx, mut rx) = crate::peer_event_channel(); handle_uninstall_game_command(&ctx, &tx, "game".to_string()).await; *ctx.game_dir.write().await = next.path().to_path_buf(); assert_eq!(*ctx.game_dir.read().await, next.path()); assert!(current.game_root().join("local/payload.txt").is_file()); assert!(next.game_root().join("local/payload.txt").is_file()); let terminal_event = tokio::time::timeout(Duration::from_secs(1), async { loop { let event = recv_event(&mut rx).await; if matches!(event, PeerEvent::UninstallGameFailed { ref id } if id == "game") { break event; } assert!(!matches!( event, PeerEvent::UninstallGameFinished { ref id } if id == "game" )); } }) .await .expect("captured-root uninstall should report failure"); assert!(matches!( terminal_event, PeerEvent::UninstallGameFailed { id } if id == "game" )); } #[tokio::test] async fn remove_downloaded_command_does_not_retarget_after_set_game_dir() { let current = TempDir::new("lanspread-handler-remove-captured-root"); let next = TempDir::new("lanspread-handler-remove-next-root"); for root in [current.game_root(), next.game_root()] { write_file(&root.join("version.ini"), b"20250101"); write_file(&root.join("game.eti"), b"archive"); } let ctx = test_ctx(current.path().to_path_buf()); crate::download::seed_download_ownership_for_test( ctx.state_dir.as_ref(), current.path(), "game", &["game.eti"], ) .await; let (tx, mut rx) = crate::peer_event_channel(); handle_remove_downloaded_game_command(&ctx, &tx, "game".to_string()).await; *ctx.game_dir.write().await = next.path().to_path_buf(); assert_eq!(*ctx.game_dir.read().await, next.path()); assert!(current.game_root().join("version.ini").is_file()); assert!(current.game_root().join("game.eti").is_file()); assert!(next.game_root().join("version.ini").is_file()); assert!(next.game_root().join("game.eti").is_file()); let terminal_event = tokio::time::timeout(Duration::from_secs(1), async { loop { let event = recv_event(&mut rx).await; if matches!( event, PeerEvent::RemoveDownloadedGameFailed { ref id } if id == "game" ) { break event; } assert!(!matches!( event, PeerEvent::RemoveDownloadedGameFinished { ref id } if id == "game" )); } }) .await .expect("captured-root removal should report failure"); assert!(matches!( terminal_event, PeerEvent::RemoveDownloadedGameFailed { id } if id == "game" )); } #[tokio::test] async fn active_operation_rejection_is_returned_in_set_game_dir_reply() { let current = TempDir::new("lanspread-handler-current-dir"); let next = TempDir::new("lanspread-handler-next-dir"); let ctx = test_ctx(current.path().to_path_buf()); ctx.active_operations .write() .await .insert("game".to_string(), OperationKind::Downloading); let (tx, _rx) = crate::peer_event_channel(); let error = handle_set_game_dir_command(&ctx, &tx, next.path().to_path_buf()) .await .expect_err("active operation should reject SetGameDir"); assert!(error.contains("operations are active")); assert_eq!(*ctx.game_dir.read().await, current.path()); } #[cfg(unix)] #[tokio::test] async fn same_directory_alias_reply_uses_existing_canonical_root() { use std::os::unix::fs::symlink; let current = TempDir::new("lanspread-handler-same-alias-target"); let aliases = TempDir::new("lanspread-handler-same-alias-parent"); let alias = aliases.path().join("games"); symlink(current.path(), &alias).expect("same-root alias should be created"); let current = std::fs::canonicalize(current.path()).expect("root should canonicalize"); let ctx = test_ctx(current.clone()); let (tx, _rx) = crate::peer_event_channel(); let accepted = handle_set_game_dir_command(&ctx, &tx, alias) .await .expect("same-root alias should be accepted"); assert_eq!(accepted, current); assert_eq!(*ctx.game_dir.read().await, current); } #[cfg(unix)] #[tokio::test] async fn path_change_alias_reply_and_storage_use_canonical_target() { use std::os::unix::fs::symlink; let current = TempDir::new("lanspread-handler-change-alias-current"); let next = TempDir::new("lanspread-handler-change-alias-target"); let aliases = TempDir::new("lanspread-handler-change-alias-parent"); let alias = aliases.path().join("games"); symlink(next.path(), &alias).expect("new-root alias should be created"); let current = std::fs::canonicalize(current.path()).expect("root should canonicalize"); let next = std::fs::canonicalize(next.path()).expect("new root should canonicalize"); let ctx = test_ctx(current); let (tx, _rx) = crate::peer_event_channel(); let accepted = handle_set_game_dir_command(&ctx, &tx, alias.clone()) .await .expect("new-root alias should be accepted"); assert_eq!(accepted, next); assert_ne!(accepted, alias); assert_eq!(*ctx.game_dir.read().await, next); } #[tokio::test] async fn invalid_directory_rejection_reply_preserves_current_root() { let current = TempDir::new("lanspread-handler-invalid-dir-current"); let missing = current.path().join("missing"); let ctx = test_ctx(current.path().to_path_buf()); let (tx, mut rx) = crate::peer_event_channel(); let error = handle_set_game_dir_command(&ctx, &tx, missing) .await .expect_err("missing root should be rejected"); assert!(error.contains("failed to canonicalize game directory")); assert_eq!(*ctx.game_dir.read().await, current.path()); assert_no_event(&mut rx).await; } #[tokio::test] async fn foreign_intent_for_absent_candidate_game_rejects_before_root_mutation() { let current = TempDir::new("lanspread-handler-foreign-intent-current"); let candidate = TempDir::new("lanspread-handler-foreign-intent-candidate"); let current_path = std::fs::canonicalize(current.path()).expect("current root should canonicalize"); let candidate_path = std::fs::canonicalize(candidate.path()).expect("candidate root should canonicalize"); let ctx = test_ctx(current_path.clone()); let (tx, mut rx) = crate::peer_event_channel(); let published = summary("game", "20250101", Availability::Ready); ctx.local_library .write() .await .update_from_scan( ¤t_path, HashMap::from([("game".to_string(), published.clone())]), 1, ) .expect("initial library revision should be available"); *ctx.local_game_db.write().await = Some(GameDB::from(vec![game_from_summary(&published)])); let before = ctx.local_library.read().await.clone(); let intent = InstallIntent::new( ¤t_path.join("orphan"), "orphan", InstallIntentState::Updating, None, ) .expect("old-root intent should be valid"); write_intent(ctx.state_dir.as_ref(), "orphan", &intent) .expect("old-root intent should be persisted"); let error = handle_set_game_dir_command(&ctx, &tx, candidate_path.clone()) .await .expect_err("foreign intent must reject the root switch"); assert!(error.contains("install intent state is unresolved")); assert!(error.contains("different configured games")); assert_eq!(*ctx.game_dir.read().await, current_path); assert!(!candidate_path.join("orphan").exists()); assert!(intent_path(ctx.state_dir.as_ref(), "orphan").is_file()); assert!(!ctx.recovery_quarantine.is_blocked(¤t_path, "game")); assert!(ctx.local_game_db.read().await.is_some()); let after = ctx.local_library.read().await; assert_eq!(after.revision, before.revision); assert_eq!(after.games, before.games); drop(after); assert_no_event(&mut rx).await; } #[tokio::test] async fn same_path_set_game_dir_drains_all_transfers_before_recovery() { let temp = TempDir::new("lanspread-handler-same-dir"); write_file(&temp.game_root().join(".version.ini.tmp"), b"tmp"); let ctx = test_ctx(temp.path().to_path_buf()); let (tx, _rx) = crate::peer_event_channel(); let first = CancellationToken::new(); let second = CancellationToken::new(); *ctx.active_outbound_transfers.write().await = HashMap::from([ ("game".to_string(), vec![(1, first.clone())]), ("other".to_string(), vec![(2, second.clone())]), ]); let command_ctx = ctx.clone(); let command_tx = tx.clone(); let game_dir = temp.path().to_path_buf(); let command = tokio::spawn(async move { handle_set_game_dir_command_with_drain_timeout( &command_ctx, &command_tx, game_dir, Duration::from_secs(1), ) .await }); tokio::time::timeout(Duration::from_secs(1), async { first.cancelled().await; second.cancelled().await; }) .await .expect("all outbound transfers should be cancelled"); assert!(temp.game_root().join(".version.ini.tmp").is_file()); ctx.active_outbound_transfers.write().await.remove("game"); tokio::task::yield_now().await; assert!(!command.is_finished()); assert!(temp.game_root().join(".version.ini.tmp").is_file()); ctx.active_outbound_transfers.write().await.remove("other"); assert_eq!( command .await .expect("SetGameDir task should not fail") .expect("same root should be accepted"), temp.path().to_path_buf() ); assert!(ctx.active_outbound_transfers.read().await.is_empty()); assert!(!temp.game_root().join(".version.ini.tmp").exists()); } #[tokio::test] async fn same_path_drain_timeout_rejection_is_returned_in_reply() { let temp = TempDir::new("lanspread-handler-same-dir-timeout"); write_file(&temp.game_root().join(".version.ini.tmp"), b"tmp"); let ctx = test_ctx(temp.path().to_path_buf()); let (tx, mut rx) = crate::peer_event_channel(); let token = CancellationToken::new(); ctx.active_outbound_transfers .write() .await .insert("game".to_string(), vec![(1, token.clone())]); let error = handle_set_game_dir_command_with_drain_timeout( &ctx, &tx, temp.path().to_path_buf(), Duration::from_millis(1), ) .await .expect_err("drain timeout should reject SetGameDir"); assert!(error.contains("did not drain")); assert!(token.is_cancelled()); assert_eq!(*ctx.game_dir.read().await, temp.path()); assert!(temp.game_root().join(".version.ini.tmp").is_file()); assert!(!ctx.recovery_quarantine.is_blocked(temp.path(), "game")); assert_no_event(&mut rx).await; } #[tokio::test] async fn same_path_set_game_dir_skips_recovery_for_active_game() { let temp = TempDir::new("lanspread-handler-same-dir-active"); write_file(&temp.game_root().join(".version.ini.tmp"), b"tmp"); let ctx = test_ctx(temp.path().to_path_buf()); ctx.active_operations .write() .await .insert("game".to_string(), OperationKind::Downloading); let (tx, _rx) = crate::peer_event_channel(); let error = handle_set_game_dir_command(&ctx, &tx, temp.path().to_path_buf()) .await .expect_err("active game should reject SetGameDir refresh"); assert!(error.contains("operations are active")); assert!(temp.game_root().join(".version.ini.tmp").is_file()); } #[tokio::test] async fn path_changing_set_game_dir_drains_all_transfers_before_switch() { let current = TempDir::new("lanspread-handler-old-dir"); let next = TempDir::new("lanspread-handler-new-dir"); write_file(&next.game_root().join(".version.ini.tmp"), b"tmp"); let ctx = test_ctx(current.path().to_path_buf()); let (tx, _rx) = crate::peer_event_channel(); let first = CancellationToken::new(); let second = CancellationToken::new(); *ctx.active_outbound_transfers.write().await = HashMap::from([ ("game".to_string(), vec![(1, first.clone())]), ("other".to_string(), vec![(2, second.clone())]), ]); let command_ctx = ctx.clone(); let command_tx = tx.clone(); let next_path = next.path().to_path_buf(); let command = tokio::spawn(async move { handle_set_game_dir_command_with_drain_timeout( &command_ctx, &command_tx, next_path, Duration::from_secs(1), ) .await }); tokio::time::timeout(Duration::from_secs(1), async { first.cancelled().await; second.cancelled().await; }) .await .expect("all outbound transfers should be cancelled"); assert_eq!(*ctx.game_dir.read().await, current.path()); assert!(next.game_root().join(".version.ini.tmp").is_file()); ctx.active_outbound_transfers.write().await.remove("game"); tokio::task::yield_now().await; assert!(!command.is_finished()); assert_eq!(*ctx.game_dir.read().await, current.path()); ctx.active_outbound_transfers.write().await.remove("other"); assert_eq!( command .await .expect("SetGameDir task should not fail") .expect("new root should be accepted"), next.path().to_path_buf() ); assert_eq!(*ctx.game_dir.read().await, next.path()); assert!(!next.game_root().join(".version.ini.tmp").exists()); } #[tokio::test] async fn path_change_drain_timeout_rejection_preserves_current_state() { let current = TempDir::new("lanspread-handler-old-dir-timeout"); let next = TempDir::new("lanspread-handler-new-dir-timeout"); let root = current.game_root(); write_file(&root.join("version.ini"), b"20250101"); write_file(&root.join("game.eti"), b"archive"); let ctx = test_ctx(current.path().to_path_buf()); let (tx, mut rx) = crate::peer_event_channel(); let catalog = ctx.catalog.catalog(); let scan = scan_local_library(current.path(), ctx.state_dir.as_ref(), catalog) .await .expect("initial scan should succeed"); update_and_announce_games(&ctx, &tx, scan).await; let _ = recv_event(&mut rx).await; assert!( ctx.recovery_quarantine .settle(current.path(), HashSet::from(["game".to_string()])) ); let library_before = ctx.local_library.read().await.clone(); assert!(ctx.local_game_db.read().await.is_some()); write_file(&root.join(".version.ini.tmp"), b"old tmp"); write_file(&next.game_root().join(".version.ini.tmp"), b"new tmp"); let token = CancellationToken::new(); ctx.active_outbound_transfers .write() .await .insert("game".to_string(), vec![(1, token.clone())]); let error = handle_set_game_dir_command_with_drain_timeout( &ctx, &tx, next.path().to_path_buf(), Duration::from_millis(1), ) .await .expect_err("drain timeout should reject SetGameDir"); assert!(error.contains("did not drain")); assert!(token.is_cancelled()); assert_eq!(*ctx.game_dir.read().await, current.path()); assert!(root.join(".version.ini.tmp").is_file()); assert!(next.game_root().join(".version.ini.tmp").is_file()); assert_eq!( ctx.recovery_quarantine.failed_ids(current.path()), HashSet::from(["game".to_string()]) ); let library_after = ctx.local_library.read().await; assert_eq!(library_after.revision, library_before.revision); assert_eq!(library_after.games, library_before.games); drop(library_after); assert!(ctx.local_game_db.read().await.is_some()); assert_no_event(&mut rx).await; } #[tokio::test] async fn path_change_fails_before_mutation_when_publication_revision_is_exhausted() { let current = TempDir::new("lanspread-handler-revision-exhausted-current"); let next = TempDir::new("lanspread-handler-revision-exhausted-next"); let ctx = test_ctx(current.path().to_path_buf()); let (tx, mut rx) = crate::peer_event_channel(); let published = summary("game", "20250101", Availability::Ready); ctx.local_library .write() .await .update_from_scan( current.path(), HashMap::from([("game".to_string(), published.clone())]), u64::MAX, ) .expect("maximum source revision may be published once"); ctx.local_library.write().await.revision = u64::MAX; *ctx.local_game_db.write().await = Some(GameDB::from(vec![game_from_summary(&published)])); let error = handle_set_game_dir_command_with_drain_timeout( &ctx, &tx, next.path().to_path_buf(), Duration::from_secs(1), ) .await .expect_err("pre-mutation recovery failure should reject SetGameDir"); assert!(error.contains("failed to begin recovery")); assert_eq!(*ctx.game_dir.read().await, current.path()); let library = ctx.local_library.read().await; assert_eq!(library.revision, u64::MAX); assert_eq!(library.games["game"], published); drop(library); assert!(ctx.local_game_db.read().await.is_some()); assert_no_event(&mut rx).await; } #[tokio::test] async fn accepted_quarantined_root_is_returned_in_canonical_reply() { let current = TempDir::new("lanspread-handler-old-dir-recovery-failure"); let next = TempDir::new("lanspread-handler-new-dir-recovery-failure"); let next_root = next.game_root(); write_file(&next_root.join("version.ini"), b"20250101"); write_file(&next_root.join("game.eti"), b"archive"); std::fs::create_dir_all(next_root.join(".version.ini.tmp")) .expect("invalid recovery scratch directory should be created"); let ctx = test_ctx(current.path().to_path_buf()); let (tx, mut rx) = crate::peer_event_channel(); assert_eq!( handle_set_game_dir_command_with_drain_timeout( &ctx, &tx, next.path().to_path_buf(), Duration::from_secs(1), ) .await .expect("accepted root should be acknowledged despite quarantine"), next.path().to_path_buf() ); assert_eq!(*ctx.game_dir.read().await, next.path()); assert!(ctx.recovery_quarantine.is_blocked(next.path(), "game")); let library = ctx.local_library.read().await; let game = library .games .get("game") .expect("failed game should be published in quarantined state"); assert!(!game.downloaded); assert!(!game.installed); drop(library); assert_eq!( begin_operation( &ctx, &tx, &operation_target(next.path()), OperationKind::Downloading, ) .await, BeginOperationResult::RecoveryBlocked ); assert!(ctx.active_operations.read().await.is_empty()); let PeerEvent::LocalLibraryChanged { games } = recv_event(&mut rx).await else { panic!("expected cleared local library snapshot"); }; assert!(games.is_empty()); let PeerEvent::LocalLibraryChanged { games } = recv_event(&mut rx).await else { panic!("expected quarantined local library snapshot"); }; assert_eq!(games.len(), 1); assert_eq!(games[0].id, "game"); assert!(!games[0].downloaded); assert_eq!(games[0].availability, Availability::LocalOnly); } #[tokio::test] async fn path_changing_set_game_dir_emits_equivalent_snapshot() { let current = TempDir::new("lanspread-handler-old-equivalent-dir"); let next = TempDir::new("lanspread-handler-new-equivalent-dir"); for root in [current.game_root(), next.game_root()] { write_file(&root.join("version.ini"), b"20250101"); write_file(&root.join("game.eti"), b"archive"); } let ctx = test_ctx(current.path().to_path_buf()); let (tx, mut rx) = crate::peer_event_channel(); let catalog = ctx.catalog.catalog(); let scan = scan_local_library(current.path(), ctx.state_dir.as_ref(), catalog) .await .expect("initial scan should succeed"); update_and_announce_games(&ctx, &tx, scan).await; assert_local_update(recv_event(&mut rx).await, false, true); let initial_snapshot = { let library = ctx.local_library.read().await; crate::library::build_library_snapshot( library.publication(ctx.catalog.catalog()), &ctx.catalog, ) .expect("initial catalog snapshot should build") }; let initial_revision = initial_snapshot.revision; assert_eq!(initial_snapshot.games.len(), 1); assert_eq!( handle_set_game_dir_command(&ctx, &tx, next.path().to_path_buf()) .await .expect("equivalent new root should be accepted"), next.path().to_path_buf() ); let PeerEvent::LocalLibraryChanged { games } = recv_event(&mut rx).await else { panic!("expected cleared local library snapshot"); }; assert!(games.is_empty()); assert_local_update(recv_event(&mut rx).await, false, true); let restored_snapshot = { let library = ctx.local_library.read().await; crate::library::build_library_snapshot( library.publication(ctx.catalog.catalog()), &ctx.catalog, ) .expect("restored catalog snapshot should build") }; assert!(restored_snapshot.revision >= initial_revision + 2); assert_eq!(restored_snapshot.games.len(), 1); assert_eq!(restored_snapshot.games[0].game_id, "game"); } }