diff --git a/crates/lanspread-peer/src/services/server.rs b/crates/lanspread-peer/src/services/server.rs index 77e22f1..19dc273 100644 --- a/crates/lanspread-peer/src/services/server.rs +++ b/crates/lanspread-peer/src/services/server.rs @@ -430,8 +430,7 @@ async fn run_server_body( if !has_child_capacity(connection_tasks.len(), MAX_ESTABLISHED_CONNECTIONS) { log::warn!( - "Closing excess peer connection from {remote_addr} at application limit {}", - MAX_ESTABLISHED_CONNECTIONS, + "Closing excess peer connection from {remote_addr} at application limit {MAX_ESTABLISHED_CONNECTIONS}", ); connection.close(application::Error::UNKNOWN); continue; diff --git a/crates/lanspread-peer/src/services/stream.rs b/crates/lanspread-peer/src/services/stream.rs index 56d1dc1..ed3819d 100644 --- a/crates/lanspread-peer/src/services/stream.rs +++ b/crates/lanspread-peer/src/services/stream.rs @@ -300,55 +300,88 @@ async fn dispatch_request( .schedule_hint(StateDomain::CallToPlay, hint, source_ip); DispatchResult::close(framed_tx) } - Request::GetGameFileChunk { - game_id, - content_id, - relative_path, - offset, - length, - } => { - match handle_file_chunk_request( - ctx, - game_id, - content_id, - relative_path, - offset, - length, - framed_tx, - stream_shutdown, - ) - .await - { - ChunkDispatch::Finished(writer) => DispatchResult::close(writer), - ChunkDispatch::Reset(writer) => DispatchResult::reset(writer), - } + request @ Request::GetGameFileChunk { .. } => { + dispatch_file_chunk(ctx, request, framed_tx, stream_shutdown).await } Request::StreamInstall { game_id, content_id, } => { - let Some(mut stream_install_admission) = admission.try_acquire_stream_install(origin) - else { - let mut tx = framed_tx.into_inner(); - let _ = tx.reset(application::Error::UNKNOWN); - return DispatchResult::reset(FramedWrite::new(tx, control_codec())); - }; - let stream_install_permit = stream_install_admission.take_global_permit(); - let writer = handle_stream_install_request( + dispatch_stream_install( ctx, game_id, content_id, framed_tx, stream_shutdown, - stream_install_permit, + admission, + origin, ) - .await; - drop(stream_install_admission); - DispatchResult::close(writer) + .await } } } +async fn dispatch_file_chunk( + ctx: &PeerCtx, + request: Request, + framed_tx: ResponseWriter, + stream_shutdown: &CancellationToken, +) -> DispatchResult { + let Request::GetGameFileChunk { + game_id, + content_id, + relative_path, + offset, + length, + } = request + else { + unreachable!("file-chunk dispatcher received a different request") + }; + match handle_file_chunk_request( + ctx, + game_id, + content_id, + relative_path, + offset, + length, + framed_tx, + stream_shutdown, + ) + .await + { + ChunkDispatch::Finished(writer) => DispatchResult::close(writer), + ChunkDispatch::Reset(writer) => DispatchResult::reset(writer), + } +} + +async fn dispatch_stream_install( + ctx: &PeerCtx, + game_id: String, + content_id: lanspread_db::content_manifest::ContentId, + framed_tx: ResponseWriter, + stream_shutdown: &CancellationToken, + admission: &ServerAdmission, + origin: ObservedOrigin, +) -> DispatchResult { + let Some(mut stream_install_admission) = admission.try_acquire_stream_install(origin) else { + let mut tx = framed_tx.into_inner(); + let _ = tx.reset(application::Error::UNKNOWN); + return DispatchResult::reset(FramedWrite::new(tx, control_codec())); + }; + let stream_install_permit = stream_install_admission.take_global_permit(); + let writer = handle_stream_install_request( + ctx, + game_id, + content_id, + framed_tx, + stream_shutdown, + stream_install_permit, + ) + .await; + drop(stream_install_admission); + DispatchResult::close(writer) +} + fn reset_response_writer(framed_tx: ResponseWriter, label: &str) -> DispatchResult { let mut tx = framed_tx.into_inner(); if let Err(error) = tx.reset(application::Error::UNKNOWN) {