MSG Chain IBC 自定义中间件开发指南
数据来源:MSG Chain 代码库核实
主网状态: No-Go — 当前 MSGChain 主网裁决为 No-Go,以下内容反映代码实际状态,不代表生产可用。
目录
1. 概述
1.1 什么是 IBC 中间件
IBC 中间件是一种位于 IBC 核心层与 IBC 应用层之间的模块化组件,它可以在不修改底层 IBC 协议的前提下,对数据包流进行拦截、检查和增强。中间件模式源于 Cosmos IBC v2+ 的设计理念——将基础传输功能与策略逻辑解耦。
在 MSG Chain 上,IBC 中间件通过包装(wrapping)标准 IBC 模块接口(IBCModule)来实现。每个中间件持有对下一个模块(可能是另一个中间件或基础应用)的引用,形成责任链模式。
1.2 ICS-20 中间件栈
标准 ICS-20 代币转账模块的中间件栈结构:
┌─────────────────────────────────────┐
│ IBC 核心 (Core) │
│ ┌───────────────────────────────┐ │
│ │ ICS-29 IBC Fee Middleware │ │
│ │ (可选,Relayer 费用) │ │
│ └───────────────────────────────┘ │
│ ┌───────────────────────────────┐ │
│ │ 自定义中间件 (你在此处) │ │
│ └───────────────────────────────┘ │
│ ┌───────────────────────────────┐ │
│ │ ICS-20 Transfer Base App │ │
│ │ (代币转账核心逻辑) │ │
│ └───────────────────────────────┘ │
└─────────────────────────────────────┘
1.3 自定义中间件用例
在 MSG Chain 上,自定义 IBC 中间件可以解决以下关键问题:
| 用例 | 说明 | 适用场景 |
|---|---|---|
| 速率限制 (Rate Limiting) | 限制单通道/单时间的跨链资金流出量 | 防止桥攻击导致资金一次性流失 |
| 费用吸收 (Fee Absorption) | 自动扣除跨链转账手续费 | 代币转账时自动收取协议费 |
| 数据包过滤 (Packet Filtering) | 按规则过滤 IBC 数据包 | 阻止可疑交易、合规过滤器 |
| 治理控制 (Governance Control) | 通过治理提案控制 IBC 通道 | 暂停/恢复通道、更新参数 |
| 元数据注入 (Metadata Injection) | 在 IBC 数据包中附加元数据 | 跨链追踪、数据标记 |
| 白名单/黑名单 | 按地址/代币类型过滤 | 合规审查、安全控制 |
1.4 MSG Chain IBC 中间件栈设计
MSG Chain 规划中的完整 IBC 中间件栈:
// X-MSG-Stub=true: MSG Chain IBC 中间件栈编排
func RegisterIBCMiddleware(app *MsgChainApp) {
transferStack := ibcfee.NewIBCMiddleware(
msgratelimit.NewIBCMiddleware(
msggov.NewIBCMiddleware(
transferIBCModule, // ICS-20 基础应用
),
),
)
// 绑定到 transfer 端口
ibcRouter.AddRoute("transfer", transferStack)
}
1.5 关键概念
| 概念 | 解释 |
|---|---|
IBCModule |
IBC 模块的标准接口,所有中间件和应用都必须实现 |
next IBCModule |
下一个链中的模块引用(另一个中间件或基础应用) |
Callbacks |
数据包生命周期的钩子函数:Send, Recv, Ack, Timeout |
Port Binding |
端口与模块的绑定关系(如 transfer 端口绑定转账模块) |
ICS-29 |
IBC Fee 中间件标准,Relayer 激励 |
2. IBC中间件架构
2.1 基础应用 vs 中间件
IBC 模块分为两类:
基础应用 (Base Application):实现 IBC 业务逻辑的终端模块。例如 ICS-20 转账模块负责代币的锁定、铸造和解锁。
中间件 (Middleware):围绕基础应用的包装器,在原始业务逻辑前后插入自定义逻辑。中间件不改变业务核心逻辑,只增加约束或增强功能。
关系对比:
| 特性 | 基础应用 | 中间件 |
|---|---|---|
| 最终处理数据包 | 是 | 否,委托给下一个模块 |
| 必须有 keeper | 是 | 通常是 |
| 可跳过 | 否(必须有一个) | 是 |
| 数量 | 每个端口一个 | 零或多个 |
2.2 中间件栈结构
中间件栈遵循洋葱模型——数据包从外层进入,逐层向内传递,结果再从内层返回:
SendPacket 方向:
┌─────────────────────────────────────────────────────────┐
│ IBC Core → Middleware1 → Middleware2 → Base App │
│ (Send) (pre-check) (pre-check) (execute) │
└─────────────────────────────────────────────────────────┘
RecvPacket 方向:
┌─────────────────────────────────────────────────────────┐
│ IBC Core → Middleware1 → Middleware2 → Base App │
│ (Recv) (pre-check) (pre-check) (execute) │
└─────────────────────────────────────────────────────────┘
Ack/Timeout 方向:
┌─────────────────────────────────────────────────────────┐
│ Base App → Middleware2 → Middleware1 → IBC Core │
│ (result) (post-process) (post-process) │
└─────────────────────────────────────────────────────────┘
2.3 多层中间件注册
在 app.go 中注册多层中间件的标准模式:
// X-MSG-Stub=true: 多层中间件注册
func RegisterIBCMiddleware(
app *MsgChainApp,
transferIBCModule ibcporttypes.IBCModule,
) ibcporttypes.IBCModule {
// 从内到外包装
var stack ibcporttypes.IBCModule
stack = transferIBCModule
// 第1层: IBC Fee 中间件
stack = ibcfee.NewIBCMiddleware(stack, app.IBCFeeKeeper)
// 第2层: 速率限制中间件
stack = msgratelimit.NewIBCMiddleware(stack, app.RateLimitKeeper, app.ChannelKeeper)
// 第3层: 治理控制中间件
stack = msggov.NewIBCMiddleware(stack, app.GovKeeper, app.ChannelKeeper)
return stack
}
2.4 数据包生命周期
IBC 数据包的完整生命周期包含 4 个核心阶段,中间件在每个阶段都有对应的回调:
发送链 (msg-chain-1) 接收链 (cosmoshub-4)
│ │
│ ┌── SendPacket ──────────────────► │
│ │ • middleware.OnSendPacket │
│ │ • base_app.SendPacket │
│ │ │
│ │ ┌── RecvPacket ◄────────────── │
│ │ │ • middleware.OnRecvPacket │
│ │ │ • base_app.OnRecvPacket │
│ │ │ │
│ ◄──┴── AcknowledgementPacket │
│ │ • middleware.OnAcknowledgement │
│ │ • base_app.OnAcknowledgement │
│ │ │
│ ◄── TimeoutPacket │
│ • middleware.OnTimeoutPacket │
│ • base_app.OnTimeoutPacket │
2.5 IBCModule 接口
所有中间件必须实现的 IBCModule 接口:
// X-MSG-Stub=true: IBCModule 接口定义 (来自 ibc-go)
type IBCModule interface {
OnChanOpenInit(
ctx sdk.Context,
order channeltypes.Order,
connectionHops []string,
portID string,
channelID string,
channelCap *capabilitytypes.Capability,
counterparty channeltypes.Counterparty,
version string,
) (string, error)
OnChanOpenTry(
ctx sdk.Context,
order channeltypes.Order,
connectionHops []string,
portID string,
channelID string,
channelCap *capabilitytypes.Capability,
counterparty channeltypes.Counterparty,
version string,
) (string, error)
OnChanOpenAck(
ctx sdk.Context,
portID string,
channelID string,
counterpartyChannelID string,
counterpartyVersion string,
) error
OnChanOpenConfirm(
ctx sdk.Context,
portID string,
channelID string,
) error
OnChanCloseInit(
ctx sdk.Context,
portID string,
channelID string,
) error
OnChanCloseConfirm(
ctx sdk.Context,
portID string,
channelID string,
) error
OnRecvPacket(
ctx sdk.Context,
packet channeltypes.Packet,
relayer sdk.AccAddress,
) ibcexported.Acknowledgement
OnAcknowledgementPacket(
ctx sdk.Context,
packet channeltypes.Packet,
acknowledgement []byte,
relayer sdk.AccAddress,
) error
OnTimeoutPacket(
ctx sdk.Context,
packet channeltypes.Packet,
relayer sdk.AccAddress,
) error
}
2.6 回调执行顺序
当数据包通过多层中间件时,回调的执行顺序如下:
| 操作 | 回调调用顺序 |
|---|---|
SendPacket |
Middleware1.SendPacket → Middleware2.SendPacket → BaseApp.SendPacket |
OnRecvPacket |
Middleware1.OnRecvPacket → Middleware2.OnRecvPacket → BaseApp.OnRecvPacket |
OnAcknowledgementPacket |
BaseApp.OnAck → Middleware2.OnAck → Middleware1.OnAck |
OnTimeoutPacket |
BaseApp.OnTimeout → Middleware2.OnTimeout → Middleware1.OnTimeout |
注意:发送和接收是从外到内,确认和超时是从内到外。
2.7 通道握手回调
通道建立时的回调顺序:
ChanOpenInit:
Middleware1.OnChanOpenInit
→ Middleware2.OnChanOpenInit
→ BaseApp.OnChanOpenInit
ChanOpenTry:
Middleware1.OnChanOpenTry
→ Middleware2.OnChanOpenTry
→ BaseApp.OnChanOpenTry
ChanOpenAck:
Middleware1.OnChanOpenAck
→ Middleware2.OnChanOpenAck
→ BaseApp.OnChanOpenAck
ChanOpenConfirm:
Middleware1.OnChanOpenConfirm
→ Middleware2.OnChanOpenConfirm
→ BaseApp.OnChanOpenConfirm
2.8 中间件 vs 消息处理器
关键区别:
| 方面 | 中间件 | 消息处理器 |
|---|---|---|
| 作用域 | 数据包级别 | 交易级别 |
| 可访问性 | IBCModule 接口 | MsgServer 接口 |
| 典型操作 | 验证、拦截、修改 | 存储、业务逻辑 |
| 用户交互 | 间接 | 直接(签名交易) |
中间件通常只负责"策略"(允许/拒绝),而消息处理器负责"机制"(具体执行)。
3. Go中间件实现
3.1 项目结构
推荐的项目结构:
x/msgratelimit/
├── module/
│ └── module.go # AppModule 定义
├── keeper/
│ ├── keeper.go # Keeper 实现
│ ├── grpc_query.go # gRPC 查询
│ └── params.go # 参数管理
├── types/
│ ├── errors.go # 错误类型
│ ├── keys.go # 存储键
│ ├── params.go # 参数类型
│ ├── codec.go # 编码注册
│ ├── expected_keepers.go # 依赖接口
│ └── messages.go # 消息类型
├── client/
│ └── cli/
│ └── tx.go # CLI 交易命令
├── ibc_module.go # IBCModule 实现
├── ibc_module_test.go # 中间件测试
├── keeper_test.go # Keeper 测试
└── genesis_test.go # Genesis 测试
3.2 基础 IBCModule 中间件骨架
最精简的中间件实现——什么也不做,直接透传:
// X-MSG-Stub=true: ibc_module.go - 基础中间件骨架
package msgratelimit
import (
sdk "github.com/cosmos/cosmos-sdk/types"
capabilitiestypes "github.com/cosmos/cosmos-sdk/x/capability/types"
channeltypes "github.com/cosmos/ibc-go/v8/modules/core/04-channel/types"
ibcexported "github.com/cosmos/ibc-go/v8/modules/core/exported"
ibcporttypes "github.com/cosmos/ibc-go/v8/modules/core/05-port/types"
)
var _ ibcporttypes.IBCModule = IBCMiddleware{}
// IBCMiddleware 是透传中间件
type IBCMiddleware struct {
keeper RateLimitKeeper
channelKeeper channelkeeper.Keeper
next ibcporttypes.IBCModule
}
func NewIBCMiddleware(
next ibcporttypes.IBCModule,
keeper RateLimitKeeper,
channelKeeper channelkeeper.Keeper,
) IBCMiddleware {
return IBCMiddleware{
keeper: keeper,
channelKeeper: channelKeeper,
next: next,
}
}
// OnChanOpenInit 透传给下一个模块
func (m IBCMiddleware) OnChanOpenInit(
ctx sdk.Context,
order channeltypes.Order,
connectionHops []string,
portID string,
channelID string,
channelCap *capabilitiestypes.Capability,
counterparty channeltypes.Counterparty,
version string,
) (string, error) {
return m.next.OnChanOpenInit(
ctx, order, connectionHops, portID, channelID,
channelCap, counterparty, version,
)
}
// OnChanOpenTry 透传给下一个模块
func (m IBCMiddleware) OnChanOpenTry(
ctx sdk.Context,
order channeltypes.Order,
connectionHops []string,
portID string,
channelID string,
channelCap *capabilitiestypes.Capability,
counterparty channeltypes.Counterparty,
version string,
) (string, error) {
return m.next.OnChanOpenTry(
ctx, order, connectionHops, portID, channelID,
channelCap, counterparty, version,
)
}
// OnChanOpenAck 透传给下一个模块
func (m IBCMiddleware) OnChanOpenAck(
ctx sdk.Context,
portID string,
channelID string,
counterpartyChannelID string,
counterpartyVersion string,
) error {
return m.next.OnChanOpenAck(
ctx, portID, channelID,
counterpartyChannelID, counterpartyVersion,
)
}
// OnChanOpenConfirm 透传给下一个模块
func (m IBCMiddleware) OnChanOpenConfirm(
ctx sdk.Context,
portID string,
channelID string,
) error {
return m.next.OnChanOpenConfirm(ctx, portID, channelID)
}
// OnChanCloseInit 透传给下一个模块
func (m IBCMiddleware) OnChanCloseInit(
ctx sdk.Context,
portID string,
channelID string,
) error {
return m.next.OnChanCloseInit(ctx, portID, channelID)
}
// OnChanCloseConfirm 透传给下一个模块
func (m IBCMiddleware) OnChanCloseConfirm(
ctx sdk.Context,
portID string,
channelID string,
) error {
return m.next.OnChanCloseConfirm(ctx, portID, channelID)
}
// OnRecvPacket 透传给下一个模块
func (m IBCMiddleware) OnRecvPacket(
ctx sdk.Context,
packet channeltypes.Packet,
relayer sdk.AccAddress,
) ibcexported.Acknowledgement {
return m.next.OnRecvPacket(ctx, packet, relayer)
}
// OnAcknowledgementPacket 透传给下一个模块
func (m IBCMiddleware) OnAcknowledgementPacket(
ctx sdk.Context,
packet channeltypes.Packet,
acknowledgement []byte,
relayer sdk.AccAddress,
) error {
return m.next.OnAcknowledgementPacket(ctx, packet, acknowledgement, relayer)
}
// OnTimeoutPacket 透传给下一个模块
func (m IBCMiddleware) OnTimeoutPacket(
ctx sdk.Context,
packet channeltypes.Packet,
relayer sdk.AccAddress,
) error {
return m.next.OnTimeoutPacket(ctx, packet, relayer)
}
3.3 发送端拦截
在数据包发送时(SendPacket)执行自定义逻辑的中间件:
// X-MSG-Stub=true: 发送端拦截中间件示例
func (m IBCMiddleware) SendPacket(
ctx sdk.Context,
chanCap *capabilitiestypes.Capability,
packet ibcexported.PacketI,
) error {
// 前置处理:检查是否可以发送
if err := m.ValidateSendPacket(ctx, packet); err != nil {
return err
}
// 调用下一个中间件/基础应用的 SendPacket
if err := m.next.SendPacket(ctx, chanCap, packet); err != nil {
return err
}
// 后置处理:发送成功后的处理
m.AfterSendPacket(ctx, packet)
return nil
}
3.4 接收端拦截
在数据包接收时(OnRecvPacket)执行自定义逻辑:
// X-MSG-Stub=true: 接收端拦截中间件示例
func (m IBCMiddleware) OnRecvPacket(
ctx sdk.Context,
packet channeltypes.Packet,
relayer sdk.AccAddress,
) ibcexported.Acknowledgement {
// 前置验证
if err := m.ValidateRecvPacket(ctx, packet); err != nil {
return channeltypes.NewErrorAcknowledgement(err)
}
// 调用下一个模块
ack := m.next.OnRecvPacket(ctx, packet, relayer)
// 后置处理
if ack.Success() {
m.AfterRecvPacket(ctx, packet)
}
return ack
}
3.5 确认端拦截
在数据包确认时(OnAcknowledgementPacket)执行自定义逻辑:
// X-MSG-Stub=true: 确认端拦截中间件示例
func (m IBCMiddleware) OnAcknowledgementPacket(
ctx sdk.Context,
packet channeltypes.Packet,
acknowledgement []byte,
relayer sdk.AccAddress,
) error {
// 前置处理
m.BeforeAckPacket(ctx, packet)
// 调用下一个模块(注意:确认是从内到外,所以后调用)
if err := m.next.OnAcknowledgementPacket(ctx, packet, acknowledgement, relayer); err != nil {
return err
}
// 解析确认结果
var ack channeltypes.Acknowledgement
if err := ack.Unmarshal(acknowledgement); err != nil {
return err
}
if ack.Success() {
// 确认成功
m.HandleSuccessfulAck(ctx, packet)
} else {
// 确认失败(接收端拒绝)
m.HandleFailedAck(ctx, packet)
}
return nil
}
3.6 超时端拦截
在数据包超时时(OnTimeoutPacket)执行自定义逻辑:
// X-MSG-Stub=true: 超时端拦截中间件示例
func (m IBCMiddleware) OnTimeoutPacket(
ctx sdk.Context,
packet channeltypes.Packet,
relayer sdk.AccAddress,
) error {
// 前置处理
m.BeforeTimeoutPacket(ctx, packet)
// 调用下一个模块
if err := m.next.OnTimeoutPacket(ctx, packet, relayer); err != nil {
return err
}
// 后置处理:超时后的清理逻辑
m.HandlePacketTimeout(ctx, packet)
return nil
}
3.7 Keeper 实现模式
中间件的存储层(keeper)遵循标准 Cosmos SDK Keeper 模式:
// X-MSG-Stub=true: keeper/keeper.go
package keeper
import (
"fmt"
"github.com/cosmos/cosmos-sdk/codec"
storetypes "github.com/cosmos/cosmos-sdk/store/types"
sdk "github.com/cosmos/cosmos-sdk/types"
paramtypes "github.com/cosmos/cosmos-sdk/x/params/types"
"github.com/cosmos/ibc-go/v8/modules/core/04-channel/types"
)
type Keeper struct {
storeKey storetypes.StoreKey
cdc codec.BinaryCodec
paramSpace paramtypes.Subspace
// 依赖的外部 keeper
bankKeeper types.BankKeeper
channelKeeper types.ChannelKeeper
// 地址
authority string // 治理模块地址
}
func NewKeeper(
storeKey storetypes.StoreKey,
cdc codec.BinaryCodec,
paramSpace paramtypes.Subspace,
bankKeeper types.BankKeeper,
channelKeeper types.ChannelKeeper,
authority string,
) Keeper {
if !paramSpace.HasKeyTable() {
paramSpace = paramSpace.WithKeyTable(ParamKeyTable())
}
return Keeper{
storeKey: storeKey,
cdc: cdc,
paramSpace: paramSpace,
bankKeeper: bankKeeper,
channelKeeper: channelKeeper,
authority: authority,
}
}
func (k Keeper) GetAuthority() string {
return k.authority
}
3.8 存储键设计
// X-MSG-Stub=true: types/keys.go
package types
const (
ModuleName = "msgratelimit"
StoreKey = ModuleName
RouterKey = ModuleName
QuerierRoute = ModuleName
)
// KVStore 键前缀
var (
// 通道速率限制配置
// Key: 0x01 | channelID
// Value: RateLimit
ChannelRateLimitPrefix = []byte{0x01}
// 通道当前周期用量
// Key: 0x02 | channelID | periodIndex
// Value: sdk.Int (累计流出量)
ChannelPeriodUsagePrefix = []byte{0x02}
// 通道暂停状态
// Key: 0x03 | channelID
// Value: bool
ChannelPausedPrefix = []byte{0x03}
// 全局参数
// Key: 0x10
// Value: Params
ParamsKey = []byte{0x10}
)
func ChannelRateLimitKey(channelID string) []byte {
return append(ChannelRateLimitPrefix, []byte(channelID)...)
}
func ChannelPeriodUsageKey(channelID string, periodIndex uint64) []byte {
return append(
ChannelPeriodUsagePrefix,
append([]byte(channelID), sdk.Uint64ToBigEndian(periodIndex)...)...,
)
}
func ChannelPausedKey(channelID string) []byte {
return append(ChannelPausedPrefix, []byte(channelID)...)
}
3.9 错误类型
// X-MSG-Stub=true: types/errors.go
package types
import (
sdkerrors "github.com/cosmos/cosmos-sdk/types/errors"
)
var (
ErrRateLimitExceeded = sdkerrors.Register(
ModuleName, 2, "rate limit exceeded for channel",
)
ErrChannelPaused = sdkerrors.Register(
ModuleName, 3, "channel is paused by governance",
)
ErrUnauthorized = sdkerrors.Register(
ModuleName, 4, "unauthorized caller",
)
ErrInvalidChannel = sdkerrors.Register(
ModuleName, 5, "invalid channel",
)
ErrInvalidPacket = sdkerrors.Register(
ModuleName, 6, "invalid packet",
)
ErrChannelNotConfigured = sdkerrors.Register(
ModuleName, 7, "channel not configured with rate limit",
)
)
3.10 Expected Keepers 接口
// X-MSG-Stub=true: types/expected_keepers.go
package types
import (
sdk "github.com/cosmos/cosmos-sdk/types"
channeltypes "github.com/cosmos/ibc-go/v8/modules/core/04-channel/types"
)
// ChannelKeeper 定义了中间件需要的 IBC 通道 keeper 接口
type ChannelKeeper interface {
GetChannel(
ctx sdk.Context,
portID string,
channelID string,
) (channeltypes.Channel, bool)
GetNextSequenceSend(
ctx sdk.Context,
portID string,
channelID string,
) (uint64, bool)
SendPacket(
ctx sdk.Context,
channelCap *capabilitytypes.Capability,
packet ibcexported.PacketI,
) error
ChanCloseInit(
ctx sdk.Context,
portID string,
channelID string,
chanCap *capabilitytypes.Capability,
) error
}
// BankKeeper 定义了需要的 bank keeper 接口
type BankKeeper interface {
SendCoins(
ctx sdk.Context,
fromAddr sdk.AccAddress,
toAddr sdk.AccAddress,
amt sdk.Coins,
) error
GetBalance(
ctx sdk.Context,
addr sdk.AccAddress,
denom string,
) sdk.Coin
}
// AccountKeeper 定义了需要的 account keeper 接口
type AccountKeeper interface {
GetModuleAddress(moduleName string) sdk.AccAddress
GetModuleAccount(ctx sdk.Context, moduleName string) sdk.ModuleAccountI
}
3.11 AppModule 注册
// X-MSG-Stub=true: module/module.go
package module
import (
"encoding/json"
"fmt"
"github.com/cosmos/cosmos-sdk/client"
"github.com/cosmos/cosmos-sdk/codec"
codectypes "github.com/cosmos/cosmos-sdk/codec/types"
sdk "github.com/cosmos/cosmos-sdk/types"
"github.com/cosmos/cosmos-sdk/types/module"
simtypes "github.com/cosmos/cosmos-sdk/types/simulation"
"github.com/gorilla/mux"
"github.com/spf13/cobra"
"github.com/grpc-ecosystem/grpc-gateway/runtime"
abci "github.com/cometbft/cometbft/abci/types"
"msgchain/x/msgratelimit/keeper"
"msgchain/x/msgratelimit/types"
)
var (
_ module.AppModule = AppModule{}
_ module.AppModuleBasic = AppModuleBasic{}
)
type AppModuleBasic struct {
cdc codec.BinaryCodec
}
func (b AppModuleBasic) Name() string { return types.ModuleName }
func (b AppModuleBasic) RegisterLegacyAminoCodec(cdc *codec.LegacyAmino) {
types.RegisterLegacyAminoCodec(cdc)
}
func (b AppModuleBasic) DefaultGenesis(cdc codec.JSONCodec) json.RawMessage {
return cdc.MustMarshalJSON(types.DefaultGenesis())
}
func (b AppModuleBasic) ValidateGenesis(
cdc codec.JSONCodec,
_ client.TxEncodingConfig,
bz json.RawMessage,
) error {
var data types.GenesisState
if err := cdc.UnmarshalJSON(bz, &data); err != nil {
return fmt.Errorf("failed to unmarshal %s genesis: %w", types.ModuleName, err)
}
return data.Validate()
}
func (b AppModuleBasic) RegisterRESTRoutes(_ client.Context, _ *mux.Router) {}
func (b AppModuleBasic) RegisterGRPCGatewayRoutes(
clientCtx client.Context,
mux *runtime.ServeMux,
) {
types.RegisterQueryHandlerClient(context.Background(), mux, types.NewQueryClient(clientCtx))
}
func (b AppModuleBasic) GetTxCmd() *cobra.Command {
return cli.GetTxCmd()
}
func (b AppModuleBasic) GetQueryCmd() *cobra.Command {
return cli.GetQueryCmd()
}
func (b AppModuleBasic) RegisterInterfaces(registry codectypes.InterfaceRegistry) {
types.RegisterInterfaces(registry)
}
type AppModule struct {
AppModuleBasic
keeper keeper.Keeper
}
func NewAppModule(keeper keeper.Keeper) AppModule {
return AppModule{
AppModuleBasic: AppModuleBasic{},
keeper: keeper,
}
}
func (am AppModule) InitGenesis(ctx sdk.Context, cdc codec.JSONCodec, data json.RawMessage) []abci.ValidatorUpdate {
var genesisState types.GenesisState
cdc.MustUnmarshalJSON(data, &genesisState)
am.keeper.InitGenesis(ctx, genesisState)
return []abci.ValidatorUpdate{}
}
func (am AppModule) ExportGenesis(ctx sdk.Context, cdc codec.JSONCodec) json.RawMessage {
gs := am.keeper.ExportGenesis(ctx)
return cdc.MustMarshalJSON(gs)
}
func (am AppModule) RegisterInvariants(_ sdk.InvariantRegistry) {}
func (am AppModule) Route() sdk.Route {
return sdk.Route{}
}
func (am AppModule) QuerierRoute() string { return types.QuerierRoute }
func (am AppModule) LegacyQuerierHandler(_ *codec.LegacyAmino) sdk.Querier {
return nil
}
func (am AppModule) RegisterServices(cfg module.Configurator) {
types.RegisterMsgServer(cfg.MsgServer(), keeper.NewMsgServerImpl(am.keeper))
types.RegisterQueryServer(cfg.QueryServer(), keeper.NewQueryServerImpl(am.keeper))
}
func (am AppModule) ConsensusVersion() uint64 { return 1 }
func (am AppModule) BeginBlock(_ sdk.Context, _ abci.RequestBeginBlock) {}
func (am AppModule) EndBlock(_ sdk.Context, _ abci.RequestEndBlock) []abci.ValidatorUpdate {
return []abci.ValidatorUpdate{}
}
4. 速率限制中间件
4.1 设计目标
速率限制中间件的核心功能:
- 按通道限制:每个 IBC 通道可以独立配置流出速率
- 时间窗口:基于滑动时间窗口(如每小时、每天、每周)累计流出量
- 治理可调:所有限制参数可通过治理提案更新
- 不阻塞接收:只限制发送方向,不影响接收方向
4.2 参数类型
// X-MSG-Stub=true: types/params.go
package types
import (
"time"
sdk "github.com/cosmos/cosmos-sdk/types"
paramtypes "github.com/cosmos/cosmos-sdk/x/params/types"
)
// Params 全局参数
type Params struct {
// 默认时间窗口(秒)
DefaultWindowDuration uint64 `json:"default_window_duration"`
// 默认最大流出量(umsg)
DefaultMaxOutflow sdk.Int `json:"default_max_outflow"`
// 是否启用速率限制
Enabled bool `json:"enabled"`
}
func DefaultParams() Params {
return Params{
DefaultWindowDuration: 86400, // 24小时
DefaultMaxOutflow: sdk.NewInt(1_000_000_000_000_000_000), // 1 MSG
Enabled: false,
}
}
func (p Params) Validate() error {
if p.DefaultWindowDuration == 0 {
return fmt.Errorf("window duration must be positive")
}
if p.DefaultMaxOutflow.IsNegative() {
return fmt.Errorf("max outflow cannot be negative")
}
return nil
}
// ChannelRateLimit 通道级速率限制配置
type ChannelRateLimit struct {
ChannelID string `json:"channel_id"`
WindowDuration uint64 `json:"window_duration"` // 时间窗口(秒)
MaxOutflow sdk.Int `json:"max_outflow"` // 窗口内最大流出量
Active bool `json:"active"` // 是否启用
}
func NewChannelRateLimit(
channelID string,
window uint64,
maxOutflow sdk.Int,
) ChannelRateLimit {
return ChannelRateLimit{
ChannelID: channelID,
WindowDuration: window,
MaxOutflow: maxOutflow,
Active: true,
}
}
// ParamKeyTable 返回参数键表
func ParamKeyTable() paramtypes.KeyTable {
return paramtypes.NewKeyTable().RegisterParamSet(&Params{})
}
func (p *Params) ParamSetPairs() paramtypes.ParamSetPairs {
return paramtypes.ParamSetPairs{
paramtypes.NewParamSetPair(
[]byte("DefaultWindowDuration"),
&p.DefaultWindowDuration,
validateWindowDuration,
),
paramtypes.NewParamSetPair(
[]byte("DefaultMaxOutflow"),
&p.DefaultMaxOutflow,
validateMaxOutflow,
),
paramtypes.NewParamSetPair(
[]byte("Enabled"),
&p.Enabled,
validateEnabled,
),
}
}
func validateWindowDuration(i interface{}) error {
v, ok := i.(uint64)
if !ok {
return fmt.Errorf("invalid parameter type: %T", i)
}
if v == 0 {
return fmt.Errorf("window duration must be positive")
}
return nil
}
func validateMaxOutflow(i interface{}) error {
v, ok := i.(sdk.Int)
if !ok {
return fmt.Errorf("invalid parameter type: %T", i)
}
if v.IsNegative() {
return fmt.Errorf("max outflow cannot be negative")
}
return nil
}
func validateEnabled(i interface{}) error {
_, ok := i.(bool)
if !ok {
return fmt.Errorf("invalid parameter type: %T", i)
}
return nil
}
4.3 Keeper 实现
// X-MSG-Stub=true: keeper/keeper.go - 速率限制 keeper
package keeper
import (
"encoding/binary"
"fmt"
"github.com/cosmos/cosmos-sdk/codec"
storetypes "github.com/cosmos/cosmos-sdk/store/types"
sdk "github.com/cosmos/cosmos-sdk/types"
paramtypes "github.com/cosmos/cosmos-sdk/x/params/types"
"msgchain/x/msgratelimit/types"
)
type Keeper struct {
storeKey storetypes.StoreKey
cdc codec.BinaryCodec
paramSpace paramtypes.Subspace
bankKeeper types.BankKeeper
channelKeeper types.ChannelKeeper
authority string
}
func NewKeeper(
storeKey storetypes.StoreKey,
cdc codec.BinaryCodec,
paramSpace paramtypes.Subspace,
bankKeeper types.BankKeeper,
channelKeeper types.ChannelKeeper,
authority string,
) Keeper {
if !paramSpace.HasKeyTable() {
paramSpace = paramSpace.WithKeyTable(types.ParamKeyTable())
}
return Keeper{
storeKey: storeKey,
cdc: cdc,
paramSpace: paramSpace,
bankKeeper: bankKeeper,
channelKeeper: channelKeeper,
authority: authority,
}
}
// GetParams 获取全局参数
func (k Keeper) GetParams(ctx sdk.Context) types.Params {
var params types.Params
k.paramSpace.GetParamSet(ctx, ¶ms)
return params
}
// SetParams 设置全局参数
func (k Keeper) SetParams(ctx sdk.Context, params types.Params) {
k.paramSpace.SetParamSet(ctx, ¶ms)
}
// GetChannelRateLimit 获取通道速率限制配置
func (k Keeper) GetChannelRateLimit(
ctx sdk.Context,
channelID string,
) (types.ChannelRateLimit, bool) {
store := ctx.KVStore(k.storeKey)
bz := store.Get(types.ChannelRateLimitKey(channelID))
if bz == nil {
return types.ChannelRateLimit{}, false
}
var limit types.ChannelRateLimit
k.cdc.MustUnmarshal(bz, &limit)
return limit, true
}
// SetChannelRateLimit 设置通道速率限制配置
func (k Keeper) SetChannelRateLimit(
ctx sdk.Context,
limit types.ChannelRateLimit,
) {
store := ctx.KVStore(k.storeKey)
bz := k.cdc.MustMarshal(&limit)
store.Set(types.ChannelRateLimitKey(limit.ChannelID), bz)
}
// DeleteChannelRateLimit 删除通道速率限制配置
func (k Keeper) DeleteChannelRateLimit(
ctx sdk.Context,
channelID string,
) {
store := ctx.KVStore(k.storeKey)
store.Delete(types.ChannelRateLimitKey(channelID))
store.Delete(types.ChannelPausedKey(channelID))
}
// GetAllChannelRateLimits 获取所有通道速率限制
func (k Keeper) GetAllChannelRateLimits(
ctx sdk.Context,
) []types.ChannelRateLimit {
store := ctx.KVStore(k.storeKey)
iter := sdk.KVStorePrefixIterator(store, types.ChannelRateLimitPrefix)
defer iter.Close()
var limits []types.ChannelRateLimit
for ; iter.Valid(); iter.Next() {
var limit types.ChannelRateLimit
k.cdc.MustUnmarshal(iter.Value(), &limit)
limits = append(limits, limit)
}
return limits
}
// GetPeriodUsage 获取当前周期累计流出量
func (k Keeper) GetPeriodUsage(
ctx sdk.Context,
channelID string,
periodIndex uint64,
) sdk.Int {
store := ctx.KVStore(k.storeKey)
bz := store.Get(types.ChannelPeriodUsageKey(channelID, periodIndex))
if bz == nil {
return sdk.ZeroInt()
}
amount := sdk.Int{}
amount.Unmarshal(bz)
return amount
}
// SetPeriodUsage 设置周期累计流出量
func (k Keeper) SetPeriodUsage(
ctx sdk.Context,
channelID string,
periodIndex uint64,
amount sdk.Int,
) {
store := ctx.KVStore(k.storeKey)
bz, _ := amount.Marshal()
store.Set(types.ChannelPeriodUsageKey(channelID, periodIndex), bz)
}
// AddToPeriodUsage 增加周期累计流出量
func (k Keeper) AddToPeriodUsage(
ctx sdk.Context,
channelID string,
amount sdk.Int,
) {
limit, found := k.GetChannelRateLimit(ctx, channelID)
if !found {
return
}
now := ctx.BlockTime().Unix()
periodIndex := uint64(now) / limit.WindowDuration
current := k.GetPeriodUsage(ctx, channelID, periodIndex)
k.SetPeriodUsage(ctx, channelID, periodIndex, current.Add(amount))
}
// CheckAndConsumeRateLimit 检查并消耗速率限制
// 如果超出限制则返回错误
func (k Keeper) CheckAndConsumeRateLimit(
ctx sdk.Context,
channelID string,
amount sdk.Int,
) error {
params := k.GetParams(ctx)
if !params.Enabled {
return nil
}
// 检查通道是否暂停
if k.IsChannelPaused(ctx, channelID) {
return types.ErrChannelPaused
}
// 获取通道限制配置
limit, found := k.GetChannelRateLimit(ctx, channelID)
if !found {
// 未配置限制,使用默认参数
limit = types.ChannelRateLimit{
ChannelID: channelID,
WindowDuration: params.DefaultWindowDuration,
MaxOutflow: params.DefaultMaxOutflow,
Active: true,
}
}
if !limit.Active {
return nil // 未启用限制
}
// 计算当前窗口
now := ctx.BlockTime().Unix()
periodIndex := uint64(now) / limit.WindowDuration
currentUsage := k.GetPeriodUsage(ctx, channelID, periodIndex)
// 检查是否超出限制
newUsage := currentUsage.Add(amount)
if newUsage.GT(limit.MaxOutflow) {
return types.ErrRateLimitExceeded
}
// 消耗额度
k.SetPeriodUsage(ctx, channelID, periodIndex, newUsage)
return nil
}
// IsChannelPaused 检查通道是否暂停
func (k Keeper) IsChannelPaused(ctx sdk.Context, channelID string) bool {
store := ctx.KVStore(k.storeKey)
return store.Has(types.ChannelPausedKey(channelID))
}
// PauseChannel 暂停通道
func (k Keeper) PauseChannel(ctx sdk.Context, channelID string) {
store := ctx.KVStore(k.storeKey)
store.Set(types.ChannelPausedKey(channelID), []byte{1})
}
// UnpauseChannel 恢复通道
func (k Keeper) UnpauseChannel(ctx sdk.Context, channelID string) {
store := ctx.KVStore(k.storeKey)
store.Delete(types.ChannelPausedKey(channelID))
}
// InitGenesis 初始化创世状态
func (k Keeper) InitGenesis(ctx sdk.Context, gs types.GenesisState) {
k.SetParams(ctx, gs.Params)
for _, limit := range gs.ChannelRateLimits {
k.SetChannelRateLimit(ctx, limit)
}
}
// ExportGenesis 导出创世状态
func (k Keeper) ExportGenesis(ctx sdk.Context) types.GenesisState {
return types.GenesisState{
Params: k.GetParams(ctx),
ChannelRateLimits: k.GetAllChannelRateLimits(ctx),
}
}
4.4 IBCModule 实现 - 速率限制
// X-MSG-Stub=true: ibc_module.go - 速率限制中间件
package msgratelimit
import (
sdk "github.com/cosmos/cosmos-sdk/types"
capabilitiestypes "github.com/cosmos/cosmos-sdk/x/capability/types"
channeltypes "github.com/cosmos/ibc-go/v8/modules/core/04-channel/types"
ibcexported "github.com/cosmos/ibc-go/v8/modules/core/exported"
ibcporttypes "github.com/cosmos/ibc-go/v8/modules/core/05-port/types"
"msgchain/x/msgratelimit/keeper"
"msgchain/x/msgratelimit/types"
)
var _ ibcporttypes.IBCModule = IBCMiddleware{}
type IBCMiddleware struct {
keeper keeper.Keeper
channelKeeper types.ChannelKeeper
next ibcporttypes.IBCModule
}
func NewIBCMiddleware(
next ibcporttypes.IBCModule,
keeper keeper.Keeper,
channelKeeper types.ChannelKeeper,
) IBCMiddleware {
return IBCMiddleware{
keeper: keeper,
channelKeeper: channelKeeper,
next: next,
}
}
// OnChanOpenInit 验证通道初始化
func (m IBCMiddleware) OnChanOpenInit(
ctx sdk.Context,
order channeltypes.Order,
connectionHops []string,
portID string,
channelID string,
channelCap *capabilitiestypes.Capability,
counterparty channeltypes.Counterparty,
version string,
) (string, error) {
// 在新通道上应用默认限制配置
params := m.keeper.GetParams(ctx)
if params.Enabled {
limit := types.ChannelRateLimit{
ChannelID: channelID,
WindowDuration: params.DefaultWindowDuration,
MaxOutflow: params.DefaultMaxOutflow,
Active: true,
}
m.keeper.SetChannelRateLimit(ctx, limit)
}
return m.next.OnChanOpenInit(
ctx, order, connectionHops, portID, channelID,
channelCap, counterparty, version,
)
}
// OnRecvPacket 接收数据包 - 不做限制(只限制发送)
func (m IBCMiddleware) OnRecvPacket(
ctx sdk.Context,
packet channeltypes.Packet,
relayer sdk.AccAddress,
) ibcexported.Acknowledgement {
return m.next.OnRecvPacket(ctx, packet, relayer)
}
// SendPacket 发送数据包 - 检查速率限制
func (m IBCMiddleware) SendPacket(
ctx sdk.Context,
chanCap *capabilitiestypes.Capability,
packet ibcexported.PacketI,
) error {
// 解析数据包获取转账金额
channelID := packet.GetSourceChannel()
transferAmount, err := m.parsePacketAmount(packet)
if err != nil {
return err
}
// 检查速率限制(会消耗额度)
if err := m.keeper.CheckAndConsumeRateLimit(
ctx, channelID, transferAmount,
); err != nil {
return err
}
// 放行
return m.next.SendPacket(ctx, chanCap, packet)
}
// OnAcknowledgementPacket 处理确认 - 失败时回退额度
func (m IBCMiddleware) OnAcknowledgementPacket(
ctx sdk.Context,
packet channeltypes.Packet,
acknowledgement []byte,
relayer sdk.AccAddress,
) error {
// 解析 ACK,如果失败则回退额度
var ack channeltypes.Acknowledgement
if err := ack.Unmarshal(acknowledgement); err == nil {
if !ack.Success() {
// 接收失败,回退已消耗的额度
channelID := packet.GetSourceChannel()
amount, err := m.parsePacketAmount(packet)
if err == nil {
m.keeper.AddToPeriodUsage(ctx, channelID, amount.Neg())
}
}
}
return m.next.OnAcknowledgementPacket(
ctx, packet, acknowledgement, relayer,
)
}
// OnTimeoutPacket 超时处理 - 回退额度
func (m IBCMiddleware) OnTimeoutPacket(
ctx sdk.Context,
packet channeltypes.Packet,
relayer sdk.AccAddress,
) error {
// 超时意味着发送未完成,回退额度
channelID := packet.GetSourceChannel()
amount, err := m.parsePacketAmount(packet)
if err == nil {
m.keeper.AddToPeriodUsage(ctx, channelID, amount.Neg())
}
return m.next.OnTimeoutPacket(ctx, packet, relayer)
}
// parsePacketAmount 从 IBC 数据包解析转账金额
// 需要处理 ICS-20 FungibleTokenPacketData
func (m IBCMiddleware) parsePacketAmount(
packet ibcexported.PacketI,
) (sdk.Int, error) {
var data channeltypes.FungibleTokenPacketData
if err := data.Unmarshal(packet.GetData()); err != nil {
return sdk.ZeroInt(), err
}
amount, ok := sdk.NewIntFromString(data.Amount)
if !ok {
return sdk.ZeroInt(), fmt.Errorf("invalid amount: %s", data.Amount)
}
return amount, nil
}
// --- 透传方法 ---
func (m IBCMiddleware) OnChanOpenTry(
ctx sdk.Context,
order channeltypes.Order,
connectionHops []string,
portID string,
channelID string,
channelCap *capabilitiestypes.Capability,
counterparty channeltypes.Counterparty,
version string,
) (string, error) {
return m.next.OnChanOpenTry(
ctx, order, connectionHops, portID, channelID,
channelCap, counterparty, version,
)
}
func (m IBCMiddleware) OnChanOpenAck(
ctx sdk.Context,
portID string,
channelID string,
counterpartyChannelID string,
counterpartyVersion string,
) error {
return m.next.OnChanOpenAck(
ctx, portID, channelID, counterpartyChannelID, counterpartyVersion,
)
}
func (m IBCMiddleware) OnChanOpenConfirm(
ctx sdk.Context,
portID string,
channelID string,
) error {
return m.next.OnChanOpenConfirm(ctx, portID, channelID)
}
func (m IBCMiddleware) OnChanCloseInit(
ctx sdk.Context,
portID string,
channelID string,
) error {
m.keeper.DeleteChannelRateLimit(ctx, channelID)
return m.next.OnChanCloseInit(ctx, portID, channelID)
}
func (m IBCMiddleware) OnChanCloseConfirm(
ctx sdk.Context,
portID string,
channelID string,
) error {
m.keeper.DeleteChannelRateLimit(ctx, channelID)
return m.next.OnChanCloseConfirm(ctx, portID, channelID)
}
4.5 治理消息
// X-MSG-Stub=true: types/messages.go
package types
import (
sdk "github.com/cosmos/cosmos-sdk/types"
sdkerrors "github.com/cosmos/cosmos-sdk/types/errors"
)
const (
TypeMsgSetChannelRateLimit = "set_channel_rate_limit"
TypeMsgRemoveChannelRateLimit = "remove_channel_rate_limit"
TypeMsgPauseChannel = "pause_channel"
TypeMsgUnpauseChannel = "unpause_channel"
)
// MsgSetChannelRateLimit 设置通道速率限制
type MsgSetChannelRateLimit struct {
Authority string `json:"authority"`
ChannelID string `json:"channel_id"`
WindowDuration uint64 `json:"window_duration"`
MaxOutflow sdk.Int `json:"max_outflow"`
}
func NewMsgSetChannelRateLimit(
authority string,
channelID string,
windowDuration uint64,
maxOutflow sdk.Int,
) *MsgSetChannelRateLimit {
return &MsgSetChannelRateLimit{
Authority: authority,
ChannelID: channelID,
WindowDuration: windowDuration,
MaxOutflow: maxOutflow,
}
}
func (msg MsgSetChannelRateLimit) Route() string { return RouterKey }
func (msg MsgSetChannelRateLimit) Type() string { return TypeMsgSetChannelRateLimit }
func (msg MsgSetChannelRateLimit) GetSigners() []sdk.AccAddress {
addr, _ := sdk.AccAddressFromBech32(msg.Authority)
return []sdk.AccAddress{addr}
}
func (msg MsgSetChannelRateLimit) ValidateBasic() error {
if _, err := sdk.AccAddressFromBech32(msg.Authority); err != nil {
return sdkerrors.ErrInvalidAddress.Wrapf("invalid authority: %s", err)
}
if msg.ChannelID == "" {
return sdkerrors.ErrInvalidRequest.Wrap("channel ID cannot be empty")
}
if msg.WindowDuration == 0 {
return sdkerrors.ErrInvalidRequest.Wrap("window duration must be positive")
}
if msg.MaxOutflow.IsNegative() {
return sdkerrors.ErrInvalidRequest.Wrap("max outflow cannot be negative")
}
return nil
}
func (msg MsgSetChannelRateLimit) GetSignBytes() []byte {
return sdk.MustSortJSON(ModuleCdc.MustMarshalJSON(&msg))
}
// MsgPauseChannel 暂停通道消息
type MsgPauseChannel struct {
Authority string `json:"authority"`
ChannelID string `json:"channel_id"`
}
func NewMsgPauseChannel(authority string, channelID string) *MsgPauseChannel {
return &MsgPauseChannel{
Authority: authority,
ChannelID: channelID,
}
}
func (msg MsgPauseChannel) Route() string { return RouterKey }
func (msg MsgPauseChannel) Type() string { return TypeMsgPauseChannel }
func (msg MsgPauseChannel) GetSigners() []sdk.AccAddress {
addr, _ := sdk.AccAddressFromBech32(msg.Authority)
return []sdk.AccAddress{addr}
}
func (msg MsgPauseChannel) ValidateBasic() error {
if _, err := sdk.AccAddressFromBech32(msg.Authority); err != nil {
return sdkerrors.ErrInvalidAddress.Wrapf("invalid authority: %s", err)
}
if msg.ChannelID == "" {
return sdkerrors.ErrInvalidRequest.Wrap("channel ID cannot be empty")
}
return nil
}
func (msg MsgPauseChannel) GetSignBytes() []byte {
return sdk.MustSortJSON(ModuleCdc.MustMarshalJSON(&msg))
}
// MsgUnpauseChannel 恢复通道消息
type MsgUnpauseChannel struct {
Authority string `json:"authority"`
ChannelID string `json:"channel_id"`
}
func NewMsgUnpauseChannel(authority string, channelID string) *MsgUnpauseChannel {
return &MsgUnpauseChannel{
Authority: authority,
ChannelID: channelID,
}
}
func (msg MsgUnpauseChannel) Route() string { return RouterKey }
func (msg MsgUnpauseChannel) Type() string { return TypeMsgUnpauseChannel }
func (msg MsgUnpauseChannel) GetSigners() []sdk.AccAddress {
addr, _ := sdk.AccAddressFromBech32(msg.Authority)
return []sdk.AccAddress{addr}
}
func (msg MsgUnpauseChannel) ValidateBasic() error {
if _, err := sdk.AccAddressFromBech32(msg.Authority); err != nil {
return sdkerrors.ErrInvalidAddress.Wrapf("invalid authority: %s", err)
}
if msg.ChannelID == "" {
return sdkerrors.ErrInvalidRequest.Wrap("channel ID cannot be empty")
}
return nil
}
func (msg MsgUnpauseChannel) GetSignBytes() []byte {
return sdk.MustSortJSON(ModuleCdc.MustMarshalJSON(&msg))
}
4.6 消息服务器
// X-MSG-Stub=true: keeper/msg_server.go
package keeper
import (
"context"
sdk "github.com/cosmos/cosmos-sdk/types"
"msgchain/x/msgratelimit/types"
)
type msgServer struct {
keeper Keeper
}
func NewMsgServerImpl(keeper Keeper) types.MsgServer {
return &msgServer{keeper: keeper}
}
var _ types.MsgServer = msgServer{}
func (s msgServer) SetChannelRateLimit(
goCtx context.Context,
msg *types.MsgSetChannelRateLimit,
) (*types.MsgSetChannelRateLimitResponse, error) {
ctx := sdk.UnwrapSDKContext(goCtx)
// 验证权限:只有治理模块地址有权限
if s.keeper.GetAuthority() != msg.Authority {
return nil, types.ErrUnauthorized
}
limit := types.ChannelRateLimit{
ChannelID: msg.ChannelID,
WindowDuration: msg.WindowDuration,
MaxOutflow: msg.MaxOutflow,
Active: true,
}
s.keeper.SetChannelRateLimit(ctx, limit)
return &types.MsgSetChannelRateLimitResponse{}, nil
}
func (s msgServer) RemoveChannelRateLimit(
goCtx context.Context,
msg *types.MsgRemoveChannelRateLimit,
) (*types.MsgRemoveChannelRateLimitResponse, error) {
ctx := sdk.UnwrapSDKContext(goCtx)
if s.keeper.GetAuthority() != msg.Authority {
return nil, types.ErrUnauthorized
}
s.keeper.DeleteChannelRateLimit(ctx, msg.ChannelID)
return &types.MsgRemoveChannelRateLimitResponse{}, nil
}
func (s msgServer) PauseChannel(
goCtx context.Context,
msg *types.MsgPauseChannel,
) (*types.MsgPauseChannelResponse, error) {
ctx := sdk.UnwrapSDKContext(goCtx)
if s.keeper.GetAuthority() != msg.Authority {
return nil, types.ErrUnauthorized
}
s.keeper.PauseChannel(ctx, msg.ChannelID)
return &types.MsgPauseChannelResponse{}, nil
}
func (s msgServer) UnpauseChannel(
goCtx context.Context,
msg *types.MsgUnpauseChannel,
) (*types.MsgUnpauseChannelResponse, error) {
ctx := sdk.UnwrapSDKContext(goCtx)
if s.keeper.GetAuthority() != msg.Authority {
return nil, types.ErrUnauthorized
}
s.keeper.UnpauseChannel(ctx, msg.ChannelID)
return &types.MsgUnpauseChannelResponse{}, nil
}
4.7 查询服务器
// X-MSG-Stub=true: keeper/grpc_query.go
package keeper
import (
"context"
sdk "github.com/cosmos/cosmos-sdk/types"
"msgchain/x/msgratelimit/types"
)
type queryServer struct {
keeper Keeper
}
func NewQueryServerImpl(keeper Keeper) types.QueryServer {
return &queryServer{keeper: keeper}
}
var _ types.QueryServer = queryServer{}
func (s queryServer) Params(
goCtx context.Context,
req *types.QueryParamsRequest,
) (*types.QueryParamsResponse, error) {
ctx := sdk.UnwrapSDKContext(goCtx)
params := s.keeper.GetParams(ctx)
return &types.QueryParamsResponse{Params: params}, nil
}
func (s queryServer) ChannelRateLimit(
goCtx context.Context,
req *types.QueryChannelRateLimitRequest,
) (*types.QueryChannelRateLimitResponse, error) {
ctx := sdk.UnwrapSDKContext(goCtx)
limit, found := s.keeper.GetChannelRateLimit(ctx, req.ChannelId)
if !found {
return nil, types.ErrChannelNotConfigured
}
// 计算当前周期用量
now := ctx.BlockTime().Unix()
periodIndex := uint64(now) / limit.WindowDuration
currentUsage := s.keeper.GetPeriodUsage(ctx, req.ChannelId, periodIndex)
return &types.QueryChannelRateLimitResponse{
ChannelRateLimit: limit,
CurrentUsage: currentUsage,
IsPaused: s.keeper.IsChannelPaused(ctx, req.ChannelId),
}, nil
}
func (s queryServer) AllChannelRateLimits(
goCtx context.Context,
req *types.QueryAllChannelRateLimitsRequest,
) (*types.QueryAllChannelRateLimitsResponse, error) {
ctx := sdk.UnwrapSDKContext(goCtx)
limits := s.keeper.GetAllChannelRateLimits(ctx)
return &types.QueryAllChannelRateLimitsResponse{
ChannelRateLimits: limits,
}, nil
}
5. 治理中间件
5.1 设计目标
治理控制中间件的核心功能:
- 治理暂停:通过治理提案暂停指定 IBC 通道的数据包传递
- 数据包过滤:通过治理设置数据包过滤规则(白名单/黑名单)
- 参数管理:通过治理更新中间件参数
- 紧急刹车:快速暂停所有 IBC 活动(安全响应)
5.2 参数类型
// X-MSG-Stub=true: types/params.go - 治理中间件参数
package types
import (
"fmt"
sdk "github.com/cosmos/cosmos-sdk/types"
paramtypes "github.com/cosmos/cosmos-sdk/x/params/types"
)
const (
ModuleName = "msggovibc"
StoreKey = ModuleName
RouterKey = ModuleName
)
// Params 治理中间件全局参数
type Params struct {
// 是否启用治理控制
Enabled bool `json:"enabled"`
// 紧急暂停开关(暂停所有 IBC 流量)
EmergencyPause bool `json:"emergency_pause"`
// 白名单模式(false=黑名单, true=白名单)
WhitelistMode bool `json:"whitelist_mode"`
}
func DefaultParams() Params {
return Params{
Enabled: true,
EmergencyPause: false,
WhitelistMode: false,
}
}
func (p Params) Validate() error {
return nil
}
// ChannelFilter 通道级过滤规则
type ChannelFilter struct {
ChannelID string `json:"channel_id"`
Enabled bool `json:"enabled"`
// 白名单的源端口(仅白名单模式生效)
AllowedPorts []string `json:"allowed_ports"`
// 黑名单的源端口(仅黑名单模式生效)
BlockedPorts []string `json:"blocked_ports"`
// 白名单的目的地址前缀
AllowedAddressPrefixes []string `json:"allowed_address_prefixes"`
// 黑名单的目的地址前缀
BlockedAddressPrefixes []string `json:"blocked_address_prefixes"`
}
// GovernanceProposal 治理提案结构
type GovernanceProposal struct {
Title string `json:"title"`
Description string `json:"description"`
Action ProposalAction `json:"action"`
ChannelID string `json:"channel_id"`
}
type ProposalAction int32
const (
ActionUnspecified ProposalAction = 0
ActionPauseChannel ProposalAction = 1
ActionResumeChannel ProposalAction = 2
ActionSetFilter ProposalAction = 3
ActionEmergencyPause ProposalAction = 4
ActionEmergencyResume ProposalAction = 5
)
5.3 治理中间件 Keeper
// X-MSG-Stub=true: keeper/keeper.go - 治理中间件 keeper
package keeper
import (
storetypes "github.com/cosmos/cosmos-sdk/store/types"
sdk "github.com/cosmos/cosmos-sdk/types"
paramtypes "github.com/cosmos/cosmos-sdk/x/params/types"
"msgchain/x/msggovibc/types"
)
type Keeper struct {
storeKey storetypes.StoreKey
cdc codec.BinaryCodec
paramSpace paramtypes.Subspace
channelKeeper types.ChannelKeeper
authority string
}
func NewKeeper(
storeKey storetypes.StoreKey,
cdc codec.BinaryCodec,
paramSpace paramtypes.Subspace,
channelKeeper types.ChannelKeeper,
authority string,
) Keeper {
if !paramSpace.HasKeyTable() {
paramSpace = paramSpace.WithKeyTable(types.ParamKeyTable())
}
return Keeper{
storeKey: storeKey,
cdc: cdc,
paramSpace: paramSpace,
channelKeeper: channelKeeper,
authority: authority,
}
}
// GetParams 获取参数
func (k Keeper) GetParams(ctx sdk.Context) types.Params {
var params types.Params
k.paramSpace.GetParamSet(ctx, ¶ms)
return params
}
// SetParams 设置参数
func (k Keeper) SetParams(ctx sdk.Context, params types.Params) {
k.paramSpace.SetParamSet(ctx, ¶ms)
}
// GetChannelFilter 获取通道过滤规则
func (k Keeper) GetChannelFilter(
ctx sdk.Context,
channelID string,
) (types.ChannelFilter, bool) {
store := ctx.KVStore(k.storeKey)
bz := store.Get(types.ChannelFilterKey(channelID))
if bz == nil {
return types.ChannelFilter{}, false
}
var filter types.ChannelFilter
k.cdc.MustUnmarshal(bz, &filter)
return filter, true
}
// SetChannelFilter 设置通道过滤规则
func (k Keeper) SetChannelFilter(
ctx sdk.Context,
filter types.ChannelFilter,
) {
store := ctx.KVStore(k.storeKey)
bz := k.cdc.MustMarshal(&filter)
store.Set(types.ChannelFilterKey(filter.ChannelID), bz)
}
// DeleteChannelFilter 删除过滤规则
func (k Keeper) DeleteChannelFilter(ctx sdk.Context, channelID string) {
store := ctx.KVStore(k.storeKey)
store.Delete(types.ChannelFilterKey(channelID))
}
// IsChannelGovPaused 检查通道是否被治理暂停
func (k Keeper) IsChannelGovPaused(ctx sdk.Context, channelID string) bool {
store := ctx.KVStore(k.storeKey)
return store.Has(types.ChannelGovPausedKey(channelID))
}
// PauseChannel 治理暂停通道
func (k Keeper) PauseChannel(ctx sdk.Context, channelID string) {
store := ctx.KVStore(k.storeKey)
store.Set(types.ChannelGovPausedKey(channelID), []byte{1})
}
// ResumeChannel 治理恢复通道
func (k Keeper) ResumeChannel(ctx sdk.Context, channelID string) {
store := ctx.KVStore(k.storeKey)
store.Delete(types.ChannelGovPausedKey(channelID))
}
// IsPacketAllowed 检查数据包是否被过滤规则允许
func (k Keeper) IsPacketAllowed(
ctx sdk.Context,
channelID string,
sourcePort string,
destAddress string,
) bool {
params := k.GetParams(ctx)
// 紧急暂停检查
if params.EmergencyPause {
return false
}
// 检查通道级暂停
if k.IsChannelGovPaused(ctx, channelID) {
return false
}
// 检查过滤规则
filter, found := k.GetChannelFilter(ctx, channelID)
if !found || !filter.Enabled {
return true // 没有过滤规则,默认允许
}
if params.WhitelistMode {
// 白名单模式:匹配白名单则允许
if !contains(filter.AllowedPorts, sourcePort) {
return false
}
if len(filter.AllowedAddressPrefixes) > 0 {
if !matchesAnyPrefix(destAddress, filter.AllowedAddressPrefixes) {
return false
}
}
return true
} else {
// 黑名单模式:匹配黑名单则拒绝
if contains(filter.BlockedPorts, sourcePort) {
return false
}
if matchesAnyPrefix(destAddress, filter.BlockedAddressPrefixes) {
return false
}
return true
}
}
// InitGenesis 创世初始化
func (k Keeper) InitGenesis(ctx sdk.Context, gs types.GenesisState) {
k.SetParams(ctx, gs.Params)
for _, filter := range gs.ChannelFilters {
k.SetChannelFilter(ctx, filter)
}
for _, paused := range gs.PausedChannels {
k.PauseChannel(ctx, paused)
}
}
// ExportGenesis 创世导出
func (k Keeper) ExportGenesis(ctx sdk.Context) types.GenesisState {
var gs types.GenesisState
gs.Params = k.GetParams(ctx)
// 导出过滤规则和暂停列表
return gs
}
func contains(slice []string, item string) bool {
for _, s := range slice {
if s == item {
return true
}
}
return false
}
func matchesAnyPrefix(s string, prefixes []string) bool {
for _, p := range prefixes {
if len(s) >= len(p) && s[:len(p)] == p {
return true
}
}
return false
}
5.4 治理提案处理
治理中间件通过 Cosmos SDK 的治理模块(gov)接受提案。提案通过后,执行对应操作:
// X-MSG-Stub=true: keeper/proposal_handler.go
package keeper
import (
sdk "github.com/cosmos/cosmos-sdk/types"
govv1 "github.com/cosmos/cosmos-sdk/x/gov/types/v1"
"msgchain/x/msggovibc/types"
)
// NewGovProposalHandler 创建治理提案处理器
func NewGovProposalHandler(k Keeper) govv1.ProposalHandler {
return func(ctx sdk.Context, proposal govv1.Proposal) error {
return k.HandleGovProposal(ctx, proposal)
}
}
// HandleGovProposal 处理治理提案
func (k Keeper) HandleGovProposal(
ctx sdk.Context,
proposal govv1.Proposal,
) error {
// 解析提案消息
for _, msg := range proposal.GetMessages() {
var govMsg types.GovernanceProposal
if err := govMsg.Unmarshal(msg.Value); err != nil {
return err
}
switch govMsg.Action {
case types.ActionPauseChannel:
k.PauseChannel(ctx, govMsg.ChannelID)
case types.ActionResumeChannel:
k.ResumeChannel(ctx, govMsg.ChannelID)
case types.ActionSetFilter:
// 设置过滤规则(需要从提案消息中解析 filter 内容)
var filter types.ChannelFilter
if err := filter.Unmarshal(msg.Value); err != nil {
return err
}
k.SetChannelFilter(ctx, filter)
case types.ActionEmergencyPause:
params := k.GetParams(ctx)
params.EmergencyPause = true
k.SetParams(ctx, params)
case types.ActionEmergencyResume:
params := k.GetParams(ctx)
params.EmergencyPause = false
k.SetParams(ctx, params)
}
}
return nil
}
5.5 治理中间件 IBCModule
// X-MSG-Stub=true: ibc_module.go - 治理中间件
package msggovibc
import (
sdk "github.com/cosmos/cosmos-sdk/types"
capabilitiestypes "github.com/cosmos/cosmos-sdk/x/capability/types"
channeltypes "github.com/cosmos/ibc-go/v8/modules/core/04-channel/types"
ibcexported "github.com/cosmos/ibc-go/v8/modules/core/exported"
ibcporttypes "github.com/cosmos/ibc-go/v8/modules/core/05-port/types"
"msgchain/x/msggovibc/keeper"
"msgchain/x/msggovibc/types"
)
var _ ibcporttypes.IBCModule = GovIBCMiddleware{}
type GovIBCMiddleware struct {
keeper keeper.Keeper
channelKeeper types.ChannelKeeper
authKeeper types.AuthKeeper
next ibcporttypes.IBCModule
}
func NewIBCMiddleware(
next ibcporttypes.IBCModule,
keeper keeper.Keeper,
channelKeeper types.ChannelKeeper,
authKeeper types.AuthKeeper,
) GovIBCMiddleware {
return GovIBCMiddleware{
keeper: keeper,
channelKeeper: channelKeeper,
authKeeper: authKeeper,
next: next,
}
}
// SendPacket 发送数据包前检查治理状态
func (m GovIBCMiddleware) SendPacket(
ctx sdk.Context,
chanCap *capabilitiestypes.Capability,
packet ibcexported.PacketI,
) error {
channelID := packet.GetSourceChannel()
// 检查治理暂停
if m.keeper.IsChannelGovPaused(ctx, channelID) {
return types.ErrChannelPaused
}
// 检查紧急暂停
params := m.keeper.GetParams(ctx)
if params.EmergencyPause {
return types.ErrChannelPaused
}
return m.next.SendPacket(ctx, chanCap, packet)
}
// OnRecvPacket 接收数据包前检查治理过滤规则
func (m GovIBCMiddleware) OnRecvPacket(
ctx sdk.Context,
packet channeltypes.Packet,
relayer sdk.AccAddress,
) ibcexported.Acknowledgement {
channelID := packet.GetDestinationChannel()
params := m.keeper.GetParams(ctx)
// 紧急暂停
if params.EmergencyPause {
return channeltypes.NewErrorAcknowledgement(
types.ErrChannelPaused,
)
}
// 通道级暂停
if m.keeper.IsChannelGovPaused(ctx, channelID) {
return channeltypes.NewErrorAcknowledgement(
types.ErrChannelPaused,
)
}
// 过滤规则检查
var data channeltypes.FungibleTokenPacketData
if err := data.Unmarshal(packet.GetData()); err != nil {
return channeltypes.NewErrorAcknowledgement(err)
}
if !m.keeper.IsPacketAllowed(
ctx,
channelID,
packet.GetSourcePort(),
data.Receiver,
) {
return channeltypes.NewErrorAcknowledgement(
types.ErrPacketFiltered,
)
}
return m.next.OnRecvPacket(ctx, packet, relayer)
}
// 以下为透传方法
func (m GovIBCMiddleware) OnChanOpenInit(
ctx sdk.Context,
order channeltypes.Order,
connectionHops []string,
portID string,
channelID string,
channelCap *capabilitiestypes.Capability,
counterparty channeltypes.Counterparty,
version string,
) (string, error) {
return m.next.OnChanOpenInit(
ctx, order, connectionHops, portID, channelID,
channelCap, counterparty, version,
)
}
func (m GovIBCMiddleware) OnChanOpenTry(
ctx sdk.Context,
order channeltypes.Order,
connectionHops []string,
portID string,
channelID string,
channelCap *capabilitiestypes.Capability,
counterparty channeltypes.Counterparty,
version string,
) (string, error) {
return m.next.OnChanOpenTry(
ctx, order, connectionHops, portID, channelID,
channelCap, counterparty, version,
)
}
func (m GovIBCMiddleware) OnChanOpenAck(
ctx sdk.Context,
portID string,
channelID string,
counterpartyChannelID string,
counterpartyVersion string,
) error {
return m.next.OnChanOpenAck(
ctx, portID, channelID, counterpartyChannelID, counterpartyVersion,
)
}
func (m GovIBCMiddleware) OnChanOpenConfirm(
ctx sdk.Context,
portID string,
channelID string,
) error {
return m.next.OnChanOpenConfirm(ctx, portID, channelID)
}
func (m GovIBCMiddleware) OnChanCloseInit(
ctx sdk.Context,
portID string,
channelID string,
) error {
return m.next.OnChanCloseInit(ctx, portID, channelID)
}
func (m GovIBCMiddleware) OnChanCloseConfirm(
ctx sdk.Context,
portID string,
channelID string,
) error {
return m.next.OnChanCloseConfirm(ctx, portID, channelID)
}
func (m GovIBCMiddleware) OnAcknowledgementPacket(
ctx sdk.Context,
packet channeltypes.Packet,
acknowledgement []byte,
relayer sdk.AccAddress,
) error {
return m.next.OnAcknowledgementPacket(
ctx, packet, acknowledgement, relayer,
)
}
func (m GovIBCMiddleware) OnTimeoutPacket(
ctx sdk.Context,
packet channeltypes.Packet,
relayer sdk.AccAddress,
) error {
return m.next.OnTimeoutPacket(ctx, packet, relayer)
}
5.6 治理 CLI 命令
// X-MSG-Stub=true: client/cli/tx.go
package cli
import (
"fmt"
"strconv"
"github.com/cosmos/cosmos-sdk/client"
"github.com/cosmos/cosmos-sdk/client/flags"
"github.com/cosmos/cosmos-sdk/client/tx"
sdk "github.com/cosmos/cosmos-sdk/types"
"github.com/spf13/cobra"
"msgchain/x/msggovibc/types"
)
func GetTxCmd() *cobra.Command {
cmd := &cobra.Command{
Use: types.ModuleName,
Short: "Governance IBC middleware transaction subcommands",
DisableOrdering: true,
}
cmd.AddCommand(
NewPauseChannelCmd(),
NewResumeChannelCmd(),
NewSetFilterCmd(),
)
return cmd
}
func NewPauseChannelCmd() *cobra.Command {
return &cobra.Command{
Use: "pause-channel [channel-id]",
Short: "Pause an IBC channel via governance",
Args: cobra.ExactArgs(1),
RunE: func(cmd *cobra.Command, args []string) error {
clientCtx, err := client.GetClientTxContext(cmd)
if err != nil {
return err
}
channelID := args[0]
authority := clientCtx.GetFromAddress().String()
msg := types.NewMsgPauseChannel(authority, channelID)
if err := msg.ValidateBasic(); err != nil {
return err
}
return tx.GenerateOrBroadcastTxCLI(clientCtx, cmd.Flags(), msg)
},
}
}
func NewResumeChannelCmd() *cobra.Command {
return &cobra.Command{
Use: "resume-channel [channel-id]",
Short: "Resume a paused IBC channel via governance",
Args: cobra.ExactArgs(1),
RunE: func(cmd *cobra.Command, args []string) error {
clientCtx, err := client.GetClientTxContext(cmd)
if err != nil {
return err
}
channelID := args[0]
authority := clientCtx.GetFromAddress().String()
msg := types.NewMsgUnpauseChannel(authority, channelID)
if err := msg.ValidateBasic(); err != nil {
return err
}
return tx.GenerateOrBroadcastTxCLI(clientCtx, cmd.Flags(), msg)
},
}
}
func NewSetFilterCmd() *cobra.Command {
return &cobra.Command{
Use: "set-filter [channel-id] [allow-ports] [block-ports]",
Short: "Set packet filter rules for a channel",
Args: cobra.ExactArgs(3),
RunE: func(cmd *cobra.Command, args []string) error {
clientCtx, err := client.GetClientTxContext(cmd)
if err != nil {
return err
}
channelID := args[0]
// 逗号分隔的端口列表
// TODO: 构建 MsgSetFilter
fmt.Println("Set filter for channel:", channelID)
return nil
},
}
}
6. 中间件测试
6.1 单元测试框架
测试使用 Cosmos SDK 的测试框架和 IBC 的 mock 模块:
// X-MSG-Stub=true: ibc_module_test.go - 速率限制中间件测试
package msgratelimit_test
import (
"testing"
"time"
sdk "github.com/cosmos/cosmos-sdk/types"
capabilitiestypes "github.com/cosmos/cosmos-sdk/x/capability/types"
channeltypes "github.com/cosmos/ibc-go/v8/modules/core/04-channel/types"
ibcexported "github.com/cosmos/ibc-go/v8/modules/core/exported"
ibcporttypes "github.com/cosmos/ibc-go/v8/modules/core/05-port/types"
"github.com/cosmos/ibc-go/v8/modules/core/04-channel/types/mock"
"github.com/stretchr/testify/require"
"github.com/stretchr/testify/suite"
"msgchain/x/msgratelimit"
"msgchain/x/msgratelimit/keeper"
"msgchain/x/msgratelimit/types"
)
type RateLimitMiddlewareTestSuite struct {
suite.Suite
ctx sdk.Context
keeper keeper.Keeper
channelKeeper *mock.ChannelKeeper
nextModule *mock.IBCModule
middleware msgratelimit.IBCMiddleware
}
func (s *RateLimitMiddlewareTestSuite) SetupTest() {
// 初始化测试上下文
s.ctx = sdk.NewContext(nil, sdk.Header{Time: time.Now()}, false, nil)
// 创建 mock keeper
s.keeper = keeper.NewKeeper(
storeKey,
cdc,
paramSpace,
bankKeeper,
channelKeeper,
"msg1govacount",
)
// 创建 mock 下一个模块
s.nextModule = &mock.IBCModule{}
// 创建中间件
s.middleware = msgratelimit.NewIBCMiddleware(
s.nextModule,
s.keeper,
s.channelKeeper,
)
}
func TestRateLimitMiddlewareSuite(t *testing.T) {
suite.Run(t, new(RateLimitMiddlewareTestSuite))
}
// TestNoLimit_BelowThreshold 测试未超出限制
func (s *RateLimitMiddlewareTestSuite) TestNoLimit_BelowThreshold() {
// 设置限制:1000 MSG/天
limit := types.ChannelRateLimit{
ChannelID: "channel-0",
WindowDuration: 86400,
MaxOutflow: sdk.NewInt(1_000_000_000_000_000_000_000), // 1000 MSG
Active: true,
}
s.keeper.SetChannelRateLimit(s.ctx, limit)
// 发送 100 MSG
packet := mock.NewMockPacket("transfer", "channel-0", "transfer", "channel-0", 1, "{\"amount\":\"100000000000000000000\",\"denom\":\"umsg\",\"receiver\":\"msg1...\"}")
err := s.middleware.SendPacket(s.ctx, nil, packet)
s.NoError(err)
// 验证存量
periodIndex := uint64(s.ctx.BlockTime().Unix()) / 86400
usage := s.keeper.GetPeriodUsage(s.ctx, "channel-0", periodIndex)
s.True(usage.Equal(sdk.NewInt(100_000_000_000_000_000_000)))
}
// TestRateLimit_Exceeded 测试超出限制
func (s *RateLimitMiddlewareTestSuite) TestRateLimit_Exceeded() {
// 设置限制:50 MSG/小时
limit := types.ChannelRateLimit{
ChannelID: "channel-0",
WindowDuration: 3600,
MaxOutflow: sdk.NewInt(50_000_000_000_000_000_000), // 50 MSG
Active: true,
}
s.keeper.SetChannelRateLimit(s.ctx, limit)
// 发送 60 MSG(超过限制)
packet := mock.NewMockPacket("transfer", "channel-0", "transfer", "channel-0", 2,
"{\"amount\":\"60000000000000000000\",\"denom\":\"umsg\",\"receiver\":\"msg1...\"}")
err := s.middleware.SendPacket(s.ctx, nil, packet)
s.Error(err)
s.ErrorIs(err, types.ErrRateLimitExceeded)
}
// TestRateLimit_ResetAfterWindow 测试窗口重置
func (s *RateLimitMiddlewareTestSuite) TestRateLimit_ResetAfterWindow() {
limit := types.ChannelRateLimit{
ChannelID: "channel-0",
WindowDuration: 3600,
MaxOutflow: sdk.NewInt(100_000_000_000_000_000_000), // 100 MSG
Active: true,
}
s.keeper.SetChannelRateLimit(s.ctx, limit)
// 发送 100 MSG
packet1 := mock.NewMockPacket("transfer", "channel-0", "transfer", "channel-0", 3,
"{\"amount\":\"100000000000000000000\",\"denom\":\"umsg\",\"receiver\":\"msg1...\"}")
err := s.middleware.SendPacket(s.ctx, nil, packet1)
s.NoError(err)
// 尝试再发 1 MSG(应该超出)
packet2 := mock.NewMockPacket("transfer", "channel-0", "transfer", "channel-0", 4,
"{\"amount\":\"1000000000000000000\",\"denom\":\"umsg\",\"receiver\":\"msg1...\"}")
err = s.middleware.SendPacket(s.ctx, nil, packet2)
s.Error(err)
// 模拟时间推进超过窗口
futureCtx := s.ctx.WithBlockTime(s.ctx.BlockTime().Add(2 * time.Hour))
packet3 := mock.NewMockPacket("transfer", "channel-0", "transfer", "channel-0", 5,
"{\"amount\":\"1000000000000000000\",\"denom\":\"umsg\",\"receiver\":\"msg1...\"}")
err = s.middleware.SendPacket(futureCtx, nil, packet3)
s.NoError(err) // 窗口重置,应该成功
}
// TestNoLimit_Disabled 测试限制禁用
func (s *RateLimitMiddlewareTestSuite) TestNoLimit_Disabled() {
limit := types.ChannelRateLimit{
ChannelID: "channel-0",
WindowDuration: 3600,
MaxOutflow: sdk.NewInt(100),
Active: false, // 禁用
}
s.keeper.SetChannelRateLimit(s.ctx, limit)
packet := mock.NewMockPacket("transfer", "channel-0", "transfer", "channel-0", 6,
"{\"amount\":\"999999999999999999999\",\"denom\":\"umsg\",\"receiver\":\"msg1...\"}")
err := s.middleware.SendPacket(s.ctx, nil, packet)
s.NoError(err) // 禁用状态不限制
}
// TestPacketAcknowledgment_Refund 测试 ACK 失败后退还额度
func (s *RateLimitMiddlewareTestSuite) TestPacketAcknowledgment_Refund() {
limit := types.ChannelRateLimit{
ChannelID: "channel-0",
WindowDuration: 86400,
MaxOutflow: sdk.NewInt(100_000_000_000_000_000_000), // 100 MSG
Active: true,
}
s.keeper.SetChannelRateLimit(s.ctx, limit)
// 发送 100 MSG
packet := mock.NewMockPacket("transfer", "channel-0", "transfer", "channel-0", 7,
"{\"amount\":\"100000000000000000000\",\"denom\":\"umsg\",\"receiver\":\"msg1...\"}")
err := s.middleware.SendPacket(s.ctx, nil, packet)
s.NoError(err)
// 模拟 ACK 失败
failAck := channeltypes.NewErrorAcknowledgement(fmt.Errorf("insufficient funds"))
ackBz, _ := failAck.Marshal()
err = s.middleware.OnAcknowledgementPacket(
s.ctx, packet, ackBz, sdk.AccAddress{},
)
s.NoError(err)
// 验证额度已回退
periodIndex := uint64(s.ctx.BlockTime().Unix()) / 86400
usage := s.keeper.GetPeriodUsage(s.ctx, "channel-0", periodIndex)
s.True(usage.IsZero())
}
// TestPacketTimeout_Refund 测试超时后退还额度
func (s *RateLimitMiddlewareTestSuite) TestPacketTimeout_Refund() {
limit := types.ChannelRateLimit{
ChannelID: "channel-0",
WindowDuration: 86400,
MaxOutflow: sdk.NewInt(100_000_000_000_000_000_000), // 100 MSG
Active: true,
}
s.keeper.SetChannelRateLimit(s.ctx, limit)
// 发送 50 MSG
packet := mock.NewMockPacket("transfer", "channel-0", "transfer", "channel-0", 8,
"{\"amount\":\"50000000000000000000\",\"denom\":\"umsg\",\"receiver\":\"msg1...\"}")
err := s.middleware.SendPacket(s.ctx, nil, packet)
s.NoError(err)
// 模拟超时
err = s.middleware.OnTimeoutPacket(s.ctx, packet, sdk.AccAddress{})
s.NoError(err)
// 验证额度已回退
periodIndex := uint64(s.ctx.BlockTime().Unix()) / 86400
usage := s.keeper.GetPeriodUsage(s.ctx, "channel-0", periodIndex)
s.True(usage.IsZero())
}
6.2 治理中间件测试
// X-MSG-Stub=true: gov_ibc_module_test.go - 治理中间件测试
package msggovibc_test
import (
"testing"
sdk "github.com/cosmos/cosmos-sdk/types"
channeltypes "github.com/cosmos/ibc-go/v8/modules/core/04-channel/types"
"github.com/cosmos/ibc-go/v8/modules/core/04-channel/types/mock"
"github.com/stretchr/testify/suite"
"msgchain/x/msggovibc"
"msgchain/x/msggovibc/keeper"
"msgchain/x/msggovibc/types"
)
type GovMiddlewareTestSuite struct {
suite.Suite
ctx sdk.Context
keeper keeper.Keeper
middleware msggovibc.GovIBCMiddleware
nextModule *mock.IBCModule
}
func (s *GovMiddlewareTestSuite) SetupTest() {
// 初始化
s.ctx = sdk.NewContext(nil, sdk.Header{}, false, nil)
s.nextModule = &mock.IBCModule{}
s.middleware = msggovibc.NewIBCMiddleware(
s.nextModule,
s.keeper,
nil,
nil,
)
}
// TestGovPause_SendBlocked 测试治理暂停后发送被阻止
func (s *GovMiddlewareTestSuite) TestGovPause_SendBlocked() {
// 暂停通道
s.keeper.PauseChannel(s.ctx, "channel-0")
packet := mock.NewMockPacket("transfer", "channel-0", "transfer", "channel-0", 1, "{}")
err := s.middleware.SendPacket(s.ctx, nil, packet)
s.Error(err)
s.ErrorIs(err, types.ErrChannelPaused)
}
// TestGovPause_Resume 测试暂停后恢复
func (s *GovMiddlewareTestSuite) TestGovPause_Resume() {
// 暂停
s.keeper.PauseChannel(s.ctx, "channel-0")
// 恢复
s.keeper.ResumeChannel(s.ctx, "channel-0")
packet := mock.NewMockPacket("transfer", "channel-0", "transfer", "channel-0", 2, "{}")
err := s.middleware.SendPacket(s.ctx, nil, packet)
s.NoError(err)
}
// TestEmergencyPause 测试紧急暂停
func (s *GovMiddlewareTestSuite) TestEmergencyPause() {
params := s.keeper.GetParams(s.ctx)
params.EmergencyPause = true
s.keeper.SetParams(s.ctx, params)
packet := mock.NewMockPacket("transfer", "channel-0", "transfer", "channel-0", 3, "{}")
err := s.middleware.SendPacket(s.ctx, nil, packet)
s.Error(err)
}
// TestWhitelistFilter 测试白名单过滤
func (s *GovMiddlewareTestSuite) TestWhitelistFilter() {
params := s.keeper.GetParams(s.ctx)
params.WhitelistMode = true
s.keeper.SetParams(s.ctx, params)
filter := types.ChannelFilter{
ChannelID: "channel-0",
Enabled: true,
AllowedPorts: []string{"transfer"},
}
s.keeper.SetChannelFilter(s.ctx, filter)
// 允许的端口
allowed := s.keeper.IsPacketAllowed(s.ctx, "channel-0", "transfer", "msg1...")
s.True(allowed)
// 不允许的端口
denied := s.keeper.IsPacketAllowed(s.ctx, "channel-0", "wasm.contract1", "msg1...")
s.False(denied)
}
// TestBlacklistFilter 测试黑名单过滤
func (s *GovMiddlewareTestSuite) TestBlacklistFilter() {
filter := types.ChannelFilter{
ChannelID: "channel-0",
Enabled: true,
BlockedPorts: []string{"wasm.suspicious"},
}
s.keeper.SetChannelFilter(s.ctx, filter)
// 黑名单中的端口
denied := s.keeper.IsPacketAllowed(s.ctx, "channel-0", "wasm.suspicious", "msg1...")
s.False(denied)
// 不在黑名单的端口
allowed := s.keeper.IsPacketAllowed(s.ctx, "channel-0", "transfer", "msg1...")
s.True(allowed)
}
// TestAddressPrefixFilter 测试地址前缀过滤
func (s *GovMiddlewareTestSuite) TestAddressPrefixFilter() {
filter := types.ChannelFilter{
ChannelID: "channel-0",
Enabled: true,
BlockedAddressPrefixes: []string{"msg1suspicious"},
}
s.keeper.SetChannelFilter(s.ctx, filter)
denied := s.keeper.IsPacketAllowed(
s.ctx, "channel-0", "transfer", "msg1suspiciousabc123",
)
s.False(denied)
allowed := s.keeper.IsPacketAllowed(
s.ctx, "channel-0", "transfer", "msg1normalkey",
)
s.True(allowed)
}
6.3 集成测试
使用 IBC 测试库进行双链集成测试:
// X-MSG-Stub=true: integration_test.go - 双链集成测试
package integration_test
import (
"testing"
"time"
sdk "github.com/cosmos/cosmos-sdk/types"
ibctesting "github.com/cosmos/ibc-go/v8/testing"
"github.com/stretchr/testify/suite"
"msgchain/x/msgratelimit/types"
)
type RateLimitIntegrationTestSuite struct {
suite.Suite
coordinator *ibctesting.Coordinator
chainA *ibctesting.TestChain
chainB *ibctesting.TestChain
path *ibctesting.Path
}
func (s *RateLimitIntegrationTestSuite) SetupTest() {
s.coordinator = ibctesting.NewCoordinator(s.T(), 2)
s.chainA = s.coordinator.GetChain(ibctesting.GetChainID(1))
s.chainB = s.coordinator.GetChain(ibctesting.GetChainID(2))
// 创建 IBC 路径
s.path = ibctesting.NewPath(s.chainA, s.chainB)
s.coordinator.SetupConnections(s.path)
// 创建通道
s.coordinator.CreateTransferChannels(s.path)
}
// TestRateLimitIntegration 集成测试:跨链转账受速率限制
func (s *RateLimitIntegrationTestSuite) TestRateLimitIntegration() {
// 在 Chain A 上设置速率限制:每小时最多 1000 MSG
limit := types.ChannelRateLimit{
ChannelID: s.path.EndpointA.ChannelID,
WindowDuration: 3600,
MaxOutflow: sdk.NewInt(1_000_000_000_000_000_000_000), // 1000 MSG
Active: true,
}
// 通过治理消息设置限制
// 注意:实际集成测试需要获取 Keeper 引用
// msgRateLimitKeeper := s.chainA.App.RateLimitKeeper
// msgRateLimitKeeper.SetChannelRateLimit(s.chainA.GetContext(), limit)
// 发送 700 MSG(应成功)
// msgTracker := ibctesting.NewMsgTracker()
// err := s.coordinator.SendMsgs(s.chainA, s.chainB, msgTracker, transferMsg)
// s.NoError(err)
// 再发送 400 MSG(应超出限制被拒绝)
// err = s.coordinator.SendMsgs(s.chainA, s.chainB, msgTracker, transferMsg2)
// s.Error(err)
}
// TestGovernanceIntegration 集成测试:治理暂停通道
func (s *RateLimitIntegrationTestSuite) TestGovernanceIntegration() {
// 通过治理提案暂停通道
// 1. 提交治理提案
// 2. 投票通过
// 3. 提案执行后通道暂停
// 4. 验证 IBC 转账被阻止
// 5. 提交恢复提案
// 6. 验证 IBC 转账恢复
}
// TestPacketAckIntegration 集成测试:ACK 处理
func (s *RateLimitIntegrationTestSuite) TestPacketAckIntegration() {
// 完整测试:发送 → 接收 → ACK → 验证额度
}
func TestRateLimitIntegrationSuite(t *testing.T) {
suite.Run(t, new(RateLimitIntegrationTestSuite))
}
6.4 Benchmark 测试
// X-MSG-Stub=true: benchmark_test.go
package msgratelimit_test
import (
"testing"
sdk "github.com/cosmos/cosmos-sdk/types"
"github.com/cosmos/ibc-go/v8/modules/core/04-channel/types/mock"
"msgchain/x/msgratelimit/types"
)
func BenchmarkSendPacket_NoLimit(b *testing.B) {
app := setupBenchmarkApp()
packet := mock.NewMockPacket("transfer", "channel-0", "transfer", "channel-0", 1, "{\"amount\":\"1000000\",\"denom\":\"umsg\"}")
b.ResetTimer()
for i := 0; i < b.N; i++ {
_ = app.middleware.SendPacket(app.ctx, nil, packet)
}
}
func BenchmarkSendPacket_WithLimit(b *testing.B) {
app := setupBenchmarkApp()
limit := types.ChannelRateLimit{
ChannelID: "channel-0",
WindowDuration: 86400,
MaxOutflow: sdk.NewInt(1_000_000_000_000_000_000_000),
Active: true,
}
app.keeper.SetChannelRateLimit(app.ctx, limit)
packet := mock.NewMockPacket("transfer", "channel-0", "transfer", "channel-0", 1, "{\"amount\":\"1000000\",\"denom\":\"umsg\"}")
b.ResetTimer()
for i := 0; i < b.N; i++ {
_ = app.middleware.SendPacket(app.ctx, nil, packet)
}
}
type benchmarkApp struct {
ctx sdk.Context
keeper keeper.Keeper
middleware msgratelimit.IBCMiddleware
}
func setupBenchmarkApp() *benchmarkApp {
// 初始化 benchmark 应用
return &benchmarkApp{}
}
7. 部署与配置
7.1 app.go 注册流程
在 MSG Chain 的 app.go 中注册 IBC 中间件栈:
// X-MSG-Stub=true: app.go - IBC 中间件注册
package app
import (
"github.com/cosmos/cosmos-sdk/runtime"
"github.com/cosmos/ibc-go/v8/modules/apps/transfer"
ibcfee "github.com/cosmos/ibc-go/v8/modules/apps/29-fee"
ibcporttypes "github.com/cosmos/ibc-go/v8/modules/core/05-port/types"
msgratelimit "msgchain/x/msgratelimit"
msgratelimitkeeper "msgchain/x/msgratelimit/keeper"
msgratelimittypes "msgchain/x/msgratelimit/types"
msggovibc "msgchain/x/msggovibc"
msggovibckeeper "msgchain/x/msggovibc/keeper"
msggovibctypes "msgchain/x/msggovibc/types"
)
// 在 MsgChainApp 结构体中新增 Keeper 字段
type MsgChainApp struct {
// ... 现有字段
// IBC 中间件 keepers
RateLimitKeeper msgratelimitkeeper.Keeper
GovIBCKeeper msggovibckeeper.Keeper
}
// InitIBCMiddleware 初始化并注册 IBC 中间件栈
func (app *MsgChainApp) InitIBCMiddleware() {
// 1. 创建 ICS-20 转账 IBC 应用
transferIBCModule := transfer.NewIBCModule(app.TransferKeeper)
// 2. 从内到外构建中间件栈
var ibcStack ibcporttypes.IBCModule
// 2a. 基础应用
ibcStack = transferIBCModule
// 2b. 速率限制中间件
ibcStack = msgratelimit.NewIBCMiddleware(
ibcStack,
app.RateLimitKeeper,
app.IBCKeeper.ChannelKeeper,
)
// 2c. 治理控制中间件
ibcStack = msggovibc.NewIBCMiddleware(
ibcStack,
app.GovIBCKeeper,
app.IBCKeeper.ChannelKeeper,
app.AccountKeeper,
)
// 2d. IBC Fee 中间件(最外层)
ibcStack = ibcfee.NewIBCMiddleware(
ibcStack,
app.IBCFeeKeeper,
)
// 3. 绑定端口
app.IBCKeeper.SetRouter(runtime.NewIBCRouter())
app.IBCKeeper.Router().AddRoute("transfer", ibcStack)
}
// 在 InitChainer 中调用
func (app *MsgChainApp) initKeepers() {
// ... 初始化现有模块 ...
// 初始化速率限制 keeper
app.RateLimitKeeper = msgratelimitkeeper.NewKeeper(
app.keys[msgratelimittypes.StoreKey],
app.appCodec,
app.GetSubspace(msgratelimittypes.ModuleName),
app.BankKeeper,
app.IBCKeeper.ChannelKeeper,
authtypes.NewModuleAddress(govtypes.ModuleName).String(),
)
// 初始化治理 IBC keeper
app.GovIBCKeeper = msggovibckeeper.NewKeeper(
app.keys[msggovibctypes.StoreKey],
app.appCodec,
app.GetSubspace(msggovibctypes.ModuleName),
app.IBCKeeper.ChannelKeeper,
authtypes.NewModuleAddress(govtypes.ModuleName).String(),
)
// 注册 IBC 中间件
app.InitIBCMiddleware()
}
// 注册模块管理器
func init() {
app.RegisterModule(
msgratelimittypes.ModuleName,
msggovibctypes.ModuleName,
)
}
7.2 Genesis 配置
// X-MSG-Stub=true: types/genesis.go
package types
import (
sdk "github.com/cosmos/cosmos-sdk/types"
)
// GenesisState 创世状态
type GenesisState struct {
Params Params `json:"params"`
ChannelRateLimits []ChannelRateLimit `json:"channel_rate_limits"`
}
func DefaultGenesis() *GenesisState {
return &GenesisState{
Params: DefaultParams(),
ChannelRateLimits: []ChannelRateLimit{},
}
}
func (gs GenesisState) Validate() error {
if err := gs.Params.Validate(); err != nil {
return err
}
for _, limit := range gs.ChannelRateLimits {
if limit.ChannelID == "" {
return fmt.Errorf("channel ID cannot be empty")
}
if limit.MaxOutflow.IsNegative() {
return fmt.Errorf("max outflow cannot be negative")
}
}
return nil
}
7.3 Genesis JSON 示例
{
"app_state": {
"msgratelimit": {
"params": {
"default_window_duration": 86400,
"default_max_outflow": "1000000000000000000000",
"enabled": true
},
"channel_rate_limits": [
{
"channel_id": "channel-0",
"window_duration": 86400,
"max_outflow": "500000000000000000000",
"active": true
},
{
"channel_id": "channel-1",
"window_duration": 3600,
"max_outflow": "100000000000000000000",
"active": true
}
]
},
"msggovibc": {
"params": {
"enabled": true,
"emergency_pause": false,
"whitelist_mode": false
},
"channel_filters": [],
"paused_channels": []
}
}
}
7.4 app.toml 配置
# app.toml - IBC 中间件配置
###############################################################################
### IBC 速率限制配置 ###
###############################################################################
[ibc-rate-limit]
# 启用速率限制
enabled = true
# 默认时间窗口(秒)
default-window-duration = 86400
# 默认最大流出量 (umsg)
# 1000 MSG 默认限制(18位小数)
default-max-outflow = "1000000000000000000000"
7.5 CLI 命令
# X-MSG-Stub=true: 速率限制 CLI
# 设置通道速率限制
msgd tx msgratelimit set-rate-limit \
--channel-id channel-0 \
--window-duration 86400 \
--max-outflow 1000000000000000000000000 \
--from gov-module-account \
--gas-prices 1000000000umsg
# 暂停通道
msgd tx msggovibc pause-channel channel-0 \
--from gov-module-account \
--gas-prices 1000000000umsg
# 恢复通道
msgd tx msggovibc resume-channel channel-0 \
--from gov-module-account \
--gas-prices 1000000000umsg
# 查询通道速率限制
msgd query msgratelimit channel-rate-limit channel-0
# 查询所有通道限制
msgd query msgratelimit all-channel-rate-limits
# 查询通道状态(含是否暂停)
msgd query msggovibc channel-status channel-0
7.6 治理提案提交
# X-MSG-Stub=true: 通过治理提案更新 IBC 中间件参数
# 提交设置速率限制的治理提案
msgd tx gov submit-proposal \
--title "Set IBC Rate Limit for channel-0" \
--description "Set daily rate limit of 10000 MSG for channel-0 (Cosmos Hub)" \
--type Text \
--deposit 100000000000000000000umsg \
--proposer msg1gov... \
-- \
--action set_channel_rate_limit \
--channel-id channel-0 \
--window-duration 86400 \
--max-outflow 10000000000000000000000umsg
# 提交暂停通道的治理提案
msgd tx gov submit-proposal \
--title "Emergency Pause channel-0" \
--description "Pause IBC channel-0 due to security incident" \
--type Text \
--deposit 100000000000000000000umsg \
--proposer msg1gov... \
-- \
--action pause_channel \
--channel-id channel-0
# 提交紧急暂停所有 IBC 流量的提案
msgd tx gov submit-proposal \
--title "Emergency Pause All IBC" \
--description "Pause all IBC traffic on MSG Chain" \
--type Text \
--deposit 100000000000000000000umsg \
--proposer msg1gov... \
-- \
--action emergency_pause
# 投票
msgd tx gov vote 1 Yes --from validator-key --gas-prices 1000000000umsg
7.7 升级现有链
如果 MSG Chain 已经运行,需要通过链升级提案添加中间件:
// X-MSG-Stub=true: upgrade_handler.go - 链升级处理器
package app
import (
sdk "github.com/cosmos/cosmos-sdk/types"
"github.com/cosmos/cosmos-sdk/types/module"
upgradekeeper "github.com/cosmos/cosmos-sdk/x/upgrade/keeper"
msgratelimittypes "msgchain/x/msgratelimit/types"
msggovibctypes "msgchain/x/msggovibc/types"
)
func (app *MsgChainApp) RegisterUpgradeHandlers() {
// v2.0.0 升级:添加 IBC 速率限制
app.UpgradeKeeper.SetUpgradeHandler("v2.0.0-ibc-rate-limit", func(
ctx sdk.Context,
plan upgrade.Plan,
fromVM module.VersionMap,
) (module.VersionMap, error) {
// 1. 初始化新模块的创世状态
app.RateLimitKeeper.InitGenesis(ctx, *msgratelimittypes.DefaultGenesis())
// 2. 重新注册 IBC 中间件栈
app.InitIBCMiddleware()
// 3. 运行模块迁移
return app.ModuleManager.RunMigrations(ctx, fromVM)
})
// v2.1.0 升级:添加治理 IBC 控制
app.UpgradeKeeper.SetUpgradeHandler("v2.1.0-gov-ibc-control", func(
ctx sdk.Context,
plan upgrade.Plan,
fromVM module.VersionMap,
) (module.VersionMap, error) {
app.GovIBCKeeper.InitGenesis(ctx, *msggovibctypes.DefaultGenesis())
app.InitIBCMiddleware()
return app.ModuleManager.RunMigrations(ctx, fromVM)
})
}
// 升级命令
// X-MSG-Stub=true:
// msgd tx upgrade software-upgrade v2.0.0-ibc-rate-limit \
// --upgrade-height 1000000 \
// --upgrade-info '{"description":"Add IBC rate limiting middleware"}' \
// --from validator-key \
// --gas-prices 1000000000umsg
7.8 Store 迁移
// X-MSG-Stub=true: keeper/migration.go
package keeper
import (
sdk "github.com/cosmos/cosmos-sdk/types"
)
// Migrate1to2 从 v1 迁移到 v2
func (k Keeper) Migrate1to2(ctx sdk.Context) error {
// 从旧存储格式迁移到新格式
store := ctx.KVStore(k.storeKey)
// 迁移所有通道限制
iterator := sdk.KVStorePrefixIterator(store, []byte{0x00}) // 旧前缀
defer iterator.Close()
for ; iterator.Valid(); iterator.Next() {
key := iterator.Key()
value := iterator.Value()
// 旧格式:0x00 | channelID → 新格式:0x01 | channelID
channelID := string(key[1:])
newKey := append([]byte{0x01}, []byte(channelID)...)
store.Set(newKey, value)
store.Delete(key)
}
return nil
}
8. 完整示例
8.1 部署速率限制中间件
以下是从零开始部署速率限制中间件的完整流程:
步骤 1: 创建模块目录
# 在 MSG Chain 代码库中创建模块
mkdir -p x/msgratelimit/{module,keeper,types,client/cli}
步骤 2: 注册 AppModule
// X-MSG-Stub=true: app.go 完整注册
package app
import (
// ... 其他 import
msgratelimitmodule "msgchain/x/msgratelimit/module"
msgratelimitkeeper "msgchain/x/msgratelimit/keeper"
msgratelimittypes "msgchain/x/msgratelimit/types"
)
func NewMsgChainApp(...) *MsgChainApp {
// ... 现有初始化 ...
// 注册模块管理器
app.ModuleManager = module.NewManager(
// ... 现有模块 ...
msgratelimitmodule.NewAppModule(app.RateLimitKeeper),
)
// 设置 order
app.ModuleManager.SetOrderPreBeginBlockers(
// ... 现有 ...
msgratelimittypes.ModuleName,
)
app.ModuleManager.SetOrderBeginBlockers(
// ... 现有 ...
msgratelimittypes.ModuleName,
)
app.ModuleManager.SetOrderEndBlockers(
// ... 现有 ...
msgratelimittypes.ModuleName,
)
app.ModuleManager.SetOrderInitGenesis(
// ... 现有 ...
msgratelimittypes.ModuleName,
)
// 初始化 IBC 中间件
app.InitIBCMiddleware()
return app
}
步骤 3: 配置通道限制
创世时配置或通过治理提交后配置:
# X-MSG-Stub=true: 配置通道限制
# 方案 A: 通过治理提案设置
msgd tx gov submit-legacy-proposal \
--title "Configure IBC Rate Limits" \
--description "Set rate limits for channel-0 (Cosmos Hub): 10000 MSG/day, channel-1 (Osmosis): 5000 MSG/day" \
--deposit 100000000000000000000umsg \
--from gov-account \
--type Text
# 方案 B: 通过专属消息(如果模块有权限设置)
msgd tx msgratelimit set-channel-rate-limit \
--channel-id channel-0 \
--window-duration 86400 \
--max-outflow 10000000000000000000000000 \
--from gov-module-account \
--gas-prices 1000000000umsg
msgd tx msgratelimit set-channel-rate-limit \
--channel-id channel-1 \
--window-duration 3600 \
--max-outflow 5000000000000000000000000 \
--from gov-module-account \
--gas-prices 1000000000umsg
步骤 4: 验证配置
# X-MSG-Stub=true: 查询所有通道限制
msgd query msgratelimit all-channel-rate-limits
# 预期输出:
# channel_rate_limits:
# - channel_id: channel-0
# window_duration: 86400
# max_outflow: "10000000000000000000000000"
# active: true
# - channel_id: channel-1
# window_duration: 3600
# max_outflow: "5000000000000000000000000"
# active: true
# 查询具体通道
msgd query msgratelimit channel-rate-limit channel-0
# 预期输出:
# channel_rate_limit:
# channel_id: channel-0
# window_duration: 86400
# max_outflow: "10000000000000000000000000"
# active: true
# current_usage: "0"
# is_paused: false
步骤 5: 测试跨链转账
# X-MSG-Stub=true: 测试跨链转账(带速率限制)
# 发送 500 MSG(在限制内)
msgd tx ibc-transfer transfer \
transfer channel-0 \
cosmos1recipient... \
500000000000000000000umsg \
--packet-timeout-height "0-100000" \
--from alice \
--gas-prices 1000000000umsg
# 预期: 成功
# code: 0
# 查询当前用量
msgd query msgratelimit channel-rate-limit channel-0
# 预期输出:
# current_usage: "500000000000000000000"
# 尝试发送 2000 MSG(超过日限 10000 - 500 = 9500,所以在限制内)
msgd tx ibc-transfer transfer \
transfer channel-0 \
cosmos1recipient... \
2000000000000000000000umsg \
--packet-timeout-height "0-100000" \
--from alice \
--gas-prices 1000000000umsg
# 预期: 成功
# 尝试发送 10000 MSG(超过剩余 7500)
msgd tx ibc-transfer transfer \
transfer channel-0 \
cosmos1recipient... \
10000000000000000000000umsg \
--packet-timeout-height "0-100000" \
--from alice \
--gas-prices 1000000000umsg
# 预期: 失败
# code: 13
# raw_log: "rate limit exceeded for channel: channel-0"
步骤 6: 更新限制通过治理
# X-MSG-Stub=true: 更新速率限制参数
# 提交修改限制的治理提案
msgd tx gov submit-proposal \
--title "Increase rate limit for channel-0" \
--description "Increase daily rate limit from 10000 MSG to 50000 MSG for Cosmos Hub channel" \
--deposit 100000000000000000000umsg \
--from gov-account \
--type Text
# 投票(需要达到通过门槛)
msgd tx gov vote 1 Yes --from validator-1 --gas-prices 1000000000umsg
msgd tx gov vote 1 Yes --from validator-2 --gas-prices 1000000000umsg
msgd tx gov vote 1 Yes --from validator-3 --gas-prices 1000000000umsg
# 提案通过后,通过治理消息更新限制
msgd tx gov submit-proposal \
msgratelimit/SetChannelRateLimit \
'{"channel_id":"channel-0","window_duration":86400,"max_outflow":"50000000000000000000000000"}' \
--title "Set channel-0 rate limit" \
--description "Set channel-0 rate limit to 50000 MSG/day" \
--deposit 100000000000000000000umsg \
--from gov-account \
--gas-prices 1000000000umsg
# 验证更新
msgd query msgratelimit channel-rate-limit channel-0
8.2 紧急暂停通道
# X-MSG-Stub=true: 安全事件响应
# 场景:检测到异常转账模式,需要紧急暂停 channel-0
# 步骤 1: 提交紧急暂停提案
msgd tx gov submit-proposal \
--title "Emergency Pause IBC Channel-0" \
--description "Suspicious activity detected on Cosmos Hub channel. Pausing all transfers." \
--deposit 500000000000000000000umsg \
--from validator-1 \
--type Text
# 步骤 2: 快速投票(所有验证者)
msgd tx gov vote 2 Yes --from validator-1 --gas-prices 1000000000umsg
msgd tx gov vote 2 Yes --from validator-2 --gas-prices 1000000000umsg
msgd tx gov vote 2 Yes --from validator-3 --gas-prices 1000000000umsg
# 步骤 3: 验证通道已暂停
msgd query msggovibc channel-status channel-0
# 预期: paused: true
# 步骤 4: 尝试转账(应被拒绝)
msgd tx ibc-transfer transfer \
transfer channel-0 \
cosmos1recipient... \
100umsg \
--packet-timeout-height "0-100000" \
--from alice \
--gas-prices 1000000000umsg
# 预期: 失败 - channel is paused by governance
# 步骤 5: 调查完成后恢复
msgd tx gov submit-proposal \
--title "Resume IBC Channel-0" \
--description "Investigation complete. Resume normal IBC operations on channel-0." \
--deposit 500000000000000000000umsg \
--from validator-1 \
--type Text
# 步骤 6: 投票通过后恢复
msgd query msggovibc channel-status channel-0
# 预期: paused: false
8.3 完整端到端测试脚本
#!/bin/bash
# X-MSG-Stub=true: end_to_end_test.sh - 端到端测试
set -euo pipefail
echo "=== MSG Chain IBC 中间件端到端测试 ==="
echo ""
# 配置
MSG_CHAIN_ID="msg-chain-1"
COSMOS_CHAIN_ID="cosmoshub-4"
CHANNEL_ID="channel-0"
GAS_PRICES="1000000000umsg"
# 1. 检查模块状态
echo "1. 检查 IBC 中间件模块..."
msgd query msgratelimit params
msgd query msggovibc params
# 2. 配置速率限制
echo ""
echo "2. 配置速率限制: 10000 MSG/天..."
msgd tx gov submit-proposal \
--title "Configure rate limit for $CHANNEL_ID" \
--description "Set daily rate limit of 10000 MSG" \
--deposit 100000000000000000000umsg \
--from validator \
--type Text
# 简化:直接通过治理消息
echo "提交治理提案设置速率限制..."
RESULT=$(msgd tx gov submit-proposal \
msgratelimit/SetChannelRateLimit \
"{\"channel_id\":\"$CHANNEL_ID\",\"window_duration\":86400,\"max_outflow\":\"10000000000000000000000000\"}" \
--title "Set rate limit for $CHANNEL_ID" \
--description "Set daily rate limit to 10000 MSG for Cosmos Hub" \
--deposit 100000000000000000000umsg \
--from validator \
--gas-prices $GAS_PRICES \
--output json 2>&1)
echo "提案提交结果: $RESULT"
# 3. 查询配置
echo ""
echo "3. 验证配置..."
msgd query msgratelimit channel-rate-limit $CHANNEL_ID
# 4. 执行跨链转账(限制内)
echo ""
echo "4. 发送 1000 MSG(限制内)..."
ALICE=$(msgd keys show alice -a)
TX_RESULT=$(msgd tx ibc-transfer transfer \
transfer $CHANNEL_ID \
cosmos1recipient... \
1000000000000000000000umsg \
--packet-timeout-height "0-200000" \
--from alice \
--gas-prices $GAS_PRICES \
--output json 2>&1)
echo "转账结果: $TX_RESULT"
# 5. 检查用量
echo ""
echo "5. 检查当前用量..."
msgd query msgratelimit channel-rate-limit $CHANNEL_ID
# 6. 执行超限转账(应失败)
echo ""
echo "6. 尝试发送 20000 MSG(超限)..."
TX_RESULT2=$(msgd tx ibc-transfer transfer \
transfer $CHANNEL_ID \
cosmos1recipient... \
20000000000000000000000000umsg \
--packet-timeout-height "0-200000" \
--from alice \
--gas-prices $GAS_PRICES \
--output json 2>&1 || true)
echo "超限转账结果: $TX_RESULT2"
echo "预期: 失败 - rate limit exceeded"
# 7. 测试治理暂停
echo ""
echo "7. 测试治理暂停通道..."
msgd tx gov submit-proposal \
msggovibc/PauseChannel \
"{\"channel_id\":\"$CHANNEL_ID\"}" \
--title "Pause $CHANNEL_ID" \
--description "Test: pause channel for maintenance" \
--deposit 100000000000000000000umsg \
--from validator \
--gas-prices $GAS_PRICES
# 投票
msgd tx gov vote 3 Yes --from validator --gas-prices $GAS_PRICES
# 8. 验证暂停
echo ""
echo "8. 验证通道暂停..."
msgd query msggovibc channel-status $CHANNEL_ID
# 9. 测试暂停后转账(应失败)
echo ""
echo "9. 尝试暂停后转账(应被阻止)..."
TX_RESULT3=$(msgd tx ibc-transfer transfer \
transfer $CHANNEL_ID \
cosmos1recipient... \
100umsg \
--packet-timeout-height "0-200000" \
--from alice \
--gas-prices $GAS_PRICES \
--output json 2>&1 || true)
echo "暂停后转账结果: $TX_RESULT3"
echo "预期: 失败 - channel is paused by governance"
# 10. 恢复通道
echo ""
echo "10. 恢复通道..."
msgd tx gov submit-proposal \
msggovibc/ResumeChannel \
"{\"channel_id\":\"$CHANNEL_ID\"}" \
--title "Resume $CHANNEL_ID" \
--description "Test: resume channel after maintenance" \
--deposit 100000000000000000000umsg \
--from validator \
--gas-prices $GAS_PRICES
# 11. 验证恢复
echo ""
echo "11. 验证通道恢复..."
msgd query msggovibc channel-status $CHANNEL_ID
# 12. 恢复后转账
echo ""
echo "12. 恢复后再次转账..."
TX_RESULT4=$(msgd tx ibc-transfer transfer \
transfer $CHANNEL_ID \
cosmos1recipient... \
100umsg \
--packet-timeout-height "0-200000" \
--from alice \
--gas-prices $GAS_PRICES \
--output json 2>&1 || true)
echo "恢复后转账结果: $TX_RESULT4"
echo "预期: 成功"
echo ""
echo "=== 测试完成 ==="
8.4 完整代码文件清单
完成中间件开发后,需要创建以下文件:
x/msgratelimit/
├── module/
│ └── module.go # AppModule/AppModuleBasic
├── keeper/
│ ├── keeper.go # 核心 Keeper
│ ├── grpc_query.go # gRPC 查询
│ ├── msg_server.go # 消息处理器
│ ├── params.go # 参数管理
│ └── migration.go # 存储迁移
├── types/
│ ├── errors.go # 错误定义
│ ├── keys.go # 存储键
│ ├── params.go # 参数类型
│ ├── messages.go # 消息类型
│ ├── genesis.go # 创世状态
│ ├── codec.go # 编码注册
│ ├── expected_keepers.go # 依赖接口
│ ├── query.pb.go # Protobuf 查询(生成)
│ └── tx.pb.go # Protobuf 交易(生成)
├── client/cli/
│ ├── tx.go # CLI 交易命令
│ └── query.go # CLI 查询命令
├── ibc_module.go # IBC 中间件实现
├── ibc_module_test.go # 单元测试
├── keeper_test.go # Keeper 测试
└── simulation.go # 模拟测试
x/msggovibc/
├── module/
│ └── module.go
├── keeper/
│ ├── keeper.go
│ ├── grpc_query.go
│ ├── msg_server.go
│ └── proposal_handler.go
├── types/
│ ├── errors.go
│ ├── keys.go
│ ├── params.go
│ ├── messages.go
│ └── genesis.go
├── client/cli/
│ └── tx.go
├── ibc_module.go
├── ibc_module_test.go
└── keeper_test.go
8.5 依赖和版本信息
| 依赖 | 版本 | 说明 |
|---|---|---|
| Cosmos SDK | v0.50.x | MSG Chain 基础框架 |
| IBC-Go | v8.x | IBC 核心和中间件标准 |
| CometBFT | v0.38.x | 共识引擎 |
| Go | 1.21+ | 开发语言 |
| Protobuf | v3 | 接口定义 |
8.6 常见问题
Q: 中间件如何影响数据包延迟?
A: 每个中间件在发送路径上增加一次 KV 存储读取和写入。速率限制中间件单次检查约消耗 50-100μs,对区块时间影响极小。
Q: 多个中间件的执行顺序如何保证?
A: 通过中间件栈的嵌套结构保证。外层中间件先执行前置逻辑,后执行后置逻辑。编写中间件时只需要信任传递给的 next 对象。
Q: 如何调试中间件问题?
A: 打开中间件的 debug 日志(--log_level info),所有 SendPacket 拦截点会记录决策日志。也可以添加自定义事件(ctx.EventManager().EmitEvent())。
Q: 速率限制的精度如何?
A: MSG 精度为 18 位小数,所有计算使用 sdk.Int 进行精确整数运算,无精度损失。
Q: 窗口重置后,之前周期的未使用额度能否累积?
A: 不累积。每个周期独立计算,结束时归零。这是为了防止用户长期累积后一次性大量转出。
Q: 治理暂停和速率限制暂停有什么区别?
A: 治理暂停由治理提案触发,显式地暂停/恢复通道。速率限制暂停是限制超出时自动触发的拒绝行为。两者的暂停状态独立存储,可以叠加生效。
注意: 本文档中标记
X-MSG-Stub=true的代码和配置内容均为规划中的设计。MSG Chain 目前 IBC 状态为规划态,尚未在生产环境实现。本文档为开发者提前了解和规划 IBC 中间件开发提供参考。实际实现可能因 IBC-Go 版本更新、Cosmos SDK 版本变化或 MSG Chain 具体实现策略而有所调整。
