dApp Docs/IBC自定义中间件开发指南
Development reference. Not independently verified for production.

MSG Chain IBC 自定义中间件开发指南

数据来源:MSG Chain 代码库核实

主网状态: No-Go — 当前 MSGChain 主网裁决为 No-Go,以下内容反映代码实际状态,不代表生产可用。


目录

  1. 概述
  2. IBC中间件架构
  3. Go中间件实现
  4. 速率限制中间件
  5. 治理中间件
  6. 中间件测试
  7. 部署与配置
  8. 完整示例

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 设计目标

速率限制中间件的核心功能:

  1. 按通道限制:每个 IBC 通道可以独立配置流出速率
  2. 时间窗口:基于滑动时间窗口(如每小时、每天、每周)累计流出量
  3. 治理可调:所有限制参数可通过治理提案更新
  4. 不阻塞接收:只限制发送方向,不影响接收方向

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, &params)
    return params
}

// SetParams 设置全局参数
func (k Keeper) SetParams(ctx sdk.Context, params types.Params) {
    k.paramSpace.SetParamSet(ctx, &params)
}

// 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 设计目标

治理控制中间件的核心功能:

  1. 治理暂停:通过治理提案暂停指定 IBC 通道的数据包传递
  2. 数据包过滤:通过治理设置数据包过滤规则(白名单/黑名单)
  3. 参数管理:通过治理更新中间件参数
  4. 紧急刹车:快速暂停所有 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, &params)
    return params
}

// SetParams 设置参数
func (k Keeper) SetParams(ctx sdk.Context, params types.Params) {
    k.paramSpace.SetParamSet(ctx, &params)
}

// 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 具体实现策略而有所调整。