Skip to content

Commit 2baa8e0

Browse files
author
linyuan.yang
committed
钉钉 qq
1 parent b021ed5 commit 2baa8e0

20 files changed

Lines changed: 1353 additions & 6 deletions
Lines changed: 4 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,4 @@
1+
/node_modules
2+
/dist
3+
/package-lock.json
4+
/tsconfig.tsbuildinfo
Lines changed: 21 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,21 @@
1+
{
2+
"name": "channel.dingtalk",
3+
"version": "0.0.1",
4+
"description": "DingTalk channel integration (Stream mode)",
5+
"main": "dist/index.js",
6+
"types": "dist/index.d.ts",
7+
"scripts": {
8+
"build": "rimraf dist && tsc -b --force"
9+
},
10+
"dependencies": {
11+
"axios": "catalog:",
12+
"channel.base": "workspace:*",
13+
"dingtalk-stream": "^2.1.6-beta.1"
14+
},
15+
"private": true,
16+
"devDependencies": {
17+
"@types/node": "catalog:",
18+
"rimraf": "catalog:",
19+
"typescript": "catalog:"
20+
}
21+
}
Lines changed: 44 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,44 @@
1+
import { DingtalkService } from './DingtalkService';
2+
import { AbstractChatProvider, parseMessages2Text, GlobalLoggerService } from 'channel.base';
3+
4+
const getLogger = () => GlobalLoggerService.getLogger('DingtalkChatProvider.ts');
5+
6+
/**
7+
* 钉钉无法原地编辑普通 markdown 消息(除非走 AI Card),因此 Provider 仅做:
8+
* - 在 onProcessStart 时不主动发"Processing..."(避免刷屏);
9+
* - 累积输出,在 onProcessEnd 时一次性 sessionWebhook 发出最终 markdown。
10+
*
11+
* 流式过程中并不更新会话——交由调用方 throttle 控制;调用方可显式 flush()。
12+
*/
13+
export class DingtalkChatProvider extends AbstractChatProvider {
14+
private sessionId = '';
15+
private sent = false;
16+
private atUsers: string[] = [];
17+
18+
constructor(private dingtalkService: DingtalkService) {
19+
super();
20+
}
21+
22+
init(sessionId: string, atUsers?: string[]): this {
23+
this.sessionId = sessionId;
24+
this.atUsers = atUsers ?? [];
25+
return this;
26+
}
27+
28+
/** 立即把当前累积内容作为最终消息发送一次(幂等) */
29+
async flush(): Promise<void> {
30+
if (this.sent || !this.sessionId) return;
31+
this.sent = true;
32+
try {
33+
const text = parseMessages2Text(this.getDisplayMessages()).trim();
34+
if (!text) return;
35+
await this.dingtalkService.sendMarkdown(this.sessionId, text, this.atUsers);
36+
} catch (e: any) {
37+
getLogger()?.error(`flush exception: ${e.message || e}`, e.stack);
38+
}
39+
}
40+
41+
protected onMessagesUpdated(): void {
42+
// 累积阶段不发送,等 flush
43+
}
44+
}
Lines changed: 241 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,241 @@
1+
import axios from 'axios';
2+
import { DWClient, EventAck, type DWClientDownStream, type EventAckData } from 'dingtalk-stream';
3+
import {
4+
IChannelService, ChannelSessionHandler, SessionService,
5+
type ILogger, type MessageContent,
6+
} from 'channel.base';
7+
import { DingtalkSessionHandler, type DingtalkMessageArgs, type DingtalkActionArgs } from './DingtalkSessionHandler';
8+
9+
const TOPIC_ROBOT = '/v1.0/im/bot/messages/get';
10+
const TOPIC_CARD = '/v1.0/card/instances/callback';
11+
12+
/** 维护 sessionWebhook 与对话信息(含过期时间) */
13+
interface SessionEntry {
14+
conversationId: string;
15+
conversationType: '1' | '2'; // 1 = 单聊, 2 = 群聊
16+
sessionWebhook: string;
17+
sessionWebhookExpiredTime: number;
18+
senderStaffId: string;
19+
robotCode: string;
20+
lastMsgId: string;
21+
}
22+
23+
/** 入站机器人消息(DingTalk Stream) */
24+
interface RobotIncoming {
25+
msgtype: string;
26+
msgId: string;
27+
text?: { content: string };
28+
conversationId: string;
29+
conversationType: '1' | '2';
30+
senderStaffId: string;
31+
senderNick: string;
32+
senderId: string;
33+
sessionWebhook: string;
34+
sessionWebhookExpiredTime: number;
35+
robotCode: string;
36+
chatbotUserId: string;
37+
isAdmin: boolean;
38+
createAt: number;
39+
}
40+
41+
export interface DingtalkServiceOptions {
42+
clientId: string;
43+
clientSecret: string;
44+
/** 群聊场景下回复时是否 @ 发送者,默认 false */
45+
atSenderOnReply?: boolean;
46+
logger?: ILogger;
47+
filterEvent: (eventId: string) => Promise<boolean>;
48+
onReceiveMessage: (
49+
userId: string,
50+
args: DingtalkMessageArgs,
51+
query: MessageContent,
52+
) => Promise<void>;
53+
onTriggerAction: (
54+
userId: string,
55+
args: DingtalkActionArgs,
56+
) => Promise<void>;
57+
}
58+
59+
export class DingtalkService implements IChannelService {
60+
private client: DWClient;
61+
private logger?: ILogger;
62+
private opts: DingtalkServiceOptions;
63+
private sessions = new Map<string, SessionEntry>();
64+
65+
constructor(options: DingtalkServiceOptions) {
66+
this.opts = options;
67+
this.logger = options.logger;
68+
this.client = new DWClient({
69+
clientId: options.clientId,
70+
clientSecret: options.clientSecret,
71+
keepAlive: true,
72+
debug: false,
73+
});
74+
}
75+
76+
createSessionHandler(session: SessionService): ChannelSessionHandler {
77+
return new DingtalkSessionHandler(session, this);
78+
}
79+
80+
async sendTextToSession(sessionId: string, text: string): Promise<void> {
81+
await this.sendMarkdown(sessionId, text);
82+
}
83+
84+
async sendTextToUser(userId: string, text: string): Promise<void> {
85+
// userId 为 senderStaffId;只能向最近发过消息的私聊用户回写。
86+
const entry = this.findSessionByStaffId(userId);
87+
if (!entry) {
88+
this.logger?.warn(`Dingtalk sendTextToUser: no session for userId=${userId}`);
89+
return;
90+
}
91+
await this.sendMarkdown(entry.conversationId, text);
92+
}
93+
94+
dispose(): void {
95+
try { this.client.disconnect(); } catch (_) {}
96+
this.sessions.clear();
97+
}
98+
99+
/** 注册 Stream 监听并连接 */
100+
async start(): Promise<void> {
101+
this.client.registerCallbackListener(TOPIC_ROBOT, (msg) => this.handleRobotMessage(msg));
102+
this.client.registerCallbackListener(TOPIC_CARD, (msg) => this.handleCardCallback(msg));
103+
await this.client.connect();
104+
this.logger?.info('Dingtalk Stream connected');
105+
}
106+
107+
/** 通过 sessionWebhook 发送 markdown 消息 */
108+
async sendMarkdown(sessionId: string, text: string, atUsers?: string[]): Promise<void> {
109+
const entry = this.sessions.get(sessionId);
110+
if (!entry) {
111+
this.logger?.warn(`Dingtalk sendMarkdown: unknown session ${sessionId}`);
112+
return;
113+
}
114+
if (entry.sessionWebhookExpiredTime && entry.sessionWebhookExpiredTime < Date.now()) {
115+
this.logger?.warn(`Dingtalk sendMarkdown: sessionWebhook expired for ${sessionId}`);
116+
return;
117+
}
118+
const isGroup = entry.conversationType === '2';
119+
const at = isGroup && atUsers?.length
120+
? { atUserIds: atUsers, isAtAll: false }
121+
: undefined;
122+
const payload: any = {
123+
msgtype: 'markdown',
124+
markdown: { title: 'AI', text },
125+
...(at ? { at } : {}),
126+
};
127+
try {
128+
await axios.post(entry.sessionWebhook, payload, {
129+
headers: { 'Content-Type': 'application/json' },
130+
timeout: 15000,
131+
});
132+
} catch (e: any) {
133+
this.logger?.error(`Dingtalk sendMarkdown failed: ${e.message}`);
134+
}
135+
}
136+
137+
/** 通过 sessionWebhook 发送 ActionCard(带按钮,用于审批 / Ask 等) */
138+
async sendActionCard(
139+
sessionId: string,
140+
title: string,
141+
text: string,
142+
buttons: Array<{ title: string; actionURL: string }>,
143+
): Promise<void> {
144+
const entry = this.sessions.get(sessionId);
145+
if (!entry) return;
146+
const payload: any = {
147+
msgtype: 'actionCard',
148+
actionCard: {
149+
title,
150+
text,
151+
btnOrientation: '0',
152+
btns: buttons.map(b => ({ title: b.title, actionURL: b.actionURL })),
153+
},
154+
};
155+
try {
156+
await axios.post(entry.sessionWebhook, payload, {
157+
headers: { 'Content-Type': 'application/json' },
158+
timeout: 15000,
159+
});
160+
} catch (e: any) {
161+
this.logger?.error(`Dingtalk sendActionCard failed: ${e.message}`);
162+
}
163+
}
164+
165+
/** 私下与 Handler 共享 session 信息 */
166+
getSession(sessionId: string): SessionEntry | undefined {
167+
return this.sessions.get(sessionId);
168+
}
169+
170+
/** 处理 IM 机器人消息 */
171+
private handleRobotMessage(msg: DWClientDownStream): EventAckData {
172+
let payload: RobotIncoming;
173+
try {
174+
payload = JSON.parse(msg.data) as RobotIncoming;
175+
} catch (e: any) {
176+
this.logger?.error(`Dingtalk robot message parse failed: ${e.message}`);
177+
return { status: EventAck.SUCCESS };
178+
}
179+
180+
// 同步处理仅做 ACK,实际逻辑放到异步执行
181+
this.processRobotMessage(payload).catch((e) => {
182+
this.logger?.error(`Dingtalk processRobotMessage error: ${e.stack || e.message}`);
183+
});
184+
return { status: EventAck.SUCCESS };
185+
}
186+
187+
private async processRobotMessage(payload: RobotIncoming): Promise<void> {
188+
if (!await this.opts.filterEvent(`dingtalk_message_${payload.msgId}`)) return;
189+
190+
const { conversationId, conversationType, senderStaffId, senderNick, sessionWebhook, sessionWebhookExpiredTime, robotCode, msgId } = payload;
191+
const sessionId = conversationId;
192+
193+
this.sessions.set(sessionId, {
194+
conversationId,
195+
conversationType,
196+
sessionWebhook,
197+
sessionWebhookExpiredTime,
198+
senderStaffId,
199+
robotCode,
200+
lastMsgId: msgId,
201+
});
202+
203+
if (payload.msgtype !== 'text' || !payload.text?.content) {
204+
this.logger?.warn(`Dingtalk: unsupported msgtype=${payload.msgtype}`);
205+
return;
206+
}
207+
208+
const query = payload.text.content.trim();
209+
if (!query) return;
210+
211+
const args: DingtalkMessageArgs = {
212+
sessionId,
213+
msgId,
214+
conversationId,
215+
conversationType,
216+
senderStaffId,
217+
senderNick,
218+
robotCode,
219+
atSenderOnReply: this.opts.atSenderOnReply ?? false,
220+
};
221+
await this.opts.onReceiveMessage(senderStaffId, args, query);
222+
}
223+
224+
/** 处理卡片按钮回调(ActionCard 不支持,AI Card 才会触发;此处保留扩展点) */
225+
private handleCardCallback(msg: DWClientDownStream): EventAckData {
226+
try {
227+
const payload = JSON.parse(msg.data || '{}');
228+
this.logger?.info(`Dingtalk card callback: ${JSON.stringify(payload).slice(0, 300)}`);
229+
} catch (e: any) {
230+
this.logger?.error(`Dingtalk card callback parse failed: ${e.message}`);
231+
}
232+
return { status: EventAck.SUCCESS };
233+
}
234+
235+
private findSessionByStaffId(staffId: string): SessionEntry | undefined {
236+
for (const entry of this.sessions.values()) {
237+
if (entry.senderStaffId === staffId && entry.conversationType === '1') return entry;
238+
}
239+
return undefined;
240+
}
241+
}

0 commit comments

Comments
 (0)