dApp Docs/数据可用性层(DA)集成指南
Development reference. Not independently verified for production.

数据可用性层(DA)集成指南

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

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


目录

  1. 数据可用性概念
    • 1.1 DA 在区块链架构中的位置
    • 1.2 DA 与存储的区别
    • 1.3 Celestia / Avail / EigenDA 生态概览
  2. MSG Chain 的 DA 策略
    • 2.1 链上存储 vs 链下 DA
    • 2.2 注册中心的 metadata 存储模式
  3. DA 层集成架构
    • 3.1 主权 rollup 模式
    • 3.2 validium 模式
    • 3.3 volition 模式
  4. 数据发布与验证
    • 4.1 blob 提交
    • 4.2 数据采样
    • 4.3 KZG 承诺
    • 4.4 欺诈证明
  5. CosmWasm 合约与 DA 交互
    • 5.1 contract 将数据发布到 DA
    • 5.2 从 DA 读取数据验证
  6. 轻客户端验证
    • 6.1 轻节点 DA 采样
    • 6.2 无需全节点即可验证数据可用性
  7. 索引器与 DA
    • 7.1 索引器从 DA 层读取数据
    • 7.2 归档存储
    • 7.3 与 MSG Chain 事件流的整合
  8. 集成实践
    • 8.1 使用 Celestia 作为 MSG Chain 的 DA 层
    • 8.2 blob 生命周期示例
  9. 安全与经济
    • 9.1 DA 层信任假设
    • 9.2 费用模型
    • 9.3 数据保留策略
  10. 总结与路线图

1. 数据可用性概念

1.1 DA 在区块链架构中的位置

数据可用性(Data Availability,简称 DA)是区块链基础设施中的核心层,负责确保区块中所有交易数据在网络中对所有参与者可访问和可验证。在传统的单体区块链架构中,DA 由主链原生提供——每个全节点下载并存储完整的区块数据,因此天然保证数据可用。然而,随着模块化区块链架构的兴起,DA 层从执行层和结算层中分离出来,成为一个独立的服务层。

在模块化区块链的堆栈中,DA 层位于以下位置:

+---------------------------------------------------+
|                   应用层 / dApp                      |
+---------------------------------------------------+
|                   执行层 (Execution)                  |
|              CosmWasm 虚拟机 / EVM                   |
+---------------------------------------------------+
|                   结算层 (Settlement)                 |
|              MSG Chain 主链 / 共识                    |
+---------------------------------------------------+
|               数据可用性层 (Data Availability)          |
|          Celestia / Avail / EigenDA                   |
+---------------------------------------------------+
|                   共识层 / 网络层                      |
|            Tendermint / CometBFT / libp2p             |
+---------------------------------------------------+

在 MSG Chain 的生态系统中,DA 层承担的关键职能包括:

MSG Chain 采用基于主权 Rollup(Sovereign Rollup)的架构设计,这意味着执行层和 DA 层之间存在明确的职责分离。DA 层仅负责数据的发布和可用性证明,而不参与交易的解释和状态转换。这种架构允许 MSG Chain 利用专用 DA 网络(如 Celestia)的高吞吐量和低成本优势,同时保持自身在结算和共识层面的主权。

1.2 DA 与存储的区别

数据可用性与数据存储在区块链架构中经常被混淆,但两者在 MSG Chain 的设计中代表截然不同的概念:

维度 数据可用性 (DA) 数据存储 (Storage)
核心目标 确保数据在特定时间点可被网络验证 确保数据长期持久化保存
时间维度 短至中期(区块确认窗口期) 长期(历史归档)
冗余机制 纠删码 + 数据采样 多副本复制
验证方式 KZG 承诺 + 欺诈证明 哈希校验 + 完整性检查
访问模式 按需采样验证 随机读写查询
典型协议 Celestia、Avail、EigenDA IPFS、Arweave、Filecoin
在 MSG Chain 中的角色 交易数据临时托管 + 可用性证明 合约状态历史 + dApp 内容

关键区别理解:

DA 层回答的问题是:"这些数据是否在区块时间内确实存在于网络中?" 存储层回答的问题是:"这些数据是否可以在任意未来时间点被检索?"

在 MSG Chain 的设计中,DA 层用于确保 Rollup 批次数据在执行后被网络确认为可用——这通常是数小时到数天的时间窗口。一旦数据在 MSG Chain 主链上通过状态承诺(State Commitment)被最终确认,DA 层的原始 blob 数据即可被安全地丢弃或归档到永续存储层。数据持久化存储则由专门的存储协议处理,这些协议通常与 DA 层解耦。

例如,一个发布在 MSG Chain 上的 NFT 项目可能使用以下分层数据策略:

  1. 交易数据: 通过 Celestia DA 层发布,仅在 Rollup 挑战期内保证可用。
  2. NFT 元数据: 存储在 IPFS 或 Arweave 中,保证长期可访问。
  3. 合约状态: 由 MSG Chain 全节点维护,通过 Tendermint 共识保证一致性。
  4. 链上事件: 由索引器(如 CosmWasm 索引器)归档到关系数据库或对象存储。

1.3 Celestia / Avail / EigenDA 生态概览

当前主流的 DA 层解决方案包括 Celestia、Avail 和 EigenDA,三者在架构设计、共识机制和生态集成方面各有特点。

1.3.1 Celestia

Celestia 是最早提出模块化区块链和 DA 分离的项目,其核心架构基于 Tendermint 共识和 Namespaced Merkle Trees(NMT)。

关键技术特性:

对 MSG Chain 的价值:

Celestia 的命名空间机制与 MSG Chain 的多合约、多链架构天然匹配。每个部署在 MSG Chain 上的 CosmWasm 合约可以被分配一个独立的 Celestia 命名空间,实现数据隔离。此外,Celestia 的 Tendermint 共识与 MSG Chain 使用的 CometBFT 同源,降低了集成复杂度。

网络参数:

参数 主网值 说明
区块时间 ~12 秒 与 CometBFT 接近
最大区块大小 8 MB 可治理调整
命名空间 ID 8 字节 支持 2^64 命名空间
DAS 采样次数 20 次 99.9% 可用性保证
验证人数量 100 DPoS 上限

1.3.2 Avail

Avail 最初是 Polygon 团队孵化的 DA 项目,现已成为独立网络。其架构核心是 Kate 承诺(KZG Commitments)结合数据可用性采样。

关键技术特性:

对 MSG Chain 的价值:

Avail 的 KZG 承诺机制提供了更强的隐私保护能力——CosmWasm 合约的数据在发布到 DA 层时可以保持加密状态,只有授权方可以解密。同时,KZG 承诺的恒定大小证明(无论数据量多大,证明始终为 48 字节)在轻客户端场景中具有显著优势。

网络参数:

参数 主网值 说明
区块时间 ~20 秒 针对 DA 优化
区块大小上限 动态调整 基于费用市场
App ID 20 字节 Ethereum 地址格式
KZG 证明大小 48 字节 BLS12-381 曲线
数据保留期 30 天 之后仅归档节点保留

1.3.3 EigenDA

EigenDA 是 EigenLayer 生态中的 DA 解决方案,利用以太坊的再质押(Restaking)安全模型。EigenDA 本身不运行独立的共识网络,而是由 EigenLayer 的再质押验证者网络提供数据可用性保证。

关键技术特性:

对 MSG Chain 的价值:

EigenDA 适合需要在 MSG Chain 和以太坊生态之间建立 DA 桥接的场景。通过 EigenDA,MSG Chain 上的 CosmWasm 合约可以利用以太坊的安全性来保证数据可用性,这对于跨链应用和混合架构尤其有价值。

网络参数:

参数 主网值 说明
吞吐量 10+ MB/s 可水平扩展
确认时间 ~1-2 秒 异步确认模式
验证者数量 动态 基于再质押总量
Slashing 条件 数据扣留 + 双重签名 经济安全保证
费用代币 ETH 通过 EigenLayer 结算

1.3.4 生态对比与选择建议

在 MSG Chain 的上下文中,选择合适的 DA 层需要综合考虑以下因素:

评估维度 Celestia Avail EigenDA
共识独立性 高(独立共识) 高(独立共识) 中(依赖以太坊)
与 CosmWasm 兼容性 高(Tendermint 同源) 中(需适配层) 中(需适配层)
数据采样效率 高(NMT) 高(KZG) 中(Disperser 架构)
经济成本 低 低 中(含再质押成本)
去中心化程度 高 高 中
主网上线状态 已主网上线 已主网上线 已主网上线

推荐策略:


2. MSG Chain 的 DA 策略

2.1 链上存储 vs 链下 DA

MSG Chain 在设计 DA 策略时面临核心权衡:是将所有数据存储在链上(由 CometBFT 共识维护),还是将数据可用性卸载到专用 DA 层。这两种策略在 MSG Chain 的架构中各有适用场景。

2.1.1 链上存储模型

在链上存储模型中,MSG Chain 的全节点存储所有交易数据,数据可用性由链自身的共识协议保证。

优势:

劣势:

在 MSG Chain 中的适用场景:

链上存储适用于 MSG Chain 上的核心资产和关键状态数据,这些数据具有以下特征:

典型示例——使用 cosmwasm-storage 的链上数据:

pub fn execute_update_balance(
    deps: DepsMut,
    info: MessageInfo,
    recipient: Addr,
    amount: Uint128,
) -> Result<Response, ContractError> {
    BALANCES.update(deps.storage, &recipient, |bal| -> StdResult<_> {
        Ok(bal.unwrap_or_default() + amount)
    })?;

    Ok(Response::new()
        .add_attribute("action", "update_balance")
        .add_attribute("recipient", recipient)
        .add_attribute("amount", amount))
}

2.1.2 链下 DA 模型

链下 DA 模型将数据发布到专用 DA 层(如 Celestia),MSG Chain 仅存储数据承诺(Data Commitment)和状态根(State Root)。

优势:

劣势:

在 MSG Chain 中的适用场景:

链下 DA 适用于 MSG Chain 上的数据密集型应用,包括:

典型示例——发布交易数据到 Celestia DA 层:

func PublishBatchToDA(
    ctx context.Context,
    daClient *celestia.Client,
    batchData []byte,
    namespaceID [8]byte,
    fee gas.Gas,
) (*celestia.Blob, error) {
    blob, err := celestia.NewBlob(namespaceID, batchData)
    if err != nil {
        return nil, fmt.Errorf("failed to create blob: %w", err)
    }

    height, err := daClient.SubmitBlob(ctx, []*celestia.Blob{blob}, fee)
    if err != nil {
        return nil, fmt.Errorf("failed to submit blob: %w", err)
    }

    commitment := sha256.Sum256(batchData)
    log.Printf("batch published to Celestia at height %d, commitment: %x", height, commitment)

    return blob, nil
}

2.1.3 混合存储策略

MSG Chain 采用混合存储策略,结合链上和链下 DA 的最佳特性。核心原则是:关键状态上链,批量数据入 DA。

混合策略的分层数据流:

用户交易
    |
    +---> 关键操作(转账、质押、治理投票)
    |         +---> MSG Chain 链上存储(CometBFT)
    |
    +---> 数据密集型操作(NFT 批量 mint、社交帖子、游戏动作)
            +---> DA 层存储(Celestia blob)
            |         +---> 全节点可选存储
            +---> MSG Chain 存储承诺(32 字节哈希)
                    +---> 全节点必须存储

混合策略的状态承诺树:

MSG Chain 区块
    |
    +-- 链上状态根 (State Root)
    |       +-- Cosmos SDK IAVL Tree
    |               +-- bank 模块余额
    |               +-- staking 模块委托
    |               +-- wasm 合约状态
    |               +-- ...
    |
    +-- DA 承诺根 (DA Commitment Root)
            +-- 稀疏 Merkle Tree
                    +-- namespace_1 -> height: 12345, commitment: 0xabcd...
                    +-- namespace_2 -> height: 12346, commitment: 0xef01...
                    +-- ...

这种混合策略允许 MSG Chain 在保持核心安全属性的同时,支持高吞吐量的数据密集型应用。

2.2 注册中心的 metadata 存储模式

MSG Chain 引入了一个创新的元数据存储模式——DA 注册中心(DA Registry),它是一个运行在 CosmWasm 虚拟机上的智能合约,负责管理链上数据与链下 DA 层数据之间的映射关系。

2.2.1 注册中心的核心功能

DA 注册中心合约的核心职责包括:

  1. 数据锚定: 记录每个 DA 发布事件与 MSG Chain 交易之间的关联。
  2. 元数据索引: 存储 blob 的描述信息,包括命名空间、数据格式、关联合约等。
  3. 验证路由: 向轻客户端和索引器提供验证 DA 数据所需的证明信息。
  4. 生命周期管理: 跟踪数据的过期时间、归档状态和访问权限。

注册中心合约接口:

#[cw_serde]
pub struct BlobMetadata {
    pub da_layer: String,
    pub da_network: String,
    pub da_height: u64,
    pub commitment: HexBinary,
    pub namespace: String,
    pub contract_addr: Addr,
    pub data_format: String,
    pub size: u64,
    pub data_uri: Option<String>,
    pub expires_at: Option<u64>,
    pub status: DataStatus,
}

#[cw_serde]
pub enum DataStatus {
    Active,
    Archived,
    Pruned,
}

#[cw_serde]
pub struct RegisterMsg {
    pub metadata: BlobMetadata,
    pub proof: Option<HexBinary>,
}

2.2.2 注册流程

当外部应用需要将数据发布到 DA 层时,遵循以下注册流程:

1. 应用构建数据
    |
2. 应用将数据提交到 DA 层(如 Celestia)
    |   <- 收到 DA 层返回的高度和承诺
    |
3. 应用构建 BlobMetadata
    |
4. 应用调用 DA 注册中心合约的 register 方法
    |   <- 注册中心验证 commitment 格式
    |
5. 注册中心存储 BlobMetadata
    |
6. MSG Chain 索引器监听注册事件
    |
7. 轻客户端 / 查询者 通过注册中心查找 DA 数据

合约实现示例:

use cosmwasm_std::{
    entry_point, to_json_binary, Binary, Deps, DepsMut, Env, HexBinary,
    MessageInfo, Response, StdError, StdResult,
};
use cw_storage_plus::Map;

pub const BLOB_REGISTRY: Map<(&[u8], &[u8]), BlobMetadata> = Map::new("blob_registry");
pub const CONTRACT_BLOBS: Map<&Addr, Vec<HexBinary>> = Map::new("contract_blobs");
pub const ACTIVE_BLOB_COUNT: Item<u64> = Item::new("active_blob_count");
pub const REGISTRY_CONFIG: Item<RegistryConfig> = Item::new("registry_config");

#[entry_point]
pub fn execute(
    deps: DepsMut,
    _env: Env,
    info: MessageInfo,
    msg: ExecuteMsg,
) -> Result<Response, ContractError> {
    match msg {
        ExecuteMsg::Register { metadata, proof } => {
            execute_register_blob(deps, info, metadata, proof)
        }
        ExecuteMsg::UpdateStatus { namespace, commitment, status } => {
            execute_update_status(deps, info, namespace, commitment, status)
        }
        ExecuteMsg::PruneExpired {} => execute_prune_expired(deps, _env),
    }
}

pub fn execute_register_blob(
    deps: DepsMut,
    info: MessageInfo,
    metadata: BlobMetadata,
    _proof: Option<HexBinary>,
) -> Result<Response, ContractError> {
    match metadata.da_layer.as_str() {
        "celestia" => {
            if metadata.commitment.len() != 32 {
                return Err(ContractError::InvalidCommitmentLength {});
            }
        }
        "avail" => {
            if metadata.commitment.len() != 48 {
                return Err(ContractError::InvalidCommitmentLength {});
            }
        }
        "eigenda" => {
            if metadata.commitment.len() != 32 {
                return Err(ContractError::InvalidCommitmentLength {});
            }
        }
        _ => return Err(ContractError::UnsupportedDALayer { layer: metadata.da_layer.clone() }),
    }

    let namespace_bytes = HexBinary::from_hex(&metadata.namespace)?;
    let commitment_bytes = metadata.commitment.to_vec();

    if BLOB_REGISTRY.has(deps.storage, (&namespace_bytes, &commitment_bytes)) {
        return Err(ContractError::BlobAlreadyRegistered {});
    }

    if info.sender != metadata.contract_addr {
        // 允许 relayer 代表合约注册,需额外权限验证
    }

    BLOB_REGISTRY.save(
        deps.storage,
        (&namespace_bytes, &commitment_bytes),
        &metadata,
    )?;

    CONTRACT_BLOBS.update(deps.storage, &metadata.contract_addr, |existing| {
        let mut list = existing.unwrap_or_default();
        list.push(metadata.commitment.clone());
        Ok::<_, StdError>(list)
    })?;

    ACTIVE_BLOB_COUNT.update(deps.storage, |count| Ok(count + 1))?;

    Ok(Response::new()
        .add_attribute("action", "register_blob")
        .add_attribute("da_layer", &metadata.da_layer)
        .add_attribute("da_height", metadata.da_height.to_string())
        .add_attribute("namespace", &metadata.namespace)
        .add_attribute("commitment", metadata.commitment.to_hex())
        .add_attribute("contract_addr", metadata.contract_addr.to_string())
        .add_attribute("size", metadata.size.to_string()))
}

#[entry_point]
pub fn query(deps: Deps, _env: Env, msg: QueryMsg) -> StdResult<Binary> {
    match msg {
        QueryMsg::GetBlob { namespace, commitment } => {
            query_blob_by_commitment(deps, namespace, commitment)
        }
        QueryMsg::ListContractBlobs { contract_addr, page, page_size } => {
            query_blobs_by_contract(deps, contract_addr, page, page_size)
        }
        QueryMsg::ListActiveBlobs { page, page_size } => {
            query_active_blobs(deps, page, page_size)
        }
        QueryMsg::GetStats {} => query_stats(deps),
    }
}

2.2.3 注册中心的经济模型

DA 注册中心的操作需要支付 MSG Chain 上的 gas 费用,此外还可以引入以下经济激励机制:

操作 gas 费用 附加费用 说明
注册 blob 元数据 标准 wasm 执行 gas 可选 Layer fee (MSG) 按数据大小阶梯收费
更新数据状态 标准 wasm 执行 gas 无 所有者操作
查询(query) 无 gas 无 免费查询
裁剪过期数据 gas 返还 无 激励清理操作
数据迁移 标准 wasm 执行 gas 迁移费用 跨 DA 层迁移

2.2.4 注册中心的权限模型

注册中心支持灵活的权限控制,确保数据注册的安全性:

#[cw_serde]
pub struct RegistryConfig {
    pub admin: Addr,
    pub allowed_contracts: Vec<Addr>,
    pub allowed_da_layers: Vec<String>,
    pub max_blobs_per_contract: u32,
    pub max_blob_size: u64,
    pub default_retention_period: u64,
}

#[cw_serde]
pub enum PrivilegeLevel {
    Admin,
    AllowedContract,
    Public,
}

权限验证逻辑:

pub fn assert_can_register(
    deps: &DepsMut,
    sender: &Addr,
    target_contract: &Addr,
    da_layer: &str,
) -> Result<(), ContractError> {
    let config = REGISTRY_CONFIG.load(deps.storage)?;

    if sender == &config.admin {
        return Ok(());
    }

    if !config.allowed_contracts.contains(sender) {
        return Err(ContractError::Unauthorized {
            reason: "sender not in allowed contracts".to_string(),
        });
    }

    if sender != target_contract {
        return Err(ContractError::Unauthorized {
            reason: "cannot register on behalf of other contract".to_string(),
        });
    }

    if !config.allowed_da_layers.contains(&da_layer.to_string()) {
        return Err(ContractError::Unauthorized {
            reason: format!("DA layer '{}' not allowed", da_layer),
        });
    }

    Ok(())
}

3. DA 层集成架构

MSG Chain 支持三种主要的 DA 集成架构模式,每种模式适用于不同的应用场景和安全需求。

3.1 主权 rollup 模式

主权 Rollup(Sovereign Rollup)是 MSG Chain 首选推荐的 DA 集成模式。在这种架构中,Rollup 将自己的交易数据发布到 DA 层,同时使用 Celestia 的共识进行数据排序,但状态转换和结算在 MSG Chain 上独立进行。

3.1.1 架构组件

                     +---------------------------------+
                     |       Celestia DA 网络            |
                     |  (数据排序 + 数据可用性保证)        |
                     +------------+--------------------+
                                  | blob 数据流
                     +------------v--------------------+
                     |     MSG Chain 结算层              |
                     |  (状态根验证 + 欺诈证明)           |
                     +------------+--------------------+
                                  | 状态承诺
                     +------------v--------------------+
                     |     Rollup 执行环境                |
                     |  (CosmWasm 合约执行)               |
                     +---------------------------------+

组件说明:

组件 职责 技术选型
Sequencer (排序器) 接收用户交易,构建批次提交到 DA 层 独立运行或去中心化排序器集
DA Client (DA 客户端) 与 Celestia 网络交互,提交和检索 blob celestia-node / go-celestia
Settlement Contract (结算合约) MSG Chain 上的验证合约,验证 DA 承诺和状态根 CosmWasm 智能合约
Full Node (全节点) 执行交易,维护状态,验证 DA 数据 msg-chaind + celestia-light-node
Light Client (轻客户端) 仅验证 DA 采样和数据承诺 celestia-light-node

3.1.2 主权 Rollup 的生命周期

阶段一: 交易收集

排序器收集用户提交的交易,构建交易批次。

type Batch struct {
    Index            uint64     `json:"index"`
    Transactions     [][]byte   `json:"transactions"`
    PrevStateRoot    [32]byte   `json:"prev_state_root"`
    PrevDACommitment [32]byte   `json:"prev_da_commitment"`
    Timestamp        int64      `json:"timestamp"`
}

阶段二: 数据发布

排序器将交易批次发布到 Celestia DA 层,获取数据承诺和区块高度。

func (s *Sequencer) SubmitBatch(ctx context.Context, batch *Batch) (*BatchSubmission, error) {
    batchData, err := json.Marshal(batch)
    if err != nil {
        return nil, fmt.Errorf("marshal batch: %w", err)
    }

    namespace := s.GetNamespace()
    blob, err := celestia.NewBlob(namespace, batchData)
    if err != nil {
        return nil, fmt.Errorf("create blob: %w", err)
    }

    height, err := s.daClient.SubmitBlob(ctx, []*celestia.Blob{blob}, s.estimateGas(batchData))
    if err != nil {
        return nil, fmt.Errorf("submit blob to celestia: %w", err)
    }

    commitment := sha256.Sum256(batchData)

    return &BatchSubmission{
        Namespace:  namespace,
        Height:     height,
        Commitment: commitment[:],
        Size:       len(batchData),
    }, nil
}

阶段三: 结算验证

排序器将 DA 提交信息提交到 MSG Chain 上的结算合约。

#[entry_point]
pub fn execute_submit_batch(
    deps: DepsMut,
    env: Env,
    info: MessageInfo,
    batch_submission: BatchSubmission,
    after_state_root: HexBinary,
) -> Result<Response, ContractError> {
    let config = SETTLEMENT_CONFIG.load(deps.storage)?;
    if !config.sequencers.contains(&info.sender) {
        return Err(ContractError::Unauthorized {
            reason: "not authorized sequencer".to_string(),
        });
    }

    let state = SETTLEMENT_STATE.load(deps.storage)?;
    if batch_submission.index != state.next_batch_index {
        return Err(ContractError::InvalidBatchIndex {
            expected: state.next_batch_index,
            actual: batch_submission.index,
        });
    }

    DA_SUBMISSIONS.save(
        deps.storage,
        batch_submission.index,
        &DaSubmissionRecord {
            sequencer: info.sender,
            da_layer: batch_submission.da_layer,
            da_height: batch_submission.da_height,
            namespace: batch_submission.namespace,
            commitment: batch_submission.commitment.clone(),
            after_state_root: after_state_root.clone(),
            submitted_at: env.block.time,
        },
    )?;

    state.latest_state_root = after_state_root;
    state.next_batch_index += 1;
    state.latest_submission_time = env.block.time;
    SETTLEMENT_STATE.save(deps.storage, &state)?;

    Ok(Response::new()
        .add_attribute("action", "submit_batch")
        .add_attribute("index", batch_submission.index.to_string())
        .add_attribute("da_layer", &batch_submission.da_layer)
        .add_attribute("da_height", batch_submission.da_height.to_string())
        .add_attribute("commitment", batch_submission.commitment.to_hex()))
}

阶段四: 挑战窗口

在挑战窗口期(如 7 天)内,任何人都可以提交欺诈证明挑战状态根的正确性。

#[entry_point]
pub fn execute_challenge(
    deps: DepsMut,
    env: Env,
    info: MessageInfo,
    batch_index: u64,
    fraud_proof: FraudProof,
) -> Result<Response, ContractError> {
    let submission = DA_SUBMISSIONS.load(deps.storage, batch_index)?;

    let config = SETTLEMENT_CONFIG.load(deps.storage)?;
    let challenge_deadline = submission.submitted_at
        .plus_nanos(config.challenge_window_nanos);
    if env.block.time > challenge_deadline {
        return Err(ContractError::ChallengeWindowExpired {});
    }

    let is_valid = verify_fraud_proof(&fraud_proof, &submission)?;
    if !is_valid {
        return Err(ContractError::InvalidFraudProof {});
    }

    // 挑战成功:回滚状态,惩罚排序器
    Ok(Response::new()
        .add_attribute("action", "challenge_succeeded")
        .add_attribute("batch_index", batch_index.to_string())
        .add_attribute("challenger", info.sender.to_string()))
}

3.1.3 主权 Rollup 的优势

主权 Rollup 模式为 MSG Chain 提供了以下关键优势:

  1. 主权独立: MSG Chain 控制自己的状态转换逻辑,不依赖 DA 层进行链的升级或治理。
  2. 硬分叉自由: 社区可以在不依赖 DA 层的情况下进行链的升级或分叉。
  3. 低费用: 利用 Celestia 的低成本 DA,交易费用远低于传统 L1。
  4. 可组合性: 同一个 MSG Chain 上的多个 Rollup 共享结算层,实现原子可组合性。

3.2 validium 模式

Validium 模式使用零知识证明(ZKP)来验证状态转换的正确性,数据本身存储在 DA 层而非 MSG Chain 链上。与 Rollup 的主要区别在于:Validium 的数据完全保存在链下,链上只保存 ZK 证明。

3.2.1 validium 架构

用户交易
    |
    v
[排序器 / 执行环境]
    |
    +---> 批量交易数据 ---> DA 层 (Celestia/Avail)
    |
    +---> ZK 证明 ---> MSG Chain 验证合约
                           |
                           v
                    +--------------+
                    | ZK 验证合约   |
                    | (验证 SNARK) |
                    +--------------+

在 MSG Chain 上部署 ZK 验证合约:

use cosmwasm_std::{
    entry_point, Binary, DepsMut, Env, MessageInfo, Response,
};

#[entry_point]
pub fn execute_verify_validium_batch(
    deps: DepsMut,
    env: Env,
    info: MessageInfo,
    da_commitment: HexBinary,
    da_height: u64,
    zk_proof: Binary,
    public_inputs: Vec<HexBinary>,
) -> Result<Response, ContractError> {
    let registry = DA_REGISTRY.load(deps.storage)?;
    registry.verify_commitment(&da_commitment, da_height)?;

    let vk = VALIDIUM_VERIFYING_KEY.load(deps.storage)?;
    let verified = verify_zk_proof(&vk, &zk_proof, &public_inputs)?;

    if !verified {
        return Err(ContractError::InvalidZkProof {});
    }

    let state_root = public_inputs[0].clone();
    VALIDIUM_STATE_ROOT.save(deps.storage, &state_root)?;

    Ok(Response::new()
        .add_attribute("action", "verify_validium_batch")
        .add_attribute("da_height", da_height.to_string())
        .add_attribute("new_state_root", state_root.to_hex()))
}

3.2.2 validium 与 Rollup 的选择

维度 Sovereign Rollup Validium
数据存储 DA 层 DA 层(完全链下)
状态验证 欺诈证明(乐观) ZK 证明(即时)
提款延迟 挑战窗口期(7 天) 即时(证明验证后)
计算复杂度 低(重执行) 高(生成 ZK 证明)
成本 低 中(证明生成成本)
隐私性 交易数据公开 可隐藏交易细节
适用场景 通用应用、高频交易 金融应用、隐私场景

3.3 volition 模式

Volition 模式是 Rollup 和 Validium 的混合体,允许用户为每笔交易选择数据存储位置。

3.3.1 volition 数据通道设计

用户交易
    |
    +-- "rollup" 模式
    |       +---> 交易数据 -> DA 层
    |       +---> 状态根 -> MSG Chain
    |
    +-- "validium" 模式
            +---> 交易数据 -> DA 层(可选私有)
            +---> ZK 证明 -> MSG Chain

Volition 合约示例:

#[cw_serde]
pub enum DataAvailabilityMode {
    Rollup,
    Validium,
    Volition,
}

#[cw_serde]
pub struct VolitionBatchSubmission {
    pub index: u64,
    pub mode: DataAvailabilityMode,
    pub da_submission: DaSubmissionInfo,
    pub after_state_root: HexBinary,
    pub zk_proof: Option<Binary>,
}

#[entry_point]
pub fn execute_submit_volition_batch(
    deps: DepsMut,
    env: Env,
    info: MessageInfo,
    submission: VolitionBatchSubmission,
) -> Result<Response, ContractError> {
    match submission.mode {
        DataAvailabilityMode::Rollup => {
            execute_submit_rollup_batch(deps, env, info, submission)
        }
        DataAvailabilityMode::Validium => {
            execute_verify_validium_batch(
                deps, env, info,
                submission.da_submission.commitment,
                submission.da_submission.height,
                submission.zk_proof.unwrap(),
                vec![submission.after_state_root],
            )
        }
        DataAvailabilityMode::Volition => {
            execute_submit_volition_mixed(deps, env, info, submission)
        }
    }
}

3.3.2 volition 的适用场景

Volition 模式适合以下 MSG Chain 应用场景:

  1. 去中心化交易所: 大额交易使用 Rollup 模式(更高安全性),小额频繁交易使用 Validium 模式(更快确认)。
  2. 游戏平台: 资产转移走 Rollup,游戏状态更新走 Validium。
  3. 社交应用: 关键身份操作走 Rollup,内容发布走 Validium。
  4. 企业合规: 公开交易走 Rollup,涉及商业机密的数据走 Validium。

4. 数据发布与验证

本章深入探讨数据在 DA 层发布和验证的技术细节。

4.1 blob 提交

Blob 是 DA 层中最基本的数据单元。

4.1.1 blob 数据结构

type Blob struct {
    Namespace   Namespace `json:"namespace"`
    Data        []byte    `json:"data"`
    ContentType string    `json:"content_type,omitempty"`
    Version     uint8     `json:"version"`
    Timestamp   int64     `json:"timestamp,omitempty"`
    Signature   []byte    `json:"signature,omitempty"`
}

type Namespace [8]byte

func NamespaceFromHex(hex string) (Namespace, error) {
    bytes, err := hex.DecodeString(hex)
    if err != nil {
        return Namespace{}, err
    }
    var ns Namespace
    copy(ns[:], bytes)
    return ns, nil
}

func GenerateNamespace(appID string) Namespace {
    hash := sha256.Sum256([]byte(appID))
    var ns Namespace
    copy(ns[:], hash[:8])
    return ns
}

4.1.2 blob 提交流程

type BlobSubmitter struct {
    daClient      *celestia.Client
    namespace     Namespace
    signer        crypto.PrivKey
    gasEstimator  func(size int) gas.Gas
}

func (bs *BlobSubmitter) SubmitBatch(
    ctx context.Context,
    txs [][]byte,
) (*BatchReceipt, error) {
    batch := Batch{
        Index:          bs.nextIndex(),
        Transactions:   txs,
        PrevStateRoot:  bs.currentStateRoot(),
        Timestamp:      time.Now().Unix(),
    }
    data, err := json.Marshal(batch)
    if err != nil {
        return nil, fmt.Errorf("marshal batch: %w", err)
    }

    compressed, err := compressData(data)
    if err != nil {
        return nil, fmt.Errorf("compress data: %w", err)
    }

    sig, err := bs.signer.Sign(compressed)
    if err != nil {
        return nil, fmt.Errorf("sign data: %w", err)
    }

    blob := &celestia.Blob{
        Namespace: bs.namespace,
        Data:      compressed,
        Version:   1,
        Signature: sig,
    }

    gas := bs.gasEstimator(len(compressed))
    height, err := bs.daClient.SubmitBlob(ctx, []*celestia.Blob{blob}, gas)
    if err != nil {
        return nil, fmt.Errorf("submit blob to celestia: %w", err)
    }

    commitment := ComputeCommitment(compressed)
    bs.updateIndex(batch.Index, height, commitment)

    return &BatchReceipt{
        Height:     height,
        Commitment: commitment,
        BlobSize:   len(compressed),
    }, nil
}

func compressData(data []byte) ([]byte, error) {
    var buf bytes.Buffer
    writer := snappy.NewBufferedWriter(&buf)
    if _, err := writer.Write(data); err != nil {
        return nil, err
    }
    if err := writer.Close(); err != nil {
        return nil, err
    }
    return buf.Bytes(), nil
}

4.1.3 多 blob 提交与原子性

func (bs *BlobSubmitter) SubmitBatchSplit(
    ctx context.Context,
    largeData []byte,
    maxBlobSize int,
) ([]*BatchReceipt, error) {
    chunks := splitIntoChunks(largeData, maxBlobSize)
    blobs := make([]*celestia.Blob, len(chunks))
    for i, chunk := range chunks {
        blobs[i] = &celestia.Blob{
            Namespace: bs.namespace,
            Data:      chunk,
            Version:   1,
        }
    }

    heights, err := bs.daClient.SubmitBlob(ctx, blobs, bs.estimateBulkGas(blobs))
    if err != nil {
        return nil, fmt.Errorf("submit blobs: %w", err)
    }

    receipts := make([]*BatchReceipt, len(blobs))
    var commitments [][]byte
    for i, height := range heights {
        commitment := ComputeCommitment(chunks[i])
        commitments = append(commitments, commitment)
        receipts[i] = &BatchReceipt{
            Height:     height,
            Commitment: commitment,
            BlobSize:   len(chunks[i]),
        }
    }

    atomicCommitment := ComputeMerkleRoot(commitments)
    _ = atomicCommitment

    return receipts, nil
}

4.1.4 DA 层 blob 费用计算

func EstimateDAFee(layer string, dataSize int, gasPrice gas.Gas) (*DAFee, error) {
    switch layer {
    case "celestia":
        gas := uint64(dataSize)*10 + 50000
        totalGas := gas * gasPrice.Uint64()
        return &DAFee{
            GasLimit:     gas,
            GasPrice:     gasPrice,
            TotalFee:     sdk.NewCoin("utia", sdk.NewIntFromUint64(totalGas)),
            Confirmation: 12,
        }, nil

    case "avail":
        pricePerByte := sdk.NewDecCoin("avail", sdk.NewInt(1))
        total := uint64(dataSize) * pricePerByte.Amount.Uint64()
        return &DAFee{
            TotalFee:     sdk.NewCoin("avail", sdk.NewIntFromUint64(total)),
            Confirmation: 20,
        }, nil

    case "eigenda":
        ratePerByte := sdk.NewDecCoin("eth", sdk.NewInt(1))
        total := uint64(dataSize) * ratePerByte.Amount.Uint64()
        return &DAFee{
            TotalFee:     sdk.NewCoin("eth", sdk.NewIntFromUint64(total)),
            Confirmation: 2,
        }, nil

    default:
        return nil, fmt.Errorf("unsupported DA layer: %s", layer)
    }
}

4.2 数据采样

数据可用性采样(Data Availability Sampling, DAS)是 DA 层的核心机制。

4.2.1 纠删码编码

type ReedSolomonEncoder struct {
    dataShards   int
    parityShards int
}

func (e *ReedSolomonEncoder) Encode(data []byte) ([][]byte, error) {
    shardSize := (len(data) + e.dataShards - 1) / e.dataShards
    shards := make([][]byte, e.dataShards+e.parityShards)
    for i := 0; i < e.dataShards; i++ {
        start := i * shardSize
        end := start + shardSize
        if end > len(data) {
            end = len(data)
        }
        shards[i] = make([]byte, shardSize)
        copy(shards[i], data[start:end])
    }

    enc, err := reedsolomon.New(e.dataShards, e.parityShards)
    if err != nil {
        return nil, err
    }
    if err := enc.Encode(shards); err != nil {
        return nil, err
    }
    return shards, nil
}

4.2.2 二维纠删码矩阵

Celestia 使用二维 Reed-Solomon 编码,将 blob 数据组织为正方形矩阵:

原始数据: [A1, A2, A3, A4, A5, A6, A7, A8, A9]

步骤 1: 重塑为 3x3 矩阵
+----+----+----+
| A1 | A2 | A3 |
+----+----+----+
| A4 | A5 | A6 |
+----+----+----+
| A7 | A8 | A9 |
+----+----+----+

步骤 2: 行 Reed-Solomon 编码(添加校验列)
+----+----+----+----+
| A1 | A2 | A3 | P1 |
+----+----+----+----+
| A4 | A5 | A6 | P2 |
+----+----+----+----+
| A7 | A8 | A9 | P3 |
+----+----+----+----+
| P4 | P5 | P6 | P7 |
+----+----+----+----+
  ^--- 列 Reed-Solomon 编码(添加校验行)

最终 4x4 矩阵:16 个分片,原始数据 9 个分片,7 个校验分片

4.2.3 轻客户端采样协议

type DasLightClient struct {
    daClient      *celestia.LightClient
    params        DasParams
    logger        log.Logger
}

type DasParams struct {
    SamplesPerBlock int
    SampleInterval  uint64
    Threshold       float64
    MaxRetries      int
}

func DefaultDasParams() DasParams {
    return DasParams{
        SamplesPerBlock: 20,
        SampleInterval:  1,
        Threshold:       0.99,
        MaxRetries:      3,
    }
}

func (lc *DasLightClient) SampleBlock(
    ctx context.Context,
    height uint64,
) (*SamplingResult, error) {
    eds, err := lc.daClient.GetErasuredData(ctx, height)
    if err != nil {
        return nil, fmt.Errorf("get erasured data: %w", err)
    }

    meta, err := lc.daClient.GetBlockMetadata(ctx, height)
    if err != nil {
        return nil, fmt.Errorf("get block metadata: %w", err)
    }

    coordinates := generateRandomSamples(meta.Rows, meta.Cols, lc.params.SamplesPerBlock)
    var successes, failures int

    for _, coord := range coordinates {
        result, err := lc.sampleCell(ctx, height, coord)
        if err != nil {
            failures++
            if failures > maxAllowedFailures(lc.params) {
                return &SamplingResult{Available: false}, nil
            }
            continue
        }
        successes++
        _ = result
    }

    available := float64(successes)/float64(successes+failures) >= lc.params.Threshold
    return &SamplingResult{
        Available: available,
        Sampled:   successes + failures,
        Successes: successes,
        Failures:  failures,
    }, nil
}

func generateRandomSamples(rows, cols, count int) []Coordinate {
    rng := rand.New(rand.NewSource(time.Now().UnixNano()))
    samples := make([]Coordinate, count)
    used := make(map[string]bool)

    for i := 0; i < count; i++ {
        for {
            row := rng.Intn(rows)
            col := rng.Intn(cols)
            key := fmt.Sprintf("%d-%d", row, col)
            if !used[key] {
                used[key] = true
                samples[i] = Coordinate{Row: row, Col: col}
                break
            }
        }
    }
    return samples
}

func maxAllowedFailures(params DasParams) int {
    total := params.SamplesPerBlock
    threshold := params.Threshold
    minSuccesses := int(math.Ceil(float64(total) * threshold))
    return total - minSuccesses
}

4.2.4 概率保证

func ComputeSamplingProbability(
    totalShards int,
    missingShards int,
    sampleCount int,
) float64 {
    if missingShards == 0 {
        return 0.0
    }
    if missingShards > totalShards {
        return 1.0
    }
    pSingleMiss := float64(totalShards-missingShards) / float64(totalShards)
    pAllMiss := math.Pow(pSingleMiss, float64(sampleCount))
    pDetection := 1.0 - pAllMiss
    return pDetection
}

不同参数下的检测概率:

总分片 (n) 缺失分片 (m) 缺失比例 采样次数 (k) 检测概率
16 1 6.25% 20 72.89%
16 2 12.5% 20 93.07%
16 4 25.0% 20 99.68%
64 4 6.25% 20 72.89%
64 8 12.5% 20 93.07%
256 4 1.56% 20 26.87%
256 4 1.56% 100 79.27%
256 4 1.56% 300 99.10%

4.3 KZG 承诺

KZG 承诺(Kate-Zaverucha-Goldberg Commitment)是一种基于椭圆曲线配对的多项式承诺方案。

4.3.1 KZG 承诺的数学基础

给定多项式 f(x) = a0 + a1x + a2x^2 + ... + anx^(n-1)

KZG 承诺 C = f(s) x G1
  其中 s 是秘密值(通过可信设置生成)
  G1 是椭圆曲线 BLS12-381 的生成元

求值证明: 对于任意点 i,可以生成证据 pi 证明 f(i) = yi
  验证: e(C - yi x G1, G2) = e(pi, s x G2 - i x G2)
  其中 e 是椭圆曲线配对运算

4.3.2 KZG 承诺在 DA 中的应用

type KZGProver struct {
    srs *kzg.SRS
}

func NewKZGProver(srsFile string) (*KZGProver, error) {
    srs, err := kzg.NewSRSFromFile(srsFile)
    if err != nil {
        return nil, fmt.Errorf("load SRS: %w", err)
    }
    return &KZGProver{srs: srs}, nil
}

func (kp *KZGProver) ComputeCommitment(data []byte) ([]byte, error) {
    coeffs := bytesToScalars(data)
    commitment, err := kzg.CommitPoly(coeffs, kp.srs)
    if err != nil {
        return nil, fmt.Errorf("compute kzg commitment: %w", err)
    }
    return commitment.Serialize(), nil
}

func (kp *KZGProver) CreateWitness(data []byte, index int) ([]byte, []byte, error) {
    coeffs := bytesToScalars(data)
    value := evaluatePolynomial(coeffs, index)
    valueBytes := scalarToBytes(value)
    witness, err := kzg.CreateWitness(coeffs, index, kp.srs)
    if err != nil {
        return nil, nil, fmt.Errorf("create witness: %w", err)
    }
    return witness.Serialize(), valueBytes, nil
}

func (kp *KZGProver) VerifyWitness(
    commitment []byte, index int, value []byte, witness []byte,
) (bool, error) {
    comm, _ := kzg.CommitmentFromBytes(commitment)
    wit, _ := kzg.WitnessFromBytes(witness)
    val := bytesToScalar(value)
    return kzg.CheckWitness(comm, index, val, wit, kp.srs)
}

func (kp *KZGProver) VerifyBatchWitness(
    commitments [][]byte, indexes []int, values [][]byte, witnesses [][]byte,
) (bool, error) {
    var comms []kzg.Commitment
    var wits []kzg.Witness
    var vals []fr.Element
    for i := range commitments {
        comm, _ := kzg.CommitmentFromBytes(commitments[i])
        wit, _ := kzg.WitnessFromBytes(witnesses[i])
        val := bytesToScalar(values[i])
        comms = append(comms, comm)
        wits = append(wits, wit)
        vals = append(vals, val)
    }
    return kzg.CheckBatchWitness(comms, indexes, vals, wits, kp.srs)
}

4.3.3 KZG vs Merkle 树对比

维度 KZG 承诺 Merkle 树
承诺大小 48 字节(恒定) 32 字节(SHA-256 根)
证明大小 48 字节(恒定) O(log n)
验证时间 恒定(1 次配对运算) O(log n)
批量验证 高效(多重配对) 低效
可信设置 需要 不需要
抗量子性 不抗量子 可抗量子
实现复杂度 高 低

4.4 欺诈证明

欺诈证明(Fraud Proof)是乐观验证系统中的关键组件。

4.4.1 欺诈证明的生命周期

1. 排序器提交批次(DA 承诺 + 状态根)
    |
2. 等待挑战窗口期
    |
    +---> 无人质疑 -> 最终确认状态根
    |
    +---> 观察者提交欺诈证明
            |
            +---> 验证欺诈证明
            |       |
            |       +---> 有效 -> 回滚状态,惩罚排序器
            |       |
            |       +---> 无效 -> 惩罚挑战者
            |
            +---> (可选) 争议解决

4.4.2 欺诈证明的数据结构

#[cw_serde]
pub enum FraudProof {
    InvalidStateTransition(InvalidStateTransitionProof),
    InvalidDACommitment(InvalidDACommitmentProof),
    DataUnavailable(DataUnavailableProof),
}

#[cw_serde]
pub struct InvalidStateTransitionProof {
    pub batch_index: u64,
    pub batch_data: HexBinary,
    pub pre_state_root: HexBinary,
    pub correct_post_state_root: HexBinary,
    pub claimed_post_state_root: HexBinary,
    pub state_diff: Vec<StateDiff>,
}

#[cw_serde]
pub struct StateDiff {
    pub contract_addr: Addr,
    pub storage_key: HexBinary,
    pub correct_value: HexBinary,
    pub claimed_value: HexBinary,
}

#[cw_serde]
pub struct InvalidDACommitmentProof {
    pub batch_index: u64,
    pub actual_data: HexBinary,
    pub correct_commitment: HexBinary,
    pub claimed_commitment: HexBinary,
}

#[cw_serde]
pub struct DataUnavailableProof {
    pub batch_index: u64,
    pub da_layer: String,
    pub da_height: u64,
    pub namespace: String,
    pub sampling_attempts: Vec<SamplingAttempt>,
}

4.4.3 欺诈证明验证

pub fn verify_fraud_proof(
    deps: &DepsMut,
    proof: &FraudProof,
    submission: &DaSubmissionRecord,
) -> Result<bool, ContractError> {
    match proof {
        FraudProof::InvalidStateTransition(istp) => {
            verify_state_transition(deps, istp, submission)
        }
        FraudProof::InvalidDACommitment(idacp) => {
            verify_da_commitment_fraud(idacp, submission)
        }
        FraudProof::DataUnavailable(dup) => {
            verify_data_unavailability(deps, dup, submission)
        }
    }
}

fn verify_state_transition(
    deps: &DepsMut,
    proof: &InvalidStateTransitionProof,
    submission: &DaSubmissionRecord,
) -> Result<bool, ContractError> {
    let settlement_state = SETTLEMENT_STATE.load(deps.storage)?;
    if settlement_state.prev_state_root != proof.pre_state_root {
        return Err(ContractError::FraudVerificationFailed {
            reason: "pre state root mismatch".to_string(),
        });
    }

    let data_hash = sha2::Sha256::digest(&proof.batch_data);
    if data_hash[..] != submission.commitment.as_slice() {
        return Err(ContractError::FraudVerificationFailed {
            reason: "data commitment mismatch".to_string(),
        });
    }

    let correct_root = reexecute_transactions(
        deps.as_ref(), &proof.pre_state_root, &proof.batch_data,
    )?;

    if correct_root != proof.correct_post_state_root {
        return Err(ContractError::FraudVerificationFailed {
            reason: "computed state root does not match proof".to_string(),
        });
    }

    if proof.claimed_post_state_root == proof.correct_post_state_root {
        return Err(ContractError::FraudVerificationFailed {
            reason: "claimed state root is actually correct".to_string(),
        });
    }
    Ok(true)
}

4.4.4 挑战窗口与经济激励

#[cw_serde]
pub struct ChallengeConfig {
    pub window_nanos: u64,
    pub sequencer_bond: Coin,
    pub challenger_reward_ratio: Decimal,
    pub challenge_deposit: Coin,
    pub invalid_challenge_penalty: Coin,
}

impl Default for ChallengeConfig {
    fn default() -> Self {
        ChallengeConfig {
            window_nanos: 7 * 24 * 60 * 60 * 1_000_000_000,
            sequencer_bond: Coin::new(10_000_000_000, "umsg"),
            challenger_reward_ratio: Decimal::percent(10),
            challenge_deposit: Coin::new(1_000_000, "umsg"),
            invalid_challenge_penalty: Coin::new(500_000, "umsg"),
        }
    }
}

5. CosmWasm 合约与 DA 交互

本章详细介绍 CosmWasm 智能合约如何与 DA 层进行交互。

5.1 contract 将数据发布到 DA

CosmWasm 合约本身无法直接访问外部 DA 网络。因此,合约通过一种称为 DA Router 的架构模式与 DA 层交互。

5.1.1 DA Router 架构

CosmWasm 合约
    |
    |   emit PublishDataEvent { namespace, data_hash, data_uri }
    |
    v
MSG Chain 事件日志
    |
    |   indexer 订阅事件
    |
    v
DA Router (链外服务)
    |
    |   1. 从 data_uri 获取原始数据
    |   2. 发布到 Celestia/Avail/EigenDA
    |   3. 将提交收据写回 MSG Chain
    |
    v
DA 层 (Celestia)

5.1.2 合约触发 DA 发布

use cosmwasm_std::{
    entry_point, Addr, Binary, DepsMut, Env, HexBinary,
    MessageInfo, Response, Uint128, to_json_binary,
};
use cw_storage_plus::{Item, Map};

#[cw_serde]
pub struct ContractConfig {
    pub admin: Addr,
    pub da_router: Addr,
    pub namespace: String,
    pub max_data_size: u64,
}

pub const CONFIG: Item<ContractConfig> = Item::new("config");

#[cw_serde]
pub struct DaPublishRequest {
    pub id: u64,
    pub data_id: String,
    pub data_hash: HexBinary,
    pub data_size: u64,
    pub data_format: String,
    pub data_uri: String,
    pub requester: Addr,
    pub requested_at: u64,
    pub status: RequestStatus,
}

#[cw_serde]
pub enum RequestStatus {
    Pending,
    Submitted { da_height: u64, commitment: HexBinary },
    Failed { reason: String },
}

pub const PUBLISH_REQUESTS: Map<u64, DaPublishRequest> = Map::new("publish_requests");
pub const NEXT_REQUEST_ID: Item<u64> = Item::new("next_request_id");

#[entry_point]
pub fn execute_request_publish(
    deps: DepsMut,
    env: Env,
    info: MessageInfo,
    data_id: String,
    data_hash: HexBinary,
    data_size: u64,
    data_format: String,
    data_uri: String,
) -> Result<Response, ContractError> {
    let config = CONFIG.load(deps.storage)?;

    if data_size > config.max_data_size {
        return Err(ContractError::DataTooLarge {
            max: config.max_data_size,
            actual: data_size,
        });
    }

    if data_hash.len() != 32 {
        return Err(ContractError::InvalidHashLength {});
    }

    let request_id = NEXT_REQUEST_ID.load(deps.storage)?;
    let request = DaPublishRequest {
        id: request_id,
        data_id,
        data_hash,
        data_size,
        data_format,
        data_uri,
        requester: info.sender.clone(),
        requested_at: env.block.time.nanos(),
        status: RequestStatus::Pending,
    };

    PUBLISH_REQUESTS.save(deps.storage, request_id, &request)?;
    NEXT_REQUEST_ID.save(deps.storage, &(request_id + 1))?;

    Ok(Response::new()
        .add_attribute("action", "request_publish")
        .add_attribute("request_id", request_id.to_string())
        .add_attribute("data_id", &data_id)
        .add_attribute("data_hash", data_hash.to_hex())
        .add_attribute("data_size", data_size.to_string())
        .add_attribute("data_format", &data_format)
        .add_attribute("data_uri", &data_uri)
        .add_attribute("namespace", &config.namespace)
        .add_attribute("requester", info.sender))
}

5.1.3 DA Router 链外服务实现

type DaRouter struct {
    msgChainClient *msgchain.Client
    daClient       *celestia.Client
    registryAddr   sdk.AccAddress
    contracts      map[string]string
    logger         log.Logger
}

func (r *DaRouter) Start(ctx context.Context) error {
    query := fmt.Sprintf(
        "wasm.action='request_publish' AND wasm._contract_address='%s'",
        r.registryAddr.String(),
    )

    subscriber, err := r.msgChainClient.SubscribeEvents(ctx, query)
    if err != nil {
        return fmt.Errorf("subscribe events: %w", err)
    }

    for {
        select {
        case <-ctx.Done():
            return ctx.Err()
        case event := <-subscriber.Events():
            r.handlePublishEvent(ctx, event)
        }
    }
}

func (r *DaRouter) handlePublishEvent(ctx context.Context, event sdk.Event) {
    attrs := make(map[string]string)
    for _, attr := range event.Attributes {
        attrs[string(attr.Key)] = string(attr.Value)
    }

    requestID := attrs["request_id"]
    dataURI := attrs["data_uri"]
    dataHash := attrs["data_hash"]
    namespace := attrs["namespace"]
    dataFormat := attrs["data_format"]

    data, err := fetchDataFromURI(ctx, dataURI)
    if err != nil {
        r.reportFailure(requestID, fmt.Sprintf("fetch failed: %v", err))
        return
    }

    hash := sha256.Sum256(data)
    expectedHash, _ := hex.DecodeString(dataHash)
    if !bytes.Equal(hash[:], expectedHash) {
        r.reportFailure(requestID, "data hash mismatch")
        return
    }

    var processedData []byte
    switch dataFormat {
    case "json", "protobuf", "binary":
        processedData = data
    default:
        processedData = data
    }

    ns := parseNamespace(namespace)
    blob, err := celestia.NewBlob(ns, processedData)
    if err != nil {
        r.reportFailure(requestID, fmt.Sprintf("create blob failed: %v", err))
        return
    }

    height, err := r.daClient.SubmitBlob(ctx, []*celestia.Blob{blob}, auto)
    if err != nil {
        r.reportFailure(requestID, fmt.Sprintf("DA submission failed: %v", err))
        return
    }

    commitment := sha256.Sum256(processedData)
    r.reportSuccess(ctx, requestID, uint64(height), commitment[:])
}

func (r *DaRouter) reportSuccess(
    ctx context.Context, requestID string, daHeight uint64, commitment []byte,
) {
    msg := &wasm.MsgExecuteContract{
        Sender:   r.operatorAddr,
        Contract: r.registryAddr.String(),
        Msg:      buildConfirmPublishMsg(requestID, daHeight, commitment),
        Funds:    sdk.NewCoins(),
    }

    _, err := r.msgChainClient.BroadcastTx(ctx, msg)
    if err != nil {
        r.logger.Error("broadcast confirm tx failed", "error", err)
    }
}

5.1.4 合约回调确认

#[entry_point]
pub fn execute_confirm_publish(
    deps: DepsMut,
    info: MessageInfo,
    request_id: u64,
    da_height: u64,
    commitment: HexBinary,
) -> Result<Response, ContractError> {
    let config = CONFIG.load(deps.storage)?;
    if info.sender != config.da_router {
        return Err(ContractError::Unauthorized {
            reason: "only DA router can confirm".to_string(),
        });
    }

    let mut request = PUBLISH_REQUESTS.load(deps.storage, request_id)?;
    request.status = RequestStatus::Submitted { da_height, commitment };
    PUBLISH_REQUESTS.save(deps.storage, request_id, &request)?;

    Ok(Response::new()
        .add_attribute("action", "confirm_publish")
        .add_attribute("request_id", request_id.to_string())
        .add_attribute("da_height", da_height.to_string())
        .add_attribute("commitment", commitment.to_hex()))
}

#[entry_point]
pub fn query(deps: Deps, _env: Env, msg: QueryMsg) -> StdResult<Binary> {
    match msg {
        QueryMsg::GetPublishRequest { request_id } => {
            let req = PUBLISH_REQUESTS.load(deps.storage, request_id)?;
            to_json_binary(&req)
        }
        QueryMsg::ListRequestsByStatus { status, page, page_size } => {
            let requests: Vec<DaPublishRequest> = PUBLISH_REQUESTS
                .range(deps.storage, None, None, cosmwasm_std::Order::Ascending)
                .filter_map(|item| {
                    if let Ok((_, req)) = item {
                        match (&req.status, &status) {
                            (RequestStatus::Submitted { .. }, _)
                            | (RequestStatus::Pending, _)
                            | (RequestStatus::Failed { .. }, _) => Some(req),
                            _ => None,
                        }
                    } else { None }
                })
                .skip(page * page_size)
                .take(page_size)
                .collect();
            to_json_binary(&requests)
        }
    }
}

5.2 从 DA 读取数据验证

5.2.1 链下数据读取

type DaDataReader struct {
    daClient      *celestia.Client
    registryAddr  string
    msgChainQuery *msgchain.QueryClient
}

func (r *DaDataReader) ReadAndVerify(
    ctx context.Context, namespace [8]byte, height uint64, commitment []byte,
) ([]byte, error) {
    blobs, err := r.daClient.GetBlobs(ctx, height, []celestia.Namespace{namespace})
    if err != nil {
        return nil, fmt.Errorf("get blobs from DA: %w", err)
    }

    if len(blobs) == 0 {
        return nil, fmt.Errorf("no blobs found at height %d", height)
    }

    dataHash := sha256.Sum256(blobs[0].Data)
    if !bytes.Equal(dataHash[:], commitment) {
        return nil, fmt.Errorf("commitment mismatch: expected %x, got %x",
            commitment, dataHash[:])
    }

    return blobs[0].Data, nil
}

5.2.2 合约内验证数据

#[cw_serde]
pub struct DaVerificationProof {
    pub da_layer: String,
    pub da_height: u64,
    pub namespace: String,
    pub commitment: HexBinary,
    pub data: HexBinary,
    pub proof: Option<HexBinary>,
}

#[entry_point]
pub fn execute_verify_da_data(
    deps: DepsMut, _env: Env, info: MessageInfo,
    verification: DaVerificationProof,
) -> Result<Response, ContractError> {
    let config = VERIFICATION_CONFIG.load(deps.storage)?;
    if verification.data.len() > config.max_verification_size {
        return Err(ContractError::DataTooLarge {
            max: config.max_verification_size,
            actual: verification.data.len() as u64,
        });
    }

    let valid = match verification.da_layer.as_str() {
        "celestia" => {
            let computed_hash = sha2::Sha256::digest(&verification.data);
            if computed_hash[..] != verification.commitment.as_slice() {
                return Err(ContractError::CommitmentVerificationFailed {
                    expected: HexBinary::from(computed_hash.as_slice()),
                    actual: verification.commitment,
                });
            }
            true
        }
        "avail" | "eigenda" => {
            return Err(ContractError::KzgVerificationNotSupported {});
        }
        _ => {
            return Err(ContractError::UnsupportedDALayer {
                layer: verification.da_layer,
            });
        }
    };

    Ok(Response::new()
        .add_attribute("action", "verify_da_data")
        .add_attribute("namespace", &verification.namespace)
        .add_attribute("commitment", verification.commitment.to_hex())
        .add_attribute("valid", valid.to_string()))
}

5.2.3 使用 IBC 跨链验证

#[cw_serde]
pub struct IbcDaVerificationPacket {
    pub da_height: u64,
    pub namespace: String,
    pub commitment: HexBinary,
    pub proof: HexBinary,
    pub requester: String,
}

#[entry_point]
pub fn ibc_packet_receive(
    deps: DepsMut, env: Env, msg: IbcPacketReceiveMsg,
) -> Result<IbcReceiveResponse, ContractError> {
    let packet: IbcDaVerificationPacket = from_json(&msg.packet.data)?;
    let client_id = IBC_CLIENT.load(deps.storage)?;

    let verified = verify_ibc_da_proof(
        deps.as_ref(), &client_id, packet.da_height, &packet.proof,
    )?;

    if !verified {
        return Err(ContractError::IbcVerificationFailed {});
    }

    CROSS_CHAIN_VERIFICATION.save(
        deps.storage,
        (&packet.namespace, &packet.commitment),
        &CrossChainVerificationRecord {
            verified: true,
            verified_at: env.block.time,
            origin_chain: msg.packet.src.channel_id,
        },
    )?;

    Ok(IbcReceiveResponse::new()
        .add_attribute("action", "ibc_da_verify")
        .add_attribute("namespace", &packet.namespace)
        .add_attribute("commitment", packet.commitment.to_hex()))
}

6. 轻客户端验证

轻客户端(Light Client)是 DA 生态中的关键组成部分。

6.1 轻节点 DA 采样

6.1.1 轻节点架构

MSG Chain 轻节点
    |
    +-- CometBFT 轻客户端
    |       +-- 验证共识状态(验证人集合变更)
    |
    +-- DA 采样模块
    |       +-- Celestia 轻客户端
    |       +-- Avail 轻客户端
    |       +-- EigenDA 轻客户端
    |
    +-- 数据验证引擎
            +-- 承诺验证
            +-- 采样验证
            +-- 状态根验证

6.1.2 轻节点 DA 采样实现

type MsgChainLightNode struct {
    tendermintClient *tmclient.Client
    daSamplers       map[string]DASampler
    verifier         *DataVerifier
    state            *LightNodeState
    logger           log.Logger
}

type DASampler interface {
    Name() string
    Sample(ctx context.Context, submission *DaSubmission) (*SamplingResult, error)
    GetSampleParams() SampleParams
}

type DaSubmission struct {
    Index       uint64
    DADataLayer string
    DAHeight    uint64
    Namespace   [8]byte
    Commitment  []byte
    StateRoot   []byte
    SubmittedAt time.Time
}

type CelestiaSampler struct {
    lightClient *celestia.LightClient
    params      SampleParams
}

func (cs *CelestiaSampler) Sample(
    ctx context.Context, submission *DaSubmission,
) (*SamplingResult, error) {
    header, err := cs.lightClient.GetHeader(ctx, submission.DAHeight)
    if err != nil {
        return nil, fmt.Errorf("get header: %w", err)
    }

    dah := header.DAH
    if dah == nil {
        return nil, fmt.Errorf("no DAH in header")
    }

    rowStart, rowEnd := dah.GetNamespaceRows(submission.Namespace)
    samplesCount := cs.params.SamplesPerBlock
    if samplesCount == 0 {
        samplesCount = 20
    }

    var successes, failures int
    rng := rand.New(rand.NewSource(time.Now().UnixNano()))

    for i := 0; i < samplesCount; i++ {
        row := rowStart + rng.Intn(rowEnd-rowStart)
        col := rng.Intn(dah.Cols)

        proof, err := cs.lightClient.GetProof(ctx, submission.DAHeight, row, col)
        if err != nil {
            failures++
            if failures > maxAllowedFailures(cs.params) {
                return &SamplingResult{Available: false}, nil
            }
            continue
        }

        if !proof.Verify(dah.Hash()) {
            failures++
            continue
        }
        successes++
    }

    available := float64(successes)/float64(successes+failures) >= 0.99
    return &SamplingResult{
        Available: available,
        Sampled:   successes + failures,
        Successes: successes,
        Failures:  failures,
    }, nil
}

type AvailSampler struct {
    lightClient *avail.LightClient
}

func (as *AvailSampler) Sample(
    ctx context.Context, submission *DaSubmission,
) (*SamplingResult, error) {
    blockData, err := as.lightClient.GetBlockData(ctx, submission.DAHeight)
    if err != nil {
        return nil, err
    }

    if !bytes.Equal(blockData.KZGCommitment, submission.Commitment) {
        return &SamplingResult{Available: false}, nil
    }

    appData := blockData.GetAppData(submission.Namespace[:])
    samples := 10
    validSamples := 0

    for i := 0; i < samples; i++ {
        idx := rand.Intn(len(appData.Cells))
        cell := appData.Cells[idx]
        valid, err := kzg.VerifyWitness(
            blockData.KZGCommitment, cell.Index, cell.Value, cell.Proof,
            as.lightClient.SRS(),
        )
        if err != nil {
            continue
        }
        if valid {
            validSamples++
        }
    }

    available := validSamples >= 8
    return &SamplingResult{
        Available: available,
        Sampled:   samples,
        Successes: validSamples,
        Failures:  samples - validSamples,
    }, nil
}

6.1.3 轻节点状态同步

func (ln *MsgChainLightNode) SyncState(ctx context.Context) error {
    trustedHeight := ln.state.GetTrustedHeight()
    if trustedHeight == 0 {
        trustedHeight = ln.state.GetCheckpoint()
    }

    latestBlock, err := ln.tendermintClient.GetLatestBlock(ctx)
    if err != nil {
        return fmt.Errorf("get latest block: %w", err)
    }

    for height := trustedHeight + 1; height <= latestBlock.Header.Height; height++ {
        header, err := ln.tendermintClient.GetHeader(ctx, height)
        if err != nil {
            return fmt.Errorf("get header %d: %w", height, err)
        }

        err = ln.verifier.VerifyHeader(header)
        if err != nil {
            return fmt.Errorf("verify header %d: %w", height, err)
        }

        events, err := ln.tendermintClient.GetBlockEvents(ctx, height,
            "wasm.action='submit_batch'")
        if err != nil {
            continue
        }

        for _, event := range events {
            submission := parseDaSubmission(event)
            sampler, ok := ln.daSamplers[submission.DADataLayer]
            if !ok {
                continue
            }

            result, err := sampler.Sample(ctx, &submission)
            if err != nil {
                ln.state.MarkSuspicious(height, submission.Index)
                continue
            }

            if !result.Available {
                ln.state.AddAlert(Alert{
                    Type:       AlertDataUnavailable,
                    Height:     height,
                    Submission: submission,
                    Detail:     result,
                })
            }
        }

        ln.state.SetSyncedHeight(height)
    }
    return nil
}

6.2 无需全节点即可验证数据可用性

6.2.1 验证流程概览

轻客户端流程(无需全节点):

1. 从可信种子获取最近的 MSG Chain 区块头
    |
2. 通过 CometBFT 轻客户端验证区块头的共识签名
    |
3. 从区块头中提取 DA 提交事件
    |
4. 对于每个 DA 提交:
    +-- a. 连接 DA 层轻节点(如 Celestia 轻节点)
    +-- b. 获取 DA 层的区块头
    +-- c. 验证 DA 层区块头的共识
    +-- d. 执行随机采样(不下载完整区块)
    +-- e. 验证每个采样的 Merkle 或 KZG 证明
    |
5. 汇总采样结果,判断数据可用性
    |
6. (可选)下载并验证状态转换

6.2.2 无全节点状态验证

type LightClientVerifier struct {
    tmClient     *tmclient.Client
    daClients    map[string]DALightClient
    stateManager *StateManager
}

type DALightClient interface {
    GetBlockHeader(ctx context.Context, height uint64) (*BlockHeader, error)
    VerifyBlockHeader(header *BlockHeader) error
    SampleCell(ctx context.Context, daHeight uint64, row, col int) (*CellSample, error)
}

func (v *LightClientVerifier) VerifyDASubmission(
    ctx context.Context, submission *DaSubmission,
) (*VerificationResult, error) {
    daClient, ok := v.daClients[submission.DADataLayer]
    if !ok {
        return nil, fmt.Errorf("unsupported DA layer: %s", submission.DADataLayer)
    }

    daHeader, err := daClient.GetBlockHeader(ctx, submission.DAHeight)
    if err != nil {
        return &VerificationResult{Verified: false, Level: VerificationLevelNone}, nil
    }

    if err := daClient.VerifyBlockHeader(daHeader); err != nil {
        return &VerificationResult{Verified: false, Level: VerificationLevelHeaderOnly}, nil
    }

    if err := v.verifyCommitment(daHeader, submission); err != nil {
        return &VerificationResult{Verified: false, Level: VerificationLevelHeaderOnly}, nil
    }

    sampleResults, err := v.performSampling(ctx, daClient, submission, daHeader)
    if err != nil {
        return &VerificationResult{Verified: false, Level: VerificationLevelCommitment}, nil
    }

    available := calculateAvailability(sampleResults)
    if !available {
        return &VerificationResult{
            Verified: false, Level: VerificationLevelSampling,
            SampleResults: sampleResults,
        }, nil
    }

    return &VerificationResult{
        Verified: true, Level: VerificationLevelFull,
        SampleResults: sampleResults,
    }, nil
}

func (v *LightClientVerifier) performSampling(
    ctx context.Context, daClient DALightClient,
    submission *DaSubmission, daHeader *BlockHeader,
) ([]*SampleResult, error) {
    samplesPerBlock := 20
    if submission.DADataLayer == "avail" {
        samplesPerBlock = 10
    }

    rowStart, rowEnd := 0, daHeader.Rows
    if submission.DADataLayer == "celestia" && len(submission.Namespace) > 0 {
        rowStart, rowEnd = daHeader.GetNamespaceRows(submission.Namespace)
    }

    rng := rand.New(rand.NewSource(time.Now().UnixNano()))
    results := make([]*SampleResult, 0, samplesPerBlock)
    used := make(map[string]bool)

    for len(results) < samplesPerBlock {
        row := rowStart + rng.Intn(rowEnd-rowStart)
        col := rng.Intn(daHeader.Cols)
        key := fmt.Sprintf("%d-%d", row, col)
        if used[key] {
            continue
        }
        used[key] = true

        cell, err := daClient.SampleCell(ctx, submission.DAHeight, row, col)
        if err != nil {
            results = append(results, &SampleResult{Row: row, Col: col, Success: false})
            continue
        }

        valid := verifyCellData(cell, daHeader, row, col)
        results = append(results, &SampleResult{Row: row, Col: col, Success: valid})
    }
    return results, nil
}

type VerificationResult struct {
    Verified       bool
    Level          VerificationLevel
    SampleResults  []*SampleResult
}

type VerificationLevel int

const (
    VerificationLevelNone        VerificationLevel = iota
    VerificationLevelHeaderOnly
    VerificationLevelCommitment
    VerificationLevelSampling
    VerificationLevelFull
)

type SampleResult struct {
    Row     int
    Col     int
    Success bool
}

func calculateAvailability(results []*SampleResult) bool {
    if len(results) == 0 {
        return false
    }
    successes := 0
    for _, r := range results {
        if r.Success {
            successes++
        }
    }
    rate := float64(successes) / float64(len(results))
    return rate >= 0.99
}

6.2.3 轻客户端的信任模型

信任等级  | 描述                                           | 需要的数据
----------|------------------------------------------------|--------------
1. 共识级  | 信任 MSG Chain 和 DA 层的验证人集合               | 验证人集合 + 区块头
2. 承诺级  | 信任数据承诺的正确性(欺骗需要碰撞哈希)      | 数据承诺 + 区块头
3. 采样级  | 以高概率验证数据可用性(概率性保证)           | 采样证明
4. 执行级  | 完全验证状态转换(确定性保证)                  | 完整区块数据

轻客户端的安全参数:

type LightClientSecurityParams struct {
    TrustedValidatorSetChangeRatio    float64
    DASamplesPerBlock                 int
    DASuccessThreshold                float64
    HeaderVerificationLevel           string
    RequireSettlementVerification     bool
    MaxUnconfirmedBlocks              int
}

func DefaultLightClientParams() LightClientSecurityParams {
    return LightClientSecurityParams{
        TrustedValidatorSetChangeRatio:  0.33,
        DASamplesPerBlock:               20,
        DASuccessThreshold:              0.99,
        HeaderVerificationLevel:         "full",
        RequireSettlementVerification:   true,
        MaxUnconfirmedBlocks:            100,
    }
}

7. 索引器与 DA

索引器(Indexer)在 MSG Chain 生态中扮演关键角色。

7.1 索引器从 DA 层读取数据

7.1.1 索引器架构

         MSG Chain                     DA 层 (Celestia)
             |                              |
             | 区块/事件                      | Blob 数据
             v                              v
    +-------------------+          +------------------+
    | CometBFT RPC      |          | Celestia RPC     |
    +--------+----------+          +--------+---------+
             |                             |
             v                             v
    +-----------------------------------------------+
    |              MSG Chain Indexer                 |
    |                                                 |
    |  +----------+  +----------+  +---------------+ |
    |  | Block    |  | Event    |  | DA Data       | |
    |  | Processor|  | Processor|  | Fetcher       | |
    |  +----+-----+  +----+-----+  +------+--------+ |
    |       |             |              |           |
    |       v             v              v           |
    |  +----------------------------------------+    |
    |  |        Data Aggregator                 |    |
    |  +-------------------+--------------------+    |
    |                      |                        |
    |                      v                        |
    |  +----------------------------------------+    |
    |  |         Storage Backend                |    |
    |  |  (PostgreSQL + Object Store)           |    |
    |  +----------------------------------------+    |
    +----------------------+------------------------+
                           |
                           v
                  GraphQL / REST API
                           |
                           v
                     dApp / 用户

7.1.2 DA 数据获取器实现

type DaDataFetcher struct {
    celestiaClients map[string]*celestia.Client
    availClients    map[string]*avail.Client
    registryClients map[string]*msgchain.Client
    store           *DataStore
    logger          log.Logger
}

func (f *DaDataFetcher) FetchAndProcessBlob(
    ctx context.Context, layer, network string, height uint64,
    namespace []byte, commitment []byte, contractAddr string,
) error {
    var data []byte
    var err error

    switch layer {
    case "celestia":
        client, ok := f.celestiaClients[network]
        if !ok {
            return fmt.Errorf("unknown celestia network: %s", network)
        }
        data, err = f.fetchCelestiaBlob(ctx, client, height, namespace)
    case "avail":
        client, ok := f.availClients[network]
        if !ok {
            return fmt.Errorf("unknown avail network: %s", network)
        }
        data, err = f.fetchAvailBlob(ctx, client, height, namespace)
    default:
        return fmt.Errorf("unsupported DA layer: %s", layer)
    }
    if err != nil {
        return fmt.Errorf("fetch blob: %w", err)
    }

    if err := verifyCommitmentMatch(data, commitment); err != nil {
        return fmt.Errorf("commitment verification failed: %w", err)
    }

    data, err = decompressData(data)
    if err != nil {
        return fmt.Errorf("decompress data: %w", err)
    }

    var batch Batch
    if err := json.Unmarshal(data, &batch); err != nil {
        return fmt.Errorf("unmarshal batch: %w", err)
    }

    for _, txBytes := range batch.Transactions {
        var tx cosmos.Tx
        if err := proto.Unmarshal(txBytes, &tx); err != nil {
            continue
        }
        for _, msg := range tx.Body.Messages {
            if msg.TypeUrl == "/cosmwasm.wasm.v1.MsgExecuteContract" {
                var execMsg wasm.MsgExecuteContract
                if err := proto.Unmarshal(msg.Value, &execMsg); err != nil {
                    continue
                }
                f.processContractData(ctx, contractAddr, execMsg)
            }
        }
    }

    indexedDA := IndexedDAData{
        Layer:            layer,
        Network:          network,
        DAHeight:         height,
        Namespace:        hex.EncodeToString(namespace),
        Commitment:       hex.EncodeToString(commitment),
        ContractAddr:     contractAddr,
        BatchIndex:       batch.Index,
        TransactionCount: len(batch.Transactions),
        IndexedAt:        time.Now(),
    }
    return f.store.SaveIndexedDA(ctx, &indexedDA)
}

func (f *DaDataFetcher) fetchCelestiaBlob(
    ctx context.Context, client *celestia.Client, height uint64, namespace []byte,
) ([]byte, error) {
    ns := celestia.Namespace{}
    copy(ns[:], namespace[:8])
    blobs, err := client.GetBlobs(ctx, height, []celestia.Namespace{ns})
    if err != nil {
        return nil, err
    }
    if len(blobs) == 0 {
        return nil, fmt.Errorf("no blobs found at height %d", height)
    }
    return blobs[0].Data, nil
}

func (f *DaDataFetcher) fetchAvailBlob(
    ctx context.Context, client *avail.Client, height uint64, appID []byte,
) ([]byte, error) {
    block, err := client.GetBlock(ctx, height)
    if err != nil {
        return nil, err
    }
    data, ok := block.GetAppData(appID)
    if !ok {
        return nil, fmt.Errorf("no data for app ID %x at height %d", appID, height)
    }
    return data, nil
}

func (f *DaDataFetcher) startDaWatcher(ctx context.Context) {
    ticker := time.NewTicker(30 * time.Second)
    defer ticker.Stop()

    for {
        select {
        case <-ctx.Done():
            return
        case <-ticker.C:
            events, err := f.registryClients["msg-chain-1"].QueryUnprocessedDAEvents(ctx)
            if err != nil {
                continue
            }
            for _, event := range events {
                err := f.FetchAndProcessBlob(
                    ctx, event.DALayer, event.DANetwork,
                    event.DAHeight, event.Namespace, event.Commitment, event.ContractAddr,
                )
                if err != nil {
                    continue
                }
                f.registryClients["msg-chain-1"].MarkEventProcessed(ctx, event.ID)
            }
        }
    }
}

7.2 归档存储

7.2.1 归档存储架构

+--------------------------------------------------+
|                  归档存储                           |
|                                                   |
|  +--------------+  +--------------+               |
|  | Hot Tier     |  | Warm Tier    |               |
|  | (30 天)      |  | (1 年)       |               |
|  | PostgreSQL   |  | S3/OSS       |               |
|  | + Redis      |  | Parquet      |               |
|  +------+-------+  +------+-------+               |
|         |                 |                       |
|         v                 v                       |
|  +-------------------------------------------+    |
|  | Cold Tier (永久)                          |    |
|  | Glacier / Arweave / Filecoin              |    |
|  +-------------------------------------------+    |
+--------------------------------------------------+

7.2.2 归档管理器实现

type ArchiveManager struct {
    hotStore  *HotDataStore
    warmStore *WarmDataStore
    coldStore *ColdDataStore
    retention RetentionPolicy
    logger    log.Logger
}

type RetentionPolicy struct {
    HotRetentionDays  int
    WarmRetentionDays int
    ColdRetention     bool
    BlobDeleteAfter   int
}

func DefaultRetentionPolicy() RetentionPolicy {
    return RetentionPolicy{
        HotRetentionDays:  30,
        WarmRetentionDays: 365,
        ColdRetention:     true,
        BlobDeleteAfter:   90,
    }
}

func (am *ArchiveManager) ArchiveIndexedData(ctx context.Context) error {
    now := time.Now()

    hotExpired, err := am.hotStore.FindExpired(ctx, am.retention.HotRetentionDays)
    if err != nil {
        return err
    }

    for _, data := range hotExpired {
        compressed, err := compressForArchive(data)
        if err != nil {
            continue
        }

        warmKey := fmt.Sprintf("da/%s/%s/height_%d.parquet",
            data.Layer, data.Namespace, data.DAHeight)
        if err := am.warmStore.Put(ctx, warmKey, compressed); err != nil {
            continue
        }
        am.hotStore.Delete(ctx, data.ID)
    }

    if am.retention.ColdRetention {
        warmExpired, err := am.warmStore.FindExpired(ctx, am.retention.WarmRetentionDays)
        if err != nil {
            return err
        }

        for _, entry := range warmExpired {
            data, err := am.warmStore.Get(ctx, entry.Key)
            if err != nil {
                continue
            }

            txID, err := am.coldStore.Archive(ctx, data)
            if err != nil {
                continue
            }

            am.hotStore.SaveArchiveRecord(ctx, &ArchiveRecord{
                OriginalKey:   entry.Key,
                ColdStoreTxID: txID,
                ArchivedAt:    now,
            })
            am.warmStore.Delete(ctx, entry.Key)
        }
    }

    return nil
}

7.2.3 归档存储 Schema

CREATE TABLE indexed_da_data (
    id              BIGSERIAL PRIMARY KEY,
    da_layer        VARCHAR(32) NOT NULL,
    da_network      VARCHAR(64) NOT NULL,
    da_height       BIGINT NOT NULL,
    namespace       VARCHAR(64) NOT NULL,
    commitment      VARCHAR(128) NOT NULL,
    contract_addr   VARCHAR(64) NOT NULL,
    batch_index     BIGINT,
    tx_count        INTEGER NOT NULL DEFAULT 0,
    blob_size       BIGINT NOT NULL DEFAULT 0,
    blob_content    BYTEA,
    data_format     VARCHAR(32) DEFAULT 'json',
    indexed_at      TIMESTAMP NOT NULL DEFAULT NOW(),
    archived        BOOLEAN DEFAULT FALSE,
    created_at      TIMESTAMP NOT NULL DEFAULT NOW(),
    updated_at      TIMESTAMP NOT NULL DEFAULT NOW()
);

CREATE INDEX idx_da_data_layer_height ON indexed_da_data(da_layer, da_height);
CREATE INDEX idx_da_data_contract ON indexed_da_data(contract_addr);
CREATE INDEX idx_da_data_namespace ON indexed_da_data(namespace);
CREATE INDEX idx_da_data_created ON indexed_da_data(created_at);

CREATE TABLE archive_records (
    id              BIGSERIAL PRIMARY KEY,
    original_key    VARCHAR(512) NOT NULL,
    cold_store_tx_id VARCHAR(128),
    cold_store_type VARCHAR(32),
    archived_at     TIMESTAMP NOT NULL DEFAULT NOW(),
    data_hash       VARCHAR(64) NOT NULL,
    verified        BOOLEAN DEFAULT FALSE
);

CREATE TABLE archive_metadata (
    id              BIGSERIAL PRIMARY KEY,
    da_layer        VARCHAR(32) NOT NULL,
    da_network      VARCHAR(64) NOT NULL,
    namespace       VARCHAR(64) NOT NULL,
    height_from     BIGINT NOT NULL,
    height_to       BIGINT NOT NULL,
    blob_count      INTEGER NOT NULL,
    total_size      BIGINT NOT NULL,
    storage_key     VARCHAR(512) NOT NULL,
    compression     VARCHAR(16),
    checksum        VARCHAR(64) NOT NULL,
    archived_at     TIMESTAMP NOT NULL DEFAULT NOW()
);

7.3 与 MSG Chain 事件流的整合

7.3.1 事件关联模型

MSG Chain 区块
    |
    +-- tx_1: 合约 A 的 execute
    |       +-- event: wasm.action="request_publish"
    |               +-- request_id: 42
    |               +-- data_id: "batch_通用维护记录_001"
    |               +-- data_hash: 0xabc...
    |               +-- data_uri: "ipfs://QmX..."
    |
    +-- tx_2: DA Router 回调
    |       +-- event: wasm.action="confirm_publish"
    |               +-- request_id: 42
    |               +-- da_height: 1234567
    |               +-- commitment: 0xdef...
    |
    +-- tx_3: 结算合约的 batch 提交
            +-- event: wasm.action="submit_batch"
                    +-- batch_index: 89
                    +-- da_layer: "celestia"
                    +-- da_height: 1234567
                    +-- namespace: "a1b2c3d4"
                    +-- commitment: 0xdef...

7.3.2 事件流整合实现

type EventIntegrator struct {
    eventSource *msgchain.EventSource
    daFetcher   *DaDataFetcher
    archiveMgr  *ArchiveManager
    store       *DataStore
}

func (ei *EventIntegrator) ProcessEvents(ctx context.Context) error {
    query := "wasm._contract_address EXISTS AND tm.event = 'Tx'"
    subscriber, err := ei.eventSource.Subscribe(ctx, query)
    if err != nil {
        return fmt.Errorf("subscribe events: %w", err)
    }

    for {
        select {
        case <-ctx.Done():
            return ctx.Err()
        case event := <-subscriber.Events():
            ei.handleEvent(ctx, event)
        }
    }
}

func (ei *EventIntegrator) handleEvent(ctx context.Context, event sdk.Event) {
    attrs := parseEventAttributes(event)
    action := attrs["wasm.action"]

    switch action {
    case "request_publish":
        ei.handlePublishRequest(ctx, attrs)
    case "confirm_publish":
        ei.handlePublishConfirm(ctx, attrs)
    case "submit_batch":
        ei.handleBatchSubmission(ctx, attrs)
    case "challenge_succeeded":
        ei.handleChallenge(ctx, attrs)
    }
}

func (ei *EventIntegrator) handleBatchSubmission(
    ctx context.Context, attrs map[string]string,
) {
    layer := attrs["da_layer"]
    heightStr := attrs["da_height"]
    namespace := attrs["namespace"]
    commitment := attrs["commitment"]
    height, _ := strconv.ParseUint(heightStr, 10, 64)

    go func() {
        err := ei.daFetcher.FetchAndProcessBlob(
            ctx, layer, "mainnet", height,
            hexDecode(namespace), hexDecode(commitment), attrs["contract_addr"],
        )
        if err != nil {
            ei.logger.Error("fetch blob", "height", height, "error", err)
            return
        }
        ei.store.UpdateDADataStatus(ctx, layer, height, commitment, "fetched")
    }()
}

7.3.3 完整的索引管道

func (ei *EventIntegrator) FullIndexingPipeline(ctx context.Context) error {
    go ei.ProcessEvents(ctx)

    go func() {
        ticker := time.NewTicker(1 * time.Minute)
        defer ticker.Stop()
        for {
            select {
            case <-ctx.Done():
                return
            case <-ticker.C:
                unprocessed, err := ei.store.GetUnprocessedDASubmissions(ctx, 100)
                if err != nil {
                    continue
                }
                for _, sub := range unprocessed {
                    ei.daFetcher.FetchAndProcessBlob(
                        ctx, sub.DALayer, sub.DANetwork,
                        sub.DAHeight, sub.Namespace, sub.Commitment, sub.ContractAddr,
                    )
                    ei.store.MarkDASubmissionProcessed(ctx, sub.ID)
                }
            }
        }
    }()

    go func() {
        archiveTicker := time.NewTicker(24 * time.Hour)
        defer archiveTicker.Stop()
        for {
            select {
            case <-ctx.Done():
                return
            case <-archiveTicker.C:
                ei.archiveMgr.ArchiveIndexedData(ctx)
            }
        }
    }()

    return nil
}

8. 集成实践

本章提供完整的集成实践指南,以 Celestia 作为 MSG Chain 的 DA 层。

8.1 使用 Celestia 作为 MSG Chain 的 DA 层

8.1.1 环境准备

前提条件:

组件 版本要求 说明
Go >= 1.22 MSG Chain 和 Celestia 客户端均使用 Go 实现
msg-chaind >= v1.0.0 MSG Chain 节点二进制
celestia-node >= v0.13.0 Celestia 轻节点或全节点
msgchain-cli >= v1.0.0 MSG Chain 命令行工具

Celestia 节点初始化:

# 初始化 Celestia 轻节点
celestia light init --p2p.network celestia

# 启动 Celestia 轻节点
celestia light start \
    --core.ip https://celestia-rpc.msgchain.org \
    --keyring.accname msgchain-da \
    --p2p.network celestia

# 创建 DA 专用密钥
celestia keys add msgchain-da-key

MSG Chain 配置:

# 在 MSG Chain 应用配置中启用 DA 模块
msg-chaind config set app.da_enabled true
msg-chaind config set app.da_layer celestia
msg-chaind config set app.da_rpc_endpoint http://localhost:26658
msg-chaind config set app.da_namespace $(celestia keys show msgchain-da-key -a)

8.1.2 部署 DA 注册中心合约

# 编译合约
docker run --rm -v "$(pwd)":/code \
    --mount type=volume,source=registry_cache,target=/code/target \
    --mount type=volume,source=registry_cache,target=/usr/local/cargo/registry \
    cosmwasm/workspace-optimizer:0.15.0

# 存储合约到 MSG Chain
RESP=$(msg-chaind tx wasm store ./artifacts/da_registry.wasm \
    --from deployer \
    --gas auto \
    --gas-adjustment 1.3 \
    --chain-id msg-chain-1 \
    --node https://rpc.msgchain.org:26657 \
    --output json)

CODE_ID=$(echo $RESP | jq -r '.logs[0].events[] | select(.type=="store_code") | .attributes[] | select(.key=="code_id") | .value')

# 实例化合约
INIT='{
    "admin": "msg1admin...",
    "allowed_contracts": [],
    "allowed_da_layers": ["celestia"],
    "max_blobs_per_contract": 1000,
    "max_blob_size": 1048576,
    "default_retention_period": 2592000000000000
}'

msg-chaind tx wasm instantiate $CODE_ID "$INIT" \
    --from deployer \
    --label "da-registry-v1" \
    --admin msg1admin... \
    --chain-id msg-chain-1 \
    --gas auto

8.1.3 DA Router 配置

// DA Router 配置文件
type DaRouterConfig struct {
    // MSG Chain 连接
    MsgChainRPC     string `yaml:"msgchain_rpc"`
    MsgChainGRPC    string `yaml:"msgchain_grpc"`
    OperatorKeyName string `yaml:"operator_key_name"`
    RegistryAddress string `yaml:"registry_address"`

    // Celestia 连接
    CelestiaRPC     string `yaml:"celestia_rpc"`
    CelestiaKeyName string `yaml:"celestia_key_name"`
    CelestiaNetwork string `yaml:"celestia_network"`

    // 运行参数
    PollInterval    int    `yaml:"poll_interval"`
    MaxRetries      int    `yaml:"max_retries"`
    GasMultiplier   float64 `yaml:"gas_multiplier"`
}

// config.yaml
// msgchain_rpc: https://rpc.msgchain.org:26657
// msgchain_grpc: https://grpc.msgchain.org:9090
// operator_key_name: da-router-operator
// registry_address: msg1registry...
// celestia_rpc: http://localhost:26658
// celestia_key_name: msgchain-da-key
// celestia_network: celestia
// poll_interval: 30
// max_retries: 3
// gas_multiplier: 1.2

8.1.4 命名空间分配策略

// NamespaceManager 命名空间管理器
type NamespaceManager struct {
    registry     *msgchain.Client
    reserved     map[string][8]byte
    nextSequence uint64
}

// AllocateNamespace 为新合约分配命名空间
func (nm *NamespaceManager) AllocateNamespace(
    ctx context.Context, contractAddr string,
) ([8]byte, error) {
    // 方法一:基于合约地址哈希
    hash := sha256.Sum256([]byte(contractAddr))
    var ns [8]byte
    copy(ns[:], hash[:8])

    // 方法二:顺序分配(适用于已知数量的合约)
    // binary.BigEndian.PutUint64(ns[:], nm.nextSequence)
    // nm.nextSequence++

    // 方法三:结合两者,确保唯一性
    // seq := nm.getNextSequence(ctx)
    // hashInput := fmt.Sprintf("%s-%d", contractAddr, seq)
    // hash := sha256.Sum256([]byte(hashInput))
    // copy(ns[:], hash[:8])

    return ns, nil
}

// 命名空间在 Celestia 中的使用限制:
// - 保留范围: 0x0000000000000000 - 0x00000000000000FF (系统使用)
// - 公共范围: 0x0000000000000100 - 0xFFFFFFFFFFFFFFFF  (应用使用)
// - 建议为每个 MSG Chain 合约或应用分配唯一命名空间
// - 命名空间在区块头中用于路由轻客户端查询

8.2 blob 生命周期示例

8.2.1 完整生命周期

阶段 1: 创建
    |
    +-- dApp 生成需要发布的数据
    +-- 数据格式化为 blob(JSON/Protobuf/Binary)
    +-- 可选:压缩和加密
    |
阶段 2: 发布到 DA
    |
    +-- 提交到 Celestia DA 层
    |   +-- 计算 gas = size * 10 + 50000
    |   +-- 签名交易
    |   +-- 广播到 Celestia 验证人
    |   +-- 等待包含(~12 秒)
    |
    +-- 获得提交收据
        +-- da_height: 区块高度
        +-- commitment: SHA-256 承诺
        +-- namespace: 命名空间 ID
    |
阶段 3: 在 MSG Chain 上注册
    |
    +-- 调用 DA 注册中心合约
    |   +-- namespace + commitment + metadata
    |   +-- 支付注册 gas 费用
    |
    +-- 注册中心发出事件
        +-- wasm.action="register_blob"
    |
阶段 4: 数据可用性验证
    |
    +-- 轻客户端执行 DAS 采样
    |   +-- 采样 20 个随机坐标
    |   +-- 验证每个坐标的 Merkle 证明
    |
    +-- 索引器获取并处理数据
        +-- 从 Celestia 读取 blob
        +-- 验证承诺匹配
        +-- 解析交易并索引
    |
阶段 5: 数据归档
    |
    +-- 热存储(PostgreSQL, 30 天)
    +-- 温存储(S3 Parquet, 1 年)
    +-- 冷存储(Arweave, 永久)
    |
阶段 6: 数据过期(可选)
    |
    +-- Chirugical 裁剪
    +-- DA 层数据自动过期(30 天)
    +-- MSG Chain 元数据标记为 pruned

8.2.2 端到端代码示例

// DaExample 展示完整的 DA blob 生命周期
func DaExample(ctx context.Context) error {
    // 1. 初始化客户端
    celestiaClient, err := celestia.NewClient("http://localhost:26658")
    if err != nil {
        return fmt.Errorf("init celestia: %w", err)
    }

    msgchainClient, err := msgchain.NewClient("https://rpc.msgchain.org:26657")
    if err != nil {
        return fmt.Errorf("init msgchain: %w", err)
    }

    registryAddr := "msg1registry..."

    // 2. 构建数据
    data := map[string]interface{}{
        "app":      "cosmwasm-example",
        "action":   "batch_mint",
        "tokens":   []string{"token1", "token2", "token3"},
        "sender":   "msg1sender...",
        "timestamp": time.Now().Unix(),
    }
    dataBytes, _ := json.Marshal(data)

    // 3. 压缩
    compressed := snappy.Encode(nil, dataBytes)

    // 4. 提交到 Celestia
    namespace := [8]byte{0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x01, 0x01}
    blob, _ := celestia.NewBlob(namespace, compressed)
    height, err := celestiaClient.SubmitBlob(ctx, []*celestia.Blob{blob}, auto)
    if err != nil {
        return fmt.Errorf("submit blob: %w", err)
    }

    commitment := sha256.Sum256(compressed)

    fmt.Printf("Submitted at height %d, commitment: %x\n", height, commitment)

    // 5. 在 MSG Chain 注册
    registerMsg := map[string]interface{}{
        "register": map[string]interface{}{
            "metadata": map[string]interface{}{
                "da_layer":      "celestia",
                "da_network":    "mainnet",
                "da_height":     height,
                "commitment":    hex.EncodeToString(commitment[:]),
                "namespace":     hex.EncodeToString(namespace[:]),
                "contract_addr": "msg1cosmwasm...",
                "data_format":   "json",
                "size":          len(dataBytes),
                "data_uri":      "",
                "status":        "active",
            },
        },
    }

    registerMsgBytes, _ := json.Marshal(registerMsg)
    tx, err := msgchainClient.ExecuteContract(
        ctx, registryAddr, registerMsgBytes, "umsg", 5000,
    )
    if err != nil {
        return fmt.Errorf("register on msgchain: %w", err)
    }

    fmt.Printf("Registered in MSG Chain tx: %s\n", tx.TxHash)

    // 6. 从 Celestia 读取验证
    time.Sleep(12 * time.Second) // 等待一个区块

    blobs, err := celestiaClient.GetBlobs(ctx, height, []celestia.Namespace{namespace})
    if err != nil {
        return fmt.Errorf("read from celestia: %w", err)
    }

    retrievedHash := sha256.Sum256(blobs[0].Data)
    if !bytes.Equal(retrievedHash[:], commitment[:]) {
        return fmt.Errorf("commitment mismatch")
    }

    fmt.Println("Data verified successfully!")
    return nil
}

8.2.3 部署检查清单

## 集成部署检查清单

### 前置检查
- [ ] Celestia 主网/测试网连接正常
- [ ] MSG Chain 节点 >= v1.0.0
- [ ] DA 模块已启用(app.da_enabled = true)
- [ ] CosmWasm 虚拟机已启用

### 合约部署
- [ ] DA 注册中心合约已编译并审计
- [ ] 合约已部署到 MSG Chain(code_id 已确认)
- [ ] 合约已实例化(contract_addr 已确认)
- [ ] 初始配置正确(管理员、白名单等)

### DA Router 部署
- [ ] DA Router 二进制已构建
- [ ] 配置文件已创建(RPC、密钥、参数)
- [ ] Celestia 密钥已创建并有余额
- [ ] MSG Chain 操作员密钥已配置
- [ ] DA Router 已启动并同步事件

### 命名空间
- [ ] 每个合约已分配唯一命名空间
- [ ] 命名空间未与系统保留范围冲突

### 监控
- [ ] blob 提交延迟监控
- [ ] MSG Chain 注册事件监控
- [ ] DA 采样成功率监控
- [ ] 费用支出监控(平衡 DA 和 MSG Chain 费用)
- [ ] 归档存储容量监控

### 测试
- [ ] 单元测试通过
- [ ] 集成测试:端到端 blob 生命周期
- [ ] 压力测试:大 blob(>1 MB)提交
- [ ] 故障测试:DA Router 重启恢复
- [ ] 降级测试:DA 层不可用时的链行为

9. 安全与经济

9.1 DA 层信任假设

9.1.1 信任模型分类

层 信任假设 安全保证 风险描述
MSG Chain (CometBFT) 假设 <= 1/3 验证人合谋 最终确定性 + 状态一致性 1/3+ 合谋可回滚或分叉
Celestia (Tendermint) 假设 <= 1/3 验证人合谋 数据排序 + 可用性保证 1/3+ 合谋可审查或作恶
Avail (KZG + 验证人) 假设 <= 1/3 验证人合谋 数据可用性 + 即时验证 KZG 安全依赖可信设置
EigenDA (Restaking) 假设 ETH 再质押的安全 经济安全 + 可用性 Slashing 机制的有效性

9.1.2 数据扣留攻击

攻击描述:

恶意排序器在 DA 层提交区块头和数据承诺,但不公布完整的原始数据。全节点无法重放交易,也无法生成状态根。

防御机制:

  1. 纠删码: 即使部分数据分片丢失,通过 Reed-Solomon 编码仍可恢复原始数据。
  2. 数据采样: 轻客户端通过随机采样以高概率检测数据扣留。
  3. 经济惩罚: 排序器需质押大量代币,数据扣留行为导致质押被罚没。
  4. 挑战窗口: 在挑战窗口期内,任何人都可以提交数据不可用的欺诈证明。

安全参数建议:

var SecurityRecommendations = struct {
    MinSamplingCount          int
    MinStakeAmount             sdk.Coin
    ChallengeWindowSeconds    uint64
    FraudProverRewardRatio    float64
}{
    MinSamplingCount:          100,
    MinStakeAmount:            sdk.NewCoin("umsg", sdk.NewInt(100_000_000)),
    ChallengeWindowSeconds:    604800,
    FraudProverRewardRatio:    0.10,
}

9.1.3 共识活跃性依赖

DA 层的安全不仅取决于验证人诚实性,还取决于网络的活跃性。如果 DA 层停止出块,MSG Chain 的 DA 层集成也将暂停。为缓解这一风险:

9.1.4 跨层合谋风险

最严重的安全威胁是 MSG Chain 和 DA 层验证人之间的合谋:

攻击场景:
1. MSG Chain 验证人接收 DA 层验证人的贿赂
2. 合谋构造虚假的区块头和数据承诺
3. 轻客户端被误导,确信不存在的状态
4. 攻击者窃取用户资产

防御:
1. 两个网络的验证人集合应尽量保持独立
2. 使用多 DA 层分散信任
3. 客户端侧增加交叉验证机制
4. 增加经济惩罚的力度(高昂的合谋成本)

9.2 费用模型

9.2.1 DA 层费用结构

费用组件 Celestia Avail EigenDA
基础费用 Gas * GasPrice 固定 + 变动 固定 + 变动
按数据大小 10 gas/byte 价格/byte 价格/byte
按区块大小 50k gas 基础 区块饱和溢价 动态调整
优先级费用 可选 tips 无 可选
费用代币 TIA (utia) AVAIL ETH

9.2.2 MSG Chain 的 DA 费用分摊

在 MSG Chain 中,DA 费用由多个参与者按比例分摊:

#[cw_serde]
pub struct DaFeeDistribution {
    /// 用户支付的 DA 费用
    pub user_fee: Coin,
    /// 排序器补贴(如果应用提供补贴)
    pub sequencer_subsidy: Coin,
    /// MSG Chain 协议补贴(可选)
    pub protocol_subsidy: Coin,
    /// 总 DA 成本
    pub total_da_cost: Coin,
    /// 支付给 DA 层的实际金额
    pub da_layer_payment: Coin,
    /// MSG Chain 注册中心费用
    pub registry_fee: Coin,
}

impl DaFeeDistribution {
    pub fn calculate(
        user_fee: Coin,
        blob_size: u64,
        sequencer_subsidy: Option<Coin>,
    ) -> Self {
        // 估算 DA 层成本
        let estimated_da_cost = estimate_celestia_cost(blob_size);

        // 注册中心费用(固定)
        let registry_fee = Coin::new(100_000, "umsg");

        // 总成本
        let total = estimated_da_cost + registry_fee;

        // 如果用户费用不够,由排序器补贴
        let shortfall = if user_fee.amount < total.amount {
            total.amount - user_fee.amount
        } else {
            sdk::Uint128::zero()
        };

        DaFeeDistribution {
            user_fee,
            sequencer_subsidy: Coin::new(shortfall, "umsg"),
            protocol_subsidy: Coin::new(0, "umsg"),
            total_da_cost: total,
            da_layer_payment: estimated_da_cost,
            registry_fee,
        }
    }
}

9.2.3 费用优化策略

// FeeOptimizer DA 费用优化器
type FeeOptimizer struct {
    daClient    *celestia.Client
    stats       *FeeStatistics
}

// FeeStatistics 费用统计数据
type FeeStatistics struct {
    AvgGasPrice          float64
    AvgBlobSize          int64
    FeePerByte           float64
    RecommendedTipRatio  float64
}

// EstimateOptimalFee 估算最优费用
func (fo *FeeOptimizer) EstimateOptimalFee(dataSize int) (*FeeEstimate, error) {
    stats, err := fo.daClient.GetFeeStats(context.Background())
    if err != nil {
        return nil, err
    }

    // 基础 gas
    baseGas := uint64(50000)
    dataGas := uint64(dataSize) * 10
    totalGas := baseGas + dataGas

    // 动态加价(根据网络拥堵)
    multiplier := 1.0
    if stats.CongestionLevel > 0.8 {
        multiplier = 1.5
    } else if stats.CongestionLevel > 0.5 {
        multiplier = 1.2
    }

    estimatedFee := uint64(float64(totalGas) * stats.AverageGasPrice * multiplier)

    return &FeeEstimate{
        DataSize:       dataSize,
        TotalGas:       totalGas,
        GasPrice:       stats.AverageGasPrice,
        EstimatedFee:   estimatedFee,
        Currency:       "utia",
        Congestion:     stats.CongestionLevel,
    }, nil
}

9.3 数据保留策略

9.3.1 多级保留策略

MSK Chain 的 DA 数据保留策略分为多个层级:

层级 保留期 存储介质 访问延迟 成本 数据完整性
L1 热存储 30 天 PostgreSQL + SSD < 10ms 高 SHA-256 校验
L2 温存储 1 年 S3 / MinIO + Parquet < 100ms 中 Parquet CRC
L3 冷存储 永久 Arweave / Filecoin > 1 秒 低 Arweave 共识
L4 精简 元数据永久 MSG Chain 注册中心 < 50ms 最低 CometBFT 共识

9.3.2 数据裁剪机制

#[cw_serde]
pub struct PrunePolicy {
    /// 自动裁剪启用
    pub auto_prune: bool,
    /// 数据过期后保留的宽限期(纳秒)
    pub grace_period_nanos: u64,
    /// 每次裁剪的最大条数
    pub max_prune_per_batch: u32,
    /// 裁剪奖励(gas 返还比例)
    pub prune_reward_ratio: Decimal,
}

#[entry_point]
pub fn execute_prune_expired(
    deps: DepsMut,
    env: Env,
) -> Result<Response, ContractError> {
    let policy = PRUNE_POLICY.load(deps.storage)?;
    if !policy.auto_prune {
        return Ok(Response::new().add_attribute("action", "prune_skipped"));
    }

    let cutoff = env.block.time.nanos() - policy.grace_period_nanos;
    let mut pruned_count = 0u32;
    let mut gas_saved = 0u64;

    // 扫描并裁剪过期数据
    let to_remove: Vec<(Vec<u8>, Vec<u8>)> = BLOB_REGISTRY
        .range(deps.storage, None, None, cosmwasm_std::Order::Ascending)
        .filter_map(|item| {
            if let Ok((key, metadata)) = item {
                if let Some(expires_at) = metadata.expires_at {
                    if expires_at < cutoff {
                        return Some(key);
                    }
                }
            }
            None
        })
        .take(policy.max_prune_per_batch as usize)
        .collect();

    for (namespace, commitment) in &to_remove {
        BLOB_REGISTRY.remove(deps.storage, (namespace.as_slice(), commitment.as_slice()));
        pruned_count += 1;
        gas_saved += 100000; // 估算的存储释放 gas
    }

    if pruned_count > 0 {
        ACTIVE_BLOB_COUNT.update(deps.storage, |c| c - pruned_count as u64)?;
    }

    Ok(Response::new()
        .add_attribute("action", "prune_expired")
        .add_attribute("pruned_count", pruned_count.to_string())
        .add_attribute("gas_saved", gas_saved.to_string()))
}

9.3.3 监管与合规

对于需要满足数据保留监管要求的场景:

#[cw_serde]
pub struct ComplianceConfig {
    /// 是否启用合规模式
    pub enabled: bool,
    /// 最短保留期限(纳秒)
    pub min_retention_period: u64,
    /// 禁止裁剪的数据标志
    pub retention_hold: bool,
    /// 合规审核地址列表
    pub auditors: Vec<Addr>,
}

impl ComplianceConfig {
    pub fn can_prune(&self, metadata: &BlobMetadata, current_time: u64) -> bool {
        if !self.enabled {
            return true; // 非合规模式,可裁剪
        }
        if self.retention_hold {
            return false; // 全局保留 hold
        }
        if let Some(expires_at) = metadata.expires_at {
            return current_time >= expires_at + self.min_retention_period;
        }
        false
    }
}

10. 总结与路线图

10.1 核心设计总结

MSG Chain 的数据可用性层集成架构以模块化为核心设计原则,实现了以下关键能力:

维度 设计决策 核心收益
存储模型 混合存储(链上 + 链下 DA) 安全与扩展的平衡
集成模式 主权 Rollup 优先 独立治理 + 低费用
数据映射 CosmWasm 注册中心合约 链上验证 DA 数据
验证机制 DAS 采样 + 欺诈证明 无需全节点验证
索引架构 监听 + 拉取模式 自动数据归档
DA 层支持 Celestia 为主,多 DA 灵活选择与降级

10.2 技术路线图

阶段一:基础集成(已完成)

阶段二:生产化(当前)

阶段三:生态扩展(2026 Q3-Q4)

阶段四:高级功能(2027 Q1+)

10.3 关键里程碑

里程碑 预期时间 描述
v1.0 DA 主网启动 2026-08 Celestia DA 集成在生产环境可用
v1.1 多 DA 支持 2026-10 添加 Avail 作为替代 DA 层
v1.2 Validium 2026-12 ZK 证明验证合约上线
v2.0 Volition 2027-03 混合模式支持上线
v2.1 去中心化排序 2027-06 排序器去中心化 + 代币激励

10.4 风险与缓解

风险 影响 概率 缓解措施
Celestia 网络拥堵 交易延迟 中 多 DA 回退 + 费用动态调整
DA 层验证人合谋 数据扣留 低 高质押门槛 + 多 DA 分散
合约安全漏洞 资产损失 低 多次独立审计 + 渐进上线
存储成本增长 运行成本高 中 归档策略 + 激励社区节点
监管要求变化 合规风险 低 合规模块 + 数据保留选项

文档维护: 本文档由 MSG Chain 核心开发团队维护。
问题反馈: https://github.com/msgchain/docs/issues
社区讨论: https://forum.msgchain.org/c/da-integration