fix(call-to-play): acknowledge live replication
Raise the wire protocol to version 7 and add explicit Call to Play delivery outcomes. Live requests now wait for an application acknowledgement, allowing the sender to distinguish applied, duplicate, obsolete, incomplete, and rejected updates instead of treating a successful write as acceptance. Remove source-IP equality from actor verification. The receiver now requires the envelope peer ID to be present in its known roster and requires every live event actor to match that envelope. This matches the cooperative-LAN trust model without misrepresenting the shared TLS identity as per-peer authentication. Transport failures, malformed responses, NeedHandshake, and NeedHistory each trigger one asynchronous Hello/HelloAck resync. Rejections are logged without retry, and local publication remains independent of remote availability. Test Plan: - `just fmt` -- passed - `just clippy` -- passed - `just test` -- passed - `git diff --cached --check` -- passed
This commit is contained in:
5 files changed
+254
-57
No files matched your search
@@ -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"));
|
||||
|
||||
Reference in new issue
Block a user