use helix_core::{Correlation, Effect, EffectSink};
use crate::channel::Channel;
use crate::error::ImError;
use crate::state::Seq;
use crate::sync_session::IncrementChannel;
use super::super::{ImWsContext, WsFrame, WsHandlerRegistration, WsMessageHandler};
const INCREMENT_CHANNEL_ACTION: &str = "increment_channel";
pub(crate) fn apply_increment(
ctx: &mut ImWsContext<'_>,
inc: &IncrementChannel,
out: &mut EffectSink,
) {
ctx.state.channel_sync_batch_pending = true;
if let Ok(data) = serde_json::from_slice::<serde_json::Value>(inc.raw.as_ref()) {
if let Some(batch_id) = data
.get("batchId")
.and_then(serde_json::Value::as_str)
.filter(|value| !value.is_empty())
{
ctx.state.pending_increment_batch_id = Some(batch_id.to_string());
}
}
for effect in apply_increment_effects(ctx, inc) {
match effect {
Effect::PersistFire { ops } => ctx.state.pending_increment_ops.extend(ops),
other => out.push(other),
}
}
ctx.state
.pending_increment_projections
.push((inc.channel_id, inc.raw.as_ref().to_vec()));
}
pub(crate) fn apply_increment_hydration(
ctx: &mut ImWsContext<'_>,
inc: &IncrementChannel,
corr: Correlation,
out: &mut EffectSink,
) -> bool {
let mut persist_ops = Vec::new();
for effect in apply_increment_effects(ctx, inc) {
match effect {
Effect::PersistFire { ops } => persist_ops.extend(ops),
other => out.push(other),
}
}
let has_persist = !persist_ops.is_empty();
if has_persist {
out.push(Effect::PersistAtomic {
corr,
ops: persist_ops,
});
}
has_persist
}
fn apply_increment_effects(ctx: &mut ImWsContext<'_>, inc: &IncrementChannel) -> Vec<Effect> {
let mut effects = Vec::with_capacity(4);
let channel_id = inc.channel_id;
ctx.state
.channels
.entry(channel_id)
.or_insert_with(|| Channel::new(channel_id, 0));
if is_terminal_projection(inc.raw.as_ref()) {
if let Some(channel) = ctx.state.channels.get_mut(&channel_id) {
channel
.mark_projection_terminal(Seq(inc.last_event_seq.0.max(channel.cursor.value().0)));
}
}
if ctx.state.increment_fetched.insert(channel_id) {
ctx.state.increment_order.push(channel_id);
}
if !inc.need_sync {
ctx.state.need_sync_skip.insert(channel_id);
}
let target = ctx
.state
.increment_target
.entry(channel_id)
.or_insert(Seq(0));
if inc.last_event_seq > *target {
*target = inc.last_event_seq;
}
if let Ok(data) = serde_json::from_slice::<serde_json::Value>(inc.raw.as_ref()) {
if let Some((cols, exclude)) =
crate::channel_write::program_full(&data, ctx.auth_user_id, ctx.now_ms)
{
effects.push(crate::acl::to_effect::upsert_channel_full(cols, exclude));
}
let members = crate::channel_write::collect_members(&data);
if let Some(eff) = crate::acl::to_effect_s1::upsert_channel_members(channel_id, members) {
effects.push(eff);
}
let leaves = crate::channel_write::collect_member_leaves(&data);
if let Some(eff) = crate::acl::to_effect_s1::delete_channel_members(channel_id, leaves) {
effects.push(eff);
}
}
if let Ok(data) = serde_json::from_slice::<serde_json::Value>(inc.raw.as_ref()) {
ctx.state
.about_me_post_ids
.extend(crate::todo::collect_about_me_ids(&data));
}
effects
}
fn is_terminal_projection(raw: &[u8]) -> bool {
let Ok(data) = serde_json::from_slice::<serde_json::Value>(raw) else {
return false;
};
let delete_at = data
.get("deleteAt")
.or_else(|| data.get("delete_at"))
.and_then(projection_i64)
.unwrap_or(0);
let is_remove = data
.get("isRemove")
.or_else(|| data.get("is_remove"))
.and_then(projection_bool)
.unwrap_or(false);
delete_at > 0 || is_remove
}
fn projection_i64(value: &serde_json::Value) -> Option<i64> {
value
.as_i64()
.or_else(|| value.as_u64().and_then(|v| i64::try_from(v).ok()))
}
fn projection_bool(value: &serde_json::Value) -> Option<bool> {
value
.as_bool()
.or_else(|| value.as_i64().map(|v| v != 0))
.or_else(|| {
value.as_str().and_then(|v| match v {
"1" | "true" | "TRUE" => Some(true),
"0" | "false" | "FALSE" => Some(false),
_ => None,
})
})
}
struct IncrementChannelHandler;
impl WsMessageHandler for IncrementChannelHandler {
fn action(&self) -> &'static str {
INCREMENT_CHANNEL_ACTION
}
fn handle(
&self,
ctx: &mut ImWsContext<'_>,
frame: &WsFrame,
out: &mut EffectSink,
) -> Result<(), ImError> {
let Ok(data) = frame.data_required() else {
tracing::warn!(
"channel-sync-ready 未触发:increment_channel 帧缺少 data,未进入批次缓冲"
);
return Ok(());
};
let Some(inc) = crate::ws::parser::parse_increment_channel(data) else {
tracing::warn!(
"channel-sync-ready 未触发:increment_channel 帧解析失败,未进入批次缓冲"
);
return Ok(());
};
apply_increment(ctx, &inc, out);
Ok(())
}
}
static INCREMENT_CHANNEL_HANDLER: IncrementChannelHandler = IncrementChannelHandler;
#[cfg(target_arch = "wasm32")]
pub(super) fn inventory_link_anchor() {
std::hint::black_box(&INCREMENT_CHANNEL_HANDLER);
}
inventory::submit! {
WsHandlerRegistration {
action: INCREMENT_CHANNEL_ACTION,
handler: &INCREMENT_CHANNEL_HANDLER,
}
}