mirror of
https://github.com/jackwener/wx-cli.git
synced 2026-10-08 15:45:46 +00:00
chore: Apache-2.0 license, Windows support, install.ps1
This commit is contained in:
+1
-2
@@ -56,8 +56,7 @@ impl DbCache {
|
||||
|
||||
fn cache_file_path(&self, rel_key: &str) -> PathBuf {
|
||||
let hash = format!("{:x}", md5::compute(rel_key.as_bytes()));
|
||||
let short = &hash[..12];
|
||||
self.cache_dir.join(format!("{}.db", short))
|
||||
self.cache_dir.join(format!("{}.db", hash))
|
||||
}
|
||||
|
||||
/// 从持久化文件加载 mtime 记录,复用未过期的解密文件
|
||||
|
||||
+1
-145
@@ -4,9 +4,7 @@ pub mod server;
|
||||
|
||||
use anyhow::Result;
|
||||
use std::collections::HashMap;
|
||||
use std::path::PathBuf;
|
||||
use std::sync::Arc;
|
||||
use tokio::sync::broadcast;
|
||||
|
||||
use crate::config;
|
||||
|
||||
@@ -78,154 +76,12 @@ async fn async_run() -> Result<()> {
|
||||
|
||||
let names_arc = Arc::new(std::sync::RwLock::new(names));
|
||||
|
||||
// 启动 WAL watcher
|
||||
let (watch_tx, _) = broadcast::channel::<crate::ipc::WatchEvent>(500);
|
||||
let session_wal = cfg.db_dir.join("session").join("session.db-wal");
|
||||
|
||||
// SAFETY: 我们确保 db 和 names_arc 在 daemon 生命周期内有效
|
||||
// 使用 Arc 传递引用避免 'static 问题
|
||||
let db_arc = Arc::clone(&db);
|
||||
let names_arc2 = Arc::clone(&names_arc);
|
||||
let tx_clone = watch_tx.clone();
|
||||
let session_wal2 = session_wal.clone();
|
||||
tokio::spawn(async move {
|
||||
run_watcher(db_arc, names_arc2, tx_clone, session_wal2).await;
|
||||
});
|
||||
|
||||
// 启动 IPC server(阻塞)
|
||||
server::serve(Arc::clone(&db), Arc::clone(&names_arc), watch_tx).await?;
|
||||
server::serve(Arc::clone(&db), Arc::clone(&names_arc)).await?;
|
||||
|
||||
Ok(())
|
||||
}
|
||||
|
||||
async fn run_watcher(
|
||||
db: Arc<cache::DbCache>,
|
||||
names: Arc<std::sync::RwLock<query::Names>>,
|
||||
tx: broadcast::Sender<crate::ipc::WatchEvent>,
|
||||
session_wal: PathBuf,
|
||||
) {
|
||||
use std::collections::HashMap;
|
||||
use std::time::Duration;
|
||||
use crate::ipc::WatchEvent;
|
||||
|
||||
let mut last_mtime = 0u64;
|
||||
let mut last_ts: HashMap<String, i64> = HashMap::new();
|
||||
let mut initialized = false;
|
||||
|
||||
loop {
|
||||
tokio::time::sleep(Duration::from_millis(500)).await;
|
||||
|
||||
if tx.receiver_count() == 0 {
|
||||
continue;
|
||||
}
|
||||
|
||||
let wal_mtime = match mtime_nanos(&session_wal) {
|
||||
0 => continue,
|
||||
m => m,
|
||||
};
|
||||
if wal_mtime == last_mtime {
|
||||
continue;
|
||||
}
|
||||
last_mtime = wal_mtime;
|
||||
|
||||
let path = match db.get("session/session.db").await {
|
||||
Ok(Some(p)) => p,
|
||||
_ => continue,
|
||||
};
|
||||
|
||||
let path2 = path.clone();
|
||||
let rows: Vec<(String, Vec<u8>, i64, i64, String)> = match tokio::task::spawn_blocking(move || {
|
||||
let conn = rusqlite::Connection::open(&path2)?;
|
||||
let mut stmt = conn.prepare(
|
||||
"SELECT username, summary, last_timestamp, last_msg_type, last_msg_sender
|
||||
FROM SessionTable WHERE last_timestamp > 0
|
||||
ORDER BY last_timestamp DESC LIMIT 50"
|
||||
)?;
|
||||
let rows = stmt.query_map([], |row| {
|
||||
Ok((
|
||||
row.get::<_, String>(0)?,
|
||||
row.get::<_, Vec<u8>>(1)
|
||||
.or_else(|_| row.get::<_, String>(1).map(|s| s.into_bytes()))
|
||||
.unwrap_or_default(),
|
||||
row.get::<_, i64>(2)?,
|
||||
row.get::<_, i64>(3).unwrap_or(0),
|
||||
row.get::<_, String>(4).unwrap_or_default(),
|
||||
))
|
||||
})?.collect::<rusqlite::Result<Vec<_>>>()?;
|
||||
Ok::<_, anyhow::Error>(rows)
|
||||
}).await {
|
||||
Ok(Ok(r)) => r,
|
||||
_ => continue,
|
||||
};
|
||||
|
||||
let names_guard = match names.read() {
|
||||
Ok(g) => g,
|
||||
Err(_) => continue,
|
||||
};
|
||||
|
||||
for (username, summary_bytes, ts, msg_type, sender) in &rows {
|
||||
if !initialized {
|
||||
last_ts.insert(username.clone(), *ts);
|
||||
continue;
|
||||
}
|
||||
let prev_ts = last_ts.get(username).copied().unwrap_or(0);
|
||||
if *ts <= prev_ts {
|
||||
continue;
|
||||
}
|
||||
last_ts.insert(username.clone(), *ts);
|
||||
|
||||
let display = names_guard.display(username);
|
||||
let is_group = username.contains("@chatroom");
|
||||
let summary = decompress_or_str(summary_bytes);
|
||||
let summary = if summary.contains(":\n") {
|
||||
summary.splitn(2, ":\n").nth(1).unwrap_or(&summary).to_string()
|
||||
} else {
|
||||
summary
|
||||
};
|
||||
let sender_display = if !sender.is_empty() {
|
||||
names_guard.map.get(sender).cloned().unwrap_or_else(|| sender.clone())
|
||||
} else {
|
||||
String::new()
|
||||
};
|
||||
|
||||
let event = WatchEvent {
|
||||
event: "message".into(),
|
||||
time: Some(fmt_hhmm(*ts)),
|
||||
chat: Some(display),
|
||||
username: Some(username.clone()),
|
||||
is_group: Some(is_group),
|
||||
sender: Some(sender_display),
|
||||
content: Some(summary),
|
||||
msg_type: Some(query::fmt_type(*msg_type)),
|
||||
timestamp: Some(*ts),
|
||||
};
|
||||
let _ = tx.send(event);
|
||||
}
|
||||
|
||||
if !initialized {
|
||||
initialized = true;
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
use cache::mtime_nanos;
|
||||
|
||||
fn decompress_or_str(data: &[u8]) -> String {
|
||||
if data.is_empty() { return String::new(); }
|
||||
if let Ok(dec) = zstd::decode_all(data) {
|
||||
if let Ok(s) = String::from_utf8(dec) { return s; }
|
||||
}
|
||||
String::from_utf8_lossy(data).into_owned()
|
||||
}
|
||||
|
||||
fn fmt_hhmm(ts: i64) -> String {
|
||||
use chrono::{Local, TimeZone};
|
||||
Local.timestamp_opt(ts, 0)
|
||||
.single()
|
||||
.map(|dt| dt.format("%H:%M").to_string())
|
||||
.unwrap_or_else(|| ts.to_string())
|
||||
}
|
||||
|
||||
/// 从 all_keys.json 提取 rel_key -> enc_key 映射
|
||||
///
|
||||
/// 兼容两种格式:
|
||||
|
||||
+668
-17
@@ -1,5 +1,5 @@
|
||||
use anyhow::{Context, Result};
|
||||
use chrono::{Local, TimeZone};
|
||||
use chrono::{Local, TimeZone, Timelike};
|
||||
use regex::Regex;
|
||||
use rusqlite::Connection;
|
||||
use serde_json::{json, Value};
|
||||
@@ -140,6 +140,7 @@ pub async fn q_history(
|
||||
offset: usize,
|
||||
since: Option<i64>,
|
||||
until: Option<i64>,
|
||||
msg_type: Option<i64>,
|
||||
) -> Result<Value> {
|
||||
let username = resolve_username(chat, names)
|
||||
.with_context(|| format!("找不到联系人: {}", chat))?;
|
||||
@@ -164,7 +165,9 @@ pub async fn q_history(
|
||||
let offset2 = offset;
|
||||
|
||||
let msgs: Vec<Value> = tokio::task::spawn_blocking(move || {
|
||||
query_messages(&path, &tname, &uname, is_group2, &names_map, since2, until2, limit2 + offset2, 0)
|
||||
// per-DB 软上限:offset + limit 已足够全局分页,避免大群全量加载
|
||||
let per_db_cap = offset2 + limit2;
|
||||
query_messages(&path, &tname, &uname, is_group2, &names_map, since2, until2, msg_type, per_db_cap, 0)
|
||||
}).await??;
|
||||
|
||||
all_msgs.extend(msgs);
|
||||
@@ -193,6 +196,7 @@ pub async fn q_search(
|
||||
limit: usize,
|
||||
since: Option<i64>,
|
||||
until: Option<i64>,
|
||||
msg_type: Option<i64>,
|
||||
) -> Result<Value> {
|
||||
let mut targets: Vec<(String, String, String, String)> = Vec::new(); // (path, table, display, uname)
|
||||
|
||||
@@ -277,7 +281,7 @@ pub async fn q_search(
|
||||
for (tname, display, uname) in &table_list {
|
||||
let is_group = uname.contains("@chatroom");
|
||||
match search_in_table(&conn, tname, &uname, is_group,
|
||||
&names_map2, &kw2, since2, until2, limit2)
|
||||
&names_map2, &kw2, since2, until2, msg_type, limit2)
|
||||
{
|
||||
Ok(rows) => {
|
||||
for mut row in rows {
|
||||
@@ -343,19 +347,21 @@ fn resolve_username(chat_name: &str, names: &Names) -> Option<String> {
|
||||
return Some(chat_name.to_string());
|
||||
}
|
||||
let low = chat_name.to_lowercase();
|
||||
// 精确匹配显示名
|
||||
for (uname, display) in &names.map {
|
||||
if low == display.to_lowercase() {
|
||||
return Some(uname.clone());
|
||||
}
|
||||
// 精确匹配显示名:排序后取第一个,保证确定性
|
||||
let mut exact: Vec<&String> = names.map.iter()
|
||||
.filter(|(_, display)| display.to_lowercase() == low)
|
||||
.map(|(uname, _)| uname)
|
||||
.collect();
|
||||
exact.sort();
|
||||
if let Some(u) = exact.into_iter().next() {
|
||||
return Some(u.clone());
|
||||
}
|
||||
// 模糊匹配
|
||||
for (uname, display) in &names.map {
|
||||
if display.to_lowercase().contains(&low) {
|
||||
return Some(uname.clone());
|
||||
}
|
||||
}
|
||||
None
|
||||
// 模糊匹配:取 display name 最短的(最精确),相同长度取字典序最小
|
||||
let mut candidates: Vec<(&String, &String)> = names.map.iter()
|
||||
.filter(|(_, display)| display.to_lowercase().contains(&low))
|
||||
.collect();
|
||||
candidates.sort_by_key(|(uname, display)| (display.len(), uname.as_str()));
|
||||
candidates.into_iter().next().map(|(uname, _)| uname.clone())
|
||||
}
|
||||
|
||||
async fn find_msg_tables(
|
||||
@@ -364,8 +370,7 @@ async fn find_msg_tables(
|
||||
username: &str,
|
||||
) -> Result<Vec<(std::path::PathBuf, String)>> {
|
||||
let table_name = format!("Msg_{:x}", md5::compute(username.as_bytes()));
|
||||
let re = Regex::new(r"^Msg_[0-9a-f]{32}$").unwrap();
|
||||
if !re.is_match(&table_name) {
|
||||
if !msg_table_re().is_match(&table_name) {
|
||||
return Ok(Vec::new());
|
||||
}
|
||||
|
||||
@@ -413,6 +418,7 @@ fn query_messages(
|
||||
names_map: &HashMap<String, String>,
|
||||
since: Option<i64>,
|
||||
until: Option<i64>,
|
||||
msg_type: Option<i64>,
|
||||
limit: usize,
|
||||
offset: usize,
|
||||
) -> Result<Vec<Value>> {
|
||||
@@ -429,6 +435,10 @@ fn query_messages(
|
||||
clauses.push("create_time <= ?");
|
||||
params.push(Box::new(u));
|
||||
}
|
||||
if let Some(t) = msg_type {
|
||||
clauses.push("local_type = ?");
|
||||
params.push(Box::new(t));
|
||||
}
|
||||
let where_clause = if clauses.is_empty() {
|
||||
String::new()
|
||||
} else {
|
||||
@@ -487,6 +497,7 @@ fn search_in_table(
|
||||
keyword: &str,
|
||||
since: Option<i64>,
|
||||
until: Option<i64>,
|
||||
msg_type: Option<i64>,
|
||||
limit: usize,
|
||||
) -> Result<Vec<Value>> {
|
||||
let id2u = load_id2u(conn);
|
||||
@@ -502,6 +513,10 @@ fn search_in_table(
|
||||
clauses.push("create_time <= ?".into());
|
||||
params.push(Box::new(u));
|
||||
}
|
||||
if let Some(t) = msg_type {
|
||||
clauses.push("local_type = ?".into());
|
||||
params.push(Box::new(t));
|
||||
}
|
||||
let where_clause = format!("WHERE {}", clauses.join(" AND "));
|
||||
let sql = format!(
|
||||
"SELECT local_id, local_type, create_time, real_sender_id,
|
||||
@@ -761,3 +776,639 @@ fn fmt_time(ts: i64, fmt: &str) -> String {
|
||||
.unwrap_or_else(|| ts.to_string())
|
||||
}
|
||||
|
||||
// ─── 新增命令查询函数 ──────────────────────────────────────────────────────────
|
||||
|
||||
/// 查询有未读消息的会话
|
||||
pub async fn q_unread(db: &DbCache, names: &Names, limit: usize) -> Result<Value> {
|
||||
let path = db.get("session/session.db").await?
|
||||
.context("无法解密 session.db")?;
|
||||
|
||||
let limit_val = limit;
|
||||
let rows: Vec<(String, i64, Vec<u8>, i64, i64, String, String)> = tokio::task::spawn_blocking(move || {
|
||||
let conn = Connection::open(&path)?;
|
||||
let mut stmt = conn.prepare(
|
||||
"SELECT username, unread_count, summary, last_timestamp,
|
||||
last_msg_type, last_msg_sender, last_sender_display_name
|
||||
FROM SessionTable
|
||||
WHERE unread_count > 0
|
||||
ORDER BY last_timestamp DESC LIMIT ?"
|
||||
)?;
|
||||
let rows = stmt.query_map([limit_val as i64], |row| {
|
||||
Ok((
|
||||
row.get::<_, String>(0)?,
|
||||
row.get::<_, i64>(1).unwrap_or(0),
|
||||
get_content_bytes(row, 2),
|
||||
row.get::<_, i64>(3).unwrap_or(0),
|
||||
row.get::<_, i64>(4).unwrap_or(0),
|
||||
row.get::<_, String>(5).unwrap_or_default(),
|
||||
row.get::<_, String>(6).unwrap_or_default(),
|
||||
))
|
||||
})?
|
||||
.collect::<rusqlite::Result<Vec<_>>>()?;
|
||||
Ok::<_, anyhow::Error>(rows)
|
||||
}).await??;
|
||||
|
||||
let mut results = Vec::new();
|
||||
for (username, unread, summary_bytes, ts, msg_type, sender, sender_name) in rows {
|
||||
let display = names.display(&username);
|
||||
let is_group = username.contains("@chatroom");
|
||||
let summary = decompress_or_str(&summary_bytes);
|
||||
let summary = strip_group_prefix(&summary);
|
||||
let sender_display = if is_group && !sender.is_empty() {
|
||||
names.map.get(&sender).cloned().unwrap_or_else(|| {
|
||||
if !sender_name.is_empty() { sender_name.clone() } else { sender.clone() }
|
||||
})
|
||||
} else {
|
||||
String::new()
|
||||
};
|
||||
results.push(json!({
|
||||
"chat": display,
|
||||
"username": username,
|
||||
"is_group": is_group,
|
||||
"unread": unread,
|
||||
"last_msg_type": fmt_type(msg_type),
|
||||
"last_sender": sender_display,
|
||||
"summary": summary,
|
||||
"timestamp": ts,
|
||||
"time": fmt_time(ts, "%m-%d %H:%M"),
|
||||
}));
|
||||
}
|
||||
let total = results.len();
|
||||
Ok(json!({ "sessions": results, "total": total }))
|
||||
}
|
||||
|
||||
/// 查询群成员:优先从 contact.db 的 chatroom_member/chat_room 表获取完整列表,
|
||||
/// 若表不存在则退化为从消息记录聚合有发言记录的成员
|
||||
pub async fn q_members(db: &DbCache, names: &Names, chat: &str) -> Result<Value> {
|
||||
let username = resolve_username(chat, names)
|
||||
.with_context(|| format!("找不到联系人: {}", chat))?;
|
||||
|
||||
if !username.contains("@chatroom") {
|
||||
anyhow::bail!("'{}' 不是群聊,无法查看群成员", names.display(&username));
|
||||
}
|
||||
|
||||
let display = names.display(&username);
|
||||
let names_map = names.map.clone();
|
||||
|
||||
// 优先路径:contact.db → chatroom_member + chat_room(完整成员列表)
|
||||
if let Some(contact_p) = db.get("contact/contact.db").await? {
|
||||
let uname2 = username.clone();
|
||||
let display2 = display.clone();
|
||||
let names_map2 = names_map.clone();
|
||||
|
||||
let members_opt: Option<Vec<Value>> = tokio::task::spawn_blocking(move || {
|
||||
let conn = Connection::open(&contact_p)?;
|
||||
|
||||
let has_table: bool = conn.query_row(
|
||||
"SELECT 1 FROM sqlite_master WHERE type='table' AND name='chatroom_member'",
|
||||
[],
|
||||
|_| Ok(true),
|
||||
).unwrap_or(false);
|
||||
|
||||
if !has_table {
|
||||
return Ok::<_, anyhow::Error>(None);
|
||||
}
|
||||
|
||||
// 从 chat_room 表获取整数 room_id 和群主
|
||||
// WeChat 不同版本列名可能不同:username / chat_room_name / name
|
||||
let (room_id, owner): (i64, String) = [
|
||||
"SELECT id, owner FROM chat_room WHERE username = ?",
|
||||
"SELECT id, owner FROM chat_room WHERE chat_room_name = ?",
|
||||
"SELECT id, owner FROM chat_room WHERE name = ?",
|
||||
].iter().find_map(|sql| {
|
||||
conn.query_row(sql, [&uname2], |row| {
|
||||
Ok((row.get::<_, i64>(0)?, row.get::<_, String>(1).unwrap_or_default()))
|
||||
}).ok()
|
||||
}).unwrap_or((0, String::new()));
|
||||
|
||||
if room_id == 0 {
|
||||
return Ok::<_, anyhow::Error>(None);
|
||||
}
|
||||
|
||||
let mut stmt = conn.prepare(
|
||||
"SELECT c.username, c.nick_name, c.remark
|
||||
FROM chatroom_member cm
|
||||
LEFT JOIN contact c ON c.id = cm.member_id
|
||||
WHERE cm.room_id = ?"
|
||||
)?;
|
||||
let raw: Vec<(String, String, String)> = stmt.query_map([room_id], |row| {
|
||||
Ok((
|
||||
row.get::<_, String>(0).unwrap_or_default(),
|
||||
row.get::<_, String>(1).unwrap_or_default(),
|
||||
row.get::<_, String>(2).unwrap_or_default(),
|
||||
))
|
||||
})?
|
||||
.filter_map(|r| r.ok())
|
||||
.filter(|(uid, _, _)| !uid.is_empty())
|
||||
.collect();
|
||||
|
||||
if raw.is_empty() {
|
||||
return Ok(None);
|
||||
}
|
||||
|
||||
let mut members: Vec<Value> = raw.iter().map(|(uid, nick, remark)| {
|
||||
let disp = if !remark.is_empty() { remark.clone() }
|
||||
else if !nick.is_empty() { nick.clone() }
|
||||
else { names_map2.get(uid).cloned().unwrap_or_else(|| uid.clone()) };
|
||||
let is_owner = uid == &owner && !owner.is_empty();
|
||||
json!({ "username": uid, "display": disp, "is_owner": is_owner })
|
||||
}).collect();
|
||||
|
||||
// 群主排首位,其余按 display 字典序
|
||||
members.sort_by(|a, b| {
|
||||
let ao = a["is_owner"].as_bool().unwrap_or(false);
|
||||
let bo = b["is_owner"].as_bool().unwrap_or(false);
|
||||
if ao != bo { return bo.cmp(&ao); }
|
||||
a["display"].as_str().unwrap_or("").cmp(b["display"].as_str().unwrap_or(""))
|
||||
});
|
||||
|
||||
let _ = display2; // 不在此 closure 内使用
|
||||
Ok(Some(members))
|
||||
}).await??;
|
||||
|
||||
if let Some(members) = members_opt {
|
||||
return Ok(json!({
|
||||
"chat": display,
|
||||
"username": username,
|
||||
"count": members.len(),
|
||||
"members": members,
|
||||
}));
|
||||
}
|
||||
}
|
||||
|
||||
// 降级路径:从消息记录中聚合发言过的成员
|
||||
let tables = find_msg_tables(db, names, &username).await?;
|
||||
if tables.is_empty() {
|
||||
return Ok(json!({
|
||||
"chat": display,
|
||||
"username": username,
|
||||
"count": 0,
|
||||
"members": [],
|
||||
}));
|
||||
}
|
||||
|
||||
let mut sender_set: std::collections::HashSet<String> = std::collections::HashSet::new();
|
||||
for (db_path, table_name) in &tables {
|
||||
let path = db_path.clone();
|
||||
let tname = table_name.clone();
|
||||
let uname = username.clone();
|
||||
|
||||
let senders: Vec<String> = tokio::task::spawn_blocking(move || {
|
||||
let conn = Connection::open(&path)?;
|
||||
let id2u = load_id2u(&conn);
|
||||
let mut stmt = conn.prepare(&format!(
|
||||
"SELECT DISTINCT real_sender_id FROM [{}] WHERE real_sender_id > 0", tname
|
||||
))?;
|
||||
let ids: Vec<i64> = stmt.query_map([], |row| row.get(0))?
|
||||
.filter_map(|r| r.ok())
|
||||
.collect();
|
||||
let senders: Vec<String> = ids.iter()
|
||||
.filter_map(|id| id2u.get(id))
|
||||
.filter(|u| *u != &uname)
|
||||
.cloned()
|
||||
.collect();
|
||||
Ok::<_, anyhow::Error>(senders)
|
||||
}).await??;
|
||||
|
||||
sender_set.extend(senders);
|
||||
}
|
||||
|
||||
let mut members: Vec<Value> = sender_set.iter().map(|u| {
|
||||
json!({
|
||||
"username": u,
|
||||
"display": names_map.get(u).cloned().unwrap_or_else(|| u.clone()),
|
||||
"is_owner": false,
|
||||
})
|
||||
}).collect();
|
||||
members.sort_by(|a, b| {
|
||||
a["display"].as_str().unwrap_or("").cmp(b["display"].as_str().unwrap_or(""))
|
||||
});
|
||||
|
||||
Ok(json!({
|
||||
"chat": display,
|
||||
"username": username,
|
||||
"count": members.len(),
|
||||
"members": members,
|
||||
}))
|
||||
}
|
||||
|
||||
/// 查询新消息:以 session.db 的 last_timestamp 作为 inbox 索引,
|
||||
/// 只查询 last_timestamp > state[username] 的会话,精确且高效
|
||||
pub async fn q_new_messages(
|
||||
db: &DbCache,
|
||||
names: &Names,
|
||||
state: Option<HashMap<String, i64>>,
|
||||
limit: usize,
|
||||
) -> Result<Value> {
|
||||
// 首次运行(state=None)或未见过的会话,用 24h 前作为起点,
|
||||
// 避免第一次运行时把全量历史消息涌入
|
||||
let fallback_ts = chrono::Utc::now().timestamp() - 86400;
|
||||
|
||||
// 1. 从 session.db 读取所有会话的当前 last_timestamp
|
||||
let session_path = db.get("session/session.db").await?
|
||||
.context("无法解密 session.db")?;
|
||||
|
||||
let all_sessions: Vec<(String, i64)> = tokio::task::spawn_blocking(move || {
|
||||
let conn = Connection::open(&session_path)?;
|
||||
let mut stmt = conn.prepare(
|
||||
"SELECT username, last_timestamp FROM SessionTable WHERE last_timestamp > 0"
|
||||
)?;
|
||||
let rows = stmt.query_map([], |row| {
|
||||
Ok((row.get::<_, String>(0)?, row.get::<_, i64>(1).unwrap_or(0)))
|
||||
})?
|
||||
.collect::<rusqlite::Result<Vec<_>>>()?;
|
||||
Ok::<_, anyhow::Error>(rows)
|
||||
}).await??;
|
||||
|
||||
// 2. 记录 session.db 的当前快照(用于构建 new_state 基础)
|
||||
let session_ts_map: HashMap<String, i64> = all_sessions.iter()
|
||||
.map(|(u, ts)| (u.clone(), *ts))
|
||||
.collect();
|
||||
|
||||
// 3. 找出有新消息的会话
|
||||
// 不在 state 中的会话(首次运行或新会话)以 fallback_ts 为基准
|
||||
let changed: Vec<(String, i64)> = all_sessions.into_iter()
|
||||
.filter(|(uname, ts)| {
|
||||
let last_known = state.as_ref()
|
||||
.and_then(|m| m.get(uname))
|
||||
.copied()
|
||||
.unwrap_or(fallback_ts);
|
||||
*ts > last_known
|
||||
})
|
||||
.collect();
|
||||
|
||||
if changed.is_empty() {
|
||||
return Ok(json!({
|
||||
"count": 0,
|
||||
"messages": [],
|
||||
"new_state": session_ts_map,
|
||||
}));
|
||||
}
|
||||
|
||||
// 4. 只查询有新消息的会话的消息表
|
||||
// per_table_limit 取 limit*5 防止单表截断,最终由全局 truncate 收尾
|
||||
let per_table_limit = limit.saturating_mul(5).max(200);
|
||||
let mut all_msgs: Vec<Value> = Vec::new();
|
||||
|
||||
for (uname, _) in &changed {
|
||||
let since_ts = state.as_ref()
|
||||
.and_then(|m| m.get(uname))
|
||||
.copied()
|
||||
.unwrap_or(fallback_ts);
|
||||
let tables = find_msg_tables(db, names, uname).await?;
|
||||
if tables.is_empty() { continue; }
|
||||
|
||||
let display = names.display(uname);
|
||||
let is_group = uname.contains("@chatroom");
|
||||
|
||||
for (db_path, table_name) in &tables {
|
||||
let path = db_path.clone();
|
||||
let tname = table_name.clone();
|
||||
let uname2 = uname.clone();
|
||||
let display2 = display.clone();
|
||||
let names_map = names.map.clone();
|
||||
let tname_for_log = tname.clone();
|
||||
|
||||
let msgs: Vec<Value> = match tokio::task::spawn_blocking(move || {
|
||||
let conn = Connection::open(&path)?;
|
||||
let id2u = load_id2u(&conn);
|
||||
|
||||
let sql = format!(
|
||||
"SELECT local_id, local_type, create_time, real_sender_id,
|
||||
message_content, WCDB_CT_message_content
|
||||
FROM [{}] WHERE create_time > ? ORDER BY create_time ASC LIMIT ?",
|
||||
tname
|
||||
);
|
||||
let rows: Vec<_> = conn.prepare(&sql)
|
||||
.and_then(|mut stmt| {
|
||||
stmt.query_map(
|
||||
rusqlite::params![since_ts, per_table_limit as i64],
|
||||
|row| Ok((
|
||||
row.get::<_, i64>(0)?,
|
||||
row.get::<_, i64>(1)?,
|
||||
row.get::<_, i64>(2)?,
|
||||
row.get::<_, i64>(3)?,
|
||||
get_content_bytes(row, 4),
|
||||
row.get::<_, i64>(5).unwrap_or(0),
|
||||
)),
|
||||
).map(|it| it.filter_map(|r| r.ok()).collect())
|
||||
})
|
||||
.unwrap_or_default();
|
||||
|
||||
let mut result = Vec::new();
|
||||
for (local_id, local_type, ts, real_sender_id, content_bytes, ct) in rows {
|
||||
let content = decompress_message(&content_bytes, ct);
|
||||
let sender = sender_label(real_sender_id, &content, is_group, &uname2, &id2u, &names_map);
|
||||
let text = fmt_content(local_id, local_type, &content, is_group);
|
||||
result.push(json!({
|
||||
"chat": display2,
|
||||
"username": uname2,
|
||||
"is_group": is_group,
|
||||
"timestamp": ts,
|
||||
"time": fmt_time(ts, "%Y-%m-%d %H:%M"),
|
||||
"sender": sender,
|
||||
"content": text,
|
||||
"type": fmt_type(local_type),
|
||||
}));
|
||||
}
|
||||
Ok::<_, anyhow::Error>(result)
|
||||
}).await {
|
||||
Ok(Ok(v)) => v,
|
||||
Ok(Err(e)) => { eprintln!("[new-messages] skip {}: {}", tname_for_log, e); continue; }
|
||||
Err(e) => { eprintln!("[new-messages] task error: {}", e); continue; }
|
||||
};
|
||||
|
||||
all_msgs.extend(msgs);
|
||||
}
|
||||
}
|
||||
|
||||
all_msgs.sort_by_key(|m| m["timestamp"].as_i64().unwrap_or(0));
|
||||
all_msgs.truncate(limit);
|
||||
|
||||
// 5. 重建 new_state,防止全局 limit 截断导致消息永久丢失:
|
||||
// - 未变化的会话:沿用 session.db 的 last_timestamp
|
||||
// - 变化但全被截断(无消息在最终结果中):保留旧 since_ts,下次重试
|
||||
// - 变化且有消息返回:推进到该会话在结果中的最大 timestamp
|
||||
let mut new_state = session_ts_map;
|
||||
// 先把 changed 会话重置回旧 since_ts
|
||||
for (uname, _) in &changed {
|
||||
let old_ts = state.as_ref()
|
||||
.and_then(|m| m.get(uname))
|
||||
.copied()
|
||||
.unwrap_or(fallback_ts);
|
||||
new_state.insert(uname.clone(), old_ts);
|
||||
}
|
||||
// 再根据实际返回的消息向前推进
|
||||
for m in &all_msgs {
|
||||
if let (Some(uname), Some(ts)) = (m["username"].as_str(), m["timestamp"].as_i64()) {
|
||||
let e = new_state.entry(uname.to_string()).or_insert(0);
|
||||
if ts > *e { *e = ts; }
|
||||
}
|
||||
}
|
||||
|
||||
Ok(json!({
|
||||
"count": all_msgs.len(),
|
||||
"messages": all_msgs,
|
||||
"new_state": new_state,
|
||||
}))
|
||||
}
|
||||
|
||||
/// 查询收藏内容(favorite/favorite.db 的 fav_db_item 表)
|
||||
pub async fn q_favorites(
|
||||
db: &DbCache,
|
||||
limit: usize,
|
||||
fav_type: Option<i64>,
|
||||
query: Option<String>,
|
||||
) -> Result<Value> {
|
||||
let path = db.get("favorite/favorite.db").await?
|
||||
.context("找不到 favorite.db,请确认微信数据目录")?;
|
||||
|
||||
let rows: Vec<Value> = tokio::task::spawn_blocking(move || {
|
||||
let conn = Connection::open(&path)?;
|
||||
|
||||
let mut clauses: Vec<&'static str> = Vec::new();
|
||||
let mut params: Vec<Box<dyn rusqlite::types::ToSql>> = Vec::new();
|
||||
|
||||
if let Some(t) = fav_type {
|
||||
clauses.push("type = ?");
|
||||
params.push(Box::new(t));
|
||||
}
|
||||
let like_str: Option<String> = query.map(|q| {
|
||||
let esc = q.replace('\\', "\\\\").replace('%', "\\%").replace('_', "\\_");
|
||||
format!("%{}%", esc)
|
||||
});
|
||||
if let Some(ref s) = like_str {
|
||||
clauses.push("content LIKE ? ESCAPE '\\'");
|
||||
params.push(Box::new(s.clone()));
|
||||
}
|
||||
|
||||
let where_clause = if clauses.is_empty() {
|
||||
String::new()
|
||||
} else {
|
||||
format!("WHERE {}", clauses.join(" AND "))
|
||||
};
|
||||
params.push(Box::new(limit as i64));
|
||||
|
||||
let sql = format!(
|
||||
"SELECT local_id, type, update_time, content, fromusr, realchatname
|
||||
FROM fav_db_item {} ORDER BY update_time DESC LIMIT ?",
|
||||
where_clause
|
||||
);
|
||||
|
||||
let params_ref: Vec<&dyn rusqlite::types::ToSql> = params.iter().map(|p| p.as_ref()).collect();
|
||||
let mut stmt = conn.prepare(&sql)?;
|
||||
let rows: Vec<Value> = stmt.query_map(params_ref.as_slice(), |row| {
|
||||
Ok((
|
||||
row.get::<_, i64>(0).unwrap_or(0),
|
||||
row.get::<_, i64>(1).unwrap_or(0),
|
||||
row.get::<_, i64>(2).unwrap_or(0),
|
||||
row.get::<_, String>(3).unwrap_or_default(),
|
||||
row.get::<_, String>(4).unwrap_or_default(),
|
||||
row.get::<_, String>(5).unwrap_or_default(),
|
||||
))
|
||||
})?
|
||||
.filter_map(|r| r.ok())
|
||||
.map(|(local_id, ftype, ts, content, fromusr, chatname)| {
|
||||
let type_str = match ftype {
|
||||
1 => "文本",
|
||||
2 => "图片",
|
||||
5 => "文章",
|
||||
19 => "名片",
|
||||
20 => "视频",
|
||||
_ => "其他",
|
||||
};
|
||||
// 安全截断(按 Unicode 字符而非字节)
|
||||
let preview: String = content.chars().take(100).collect();
|
||||
let preview = if content.chars().count() > 100 {
|
||||
format!("{}...", preview)
|
||||
} else {
|
||||
preview
|
||||
};
|
||||
// WeChat 部分版本的 update_time 为毫秒,10位以上判定为毫秒后转秒
|
||||
let ts_secs = if ts > 9_999_999_999 { ts / 1000 } else { ts };
|
||||
json!({
|
||||
"id": local_id,
|
||||
"type": type_str,
|
||||
"type_num": ftype,
|
||||
"time": fmt_time(ts_secs, "%Y-%m-%d %H:%M"),
|
||||
"timestamp": ts_secs,
|
||||
"preview": preview,
|
||||
"from": fromusr,
|
||||
"chat": chatname,
|
||||
})
|
||||
})
|
||||
.collect();
|
||||
|
||||
Ok::<_, anyhow::Error>(rows)
|
||||
}).await??;
|
||||
|
||||
Ok(json!({
|
||||
"count": rows.len(),
|
||||
"items": rows,
|
||||
}))
|
||||
}
|
||||
|
||||
/// 聊天统计:消息总数、类型分布、发言排行、24小时分布
|
||||
pub async fn q_stats(
|
||||
db: &DbCache,
|
||||
names: &Names,
|
||||
chat: &str,
|
||||
since: Option<i64>,
|
||||
until: Option<i64>,
|
||||
) -> Result<Value> {
|
||||
let username = resolve_username(chat, names)
|
||||
.with_context(|| format!("找不到联系人: {}", chat))?;
|
||||
let display = names.display(&username);
|
||||
let is_group = username.contains("@chatroom");
|
||||
|
||||
let tables = find_msg_tables(db, names, &username).await?;
|
||||
if tables.is_empty() {
|
||||
anyhow::bail!("找不到 {} 的消息记录", display);
|
||||
}
|
||||
|
||||
// 跨所有分片 DB 累计统计
|
||||
let mut total: i64 = 0;
|
||||
let mut type_counts: HashMap<String, i64> = HashMap::new();
|
||||
let mut sender_counts: HashMap<String, i64> = HashMap::new();
|
||||
let mut hour_counts = [0i64; 24];
|
||||
|
||||
for (db_path, table_name) in &tables {
|
||||
let path = db_path.clone();
|
||||
let tname = table_name.clone();
|
||||
let uname = username.clone();
|
||||
let is_group2 = is_group;
|
||||
let names_map = names.map.clone();
|
||||
|
||||
// 用 SQL GROUP BY 在数据库侧聚合,避免把全量消息内容加载进内存
|
||||
let result: (i64, HashMap<String, i64>, HashMap<String, i64>, [i64; 24]) =
|
||||
tokio::task::spawn_blocking(move || {
|
||||
let conn = Connection::open(&path)?;
|
||||
let id2u = load_id2u(&conn);
|
||||
|
||||
let mut clauses = Vec::new();
|
||||
let mut params: Vec<Box<dyn rusqlite::types::ToSql>> = Vec::new();
|
||||
if let Some(s) = since {
|
||||
clauses.push("create_time >= ?");
|
||||
params.push(Box::new(s));
|
||||
}
|
||||
if let Some(u) = until {
|
||||
clauses.push("create_time <= ?");
|
||||
params.push(Box::new(u));
|
||||
}
|
||||
let where_clause = if clauses.is_empty() {
|
||||
String::new()
|
||||
} else {
|
||||
format!("WHERE {}", clauses.join(" AND "))
|
||||
};
|
||||
let params_ref: Vec<&dyn rusqlite::types::ToSql> =
|
||||
params.iter().map(|p| p.as_ref()).collect();
|
||||
|
||||
// 1. 总数
|
||||
let count: i64 = conn.query_row(
|
||||
&format!("SELECT COUNT(*) FROM [{}] {}", tname, where_clause),
|
||||
params_ref.as_slice(),
|
||||
|row| row.get(0),
|
||||
).unwrap_or(0);
|
||||
|
||||
// 2. 类型分布:SQL GROUP BY,不加载消息内容
|
||||
let type_sql = format!(
|
||||
"SELECT (local_type & 0xFFFFFFFF), COUNT(*) FROM [{}] {} GROUP BY (local_type & 0xFFFFFFFF)",
|
||||
tname, where_clause
|
||||
);
|
||||
let mut type_c: HashMap<String, i64> = HashMap::new();
|
||||
if let Ok(mut stmt) = conn.prepare(&type_sql) {
|
||||
let _ = stmt.query_map(params_ref.as_slice(), |row| {
|
||||
Ok((row.get::<_, i64>(0)?, row.get::<_, i64>(1)?))
|
||||
}).map(|rows| {
|
||||
for r in rows.flatten() {
|
||||
*type_c.entry(fmt_type(r.0)).or_insert(0) += r.1;
|
||||
}
|
||||
});
|
||||
}
|
||||
|
||||
// 3. 小时分布:只取时间戳,不加载消息内容
|
||||
let hour_sql = format!(
|
||||
"SELECT create_time FROM [{}] {}",
|
||||
tname, where_clause
|
||||
);
|
||||
let mut hour_c = [0i64; 24];
|
||||
if let Ok(mut stmt) = conn.prepare(&hour_sql) {
|
||||
let _ = stmt.query_map(params_ref.as_slice(), |row| row.get::<_, i64>(0))
|
||||
.map(|rows| {
|
||||
for ts in rows.flatten() {
|
||||
if let Some(dt) = Local.timestamp_opt(ts, 0).single() {
|
||||
let h = dt.hour() as usize;
|
||||
if h < 24 { hour_c[h] += 1; }
|
||||
}
|
||||
}
|
||||
});
|
||||
}
|
||||
|
||||
// 4. 发言排行:只取 real_sender_id,不加载消息内容
|
||||
// where_clause 可能已含 WHERE,用 AND 追加而非重复写 WHERE
|
||||
let sender_filter = if where_clause.is_empty() {
|
||||
"WHERE real_sender_id > 0".to_string()
|
||||
} else {
|
||||
format!("{} AND real_sender_id > 0", where_clause)
|
||||
};
|
||||
let sender_sql = format!(
|
||||
"SELECT real_sender_id, COUNT(*) FROM [{}] {} GROUP BY real_sender_id",
|
||||
tname, sender_filter
|
||||
);
|
||||
let mut sender_c: HashMap<String, i64> = HashMap::new();
|
||||
if is_group2 {
|
||||
if let Ok(mut stmt) = conn.prepare(&sender_sql) {
|
||||
let _ = stmt.query_map(params_ref.as_slice(), |row| {
|
||||
Ok((row.get::<_, i64>(0)?, row.get::<_, i64>(1)?))
|
||||
}).map(|rows| {
|
||||
for (id, cnt) in rows.flatten() {
|
||||
if let Some(u) = id2u.get(&id) {
|
||||
if u != &uname {
|
||||
let name = names_map.get(u).cloned().unwrap_or_else(|| u.clone());
|
||||
*sender_c.entry(name).or_insert(0) += cnt;
|
||||
}
|
||||
}
|
||||
}
|
||||
});
|
||||
}
|
||||
}
|
||||
|
||||
Ok::<_, anyhow::Error>((count, type_c, sender_c, hour_c))
|
||||
}).await??;
|
||||
|
||||
let (count, type_c, sender_c, hour_c) = result;
|
||||
total += count;
|
||||
for (k, v) in type_c { *type_counts.entry(k).or_insert(0) += v; }
|
||||
for (k, v) in sender_c { *sender_counts.entry(k).or_insert(0) += v; }
|
||||
for i in 0..24 { hour_counts[i] += hour_c[i]; }
|
||||
}
|
||||
|
||||
// 类型分布,按数量降序
|
||||
let mut by_type: Vec<Value> = type_counts.iter()
|
||||
.map(|(t, c)| json!({ "type": t, "count": c }))
|
||||
.collect();
|
||||
by_type.sort_by_key(|v| std::cmp::Reverse(v["count"].as_i64().unwrap_or(0)));
|
||||
|
||||
// 发言排行,Top 10
|
||||
let mut top_senders: Vec<Value> = sender_counts.iter()
|
||||
.map(|(s, c)| json!({ "sender": s, "count": c }))
|
||||
.collect();
|
||||
top_senders.sort_by_key(|v| std::cmp::Reverse(v["count"].as_i64().unwrap_or(0)));
|
||||
top_senders.truncate(10);
|
||||
|
||||
// 24小时分布
|
||||
let by_hour: Vec<Value> = hour_counts.iter().enumerate()
|
||||
.map(|(h, c)| json!({ "hour": h, "count": c }))
|
||||
.collect();
|
||||
|
||||
Ok(json!({
|
||||
"chat": display,
|
||||
"username": username,
|
||||
"is_group": is_group,
|
||||
"total": total,
|
||||
"by_type": by_type,
|
||||
"top_senders": top_senders,
|
||||
"by_hour": by_hour,
|
||||
}))
|
||||
}
|
||||
|
||||
|
||||
+57
-54
@@ -1,10 +1,8 @@
|
||||
use anyhow::Result;
|
||||
use std::sync::Arc;
|
||||
use std::time::Duration;
|
||||
use tokio::io::{AsyncBufReadExt, AsyncWriteExt, BufReader};
|
||||
use tokio::sync::broadcast;
|
||||
|
||||
use crate::ipc::{Request, Response, WatchEvent};
|
||||
use crate::ipc::{Request, Response};
|
||||
use super::cache::DbCache;
|
||||
use super::query::Names;
|
||||
|
||||
@@ -12,12 +10,11 @@ use super::query::Names;
|
||||
pub async fn serve(
|
||||
db: Arc<DbCache>,
|
||||
names: Arc<std::sync::RwLock<Names>>,
|
||||
watch_tx: broadcast::Sender<WatchEvent>,
|
||||
) -> Result<()> {
|
||||
#[cfg(unix)]
|
||||
serve_unix(db, names, watch_tx).await?;
|
||||
serve_unix(db, names).await?;
|
||||
#[cfg(windows)]
|
||||
serve_windows(db, names, watch_tx).await?;
|
||||
serve_windows(db, names).await?;
|
||||
Ok(())
|
||||
}
|
||||
|
||||
@@ -25,7 +22,6 @@ pub async fn serve(
|
||||
async fn serve_unix(
|
||||
db: Arc<DbCache>,
|
||||
names: Arc<std::sync::RwLock<Names>>,
|
||||
watch_tx: broadcast::Sender<WatchEvent>,
|
||||
) -> Result<()> {
|
||||
use tokio::net::UnixListener;
|
||||
let sock_path = crate::config::sock_path();
|
||||
@@ -49,10 +45,9 @@ async fn serve_unix(
|
||||
let (stream, _) = listener.accept().await?;
|
||||
let db2 = Arc::clone(&db);
|
||||
let names2 = Arc::clone(&names);
|
||||
let tx2 = watch_tx.clone();
|
||||
|
||||
tokio::spawn(async move {
|
||||
if let Err(e) = handle_connection_unix(stream, db2, names2, tx2).await {
|
||||
if let Err(e) = handle_connection_unix(stream, db2, names2).await {
|
||||
eprintln!("[server] 连接处理错误: {}", e);
|
||||
}
|
||||
});
|
||||
@@ -64,7 +59,6 @@ async fn handle_connection_unix(
|
||||
stream: tokio::net::UnixStream,
|
||||
db: Arc<DbCache>,
|
||||
names: Arc<std::sync::RwLock<Names>>,
|
||||
watch_tx: broadcast::Sender<WatchEvent>,
|
||||
) -> Result<()> {
|
||||
let (reader, mut writer) = stream.into_split();
|
||||
let mut lines = BufReader::new(reader).lines();
|
||||
@@ -84,41 +78,8 @@ async fn handle_connection_unix(
|
||||
}
|
||||
};
|
||||
|
||||
match req {
|
||||
Request::Watch => {
|
||||
// 流式模式:持续推送事件
|
||||
let mut rx = watch_tx.subscribe();
|
||||
let connected = WatchEvent::connected();
|
||||
writer.write_all(connected.to_json_line()?.as_bytes()).await?;
|
||||
|
||||
loop {
|
||||
tokio::select! {
|
||||
event = rx.recv() => {
|
||||
match event {
|
||||
Ok(e) => {
|
||||
if writer.write_all(e.to_json_line()?.as_bytes()).await.is_err() {
|
||||
break;
|
||||
}
|
||||
}
|
||||
Err(broadcast::error::RecvError::Lagged(_)) => continue,
|
||||
Err(_) => break,
|
||||
}
|
||||
}
|
||||
_ = tokio::time::sleep(Duration::from_secs(30)) => {
|
||||
// 心跳
|
||||
let hb = WatchEvent::heartbeat();
|
||||
if writer.write_all(hb.to_json_line()?.as_bytes()).await.is_err() {
|
||||
break;
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
other => {
|
||||
let resp = dispatch(other, &db, &names).await;
|
||||
writer.write_all(resp.to_json_line()?.as_bytes()).await?;
|
||||
}
|
||||
}
|
||||
let resp = dispatch(req, &db, &names).await;
|
||||
writer.write_all(resp.to_json_line()?.as_bytes()).await?;
|
||||
Ok(())
|
||||
}
|
||||
|
||||
@@ -126,7 +87,6 @@ async fn handle_connection_unix(
|
||||
async fn serve_windows(
|
||||
db: Arc<DbCache>,
|
||||
names: Arc<std::sync::RwLock<Names>>,
|
||||
watch_tx: broadcast::Sender<WatchEvent>,
|
||||
) -> Result<()> {
|
||||
use interprocess::local_socket::{
|
||||
tokio::prelude::*, GenericNamespaced, ListenerOptions,
|
||||
@@ -143,10 +103,9 @@ async fn serve_windows(
|
||||
let conn = listener.accept().await?;
|
||||
let db2 = Arc::clone(&db);
|
||||
let names2 = Arc::clone(&names);
|
||||
let tx2 = watch_tx.clone();
|
||||
|
||||
tokio::spawn(async move {
|
||||
if let Err(e) = handle_connection_generic(conn, db2, names2, tx2).await {
|
||||
if let Err(e) = handle_connection_generic(conn, db2, names2).await {
|
||||
eprintln!("[server] 连接处理错误: {}", e);
|
||||
}
|
||||
});
|
||||
@@ -164,7 +123,6 @@ async fn dispatch(
|
||||
match req {
|
||||
Ping => Response::ok(serde_json::json!({ "pong": true })),
|
||||
Sessions { limit } => {
|
||||
// 在 await 前获取并复制所需数据,避免 RwLockGuard 跨 await
|
||||
let names_snapshot = match clone_names(names) {
|
||||
Ok(n) => n,
|
||||
Err(e) => return Response::err(e),
|
||||
@@ -174,22 +132,22 @@ async fn dispatch(
|
||||
Err(e) => Response::err(e.to_string()),
|
||||
}
|
||||
}
|
||||
History { chat, limit, offset, since, until } => {
|
||||
History { chat, limit, offset, since, until, msg_type } => {
|
||||
let names_snapshot = match clone_names(names) {
|
||||
Ok(n) => n,
|
||||
Err(e) => return Response::err(e),
|
||||
};
|
||||
match query::q_history(db, &names_snapshot, &chat, limit, offset, since, until).await {
|
||||
match query::q_history(db, &names_snapshot, &chat, limit, offset, since, until, msg_type).await {
|
||||
Ok(v) => Response::ok(v),
|
||||
Err(e) => Response::err(e.to_string()),
|
||||
}
|
||||
}
|
||||
Search { keyword, chats, limit, since, until } => {
|
||||
Search { keyword, chats, limit, since, until, msg_type } => {
|
||||
let names_snapshot = match clone_names(names) {
|
||||
Ok(n) => n,
|
||||
Err(e) => return Response::err(e),
|
||||
};
|
||||
match query::q_search(db, &names_snapshot, &keyword, chats, limit, since, until).await {
|
||||
match query::q_search(db, &names_snapshot, &keyword, chats, limit, since, until, msg_type).await {
|
||||
Ok(v) => Response::ok(v),
|
||||
Err(e) => Response::err(e.to_string()),
|
||||
}
|
||||
@@ -204,7 +162,52 @@ async fn dispatch(
|
||||
Err(e) => Response::err(e.to_string()),
|
||||
}
|
||||
}
|
||||
Watch => Response::err("Watch 命令不应通过 dispatch 处理"),
|
||||
Unread { limit } => {
|
||||
let names_snapshot = match clone_names(names) {
|
||||
Ok(n) => n,
|
||||
Err(e) => return Response::err(e),
|
||||
};
|
||||
match query::q_unread(db, &names_snapshot, limit).await {
|
||||
Ok(v) => Response::ok(v),
|
||||
Err(e) => Response::err(e.to_string()),
|
||||
}
|
||||
}
|
||||
Members { chat } => {
|
||||
let names_snapshot = match clone_names(names) {
|
||||
Ok(n) => n,
|
||||
Err(e) => return Response::err(e),
|
||||
};
|
||||
match query::q_members(db, &names_snapshot, &chat).await {
|
||||
Ok(v) => Response::ok(v),
|
||||
Err(e) => Response::err(e.to_string()),
|
||||
}
|
||||
}
|
||||
NewMessages { state, limit } => {
|
||||
let names_snapshot = match clone_names(names) {
|
||||
Ok(n) => n,
|
||||
Err(e) => return Response::err(e),
|
||||
};
|
||||
match query::q_new_messages(db, &names_snapshot, state, limit).await {
|
||||
Ok(v) => Response::ok(v),
|
||||
Err(e) => Response::err(e.to_string()),
|
||||
}
|
||||
}
|
||||
Favorites { limit, fav_type, query } => {
|
||||
match query::q_favorites(db, limit, fav_type, query).await {
|
||||
Ok(v) => Response::ok(v),
|
||||
Err(e) => Response::err(e.to_string()),
|
||||
}
|
||||
}
|
||||
Stats { chat, since, until } => {
|
||||
let names_snapshot = match clone_names(names) {
|
||||
Ok(n) => n,
|
||||
Err(e) => return Response::err(e),
|
||||
};
|
||||
match query::q_stats(db, &names_snapshot, &chat, since, until).await {
|
||||
Ok(v) => Response::ok(v),
|
||||
Err(e) => Response::err(e.to_string()),
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
Reference in New Issue
Block a user