use helix_core::effect::StorageOp;
use helix_core::{Effect, EffectSink};
use crate::error::ImError;
use crate::state::{ChannelId, Seq};
use crate::sync_session::{EventEnvelope, EventKind, PostFields};
use super::super::{ImWsContext, WsFrame};
const GAP_IMMEDIATE_RESYNC_THRESHOLD: u64 = 2;
pub(super) fn recent_effects_persist_message(effects: &[Effect]) -> bool {
effects.iter().any(|effect| {
let Effect::PersistFire { ops } = effect else {
return false;
};
ops.iter().any(|op| match op {
StorageOp::BatchUpsert(spec) => spec.table == "message",
StorageOp::BatchUpdate(spec) => spec.table == "message",
StorageOp::MonotonicUpsert(spec) => spec.table == "message",
StorageOp::GuardedBump(spec) => spec.table == "message",
StorageOp::ScopedGuardedBump(spec) => spec.table == "message",
StorageOp::Get(spec) => spec.table == "message",
StorageOp::ScopedGet(spec) => spec.table == "message",
StorageOp::Scan(spec) => spec.table == "message",
StorageOp::BatchDelete(spec) => spec.table == "message",
})
})
}
pub(super) fn trigger_backfill_if_large_gap(
ctx: &mut ImWsContext<'_>,
channel_id: ChannelId,
seq: Seq,
out: &mut EffectSink,
) {
let should_resync = match ctx.state.channels.get(&channel_id) {
Some(ch) => {
let cursor = ch.cursor.value().0;
ch.gate.is_some()
&& ch.inflight_sync.is_none()
&& seq.0 > cursor.saturating_add(GAP_IMMEDIATE_RESYNC_THRESHOLD)
}
None => false,
};
if should_resync {
super::increment_channel_end::emit_proactive_resync_for(ctx, &[channel_id], out);
}
}
fn recent_effects_emitted_sync_too_long(effects: &[Effect], channel_id: ChannelId) -> bool {
effects.iter().any(|effect| {
matches!(effect, Effect::Emit { event } if {
let s = String::from_utf8_lossy(event.0.as_ref());
s.contains("\"event\":\"im:sync:too_long\"")
&& s.contains("\"mode\":\"drop_and_reload\"")
&& s.contains(channel_id.as_str())
})
})
}
fn recent_effects_already_request_latest_post(effects: &[Effect], channel_id: ChannelId) -> bool {
effects.iter().any(|effect| {
let Effect::Http { req, .. } = effect else {
return false;
};
if !req.url.contains("posts/getLatestPost") {
return false;
}
let Some(body) = req.body.as_ref() else {
return false;
};
let Ok(value) = serde_json::from_slice::<serde_json::Value>(body.as_ref()) else {
return false;
};
value["channelId"] == channel_id.as_str()
})
}
pub(super) fn schedule_reload_if_recent_too_long(
ctx: &mut ImWsContext<'_>,
channel_id: ChannelId,
recent_start: usize,
out: &mut EffectSink,
) -> Result<(), ImError> {
let Some(recent) = out.as_slice().get(recent_start..) else {
return Ok(());
};
if !recent_effects_emitted_sync_too_long(recent, channel_id)
|| recent_effects_already_request_latest_post(recent, channel_id)
{
return Ok(());
}
let Some(reset_to) = ctx
.state
.channels
.get(&channel_id)
.map(|ch| Seq(ch.cursor.value().0 + 1))
else {
return Ok(());
};
let connection_id = ctx.state.connection_id.clone();
let corr = ctx.alloc_corr();
let payload = serde_json::json!({
"channel_id": channel_id.as_str(),
"timestamp": 0,
});
let payload_bytes = serde_json::to_vec(&payload).unwrap_or_default();
let effects = crate::commands::handle_outbound(
"im_get_latest_post",
payload_bytes.as_ref(),
ctx.api_base_url,
ctx.api_base_url,
connection_id.as_deref(),
corr,
)?;
ctx.state.corr_map.insert(
corr,
crate::state::CorrelationContext::TooLongReload {
channel_id,
reset_to,
},
);
for effect in effects {
out.push(effect);
}
Ok(())
}
pub(super) fn gate_ingest_content_event(
ctx: &mut ImWsContext<'_>,
frame: &WsFrame,
channel_id: ChannelId,
kind: EventKind,
fields: PostFields,
msg_id: Option<String>,
out: &mut EffectSink,
) -> Result<(), ImError> {
let Some(seq) = frame.event_seq() else {
return Ok(());
};
let ev = EventEnvelope::new(channel_id, seq, kind, fields)
.with_msg_id(msg_id)
.with_viewer_user_id(ctx.auth_user_id);
if let Some(ch) = ctx.state.channels.get_mut(&channel_id) {
let recent_start = out.as_slice().len();
let edit_settle =
if matches!(ev.kind, EventKind::PostEdit) && ev.seq.0 != ch.cursor.value().0 + 1 {
let msg_id_ref = ev
.msg_id
.as_deref()
.filter(|s| !s.is_empty())
.unwrap_or(ev.fields.id.as_str());
Some((
helix_core::Effect::PersistFire {
ops: vec![crate::channel::event_to_storage_op(&ev)],
},
crate::acl::to_effect::emit_post_updated_for_viewer(
channel_id,
ev.seq.0,
msg_id_ref,
&ev.fields,
ev.viewer_user_id.as_str(),
),
))
} else {
None
};
ch.ingest(ev, out, ctx.now_ms)?;
schedule_reload_if_recent_too_long(ctx, channel_id, recent_start, out)?;
if let Some((persist, emit)) = edit_settle {
out.push(persist);
out.push(emit);
}
}
trigger_backfill_if_large_gap(ctx, channel_id, seq, out);
Ok(())
}