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 unread_authority = data
.get("unreadCount")
.or_else(|| data.get("unread_count"))
.and_then(|value| {
value
.as_i64()
.or_else(|| value.as_u64().and_then(|value| i64::try_from(value).ok()))
});
let event_seq = i64::try_from(inc.last_event_seq.0).ok();
if let (Some(unread_count), Some(event_seq)) = (unread_authority, event_seq) {
if let Some((mut row, _)) = crate::channel_update::member_channel_from_update_channel(
&data,
channel_id,
ctx.auth_user_id,
ctx.now_ms,
)
.filter(|_| !ctx.auth_user_id.is_empty())
{
if let Some((_, user_id)) = row.iter_mut().find(|(column, _)| column == "user_id") {
*user_id = helix_core::effect::SqlValue::Text(ctx.auth_user_id.to_string());
}
row.push((
"last_unread_event_seq".to_string(),
helix_core::effect::SqlValue::Integer(event_seq),
));
let seed =
helix_core::effect::StorageOp::BatchUpsert(helix_core::effect::UpsertSpec {
table: "channel_member",
rows: vec![row],
conflict_key: Some("channel_id,user_id"),
exclude_from_update: vec!["unread_count", "last_unread_event_seq"],
});
let guarded_set = helix_core::effect::StorageOp::ScopedGuardedBump(
helix_core::effect::ScopedGuardedBumpSpec {
table: "channel_member",
scope_col: "channel_id",
scope_val: helix_core::effect::SqlValue::Text(
channel_id.as_str().to_string(),
),
key_col: "user_id",
key_val: helix_core::effect::SqlValue::Text(ctx.auth_user_id.to_string()),
bump_col: "last_post_at",
bump_delta: 0,
set_cols: vec![
(
"unread_count".to_string(),
helix_core::effect::SqlValue::Integer(unread_count),
),
(
"last_unread_event_seq".to_string(),
helix_core::effect::SqlValue::Integer(event_seq),
),
],
guard_col: "last_unread_event_seq",
guard_val: event_seq,
},
);
effects.push(Effect::PersistFire {
ops: vec![seed, guarded_set],
});
}
}
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,
}
}