openevent

package module
v1.0.1 Latest Latest
Warning

This package is not in the latest version of its module.

Go to latest
Published: Feb 3, 2026 License: MIT Imports: 3 Imported by: 0

README

open-event-sdk-go

开放平台事件订阅 SDK(Go 语言版),支持通过 WebSocket 长连接接收和处理事件。

特性

  • WebSocket 长连接:无需公网 IP,延迟更低、实时性更好
  • 自动重连:网络断开时自动重连,支持指数退避策略
  • KSO-1 签名认证:安全的认证机制
  • 灵活的事件处理:支持单一 Handler 和 Dispatcher 分发两种模式
  • 开箱即用:内置默认配置,无需额外设置即可使用

安装

go get github.com/GongchuangSu/open-event-sdk-go

快速开始

package main

import (
    "context"
    "log"

    openevent "github.com/GongchuangSu/open-event-sdk-go"
)

func main() {
    client := openevent.NewClient("your_app_id", "your_app_secret",
        openevent.WithEventHandlerFunc(func(ctx context.Context, e *openevent.Event) error {
            log.Printf("收到事件: event_code=%s", e.EventCode())
            log.Printf("事件数据: %s", e.Data)
            return nil
        }),
    )

    if err := client.Start(context.Background()); err != nil {
        log.Fatal(err)
    }
}
Dispatcher 分发模式

按事件编码(event_code)分别处理不同事件:

package main

import (
    "context"
    "log"

    openevent "github.com/GongchuangSu/open-event-sdk-go"
)

func main() {
    // 创建分发器
    dispatcher := openevent.NewDispatcher()

    // 注册不同事件编码的处理器
    // 事件编码 = topic.operation,如 "kso.app_chat.message.create"
    dispatcher.RegisterFunc("kso.app_chat.message.create", func(ctx context.Context, e *openevent.Event) error {
        log.Printf("处理聊天消息事件: %s", e.Data)
        return nil
    })

    dispatcher.RegisterFunc("kso.user.status.update", func(ctx context.Context, e *openevent.Event) error {
        log.Printf("处理用户状态变更事件: %s", e.Data)
        return nil
    })

    // 注册兜底处理器
    dispatcher.RegisterFallbackFunc(func(ctx context.Context, e *openevent.Event) error {
        log.Printf("未知事件: event_code=%s", e.EventCode())
        return nil
    })

    // 创建客户端
    client := openevent.NewClient("your_app_id", "your_app_secret",
        openevent.WithDispatcher(dispatcher),
    )

    if err := client.Start(context.Background()); err != nil {
        log.Fatal(err)
    }
}
类型化事件处理

目前部分事件已支持 OnV7XXX 方法,可使用链式调用注册类型化处理器,事件数据会自动解析为对应的结构体。其他事件请使用 RegisterFunc 方法处理。

已支持 OnV7XXX 方法的事件:

方法 事件编码 说明
OnV7AppChatMessageCreate kso.app_chat.message.create 用户给应用发送消息
OnV7AppChatCreate kso.app_chat.create 首次创建用户和机器人的会话
OnV7AppGroupChatDelete kso.xz.app.group_chat.delete 群聊解散
OnV7AppGroupChatMemberUserCreate kso.xz.app.group_chat.member.user.create 用户进群
OnV7AppGroupChatMemberUserDelete kso.xz.app.group_chat.member.user.delete 用户退群
OnV7AppGroupChatMemberRobotCreate kso.xz.app.group_chat.member.robot.create 机器人进群
OnV7AppGroupChatMemberRobotDelete kso.xz.app.group_chat.member.robot.delete 机器人退群

组合使用示例:

package main

import (
    "context"
    "encoding/json"
    "log"

    openevent "github.com/GongchuangSu/open-event-sdk-go"
)

func main() {
    dispatcher := openevent.NewDispatcher()

    // ========== 方式一: OnV7XXX 方法(类型安全,推荐) ==========
    dispatcher.
        OnV7AppChatMessageCreate(func(ctx context.Context, e *openevent.V7AppChatMessageCreateEvent) error {
            log.Printf("收到消息: chat_id=%s, sender=%s", e.Data.Chat.Id, e.Data.Sender.Id)
            return nil
        }).
        OnV7AppChatCreate(func(ctx context.Context, e *openevent.V7AppChatCreateEvent) error {
            log.Printf("会话创建: chat_id=%s", e.Data.ChatId)
            return nil
        }).
        OnV7AppGroupChatMemberUserCreate(func(ctx context.Context, e *openevent.V7AppGroupChatMemberUserCreateEvent) error {
            log.Printf("用户进群: chat_id=%s", e.Data.ChatId)
            return nil
        })

    // ========== 方式二: RegisterFunc(处理其他事件,需自行解析 Data) ==========
    dispatcher.RegisterFunc("kso.user.status.update", func(ctx context.Context, e *openevent.Event) error {
        log.Printf("用户状态变更: %s", e.EventCode())
        var data map[string]any
        json.Unmarshal([]byte(e.Data), &data)
        log.Printf("数据: %+v", data)
        return nil
    })

    // ========== 兜底处理器 ==========
    dispatcher.RegisterFallbackFunc(func(ctx context.Context, e *openevent.Event) error {
        log.Printf("未处理的事件: %s", e.EventCode())
        return nil
    })

    client := openevent.NewClient("your_app_id", "your_app_secret",
        openevent.WithDispatcher(dispatcher),
    )

    if err := client.Start(context.Background()); err != nil {
        log.Fatal(err)
    }
}

使用事件数据模型:

如需直接使用事件数据模型,请导入 model 包:

import "github.com/GongchuangSu/open-event-sdk-go/event/model"

var sender model.V7Identity
var data model.V7NotificationAppChatMessageCreateData

配置选项

基础配置
client := openevent.NewClient(appId, appSecret,
    // 自定义 WebSocket 端点(可选,默认 wss://openapi.wps.cn/v7/event/ws)
    openevent.WithEndpoint("wss://custom-endpoint.com/event/ws"),

    // 设置日志级别
    openevent.WithLogLevel(openevent.LogLevelDebug),

    // 使用自定义日志
    openevent.WithLogger(customLogger),
)
重连配置(指数退避策略)

SDK 采用指数退避(Exponential Backoff)策略进行重连,避免在网络恢复时产生惊群效应。

重连间隔计算公式interval = min(baseInterval * multiplier^(retryCount-1), maxInterval) * (1 ± jitter)

client := openevent.NewClient(appId, appSecret,
    // 开启/关闭自动重连(默认开启)
    openevent.WithAutoReconnect(true),

    // 重连基础间隔(默认 1 秒)
    openevent.WithReconnectBaseInterval(1 * time.Second),

    // 重连最大间隔(默认 60 秒)
    openevent.WithReconnectMaxInterval(60 * time.Second),

    // 重连间隔倍数(默认 2.0)
    openevent.WithReconnectMultiplier(2.0),

    // 最大重试次数(-1 表示无限重试,默认 -1)
    openevent.WithReconnectMaxRetry(10),

    // 重连抖动系数(默认 0.2,表示 ±20% 随机抖动)
    openevent.WithReconnectJitter(0.2),
)

默认重连时间序列示例(baseInterval=1s, multiplier=2, maxInterval=60s, jitter=0.2):

重试次数 基础间隔 实际间隔范围
1 1s 0.8s ~ 1.2s
2 2s 1.6s ~ 2.4s
3 4s 3.2s ~ 4.8s
4 8s 6.4s ~ 9.6s
5 16s 12.8s ~ 19.2s
6 32s 25.6s ~ 38.4s
7+ 60s 48s ~ 72s
超时配置
client := openevent.NewClient(appId, appSecret,
    // 写操作超时(默认 10 秒)
    openevent.WithWriteWait(10 * time.Second),

    // Pong 等待超时(默认 90 秒)
    openevent.WithPongWait(90 * time.Second),
)

事件结构

原始事件消息(加密)

SDK 接收到的原始消息包含以下字段:

字段 类型 说明
topic string 消息主题(根据不同事件而定)
operation string 消息变更动作(根据不同事件而定)
time int64 时间(秒为单位的时间戳)
nonce string iv 向量(解密时使用)
signature string 消息签名
encrypted_data string 加密的消息数据
解密后的事件结构

SDK 会自动验证签名并解密数据,处理器接收到的是解密后的事件:

type Event struct {
    Topic     string `json:"topic"`      // 消息主题
    Operation string `json:"operation"`  // 变更动作
    Time      int64  `json:"time"`       // 时间戳(秒)
    Data      string `json:"data"`       // 解密后的事件数据(JSON 字符串)
}

// 获取事件编码
func (e *Event) EventCode() string

事件编码(event_code)说明

  • 事件编码 = topic + . + operation,全局唯一
  • 例如:topic="kso.app_chat.message", operation="create"event_code="kso.app_chat.message.create"
  • 通过 e.EventCode() 方法动态获取事件编码
  • Dispatcher 按事件编码进行事件分发
签名验证

签名计算方式:

  1. 构建签名原文:content = access_key:topic:nonce:time:encrypted_data
  2. 计算签名:signature = HMAC-SHA256(content, secret_key)
  3. 签名使用 URL 安全的无填充 base64 编码
数据解密

解密方式:

  1. encrypted_data 使用标准的有填充 base64 编码
  2. 密钥 cipher = MD5(secret_key)
  3. 使用 AES-CBC 模式解密,iv 为 nonce 的前 16 字节
  4. 解密后移除 PKCS7 填充

事件处理

处理成功
handler := openevent.HandlerFunc(func(ctx context.Context, e *openevent.Event) error {
    log.Printf("处理事件: event_code=%s", e.EventCode())
    
    // 解析事件数据
    var data map[string]interface{}
    if err := json.Unmarshal([]byte(e.Data), &data); err != nil {
        return err
    }
    
    // 处理业务逻辑...
    return nil
})
处理失败

返回 error 表示处理失败:

handler := openevent.HandlerFunc(func(ctx context.Context, e *openevent.Event) error {
    if err := processEvent(e); err != nil {
        return err // 处理失败
    }
    return nil
})

优雅关闭

使用 context 控制连接生命周期:

ctx, cancel := context.WithCancel(context.Background())

// 监听退出信号
go func() {
    sigChan := make(chan os.Signal, 1)
    signal.Notify(sigChan, syscall.SIGINT, syscall.SIGTERM)
    <-sigChan

    cancel()       // 取消 context
    client.Stop()  // 停止客户端
}()

client.Start(ctx)

自定义日志

实现 Logger 接口:

type Logger interface {
    Debug(ctx context.Context, args ...interface{})
    Info(ctx context.Context, args ...interface{})
    Warn(ctx context.Context, args ...interface{})
    Error(ctx context.Context, args ...interface{})
}

示例:

type MyLogger struct{}

func (l *MyLogger) Debug(ctx context.Context, args ...interface{}) {
    // 自定义实现
}

// ... 其他方法

client := openevent.NewClient(appId, appSecret,
    openevent.WithLogger(&MyLogger{}),
)

目录结构

open-event-sdk-go/
├── openevent.go            # 根包入口(推荐使用)
├── ws/                     # WebSocket 客户端
│   ├── client.go           # 客户端主逻辑
│   ├── option.go           # 配置选项
│   └── error.go            # 错误定义
├── event/                  # 事件处理
│   ├── dispatcher.go       # 事件分发器
│   ├── dispatcher_typed.go # OnXXX 类型化处理方法
│   ├── typed_event.go      # 类型化事件定义
│   ├── handler.go          # Handler 接口
│   ├── event.go            # 事件实体
│   └── model/              # 事件数据模型
│       ├── common.go       # 通用类型(Identity 等)
│       ├── im.go           # IM 相关类型
│       ├── im_event.go     # IM 事件数据结构
│       └── event_code.go   # 事件编码常量
├── core/                   # 核心公共组件
│   └── logger.go           # 日志接口
├── internal/               # 内部实现(不对外暴露)
│   ├── kso/                # KSO-1 签名和加解密
│   └── protocol/           # WebSocket 协议定义
└── examples/               # 使用示例
    ├── simple/             # 简单示例
    └── dispatcher/         # Dispatcher 模式示例

协议说明

消息类型

服务端通过 WebSocket 向客户端推送两种类型的消息:

1. 事件消息

携带加密的事件数据,SDK 会自动验签和解密。

{
    "topic": "kso.app_chat.message",
    "operation": "create",
    "time": 1704067200,
    "nonce": "b38c06ba3330d2a3",
    "signature": "EzpccL5eAnDbOH2qtZK3fHNBaO0UV3xvYvhbWLp1wuQ",
    "encrypted_data": "oadgi2+nWGZal2EfCSlxbLrr2Aog..."
}
2. 关闭通知(goaway)

服务端主动关闭连接时发送,包含 type: "goaway" 和关闭原因:

{
    "type": "goaway",
    "reason": "server_shutdown",
    "message": "服务器维护中",
    "reconnect_ms": 5000
}

GoAway 原因类型

原因 说明 是否重连
server_shutdown 服务器关闭(如维护升级) 是,按 reconnect_ms 延迟重连
connection_replaced 连接被新连接替换(同一应用重复连接)
heartbeat_timeout 心跳超时
心跳机制
  • 服务端每 30 秒发送 WebSocket Ping
  • 客户端自动回复 Pong
  • 90 秒内未收到 Pong 则断开连接
认证方式

使用 KSO-1 签名认证,需要在 WebSocket 握手时携带以下 HTTP 头:

  • X-Kso-Date: 请求时间(RFC1123 格式)
  • X-Kso-Authorization: 签名(格式:KSO-1 {app_id}:{signature}

许可证

MIT License

Documentation

Overview

Package openevent 提供事件订阅 SDK

该 SDK 用于接收和处理来自开放平台的事件推送。

快速开始:

client := openevent.NewClient("app_id", "app_secret",
    openevent.WithEventHandlerFunc(func(ctx context.Context, e *openevent.Event) error {
        log.Printf("收到事件: %s", e.EventCode())
        return nil
    }),
)
client.Start(context.Background())

使用 Dispatcher 分发事件:

dispatcher := openevent.NewDispatcher().
    // 已支持的事件使用 OnV7XXX 方法(类型安全)
    OnV7AppChatMessageCreate(func(ctx context.Context, e *openevent.V7AppChatMessageCreateEvent) error {
        log.Printf("收到消息: %s", e.Data.Message.Id)
        return nil
    })

// 其他事件使用 RegisterFunc 方法
dispatcher.RegisterFunc("kso.other.event", func(ctx context.Context, e *openevent.Event) error {
    log.Printf("其他事件: %s", e.Data)
    return nil
})

client := openevent.NewClient("app_id", "app_secret",
    openevent.WithDispatcher(dispatcher),
)

Index

Constants

View Source
const (
	// LogLevelDebug 调试级别
	LogLevelDebug = core.LogLevelDebug

	// LogLevelInfo 信息级别
	LogLevelInfo = core.LogLevelInfo

	// LogLevelWarn 警告级别
	LogLevelWarn = core.LogLevelWarn

	// LogLevelError 错误级别
	LogLevelError = core.LogLevelError
)

Variables

View Source
var (
	// ErrHandlerNotSet 事件处理器未设置
	ErrHandlerNotSet = ws.ErrHandlerNotSet

	// ErrClientClosed 客户端已关闭
	ErrClientClosed = ws.ErrClientClosed

	// ErrReconnectExceeded 超过最大重连次数
	ErrReconnectExceeded = ws.ErrReconnectExceeded
)
View Source
var NewClient = ws.NewClient

NewClient 创建 WebSocket 客户端

参数:

  • appId: 应用 ID
  • appSecret: 应用密钥
  • opts: 可选配置项
View Source
var NewDefaultLogger = core.NewDefaultLogger

NewDefaultLogger 创建默认日志实例

View Source
var NewDispatcher = event.NewDispatcher

NewDispatcher 创建事件分发器

View Source
var NewNopLogger = core.NewNopLogger

NewNopLogger 创建空日志实例(不输出任何日志)

View Source
var WithAckMode = ws.WithAckMode

WithAckMode 设置是否启用 ACK 模式(默认: true) 启用后,事件处理结果会发送给服务端,处理失败时服务端会触发重试

View Source
var WithAutoReconnect = ws.WithAutoReconnect

WithAutoReconnect 设置是否开启自动重连(默认: true)

View Source
var WithDispatcher = ws.WithDispatcher

WithDispatcher 设置事件分发器

View Source
var WithEndpoint = ws.WithEndpoint

WithEndpoint 设置 WebSocket 连接端点 默认使用 SDK 内置端点,一般无需设置

View Source
var WithEventHandler = ws.WithEventHandler

WithEventHandler 设置单一事件处理器

View Source
var WithEventHandlerFunc = ws.WithEventHandlerFunc

WithEventHandlerFunc 设置函数类型的事件处理器

View Source
var WithLogLevel = ws.WithLogLevel

WithLogLevel 设置日志级别

View Source
var WithLogger = ws.WithLogger

WithLogger 设置自定义日志实例

View Source
var WithPongWait = ws.WithPongWait

WithPongWait 设置等待 Pong 响应超时时间(默认: 90秒)

View Source
var WithReconnectBaseInterval = ws.WithReconnectBaseInterval

WithReconnectBaseInterval 设置重连基础间隔(默认: 1秒)

View Source
var WithReconnectJitter = ws.WithReconnectJitter

WithReconnectJitter 设置重连抖动系数(默认: 0.2)

View Source
var WithReconnectMaxInterval = ws.WithReconnectMaxInterval

WithReconnectMaxInterval 设置重连最大间隔(默认: 60秒)

View Source
var WithReconnectMaxRetry = ws.WithReconnectMaxRetry

WithReconnectMaxRetry 设置最大重试次数(默认: -1 无限重试)

View Source
var WithReconnectMultiplier = ws.WithReconnectMultiplier

WithReconnectMultiplier 设置重连间隔倍数(默认: 2.0)

View Source
var WithWriteWait = ws.WithWriteWait

WithWriteWait 设置写操作超时时间(默认: 10秒)

Functions

This section is empty.

Types

type Client

type Client = ws.Client

Client WebSocket 长连接客户端

type ClientError

type ClientError = ws.ClientError

ClientError 客户端错误

type Dispatcher

type Dispatcher = event.Dispatcher

Dispatcher 事件分发器

type Event

type Event = event.Event

Event 事件实体

type Handler

type Handler = event.Handler

Handler 事件处理器接口

type HandlerFunc

type HandlerFunc = event.HandlerFunc

HandlerFunc 函数类型的事件处理器

type LogLevel

type LogLevel = core.LogLevel

LogLevel 日志级别

type Logger

type Logger = core.Logger

Logger 日志接口

type Option

type Option = ws.Option

Option 客户端配置选项

type ServerError

type ServerError = ws.ServerError

ServerError 服务端错误

type V7AppChatCreateEvent added in v1.0.1

type V7AppChatCreateEvent = event.V7AppChatCreateEvent

V7AppChatCreateEvent 应用会话创建事件

type V7AppChatMessageCreateEvent added in v1.0.1

type V7AppChatMessageCreateEvent = event.V7AppChatMessageCreateEvent

V7AppChatMessageCreateEvent 应用收到消息事件

type V7AppGroupChatDeleteEvent added in v1.0.1

type V7AppGroupChatDeleteEvent = event.V7AppGroupChatDeleteEvent

V7AppGroupChatDeleteEvent 群聊解散事件

type V7AppGroupChatMemberRobotCreateEvent added in v1.0.1

type V7AppGroupChatMemberRobotCreateEvent = event.V7AppGroupChatMemberRobotCreateEvent

V7AppGroupChatMemberRobotCreateEvent 机器人进群事件

type V7AppGroupChatMemberRobotDeleteEvent added in v1.0.1

type V7AppGroupChatMemberRobotDeleteEvent = event.V7AppGroupChatMemberRobotDeleteEvent

V7AppGroupChatMemberRobotDeleteEvent 机器人退群事件

type V7AppGroupChatMemberUserCreateEvent added in v1.0.1

type V7AppGroupChatMemberUserCreateEvent = event.V7AppGroupChatMemberUserCreateEvent

V7AppGroupChatMemberUserCreateEvent 用户进群事件

type V7AppGroupChatMemberUserDeleteEvent added in v1.0.1

type V7AppGroupChatMemberUserDeleteEvent = event.V7AppGroupChatMemberUserDeleteEvent

V7AppGroupChatMemberUserDeleteEvent 用户退群事件

Directories

Path Synopsis
Package core 提供 SDK 的核心公共组件
Package core 提供 SDK 的核心公共组件
Package event 提供事件处理相关的类型和接口
Package event 提供事件处理相关的类型和接口
model
Package model 定义事件数据模型 该包中的类型定义可通过代码生成器自动生成 类型命名与 xidl 定义保持一致,使用 V7 前缀
Package model 定义事件数据模型 该包中的类型定义可通过代码生成器自动生成 类型命名与 xidl 定义保持一致,使用 V7 前缀
examples
dispatcher command
Package main 演示使用 Dispatcher 处理事件
Package main 演示使用 Dispatcher 处理事件
simple command
Package main 演示使用单一 Handler 接收和处理事件
Package main 演示使用单一 Handler 接收和处理事件
internal
kso
Package kso 提供 KSO-1 签名相关工具(内部使用)
Package kso 提供 KSO-1 签名相关工具(内部使用)
protocol
Package protocol 定义 WebSocket 协议常量(内部使用)
Package protocol 定义 WebSocket 协议常量(内部使用)
Package ws 提供 WebSocket 长连接客户端
Package ws 提供 WebSocket 长连接客户端

Jump to

Keyboard shortcuts

? : This menu
/ : Search site
f or F : Jump to
y or Y : Canonical URL