From 1180d8d742a75281a5667f49afa2a3688caf4fc3 Mon Sep 17 00:00:00 2001 From: cy1223710788 Date: Wed, 8 Jul 2026 20:54:12 +0800 Subject: [PATCH] fix: add dedup, idempotent start, and ensurePrivateSession mutex for Feishu WebSocket Three concurrency fixes for duplicate message handling on Feishu/Lark: 1. **client.ts - message dedup**: WebSocket may deliver the same im.message.receive_v1 event twice. Cache message_id in a Map with 30s TTL and skip duplicates. 2. **client.ts - idempotent start()**: Guard against repeated start() calls that could create multiple WSClient connections. 3. **p2p.ts - ensurePrivateSession mutex**: Per-chatId promise lock so concurrent handleMessage calls don't both enter the session-creation critical section (TOCTOU race), creating duplicate sessions. Without these fixes, every message on Feishu could trigger two replies and two OpenCode sessions. --- src/feishu/client.ts | 28 ++++++++++++++++++++++++++++ src/handlers/p2p.ts | 20 ++++++++++++++++++++ 2 files changed, 48 insertions(+) diff --git a/src/feishu/client.ts b/src/feishu/client.ts index b561ede..beeda61 100644 --- a/src/feishu/client.ts +++ b/src/feishu/client.ts @@ -42,6 +42,10 @@ class FeishuClient extends EventEmitter { // 机器人自身信息 private botOpenId: string | null = null; + // message_id 去重缓存:WebSocket 可能重复投递同一条消息 + private _recentMessageIds = new Map(); + private readonly _RECENT_MESSAGE_TTL_MS = 30000; + constructor() { super(); this.client = new lark.Client({ @@ -153,6 +157,10 @@ class FeishuClient extends EventEmitter { // 启动长连接 async start(): Promise { + if (this.connectionState === 'connected' || this.connectionState === 'connecting') { + console.log('[飞书] 长连接已连接/连接中,忽略重复调用'); + return; + } console.log('[飞书] 正在启动长连接...'); this.connectionState = 'connecting'; @@ -243,6 +251,26 @@ class FeishuClient extends EventEmitter { return; } + // 消息去重:飞书 WebSocket 可能重复投递同一条消息 + const msgId = message.message_id; + if (msgId) { + const now = Date.now(); + const recent = this._recentMessageIds.get(msgId); + if (recent !== undefined && now - recent < this._RECENT_MESSAGE_TTL_MS) { + console.log(`[飞书] 忽略重复消息: messageId=${msgId.slice(0, 16)}...`); + return; + } + this._recentMessageIds.set(msgId, now); + // 惰性清理:每收到 500 条消息后扫一遍过期记录 + if (this._recentMessageIds.size >= 500) { + for (const [id, ts] of this._recentMessageIds) { + if (now - ts >= this._RECENT_MESSAGE_TTL_MS) { + this._recentMessageIds.delete(id); + } + } + } + } + const msgType = message.message_type; let content = ''; let parsedContent: Record | null = null; diff --git a/src/handlers/p2p.ts b/src/handlers/p2p.ts index 7f239ff..d7a775b 100644 --- a/src/handlers/p2p.ts +++ b/src/handlers/p2p.ts @@ -25,6 +25,9 @@ export class P2PHandler { private createChatDirectoryInputMap: Map = new Map(); private createChatNameInputMap: Map = new Map(); + // per-chatId 互斥锁,防 ensurePrivateSession TOCTOU 竞态 + private _ensureSessionLocks = new Map>(); + private async safeReply( messageId: string | undefined, chatId: string | undefined, @@ -470,6 +473,23 @@ export class P2PHandler { } private async ensurePrivateSession(chatId: string, senderId: string): Promise { + // per-chatId 互斥锁:等第一个完成,防止并发的 handleMessage 都走到创建 session + const pending = this._ensureSessionLocks.get(chatId); + if (pending) { + return await pending; + } + const promise = this._ensurePrivateSessionImpl(chatId, senderId); + this._ensureSessionLocks.set(chatId, promise); + try { + return await promise; + } finally { + if (this._ensureSessionLocks.get(chatId) === promise) { + this._ensureSessionLocks.delete(chatId); + } + } + } + + private async _ensurePrivateSessionImpl(chatId: string, senderId: string): Promise { const current = chatSessionStore.getSession(chatId); if (current?.sessionId) { const missing = await this.isSessionMissingInOpenCode(current.sessionId);