use crate::state::ChannelId;
use bytes::Bytes;
use helix_core::effect::DomainEventBytes;
use helix_core::Effect;
mod data;
use data::message_item_data;
pub fn emit_post_received(
channel_id: ChannelId,
event_seq: u64,
msg_id: &str,
fields: &crate::sync_session::PostFields,
) -> Effect {
emit_post_received_for_viewer(channel_id, event_seq, msg_id, fields, "")
}
pub fn emit_post_received_for_viewer(
channel_id: ChannelId,
event_seq: u64,
msg_id: &str,
fields: &crate::sync_session::PostFields,
viewer_user_id: &str,
) -> Effect {
emit_post_received_for_viewer_with_path(
channel_id,
event_seq,
msg_id,
fields,
viewer_user_id,
"live_ws",
)
}
pub fn emit_post_received_for_canonical_viewer(
channel_id: ChannelId,
event_seq: u64,
msg_id: &str,
fields: &crate::sync_session::PostFields,
viewer_user_id: &str,
) -> Effect {
emit_post_received_for_viewer_with_path_and_event(
channel_id,
event_seq,
msg_id,
fields,
viewer_user_id,
"live_ws",
"im:post:received",
)
}
pub fn emit_post_received_for_viewer_with_path(
channel_id: ChannelId,
event_seq: u64,
msg_id: &str,
fields: &crate::sync_session::PostFields,
viewer_user_id: &str,
telemetry_path: &'static str,
) -> Effect {
let event = if !viewer_user_id.is_empty() && fields.user_id == viewer_user_id {
"im:post:sent"
} else {
"im:post:received"
};
emit_post_received_for_viewer_with_path_and_event(
channel_id,
event_seq,
msg_id,
fields,
viewer_user_id,
telemetry_path,
event,
)
}
fn emit_post_received_for_viewer_with_path_and_event(
channel_id: ChannelId,
event_seq: u64,
msg_id: &str,
fields: &crate::sync_session::PostFields,
viewer_user_id: &str,
telemetry_path: &'static str,
event: &'static str,
) -> Effect {
use serde_json::json;
let mut data = message_item_data(channel_id, event_seq, msg_id, fields, viewer_user_id, false);
if let Some(object) = data.as_object_mut() {
object.insert("telemetryPath".to_string(), json!(telemetry_path));
}
let payload = json!({
"event": event,
"data": data,
});
let bytes = Bytes::from(
serde_json::to_vec(&payload)
.expect("emit_post_received: static JSON shape must not fail to serialize"),
);
Effect::Emit {
event: DomainEventBytes(bytes),
}
}
pub fn emit_post_read(
channel_id: ChannelId,
event_seq: u64,
msg_id: &str,
fields: &crate::sync_session::PostFields,
) -> Effect {
emit_post_read_for_viewer(channel_id, event_seq, msg_id, fields, "")
}
pub fn emit_post_read_for_viewer(
channel_id: ChannelId,
event_seq: u64,
msg_id: &str,
fields: &crate::sync_session::PostFields,
viewer_user_id: &str,
) -> Effect {
emit_post_read_with_receipt_revision_for_viewer(
channel_id,
event_seq,
msg_id,
fields,
0,
viewer_user_id,
)
}
pub fn emit_post_read_with_receipt_revision(
channel_id: ChannelId,
event_seq: u64,
msg_id: &str,
fields: &crate::sync_session::PostFields,
receipt_revision: i64,
) -> Effect {
emit_post_read_with_receipt_revision_for_viewer(
channel_id,
event_seq,
msg_id,
fields,
receipt_revision,
"",
)
}
pub fn emit_post_read_with_receipt_revision_for_viewer(
channel_id: ChannelId,
event_seq: u64,
msg_id: &str,
fields: &crate::sync_session::PostFields,
receipt_revision: i64,
viewer_user_id: &str,
) -> Effect {
use serde_json::json;
let mut data = message_item_data(channel_id, event_seq, msg_id, fields, viewer_user_id, false);
data["receiptRevision"] = json!(receipt_revision);
let payload = json!({
"event": "im:post:read",
"data": data,
});
let bytes = Bytes::from(
serde_json::to_vec(&payload)
.expect("emit_post_read: static JSON shape must not fail to serialize"),
);
Effect::Emit {
event: DomainEventBytes(bytes),
}
}
pub fn emit_sync_post_read(
channel_id: ChannelId,
event_seq: u64,
msg_id: &str,
fields: &crate::sync_session::PostFields,
reader_id: &str,
receipt_revision: i64,
) -> Effect {
emit_sync_post_read_for_viewer(
channel_id,
event_seq,
msg_id,
fields,
reader_id,
receipt_revision,
"",
)
}
pub fn emit_sync_post_read_for_viewer(
channel_id: ChannelId,
event_seq: u64,
msg_id: &str,
fields: &crate::sync_session::PostFields,
reader_id: &str,
receipt_revision: i64,
viewer_user_id: &str,
) -> Effect {
use serde_json::json;
let mut data = message_item_data(channel_id, event_seq, msg_id, fields, viewer_user_id, false);
data["postId"] = json!(msg_id);
data["readerId"] = json!(reader_id);
data["receiptRevision"] = json!(receipt_revision);
let payload = json!({
"event": "im:post:read",
"data": data,
});
let bytes = Bytes::from(
serde_json::to_vec(&payload)
.expect("emit_sync_post_read: static JSON shape must not fail to serialize"),
);
Effect::Emit {
event: DomainEventBytes(bytes),
}
}
pub fn emit_channel_read_echo(
channel_id: ChannelId,
event_seq: u64,
msg_id: &str,
fields: &crate::sync_session::PostFields,
) -> Effect {
emit_channel_read_echo_for_viewer(channel_id, event_seq, msg_id, fields, "")
}
pub fn emit_channel_read_echo_for_viewer(
channel_id: ChannelId,
event_seq: u64,
msg_id: &str,
fields: &crate::sync_session::PostFields,
viewer_user_id: &str,
) -> Effect {
use serde_json::json;
let payload = json!({
"event": "im:channel:read_echo",
"data": message_item_data(channel_id, event_seq, msg_id, fields, viewer_user_id, false),
});
let bytes = Bytes::from(
serde_json::to_vec(&payload)
.expect("emit_channel_read_echo: static JSON shape must not fail to serialize"),
);
Effect::Emit {
event: DomainEventBytes(bytes),
}
}
pub fn emit_post_updated(
channel_id: ChannelId,
event_seq: u64,
msg_id: &str,
fields: &crate::sync_session::PostFields,
) -> Effect {
emit_post_updated_for_viewer(channel_id, event_seq, msg_id, fields, "")
}
pub fn emit_post_updated_for_viewer(
channel_id: ChannelId,
event_seq: u64,
msg_id: &str,
fields: &crate::sync_session::PostFields,
viewer_user_id: &str,
) -> Effect {
use serde_json::json;
let payload = json!({
"event": "im:post:updated",
"data": message_item_data(channel_id, event_seq, msg_id, fields, viewer_user_id, false),
});
let bytes = Bytes::from(
serde_json::to_vec(&payload)
.expect("emit_post_updated: static JSON shape must not fail to serialize"),
);
Effect::Emit {
event: DomainEventBytes(bytes),
}
}
pub fn emit_post_deleted(
channel_id: ChannelId,
event_seq: u64,
msg_id: &str,
fields: &crate::sync_session::PostFields,
) -> Effect {
emit_post_deleted_for_viewer(channel_id, event_seq, msg_id, fields, "", "", 0, "")
}
#[allow(clippy::too_many_arguments)]
pub fn emit_post_deleted_for_viewer(
channel_id: ChannelId,
event_seq: u64,
msg_id: &str,
fields: &crate::sync_session::PostFields,
event_id: &str,
actor_id: &str,
occurred_at: i64,
viewer_user_id: &str,
) -> Effect {
use serde_json::json;
let actor_available = !actor_id.is_empty();
let is_self = actor_available && actor_id == viewer_user_id;
let display_text = if is_self {
"你撤回了一条消息"
} else {
"某人撤回了一条消息"
};
let stable_event_id = if event_id.starts_with("revoke:") {
event_id.to_string()
} else {
format!("revoke:{}:{}", channel_id.as_str(), event_seq)
};
let mut data = message_item_data(channel_id, event_seq, msg_id, fields, "", true);
let object = data
.as_object_mut()
.expect("recall projection starts from a static JSON object");
object.insert("recalledText".to_string(), json!(fields.message));
object.insert("type".to_string(), json!("system"));
object.insert("message".to_string(), json!(display_text));
object.insert("text".to_string(), json!(display_text));
object.insert("systemNotice".to_string(), json!(true));
object.insert("isSelf".to_string(), json!(is_self));
object.insert("kind".to_string(), json!("system"));
object.insert("system".to_string(), json!(true));
object.insert("systemEvent".to_string(), json!("message-recalled"));
object.insert("eventId".to_string(), json!(stable_event_id));
object.insert("actorId".to_string(), json!(actor_id));
object.insert("actorAvailable".to_string(), json!(actor_available));
object.insert("subjectMemberIds".to_string(), json!([]));
object.insert("occurredAt".to_string(), json!(occurred_at.max(0)));
object.insert("displayText".to_string(), json!(display_text));
object.insert("previewText".to_string(), json!(display_text));
object.insert("recalledMsgId".to_string(), json!(msg_id));
object.insert("targetMsgId".to_string(), json!(msg_id));
let payload = json!({
"event": "im:post:deleted",
"data": data,
});
let bytes = Bytes::from(
serde_json::to_vec(&payload)
.expect("emit_post_deleted: static JSON shape must not fail to serialize"),
);
Effect::Emit {
event: DomainEventBytes(bytes),
}
}
pub fn emit_post_revoke_for_canonical_viewer(
channel_id: ChannelId,
event_seq: u64,
msg_id: &str,
fields: &crate::sync_session::PostFields,
) -> Effect {
use serde_json::json;
let mut data = json!({
"id": msg_id,
"channelId": channel_id.as_str(),
"eventSeq": event_seq,
"revoke": true,
});
let object = data
.as_object_mut()
.expect("canonical revoke projection starts from a static JSON object");
for (key, value) in [
("temporaryId", fields.temporary_id.as_str()),
("userId", fields.user_id.as_str()),
("type", fields.msg_type.as_str()),
("message", fields.message.as_str()),
] {
if !value.is_empty() {
object.insert(key.to_string(), json!(value));
}
}
if fields.create_at > 0 {
object.insert("createAt".to_string(), json!(fields.create_at));
}
if fields.update_at > 0 {
object.insert("updateAt".to_string(), json!(fields.update_at));
}
crate::event::post::revoke(data)
.expect("emit_post_revoke_for_canonical_viewer: static JSON shape must serialize")
.into_effect()
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn canonical_sender_projection_is_received_with_self_marker() {
let channel_id = crate::state::test_channel_id(41);
let fields = crate::sync_session::PostFields {
id: "post-41".to_string(),
user_id: "viewer-41".to_string(),
temporary_id: "tmp-41".to_string(),
message: "hello".to_string(),
..Default::default()
};
let Effect::Emit { event } =
emit_post_received_for_canonical_viewer(channel_id, 7, "post-41", &fields, "viewer-41")
else {
panic!("canonical projection must emit a domain event");
};
let payload: serde_json::Value = serde_json::from_slice(event.0.as_ref()).unwrap();
assert_eq!(payload["event"], "im:post:received");
assert_eq!(payload["data"]["isSelf"], true);
assert_eq!(payload["data"]["temporaryId"], "tmp-41");
}
#[test]
fn sync_sender_projection_remains_sent() {
let channel_id = crate::state::test_channel_id(42);
let fields = crate::sync_session::PostFields {
user_id: "viewer-42".to_string(),
..Default::default()
};
let Effect::Emit { event } = emit_post_received_for_viewer_with_path(
channel_id,
8,
"post-42",
&fields,
"viewer-42",
"sync_replay",
) else {
panic!("sync projection must emit a domain event");
};
let payload: serde_json::Value = serde_json::from_slice(event.0.as_ref()).unwrap();
assert_eq!(payload["event"], "im:post:sent");
assert_eq!(payload["data"]["telemetryPath"], "sync_replay");
}
#[test]
fn canonical_revoke_projection_uses_online_revoke_event() {
let channel_id = crate::state::test_channel_id(43);
let fields = crate::sync_session::PostFields {
id: "post-43".to_string(),
message: "hello".to_string(),
..Default::default()
};
let Effect::Emit { event } =
emit_post_revoke_for_canonical_viewer(channel_id, 9, "post-43", &fields)
else {
panic!("canonical revoke projection must emit a domain event");
};
let payload: serde_json::Value = serde_json::from_slice(event.0.as_ref()).unwrap();
assert_eq!(payload["event"], "im:post:revoke");
assert_eq!(payload["data"]["id"], "post-43");
assert_eq!(payload["data"]["channelId"], channel_id.as_str());
assert_eq!(payload["data"]["eventSeq"], 9);
assert_eq!(payload["data"]["revoke"], true);
assert_eq!(payload["data"]["message"], "hello");
assert!(payload["data"].get("systemEvent").is_none());
}
}