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 post_echo_gap_settle(
event: &EventEnvelope,
is_echo: bool,
expected: Seq,
) -> Option<(helix_core::Effect, helix_core::Effect)> {
if !is_echo || event.seq.0 <= expected.0 {
return None;
}
let msg_id_ref = event
.msg_id
.as_deref()
.filter(|value| !value.is_empty())
.unwrap_or(event.fields.id.as_str());
Some((
helix_core::Effect::PersistFire {
ops: vec![crate::channel::event_to_storage_op(event)],
},
crate::acl::to_effect::emit_post_received_for_viewer(
event.channel_id,
event.seq.0,
msg_id_ref,
&event.fields,
event.viewer_user_id.as_str(),
),
))
}
pub(super) fn reconcile_post_echo(
ctx: &mut ImWsContext<'_>,
temporary_id: &str,
server_id: Option<ServerId>,
out: &mut EffectSink,
) {
if temporary_id.is_empty() {
return;
}
let Some(server_id) = server_id else {
return;
};
let tmp = TemporaryId(temporary_id.to_string());
if !ctx.state.pending_sends.contains_key(&tmp) {
return;
}
let reconcile_corr = ctx.alloc_corr();
let persist_corr = ctx
.state
.pending_sends
.get(&tmp)
.map(|ps| ps.persist_corr)
.unwrap_or_default();
let authoritative_readback_timer = ctx
.state
.pending_sends
.get(&tmp)
.and_then(|ps| ps.authoritative_readback_timer);
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);
}
if let Some(timer_id) = authoritative_readback_timer {
out.push(helix_core::Effect::CancelTimer { id: timer_id });
}
ctx.state.corr_map.insert(
reconcile_corr,
CorrelationContext::AuthoritativeSendReconcilePersist { temporary_id: tmp },
);
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn out_of_order_self_echo_settles_projection_without_advancing_cursor() {
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");
let (persist, emit) = post_echo_gap_settle(&event, true, Seq(2))
.expect("本人乱序 echo 必须先结算 message 与 projection");
assert!(
matches!(persist, helix_core::Effect::PersistFire { ops } if ops.iter().any(|op| {
format!("{op:?}").contains("message")
}))
);
assert!(matches!(emit, helix_core::Effect::Emit { event }
if String::from_utf8_lossy(event.0.as_ref()).contains("im:post:received")
&& String::from_utf8_lossy(event.0.as_ref()).contains("post-gap-echo")));
}
#[test]
fn non_echo_or_contiguous_post_does_not_bypass_gate() {
let event = EventEnvelope::new(
crate::state::test_channel_id(10),
Seq(3),
crate::sync_session::EventKind::PostUpsert,
crate::sync_session::PostFields::default(),
);
assert!(post_echo_gap_settle(&event, false, Seq(1)).is_none());
assert!(post_echo_gap_settle(&event, true, Seq(3)).is_none());
}
#[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,
);
}
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,
);
}
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(),
})
);
}
}