use bytes::Bytes;
use helix_core::effect::{
Correlation, DomainEventBytes, Effect, ScanOrder, ScanSpec, SqlValue, StorageOp,
};
use crate::error::ImError;
use crate::state::ChannelId;
pub const QUERY_COMMAND_NAMES: &[&str] = &[
"im_query_messages_by_channel",
"im_delete_all_dialogs",
"im_query_dialog_list",
crate::query::channel_view_snapshot::QUERY_CHANNEL_VIEW_SNAPSHOT,
"im_query_channel_sync_page",
"im_complete_channel_sync",
"im_query_subtopics",
crate::query::pinned_projection::QUERY_PINNED_PROJECTION,
crate::older_context::LOAD_OLDER_CONTEXT,
crate::timeline_navigation::LOAD_NEWER_CONTEXT,
crate::timeline_navigation::LOCATE_MESSAGE,
crate::timeline_navigation::LOCATE_CONTEXT,
];
pub fn query_command_names() -> &'static [&'static str] {
QUERY_COMMAND_NAMES
}
pub fn is_query(name: &str) -> bool {
QUERY_COMMAND_NAMES.contains(&name)
}
pub(crate) const QUERY_MESSAGES_DEFAULT: u32 = crate::timeline_state::DEFAULT_TIMELINE_PAGE_SIZE;
pub(crate) const QUERY_DIALOGS_DEFAULT: u32 = 500;
pub(crate) const QUERY_DIALOGS_MAX: u32 = 2000;
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct MessageQueryRequest {
pub channel_id: ChannelId,
pub limit: u32,
pub window_token: String,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub struct DialogListQueryRequest {
pub limit: u32,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct SubtopicsQueryRequest {
pub parent_channel_id: Option<ChannelId>,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct DialogListScope {
pub viewer_user_id: String,
pub company_id: String,
}
impl DialogListScope {
pub fn new(viewer_user_id: &str, company_id: &str) -> Self {
Self {
viewer_user_id: viewer_user_id.to_string(),
company_id: company_id.to_string(),
}
}
}
pub fn parse_message_query(payload: &[u8]) -> Result<MessageQueryRequest, ImError> {
let value: serde_json::Value = serde_json::from_slice(payload)
.map_err(|e| ImError::Parse(format!("im_query_messages_by_channel payload: {e}")))?;
let channel = value
.get("channel_id")
.or_else(|| value.get("channelId"))
.and_then(serde_json::Value::as_str)
.filter(|value| !value.is_empty())
.ok_or_else(|| ImError::Parse("im_query_messages_by_channel 缺 channel_id".into()))?;
let channel_id = ChannelId::from_str(channel)
.ok_or_else(|| ImError::Parse(format!("非法 channel_id: {channel}")))?;
let limit = crate::timeline_state::TimelinePageSize::parse(
value.get("pageSize").or_else(|| value.get("limit")),
)
.map_err(|error| ImError::Parse(format!("im_query_messages_by_channel pageSize: {error}")))?
.get();
let window_token = match value
.get("windowToken")
.or_else(|| value.get("window_token"))
{
Some(serde_json::Value::String(value)) => value.as_str(),
Some(_) => {
return Err(ImError::Parse(
"im_query_messages_by_channel windowToken must be string".into(),
));
}
None => "latest",
};
if window_token.is_empty() {
return Err(ImError::Parse(
"im_query_messages_by_channel windowToken must not be empty".into(),
));
}
Ok(MessageQueryRequest {
channel_id,
limit,
window_token: window_token.to_string(),
})
}
pub fn message_scan_spec(request: &MessageQueryRequest) -> ScanSpec {
ScanSpec {
table: "message",
limit: Some(request.limit.saturating_add(1)),
filter: Some((
"channel_id",
SqlValue::Text(request.channel_id.as_str().to_string()),
)),
order_by: MESSAGE_QUERY_ORDER,
}
}
pub fn parse_dialog_list_query(payload: &[u8]) -> Result<DialogListQueryRequest, ImError> {
let value = serde_json::from_slice::<serde_json::Value>(payload)
.map_err(|error| ImError::Parse(format!("im_query_dialog_list payload: {error}")))?;
let object = value.as_object().ok_or_else(|| {
ImError::Parse("im_query_dialog_list payload must be an object".to_string())
})?;
for key in object.keys() {
if !matches!(key.as_str(), "limit" | "req_id") {
return Err(ImError::Parse(format!(
"im_query_dialog_list field is not caller-owned: {key}"
)));
}
}
let limit = match object.get("limit") {
None => QUERY_DIALOGS_DEFAULT,
Some(serde_json::Value::Number(value)) => value
.as_u64()
.map(|value| value.clamp(1, u64::from(QUERY_DIALOGS_MAX)) as u32)
.ok_or_else(|| ImError::Parse("im_query_dialog_list limit must be integer".into()))?,
Some(_) => {
return Err(ImError::Parse(
"im_query_dialog_list limit must be integer".into(),
));
}
};
if let Some(req_id) = object.get("req_id") {
if !req_id.is_null() && !req_id.is_string() {
return Err(ImError::Parse(
"im_query_dialog_list req_id must be string".into(),
));
}
}
Ok(DialogListQueryRequest { limit })
}
pub fn parse_subtopics_query(payload: &[u8]) -> Result<SubtopicsQueryRequest, ImError> {
let value = serde_json::from_slice::<serde_json::Value>(payload)
.map_err(|error| ImError::Parse(format!("im_query_subtopics payload: {error}")))?;
let object = value.as_object().ok_or_else(|| {
ImError::Parse("im_query_subtopics payload must be an object".to_string())
})?;
for key in object.keys() {
if !matches!(
key.as_str(),
"parentChannelId" | "parent_channel_id" | "req_id"
) {
return Err(ImError::Parse(format!(
"im_query_subtopics field is not caller-owned: {key}"
)));
}
}
let parent = object
.get("parentChannelId")
.or_else(|| object.get("parent_channel_id"));
let parent_channel_id = match parent {
None | Some(serde_json::Value::Null) => None,
Some(serde_json::Value::String(value)) if value.is_empty() => None,
Some(serde_json::Value::String(value)) => {
ChannelId::from_str(value).map(Some).ok_or_else(|| {
ImError::Parse("im_query_subtopics parentChannelId is invalid".into())
})?
}
Some(_) => {
return Err(ImError::Parse(
"im_query_subtopics parentChannelId must be string".into(),
));
}
};
if let Some(req_id) = object.get("req_id") {
if !req_id.is_null() && !req_id.is_string() {
return Err(ImError::Parse(
"im_query_subtopics req_id must be string".into(),
));
}
}
Ok(SubtopicsQueryRequest { parent_channel_id })
}
pub fn dialog_list_scan_spec(request: DialogListQueryRequest) -> ScanSpec {
ScanSpec {
table: "channel",
limit: Some(request.limit),
filter: None,
order_by: DIALOG_QUERY_ORDER,
}
}
pub(crate) fn dialog_member_scan_spec(auth_user_id: &str) -> ScanSpec {
ScanSpec {
table: "channel_member",
limit: None,
filter: Some(("user_id", SqlValue::Text(auth_user_id.to_string()))),
order_by: DIALOG_MEMBER_QUERY_ORDER,
}
}
pub fn dialog_list_scan_spec_for_scope(
request: DialogListQueryRequest,
scope: &DialogListScope,
) -> Result<ScanSpec, ImError> {
if scope.company_id.is_empty() || scope.viewer_user_id.is_empty() {
return Err(ImError::Parse(
"im_query_dialog_list requires RuntimeAuth user and company".into(),
));
}
Ok(ScanSpec {
table: "channel",
limit: Some(request.limit),
filter: Some(("team_id", SqlValue::Text(scope.company_id.clone()))),
order_by: DIALOG_QUERY_ORDER,
})
}
pub fn subtopics_scan_spec_for_scope(scope: &DialogListScope) -> Result<ScanSpec, ImError> {
if scope.company_id.is_empty() || scope.viewer_user_id.is_empty() {
return Err(ImError::Parse(
"im_query_subtopics requires RuntimeAuth user and company".into(),
));
}
Ok(ScanSpec {
table: "channel",
limit: Some(QUERY_DIALOGS_MAX),
filter: Some(("team_id", SqlValue::Text(scope.company_id.clone()))),
order_by: DIALOG_QUERY_ORDER,
})
}
const MESSAGE_QUERY_ORDER: &[ScanOrder] = &[
ScanOrder::desc("create_at"),
ScanOrder::desc("temporary_id"),
];
pub(crate) const DIALOG_QUERY_ORDER: &[ScanOrder] = &[
ScanOrder::desc("is_top"),
ScanOrder::desc("last_post_at"),
ScanOrder::desc("created_at"),
];
const DIALOG_MEMBER_QUERY_ORDER: &[ScanOrder] = &[];
pub fn build_message_query(
payload: &[u8],
corr: Correlation,
) -> Result<(ChannelId, Effect), ImError> {
let request = parse_message_query(payload)?;
let effect = build_message_query_from_request(&request, corr);
Ok((request.channel_id, effect))
}
pub(crate) fn build_message_query_from_request(
request: &MessageQueryRequest,
corr: Correlation,
) -> Effect {
Effect::Persist {
corr,
ops: vec![StorageOp::Scan(message_scan_spec(request))],
}
}
pub fn emit_message_query_result(channel_id: &ChannelId, reply_bytes: &[u8]) -> Effect {
emit_message_query_result_for_viewer(channel_id, reply_bytes, "")
}
pub fn emit_message_query_result_for_viewer(
channel_id: &ChannelId,
reply_bytes: &[u8],
viewer_user_id: &str,
) -> Effect {
Effect::Emit {
event: DomainEventBytes(Bytes::from(message_query_result_bytes_for_viewer(
channel_id,
reply_bytes,
viewer_user_id,
))),
}
}
pub(crate) fn message_query_result_bytes_for_viewer(
channel_id: &ChannelId,
reply_bytes: &[u8],
viewer_user_id: &str,
) -> Vec<u8> {
message_query_result_bytes(channel_id.as_str(), reply_bytes, viewer_user_id)
}
pub(crate) fn message_query_result_bytes(
channel_id: &str,
reply_bytes: &[u8],
viewer_user_id: &str,
) -> Vec<u8> {
let rows = match serde_json::from_slice::<serde_json::Value>(reply_bytes) {
Ok(serde_json::Value::Array(mut items)) => {
items.reverse();
serde_json::Value::Array(items)
}
_ => serde_json::json!([]),
};
let payload = serde_json::json!({
"event": "im:messages:query_result",
"data": {
"channel_id": channel_id,
"messages": crate::render_ready::shape_message_rows_for_viewer(&rows, viewer_user_id),
}
});
serde_json::to_vec(&payload).expect("message query static JSON shape must serialize")
}
pub fn build_dialog_list_query(payload: &[u8], corr: Correlation) -> Result<Effect, ImError> {
let request = parse_dialog_list_query(payload)?;
Ok(Effect::Persist {
corr,
ops: vec![StorageOp::Scan(dialog_list_scan_spec(request))],
})
}
pub fn build_dialog_list_query_for_scope(
payload: &[u8],
corr: Correlation,
scope: &DialogListScope,
) -> Result<Effect, ImError> {
let request = parse_dialog_list_query(payload)?;
let scan = dialog_list_scan_spec_for_scope(request, scope)?;
Ok(Effect::Persist {
corr,
ops: vec![StorageOp::Scan(scan)],
})
}
pub fn build_subtopics_query_for_scope(
request: &SubtopicsQueryRequest,
corr: Correlation,
scope: &DialogListScope,
) -> Result<Effect, ImError> {
if request.parent_channel_id.is_none() {
return Err(ImError::Parse(
"im_query_subtopics requires parentChannelId".into(),
));
}
let scan = subtopics_scan_spec_for_scope(scope)?;
Ok(Effect::Persist {
corr,
ops: vec![StorageOp::Scan(scan)],
})
}
pub fn project_dialog_list_items(
reply_bytes: &[u8],
scope: &DialogListScope,
) -> Vec<serde_json::Value> {
let Ok(serde_json::Value::Array(rows)) = serde_json::from_slice(reply_bytes) else {
return Vec::new();
};
let mut seen = std::collections::HashSet::new();
rows.into_iter()
.filter(|row| dialog_row_visible(row, scope))
.filter_map(|row| {
let id = row
.get("id")
.or_else(|| row.get("channel_id"))
.and_then(serde_json::Value::as_str)
.filter(|id| !id.is_empty())?;
seen.insert(id.to_string())
.then(|| super::render_ready::channel::shape_channel_row(&row))
})
.collect()
}
pub fn project_channel_view_snapshot(
reply_bytes: &[u8],
scope: &DialogListScope,
channel_id: ChannelId,
) -> Option<serde_json::Value> {
let serde_json::Value::Array(rows) = serde_json::from_slice(reply_bytes).ok()? else {
return None;
};
rows.into_iter().find_map(|row| {
let id = row
.get("id")
.or_else(|| row.get("channel_id"))
.and_then(serde_json::Value::as_str)?;
(id == channel_id.as_str() && dialog_row_visible(&row, scope))
.then(|| super::render_ready::channel::shape_channel_row(&row))
})
}
pub fn emit_dialog_list_result(
req_id: &str,
reply_bytes: &[u8],
scope: &DialogListScope,
) -> Effect {
crate::read_relay::emit_read_body(
req_id,
serde_json::json!({
"items": project_dialog_list_items(reply_bytes, scope),
}),
)
}
pub fn project_subtopic_items(
reply_bytes: &[u8],
scope: &DialogListScope,
parent_channel_id: Option<&str>,
) -> Vec<serde_json::Value> {
let Some(parent_channel_id) = parent_channel_id.filter(|value| !value.is_empty()) else {
return Vec::new();
};
let Ok(serde_json::Value::Array(rows)) = serde_json::from_slice(reply_bytes) else {
return Vec::new();
};
let Some(parent_row) = rows.iter().find(|row| {
channel_row_id(row) == Some(parent_channel_id) && dialog_row_visible(row, scope)
}) else {
return Vec::new();
};
let parent_sync_watermark = parent_row
.get("subtopics_loaded_at")
.or_else(|| parent_row.get("subtopicsLoadedAt"))
.cloned();
let mut seen = std::collections::HashSet::new();
rows.into_iter()
.filter(|row| dialog_row_visible(row, scope))
.filter(|row| channel_row_type(row) == Some("T"))
.filter(|row| channel_row_root_id(row) == Some(parent_channel_id))
.filter_map(|row| {
let id = channel_row_id(&row)?;
seen.insert(id.to_string()).then(|| {
let mut rendered = render_subtopic_row(row);
if let (Some(watermark), Some(object)) =
(&parent_sync_watermark, rendered.as_object_mut())
{
object.insert("subtopicsLoadedAt".to_string(), watermark.clone());
}
rendered
})
})
.collect()
}
pub fn emit_subtopics_result(
req_id: &str,
reply_bytes: &[u8],
scope: &DialogListScope,
parent_channel_id: Option<&str>,
) -> Effect {
crate::read_relay::emit_read_body(
req_id,
serde_json::json!({
"items": project_subtopic_items(reply_bytes, scope, parent_channel_id),
}),
)
}
fn channel_row_id(row: &serde_json::Value) -> Option<&str> {
row.get("id")
.or_else(|| row.get("channel_id"))
.or_else(|| row.get("channelId"))
.and_then(serde_json::Value::as_str)
.filter(|value| !value.is_empty())
}
fn channel_row_type(row: &serde_json::Value) -> Option<&str> {
row.get("type")
.or_else(|| row.get("channel_type"))
.or_else(|| row.get("channelType"))
.and_then(serde_json::Value::as_str)
}
fn channel_row_root_id(row: &serde_json::Value) -> Option<&str> {
row.get("root_id")
.or_else(|| row.get("rootId"))
.and_then(serde_json::Value::as_str)
}
fn render_subtopic_row(row: serde_json::Value) -> serde_json::Value {
let owner_id = row
.get("user_id")
.or_else(|| row.get("userId"))
.and_then(serde_json::Value::as_str)
.map(str::to_string);
let creator_id = row
.get("create_by")
.or_else(|| row.get("createBy"))
.and_then(serde_json::Value::as_str)
.map(str::to_string);
let mut rendered = render_dialog_row(row);
let Some(owner_id) = owner_id else {
return rendered;
};
let mut members = rendered
.get("members")
.map(parse_json_column)
.and_then(|value| value.as_array().cloned())
.unwrap_or_default();
if !members.iter().any(|member| {
member
.get("userId")
.or_else(|| member.get("user_id"))
.or_else(|| member.get("id"))
.and_then(serde_json::Value::as_str)
== Some(owner_id.as_str())
}) {
members.push(serde_json::json!({ "userId": owner_id.clone() }));
}
if let Some(object) = rendered.as_object_mut() {
object.insert("userId".to_string(), serde_json::Value::String(owner_id));
if let Some(creator_id) = creator_id {
object.insert(
"createBy".to_string(),
serde_json::Value::String(creator_id),
);
}
object.insert("members".to_string(), serde_json::Value::Array(members));
let member_count = object
.get("members")
.and_then(serde_json::Value::as_array)
.map_or(0, Vec::len);
object.insert(
"memberCount".to_string(),
serde_json::Value::from(member_count as u64),
);
if object.contains_key("member_count") {
object.insert(
"member_count".to_string(),
serde_json::Value::from(member_count as u64),
);
}
}
rendered
}
pub(crate) fn channel_row_is_terminal(row: &serde_json::Value) -> bool {
["is_remove", "isRemove"].iter().any(|key| {
row.get(*key).is_some_and(|value| {
value.as_bool().unwrap_or(false)
|| value.as_i64().is_some_and(|value| value != 0)
|| value
.as_str()
.is_some_and(|value| matches!(value, "1" | "true" | "TRUE"))
})
})
}
fn dialog_row_visible(row: &serde_json::Value, scope: &DialogListScope) -> bool {
if channel_row_is_terminal(row) {
return false;
}
let Some(object) = row.as_object() else {
return false;
};
let team_id = object
.get("team_id")
.or_else(|| object.get("teamId"))
.and_then(serde_json::Value::as_str);
if team_id != Some(scope.company_id.as_str()) {
return false;
}
let owner_id = object
.get("user_id")
.or_else(|| object.get("userId"))
.and_then(serde_json::Value::as_str);
if owner_id == Some(scope.viewer_user_id.as_str()) {
return true;
}
["members", "target_users", "targetUsers"]
.iter()
.filter_map(|key| object.get(*key))
.map(parse_json_column)
.any(|members| member_list_contains(&members, scope.viewer_user_id.as_str()))
}
fn parse_json_column(value: &serde_json::Value) -> serde_json::Value {
match value {
serde_json::Value::String(raw) => {
serde_json::from_str(raw).unwrap_or_else(|_| value.clone())
}
_ => value.clone(),
}
}
fn member_list_contains(value: &serde_json::Value, viewer_user_id: &str) -> bool {
value.as_array().is_some_and(|members| {
members.iter().any(|member| {
member
.get("userId")
.or_else(|| member.get("user_id"))
.or_else(|| member.get("id"))
.and_then(serde_json::Value::as_str)
== Some(viewer_user_id)
})
})
}
fn render_dialog_row(row: serde_json::Value) -> serde_json::Value {
let serde_json::Value::Object(mut object) = row else {
return serde_json::Value::Null;
};
const ALIASES: &[(&str, &str)] = &[
("display_name", "displayName"),
("user_id", "userId"),
("team_id", "teamId"),
("root_id", "rootId"),
("root_post_id", "rootPostId"),
("member_count", "memberCount"),
("unread_count", "unreadCount"),
("mention_count", "mentionCount"),
("mention_count_root", "mentionCountRoot"),
("last_post", "lastPost"),
("last_post_at", "lastPostAt"),
("subtopics_loaded_at", "subtopicsLoadedAt"),
("created_at", "createdAt"),
("is_top", "isTop"),
("admin_users", "adminUsers"),
("target_users", "targetUsers"),
];
for (source, target) in ALIASES {
if let Some(value) = object.get(*source).cloned() {
object.insert((*target).to_string(), parse_json_column(&value));
}
}
for key in [
"members",
"admin_users",
"adminUsers",
"boss",
"owner",
"target_users",
] {
if let Some(value) = object.get(key).cloned() {
object.insert(key.to_string(), parse_json_column(&value));
}
}
serde_json::Value::Object(object)
}
pub fn emit_dialogs_cleared() -> Effect {
let payload = serde_json::json!({
"event": "im:channels:loaded",
"data": { "reason": "delete_all_dialogs" }
});
let bytes = Bytes::from(
serde_json::to_vec(&payload)
.expect("emit_dialogs_cleared: static JSON shape must serialize"),
);
Effect::Emit {
event: DomainEventBytes(bytes),
}
}
#[cfg(test)]
#[path = "core_message_tests.rs"]
mod message_tests;
#[cfg(test)]
mod phase2_contract_tests {
#[test]
fn timeline_page_size_contract_accepts_twenty_and_sixty() {
assert_eq!(
crate::timeline_state::TimelinePageSize::new(20)
.expect("default page size is valid")
.get(),
20
);
assert_eq!(
crate::timeline_state::TimelinePageSize::new(60)
.expect("maximum page size is valid")
.get(),
60
);
assert!(crate::timeline_state::TimelinePageSize::new(0).is_err());
assert!(crate::timeline_state::TimelinePageSize::new(61).is_err());
}
}