diff --git a/crates/lanspread-peer/ARCHITECTURE.md b/crates/lanspread-peer/ARCHITECTURE.md index 274fa2a..6a4096f 100644 --- a/crates/lanspread-peer/ARCHITECTURE.md +++ b/crates/lanspread-peer/ARCHITECTURE.md @@ -63,17 +63,22 @@ local action is applied to that history, sent to the UI, and broadcast to every currently known peer. An incoming live event is applied once and sent to the UI without being rebroadcast, which prevents forwarding loops. -If a live Call to Play delivery fails, the sender immediately falls back to a -normal `Hello` / `HelloAck` exchange with that peer. The handshake carries the -full active history in both directions, so a transient request failure heals -without waiting for mDNS rediscovery or a later reconnect. +Live Call to Play delivery is acknowledged by the receiver. Applied, duplicate, +and obsolete events need no follow-up. An unknown envelope peer, missing call +root, transport failure, or malformed acknowledgement makes the sender perform +one normal `Hello` / `HelloAck` exchange with that peer. The handshake carries +the full retained history in both directions, so a transient request failure +heals without waiting for mDNS rediscovery or a later reconnect. A rejected +event is logged without retry. Local publication remains successful while this +healing happens asynchronously, so an offline peer cannot block an action. Actors are keyed by the peer's stable ID and carry a separate display name. The -origin peer overwrites the actor ID on local actions, and live-event envelopes -must match the known sending peer. This prevents duplicate default usernames -from merging participants and protects creator controls from other normal -clients. It is not authentication against a hostile LAN peer; the QUIC setup -uses the project's trusted-LAN identity model. +origin peer overwrites the actor ID on local actions. A live-event envelope must +name a peer already in the receiver's roster, and every enclosed actor ID must +match that envelope. This prevents accidental identity mixing and protects +creator controls from other normal clients. It is not authentication against a +hostile LAN peer: all peers use the shared application TLS identity, and stable +peer IDs are self-asserted under the project's trusted-LAN model. `Hello` and `HelloAck` include each side's event history. This lets peers that join after a call was created reconstruct the same nominations, responses, diff --git a/crates/lanspread-peer/src/call_to_play.rs b/crates/lanspread-peer/src/call_to_play.rs index 461880f..fa7a64f 100644 --- a/crates/lanspread-peer/src/call_to_play.rs +++ b/crates/lanspread-peer/src/call_to_play.rs @@ -6,7 +6,7 @@ use std::{ time::{SystemTime, UNIX_EPOCH}, }; -use lanspread_proto::{CallToPlayAction, CallToPlayEvent}; +use lanspread_proto::{CallToPlayAck, CallToPlayAction, CallToPlayEvent}; use tokio::sync::mpsc::UnboundedSender; use crate::{ @@ -416,18 +416,7 @@ pub(crate) async fn publish( let peer_id = peer_id.clone(); let handshake_ctx = handshake_ctx.clone(); async move { - if let Err(err) = - send_call_to_play_events(peer_addr, peer_id.as_ref(), vec![event]).await - { - log::warn!("Failed to send Call to Play event to {peer_addr}: {err}"); - if let Err(resync_err) = - perform_handshake_with_peer(handshake_ctx, peer_addr, None).await - { - log::warn!( - "Failed to resync Call to Play history with {peer_addr}: {resync_err}" - ); - } - } + deliver_to_peer(handshake_ctx, peer_addr, peer_id.as_ref(), event).await; } }); futures::future::join_all(deliveries).await; @@ -435,6 +424,53 @@ pub(crate) async fn publish( Ok(()) } +async fn deliver_to_peer( + handshake_ctx: HandshakeCtx, + peer_addr: std::net::SocketAddr, + peer_id: &str, + event: CallToPlayEvent, +) { + let delivery = send_call_to_play_events(peer_addr, peer_id, vec![event]).await; + match &delivery { + Ok(CallToPlayAck::Rejected { reason }) => { + log::warn!("Peer {peer_addr} rejected a Call to Play event: {reason}"); + } + Ok(CallToPlayAck::Obsolete) => { + log::debug!("Peer {peer_addr} already retired the Call to Play event"); + } + Err(err) => { + log::warn!("Failed to deliver a Call to Play event to {peer_addr}: {err}"); + } + Ok( + CallToPlayAck::Applied + | CallToPlayAck::Duplicate + | CallToPlayAck::NeedHandshake + | CallToPlayAck::NeedHistory, + ) => {} + } + + let Some(reason) = delivery_resync_reason(delivery.as_ref().map_err(|_| ())) else { + return; + }; + if let Err(err) = perform_handshake_with_peer(handshake_ctx, peer_addr, None).await { + log::warn!("Failed to {reason} with {peer_addr}: {err}"); + } +} + +fn delivery_resync_reason(delivery: Result<&CallToPlayAck, ()>) -> Option<&'static str> { + match delivery { + Err(()) => Some("heal a failed Call to Play delivery"), + Ok(CallToPlayAck::NeedHandshake) => Some("complete a requested Call to Play handshake"), + Ok(CallToPlayAck::NeedHistory) => Some("restore missing Call to Play history"), + Ok( + CallToPlayAck::Applied + | CallToPlayAck::Duplicate + | CallToPlayAck::Obsolete + | CallToPlayAck::Rejected { .. }, + ) => None, + } +} + fn validate_event(event: &CallToPlayEvent) -> Result<(), &'static str> { validate_nonempty(&event.id, MAX_ID_CHARS, "invalid event id")?; validate_nonempty(&event.call_id, MAX_ID_CHARS, "invalid call id")?; @@ -498,9 +534,9 @@ fn validate_nonempty( #[cfg(test)] mod tests { - use lanspread_proto::{CallToPlayAction, CallToPlayEvent}; + use lanspread_proto::{CallToPlayAck, CallToPlayAction, CallToPlayEvent}; - use super::{CallToPlayStore, MAX_EVENTS, MergeError}; + use super::{CallToPlayStore, MAX_EVENTS, MergeError, delivery_resync_reason}; const TEST_NOW: i64 = 8_000_000_000_000; @@ -768,6 +804,24 @@ mod tests { assert!(store.snapshot_at(TEST_NOW).contains(&message)); } + #[test] + fn only_failed_or_incomplete_deliveries_request_resync() { + assert!(delivery_resync_reason(Err(())).is_some()); + assert!(delivery_resync_reason(Ok(&CallToPlayAck::NeedHandshake)).is_some()); + assert!(delivery_resync_reason(Ok(&CallToPlayAck::NeedHistory)).is_some()); + + for ack in [ + CallToPlayAck::Applied, + CallToPlayAck::Duplicate, + CallToPlayAck::Obsolete, + CallToPlayAck::Rejected { + reason: "invalid".to_string(), + }, + ] { + assert!(delivery_resync_reason(Ok(&ack)).is_none()); + } + } + fn full_active_store() -> CallToPlayStore { let mut history = Vec::with_capacity(MAX_EVENTS); history.push(create_event("create")); diff --git a/crates/lanspread-peer/src/network.rs b/crates/lanspread-peer/src/network.rs index c491fa1..0727f44 100644 --- a/crates/lanspread-peer/src/network.rs +++ b/crates/lanspread-peer/src/network.rs @@ -9,7 +9,16 @@ use bytes::BytesMut; use futures::{SinkExt, StreamExt}; use if_addrs::{IfAddr, Interface, get_if_addrs}; use lanspread_db::db::GameFileDescription; -use lanspread_proto::{CallToPlayEvent, Hello, HelloAck, LibraryDelta, Message, Request, Response}; +use lanspread_proto::{ + CallToPlayAck, + CallToPlayEvent, + Hello, + HelloAck, + LibraryDelta, + Message, + Request, + Response, +}; use s2n_quic::{ Client as QuicClient, Connection, @@ -132,6 +141,14 @@ pub async fn send_oneway_request(peer_addr: SocketAddr, request: Request) -> eyr /// Performs a hello/ack handshake with a peer. pub async fn exchange_hello(peer_addr: SocketAddr, hello: Hello) -> eyre::Result { + let response = exchange_request(peer_addr, Request::Hello(hello)).await?; + match response { + Response::HelloAck(ack) => Ok(ack), + other => eyre::bail!("Unexpected response from peer {peer_addr}: {other:?}"), + } +} + +async fn exchange_request(peer_addr: SocketAddr, request: Request) -> eyre::Result { let mut conn = connect_to_peer(peer_addr).await?; let stream = conn.open_bidirectional_stream().await?; @@ -139,19 +156,15 @@ pub async fn exchange_hello(peer_addr: SocketAddr, hello: Hello) -> eyre::Result let mut framed_rx = FramedRead::new(rx, LengthDelimitedCodec::new()); let mut framed_tx = FramedWrite::new(tx, LengthDelimitedCodec::new()); - framed_tx.send(Request::Hello(hello).encode()).await?; - let _ = framed_tx.close().await; + framed_tx.send(request.encode()).await?; + framed_tx.close().await?; let mut data = BytesMut::new(); - while let Some(Ok(bytes)) = framed_rx.next().await { - data.extend_from_slice(&bytes); + while let Some(frame) = framed_rx.next().await { + data.extend_from_slice(&frame?); } - let response = Response::decode(data.freeze()); - match response { - Response::HelloAck(ack) => Ok(ack), - other => eyre::bail!("Unexpected response from peer {peer_addr}: {other:?}"), - } + Ok(Response::decode(data.freeze())) } pub async fn send_library_delta( @@ -177,15 +190,19 @@ pub async fn send_call_to_play_events( peer_addr: SocketAddr, peer_id: &str, events: Vec, -) -> eyre::Result<()> { - send_oneway_request( +) -> eyre::Result { + let response = exchange_request( peer_addr, Request::CallToPlayEvents { peer_id: peer_id.to_string(), events, }, ) - .await + .await?; + match response { + Response::CallToPlayAck(ack) => Ok(ack), + other => eyre::bail!("Unexpected Call to Play response from peer {peer_addr}: {other:?}"), + } } /// Requests game file details from a peer. diff --git a/crates/lanspread-peer/src/services/stream.rs b/crates/lanspread-peer/src/services/stream.rs index 5c2d73d..240cf4f 100644 --- a/crates/lanspread-peer/src/services/stream.rs +++ b/crates/lanspread-peer/src/services/stream.rs @@ -4,7 +4,7 @@ use std::net::SocketAddr; use futures::{SinkExt, StreamExt}; use lanspread_db::db::{Game, GameFileDescription}; -use lanspread_proto::{LibraryDelta, Message, Request, Response}; +use lanspread_proto::{CallToPlayAck, LibraryDelta, Message, Request, Response}; use s2n_quic::stream::{BidirectionalStream, SendStream}; use tokio_util::codec::{FramedRead, FramedWrite, LengthDelimitedCodec}; @@ -94,8 +94,8 @@ async fn dispatch_request( peer_id, events: incoming, } => { - handle_call_to_play_events(ctx, remote_addr, &peer_id, incoming).await; - framed_tx + let ack = handle_call_to_play_events(ctx, &peer_id, incoming).await; + send_response(framed_tx, Response::CallToPlayAck(ack), "CallToPlayAck").await } Request::GetGame { id } => handle_get_game(ctx, id, framed_tx).await, Request::GetGameFileData(desc) => handle_file_data_request(ctx, desc, framed_tx).await, @@ -123,31 +123,35 @@ async fn dispatch_request( async fn handle_call_to_play_events( ctx: &PeerCtx, - remote_addr: Option, peer_id: &str, incoming: Vec, -) { +) -> CallToPlayAck { let peer_id = peer_id.to_string(); - let sender_matches = if let Some(remote_addr) = remote_addr { - ctx.peer_game_db - .read() - .await - .peer_addr(&peer_id) - .is_some_and(|listen_addr| listen_addr.ip() == remote_addr.ip()) - } else { - false - }; - if !sender_matches { - log::warn!("Ignoring Call to Play events from unverified peer {peer_id}"); - return; + if ctx.peer_game_db.read().await.peer_addr(&peer_id).is_none() { + log::debug!("Requesting a handshake before accepting Call to Play events from {peer_id}"); + return CallToPlayAck::NeedHandshake; } if incoming.iter().any(|event| event.actor_id != peer_id) { - log::warn!("Ignoring Call to Play events with an actor that does not match {peer_id}"); - return; + let reason = format!("event actor does not match envelope peer {peer_id}"); + log::warn!("Rejecting Call to Play events: {reason}"); + return CallToPlayAck::Rejected { reason }; } match ctx.call_to_play.write().await.merge_batch(incoming) { Ok(merged) => { + let ack = if merged.needs_history() { + CallToPlayAck::NeedHistory + } else if !merged.applied.is_empty() { + CallToPlayAck::Applied + } else if merged.obsolete > 0 { + CallToPlayAck::Obsolete + } else if merged.duplicates > 0 { + CallToPlayAck::Duplicate + } else { + CallToPlayAck::Rejected { + reason: "empty Call to Play event batch".to_string(), + } + }; if merged.needs_history() { log::warn!( "Ignoring Call to Play actions without history from {peer_id}: {}", @@ -160,9 +164,13 @@ async fn handle_call_to_play_events( crate::PeerEvent::CallToPlayEvents(merged.applied), ); } + ack } Err(err) => { log::warn!("Rejecting Call to Play events from {peer_id}: {err}"); + CallToPlayAck::Rejected { + reason: err.to_string(), + } } } } @@ -503,6 +511,7 @@ mod tests { }; use lanspread_db::db::GameCatalog; + use lanspread_proto::{CallToPlayAction, CallToPlayEvent}; use tokio::sync::{RwLock, mpsc}; use tokio_util::{sync::CancellationToken, task::TaskTracker}; @@ -548,6 +557,29 @@ mod tests { .to_peer_ctx(tx_notify_ui) } + fn call_to_play_event(actor_id: &str, action: CallToPlayAction) -> CallToPlayEvent { + CallToPlayEvent { + id: "event-1".to_string(), + call_id: "call-1".to_string(), + actor_id: actor_id.to_string(), + actor_name: "Alice".to_string(), + at: 8_000_000_000_000, + action, + } + } + + fn call_to_play_create(actor_id: &str) -> CallToPlayEvent { + call_to_play_event( + actor_id, + CallToPlayAction::Create { + game_id: "game".to_string(), + max_players: 4, + scheduled_for: None, + deadline: 8_000_000_060_000, + }, + ) + } + #[test] fn local_relative_paths_are_never_transferable() { assert!(path_points_inside_local("game", "game/local/save.dat")); @@ -570,6 +602,84 @@ mod tests { )); } + #[tokio::test] + async fn known_peer_id_accepts_live_events_without_transport_ip_matching() { + let temp = TempDir::new("lanspread-call-to-play-known-peer"); + let ctx = test_ctx(temp.path().to_path_buf(), GameCatalog::empty()); + ctx.peer_game_db.write().await.upsert_peer( + "peer-alice".to_string(), + SocketAddr::from(([10, 66, 0, 2], 40000)), + ); + + let ack = + handle_call_to_play_events(&ctx, "peer-alice", vec![call_to_play_create("peer-alice")]) + .await; + + assert_eq!(ack, CallToPlayAck::Applied); + assert_eq!(ctx.call_to_play.write().await.snapshot().len(), 1); + } + + #[tokio::test] + async fn unknown_peer_and_mismatched_actor_receive_explicit_acks() { + let temp = TempDir::new("lanspread-call-to-play-identity"); + let ctx = test_ctx(temp.path().to_path_buf(), GameCatalog::empty()); + + assert_eq!( + handle_call_to_play_events( + &ctx, + "peer-alice", + vec![call_to_play_create("peer-alice")], + ) + .await, + CallToPlayAck::NeedHandshake + ); + + ctx.peer_game_db.write().await.upsert_peer( + "peer-alice".to_string(), + SocketAddr::from(([10, 66, 0, 2], 40000)), + ); + assert!(matches!( + handle_call_to_play_events( + &ctx, + "peer-alice", + vec![call_to_play_create("peer-mallory")], + ) + .await, + CallToPlayAck::Rejected { reason } + if reason.contains("does not match envelope peer") + )); + } + + #[tokio::test] + async fn live_event_ack_reports_missing_history_and_duplicates() { + let temp = TempDir::new("lanspread-call-to-play-outcomes"); + let ctx = test_ctx(temp.path().to_path_buf(), GameCatalog::empty()); + ctx.peer_game_db.write().await.upsert_peer( + "peer-alice".to_string(), + SocketAddr::from(([10, 66, 0, 2], 40000)), + ); + let orphan = call_to_play_event( + "peer-alice", + CallToPlayAction::AddTime { + deadline: 8_000_000_600_000, + }, + ); + assert_eq!( + handle_call_to_play_events(&ctx, "peer-alice", vec![orphan]).await, + CallToPlayAck::NeedHistory + ); + + let create = call_to_play_create("peer-alice"); + assert_eq!( + handle_call_to_play_events(&ctx, "peer-alice", vec![create.clone()]).await, + CallToPlayAck::Applied + ); + assert_eq!( + handle_call_to_play_events(&ctx, "peer-alice", vec![create]).await, + CallToPlayAck::Duplicate + ); + } + #[tokio::test] async fn get_game_response_respects_serve_gates() { let temp = TempDir::new("lanspread-stream"); diff --git a/crates/lanspread-proto/src/lib.rs b/crates/lanspread-proto/src/lib.rs index 412a40f..040698d 100644 --- a/crates/lanspread-proto/src/lib.rs +++ b/crates/lanspread-proto/src/lib.rs @@ -4,7 +4,7 @@ use bytes::Bytes; use lanspread_db::db::{Game, GameFileDescription}; use serde::{Deserialize, Serialize}; -pub const PROTOCOL_VERSION: u32 = 6; +pub const PROTOCOL_VERSION: u32 = 7; pub use lanspread_db::db::Availability; @@ -74,6 +74,16 @@ pub enum CallToPlayAction { }, } +#[derive(Clone, Debug, Serialize, Deserialize, PartialEq, Eq)] +pub enum CallToPlayAck { + Applied, + Duplicate, + NeedHandshake, + NeedHistory, + Obsolete, + Rejected { reason: String }, +} + #[derive(Clone, Debug, Serialize, Deserialize)] pub struct LibrarySnapshot { pub library_rev: u64, @@ -130,6 +140,7 @@ pub enum Response { file_descriptions: Vec, }, HelloAck(HelloAck), + CallToPlayAck(CallToPlayAck), GameNotFound(String), InvalidRequest(Bytes, String), EncodingError(String),