mirror of
https://github.com/jackwener/wx-cli.git
synced 2026-10-08 21:05:44 +00:00
feat(meta): expose freshness coverage in query output
This commit is contained in:
+151
-52
@@ -23,6 +23,29 @@ struct CacheEntry {
|
||||
decrypted_path: PathBuf,
|
||||
}
|
||||
|
||||
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
|
||||
pub enum CacheMode {
|
||||
CacheHit,
|
||||
WalIncremental,
|
||||
FullDecrypt,
|
||||
}
|
||||
|
||||
impl CacheMode {
|
||||
pub fn as_str(self) -> &'static str {
|
||||
match self {
|
||||
CacheMode::CacheHit => "cache_hit",
|
||||
CacheMode::WalIncremental => "wal_incremental",
|
||||
CacheMode::FullDecrypt => "full_decrypt",
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
#[derive(Debug, Clone)]
|
||||
pub struct CacheResolve {
|
||||
pub path: PathBuf,
|
||||
pub mode: CacheMode,
|
||||
}
|
||||
|
||||
/// 解密后数据库的 mtime-aware 缓存
|
||||
///
|
||||
/// 当数据库文件(.db)或 WAL 文件(.db-wal)的 mtime 发生变化时,
|
||||
@@ -33,13 +56,11 @@ pub struct DbCache {
|
||||
mtime_file: PathBuf,
|
||||
all_keys: HashMap<String, String>, // rel_key -> enc_key(hex)
|
||||
inner: Arc<Mutex<HashMap<String, CacheEntry>>>,
|
||||
last_modes: Arc<Mutex<HashMap<String, CacheMode>>>,
|
||||
}
|
||||
|
||||
impl DbCache {
|
||||
pub async fn new(
|
||||
db_dir: PathBuf,
|
||||
all_keys: HashMap<String, String>,
|
||||
) -> Result<Self> {
|
||||
pub async fn new(db_dir: PathBuf, all_keys: HashMap<String, String>) -> Result<Self> {
|
||||
Self::with_dirs(db_dir, config::cache_dir(), config::mtime_file(), all_keys).await
|
||||
}
|
||||
|
||||
@@ -58,6 +79,7 @@ impl DbCache {
|
||||
mtime_file,
|
||||
all_keys,
|
||||
inner: Arc::new(Mutex::new(HashMap::new())),
|
||||
last_modes: Arc::new(Mutex::new(HashMap::new())),
|
||||
};
|
||||
|
||||
cache.load_persistent().await;
|
||||
@@ -70,6 +92,18 @@ impl DbCache {
|
||||
&self.db_dir
|
||||
}
|
||||
|
||||
#[allow(dead_code)]
|
||||
pub async fn last_mode(&self, rel_key: &str) -> Option<CacheMode> {
|
||||
self.last_modes.lock().await.get(rel_key).copied()
|
||||
}
|
||||
|
||||
async fn stamp_mode(&self, rel_key: &str, mode: CacheMode) {
|
||||
self.last_modes
|
||||
.lock()
|
||||
.await
|
||||
.insert(rel_key.to_string(), mode);
|
||||
}
|
||||
|
||||
fn cache_file_path(&self, rel_key: &str) -> PathBuf {
|
||||
let hash = format!("{:x}", md5::compute(rel_key.as_bytes()));
|
||||
self.cache_dir.join(format!("{}.db", hash))
|
||||
@@ -94,23 +128,34 @@ impl DbCache {
|
||||
if !dec_path.exists() {
|
||||
continue;
|
||||
}
|
||||
let db_path = self.db_dir.join(rel_key.replace('\\', std::path::MAIN_SEPARATOR_STR).replace('/', std::path::MAIN_SEPARATOR_STR));
|
||||
let db_path = self.db_dir.join(
|
||||
rel_key
|
||||
.replace('\\', std::path::MAIN_SEPARATOR_STR)
|
||||
.replace('/', std::path::MAIN_SEPARATOR_STR),
|
||||
);
|
||||
let wal_path = wal_path_for(&db_path);
|
||||
|
||||
let db_mt = mtime_nanos(&db_path);
|
||||
let _wal_mt = if wal_path.exists() { mtime_nanos(&wal_path) } else { 0 };
|
||||
let _wal_mt = if wal_path.exists() {
|
||||
mtime_nanos(&wal_path)
|
||||
} else {
|
||||
0
|
||||
};
|
||||
|
||||
// 只要主 .db 没变,就把 cached 产物载回来。
|
||||
// 如果 WAL mtime 变了,后续 `get()` 会自动走 Path 2:在已有 cached DB 上增量 apply_wal,
|
||||
// 而不是 daemon 重启后第一条请求又退回全量解密。
|
||||
if db_mt == entry.db_mt {
|
||||
inner.insert(rel_key.clone(), CacheEntry {
|
||||
db_mtime: db_mt,
|
||||
// 保留"cached 产物构建时看到的 wal_mtime",让 `get()` 去比较当前 WAL
|
||||
// 是否发生了变化,从而决定 exact-hit 还是 WAL 增量。
|
||||
wal_mtime: entry.wal_mt,
|
||||
decrypted_path: dec_path,
|
||||
});
|
||||
inner.insert(
|
||||
rel_key.clone(),
|
||||
CacheEntry {
|
||||
db_mtime: db_mt,
|
||||
// 保留"cached 产物构建时看到的 wal_mtime",让 `get()` 去比较当前 WAL
|
||||
// 是否发生了变化,从而决定 exact-hit 还是 WAL 增量。
|
||||
wal_mtime: entry.wal_mt,
|
||||
decrypted_path: dec_path,
|
||||
},
|
||||
);
|
||||
reused += 1;
|
||||
}
|
||||
}
|
||||
@@ -123,13 +168,19 @@ impl DbCache {
|
||||
async fn save_persistent(&self) {
|
||||
let mtime_file = &self.mtime_file;
|
||||
let inner = self.inner.lock().await;
|
||||
let data: HashMap<String, MtimeEntry> = inner.iter().map(|(k, v)| {
|
||||
(k.clone(), MtimeEntry {
|
||||
db_mt: v.db_mtime,
|
||||
wal_mt: v.wal_mtime,
|
||||
path: v.decrypted_path.to_string_lossy().into_owned(),
|
||||
let data: HashMap<String, MtimeEntry> = inner
|
||||
.iter()
|
||||
.map(|(k, v)| {
|
||||
(
|
||||
k.clone(),
|
||||
MtimeEntry {
|
||||
db_mt: v.db_mtime,
|
||||
wal_mt: v.wal_mtime,
|
||||
path: v.decrypted_path.to_string_lossy().into_owned(),
|
||||
},
|
||||
)
|
||||
})
|
||||
}).collect();
|
||||
.collect();
|
||||
drop(inner);
|
||||
|
||||
if let Ok(json) = serde_json::to_string_pretty(&data) {
|
||||
@@ -148,14 +199,19 @@ impl DbCache {
|
||||
/// WeChat 在写消息时只 append WAL(除非触发 checkpoint),因此 path 2 是常态;
|
||||
/// 这条路径把"每次请求都全量解密 ~1.8GB DB(~120s)"压到"只解 WAL 帧(典型 < 10s)"。
|
||||
pub async fn get(&self, rel_key: &str) -> Result<Option<PathBuf>> {
|
||||
Ok(self.get_with_mode(rel_key).await?.map(|r| r.path))
|
||||
}
|
||||
|
||||
pub async fn get_with_mode(&self, rel_key: &str) -> Result<Option<CacheResolve>> {
|
||||
let enc_key_hex = match self.all_keys.get(rel_key) {
|
||||
Some(k) => k.clone(),
|
||||
None => return Ok(None),
|
||||
};
|
||||
|
||||
let db_path = self.db_dir.join(
|
||||
rel_key.replace('\\', std::path::MAIN_SEPARATOR_STR)
|
||||
.replace('/', std::path::MAIN_SEPARATOR_STR)
|
||||
rel_key
|
||||
.replace('\\', std::path::MAIN_SEPARATOR_STR)
|
||||
.replace('/', std::path::MAIN_SEPARATOR_STR),
|
||||
);
|
||||
if !db_path.exists() {
|
||||
return Ok(None);
|
||||
@@ -163,21 +219,29 @@ impl DbCache {
|
||||
|
||||
let wal_path = wal_path_for(&db_path);
|
||||
let db_mt = mtime_nanos(&db_path);
|
||||
let wal_mt = if wal_path.exists() { mtime_nanos(&wal_path) } else { 0 };
|
||||
let wal_mt = if wal_path.exists() {
|
||||
mtime_nanos(&wal_path)
|
||||
} else {
|
||||
0
|
||||
};
|
||||
|
||||
let cached = {
|
||||
let inner = self.inner.lock().await;
|
||||
inner.get(rel_key).cloned()
|
||||
};
|
||||
|
||||
let enc_key_bytes = hex_to_32bytes(&enc_key_hex)
|
||||
.with_context(|| format!("密钥格式错误: {}", rel_key))?;
|
||||
let enc_key_bytes =
|
||||
hex_to_32bytes(&enc_key_hex).with_context(|| format!("密钥格式错误: {}", rel_key))?;
|
||||
|
||||
// Path 1 / Path 2:主 .db mtime 未变且 cached 产物仍在
|
||||
if let Some(entry) = cached.as_ref() {
|
||||
if entry.db_mtime == db_mt && entry.decrypted_path.exists() {
|
||||
if entry.wal_mtime == wal_mt {
|
||||
return Ok(Some(entry.decrypted_path.clone()));
|
||||
self.stamp_mode(rel_key, CacheMode::CacheHit).await;
|
||||
return Ok(Some(CacheResolve {
|
||||
path: entry.decrypted_path.clone(),
|
||||
mode: CacheMode::CacheHit,
|
||||
}));
|
||||
}
|
||||
|
||||
// Path 2: WAL-only 变化 → 在 cached 产物上重新 apply_wal
|
||||
@@ -190,20 +254,32 @@ impl DbCache {
|
||||
let key_copy = enc_key_bytes;
|
||||
tokio::task::spawn_blocking(move || {
|
||||
wal::apply_wal(&wal_path2, &out_path2, &key_copy)
|
||||
}).await??;
|
||||
})
|
||||
.await??;
|
||||
}
|
||||
eprintln!("[cache] WAL 增量 {} ({}ms)", rel_key, t0.elapsed().as_millis());
|
||||
eprintln!(
|
||||
"[cache] WAL 增量 {} ({}ms)",
|
||||
rel_key,
|
||||
t0.elapsed().as_millis()
|
||||
);
|
||||
|
||||
{
|
||||
let mut inner = self.inner.lock().await;
|
||||
inner.insert(rel_key.to_string(), CacheEntry {
|
||||
db_mtime: db_mt,
|
||||
wal_mtime: wal_mt,
|
||||
decrypted_path: out_path.clone(),
|
||||
});
|
||||
inner.insert(
|
||||
rel_key.to_string(),
|
||||
CacheEntry {
|
||||
db_mtime: db_mt,
|
||||
wal_mtime: wal_mt,
|
||||
decrypted_path: out_path.clone(),
|
||||
},
|
||||
);
|
||||
}
|
||||
self.stamp_mode(rel_key, CacheMode::WalIncremental).await;
|
||||
self.save_persistent().await;
|
||||
return Ok(Some(out_path));
|
||||
return Ok(Some(CacheResolve {
|
||||
path: out_path,
|
||||
mode: CacheMode::WalIncremental,
|
||||
}));
|
||||
}
|
||||
}
|
||||
|
||||
@@ -213,39 +289,52 @@ impl DbCache {
|
||||
let db_path2 = db_path.clone();
|
||||
let out_path2 = out_path.clone();
|
||||
let key_copy = enc_key_bytes;
|
||||
tokio::task::spawn_blocking(move || {
|
||||
crypto::full_decrypt(&db_path2, &out_path2, &key_copy)
|
||||
}).await??;
|
||||
tokio::task::spawn_blocking(move || crypto::full_decrypt(&db_path2, &out_path2, &key_copy))
|
||||
.await??;
|
||||
|
||||
if wal_path.exists() {
|
||||
let out_path3 = out_path.clone();
|
||||
let wal_path3 = wal_path.clone();
|
||||
let key_copy2 = enc_key_bytes;
|
||||
tokio::task::spawn_blocking(move || {
|
||||
wal::apply_wal(&wal_path3, &out_path3, &key_copy2)
|
||||
}).await??;
|
||||
tokio::task::spawn_blocking(move || wal::apply_wal(&wal_path3, &out_path3, &key_copy2))
|
||||
.await??;
|
||||
}
|
||||
|
||||
eprintln!("[cache] 全量解密 {} ({}ms)", rel_key, t0.elapsed().as_millis());
|
||||
eprintln!(
|
||||
"[cache] 全量解密 {} ({}ms)",
|
||||
rel_key,
|
||||
t0.elapsed().as_millis()
|
||||
);
|
||||
|
||||
{
|
||||
let mut inner = self.inner.lock().await;
|
||||
inner.insert(rel_key.to_string(), CacheEntry {
|
||||
db_mtime: db_mt,
|
||||
wal_mtime: wal_mt,
|
||||
decrypted_path: out_path.clone(),
|
||||
});
|
||||
inner.insert(
|
||||
rel_key.to_string(),
|
||||
CacheEntry {
|
||||
db_mtime: db_mt,
|
||||
wal_mtime: wal_mt,
|
||||
decrypted_path: out_path.clone(),
|
||||
},
|
||||
);
|
||||
}
|
||||
self.stamp_mode(rel_key, CacheMode::FullDecrypt).await;
|
||||
|
||||
self.save_persistent().await;
|
||||
Ok(Some(out_path))
|
||||
Ok(Some(CacheResolve {
|
||||
path: out_path,
|
||||
mode: CacheMode::FullDecrypt,
|
||||
}))
|
||||
}
|
||||
}
|
||||
|
||||
pub(super) fn mtime_nanos(path: &Path) -> u64 {
|
||||
std::fs::metadata(path)
|
||||
.and_then(|m| m.modified())
|
||||
.map(|t| t.duration_since(std::time::UNIX_EPOCH).unwrap_or_default().as_nanos() as u64)
|
||||
.map(|t| {
|
||||
t.duration_since(std::time::UNIX_EPOCH)
|
||||
.unwrap_or_default()
|
||||
.as_nanos() as u64
|
||||
})
|
||||
.unwrap_or(0)
|
||||
}
|
||||
|
||||
@@ -273,8 +362,7 @@ mod tests {
|
||||
use super::*;
|
||||
|
||||
/// 64 字符 hex(不需要是真 SQLCipher key — 仅用来证明"是否触发了 full_decrypt")
|
||||
const FAKE_KEY_HEX: &str =
|
||||
"0000000000000000000000000000000000000000000000000000000000000000";
|
||||
const FAKE_KEY_HEX: &str = "0000000000000000000000000000000000000000000000000000000000000000";
|
||||
|
||||
/// 路径区分约定:
|
||||
/// - 完全 hit / WAL 增量 → `decrypted_path` **内容不变**
|
||||
@@ -337,7 +425,11 @@ mod tests {
|
||||
let (cache, _db_path, decrypted_path, _mtime_file, rel_key) =
|
||||
setup_seeded_cache("exact").await;
|
||||
|
||||
let p = cache.get(&rel_key).await.unwrap().expect("cache should hit");
|
||||
let p = cache
|
||||
.get(&rel_key)
|
||||
.await
|
||||
.unwrap()
|
||||
.expect("cache should hit");
|
||||
assert_eq!(p, decrypted_path);
|
||||
|
||||
// 完全 hit → cached file 内容不应被改
|
||||
@@ -387,7 +479,10 @@ mod tests {
|
||||
// 第一次:完全 hit
|
||||
let p1 = cache.get(&rel_key).await.unwrap().expect("first get hits");
|
||||
assert_eq!(p1, decrypted_path);
|
||||
assert_eq!(std::fs::read(&decrypted_path).unwrap(), ORIGINAL_CACHED_BYTES);
|
||||
assert_eq!(
|
||||
std::fs::read(&decrypted_path).unwrap(),
|
||||
ORIGINAL_CACHED_BYTES
|
||||
);
|
||||
|
||||
// bump WAL mtime(重写仍 31 bytes,apply_wal 仍 noop)
|
||||
std::thread::sleep(std::time::Duration::from_millis(20));
|
||||
@@ -486,7 +581,11 @@ mod tests {
|
||||
.await
|
||||
.unwrap();
|
||||
|
||||
let p = cache.get(&rel_key).await.unwrap().expect("cache should reuse persisted DB");
|
||||
let p = cache
|
||||
.get(&rel_key)
|
||||
.await
|
||||
.unwrap()
|
||||
.expect("cache should reuse persisted DB");
|
||||
assert_eq!(p, decrypted_path);
|
||||
let body = std::fs::read(&decrypted_path).unwrap();
|
||||
assert_eq!(
|
||||
|
||||
@@ -0,0 +1,263 @@
|
||||
//! Freshness metadata appended to every q_* response.
|
||||
//!
|
||||
//! 背景:`all_keys.json` 是 `wx init` 时的快照。WeChat 在 daemon 启动后随时可能创建
|
||||
//! 新的 `message_N.db` 分片;如果只信任 init 时收到的 `msg_db_keys` 列表,新分片里
|
||||
//! 的数据对 daemon 完全不可见 → 调用方拿到的是看似正常但缺数据的结果("stale")。
|
||||
//!
|
||||
//! 本模块的职责:
|
||||
//! 1. 提供 `Meta` 结构体,由各 `q_*` 函数填充后塞进 response(顶层 `meta` 字段)。
|
||||
//! 2. 提供 `discover_unknown_shards(db_dir, msg_db_keys)`:扫描磁盘上当前真实存在的
|
||||
//! `message/message_*.db` 文件,diff 出 daemon 未持有 enc_key 的"未知分片"列表。
|
||||
//! 3. 集中 `MetaStatus` 的判定规则,避免 8 个 q_* 各自判,规则漂移。
|
||||
|
||||
use serde::Serialize;
|
||||
use std::collections::HashMap;
|
||||
use std::path::Path;
|
||||
|
||||
/// 每条 q_* 响应附带的"新鲜度元数据"。
|
||||
///
|
||||
/// 序列化为 JSON 时,所有 `Option` 字段在 `None` 时省略,让最常见的命令调用
|
||||
/// 输出尽量短;重负载字段(per_shard_*、shard_paths)默认不填,由 CLI 层
|
||||
/// 通过 `--debug-source` 等开关显式请求时才放进来。
|
||||
#[derive(Debug, Clone, Serialize, Default)]
|
||||
pub struct Meta {
|
||||
/// 命中数据中最新一条的 create_time(unix 秒)。
|
||||
/// `q_history` / `q_search` / `q_new_messages` 等基于 Msg_ 表的查询都应填。
|
||||
/// `q_sessions` / `q_unread` 这类基于 SessionTable 的查询填会话维度的最新 ts。
|
||||
#[serde(skip_serializing_if = "Option::is_none")]
|
||||
pub chat_latest_timestamp: Option<i64>,
|
||||
|
||||
/// 上面那条最新消息所在的分片 rel_key(`message/message_3.db`)。
|
||||
/// 让 agent 一眼看出"当前命中的最新数据来自哪个分片"。
|
||||
#[serde(skip_serializing_if = "Option::is_none")]
|
||||
pub chat_latest_db: Option<String>,
|
||||
|
||||
/// 该 chat 在 `session.db.SessionTable.last_timestamp` 里的值(如果可读)。
|
||||
/// 这是 WeChat 自己写的"最近一条消息时间",与上面 `chat_latest_timestamp` 比较
|
||||
/// 即可发现"session 说有更新但 history 没读到" → 漏分片。
|
||||
#[serde(skip_serializing_if = "Option::is_none")]
|
||||
pub session_last_timestamp: Option<i64>,
|
||||
|
||||
/// 本次查询实际遍历的分片数(即 `names.msg_db_keys.len()` 的子集;包括命中 0 行的)。
|
||||
pub shards_scanned: usize,
|
||||
|
||||
/// 本次查询里至少返回了 1 行的分片数。
|
||||
pub shards_hit: usize,
|
||||
|
||||
/// 磁盘上存在但 daemon 没有 enc_key 的分片 rel_key 列表。
|
||||
/// 非空 ⇒ `wx init` 之后 WeChat 又分裂了新分片 → 必须重跑 `wx init`。
|
||||
pub unknown_shards: Vec<String>,
|
||||
|
||||
/// 由上述字段派生出的总体状态,CLI / agent 主要看这一个。
|
||||
pub status: MetaStatus,
|
||||
|
||||
// 重负载/调试字段:默认不填,CLI 层显式开启
|
||||
#[serde(skip_serializing_if = "Option::is_none")]
|
||||
pub per_shard_latest: Option<HashMap<String, i64>>,
|
||||
#[serde(skip_serializing_if = "Option::is_none")]
|
||||
pub cache_mode_per_shard: Option<HashMap<String, String>>,
|
||||
#[serde(skip_serializing_if = "Option::is_none")]
|
||||
pub shard_paths: Option<HashMap<String, String>>,
|
||||
}
|
||||
|
||||
#[derive(Debug, Clone, Copy, Serialize, PartialEq, Eq, Default)]
|
||||
#[serde(rename_all = "snake_case")]
|
||||
pub enum MetaStatus {
|
||||
#[default]
|
||||
Ok,
|
||||
PossiblyStale,
|
||||
PossiblyStaleUnknownShards,
|
||||
Windowed,
|
||||
}
|
||||
|
||||
impl MetaStatus {
|
||||
#[allow(dead_code)]
|
||||
pub fn as_str(&self) -> &'static str {
|
||||
match self {
|
||||
MetaStatus::Ok => "ok",
|
||||
MetaStatus::PossiblyStale => "possibly_stale",
|
||||
MetaStatus::PossiblyStaleUnknownShards => "possibly_stale_unknown_shards",
|
||||
MetaStatus::Windowed => "windowed",
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/// session 领先 history 多少秒就报 PossiblyStale。
|
||||
pub const STALE_THRESHOLD_SECS: i64 = 24 * 3600;
|
||||
|
||||
pub fn derive_status(
|
||||
chat_latest: Option<i64>,
|
||||
session_last: Option<i64>,
|
||||
unknown_shards: &[String],
|
||||
windowed: bool,
|
||||
) -> MetaStatus {
|
||||
if !unknown_shards.is_empty() {
|
||||
return MetaStatus::PossiblyStaleUnknownShards;
|
||||
}
|
||||
if windowed {
|
||||
return MetaStatus::Windowed;
|
||||
}
|
||||
match (chat_latest, session_last) {
|
||||
(Some(c), Some(s)) if s - c > STALE_THRESHOLD_SECS => MetaStatus::PossiblyStale,
|
||||
_ => MetaStatus::Ok,
|
||||
}
|
||||
}
|
||||
|
||||
pub fn discover_unknown_shards(db_dir: &Path, known: &[String]) -> Vec<String> {
|
||||
let known_set: std::collections::HashSet<String> =
|
||||
known.iter().map(|k| k.replace('\\', "/")).collect();
|
||||
|
||||
let msg_dir = db_dir.join("message");
|
||||
let entries = match std::fs::read_dir(&msg_dir) {
|
||||
Ok(it) => it,
|
||||
Err(_) => return Vec::new(),
|
||||
};
|
||||
|
||||
let mut unknown: Vec<String> = Vec::new();
|
||||
for entry in entries.flatten() {
|
||||
let name = entry.file_name();
|
||||
let Some(name_str) = name.to_str() else {
|
||||
continue;
|
||||
};
|
||||
if !is_message_shard(name_str) {
|
||||
continue;
|
||||
}
|
||||
let rel = format!("message/{}", name_str);
|
||||
if !known_set.contains(&rel) {
|
||||
unknown.push(rel);
|
||||
}
|
||||
}
|
||||
unknown.sort();
|
||||
unknown
|
||||
}
|
||||
|
||||
fn is_message_shard(file_name: &str) -> bool {
|
||||
if !file_name.starts_with("message_") || !file_name.ends_with(".db") {
|
||||
return false;
|
||||
}
|
||||
if file_name.contains("_fts") || file_name.contains("_resource") {
|
||||
return false;
|
||||
}
|
||||
let stem = &file_name["message_".len()..file_name.len() - ".db".len()];
|
||||
!stem.is_empty() && stem.chars().all(|c| c.is_ascii_digit())
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
use super::*;
|
||||
|
||||
#[test]
|
||||
fn is_message_shard_accepts_normal_shards() {
|
||||
assert!(is_message_shard("message_0.db"));
|
||||
assert!(is_message_shard("message_12.db"));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn is_message_shard_rejects_fts_and_resource() {
|
||||
assert!(!is_message_shard("message_0_fts.db"));
|
||||
assert!(!is_message_shard("message_fts.db"));
|
||||
assert!(!is_message_shard("message_0_resource.db"));
|
||||
assert!(!is_message_shard("message_resource.db"));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn is_message_shard_rejects_non_digits() {
|
||||
assert!(!is_message_shard("message_a.db"));
|
||||
assert!(!is_message_shard("message_.db"));
|
||||
assert!(!is_message_shard("session.db"));
|
||||
assert!(!is_message_shard("message_0.db.bak"));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn discover_unknown_shards_finds_disk_only_shards() {
|
||||
let dir = tempdir();
|
||||
let msg_dir = dir.join("message");
|
||||
std::fs::create_dir_all(&msg_dir).unwrap();
|
||||
for f in [
|
||||
"message_0.db",
|
||||
"message_1.db",
|
||||
"message_2.db",
|
||||
"message_0_fts.db",
|
||||
] {
|
||||
std::fs::write(msg_dir.join(f), b"").unwrap();
|
||||
}
|
||||
let known = vec![
|
||||
"message/message_0.db".to_string(),
|
||||
"message/message_1.db".to_string(),
|
||||
];
|
||||
let unknown = discover_unknown_shards(&dir, &known);
|
||||
assert_eq!(unknown, vec!["message/message_2.db".to_string()]);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn discover_unknown_shards_normalizes_backslash_in_known_keys() {
|
||||
let dir = tempdir();
|
||||
let msg_dir = dir.join("message");
|
||||
std::fs::create_dir_all(&msg_dir).unwrap();
|
||||
std::fs::write(msg_dir.join("message_0.db"), b"").unwrap();
|
||||
|
||||
let known = vec!["message\\message_0.db".to_string()];
|
||||
assert!(discover_unknown_shards(&dir, &known).is_empty());
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn discover_unknown_shards_returns_empty_when_message_dir_missing() {
|
||||
let dir = tempdir();
|
||||
assert!(discover_unknown_shards(&dir, &[]).is_empty());
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn derive_status_unknown_shards_overrides_windowed() {
|
||||
let unknown = vec!["message/message_3.db".to_string()];
|
||||
assert_eq!(
|
||||
derive_status(Some(100), Some(100), &unknown, true),
|
||||
MetaStatus::PossiblyStaleUnknownShards
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn derive_status_windowed_when_user_paginates() {
|
||||
assert_eq!(
|
||||
derive_status(Some(100), Some(999_999), &[], true),
|
||||
MetaStatus::Windowed,
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn derive_status_possibly_stale_when_session_far_ahead() {
|
||||
let chat = Some(1_000_000);
|
||||
let session = Some(1_000_000 + STALE_THRESHOLD_SECS + 1);
|
||||
assert_eq!(
|
||||
derive_status(chat, session, &[], false),
|
||||
MetaStatus::PossiblyStale
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn derive_status_ok_when_within_threshold() {
|
||||
let chat = Some(1_000_000);
|
||||
let session = Some(1_000_000 + STALE_THRESHOLD_SECS - 1);
|
||||
assert_eq!(derive_status(chat, session, &[], false), MetaStatus::Ok);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn derive_status_ok_when_either_side_unknown() {
|
||||
assert_eq!(
|
||||
derive_status(None, Some(999_999_999), &[], false),
|
||||
MetaStatus::Ok
|
||||
);
|
||||
assert_eq!(derive_status(Some(1), None, &[], false), MetaStatus::Ok);
|
||||
assert_eq!(derive_status(None, None, &[], false), MetaStatus::Ok);
|
||||
}
|
||||
|
||||
fn tempdir() -> std::path::PathBuf {
|
||||
let pid = std::process::id();
|
||||
let nanos = std::time::SystemTime::now()
|
||||
.duration_since(std::time::UNIX_EPOCH)
|
||||
.unwrap()
|
||||
.as_nanos();
|
||||
let p = std::env::temp_dir().join(format!("wx-cli-meta-test-{}-{}", pid, nanos));
|
||||
std::fs::create_dir_all(&p).unwrap();
|
||||
p
|
||||
}
|
||||
}
|
||||
@@ -1,4 +1,5 @@
|
||||
pub mod cache;
|
||||
pub mod meta;
|
||||
pub mod query;
|
||||
pub mod server;
|
||||
|
||||
|
||||
+1325
-556
File diff suppressed because it is too large
Load Diff
+178
-73
@@ -2,15 +2,12 @@ use anyhow::Result;
|
||||
use std::sync::Arc;
|
||||
use tokio::io::{AsyncBufReadExt, AsyncWriteExt, BufReader};
|
||||
|
||||
use crate::ipc::{Request, Response};
|
||||
use super::cache::DbCache;
|
||||
use super::query::Names;
|
||||
use crate::ipc::{Request, Response};
|
||||
|
||||
/// 启动 IPC server(Unix socket / Windows named pipe)
|
||||
pub async fn serve(
|
||||
db: Arc<DbCache>,
|
||||
names: Arc<tokio::sync::RwLock<Arc<Names>>>,
|
||||
) -> Result<()> {
|
||||
pub async fn serve(db: Arc<DbCache>, names: Arc<tokio::sync::RwLock<Arc<Names>>>) -> Result<()> {
|
||||
#[cfg(unix)]
|
||||
serve_unix(db, names).await?;
|
||||
#[cfg(windows)]
|
||||
@@ -19,10 +16,7 @@ pub async fn serve(
|
||||
}
|
||||
|
||||
#[cfg(unix)]
|
||||
async fn serve_unix(
|
||||
db: Arc<DbCache>,
|
||||
names: Arc<tokio::sync::RwLock<Arc<Names>>>,
|
||||
) -> Result<()> {
|
||||
async fn serve_unix(db: Arc<DbCache>, names: Arc<tokio::sync::RwLock<Arc<Names>>>) -> Result<()> {
|
||||
use tokio::net::UnixListener;
|
||||
let sock_path = crate::config::sock_path();
|
||||
|
||||
@@ -88,9 +82,7 @@ async fn serve_windows(
|
||||
db: Arc<DbCache>,
|
||||
names: Arc<tokio::sync::RwLock<Arc<Names>>>,
|
||||
) -> Result<()> {
|
||||
use interprocess::local_socket::{
|
||||
tokio::prelude::*, GenericNamespaced, ListenerOptions,
|
||||
};
|
||||
use interprocess::local_socket::{tokio::prelude::*, GenericNamespaced, ListenerOptions};
|
||||
|
||||
// interprocess 的 GenericNamespaced 在 Windows 上会自动拼接 `\\.\pipe\` 前缀,
|
||||
// 这里必须传相对名;client 端用 `\\.\pipe\wx-cli-daemon` 直接打开可以对上
|
||||
@@ -141,13 +133,9 @@ async fn handle_connection_windows(
|
||||
Ok(())
|
||||
}
|
||||
|
||||
async fn dispatch(
|
||||
req: Request,
|
||||
db: &DbCache,
|
||||
names: &tokio::sync::RwLock<Arc<Names>>,
|
||||
) -> Response {
|
||||
use crate::ipc::Request::*;
|
||||
async fn dispatch(req: Request, db: &DbCache, names: &tokio::sync::RwLock<Arc<Names>>) -> Response {
|
||||
use super::query;
|
||||
use crate::ipc::Request::*;
|
||||
|
||||
// 取 guard → O(1) clone Arc → 立即 drop 锁。后续 await 期间不持有锁,
|
||||
// 多个并发 IPC 请求可以真正并行。Names 本身不可变(由 daemon 启动时
|
||||
@@ -159,20 +147,66 @@ async fn dispatch(
|
||||
|
||||
match req {
|
||||
Ping => Response::ok(serde_json::json!({ "pong": true })),
|
||||
Sessions { limit } => {
|
||||
match query::q_sessions(db, &names_arc, limit).await {
|
||||
Sessions {
|
||||
limit,
|
||||
with_meta,
|
||||
debug_source,
|
||||
} => match query::q_sessions(db, &names_arc, limit, with_meta, debug_source).await {
|
||||
Ok(v) => Response::ok(v),
|
||||
Err(e) => Response::err(e.to_string()),
|
||||
},
|
||||
History {
|
||||
chat,
|
||||
limit,
|
||||
offset,
|
||||
since,
|
||||
until,
|
||||
msg_type,
|
||||
with_meta,
|
||||
debug_source,
|
||||
} => {
|
||||
match query::q_history(
|
||||
db,
|
||||
&names_arc,
|
||||
&chat,
|
||||
limit,
|
||||
offset,
|
||||
since,
|
||||
until,
|
||||
msg_type,
|
||||
with_meta,
|
||||
debug_source,
|
||||
)
|
||||
.await
|
||||
{
|
||||
Ok(v) => Response::ok(v),
|
||||
Err(e) => Response::err(e.to_string()),
|
||||
}
|
||||
}
|
||||
History { chat, limit, offset, since, until, msg_type } => {
|
||||
match query::q_history(db, &names_arc, &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, msg_type } => {
|
||||
match query::q_search(db, &names_arc, &keyword, chats, limit, since, until, msg_type).await {
|
||||
Search {
|
||||
keyword,
|
||||
chats,
|
||||
limit,
|
||||
since,
|
||||
until,
|
||||
msg_type,
|
||||
with_meta,
|
||||
debug_source,
|
||||
} => {
|
||||
match query::q_search(
|
||||
db,
|
||||
&names_arc,
|
||||
&keyword,
|
||||
chats,
|
||||
limit,
|
||||
since,
|
||||
until,
|
||||
msg_type,
|
||||
with_meta,
|
||||
debug_source,
|
||||
)
|
||||
.await
|
||||
{
|
||||
Ok(v) => Response::ok(v),
|
||||
Err(e) => Response::err(e.to_string()),
|
||||
}
|
||||
@@ -183,74 +217,145 @@ async fn dispatch(
|
||||
Err(e) => Response::err(e.to_string()),
|
||||
}
|
||||
}
|
||||
Unread { limit, filter } => {
|
||||
match query::q_unread(db, &names_arc, limit, filter).await {
|
||||
Unread {
|
||||
limit,
|
||||
filter,
|
||||
with_meta,
|
||||
debug_source,
|
||||
} => match query::q_unread(db, &names_arc, limit, filter, with_meta, debug_source).await {
|
||||
Ok(v) => Response::ok(v),
|
||||
Err(e) => Response::err(e.to_string()),
|
||||
},
|
||||
Members { chat } => match query::q_members(db, &names_arc, &chat).await {
|
||||
Ok(v) => Response::ok(v),
|
||||
Err(e) => Response::err(e.to_string()),
|
||||
},
|
||||
NewMessages {
|
||||
state,
|
||||
limit,
|
||||
with_meta,
|
||||
debug_source,
|
||||
} => {
|
||||
match query::q_new_messages(db, &names_arc, state, limit, with_meta, debug_source).await
|
||||
{
|
||||
Ok(v) => Response::ok(v),
|
||||
Err(e) => Response::err(e.to_string()),
|
||||
}
|
||||
}
|
||||
Members { chat } => {
|
||||
match query::q_members(db, &names_arc, &chat).await {
|
||||
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,
|
||||
with_meta,
|
||||
debug_source,
|
||||
} => {
|
||||
match query::q_stats(db, &names_arc, &chat, since, until, with_meta, debug_source).await
|
||||
{
|
||||
Ok(v) => Response::ok(v),
|
||||
Err(e) => Response::err(e.to_string()),
|
||||
}
|
||||
}
|
||||
NewMessages { state, limit } => {
|
||||
match query::q_new_messages(db, &names_arc, state, limit).await {
|
||||
SnsNotifications {
|
||||
limit,
|
||||
since,
|
||||
until,
|
||||
include_read,
|
||||
} => {
|
||||
match query::q_sns_notifications(db, &names_arc, limit, since, until, include_read)
|
||||
.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 {
|
||||
SnsFeed {
|
||||
limit,
|
||||
since,
|
||||
until,
|
||||
user,
|
||||
} => match query::q_sns_feed(db, &names_arc, limit, since, until, user.as_deref()).await {
|
||||
Ok(v) => Response::ok(v),
|
||||
Err(e) => Response::err(e.to_string()),
|
||||
},
|
||||
SnsSearch {
|
||||
keyword,
|
||||
limit,
|
||||
since,
|
||||
until,
|
||||
user,
|
||||
} => {
|
||||
match query::q_sns_search(
|
||||
db,
|
||||
&names_arc,
|
||||
&keyword,
|
||||
limit,
|
||||
since,
|
||||
until,
|
||||
user.as_deref(),
|
||||
)
|
||||
.await
|
||||
{
|
||||
Ok(v) => Response::ok(v),
|
||||
Err(e) => Response::err(e.to_string()),
|
||||
}
|
||||
}
|
||||
Stats { chat, since, until } => {
|
||||
match query::q_stats(db, &names_arc, &chat, since, until).await {
|
||||
ReloadConfig => Response::ok(serde_json::json!({ "reloading": true })),
|
||||
BizArticles {
|
||||
limit,
|
||||
account,
|
||||
since,
|
||||
until,
|
||||
unread,
|
||||
} => {
|
||||
match query::q_biz_articles(db, &names_arc, limit, account, since, until, unread).await
|
||||
{
|
||||
Ok(v) => Response::ok(v),
|
||||
Err(e) => Response::err(e.to_string()),
|
||||
}
|
||||
}
|
||||
SnsNotifications { limit, since, until, include_read } => {
|
||||
match query::q_sns_notifications(db, &names_arc, limit, since, until, include_read).await {
|
||||
Ok(v) => Response::ok(v),
|
||||
Err(e) => Response::err(e.to_string()),
|
||||
}
|
||||
}
|
||||
SnsFeed { limit, since, until, user } => {
|
||||
match query::q_sns_feed(db, &names_arc, limit, since, until, user.as_deref()).await {
|
||||
Ok(v) => Response::ok(v),
|
||||
Err(e) => Response::err(e.to_string()),
|
||||
}
|
||||
}
|
||||
SnsSearch { keyword, limit, since, until, user } => {
|
||||
match query::q_sns_search(db, &names_arc, &keyword, limit, since, until, user.as_deref()).await {
|
||||
Ok(v) => Response::ok(v),
|
||||
Err(e) => Response::err(e.to_string()),
|
||||
}
|
||||
}
|
||||
ReloadConfig => {
|
||||
Response::ok(serde_json::json!({ "reloading": true }))
|
||||
}
|
||||
BizArticles { limit, account, since, until, unread } => {
|
||||
match query::q_biz_articles(db, &names_arc, limit, account, since, until, unread).await {
|
||||
Ok(v) => Response::ok(v),
|
||||
Err(e) => Response::err(e.to_string()),
|
||||
}
|
||||
}
|
||||
Attachments { chat, kinds, limit, offset, since, until } => {
|
||||
match query::q_attachments(db, &names_arc, &chat, kinds, limit, offset, since, until).await {
|
||||
Ok(v) => Response::ok(v),
|
||||
Err(e) => Response::err(e.to_string()),
|
||||
}
|
||||
}
|
||||
Extract { attachment_id, output, overwrite } => {
|
||||
match query::q_extract(db, &names_arc, &attachment_id, &output, overwrite).await {
|
||||
Attachments {
|
||||
chat,
|
||||
kinds,
|
||||
limit,
|
||||
offset,
|
||||
since,
|
||||
until,
|
||||
with_meta,
|
||||
debug_source,
|
||||
} => {
|
||||
match query::q_attachments(
|
||||
db,
|
||||
&names_arc,
|
||||
&chat,
|
||||
kinds,
|
||||
limit,
|
||||
offset,
|
||||
since,
|
||||
until,
|
||||
with_meta,
|
||||
debug_source,
|
||||
)
|
||||
.await
|
||||
{
|
||||
Ok(v) => Response::ok(v),
|
||||
Err(e) => Response::err(e.to_string()),
|
||||
}
|
||||
}
|
||||
Extract {
|
||||
attachment_id,
|
||||
output,
|
||||
overwrite,
|
||||
} => match query::q_extract(db, &names_arc, &attachment_id, &output, overwrite).await {
|
||||
Ok(v) => Response::ok(v),
|
||||
Err(e) => Response::err(e.to_string()),
|
||||
},
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user