use helix_core::effect::SqlValue;
use helix_core::{Effect, EffectSink};
use crate::error::ImError;
use crate::state::ChannelId;
use super::super::{ImWsContext, WsFrame, WsHandlerRegistration, WsMessageHandler};
const CHANNEL_CLOSE_ACTION: &str = "channel_close";
struct ChannelCloseHandler;
impl WsMessageHandler for ChannelCloseHandler {
fn action(&self) -> &'static str {
CHANNEL_CLOSE_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("channelId")
.or_else(|| data.get("channel_id"))
.and_then(serde_json::Value::as_str)
.and_then(ChannelId::from_str)
else {
return Ok(());
};
let delete_at = data
.get("deleteAt")
.or_else(|| data.get("delete_at"))
.and_then(serde_json::Value::as_i64)
.unwrap_or(0);
let is_viewer_left = data
.get("terminalReason")
.or_else(|| data.get("terminal_reason"))
.and_then(serde_json::Value::as_str)
== Some("left");
let terminal_seq = data
.get("eventSeq")
.or_else(|| data.get("event_seq"))
.and_then(serde_json::Value::as_u64)
.filter(|seq| *seq > 0 && *seq <= i64::MAX as u64);
let terminal_type = data
.get("eventType")
.or_else(|| data.get("event_type"))
.and_then(serde_json::Value::as_u64);
if !is_viewer_left && terminal_type == Some(7) {
if let Some(terminal_seq) = terminal_seq {
let event_id = data
.get("id")
.and_then(serde_json::Value::as_str)
.filter(|value| !value.is_empty())
.map(str::to_string);
let event = crate::sync_session::EventEnvelope::new(
channel_id,
crate::state::Seq(terminal_seq),
crate::sync_session::EventKind::ChannelTerminalClosed,
crate::sync_session::PostFields::default(),
)
.with_event_identity(event_id, None, delete_at.max(0), String::new())
.with_viewer_user_id(ctx.auth_user_id);
let admitted = ctx
.state
.channels
.entry(channel_id)
.or_insert_with(|| crate::channel::Channel::new(channel_id, 0))
.admit_message_v3_post(event, out)?;
if let Some(event) = admitted {
let corr = ctx.alloc_corr();
super::channel_stream_event::queue_stream_commit(ctx.state, corr, event, out);
}
return Ok(());
}
}
let (transition, cols) = if is_viewer_left {
(
crate::state::ChannelLifecycleTransition::Left,
vec![("is_remove", SqlValue::Integer(1))],
)
} else {
(
crate::state::ChannelLifecycleTransition::Closed { delete_at },
vec![
("delete_at", SqlValue::Integer(delete_at)),
("is_active", SqlValue::Integer(0)),
],
)
};
let corr = ctx.alloc_corr();
ctx.state.corr_map.insert(
corr,
crate::state::CorrelationContext::ChannelLifecyclePersist {
channel_id,
transition,
causation_id: None,
},
);
let mut ops = vec![crate::acl::to_effect_s1::channel_set_cols_op(
channel_id, cols,
)];
if is_viewer_left && !ctx.auth_user_id.is_empty() {
if let Some(delete_member) = crate::acl::to_effect_s1::delete_channel_members_op(
channel_id,
vec![ctx.auth_user_id.to_string()],
) {
ops.push(delete_member);
}
}
out.push(Effect::PersistAtomic { corr, ops });
Ok(())
}
}
static CHANNEL_CLOSE_HANDLER: ChannelCloseHandler = ChannelCloseHandler;
#[cfg(target_arch = "wasm32")]
pub(super) fn inventory_link_anchor() {
std::hint::black_box(&CHANNEL_CLOSE_HANDLER);
}
inventory::submit! {
WsHandlerRegistration {
action: CHANNEL_CLOSE_ACTION,
handler: &CHANNEL_CLOSE_HANDLER,
}
}