use helix_core::effect::DomainEventBytes;
use crate::event_sink::to_bus_envelope;
const BUS_SELF_PREFIX: &str = "im:__";
const IM_PREFIX: &str = "im:";
#[derive(Debug, Clone, PartialEq)]
pub struct GuardedEnvelope {
pub channel: String,
pub envelope: serde_json::Value,
pub channel_ok: bool,
}
pub fn guard_bus_envelope(ev: &DomainEventBytes) -> GuardedEnvelope {
let mut envelope = to_bus_envelope(ev);
let channel = envelope
.get("channel")
.and_then(|c| c.as_str())
.unwrap_or("")
.to_string();
let channel_ok = is_emittable_channel(&channel);
if let Some(data) = envelope
.get_mut("payload")
.and_then(|p| p.get_mut("data"))
.filter(|d| d.is_array())
{
let arr = data.take();
*data = serde_json::json!({ "items": arr });
}
GuardedEnvelope {
channel,
envelope,
channel_ok,
}
}
pub fn is_emittable_channel(channel: &str) -> bool {
channel.starts_with(IM_PREFIX) && !channel.starts_with(BUS_SELF_PREFIX)
}
#[cfg(test)]
mod tests {
use super::*;
use bytes::Bytes;
fn ev(json: serde_json::Value) -> DomainEventBytes {
DomainEventBytes(Bytes::from(serde_json::to_vec(&json).unwrap()))
}
#[test]
fn object_data_passes_through_with_ok_channel() {
let g = guard_bus_envelope(&ev(serde_json::json!({
"event": "im:post:received",
"data": { "channel_id": "ch_1", "event_seq": 5 }
})));
assert_eq!(g.channel, "im:post:received");
assert!(g.channel_ok);
assert_eq!(
g.envelope["payload"]["data"],
serde_json::json!({ "channel_id": "ch_1", "event_seq": 5 })
);
}
#[test]
fn array_data_wrapped_in_items() {
let g = guard_bus_envelope(&ev(serde_json::json!({
"event": "im:messages:query_result",
"data": [ {"id": "m1"}, {"id": "m2"} ]
})));
assert!(g.channel_ok);
assert_eq!(
g.envelope["payload"]["data"],
serde_json::json!({ "items": [ {"id": "m1"}, {"id": "m2"} ] }),
"数组 data 必须 {{items:[]}} 包装,前端 dispatcher 统一对象解构"
);
}
#[test]
fn bus_self_prefix_channel_rejected() {
let g = guard_bus_envelope(&ev(serde_json::json!({
"event": "im:__bus__",
"data": {}
})));
assert_eq!(g.channel, "im:__bus__");
assert!(!g.channel_ok, "im:__ 前缀 channel 必须被拒(防自环)");
}
#[test]
fn non_im_prefix_channel_rejected() {
let g = guard_bus_envelope(&ev(serde_json::json!({
"event": "store:setItem",
"data": {}
})));
assert!(!g.channel_ok);
}
#[test]
fn invalid_json_yields_empty_channel_not_ok() {
let g = guard_bus_envelope(&DomainEventBytes(Bytes::from_static(b"not json{")));
assert_eq!(g.channel, "");
assert!(!g.channel_ok);
}
#[test]
fn emittable_channel_boundary_table() {
assert!(is_emittable_channel("im:post:received"));
assert!(is_emittable_channel("im:channel:update"));
assert!(!is_emittable_channel("im:__bus__"));
assert!(!is_emittable_channel("im:__internal"));
assert!(!is_emittable_channel("store:get"));
assert!(!is_emittable_channel(""));
}
}