use helix_core::{Effect, EffectSink};
use serde_json::Value;
use std::collections::HashSet;
use crate::error::ImError;
use crate::state::ChannelId;
use crate::ws::parser::extract_post_fields;
use super::super::{ImWsContext, WsFrame, WsHandlerRegistration, WsMessageHandler};
const POST_READ_ACTION: &str = "post_read";
struct PostReadHandler;
impl WsMessageHandler for PostReadHandler {
fn action(&self) -> &'static str {
POST_READ_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(post_id) = data
.get("postId")
.or_else(|| data.get("post_id"))
.and_then(serde_json::Value::as_str)
.filter(|s| !s.is_empty())
.map(str::to_string)
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(());
};
if !has_receipt_bit_fields(data) || !has_receipt_member_fields(data) {
return Ok(());
}
ctx.state.invalidate_recent_message_coverage(channel_id);
let mut fields = extract_post_fields(data);
if fields.read_bits.is_empty() {
if let Some(rm) = data
.get("readMap")
.or_else(|| data.get("read_map"))
.and_then(serde_json::Value::as_str)
{
fields.read_bits = rm.to_string();
}
}
let receipt_revision = data
.get("updateAt")
.or_else(|| data.get("update_at"))
.and_then(serde_json::Value::as_i64)
.unwrap_or_default();
if let Some((committed_revision, committed_bits)) =
ctx.state.committed_post_reads.get(post_id.as_str())
{
if (*committed_revision == receipt_revision && committed_bits == &fields.read_bits)
|| (receipt_revision > 0
&& *committed_revision > 0
&& receipt_revision < *committed_revision)
{
return Ok(());
}
}
let Some(reader_id) = data
.get("readerUserId")
.or_else(|| data.get("reader_user_id"))
.or_else(|| data.get("userId"))
.or_else(|| data.get("user_id"))
.and_then(serde_json::Value::as_str)
.filter(|value| !value.trim().is_empty())
else {
return Ok(());
};
let Some(author_user_id) = data
.get("authorUserId")
.or_else(|| data.get("author_user_id"))
.and_then(serde_json::Value::as_str)
.filter(|value| !value.trim().is_empty())
else {
return Ok(());
};
let Some(snapshot_id) = data
.get("snapshotId")
.or_else(|| data.get("snapshot_id"))
.and_then(serde_json::Value::as_str)
.filter(|value| !value.trim().is_empty())
else {
return Ok(());
};
let member_ids: Vec<String> = data
.get("memberIds")
.or_else(|| data.get("member_ids"))
.and_then(serde_json::Value::as_array)
.map(|values| {
values
.iter()
.filter_map(serde_json::Value::as_str)
.map(str::to_string)
.collect()
})
.unwrap_or_default();
let mut member_set = HashSet::with_capacity(member_ids.len());
let valid_receipt_shape = !member_ids.is_empty()
&& member_ids
.iter()
.all(|member_id| !member_id.is_empty() && member_id.trim() == member_id)
&& member_ids
.iter()
.all(|member_id| member_set.insert(member_id))
&& fields.read_bits.len() == member_ids.len()
&& fields
.read_bits
.bytes()
.all(|bit| bit == b'0' || bit == b'1');
if !valid_receipt_shape {
return Ok(());
}
let terminal_event = match crate::event::read::post_for_viewer(
channel_id.as_str(),
post_id.as_str(),
author_user_id,
reader_id,
snapshot_id,
&member_ids,
fields.read_bits.as_str(),
receipt_revision,
ctx.auth_user_id,
) {
Ok(event) => event.into_bytes(),
Err(_) => return Ok(()),
};
let read_op = crate::channel::apply_read_op(&post_id, fields.read_bits.as_str());
let corr = ctx.alloc_corr();
ctx.state.corr_map.insert(
corr,
crate::state::CorrelationContext::PostReadPersist {
channel_id,
message_id: post_id,
receipt_revision,
read_bits: fields.read_bits,
terminal_event,
},
);
out.push(Effect::Persist {
corr,
ops: vec![read_op],
});
Ok(())
}
}
fn has_receipt_bit_fields(data: &Value) -> bool {
["readMap", "read_map", "readBits", "read_bits"]
.iter()
.any(|key| data.get(*key).is_some())
}
fn has_receipt_member_fields(data: &Value) -> bool {
let has_read_projection = ["readMap", "read_map", "readBits", "read_bits"]
.iter()
.any(|key| data.get(*key).is_some());
let has_member_ids = ["memberIds", "member_ids"]
.iter()
.any(|key| data.get(*key).is_some());
has_read_projection && has_member_ids
}
static POST_READ_HANDLER: PostReadHandler = PostReadHandler;
#[cfg(target_arch = "wasm32")]
pub(super) fn inventory_link_anchor() {
std::hint::black_box(&POST_READ_HANDLER);
}
inventory::submit! {
WsHandlerRegistration {
action: POST_READ_ACTION,
handler: &POST_READ_HANDLER,
}
}