use helix_core::EffectSink;
use crate::state::{CorrelationContext, Seq, ServerId, TemporaryId};
use crate::sync_session::EventEnvelope;
use super::super::{ImWsContext, WsFrame};
pub(super) fn is_echo_frame(
ctx: &ImWsContext<'_>,
temporary_id: &str,
ev: &EventEnvelope,
frame: &WsFrame,
) -> bool {
if !temporary_id.is_empty() {
let tmp = TemporaryId(temporary_id.to_string());
if ctx.state.pending_sends.contains_key(&tmp) {
return true;
}
}
if !ctx.auth_user_id.is_empty() {
let sender = if !ev.fields.user_id.is_empty() {
ev.fields.user_id.as_str()
} else {
frame
.data()
.and_then(|d| {
d.get("userSnapshot")
.and_then(|s| s.get("userId"))
.or_else(|| d.get("userId"))
})
.and_then(serde_json::Value::as_str)
.unwrap_or("")
};
if !sender.is_empty() && sender == ctx.auth_user_id {
return true;
}
}
false
}
pub(super) fn reconcile_post_echo(
ctx: &mut ImWsContext<'_>,
temporary_id: &str,
server_id: Option<ServerId>,
out: &mut EffectSink,
) -> Result<(), crate::ImError> {
if temporary_id.is_empty() {
return Ok(());
}
let Some(server_id) = server_id else {
return Ok(());
};
let tmp = TemporaryId(temporary_id.to_string());
if !ctx.state.pending_sends.contains_key(&tmp) {
return Ok(());
}
let terminal_event = ctx
.state
.pending_sends
.get(&tmp)
.and_then(|pending| pending.body.as_ref().map(|body| (pending, body)))
.map(|(pending, body)| {
crate::event::post::sent_from_local_body(
pending.timeline_readback.causation_id.as_deref(),
temporary_id,
server_id.as_str(),
ctx.auth_user_id,
body,
)
.map(|event| event.into_bytes())
})
.transpose()?;
let reconcile_corr = ctx.alloc_corr();
let persist_corr = ctx
.state
.pending_sends
.get(&tmp)
.map(|ps| ps.persist_corr)
.unwrap_or_default();
if let Some(ps) = ctx.state.pending_sends.get_mut(&tmp) {
ps.reconcile(server_id, reconcile_corr, out);
}
ctx.state.pending_sends.remove(&tmp);
if let Some(corr) = persist_corr {
ctx.state.corr_map.remove(&corr);
}
let continuation = match terminal_event {
Some(terminal_event) => CorrelationContext::AuthoritativeSendTerminalPersist {
temporary_id: tmp,
terminal_event: bytes::Bytes::from(terminal_event),
},
None => CorrelationContext::AuthoritativeSendReconcilePersist { temporary_id: tmp },
};
ctx.state.corr_map.insert(reconcile_corr, continuation);
Ok(())
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn out_of_order_self_echo_keeps_its_sequence_for_the_shared_gate() {
let event = EventEnvelope::new(
crate::state::test_channel_id(9),
Seq(22),
crate::sync_session::EventKind::PostUpsert,
crate::sync_session::PostFields {
id: "post-gap-echo".to_string(),
msg_type: "VOTE".to_string(),
props: r#"{"vote":{"state":0,"items":[],"options":[]}}"#.to_string(),
..Default::default()
},
)
.with_msg_id(Some("post-gap-echo".to_string()))
.with_viewer_user_id("viewer-1");
assert_eq!(event.seq, Seq(22));
}
#[test]
fn self_echo_terminal_persist_does_not_schedule_latest_timeline_refresh() {
let temporary_id = TemporaryId("temporary-echo-1".to_string());
let mut state = crate::state::ImState::new();
let mut pending = crate::pending_send::PendingSend::new(
temporary_id.clone(),
helix_core::TimerId::from_raw(91),
None,
);
pending.timeline_readback.causation_id = Some("send-request-echo".to_string());
state.pending_sends.insert(temporary_id.clone(), pending);
let mut next_corr = 700_u64;
let mut alloc_corr = || {
let corr = helix_core::Correlation::from_raw(next_corr);
next_corr += 1;
corr
};
let mut sink = EffectSink::new();
{
let mut ctx = ImWsContext::new(&mut state, 1_000, "", "", &mut alloc_corr);
reconcile_post_echo(
&mut ctx,
temporary_id.0.as_str(),
Some(crate::state::test_server_id(41)),
&mut sink,
)
.expect("echo reconciliation");
}
let terminal_corr = sink
.as_slice()
.iter()
.find_map(|effect| match effect {
helix_core::Effect::Persist { corr, .. } => Some(*corr),
_ => None,
})
.expect("self echo writes the terminal sent fact");
assert_eq!(
state.corr_map.get(&terminal_corr),
Some(&CorrelationContext::AuthoritativeSendReconcilePersist {
temporary_id: temporary_id.clone(),
}),
"普通 WS echo 只完成发送对账,timeline 由 post 权威投影更新"
);
}
#[test]
fn retry_self_echo_does_not_schedule_a_second_latest_query() {
let temporary_id = TemporaryId("temporary-retry-echo".to_string());
let mut state = crate::state::ImState::new();
let mut pending = crate::pending_send::PendingSend::new(
temporary_id.clone(),
helix_core::TimerId::from_raw(92),
None,
);
pending.timeline_readback = crate::pending_send::TimelineReadbackContext {
window_token: Some("retry-echo-window".to_string()),
causation_id: Some("retry-echo-request".to_string()),
};
state.pending_sends.insert(temporary_id.clone(), pending);
let mut next_corr = 210;
let mut alloc_corr = || {
let corr = helix_core::Correlation::from_raw(next_corr);
next_corr += 1;
corr
};
let mut sink = EffectSink::new();
{
let mut ctx = ImWsContext::new(&mut state, 1_000, "", "", &mut alloc_corr);
reconcile_post_echo(
&mut ctx,
temporary_id.0.as_str(),
Some(crate::state::test_server_id(42)),
&mut sink,
)
.expect("echo reconciliation");
}
let terminal_corr = sink
.as_slice()
.iter()
.find_map(|effect| match effect {
helix_core::Effect::Persist { corr, .. } => Some(*corr),
_ => None,
})
.expect("retry self echo writes the terminal sent fact");
assert_eq!(
state.corr_map.get(&terminal_corr),
Some(&CorrelationContext::AuthoritativeSendReconcilePersist {
temporary_id: temporary_id.clone(),
})
);
}
}