会话采集与导出协议
前七篇的 Agent 一直在查工单。从这一篇开始喂真实语料,第一个来源是群聊。
先把它打开看看
聊天客户端的消息就在本机,一个 995 MiB 的 SQLite 文件。直接打开:
$ sqlite3 dingtalk.db 'select count(*) from sqlite_schema;'
Error: file is not a database (26)
这句报错只说明「不能按标准 SQLite 打开」,不说明业务数据不存在。文件头被逐页加密了。
这里要先划一条线。本篇可运行的起点,是一份已经解密并通过完整性检查的重放副本,前面那两步属于本机私有前置操作:
完全退出客户端
→ 复制主库、WAL、SHM 一致快照
→ 用账号分区 UID 在本机解密
→ 校验并重放 WAL 到最后有效提交
→ integrity_check 通过的工作副本 ← 本篇从这里开始
解密那步的算法很短:UID 做 MD5,取十六进制串的前 16 字节作 AES key,按 4096 字节页解密,再用 SQLite 页头和 schema 校验这个 UID 猜对了没有。
WAL 那步才是容易被跳过、又必须做的。 只解密主库会漏掉尚未 checkpoint 的已提交消息——也就是最近这批。按 SQLite 的规则校验 WAL 头魔数、版本、页大小、盐值和头校验和,再逐帧验证滚动校验和,只采用最后一个有效提交帧及之前的数据,按最终页数写入并截断。实测一次典型重放:
漏掉 WAL 的后果不是报错,是静默地少一批最近消息——而你不会发现,因为查出来的东西看起来完全正常。
采集器只做只读
从重放副本开始,采集要有一条硬约束:绝不写入。原件和快照都不能被改。
function runSql(sql) {
if (!dbPath || !existsSync(dbPath)) throw new Error('缺少可读取的 --db 重放数据库')
return execFileSync('/usr/bin/sqlite3', [`file:${dbPath}?immutable=1`], {
input: `PRAGMA query_only=ON;\n${sql}\n`,
encoding: 'utf8',
maxBuffer: 512 * 1024 * 1024,
})
}
两道保险叠着用:immutable=1 告诉 SQLite 这个文件不会被改动,跳过锁和 WAL 恢复;query_only=ON 让任何写语句直接失败。少了 immutable=1,SQLite 打开时可能自己发起一次 WAL 恢复,那就已经改了文件。
同时校验 schema 指纹,版本不对就停下:
function messageTables() {
const tables = runSql(
".mode list\nSELECT name FROM sqlite_schema WHERE type='table' AND name LIKE 'tbmsg_%' ORDER BY name;"
).trim().split('\n').filter(Boolean)
if (!tables.length || tables.some((name) => !/^tbmsg_\d{3}$/.test(name))) {
throw new Error('不兼容的消息表结构')
}
return tables
}
消息按 128 张 tbmsg_000 ~ tbmsg_127 分表存。分表数量或命名不符就采集失败,而不是尽力而为——客户端升级可能改 schema,这时候拿到的数据不可信,宁可停下来重新做兼容验证。
会话发现:找不准就不猜
按名称搜索会话,但名称只用于发现,不用于后续同步。
// INVISIBLE 匹配 U+200B ~ U+200D 和 U+FEFF 这几个零宽字符
export function normalizeConversationName(value) {
return String(value ?? '')
.normalize('NFKC')
.toLocaleLowerCase('zh-CN')
.replace(INVISIBLE, '')
.replace(/\s+/gu, ' ')
.trim()
}
NFKC 归一化处理全角半角,零宽字符要单独清掉——群名里粘进一个 U+200B 是很常见的事,肉眼看不出来,=== 却匹配不上。
同名的全部返回,绝不自动选一个。 匹配数不是 1 时由人来确认,采集器不猜。选错会话意味着把另一个群的内容当成语料发布出去,这个错误没法在下游被发现。
稳定引用:群改名不影响它
确认会话之后,一切都用引用,不再用名称,也不用原始 ID。
export function stableConversationRef(secret, provider, providerId) {
const digest = createHmac('sha256', secret)
.update(`${provider}:${providerId}`)
.digest('hex').slice(0, 20)
return `conv_${digest}`
}
export function stableMessageId(conversationRef, providerMessageId) {
return `msg_${createHash('sha256')
.update(`${conversationRef}:${providerMessageId}`)
.digest('hex').slice(0, 20)}`
}
三点值得说清楚:
会话和发送者用 HMAC,消息 ID 用 SHA-256。 会话 ID 和用户 ID 是需要防还原的敏感标识,加密钥;消息 ID 已经建立在 conversationRef 之上,本身不泄露原始信息,用普通哈希即可,好处是同一份导出包在任何机器上都能算出一致的消息 ID。
密钥只在本机。 32 字节随机数,0o600,写在私有目录里,永不进导出包:
function secret() {
mkdirSync(root, { recursive: true })
if (!existsSync(secretPath)) writeFileSync(secretPath, randomBytes(32), { mode: 0o600 })
return readFileSync(secretPath)
}
引用只在同一份数据集内有意义。 它们是给证据回查用的,不能反推回真实会话或用户。换一个密钥重新采集,所有引用都会变——这是设计上的取舍:可回查性优先于跨数据集可比性。
确定性脱敏
脱敏必须是确定性的:同样的输入永远得到同样的输出。否则增量同步时内容哈希会漂移,每次都判定成「变了」。
function redactText(value, names, secretValue) {
let text = String(value ?? '').normalize('NFKC').replace(INVISIBLE, '')
// 用户表里的昵称、实名、别名,按长度倒序替换,避免短名先命中把长名切碎
for (const name of names) {
if (name.length < 2) continue
const ref = `user_${createHmac('sha256', secretValue).update(`person:${name}`).digest('hex').slice(0, 10)}`
text = text.replaceAll(name, ref)
}
return text
.replace(/@[^\s,。,::;;]+/gu, '@用户')
.replace(/https?:\/\/[^\s<>'"))]+/giu, '[链接已脱敏]')
.replace(/[A-Z0-9._%+-]+@[A-Z0-9.-]+\.[A-Z]{2,}/giu, '[邮箱已脱敏]')
.replace(/(?<!\d)1[3-9]\d{9}(?!\d)/g, '[手机号已脱敏]')
.replace(/\b(?:10|127|169\.254|172\.(?:1[6-9]|2\d|3[01])|192\.168)(?:\.\d{1,3}){2,3}\b/g, '[网络地址已脱敏]')
.replace(/(?:\/Users\/|[A-Z]:\\Users\\)[^\\/\s]+/giu, '[用户目录已脱敏]')
.replace(/\b(token|cookie|authorization|password|secret|api[_-]?key)\s*[:=]\s*[^\s,;]+/giu, '$1=[凭据已脱敏]')
}
几个容易踩的点:
姓名按长度倒序替换。 群里同时有「张伟」和「张伟明」时,先替换短的会 把长的切成 user_xxx明。names 数组必须先 sort((a, b) => b.length - a.length)。
两字以下的名字跳过。 「小李」还好,单字昵称会在正文里大面积误伤——「问题在于内存」里的「存」如果恰好是某人的昵称,整段话就毁了。这条是有意的漏网,宁可漏一个短昵称,不要毁掉语料。
组织和产品名要单独维护一张泛化表,把具体品牌换成「桌面客户端」「安全软件」「第三方服务方」这类通名。这张表是项目相关的,没法通用,但它是发布前扫描的主要拦截对象。
内网地址只处理 RFC1918。 公网 IP 不脱,因为它们常常是问题现场的关键信息(比如某个 CDN 节点异常)。
消息解析:不认识就说不认识
function parseText(row) {
let content = {}
try { content = JSON.parse(row.content || '{}') } catch {}
return String(
content.text ?? content.title ?? content.desc ?? content.description ??
({ 2: '[图片]', 4: '[文件]', 1200: '[视频]' }[row.contentType] ?? '[消息]')
)
}
优先读正文字段,媒体类型落固定占位符。未知类型返回 [消息],不伪装成 text。
这条很重要。真实分布里非文本占了三分之一——618 条消息中 59 条视频、46 条图片、30 条卡片、12 条文件。如果把解析不出来的东西塞一段空字符串进去当正文,下游会拿它去做 embedding,得到一堆语义为空的向量。老实标注占位符,第 10 篇的线程重建才知道这里有个洞。
非零 recallStatus 标成 recalled,不进可用消息层——不恢复用户当前已经看不到的撤回消息。
增量协议:三种操作,一个哈希
采集不是一次性的。群还在说话,下次同步要能算出差异。
export function contentHash(message) {
return createHash('sha256').update(JSON.stringify({
conversationRef: message.conversationRef,
sentAt: message.sentAt,
type: message.type,
text: message.text,
attachments: message.attachments ?? [],
status: message.status ?? 'active',
})).digest('hex')
}
这是一个显式白名单,不是把整个对象哈希。 直接 JSON.stringify(message) 会把 messageId 也算进去,而 messageId 恒定,等于白算;更糟的是将来给消息加一个 syncedAt 之类的字段,所有消息的哈希会一起变化,一次同步就把整个语料库重建一遍。字段进哈希要是一个有意的决定。
export function diffMessages(previous = [], current = []) {
const oldById = new Map(previous.map((item) => [item.messageId, item]))
const newById = new Map(current.map((item) => [item.messageId, item]))
const delta = []
for (const item of current) {
const before = oldById.get(item.messageId)
if (!before || before.contentHash !== item.contentHash) {
delta.push({ operation: item.status === 'recalled' ? 'recalled' : 'upsert', messageId: item.messageId, item })
}
}
for (const item of previous) {
if (newById.has(item.messageId)) continue
delta.push({ operation: 'missing', messageId: item.messageId, item })
}
return delta
}
upsert、recalled、missing 三种操作的判定路径与各自的语义边界
upsert 和 recalled 好理解。missing 这个词是刻意选的,它不叫 deleted。
missing 只表示「后一次本地快照里没找到这条」。可能的原因有好几种:本地库做了清理、消息滚出了本地缓存范围、这次快照的时间窗不同。它不能被解释为服务端删除。 下游拿到 missing 时的正确行为是标记待确认,而不是直接从语料里删掉——否则本地缓存一次正常清理,就会把一批有效知识从知识库里抹掉。
同样要写进文档的一句:本地库的覆盖范围不等于服务端的全量历史。 采集报告只能声称「这份副本里的记录全部经过确定性处理」 。
导出包与协议版本
.conversation-exporter/
├── manifest.json 协议版本、导出 ID、会话列表、计数、文件名
├── conversations.jsonl 会话引用、当前名称、类型、最后活动时间、本地消息数
├── messages.jsonl 消息 ID、会话引用、发送者引用、时间、类型、正文、状态、内容哈希
├── delta.jsonl upsert / recalled / missing
└── checkpoint.json 本次导出 ID 与消息状态摘要
采集侧与下游之间的协议边界,以及每一层不得向下传递的内容
这个包就是整个专栏的语言分界线和合规分界线。采集侧是 Node,下游全是 Python;真实会话 ID、用 户 ID、媒体地址、解密密钥全部止步于此。下游只认这五个文件,不知道上游是钉钉、飞书还是别的什么。
版本号要在 manifest 里显式声明,并且消费方必须校验:
export const EXPORT_PROTOCOL_VERSION = 1
export function validateExportManifest(manifest) {
return Boolean(manifest &&
manifest.protocolVersion === EXPORT_PROTOCOL_VERSION &&
typeof manifest.exportId === 'string' &&
Array.isArray(manifest.conversationRefs) &&
typeof manifest.files?.messages === 'string')
}
verify 命令在版本不符时以退出码 2 结束。
为什么不做「向前兼容、忽略不认识的字段」:V2 给消息加了图片资产(assets.jsonl 和 attachments 里的 assetRef)。如果 V1 的消费方宽容地忽略这些字段,它会安静地把所有图片语料丢掉,然后报告一切正常。宁可硬失败:V1 消费方必须拒绝 V2 包。协议升级本来就应该是一次需要人过一眼的事件。
失败处理
这个表比代码更重要。
一条贯穿的原则:所有失败都是停下来,没有一条是「尽力而为」。 采集这层出的错,下游发现不了——一份少了 200 条消息的语料和一份完整语料,检索结果都「看起来正常」。
测试与验收
conversation-exporter doctor --db fixtures/synthetic.replayed.db
conversation-exporter sync --db fixtures/synthetic.replayed.db --output .conversation-exporter
conversation-exporter verify --output .conversation-exporter
仓库里的 fixtures/synthetic.replayed.db 是一份合成的明文 SQLite,schema 和真实结构一致(tbconversation、tbuser_profile_v2、128 张 tbmsg_NNN)。你不需要任何真实聊天数据就能跑通这一篇;跑好的导出包也一并提交在 fixtures/export-bundle/,下一篇可以直接用。
验收标准:
- 同 一份输入跑两次,
messages.jsonl 逐字节相同(脱敏是确定性的)。
- 第二次
sync 的 delta.jsonl 为空。
- 手工改动 fixture 里一条消息的正文,delta 恰好一条
upsert。
- 删掉 fixture 里一条消息,delta 恰好一条
missing,且语料不被自动删除。
- 导出包里 grep 不到原始会话 ID、用户 ID、URL、手机号、邮箱和本地用户名。
- 篡改 manifest 的
protocolVersion,verify 退出码为 2。
下一篇处理另一个来源:上传的 PDF、Word 和 TXT。它有完整的标题层级,切分方式和群聊完全不同——这个差异会在第 10 篇变成本专栏最核心的一个设计决定。