use helix_core::{Effect, EffectSink};
use crate::error::ImError;
use crate::state::{ChannelId, Seq};
use super::super::{ImWsContext, WsFrame, WsHandlerRegistration, WsMessageHandler};
const CHANNEL_MEMBER_UPDATE_ACTION: &str = "channel_member_update";
struct ChannelMemberUpdateHandler;
impl WsMessageHandler for ChannelMemberUpdateHandler {
fn action(&self) -> &'static str {
CHANNEL_MEMBER_UPDATE_ACTION
}
fn handle(
&self,
ctx: &mut ImWsContext<'_>,
frame: &WsFrame,
out: &mut EffectSink,
) -> Result<(), ImError> {
let Ok(data) = frame.data_required() else {
return Ok(());
};
let Some(channel_id) = data
.get("id")
.and_then(serde_json::Value::as_str)
.and_then(ChannelId::from_str)
else {
return Ok(());
};
let event_seq = frame.event_seq();
if event_seq.is_some_and(|seq| {
ctx.state
.committed_member_update_seqs
.contains(&(channel_id, seq))
|| ctx
.state
.inflight_member_update_seqs
.contains(&(channel_id, seq))
}) {
return Ok(());
}
let viewer_rejoined = member_change_joins_user(data, ctx.auth_user_id);
let mut ops = channel_full_and_member_ops(ctx, channel_id, data);
if viewer_rejoined {
ops.push(crate::acl::to_effect_s1::channel_set_cols_op(
channel_id,
vec![("is_remove", helix_core::effect::SqlValue::Integer(0))],
));
}
if ops.is_empty() {
return Ok(());
}
if let Some(seq) = event_seq {
ctx.state
.inflight_member_update_seqs
.insert((channel_id, seq));
}
let corr = ctx.alloc_corr();
ctx.state.corr_map.insert(
corr,
crate::state::CorrelationContext::ChannelMemberUpdatePersist {
channel_id,
event_seq,
viewer_rejoined,
causation_id: event_seq.map(|seq| format!("channel-member-update:{}", seq.0)),
},
);
out.push(Effect::PersistAtomic { corr, ops });
Ok(())
}
}
fn member_change_joins_user(data: &serde_json::Value, user_id: &str) -> bool {
!user_id.is_empty()
&& data
.get("memberChange")
.and_then(|change| change.get("join"))
.and_then(serde_json::Value::as_array)
.is_some_and(|join| {
join.iter().any(|member| {
member
.get("id")
.or_else(|| member.get("userId"))
.and_then(serde_json::Value::as_str)
== Some(user_id)
})
})
}
pub(crate) fn apply_channel_full_and_members(
ctx: &mut ImWsContext<'_>,
channel_id: ChannelId,
data: &serde_json::Value,
out: &mut EffectSink,
) {
let ops = channel_full_and_member_ops(ctx, channel_id, data);
if !ops.is_empty() {
out.push(Effect::PersistFire { ops });
}
}
pub(crate) fn queue_channel_create_persist(
ctx: &mut ImWsContext<'_>,
channel_id: ChannelId,
mut channel: serde_json::Value,
causation_id: Option<String>,
out: &mut EffectSink,
) {
if ctx.state.committed_channel_creates.contains(&channel_id)
|| ctx.state.inflight_channel_creates.contains(&channel_id)
{
return;
}
let member_rows = crate::channel_write::collect_members(&channel)
.into_iter()
.map(|member| {
serde_json::json!({
"channel_id": channel_id.as_str(),
"user_id": member.user_id,
"team_id": member.team_id,
"role": member.role,
"nick_name": member.nick_name,
})
})
.collect::<Vec<_>>();
if let Some(object) = channel.as_object_mut() {
object.insert(
"memberCount".to_string(),
serde_json::Value::from(member_rows.len() as u64),
);
}
let ops = channel_full_and_member_ops(ctx, channel_id, &channel);
if ops.is_empty() {
return;
}
let corr = ctx.alloc_corr();
ctx.state.inflight_channel_creates.insert(channel_id);
ctx.state.corr_map.insert(
corr,
crate::state::CorrelationContext::ChannelCreatePersist {
channel_id,
channel: Box::new(channel),
member_rows,
causation_id,
},
);
out.push(Effect::PersistAtomic { corr, ops });
}
fn channel_full_and_member_ops(
ctx: &mut ImWsContext<'_>,
channel_id: ChannelId,
data: &serde_json::Value,
) -> Vec<helix_core::effect::StorageOp> {
ctx.state
.channels
.entry(channel_id)
.or_insert_with(|| crate::channel::Channel::new(channel_id, 0));
if let Some(last_event_seq) = data.get("lastEventSeq").and_then(serde_json::Value::as_u64) {
let target = ctx
.state
.increment_target
.entry(channel_id)
.or_insert(Seq(0));
if Seq(last_event_seq) > *target {
*target = Seq(last_event_seq);
}
}
let mut ops = Vec::new();
if let Some((cols, exclude)) =
crate::channel_write::program_full(data, ctx.auth_user_id, ctx.now_ms)
{
ops.push(crate::acl::to_effect::upsert_channel_full_op(cols, exclude));
}
let members = crate::channel_write::collect_members(data);
if let Some(op) = crate::acl::to_effect_s1::upsert_channel_members_op(channel_id, members) {
ops.push(op);
}
let leaves = crate::channel_write::collect_member_leaves(data);
if let Some(op) = crate::acl::to_effect_s1::delete_channel_members_op(channel_id, leaves) {
ops.push(op);
}
ops
}
static CHANNEL_MEMBER_UPDATE_HANDLER: ChannelMemberUpdateHandler = ChannelMemberUpdateHandler;
#[cfg(target_arch = "wasm32")]
pub(super) fn inventory_link_anchor() {
std::hint::black_box(&CHANNEL_MEMBER_UPDATE_HANDLER);
}
inventory::submit! {
WsHandlerRegistration {
action: CHANNEL_MEMBER_UPDATE_ACTION,
handler: &CHANNEL_MEMBER_UPDATE_HANDLER,
}
}