数据可用性层(DA)集成指南
数据来源:MSG Chain 代码库核实
主网状态: No-Go — 当前 MSGChain 主网裁决为 No-Go,以下内容反映代码实际状态,不代表生产可用。
目录
- 数据可用性概念
- 1.1 DA 在区块链架构中的位置
- 1.2 DA 与存储的区别
- 1.3 Celestia / Avail / EigenDA 生态概览
- MSG Chain 的 DA 策略
- 2.1 链上存储 vs 链下 DA
- 2.2 注册中心的 metadata 存储模式
- DA 层集成架构
- 3.1 主权 rollup 模式
- 3.2 validium 模式
- 3.3 volition 模式
- 数据发布与验证
- 4.1 blob 提交
- 4.2 数据采样
- 4.3 KZG 承诺
- 4.4 欺诈证明
- CosmWasm 合约与 DA 交互
- 5.1 contract 将数据发布到 DA
- 5.2 从 DA 读取数据验证
- 轻客户端验证
- 6.1 轻节点 DA 采样
- 6.2 无需全节点即可验证数据可用性
- 索引器与 DA
- 7.1 索引器从 DA 层读取数据
- 7.2 归档存储
- 7.3 与 MSG Chain 事件流的整合
- 集成实践
- 8.1 使用 Celestia 作为 MSG Chain 的 DA 层
- 8.2 blob 生命周期示例
- 安全与经济
- 9.1 DA 层信任假设
- 9.2 费用模型
- 9.3 数据保留策略
- 总结与路线图
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 层承担的关键职能包括:
- 交易数据托管: 将执行层产生的交易数据发布到专用的 DA 网络,而非全部存储在 MSG Chain 主链上。
- 数据完整性保证: 通过纠删码(Erasure Coding)和数据采样(Data Availability Sampling,DAS)机制,确保数据在网络中完整可用。
- 可扩展性提升: 将数据存储从共识层卸载到 DA 层,显著降低主链的存储负担和验证开销。
- 跨链互操作性: 标准化的 DA 接口允许不同执行层(如 CosmWasm 合约、EVM Rollup)共享同一数据可用性基础设施。
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 项目可能使用以下分层数据策略:
- 交易数据: 通过 Celestia DA 层发布,仅在 Rollup 挑战期内保证可用。
- NFT 元数据: 存储在 IPFS 或 Arweave 中,保证长期可访问。
- 合约状态: 由 MSG Chain 全节点维护,通过 Tendermint 共识保证一致性。
- 链上事件: 由索引器(如 CosmWasm 索引器)归档到关系数据库或对象存储。
1.3 Celestia / Avail / EigenDA 生态概览
当前主流的 DA 层解决方案包括 Celestia、Avail 和 EigenDA,三者在架构设计、共识机制和生态集成方面各有特点。
1.3.1 Celestia
Celestia 是最早提出模块化区块链和 DA 分离的项目,其核心架构基于 Tendermint 共识和 Namespaced Merkle Trees(NMT)。
关键技术特性:
- 共识机制: 基于 Tendermint 的 DPoS(Delegated Proof of Stake),验证人集合通过质押 $TIA 代币参与共识。
- 数据分片: 使用二维 Reed-Solomon 纠删码将 blob 数据编码为矩形矩阵,支持数据采样。
- 命名空间(Namespace): NMT 允许各执行链(Rollup)在不扫描全局数据树的情况下,仅检索与其命名空间相关的数据。
- 数据可用性采样(DAS): 轻客户端通过对随机选择的坐标进行采样,以高概率验证数据可用性。
- 量子引力桥(Quantum Gravity Bridge): Celestia 与以太坊之间的跨链桥,允许以太坊 L2 使用 Celestia 作为 DA。
对 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)结合数据可用性采样。
关键技术特性:
- Kate 多项式承诺: 使用 KZG 承诺(基于椭圆曲线配对)替代 Merkle 树,实现恒定大小的数据承诺证明。
- 有效性证明(Validity Proofs): 使用零知识证明验证纠删码编码的正确性,消除欺诈证明的延迟窗口。
- App ID: 每个应用程序(Rollup)被分配一个 App ID,类似于 Celestia 的命名空间。
- 轻客户端: Avail 的轻客户端实现数据可用性采样,使用 KZG 承诺的简洁性减少带宽开销。
- 数据证明(Data Attestation): 验证人节点对数据可用性签署证明,形成数据可用性证书(Data Attestation Certificate)。
对 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 的再质押验证者网络提供数据可用性保证。
关键技术特性:
- 再质押安全模型: 验证者将以太坊上的 ETH 通过 EigenLayer 再质押到 EigenDA,通过经济惩罚(Slashing)保证诚实行为。
- 水平扩展: EigenDA 采用分散式验证者集,支持水平扩展以处理更大的 blob 吞吐量。
- Blob 生命周期: 客户端提交 blob 到 Disperser 服务,Disperser 将 blob 分发到验证者网络,验证者签署可用性证明。
- 以太坊集成: 与以太坊生态深度集成,尤其适合基于以太坊的应用链和 Rollup。
对 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 架构) |
| 经济成本 | 低 | 低 | 中(含再质押成本) |
| 去中心化程度 | 高 | 高 | 中 |
| 主网上线状态 | 已主网上线 | 已主网上线 | 已主网上线 |
推荐策略:
- 优先集成 Celestia: 作为 MSG Chain 的首要 DA 层,因为两者共享 Tendermint 共识基础,集成复杂度最低。
- 渐进支持 Avail: 当需要 KZG 承诺的隐私特性或零知识证明集成时,引入 Avail 作为替代或补充 DA 层。
- 战略储备 EigenDA: 当需要与以太坊生态深度桥接时,考虑 EigenDA,但需评估再质押的额外信任假设。
2. MSG Chain 的 DA 策略
2.1 链上存储 vs 链下 DA
MSG Chain 在设计 DA 策略时面临核心权衡:是将所有数据存储在链上(由 CometBFT 共识维护),还是将数据可用性卸载到专用 DA 层。这两种策略在 MSG Chain 的架构中各有适用场景。
2.1.1 链上存储模型
在链上存储模型中,MSG Chain 的全节点存储所有交易数据,数据可用性由链自身的共识协议保证。
优势:
- 单信任域: 所有数据在同一个共识网络中验证,没有外部信任假设。
- 简化架构: 无需额外的 DA 客户端、桥接和验证逻辑。
- 低延迟: 数据立即可用,无需等待外部 DA 层的确认。
- 固有安全: 利用 CometBFT 的拜占庭容错特性,不存在 DA 层欺诈的风险。
劣势:
- 存储爆炸: 每个全节点必须存储所有历史数据,限制了链的吞吐量和参与节点数量。
- 高成本: 验证人需支付高昂的存储成本,这些成本最终转嫁到用户交易费用上。
- 扩展瓶颈: 区块大小受限于全节点的带宽和存储能力,无法支持大规模数据密集型应用。
- 数据冗余: 所有数据在所有全节点间完全复制,造成资源浪费。
在 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 验证人存储限制,可支持 GB 级数据块。
- 低成本: 用户仅支付 DA 层的发布费用,远低于链上存储成本。
- 数据隔离: 不同应用的数据在 DA 层通过命名空间隔离,互不干扰。
- 灵活验证: 轻客户端通过数据采样验证可用性,无需运行全节点。
- 模块化升级: DA 层可与执行层独立升级,支持热插拔式替换。
劣势:
- 信任扩展: 引入 DA 层的信任假设(DA 层验证人集的安全性和活跃性)。
- 额外延迟: 数据需先发布到 DA 层并获得确认,然后才能在 MSG Chain 上提交承诺。
- 桥接复杂度: 需要 DA 层桥接客户端来验证和路由数据承诺。
- 数据可用性风险: 如果 DA 层发生数据扣留攻击,用户可能无法访问交易数据。
在 MSG Chain 中的适用场景:
链下 DA 适用于 MSG Chain 上的数据密集型应用,包括:
- 大规模 NFT 铸造(mint)事件的交易数据
- 去中心化社交媒体平台的内容发布
- 游戏状态更新和动作记录
- 大数据集的预言机价格馈送
- Rollup 批次数据提交
典型示例——发布交易数据到 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 注册中心合约的核心职责包括:
- 数据锚定: 记录每个 DA 发布事件与 MSG Chain 交易之间的关联。
- 元数据索引: 存储 blob 的描述信息,包括命名空间、数据格式、关联合约等。
- 验证路由: 向轻客户端和索引器提供验证 DA 数据所需的证明信息。
- 生命周期管理: 跟踪数据的过期时间、归档状态和访问权限。
注册中心合约接口:
#[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 提供了以下关键优势:
- 主权独立: MSG Chain 控制自己的状态转换逻辑,不依赖 DA 层进行链的升级或治理。
- 硬分叉自由: 社区可以在不依赖 DA 层的情况下进行链的升级或分叉。
- 低费用: 利用 Celestia 的低成本 DA,交易费用远低于传统 L1。
- 可组合性: 同一个 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 应用场景:
- 去中心化交易所: 大额交易使用 Rollup 模式(更高安全性),小额频繁交易使用 Validium 模式(更快确认)。
- 游戏平台: 资产转移走 Rollup,游戏状态更新走 Validium。
- 社交应用: 关键身份操作走 Rollup,内容发布走 Validium。
- 企业合规: 公开交易走 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 层提交区块头和数据承诺,但不公布完整的原始数据。全节点无法重放交易,也无法生成状态根。
防御机制:
- 纠删码: 即使部分数据分片丢失,通过 Reed-Solomon 编码仍可恢复原始数据。
- 数据采样: 轻客户端通过随机采样以高概率检测数据扣留。
- 经济惩罚: 排序器需质押大量代币,数据扣留行为导致质押被罚没。
- 挑战窗口: 在挑战窗口期内,任何人都可以提交数据不可用的欺诈证明。
安全参数建议:
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 层集成也将暂停。为缓解这一风险:
- MSG Chain 可以配置多 DA 层回退机制
- 在设置降级模式下,当 DA 层不可用时,允许紧急回到链上存储模式
- 设置 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 技术路线图
阶段一:基础集成(已完成)
- [x] MSG Chain 与 Celestia DA 层的连接建立
- [x] DA 注册中心合约的设计和基础实现
- [x] DA Router 链外服务原型
- [x] blob 提交流程和基本验证
阶段二:生产化(当前)
- [ ] DA 注册中心合约的安全审计
- [ ] DA Router 的高可用部署(多实例 + 负载均衡)
- [ ] 轻客户端 DA 采样模块集成到 msg-chaind
- [ ] 索引器 DA 数据获取管道
- [ ] 归档存储基础架构(热 + 温 + 冷)
阶段三:生态扩展(2026 Q3-Q4)
- [ ] Avail DA 层的集成适配
- [ ] 多 DA 层自动故障切换
- [ ] CosmWasm 合约直接 DA 交互的 SDK 发布
- [ ] DA 费用分摊机制上线
- [ ] 社区运行的归档节点网络
阶段四:高级功能(2027 Q1+)
- [ ] Validium 模式支持(ZK 证明集成)
- [ ] Volition 模式支持(混合数据可用性)
- [ ] KZG 承诺在 CosmWasm 中的原生预编译支持
- [ ] DA 数据压缩标准(降低费用)
- [ ] 跨链 DA 验证桥(IBC 集成)
- [ ] 去中心化排序器网络
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
