接收回调
IM 服务向你配置的回调地址发送三种请求:事件回调(含消息抄送)、同步回调和 URL 验证,见事件回调概述。callback 包处理这些请求:校验签名、防止重放、按事件 ID 去重,把事件和同步回调分派给你的处理函数,并按服务端的要求返回状态码。你只需要写处理函数。
import "github.com/deeprespond/im/sdk/server/go/callback"创建处理器
package main
import (
"context"
"encoding/json"
"log"
"log/slog"
"net/http"
"os"
"strings"
"time"
"github.com/deeprespond/im/sdk/server/go/callback"
)
func main() {
fallback := callback.Allow() // 同步回调的处理函数出错或超时时放行
h, err := callback.NewHandler(callback.Config{
// 回调地址 ID → 密钥(控制台新建地址时显示;轮换期间新旧两个都配置,新的在前)
Secrets: callback.StaticSecrets{
os.Getenv("IM_CALLBACK_ENDPOINT_ID"): strings.Split(os.Getenv("IM_CALLBACK_SECRETS"), ","),
},
Nonces: callback.NewMemoryNonceStore(), // 多个实例部署时换成共享存储
Dedup: callback.NewMemoryDedupStore(24 * time.Hour), // 同上
Events: callback.EventHandlers{
MessageSent: func(ctx context.Context, ev *callback.Event, d *callback.MessageSent) error {
if d.Truncated || !d.ContentAvailable {
return nil // 内容被截短或已不可用:需要时按 MessageID 调用 OpenAPI 读取
}
// d.Body 与请求体中的逐字节相同,可以直接存档
log.Printf("消息 %s(会话 %s 序号 %d):%s", d.MessageID, d.ConversationID, d.Seq, d.Body)
return nil
},
UserCreated: func(ctx context.Context, ev *callback.Event, d *callback.UserCreated) error {
log.Printf("新用户 %s(%s)", d.Username, d.CreatedVia)
return nil
},
},
Hooks: callback.HookHandlers{
MessageBeforeSend: func(ctx context.Context, req *callback.MessageBeforeSend) (callback.Verdict, error) {
var body struct {
Text string `json:"text"`
}
_ = json.Unmarshal(req.Body, &body)
if strings.Contains(body.Text, "加微信") {
return callback.Reject("forbidden_word", "消息中有不允许的内容"), nil
}
return callback.Allow(), nil
},
},
HookFallback: &fallback,
Logger: slog.Default(),
})
if err != nil {
log.Fatal(err)
}
mux := http.NewServeMux()
mux.Handle("POST /im/callback", h)
log.Fatal(http.ListenAndServe(":8090", mux))
}callback.Handler 实现了 http.Handler,对每个请求依次:
- 读取原始请求体(最多
MaxBodyBytes,默认 2 MB),签名按原始字节计算; - 按请求头
X-IM-Endpoint-Id找到这个地址的密钥,校验签名和时间戳(与当前时间相差不超过 5 分钟),再核对请求体中的endpoint_id与请求头一致,见签名与安全; - 签名通过后记录随机数
X-IM-Nonce,10 分钟内重复的拒绝(防重放); - 按请求体中的
kind分派:URL 验证直接回应,同步回调交给Hooks,事件回调逐个事件去重后交给Events(或Queue)。
配置项
| 字段 | 类型 | 说明 |
|---|---|---|
Secrets | callback.SecretProvider | 必填。按回调地址 ID 返回这个地址的密钥列表,见密钥 |
Nonces | callback.NonceStore | 必填。随机数的存储,见随机数和去重的存储 |
Dedup | callback.DedupStore | 必填。事件 ID 去重的存储 |
Events | callback.EventHandlers | 事件的处理函数(同步方式),见处理事件回调 |
Queue | callback.EventQueue | 把事件交给你的持久队列(异步方式),与 Events 二选一,见异步处理 |
Hooks | callback.HookHandlers | 同步回调的处理函数,见处理同步回调 |
HookFallback | *callback.Verdict | 同步回调的处理函数超时、出错或结果不合规时返回的结果;nil 表示返回 503,由服务端按应用设置的失败策略处理 |
ClaimTTL | time.Duration | 处理一个事件的占用时限,默认 30 秒;超过后同一事件可以在别的请求中重新处理 |
MaxBodyBytes | int64 | 请求体的上限,默认 2 MB,超过的返回 413 |
MaxConcurrentHooks | int | 同时运行的同步回调处理函数的上限,默认 100;已满时直接按 HookFallback(或 503)返回 |
Clock | func() time.Time | 当前时间,测试用,默认 time.Now |
Logger | *slog.Logger | 日志,默认 slog.Default()。签名校验失败、同步回调的结果不合规记错误日志,请据此告警 |
密钥
每个回调地址的密钥不同(cbs_ 开头),在控制台新建地址时显示。callback.StaticSecrets 是固定的“地址 ID → 密钥列表”;密钥保存在配置中心、需要随时更新的,用 callback.SecretFunc 或自己实现 SecretProvider:
// 作为 callback.Config 的 Secrets
func secretsFromConfigCenter() callback.SecretProvider {
return callback.SecretFunc(func(ctx context.Context, endpointID string) ([]string, error) {
return loadCallbackSecrets(ctx, endpointID) // 从配置中心读取;找不到时返回空列表
})
}- 轮换密钥:控制台轮换后,服务端在宽限期内同时用新旧两个密钥签名。先把新密钥加入这个地址的列表(新的在前),宽限期结束后去掉旧的,见轮换密钥。
- 找不到地址的请求返回
401;SecretProvider返回错误时返回503,服务端稍后重发。 - 地址 ID 在全平台唯一,多个应用共用一个接收端时,一个处理器即可:把各应用的地址都配置进去,处理函数中用
callback.RequestFrom(ctx).App(AppKey)区分应用。
挂到 Web 框架上
Handler 是标准的 http.Handler,可以直接挂到 net/http 的 ServeMux、chi 等路由上。签名按原始请求体计算,不要在它之前加会读取或改写请求体的中间件(如解压、把请求体读出来记日志后不放回的中间件),也不要让网关或反向代理改动请求体。
Gin 用 gin.WrapH:
package main
import (
"log"
"os"
"time"
"github.com/gin-gonic/gin"
"github.com/deeprespond/im/sdk/server/go/callback"
)
func main() {
h, err := callback.NewHandler(callback.Config{
Secrets: callback.StaticSecrets{os.Getenv("IM_CALLBACK_ENDPOINT_ID"): {os.Getenv("IM_CALLBACK_SECRET")}},
Nonces: callback.NewMemoryNonceStore(),
Dedup: callback.NewMemoryDedupStore(24 * time.Hour),
Events: callback.EventHandlers{ /* …… */ },
})
if err != nil {
log.Fatal(err)
}
r := gin.New()
r.POST("/im/callback", gin.WrapH(h))
log.Fatal(r.Run(":8090"))
}Echo 用 e.POST("/im/callback", echo.WrapHandler(h))。
在云函数等自己读取请求的环境中,用与框架无关的 HandleRequest:传入请求头和原始请求体,得到状态码、响应头和响应体。
// 云函数的入口:headers 为请求头,body 为原始请求体(如果平台给的是 base64,先解码)
func handleCallback(ctx context.Context, h *callback.Handler, headers map[string]string, body []byte) (int, map[string]string, string) {
header := http.Header{}
for k, v := range headers {
header.Set(k, v)
}
status, respHeader, respBody := h.HandleRequest(ctx, header, body)
out := map[string]string{}
for k := range respHeader {
out[k] = respHeader.Get(k)
}
return status, out, string(respBody)
}返回的状态码
| 情况 | 状态码 |
|---|---|
事件全部处理成功(含重复的、交给 Unknown 或被忽略的) | 204 |
| URL 验证 | 200,{"challenge":"..."} |
| 同步回调 | 200 和结果;处理函数出错、超时、结果不合规且没有设置 HookFallback 时 503 |
签名、时间戳、随机数校验失败,地址 ID 不在配置中,请求体中的 endpoint_id 与请求头不同 | 401 |
请求体不是合法的回调请求(不是 JSON、kind 不认识、事件缺少 id 或 type) | 400 |
请求体超过 MaxBodyBytes | 413 |
不是 POST | 405 |
有事件处理失败、正被另一个请求处理、入队失败;存储或 SecretProvider 出错 | 503,服务端稍后重发 |
服务端对任何非 2xx 都按失败重试,密钥配置错误时同一批事件会重试 24 小时,持续失败的地址会被自动停用(见投递与重试)。签名校验失败时 SDK 记错误日志,日志中写明可能的原因(密钥配置错误、请求体被改动、本机时钟偏差),请据此告警。
处理事件回调
callback.EventHandlers 中每种事件一个字段,签名为 func(ctx context.Context, ev *callback.Event, d *类型) error,d 是按事件类型解析好的 data:
func eventHandlers() callback.EventHandlers {
return callback.EventHandlers{
GroupMembersAdded: func(ctx context.Context, ev *callback.Event, d *callback.GroupMembersAdded) error {
// 不保证顺序:用版本号判断新旧,旧事件不覆盖新状态
return saveMemberCount(ctx, d.GroupID, d.MemberVersion, d.MemberCount)
},
UserStatusChanged: func(ctx context.Context, ev *callback.Event, d *callback.UserStatusChanged) error {
if ev.Test {
return nil // 控制台“测试发送”的事件,不执行业务操作
}
if d.Status == "deleted" {
return deleteLocalProfile(ctx, d.Username)
}
return nil
},
// 所有事件(含不认识的)先交给 Any,返回错误即失败;用于统一记录或转发
Any: func(ctx context.Context, ev *callback.Event) error {
slog.DebugContext(ctx, "收到事件", "id", ev.ID, "type", ev.Type, "attempt", ev.Attempt)
return nil
},
// SDK 还不认识的事件类型(服务端新增的);不设置时忽略
Unknown: func(ctx context.Context, ev *callback.Event) error {
slog.WarnContext(ctx, "SDK 还不认识的事件", "type", ev.Type, "data", string(ev.Data))
return nil
},
}
}- 处理成功返回
nil;返回错误时这一批返回503,服务端稍后重发整批(已处理成功的事件会被跳过)。处理函数 panic 时 SDK 恢复并按失败处理。 - 服务端默认只等待 5 秒(地址可设 1 到 10 秒),超时算失败。处理函数要快:写一条记录、更新缓存;耗时的处理用异步方式。
- 需要推迟重发时,返回
callback.RetryAfter(d, err)包装的错误:响应带上Retry-After(取一批中最长的,不超过 1 小时),服务端按它推迟重发。 - 没有设置处理函数的事件类型直接算作成功。
callback.Event 的字段:
| 字段 | 说明 |
|---|---|
ID | 事件 ID:同一事件多次到达(重试、重新投递)时不变 |
Type | 事件类型,如 message.sent |
OccurredAt | 发生的时刻 |
Attempt | 第几次发送,从 1 开始(重新投递后重新计数) |
Test | 是否为控制台“测试发送”产生的事件 |
Data | data 的原始 JSON,与请求体中的逐字节相同 |
Raw | 这个事件在请求体中的原始 JSON |
每种事件的数据类型都嵌入 callback.EventBase:Origin(谁引起的:client、server、system、platform 等,在线状态变化为 nil)、Truncated(data 超过 70 KB 被截短,被去掉的字段为零值,需要时按 ID 调用 OpenAPI 读取)、Raw(data 的原始 JSON,SDK 没有解析的字段从中读取)。字段名由 JSON 字段名转换(sender_username 为 SenderUsername),可以为 null 的字段是指针,消息的 Body、Ext 等自由内容是 json.RawMessage。各字段的含义见事件。
ev.Decode() 按 Type 把 Data 解析为对应的类型(指针,如 *callback.MessageSent),不认识的类型返回 callback.ErrUnknownEventType;在 Any 或异步方式的工作者中按类型分派时使用。callback.EventTypes 是 SDK 认识的全部事件类型。
按事件 ID 去重
投递是“至少一次”:同一事件可能因重试、重新投递多次到达,事件 ID 不变(同一事件发给同一应用的多个地址时 ID 也相同)。SDK 在处理每个事件之前检查它是否已处理过:
- 已处理过的跳过,算作成功(返回
2xx,不能返回错误,否则服务端会一直重试); - 处理之前先占用这个事件 ID(
ClaimTTL内有效),处理成功后记为已完成,失败时释放; - 同一事件正在另一个请求中处理时(前一次处理较慢、服务端超时后又重发),这一批返回
503,服务端稍后重发。
NewMemoryDedupStore(retention) 在进程内存中保留已完成的事件 ID(retention 为 0 时 24 小时),最多 100 万个,超出时淘汰最早的并记警告日志。多个实例部署时,或需要在控制台重新投递几天前的事件时,换成共享存储并保留 7 天(失败的投递在控制台保留 7 天),见随机数和去重的存储。
顺序与截短
- 事件不保证顺序,同一批中的事件也不按发生先后排列。用版本号(
Version、MemberVersion、InfoVersion、消息的Seq等)和OccurredAt判断新旧,旧事件不能覆盖新状态。SDK 不缓存、不重排事件。 - 先看
Truncated:为true时对象、数组和较长的文本字段被去掉了。消息类事件再看ContentAvailable:为false时消息已撤回、擦除或过期,Body、Ext为nil。 - 在线状态变化、聊天室消息抄送和撤回是尽力而为的事件,可能丢失,不能作为对账依据。
异步处理
量大的消息抄送、耗时的处理,用异步方式:设置 Queue 而不是 Events,SDK 只做校验、去重和占用,把事件交给你的队列;入队成功即把这个事件记为已完成,全部入队后返回 2xx。
// 把事件写入数据库的待处理表(必须是持久的:入队后服务端不会再为这个事件重发)
type tableQueue struct{ db *sql.DB }
func (q tableQueue) Enqueue(ctx context.Context, ev *callback.Event) error {
_, err := q.db.ExecContext(ctx,
"INSERT INTO im_events (id, type, occurred_at, data) VALUES (?, ?, ?, ?) ON DUPLICATE KEY UPDATE id = id",
ev.ID, ev.Type, ev.OccurredAt, []byte(ev.Raw))
return err
}
func newAsyncHandler(db *sql.DB, secrets callback.SecretProvider) (*callback.Handler, error) {
return callback.NewHandler(callback.Config{
Secrets: secrets,
Nonces: callback.NewMemoryNonceStore(),
Dedup: callback.NewMemoryDedupStore(24 * time.Hour),
Queue: tableQueue{db: db},
})
}- 入队失败的事件释放占用,这一批返回
503,服务端稍后重发。 - 之后的处理和重试由你的队列负责(至少一次),进程内存中的通道不行。
- 一个处理器中
Events和Queue二选一。需要按类型分别处理的,在你的工作者中用ev.Decode()分派。
事件类型与处理函数
| 事件 | 处理函数的字段和数据类型 | 说明 |
|---|---|---|
user.created | UserCreated | 用户注册(客户端注册、服务端或控制台创建) |
user.status_changed | UserStatusChanged | 用户被封禁、解封或删除 |
user.profile_changed | UserProfileChanged | 用户修改了昵称、头像或自定义属性;不带变化的内容,需要时用 OpenAPI 读取 |
user.mute_changed | UserMuteChanged | 设置或解除全局禁言 |
user.session_created | UserSessionCreated | 用户在一台设备上登录 |
user.sessions_revoked | UserSessionsRevoked | 登录会话被吊销:退出、被踢、被挤下线、会话过期等 |
user.relogin_required | UserReloginRequired | 要求应用内全部用户重新登录 |
presence.changed | PresenceChanged | 用户的在线平台变化:上线、下线、多了或少了一个平台;尽力而为 |
friend.added | FriendAdded | 成为好友(双方各一个事件) |
friend.removed | FriendRemoved | 解除好友关系(双方各一个事件;对方账号被删除时 reason 为 user_deleted) |
friend.updated | FriendUpdated | 修改了好友的备注或自定义属性 |
friend_request.received | FriendRequestReceived | 收到需要同意的好友申请 |
friend_request.declined | FriendRequestDeclined | 好友申请被拒绝 |
blacklist.changed | BlacklistChanged | 拉黑或移出黑名单 |
group.created | GroupCreated | 建群 |
group.info_changed | GroupInfoChanged | 群资料、设置、全员禁言、群主变化,群被封禁或解封 |
group.dismissed | GroupDismissed | 群被解散 |
group.members_added | GroupMembersAdded | 成员加入(一次操作一个事件) |
group.members_removed | GroupMembersRemoved | 成员退出、被移出、被拉入群黑名单或账号被删除 |
group.members_updated | GroupMembersUpdated | 成员的角色、群昵称、成员属性、禁言变化 |
group.request_created | GroupRequestCreated | 新的入群申请或邀请 |
group.request_handled | GroupRequestHandled | 入群申请或邀请被同意、拒绝、撤回或自动失效 |
message.sent | MessageSent | 消息抄送:每写入一条单聊、群聊消息一个事件(群提示默认不发送) |
message.recalled | MessageRecalled | 消息被撤回(不受“不接收由租户服务端引起的事件”的影响) |
message.edited | MessageEdited | 消息被服务端编辑 |
chatroom.created | ChatroomCreated | 聊天室被创建 |
chatroom.dismissed | ChatroomDismissed | 聊天室被解散 |
chatroom.info_changed | ChatroomInfoChanged | 聊天室的资料、设置、所有者和管理员变化,被封禁或解封 |
chatroom.members_muted | ChatroomMembersMuted | 聊天室成员被禁言 |
chatroom.members_unmuted | ChatroomMembersUnmuted | 聊天室成员被解除禁言 |
chatroom.members_kicked | ChatroomMembersKicked | 聊天室成员被移出 |
chatroom.members_banned | ChatroomMembersBanned | 聊天室成员被封禁 |
chatroom.members_unbanned | ChatroomMembersUnbanned | 聊天室成员被解封 |
chatroom.allowlist_changed | ChatroomAllowlistChanged | 聊天室的白名单变化 |
chatroom.attributes_changed | ChatroomAttributesChanged | 聊天室的属性变化 |
chatroom.message_sent | ChatroomMessageSent | 聊天室的消息抄送:每条推送出去的消息一个事件;尽力而为 |
chatroom.message_recalled | ChatroomMessageRecalled | 聊天室的消息被撤回;尽力而为,不受“不接收由租户服务端引起的事件”的影响 |
rtc.call_created | RtcCallCreated | 发起通话 |
rtc.call_answered | RtcCallAnswered | 通话接通 |
rtc.member_joined | RtcMemberJoined | 成员加入群通话 |
rtc.member_left | RtcMemberLeft | 成员离开群通话 |
rtc.call_ended | RtcCallEnded | 通话结束,带时长 |
moderation.item_created | ModerationItemCreated | 有新的待审核记录(租户队列);不带被审核的内容,按 item_id 用 OpenAPI 读取 |
moderation.item_decided | ModerationItemDecided | 审核记录有了结论 |
moderation.report_created | ModerationReportCreated | 用户举报 |
moderation.violation_recorded | ModerationViolationRecorded | 记下一次违规(可能触发处罚) |
media.file_uploaded | MediaFileUploaded | 文件上传完成(企业认证材料除外) |
media.file_blocked | MediaFileBlocked | 文件被屏蔽 |
media.file_unblocked | MediaFileUnblocked | 文件被解除屏蔽 |
处理同步回调
同步回调在操作完成之前调用你的接收端,由你决定放行、拒绝或修改,见同步回调。callback.HookHandlers 中每种同步回调一个处理函数,返回 callback.Verdict:
| 处理函数 | 时机 | 可用的结果 |
|---|---|---|
MessageBeforeSend | 单聊、群聊消息发送和编辑之前 | Allow、Reject、Modify |
ChatroomBeforeSend | 聊天室消息发送之前 | Allow、Reject、Modify |
FriendBeforeAdd | 客户端发送好友申请、同意申请之前 | Allow、Reject |
GroupBeforeJoin | 客户端建群、邀请、申请、通过链接入群之前 | Allow、Reject;建群和邀请另有 Partial |
func hookHandlers() callback.HookHandlers {
return callback.HookHandlers{
ChatroomBeforeSend: func(ctx context.Context, req *callback.ChatroomBeforeSend) (callback.Verdict, error) {
var body struct {
Text string `json:"text"`
}
if err := json.Unmarshal(req.Body, &body); err != nil || body.Text == "" {
return callback.Allow(), nil
}
masked := strings.ReplaceAll(body.Text, "傻", "*")
if masked == body.Text {
return callback.Allow(), nil
}
// 替换消息内容;ext 不变
newBody, err := json.Marshal(map[string]string{"text": masked})
if err != nil {
return callback.Verdict{}, err
}
return callback.Modify(newBody, drim.Field[json.RawMessage]{}), nil
},
GroupBeforeJoin: func(ctx context.Context, req *callback.GroupBeforeJoin) (callback.Verdict, error) {
if req.JoinedVia != "invite" && req.JoinedVia != "create" {
return callback.Allow(), nil
}
// 邀请和建群可以只拒绝其中一部分人
var rejected []callback.RejectedTarget
for _, username := range req.Targets {
if strings.HasPrefix(username, "guest_") {
rejected = append(rejected, callback.RejectedTarget{Username: username, Reason: "guest_not_allowed", Message: "访客不能加入群聊"})
}
}
if len(rejected) == 0 {
return callback.Allow(), nil
}
return callback.Partial(rejected...), nil
},
}
}结果的构造函数:
| 函数 | 说明 |
|---|---|
callback.Allow() | 放行 |
callback.Reject(reason, message) | 拒绝。reason 为 1 到 64 个小写字母、数字或 _、.、-,客户端在错误的 details.app_reason 中得到它;message 为给用户看的提示,不超过 256 个字符、不能含控制字符(换行、制表符也是)。两者都可以为空 |
callback.Modify(body, ext) | 修改消息,只用于两种发消息前的回调。body 为 nil 表示不改,否则必须是 JSON 对象;ext 为 drim.Field[json.RawMessage]:零值不改,drim.Null[json.RawMessage]() 去掉 ext,drim.Set(v) 替换;至少改一个 |
callback.Partial(rejected...) | 部分拒绝,只用于 GroupBeforeJoin 的邀请和建群(JoinedVia 为 invite、create);至少一项,用户名必须在 req.Targets 中 |
Modify 的 ext 参数类型来自 drim 包,使用时导入 drim "github.com/deeprespond/im/sdk/server/go"。修改后的消息不再经过内容审核,由服务端按正常规则重新校验,不能改变类型。
时限与失败
- 处理函数收到的
ctx带截止时间:请求中的timeout_ms(服务端本次实际等待的时间,默认 1000 毫秒)减 50 毫秒。req.Timeout为服务端给出的等待时间。查询数据库、调用其他服务时使用这个ctx。 - 处理函数超时、返回错误、panic,或返回的结果不合规时,返回
HookFallback;没有设置HookFallback时返回503,由服务端按应用为这种回调设置的失败策略(放行或拒绝)处理。超时后 SDK 不再等待处理函数,它的结果被丢弃。 HookFallback为Allow()时,出错和超时都放行,不依赖服务端的失败策略;返回503则计入服务端的熔断统计,持续出错时服务端会暂时不再调用你的接收端。- 处理函数返回
Verdict的零值(忘了构造结果)不当作放行,按出错处理。 - 收到没有设置处理函数的同步回调时,同样返回
HookFallback(或503)并记错误日志。只在控制台开启已设置了处理函数的同步回调。 - SDK 在返回前检查结果:
reason、message的格式,Modify的内容,Partial的用户名,动作是否适用于这种回调,响应不超过 64 KB。不合规的不发出,按出错处理并记错误日志,便于发现自己的缺陷。 - 同时运行的处理函数超过
MaxConcurrentHooks(默认 100)时,新的请求直接按HookFallback(或503)返回。
同步回调直接影响用户发消息、加好友的速度,接收端慢或出错时影响所有用户。处理函数中只做本地的快速判断;接入时先在控制台把失败策略设为放行,稳定后再按需改为拒绝。
请求的类型(MessageBeforeSend 等)都嵌入 callback.HookBase:Test(控制台“测试发送”产生的请求)、Timeout、Raw(原始的 data)。各字段的含义见同步回调。
URL 验证
在控制台新建或修改回调地址时,服务端发送 URL 验证请求。处理器校验签名后自动返回 200 和 {"challenge":"<请求中的值>"},不调用你的代码,见 URL 验证。所以完成验证之前,接收端要先部署好并配置了这个地址的密钥。
接入的顺序:部署接收端(先只配置存储和密钥) → 在控制台新建地址,把显示的密钥配置到 Secrets → 用控制台的“测试发送”确认签名通过 → 完成 URL 验证 → 订阅事件 → 同步回调先以失败时放行开启。
随机数和去重的存储
NewMemoryNonceStore 和 NewMemoryDedupStore 只在一个进程内有效。多个实例、多台服务器部署时,同一事件或重放的请求可能到达另一个实例,要换成共享存储,实现两个接口:
// NonceStore 记录 10 分钟内见过的随机数:已存在返回 true,否则记下(ttl 后过期)并返回 false
type NonceStore interface {
Seen(ctx context.Context, nonce string, ttl time.Duration) (bool, error)
}
// DedupStore 按事件 ID 去重:开始处理时占用(得到占用的令牌),成功后改为已完成,失败时释放;
// 完成和释放只在令牌相同时生效(占用过期后被别的请求重新占用时,不改掉别人的占用)
type DedupStore interface {
Claim(ctx context.Context, eventID string, ttl time.Duration) (state callback.ClaimState, token string, err error)
Complete(ctx context.Context, eventID, token string) error
Release(ctx context.Context, eventID, token string) error
}Claim 的结果:callback.Acquired(占用成功,可以处理)、callback.Done(已处理过,跳过)、callback.Busy(正在另一个请求中处理)。
存储出错时这一批返回 503,服务端稍后重发:随机数存储出错时不能确认是否重放,宁可让服务端重发。内存随机数存储最多 100 万个,已满时返回 callback.ErrStoreFull(同样 503),不淘汰未到期的随机数。
用 go-redis 实现的写法(已完成的事件 ID 保留 7 天):
package callbackstore
import (
"context"
"crypto/rand"
"encoding/hex"
"time"
"github.com/redis/go-redis/v9"
"github.com/deeprespond/im/sdk/server/go/callback"
)
// RedisStores 同时实现 callback.NonceStore 和 callback.DedupStore。
type RedisStores struct {
R *redis.Client
Prefix string // 如 "im-callback:"
}
var (
_ callback.NonceStore = RedisStores{}
_ callback.DedupStore = RedisStores{}
)
func (s RedisStores) Seen(ctx context.Context, nonce string, ttl time.Duration) (bool, error) {
ok, err := s.R.SetNX(ctx, s.Prefix+"nonce:"+nonce, 1, ttl).Result()
return !ok, err
}
var claimScript = redis.NewScript(`
local v = redis.call('GET', KEYS[1])
if not v then redis.call('SET', KEYS[1], 'busy:' .. ARGV[1], 'PX', ARGV[2]) return 'acquired' end
if v == 'done' then return 'done' end
return 'busy'`)
var completeScript = redis.NewScript(`
if redis.call('GET', KEYS[1]) == 'busy:' .. ARGV[1] then
redis.call('SET', KEYS[1], 'done', 'PX', ARGV[2]) return 1
end
return 0`)
var releaseScript = redis.NewScript(`
if redis.call('GET', KEYS[1]) == 'busy:' .. ARGV[1] then return redis.call('DEL', KEYS[1]) end
return 0`)
func (s RedisStores) Claim(ctx context.Context, eventID string, ttl time.Duration) (callback.ClaimState, string, error) {
var b [16]byte
if _, err := rand.Read(b[:]); err != nil {
return 0, "", err
}
token := hex.EncodeToString(b[:])
r, err := claimScript.Run(ctx, s.R, []string{s.Prefix + "event:" + eventID}, token, ttl.Milliseconds()).Text()
switch {
case err != nil:
return 0, "", err
case r == "acquired":
return callback.Acquired, token, nil
case r == "done":
return callback.Done, "", nil
}
return callback.Busy, "", nil
}
func (s RedisStores) Complete(ctx context.Context, eventID, token string) error {
keep := 7 * 24 * time.Hour
return completeScript.Run(ctx, s.R, []string{s.Prefix + "event:" + eventID}, token, keep.Milliseconds()).Err()
}
func (s RedisStores) Release(ctx context.Context, eventID, token string) error {
return releaseScript.Run(ctx, s.R, []string{s.Prefix + "event:" + eventID}, token).Err()
}创建处理器时传入 Nonces: stores, Dedup: stores(stores := RedisStores{R: rdb, Prefix: "im-callback:"})。
不使用处理器
需要完全自己控制流程时,可以只用校验和解析的函数,随机数和事件去重由你自己完成:
| 函数 | 说明 |
|---|---|
callback.Verify(secrets, header, body, now) | 校验签名和时间戳,secrets 为这个地址的密钥列表;不通过时返回 *callback.VerifyError,Reason 为 bad_signature、bad_timestamp、stale_timestamp、bad_format |
callback.Parse(body) | 解析请求体为 *callback.Request(不校验签名,先调用 Verify) |
callback.Sign(secret, timestamp, nonce, body) | 计算签名(小写十六进制,不含 v1=) |
callback.SignatureHeader(secrets, timestamp, nonce, body) | X-IM-Signature 的值,每个密钥一个 v1= |
func verifyAndParse(r *http.Request, secrets []string) (*callback.Request, error) {
body, err := io.ReadAll(io.LimitReader(r.Body, 2<<20))
if err != nil {
return nil, err
}
if err := callback.Verify(secrets, r.Header, body, time.Now()); err != nil {
var verr *callback.VerifyError
if errors.As(err, &verr) {
log.Printf("签名校验失败:%s", verr.Reason)
}
return nil, err
}
// 校验通过后:记录 X-IM-Nonce(10 分钟内重复的拒绝)、核对 endpoint_id,再按 Kind 处理
return callback.Parse(body)
}callback.Request 的 Kind 为 event、hook 或 verify,Events 为事件列表,Hook、Data、TimeoutMs 为同步回调的内容,Challenge 为 URL 验证的随机串。在处理器的处理函数中,callback.RequestFrom(ctx) 返回这个请求的公共信息(RequestID、App、EndpointID、SentAt 等,不含事件和数据)。
本地测试
drimtest 包可以构造带正确签名的回调请求,不需要真实的服务端就能测试你的处理函数:
package receiver_test
import (
"context"
"net/http"
"net/http/httptest"
"strings"
"testing"
"time"
"github.com/deeprespond/im/sdk/server/go/callback"
"github.com/deeprespond/im/sdk/server/go/drimtest"
)
const secret = "cbs_test_secret_for_unit_tests_0000000000000000"
func newHandler(t *testing.T, events callback.EventHandlers, hooks callback.HookHandlers) *callback.Handler {
t.Helper()
h, err := callback.NewHandler(callback.Config{
Secrets: callback.StaticSecrets{drimtest.DefaultEndpointID: {secret}},
Nonces: callback.NewMemoryNonceStore(),
Dedup: callback.NewMemoryDedupStore(time.Hour),
Events: events,
Hooks: hooks,
})
if err != nil {
t.Fatal(err)
}
return h
}
func TestUserCreatedIsHandledOnce(t *testing.T) {
created := 0
h := newHandler(t, callback.EventHandlers{
UserCreated: func(ctx context.Context, ev *callback.Event, d *callback.UserCreated) error {
created++
return nil
},
}, callback.HookHandlers{})
body := drimtest.EventRequest(drimtest.Event{
ID: "1840300000000000001",
Type: "user.created",
Data: map[string]any{"origin": "server", "username": "zhangsan", "created_via": "openapi"},
})
for i := 0; i < 2; i++ { // 同一事件到达两次:只处理一次,两次都返回 204
rec := httptest.NewRecorder()
h.ServeHTTP(rec, drimtest.SignedRequest(drimtest.DefaultEndpointID, body, secret))
if rec.Code != http.StatusNoContent {
t.Fatalf("status %d", rec.Code)
}
}
if created != 1 {
t.Fatalf("handled %d times", created)
}
}
func TestRejectBannedWord(t *testing.T) {
h := newHandler(t, callback.EventHandlers{}, callback.HookHandlers{
MessageBeforeSend: func(ctx context.Context, req *callback.MessageBeforeSend) (callback.Verdict, error) {
if strings.Contains(string(req.Body), "加微信") {
return callback.Reject("forbidden_word", "消息中有不允许的内容"), nil
}
return callback.Allow(), nil
},
})
body := drimtest.HookRequest("message.before_send", map[string]any{
"via": "client", "client_msg_id": "c1", "conversation_type": "single",
"sender_username": "zhangsan", "recipient_username": "lisi",
"type": "text", "body": map[string]any{"text": "加微信 123"},
})
rec := httptest.NewRecorder()
h.ServeHTTP(rec, drimtest.SignedRequest(drimtest.DefaultEndpointID, body, secret))
if rec.Code != http.StatusOK || !strings.Contains(rec.Body.String(), `"action":"reject"`) {
t.Fatalf("status %d, body %s", rec.Code, rec.Body)
}
}
func TestWrongSecretIsRejected(t *testing.T) {
h := newHandler(t, callback.EventHandlers{}, callback.HookHandlers{})
body := drimtest.VerifyRequest("challenge-1")
rec := httptest.NewRecorder()
h.ServeHTTP(rec, drimtest.SignedRequest(drimtest.DefaultEndpointID, body, "cbs_another_secret"))
if rec.Code != http.StatusUnauthorized {
t.Fatalf("status %d", rec.Code)
}
}| 函数 | 说明 |
|---|---|
drimtest.EventRequest(events...) | 事件回调的请求体,地址 ID 为 drimtest.DefaultEndpointID;drimtest.Event 的 ID 不给时随机生成,Data 为 string、[]byte、json.RawMessage 时原样使用,其他值编码为 JSON |
drimtest.EventRequestFor(endpointID, events...) | 指定地址 ID 的事件回调请求体 |
drimtest.HookRequest(hook, data)、HookRequestFor(endpointID, hook, timeoutMs, data) | 同步回调的请求体,默认 timeout_ms 为 1000 |
drimtest.VerifyRequest(challenge) | URL 验证的请求体 |
drimtest.SignedRequest(endpointID, body, secrets...) | 带正确签名的 *http.Request(当前时间、新的随机数);轮换期间可以给新旧两个密钥 |
drimtest.SignedHeader(endpointID, body, at, nonce, secrets...) | 指定时间戳和随机数的请求头,用于测试过期和重放 |
drimtest.SignatureVectors() | 签名的测试向量 |
签名测试向量
drimtest.SignatureVectors() 返回一组签名的测试向量,与服务端和签名与安全中的示例使用同一套数据:有效的签名、轮换期间的两个签名、多字节字符和不转义的 <>&、时间戳带前导零或加号、过期、请求体改了一个字节等。自己实现签名校验(或在网关上校验)时,可以用它确认实现正确:
func TestSignatureVectors(t *testing.T) {
for _, v := range drimtest.SignatureVectors() {
err := callback.Verify([]string{v.Secret}, v.Header(), v.Body(), time.Unix(v.Now, 0))
var verr *callback.VerifyError
switch {
case v.Valid && err != nil:
t.Errorf("%s: want valid, got %v", v.Name, err)
case !v.Valid && (!errors.As(err, &verr) || verr.Reason != v.Error):
t.Errorf("%s: want %s, got %v", v.Name, v.Error, err)
}
}
}每个向量有 Name、Secret(接收端配置的密钥)、Timestamp、Nonce、Body()(原始请求体)、Header()(请求头)、Now(校验时的当前时间)、Valid 和 Error(不通过时期望的原因)。
接入真实的服务端之后,用控制台回调地址的“测试发送”确认签名能通过:签名失败时 SDK 的错误日志写明了可能的原因,常见的是密钥配置错误,或网关、框架改动了请求体,见调试签名。
独立频道回调
| 事件 | 强类型模型 |
|---|---|
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 |
这九种事件可通过现有分派器处理,沿用签名验证、去重和错误处理。字段见频道事件;部分字段可缺省,不应假设非会话事件包含目标用户。频道与原通话的事件名和模型分开,配置内部广播不作为回调。
