Files
kefu/wechat_rpa/wxwork_db.py
2026-08-27 14:04:28 +08:00

918 lines
35 KiB
Python
Raw Permalink Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
# -*- coding: utf-8 -*-
"""
wxwork_db.py - 企业微信本地数据库读取模块(引擎 B 升级数据源)
=================================================================
基于 wxwork_crypto 解密 message.db / session.db / user.db 为标准 SQLite
提供增量消息轮询接口,供 engine_b 使用。
参考实现(C:\\Users\\pc\\WorkBuddy\\2026-08-19-18-27-34\\,密钥已验证正确):
- wxwork_keys.json 已复制的验证过密钥(本项目 wxwork_keys.json
- wxwork_export_final.py 移植:parse_content / extract_protobuf_texts /
clean_text / MSG_TYPE_MAP / is_personal_chat /
is_blocked_conv_name / format_timestamp /
load_keys / detect_wxwork_dir / connect_sqlite /
decrypt_with_keysWAL 感知增量解密缓存)/
会话名解析(user_cache / conv_cache / S: 对端推导)
- wxwork_crypto.py 解密算法与本项目 wxwork_crypto.py 完全一致(互相印证)
相对参考实现的增强:
- WAL 帧合并:企微运行时新消息在 message.db-wal(未 checkpoint),参考实现仅
以 wal_mtime 触发重解密主库(仍缺 WAL 内最新消息)。本项目直接解密 WAL 帧
(帧数据与主库同算法加密,帧头明文),按页号覆盖合并进主库解密副本,
经 PRAGMA quick_check 验证(见 test_wal_merge.py)。
- 会话指纹:复用 wechat_bot.WeChatBot._fp_from_name 的名称哈希算法,
使 DB 直读的 fp_hex 与引擎 Aidentity_by_name 模式)完全一致,
引擎 A 可在会话列表按名称重新定位窗口。
典型流程:
keys = load_keys()
db = WXWorkDB(detect_wxwork_dir(), keys)
msgs = db.get_new_messages(since_ts=...)
db.close()
"""
from __future__ import annotations
import hashlib
import json
import os
import re
import sqlite3
import struct
import time
import unicodedata
from datetime import datetime
from wxwork_crypto import (
decrypt_db_to_file,
decrypt_page,
read_page_header,
SQLITE_FILE_HEADER,
verify_key,
)
_SCRIPT_DIR = os.path.dirname(os.path.abspath(__file__))
KEYS_FILE = os.path.join(_SCRIPT_DIR, "wxwork_keys.json")
DEFAULT_CACHE_DIR = os.path.join(_SCRIPT_DIR, "wxwork_decrypted")
# 企微数据库文件名(与聊天相关的库)
_RELEVANT_DBS = (
"message.db", "session.db", "user.db", "company.db",
"message_lookup.db", "user_extend.db",
)
# 引擎 A 会话指纹字节数(wechat_bot.py: _AVATAR_FP_BYTES=8 + _NAME_FP_BYTES=32
SESSION_FP_BYTES = 40
# ---------------------------------------------------------------------------
# 消息类型映射 / 会话过滤(移植自 wxwork_export_final.py
# ---------------------------------------------------------------------------
MSG_TYPE_MAP = {
0: "文本", 1: "文本", 2: "文本",
3: "图片", 4: "语音", 5: "表情",
6: "链接", 7: "文件", 8: "视频",
9: "位置", 10: "名片", 11: "系统",
12: "引用", 13: "红包", 14: "图片",
15: "转账", 16: "语音", 17: "视频号",
20: "合并转发", 21: "日程", 22: "接龙",
23: "文件回复", 26: "位置共享", 29: "视频通话",
31: "图文链接", 38: "系统通知", 40: "待办",
42: "文件预览", 46: "话题", 47: "图文",
49: "待办", 51: "收藏", 53: "订阅通知",
55: "文件分享", 56: "群公告", 57: "投票",
59: "回执", 61: "文件编辑", 62: "文件评论",
63: "收藏合并", 64: "卡片", 65: "群文件",
66: "工作台", 67: "位置共享", 68: "视频通话",
69: "音频通话", 70: "企业微信应用", 73: "微信好友",
74: "回复", 76: "邀请", 77: "移除",
78: "修改群名", 79: "修改群公告", 80: "加入群聊",
81: "退出群聊", 82: "解散群聊", 83: "群主转让",
101: "系统", 111: "文件", 123: "截图",
132: "系统",
503: "应用消息", 529: "朋友圈", 561: "应用消息",
565: "应用消息", 573: "应用消息",
1002: "安全通知", 1011: "系统消息", 1012: "系统消息",
1017: "系统消息", 1022: "系统消息", 1025: "系统消息",
1043: "系统消息",
1988: "系统消息",
}
def get_msg_type_name(ct):
"""消息类型名"""
try:
ct = int(ct)
except (ValueError, TypeError):
return f"类型{ct}"
return MSG_TYPE_MAP.get(ct, f"类型{ct}")
def is_personal_chat(conv_id):
"""判断会话是否为人与人之间的单聊(移植自 wxwork_export_final.py)。
保留: M:xxx(微信单聊)、S:UID_UID(企微单聊)
过滤: R:xxx 群聊 / Y:xxx 应用 / O:xxx 第三方 / MAIL/APPROVAL 系统虚拟会话
"""
c = str(conv_id or "")
if c.startswith("M:"):
return True
if c.startswith("S:"):
parts = c[2:].split("_")
if len(parts) == 2 and parts[0].isdigit() and parts[1].isdigit():
return True
return False
return False
# 按会话名称过滤的官方/系统账号
DEFAULT_BLOCKED_CONV_NAMES = ("企业微信团队", "微信团队", "微信支付", "腾讯客服", "腾讯新闻")
def is_blocked_conv_name(conv_name, blocked=None):
"""按会话名称判断是否命中过滤名单(官方/系统账号)。"""
if not conv_name:
return False
blocked = DEFAULT_BLOCKED_CONV_NAMES if blocked is None else tuple(blocked)
return any(kw and kw in str(conv_name) for kw in blocked)
# ---------------------------------------------------------------------------
# 消息内容解析(移植自 wxwork_export_final.py
# ---------------------------------------------------------------------------
def clean_text(text):
"""清洗提取的文本: 去除控制字符和 protobuf 残留。"""
if not text:
return ""
cleaned = "".join(c for c in text if c.isprintable() or c in "\n\t")
cleaned = cleaned.strip()
cleaned = re.sub(r"[ \t]+", " ", cleaned)
cleaned = re.sub(r"\n{3,}", "\n\n", cleaned)
return cleaned
def extract_protobuf_texts(data):
"""从 protobuf 编码的 content 中递归提取可读文本,返回去重后的文本列表。"""
if not data or not isinstance(data, bytes):
return []
texts = []
seen = set()
def add_text(t):
t = clean_text(t)
if len(t) >= 2 and t not in seen:
has_chinese = any("\u4e00" <= c <= "\u9fff" for c in t)
alpha_count = sum(1 for c in t if c.isalpha())
if has_chinese or alpha_count >= 6:
seen.add(t)
texts.append(t)
def parse(buf, depth):
if depth > 6 or not buf:
return
i = 0
while i < len(buf):
try:
tag = 0
shift = 0
while i < len(buf) and shift < 70:
b = buf[i]
tag |= (b & 0x7f) << shift
i += 1
if not (b & 0x80):
break
shift += 7
if shift >= 70:
break
field_num = tag >> 3
wire_type = tag & 7
if field_num == 0:
break
if wire_type == 0: # varint
while i < len(buf):
b = buf[i]
i += 1
if not (b & 0x80):
break
elif wire_type == 1: # 64-bit
i += 8
elif wire_type == 2: # length-delimited
length = 0
shift = 0
while i < len(buf) and shift < 70:
b = buf[i]
length |= (b & 0x7f) << shift
i += 1
if not (b & 0x80):
break
shift += 7
if shift >= 70:
break
chunk = buf[i:i + length]
i += length
try:
text = chunk.decode("utf-8")
if all(ord(c) >= 32 or c in "\n\t\r" for c in text):
add_text(text)
except (UnicodeDecodeError, ValueError):
pass
if length < 8192:
parse(chunk, depth + 1)
elif wire_type == 5: # 32-bit
i += 4
else:
break
except Exception:
break
parse(data, 0)
if not any("\u4e00" <= c <= "\u9fff" for t in texts for c in t[:50]):
try:
text = data.decode("utf-8", errors="ignore")
readable = re.findall(
r"[\u4e00-\u9fff\w\s,。!?、;:“”‘’()《》【】\-—.,;:!?()\[\]/\\+*=<>@#$%^&~`]{2,}",
text)
for r in readable:
add_text(r.strip())
except Exception:
pass
return texts
def parse_content(content):
"""解析消息内容, 返回可读文本(移植自 wxwork_export_final.py)。"""
if content is None:
return ""
if not isinstance(content, bytes):
content = str(content).encode("utf-8", errors="replace")
if not content:
return ""
# 情况 1: 直接是 UTF-8 文本
try:
text = content.decode("utf-8")
if text and all(ord(c) >= 32 or c in "\n\r\t" for c in text[:500]):
return clean_text(text)
except (UnicodeDecodeError, ValueError):
pass
# 情况 2: protobuf 编码
texts = extract_protobuf_texts(content)
if texts:
texts.sort(key=len, reverse=True)
meaningful = [t for t in texts
if not (t.startswith("{") and "}" in t)
and not t.startswith("http")
and len(t) > 1]
if meaningful:
cleaned_msgs = []
for t in meaningful:
t = t.strip()
m = re.match(r"^(.{6,}?)[?zxa-zA-Z]{1,4}\1$", t)
if m:
t = m.group(1)
t = re.sub(
r"^[?zxa-zA-Z]{1,4}(?=[\u4e00-\u9fff\u3000-\u303f\uff00-\uffef])", "", t)
t = re.sub(
r"^[0-9](?=[\u4e00-\u9fff])(?![\u4e00-\u9fff]{0,2}"
r"(?:月|日|年|个|天|点|号|楼|人|次|分|秒|周|元|块|台|家|辆|条|件|张|本|层|岁|万|亿|美元|室|房|折))",
"", t)
t = t.strip()
if len(t) >= 3:
dup = False
for c in cleaned_msgs:
if len(c) >= 4 and (t in c or c in t):
dup = True
break
if dup:
continue
if t and t not in cleaned_msgs:
cleaned_msgs.append(t)
if cleaned_msgs:
return " | ".join(cleaned_msgs[:2])
return texts[0]
# 情况 3: JSON
try:
obj = json.loads(content.decode("utf-8", errors="replace"))
return json.dumps(obj, ensure_ascii=False)[:2000]
except Exception:
pass
# 情况 4: 二进制数据
if len(content) <= 64:
return content.hex()
return f"[二进制数据 {len(content)} 字节]"
def format_timestamp(ts):
"""Unix 时间戳转可读格式(>10^12 视为毫秒)。"""
if ts is None or ts == "":
return ""
try:
ts = int(ts)
if ts > 10 ** 12:
ts = ts // 1000
if ts < 0:
return str(ts)
return datetime.fromtimestamp(ts).strftime("%Y-%m-%d %H:%M:%S")
except (ValueError, OSError, OverflowError):
return str(ts)
# ---------------------------------------------------------------------------
# 会话指纹(与 wechat_bot.WeChatBot._fp_from_name 算法一致)
# ---------------------------------------------------------------------------
_WHITESPACE_RE = re.compile(r"\s+")
def normalize_name(value: str) -> str:
"""昵称归一化(与 session_name.normalize_name 完全一致)。
实测企微标题读出来是「高瑞 @微信」、列表行是「高瑞@微信」,中间多个空格;
不归一的话同一个人会算出两个 md5。仅用标准库实现(session_name 顶层
依赖 numpy,DB 模块不引入视觉依赖)。
"""
text = unicodedata.normalize("NFKC", str(value or ""))
text = _WHITESPACE_RE.sub("", text)
return text.strip()
def session_fp_from_name(name: str) -> str:
"""会话名 → 40 字节会话指纹 hex。
与 wechat_bot.py::_fp_from_name 完全一致:前 16 字节 = md5(归一化昵称),
后 24 字节 = blake2b(md5)。引擎 A 在 identity_by_name 模式下用同一算法
从屏幕 OCR 的昵称生成指纹,因此 DB 直读投递的 fp_hex 可被引擎 A
在会话列表中按名称重新定位窗口。
"""
normalized = normalize_name(name)
if not normalized:
return ""
digest = hashlib.md5(normalized.encode("utf-8")).digest()
tail = hashlib.blake2b(digest, digest_size=SESSION_FP_BYTES - len(digest)).digest()
return (digest + tail).hex()
# ---------------------------------------------------------------------------
# 密钥与解密(移植自 wxwork_export_final.py + WAL 帧合并增强)
# ---------------------------------------------------------------------------
def load_keys() -> dict:
"""加载密钥映射 {user_dir: key_hex},支持 global_key/* 通配。"""
if not os.path.exists(KEYS_FILE):
return {}
with open(KEYS_FILE, encoding="utf-8") as f:
data = json.load(f)
keys = data.get("keys", {})
if "global_key" in data:
keys.setdefault("*", data["global_key"])
return keys
def detect_wxwork_dir(candidates=None) -> str | None:
"""自动检测企业微信数据根目录(移植自 wxwork_export_final.py)。
依次尝试: 配置候选 → 当前用户 Documents\\WXWork →
注册表 Personal 下的 WXWork → 扫描所有用户。
返回第一个含账号 Data/message.db 的有效目录。
"""
def _valid(base):
if not base or not os.path.isdir(base):
return False
try:
for d in os.listdir(base):
if os.path.exists(os.path.join(base, d, "Data", "message.db")):
return True
except OSError:
pass
return False
for c in (candidates or []):
if _valid(c):
return c
default = os.path.join(os.path.expanduser("~"), "Documents", "WXWork")
if _valid(default):
return default
try:
import winreg
with winreg.OpenKey(
winreg.HKEY_CURRENT_USER,
r"Software\Microsoft\Windows\CurrentVersion\Explorer\Shell Folders",
) as k:
personal, _ = winreg.QueryValueEx(k, "Personal")
if personal:
reg_path = os.path.join(personal, "WXWork")
if _valid(reg_path):
return reg_path
except Exception:
pass
users_root = os.path.join(os.environ.get("SystemDrive", "C:"), os.sep, "Users")
if os.path.isdir(users_root):
try:
for u in sorted(os.listdir(users_root)):
p = os.path.join(users_root, u, "Documents", "WXWork")
if _valid(p):
return p
except OSError:
pass
return None
def _is_plain_sqlite(path: str) -> bool:
"""文件头是否为明文 SQLite。"""
try:
with open(path, "rb") as f:
return f.read(16) == SQLITE_FILE_HEADER
except OSError:
return False
def wal_read_frames(wal_path: str, key: bytes) -> tuple[dict, dict]:
"""读取 WAL 帧并解密,返回 ({pgno: bytes}, meta)。
帧数据与主库同算法加密(wxSQLite3 aes128cbc,帧头 24 字节明文)。
仅应用 salt 与 WAL 头一致的帧(跳过已回滚/旧事务帧),后帧覆盖前帧。
"""
pages: dict[int, bytes] = {}
meta = {"pgsz": 0, "frames": 0, "applied": 0, "salt_ok": 0}
try:
with open(wal_path, "rb") as f:
hdr = f.read(32)
if len(hdr) < 32:
return pages, meta
magic, ver, pgsz, ckpt, salt1, salt2, _, _ = struct.unpack(">IIIIIIII", hdr)
meta["pgsz"] = pgsz
while True:
fh = f.read(24)
if len(fh) < 24:
break
pgno, commit, fs1, fs2, _, _ = struct.unpack(">IIIIII", fh)
data = f.read(pgsz)
if len(data) < pgsz:
break
meta["frames"] += 1
if fs1 != salt1 or fs2 != salt2:
continue
meta["salt_ok"] += 1
if pgno == 0 or pgno > (1 << 31):
continue
dec = decrypt_page(key, pgno, data, pgsz)
if dec is None:
continue
pages[pgno] = dec
meta["applied"] += 1
except OSError:
pass
return pages, meta
def merge_wal_into_db(db_path: str, wal_path: str, key: bytes, out_path: str) -> bool:
"""解密主库 + 合并 WAL 帧 → out_path。返回是否成功。
主库与 WAL 帧都解密后,按页号覆盖(等价于 checkpoint 后的状态),
并更新 SQLite 头的页数字段。WAL 无有效帧时退化为仅主库。
"""
if not decrypt_db_to_file(db_path, key, out_path):
return False
pages, meta = wal_read_frames(wal_path, key)
if not pages:
return True
pgsz = meta["pgsz"] or 4096
try:
with open(out_path, "r+b") as f:
size = os.path.getsize(out_path)
n_pages = (size + pgsz - 1) // pgsz
for pgno, data in pages.items():
if pgno > n_pages:
f.seek(0, os.SEEK_END)
f.write(b"\x00" * ((pgno - n_pages) * pgsz))
n_pages = pgno
f.seek((pgno - 1) * pgsz)
f.write(data)
if n_pages > 0:
f.seek(28)
f.write(struct.pack(">I", n_pages))
except OSError:
return False
return True
def decrypt_with_keys(db_base: str, out_dir: str, keys_map: dict, use_cache: bool = True) -> list:
"""解密所有可解密的账号数据库(移植自 wxwork_export_final.py)。
返回 [(out_path, db_name, user_dir), ...]。
- 增量缓存: 解密副本 mtime >= max(db, wal) 且大小正常则复用
- WAL 增强: message.db 的 WAL 非空时合并帧(拿 checkpoint 前的最新消息)
- 密钥: 优先目录密钥,失败尝试全部密钥
"""
decrypted_dbs = []
os.makedirs(out_dir, exist_ok=True)
all_keys = set()
for k in keys_map.values():
try:
all_keys.add(bytes.fromhex(k.replace("x'", "").replace("'", "")))
except ValueError:
pass
if not os.path.isdir(db_base):
return decrypted_dbs
for user_dir in sorted(os.listdir(db_base)):
data_dir = os.path.join(db_base, user_dir, "Data")
if not os.path.isdir(data_dir):
continue
db_files = [f for f in os.listdir(data_dir)
if f.endswith(".db") and not f.endswith("-wal") and not f.endswith("-shm")]
if not db_files:
continue
user_out = os.path.join(out_dir, user_dir)
os.makedirs(user_out, exist_ok=True)
dir_key = None
if user_dir in keys_map:
try:
dir_key = bytes.fromhex(keys_map[user_dir])
except ValueError:
dir_key = None
for db_name in _RELEVANT_DBS:
db_path = os.path.join(data_dir, db_name)
if not os.path.isfile(db_path):
continue
out_path = os.path.join(user_out, db_name)
# 增量缓存判断(WAL 感知)
wal_path = db_path + "-wal"
wal_mtime = 0
try:
if os.path.exists(wal_path) and os.path.getsize(wal_path) > 0:
wal_mtime = os.path.getmtime(wal_path)
except OSError:
wal_mtime = 0
src_mtime = max(os.path.getmtime(db_path), wal_mtime)
if use_cache and os.path.exists(out_path):
try:
if (os.path.getmtime(out_path) >= src_mtime
and os.path.getsize(out_path) > 4096):
decrypted_dbs.append((out_path, db_name, user_dir))
continue
except OSError:
pass
# 明文库直接复制
if _is_plain_sqlite(db_path):
try:
with open(db_path, "rb") as fin, open(out_path, "wb") as fout:
fout.write(fin.read())
decrypted_dbs.append((out_path, db_name, user_dir))
continue
except OSError:
continue
# 加密库:目录密钥 → 全部密钥
if read_page_header(db_path) is None:
continue # 不是企微加密格式
candidate_keys = [dir_key] if dir_key else []
candidate_keys += [k for k in all_keys if k != dir_key]
decrypted_ok = False
for key in candidate_keys:
if key is None:
continue
try:
if not verify_key(key, db_path):
continue
# message.db 且 WAL 非空 → 合并;否则整库解密
if db_name == "message.db" and wal_mtime > 0:
ok = merge_wal_into_db(db_path, wal_path, key, out_path)
else:
ok = decrypt_db_to_file(db_path, key, out_path)
if ok:
decrypted_dbs.append((out_path, db_name, user_dir))
decrypted_ok = True
break
except Exception:
continue
if not decrypted_ok:
try:
os.remove(out_path)
except OSError:
pass
return decrypted_dbs
def connect_sqlite(path: str) -> sqlite3.Connection:
conn = sqlite3.connect(f"file:{path}?mode=ro", uri=True)
conn.text_factory = lambda b: b.decode("utf-8", errors="replace")
return conn
# ---------------------------------------------------------------------------
# WXWorkDB:多账号增量消息读取器
# ---------------------------------------------------------------------------
class WXWorkDB:
"""企微数据根目录下的加密数据库读取器(只读,多账号)。
- 构造时: 验证密钥 → 确保解密缓存(WAL 感知)→ 加载元数据 → 打开各账号
message.db 解密副本连接
- 轮询: get_new_messages(since_ts) 返回 since_ts 之后的新客户消息;
内部按 mtime 检查 WAL/主库变化,必要时重解密并重开连接
"""
def __init__(self, db_base: str, keys_map: dict | None = None,
cache_dir: str | None = None, poll_interval: float = 2.0):
self.db_base = db_base
self.keys_map = keys_map if keys_map is not None else load_keys()
self.cache_dir = cache_dir or DEFAULT_CACHE_DIR
self.poll_interval = max(0.5, float(poll_interval))
self.last_error = ""
# (user_dir) -> sqlite3.Connectionmessage.db 解密副本)
self._conns: dict[str, sqlite3.Connection] = {}
# (user_dir) -> 解密副本路径
self._paths: dict[str, str] = {}
# 元数据缓存
self.user_cache: dict = {}
self.conv_cache: dict = {}
self._open()
# ── 生命周期 ────────────────────────────────────────────────────────────
def _open(self):
decrypted = decrypt_with_keys(self.db_base, self.cache_dir, self.keys_map)
self._decrypted = decrypted
if not decrypted:
self.last_error = "没有可用密钥解密的账号数据库"
return
self._load_metadata()
for out_path, db_name, user_dir in decrypted:
if db_name != "message.db":
continue
try:
conn = connect_sqlite(out_path)
self._conns[user_dir] = conn
self._paths[user_dir] = out_path
except sqlite3.Error as exc:
self.last_error = f"打开 {user_dir}/message.db 失败: {exc}"
def _load_metadata(self):
"""加载 user_cache / conv_cache(移植自 wxwork_export_final.py 阶段 1)。"""
for out_path, db_name, user_dir in self._decrypted:
try:
conn = connect_sqlite(out_path)
cur = conn.cursor()
if db_name == "user.db":
try:
cur.execute("PRAGMA table_info(user_table)")
cols = [r[1] for r in cur.fetchall()]
if cols:
cur.execute(
"SELECT id, name, real_name, account FROM user_table")
for uid, name, real_name, account in cur.fetchall():
if uid is not None:
nm = name or real_name or account or str(uid)
self.user_cache[(user_dir, str(uid))] = str(nm)
except sqlite3.Error:
pass
elif db_name == "session.db":
try:
cur.execute("PRAGMA table_info(conversation_table)")
cols = [r[1] for r in cur.fetchall()]
if cols:
cur.execute(
"SELECT id, name, roomname_remark, session_id "
"FROM conversation_table")
for cid, name, remark, sid in cur.fetchall():
if cid:
nm = remark or name or sid
self.conv_cache[(user_dir, str(cid))] = (
str(nm) if nm else "")
except sqlite3.Error:
pass
conn.close()
except sqlite3.Error:
continue
def _conv_display_name(self, user_dir: str, conv_id) -> str:
"""会话显示名(移植自 wxwork_export_final.py 的 conv_name 解析)。"""
conv_id = str(conv_id or "")
name = self.conv_cache.get((user_dir, conv_id), "")
if name:
return name
if conv_id.startswith("M:"):
uid = conv_id[2:]
return self.user_cache.get((user_dir, uid), uid)
if conv_id.startswith("S:"):
parts = conv_id[2:].split("_")
peer_uid = ""
for p in parts:
if p != user_dir:
peer_uid = p
break
if not peer_uid:
peer_uid = parts[0] if parts else ""
peer_name = self.user_cache.get((user_dir, peer_uid), "")
return peer_name if peer_name else ("单聊 " + conv_id)
if conv_id.startswith("O:"):
return conv_id[2:]
if conv_id.startswith("Y:"):
return "应用 " + conv_id[2:]
return conv_id
def _refresh_if_needed(self):
"""WAL/主库比解密副本新 → 重解密并重开连接。"""
for out_path, db_name, user_dir in self._decrypted:
if db_name != "message.db":
continue
src = os.path.join(self.db_base, user_dir, "Data", db_name)
wal = src + "-wal"
try:
src_mtime = os.path.getmtime(src)
wal_mtime = 0
if os.path.exists(wal) and os.path.getsize(wal) > 0:
wal_mtime = os.path.getmtime(wal)
src_mtime = max(src_mtime, wal_mtime)
if os.path.getmtime(out_path) >= src_mtime:
continue
except OSError:
continue
# 重解密
key = None
for candidate in self._keys_for(user_dir):
try:
if verify_key(candidate, src):
key = candidate
break
except Exception:
continue
if key is None:
continue
try:
if wal_mtime > 0:
ok = merge_wal_into_db(src, wal, key, out_path)
else:
ok = decrypt_db_to_file(src, key, out_path)
if ok:
old = self._conns.pop(user_dir, None)
if old is not None:
old.close()
self._conns[user_dir] = connect_sqlite(out_path)
self._paths[user_dir] = out_path
except Exception as exc:
self.last_error = f"重解密 {user_dir} 失败: {exc}"
def _keys_for(self, user_dir: str) -> list:
keys = []
if user_dir in self.keys_map:
try:
keys.append(bytes.fromhex(self.keys_map[user_dir]))
except ValueError:
pass
if "*" in self.keys_map:
try:
keys.append(bytes.fromhex(self.keys_map["*"]))
except ValueError:
pass
return keys
# ── 消息读取 ────────────────────────────────────────────────────────────
def get_new_messages(self, since_ts: float) -> list[dict]:
"""返回 since_ts(秒)之后的新客户消息列表。
返回字段:
account 账号目录名
conv_id 企微会话 IDM:/S:/R:...
fp_hex 会话名派生的 40 字节指纹 hex(引擎 A identity_by_name 兼容)
display_name 会话显示名
sender_id 发送者 UID
is_self 是否本账号自己发送
content 解析后的可读文本(空文本不返回)
content_type 原始消息类型
send_time 秒级时间戳
server_id 服务器消息 ID(去重键)
dedup_key f"{fp_hex}:{send_time}"
"""
self._refresh_if_needed()
result: list[dict] = []
try:
since = float(since_ts or 0.0)
except (TypeError, ValueError):
since = 0.0
for user_dir, conn in self._conns.items():
try:
cur = conn.cursor()
cur.execute(
"PRAGMA table_info(message_table)")
columns = [row[1] for row in cur.fetchall()]
if "send_time" not in columns or "conversation_id" not in columns:
continue
# 增量查询(send_time 为秒级 Unix 时间戳)
cur.execute(
"SELECT * FROM message_table WHERE send_time > ? "
"ORDER BY send_time ASC LIMIT 2000",
(since,))
for row in cur.fetchall():
msg = dict(zip(columns, row))
m = self._parse_message(user_dir, msg)
if m:
result.append(m)
except sqlite3.Error as exc:
self.last_error = f"查询 {user_dir} 失败: {exc}"
result.sort(key=lambda m: m["send_time"])
return result
def _parse_message(self, user_dir: str, msg: dict) -> dict | None:
conv_id = str(msg.get("conversation_id") or "")
# 仅人与人单聊(过滤群聊/应用/第三方/系统虚拟会话)
if not is_personal_chat(conv_id):
return None
sender_id = str(msg.get("sender_id") or "")
send_time = msg.get("send_time") or 0
try:
send_time = float(send_time)
except (TypeError, ValueError):
return None
if send_time <= 0:
return None
# 本账号自己发送的消息不投递(引擎 B 只管客户消息)
is_self = bool(sender_id) and sender_id == str(user_dir)
content = parse_content(msg.get("content"))
if not content.strip():
return None
display_name = self._conv_display_name(user_dir, conv_id)
if is_blocked_conv_name(display_name) or is_blocked_conv_name(conv_id):
return None
fp_hex = session_fp_from_name(display_name)
if not fp_hex:
return None
server_id = str(msg.get("server_id") or "")
return {
"account": user_dir,
"conv_id": conv_id,
"fp_hex": fp_hex,
"display_name": display_name,
"sender_id": sender_id,
"is_self": is_self,
"content": content,
"content_type": msg.get("content_type"),
"send_time": send_time,
"server_id": server_id,
"dedup_key": f"{fp_hex}:{send_time}",
}
def health_check(self) -> bool:
"""连接与密钥有效性自检。"""
if not self._conns:
return False
try:
for conn in self._conns.values():
cur = conn.cursor()
cur.execute("SELECT COUNT(*) FROM sqlite_master")
cur.fetchone()
return True
except Exception:
return False
def close(self):
for conn in self._conns.values():
try:
conn.close()
except Exception:
pass
self._conns.clear()
def open_account_db(account_dir: str, key: bytes) -> WXWorkDB | None:
"""便捷工厂(兼容旧接口):单账号目录 + 单密钥。"""
try:
user_dir = os.path.basename(os.path.normpath(account_dir))
return WXWorkDB(
os.path.dirname(os.path.normpath(account_dir)),
{user_dir: key.hex()},
)
except Exception:
return None
if __name__ == "__main__":
# 冒烟测试:真实数据目录 + 密钥,打印可解密账号与最近消息
keys = load_keys()
print(f"[*] 密钥文件: {KEYS_FILE},账号密钥 {len(keys)} 个")
base = detect_wxwork_dir()
print(f"[*] 企微数据目录: {base}")
if not keys or not base:
print("[-] 缺少密钥或数据目录,无法冒烟")
raise SystemExit(1)
db = WXWorkDB(base, keys)
print(f"[*] 可读账号: {sorted(db._conns.keys())}")
print(f"[*] health_check: {db.health_check()}")
msgs = db.get_new_messages(0)
print(f"[*] 全量消息 {len(msgs)} 条,最近 5 条:")
for m in msgs[-5:]:
print(f" {m['send_time']} [{m['account']}] {m['display_name']}: "
f"{m['content'][:40]}")
db.close()