接收回调
Node.js SDK 支持事件回调、同步回调和 URL 验证,处理签名、随机数、防重放、事件去重和分派。控制台配置见回调,协议见回调概述。
创建处理器
单进程测试或单实例部署可以显式使用内存存储:
import {
CallbackHandler, MemoryNonceStore, MemoryDedupStore, allow, reject,
} from '@deeprespond/im-server-sdk/callback';
const handler = new CallbackHandler({
secrets: {
[process.env.IM_CALLBACK_ENDPOINT_ID!]: [process.env.IM_CALLBACK_SECRET!],
},
nonceStore: new MemoryNonceStore(),
dedupStore: new MemoryDedupStore(),
});
handler.on('user.created', async (event, data, ctx) => {
console.log(event.id, data.username);
// 业务数据库操作可以使用 ctx.signal。
});
handler.onHook('message.before_send', async (req, ctx) => {
if (typeof req.body.text === 'string' && req.body.text.includes('spam')) return reject('spam', '消息不符合发送规则');
return allow();
});secrets 的键是回调地址的 endpoint ID,值是密钥数组,不能传单个字符串。轮换期间同时提供新旧密钥。也可用 async (endpointId, { signal }) => readonly string[] 动态读取;找不到返回空数组。取密钥失败或超时返回 503。
| 选项 | 默认值 | 说明 |
|---|---|---|
secrets | 必填 | endpoint ID 对应的密钥数组,或密钥查询函数 |
nonceStore、dedupStore | 必填 | 随机数存储、事件去重存储 |
queue | 无 | 持久队列,事件收到后交给 enqueue,见事件处理 |
hookFallback | null | 同步回调失败时返回 503;可显式设置 allow() 或 reject() 等结果 |
maxConcurrentHooks | 100 | 同时运行的同步回调处理函数上限,包含超时后仍在运行的函数 |
claimTtlMs | 30000 | 事件处理的占用时限(毫秒) |
maxBodyBytes | 2097152 | 最大请求体字节数,默认 2 MiB |
logger | error、warn 输出到 console | 日志级别和 sink |
共享存储
多进程、集群、Serverless 或多台服务器部署时使用共享 Redis;内存存储随进程退出丢失,也不能跨实例防重放和去重。
import { Redis } from 'ioredis';
import { CallbackHandler } from '@deeprespond/im-server-sdk/callback';
import { ioredisStores } from '@deeprespond/im-server-sdk/redis';
const redis = new Redis(process.env.REDIS_URL!);
const { nonceStore, dedupStore } = ioredisStores(redis, {
keyPrefix: 'my-service:im:',
retentionMs: 7 * 24 * 60 * 60 * 1000,
});
const handler = new CallbackHandler({
secrets: { [process.env.IM_CALLBACK_ENDPOINT_ID!]: [process.env.IM_CALLBACK_SECRET!] },
nonceStore,
dedupStore,
});使用 node-redis 时调用 nodeRedisStores(redis, options),先连接 Redis。默认 keyPrefix 为 drim:,事件去重记录默认保留 7 天。所有接收相同回调的实例使用相同存储和前缀,Redis 连接由业务应用管理。
Express
import express from 'express';
import { expressHandler } from '@deeprespond/im-server-sdk/callback/express';
const app = express();
app.post('/im/callback', expressHandler(handler));
app.use(express.json());
app.listen(3000);适用于 Express 4、5。回调路由放在 JSON 解析中间件之前,适配器直接读取原始字节。如果先使用 express.raw({ type: 'application/json', limit: '2mb' }),也可将原始 Buffer 交给适配器。已经解析成对象且未保留原文的请求无法验签;不要 JSON.stringify 后再验签。
NestJS
使用 Express 平台,以 rawBody: true 创建应用,保留原始请求体:
import { NestFactory } from '@nestjs/core';
import type { NestExpressApplication } from '@nestjs/platform-express';
import { expressHandler } from '@deeprespond/im-server-sdk/callback/express';
// AppModule 为业务应用已有的 NestJS 模块。
const app = await NestFactory.create<NestExpressApplication>(AppModule, { rawBody: true });
app.useBodyParser('json', { limit: '2mb' });
app.use('/im/callback', expressHandler(handler));
await app.listen(3000);Fastify
import Fastify from 'fastify';
import { fastifyCallback } from '@deeprespond/im-server-sdk/callback/fastify';
const app = Fastify();
await app.register(fastifyCallback(handler), { path: '/im/callback' });
await app.listen({ port: 3000 });适用于 Fastify 5,插件在自己的路由作用域中保存原始字节,不改变其他路由的 JSON 解析。
Koa
import Koa from 'koa';
import { koaHandler, type KoaContextLike } from '@deeprespond/im-server-sdk/callback/koa';
const app = new Koa();
const receive = koaHandler(handler);
app.use(async (ctx, next) => {
if (ctx.path === '/im/callback') await receive(ctx as KoaContextLike);
else await next();
});
app.listen(3000);适用于 Koa 2、3。将此路由放在 body parser 之前;如已使用 koa-bodyparser 或 @koa/bodyparser,需要保留 ctx.request.rawBody 并将 jsonLimit 设为 2mb。
Next.js 与 Web 标准 Request
适配器使用 Request 和 Response,但必须在 Node.js 运行时运行。Next.js App Router 的 app/im/callback/route.ts:
import { fetchHandler } from '@deeprespond/im-server-sdk/callback/fetch';
import { handler } from '@/lib/im-callback'; // 导出上文创建并注册好的处理器。
export const runtime = 'nodejs';
export const POST = fetchHandler(handler);不要预先调用 request.json() 消耗请求体。Serverless 的多个实例使用共享存储;该适配器也可接入使用 Node.js 的 Hono 等框架。
原生 node:http 与自定义框架
import { createServer } from 'node:http';
import { nodeHandler } from '@deeprespond/im-server-sdk/callback/node';
const receive = nodeHandler(handler);
createServer((req, res) => {
if (req.url === '/im/callback') {
void receive(req, res).catch(() => { res.statusCode = 500; res.end(); });
} else {
res.statusCode = 404;
res.end();
}
}).listen(3000);自定义框架可以调用 await handler.handle({ headers, body }),body 为原始 Uint8Array 或 ArrayBuffer;把返回的 status、headers、body 原样写入 HTTP 响应。
事件处理
handler.on(type, fn) 按事件名推断 data 的类型,fn 参数是 (event, data, ctx)。同一事件只能注册一次,事件 ID 为 event.id。已成功处理的事件不重复分派;整批中部分失败时成功项仍保留去重记录,重投时只处理未完成项。
处理函数抛出错误时返回 503,服务端按投递规则重试。需要指定等待时间时,抛出 new RetryLater(seconds),从 callback 子路径导入。未知事件可用 handler.onUnknown((event, ctx) => ...) 处理;具体事件结构见事件回调。
长任务可以传入 queue: { async enqueue(event, ctx) { /* 写入持久队列 */ } }。enqueue 成功表示已可靠接收,失败或超时返回 503。配置 queue 后不能再用 on、onAny、onUnknown 注册事件处理函数;同步回调仍可用 onHook。队列消费者仍应按事件 ID 保证业务操作幂等。
SDK 去重不能使业务写入和去重提交成为同一个事务。进程在业务写入后、去重提交前退出时,事件可能再次处理,因此业务数据库也应以事件 ID 防止重复副作用。
同步回调
handler.onHook(name, fn) 的参数为 (req, ctx)。支持 message.before_send、friend.before_add、group.before_join、chatroom.message_before_send;返回值必须使用 SDK 的结果函数创建:
| 结果 | 用途 |
|---|---|
allow() | 放行 |
reject(reason?, message?) | 拒绝;reason 为 1~64 个小写字母、数字或 _.-,message 为最多 256 个字符的单行提示 |
modify({ body?, ext? }) | 替换消息正文或扩展字段,适用范围见同步回调 |
partial([{ username, reason?, message? }]) | 部分拒绝,仅用于 group.before_join 的邀请和建群目标 |
ctx.timeoutMs 是服务端给出的等待时间,ctx.deadline 是 SDK 预留 50 毫秒后的截止时刻,ctx.signal 到时触发。将 signal 传给自己的数据库或 fetch 调用,避免超时后任务仍长期占用资源。未注册、超时、抛错或返回不合法结果时,默认返回 503,服务端按配置的失败策略处理;hookFallback 是业务显式选择的替代结果。
验证、签名与排查
处理器自动响应 URL 验证,业务无需单独注册验证函数。签名检查使用原始请求体,校验时间戳及随机数;算法和密钥轮换见回调安全。
常见接入问题:
- 签名失败:核对 endpoint ID、密钥数组和原始字节是否被中间件改动。
- 时间戳失效:核对服务端和接收端的系统时钟。
- 请求体超过限制:同时检查 SDK、框架和反向代理的请求体上限。
- 返回 503:查看密钥查询、Redis、处理函数异常和同步回调时限。
本地测试
import assert from 'node:assert/strict';
import { CallbackHandler, MemoryNonceStore, MemoryDedupStore } from '@deeprespond/im-server-sdk/callback';
import { eventBody, signedRequest } from '@deeprespond/im-server-sdk/testing';
const body = eventBody(['user.created', { username: 'alice' }]);
const endpointId = JSON.parse(Buffer.from(body).toString('utf8')).endpoint_id as string;
const handler = new CallbackHandler({
secrets: { [endpointId]: ['test-secret'] },
nonceStore: new MemoryNonceStore(),
dedupStore: new MemoryDedupStore(),
});
let username = '';
handler.on('user.created', (_event, data) => { username = data.username; });
const response = await handler.handle(signedRequest('test-secret', endpointId, body));
assert.equal(response.status, 204);
assert.equal(username, 'alice');
await handler.close();hookBody(name, data, { timeoutMs?, test? }) 构造同步回调请求体,verifyBody(challenge) 构造 URL 验证请求体,再用 signedRequest 添加签名。也可以把构造出的原始字节发给本地框架路由,验证中间件顺序。
停止服务
先停止接收新请求,再 await handler.close({ timeoutMs: 5000 }) 等待在途任务;超时后通过处理函数的 ctx.signal 发出取消信号。最后关闭由业务管理的 Redis 连接。
独立频道回调
| 事件 | 强类型模型 |
|---|---|
rtc_channel.created | RTCChannelCreated |
rtc_channel.updated | RTCChannelUpdated |
rtc_channel.closing | RTCChannelClosing |
rtc_channel.closed | RTCChannelClosed |
rtc_channel.session_joined | RTCChannelSessionJoined |
rtc_channel.session_leaving | RTCChannelSessionLeaving |
rtc_channel.session_left | RTCChannelSessionLeft |
rtc_channel.member_banned | RTCChannelMemberBanned |
rtc_channel.member_unbanned | RTCChannelMemberUnbanned |
这九种事件可通过现有分派器处理,沿用签名验证、去重和错误处理。字段见频道事件;部分字段可缺省,不应假设非会话事件包含目标用户。频道与原通话的事件名和模型分开,配置内部广播不作为回调。
