MSG Chain AI Agent 去中心化存储集成指南
⚠️ No-Go Disclaimer: MSGChain 主网裁决为 No-Go。本文件所有内容反映的是开发阶段的技术设计,不代表主网未独立核验上线状态。生产部署状态请以白皮书为准:https://msgchain.org/whitepaper/
目录
1. 概述
1.1 为什么 AI Agent 需要去中心化存储
AI Agent 在运行过程中产生大量数据:对话历史、知识库、模型状态、任务日志、用户偏好等。传统中心化存储方案存在以下问题:
- 单点故障:服务器宕机导致 Agent 无法访问关键数据
- 审查风险:中心化服务商可删除或限制数据访问
- 数据主权:用户和 Agent 无法真正拥有自己的数据
- 可验证性:缺乏对数据完整性的密码学证明
去中心化存储解决了这些问题,为 AI Agent 提供了:
- 抗审查:数据一旦存储,无法被单一方删除
- 内容寻址:通过内容哈希(CID)而非位置寻址,保证数据完整性
- 持久性保证:通过经济激励确保数据长期可用
- 可验证性:任何节点可独立验证数据的完整性
1.2 MSG Chain 与去中心化存储的关系
MSG Chain(msg-chain-1,bech32 前缀 msg)是专为 AI Agent 设计的 Layer 1 区块链。它在去中心化存储生态中扮演以下角色:
| 角色 | 说明 |
|---|---|
| 锚定层 | 将链下存储的 CID 锚定到链上,建立链下数据的永久引用 |
| 索引层 | 维护链下数据的元数据索引(所有者、类型、时间戳等) |
| 激励层 | 通过代币经济激励存储节点和数据提供者 |
| 验证层 | 提供链上验证机制以验证链下数据的完整性和可用性 |
1.3 三大存储方案对比
| 特性 | IPFS | Arweave | Filecoin |
|---|---|---|---|
| 寻址方式 | 内容寻址(CID) | 内容寻址(TX ID) | 内容寻址(CID) |
| 持久性 | 取决于 pinning 服务 | 一次性付费,永久存储 | 通过存储合约(Deal)保证 |
| 成本模型 | 按带宽/存储付费 | 一次性预付费 | 按时间/空间付费 |
| 访问速度 | 取决于节点距离 | 中等到快速 | 需要检索交易 |
| 数据新鲜度 | 可变内容支持(MFS) | 不可变 | 不可变 |
| 最适合 | 热/温数据访问 | 归档/永久保存 | 大规模归档+可证明存储 |
1.4 选型决策流程
Agent 产生数据
│
├── 需要频繁访问 (< 1小时) → 本地缓存 (Hot Tier)
│
├── 需要定期访问 + 持久 → IPFS + Pin (Warm Tier)
│
├── 需要永久保存 + 不变 → Arweave (Cold Tier)
│
└── 大规模归档 + 可证明 → Filecoin (Archive Tier)
1.5 前置条件
在开始本指南之前,请确保:
- 安装 Python 3.10+ 和 Rust 1.70+
- 拥有 MSG Chain 账户(使用
msg1前缀地址) - 安装
msgchaindCLI 工具(来自 msg-chain-1 网络) - 基本的智能合约开发知识
- 熟悉 IPFS 节点操作基础概念
2. IPFS 集成
2.1 IPFS 简介
IPFS(InterPlanetary File System)是一种点对点的分布式文件系统,通过内容寻址(CID)来标识文件。对于 AI Agent 而言,IPFS 提供了:
- 内容完整性:CID 是内容的哈希值,任何修改都会导致 CID 变化
- 去重:相同内容只存储一次
- 分布式访问:数据可从任何拥有它的节点获取
2.2 Python IPFS 客户端
import asyncio
import json
import hashlib
import ipfshttpclient
from typing import Optional, List, Dict, Any
from dataclasses import dataclass, asdict
from datetime import datetime
import aiohttp
# MSG Chain bech32 前缀
MSG_CHAIN_PREFIX = 'msg'
@dataclass
class IPFSContent:
"""IPFS 内容元数据"""
cid: str
size: int
content_type: str
created_at: str
checksum: str
class IPFSClient:
"""
IPFS 基础客户端
连接到本地或远程 IPFS 节点
"""
def __init__(self, node_url: str = 'http://localhost:5001',
timeout: int = 120):
self.node_url = node_url
self.timeout = timeout
self._client = None
async def connect(self):
"""异步连接到 IPFS 节点"""
loop = asyncio.get_event_loop()
self._client = await loop.run_in_executor(
None,
lambda: ipfshttpclient.connect(self.node_url)
)
return self._client
async def store_json(self, data: dict) -> IPFSContent:
"""
存储 JSON 数据到 IPFS
参数:
data: 要存储的字典数据
返回:
IPFSContent 对象包含 CID 和元数据
"""
if self._client is None:
await self.connect()
serialized = json.dumps(data, ensure_ascii=False, default=str)
content_bytes = serialized.encode('utf-8')
loop = asyncio.get_event_loop()
result = await loop.run_in_executor(
None,
lambda: self._client.add_bytes(content_bytes)
)
cid = result.decode() if isinstance(result, bytes) else result
checksum = hashlib.sha256(content_bytes).hexdigest()
return IPFSContent(
cid=cid,
size=len(content_bytes),
content_type='application/json',
created_at=datetime.utcnow().isoformat(),
checksum=checksum,
)
async def store_agent_data(self, data: dict) -> str:
"""
存储 Agent 数据并返回 CID
参数:
data: Agent 数据字典
返回:
内容标识符 CID 字符串
"""
content = await self.store_json(data)
enriched_data = {
'msg_chain_version': 'msg-chain-1',
'bech32_prefix': MSG_CHAIN_PREFIX,
'content': data,
'metadata': {
'cid': content.cid,
'checksum': content.checksum,
'stored_at': content.created_at,
},
}
serialized = json.dumps(enriched_data, ensure_ascii=False, default=str)
content_bytes = serialized.encode('utf-8')
loop = asyncio.get_event_loop()
result = await loop.run_in_executor(
None,
lambda: self._client.add_bytes(content_bytes)
)
return result.decode() if isinstance(result, bytes) else result
async def store_bytes(self, data: bytes, filename: str = '') -> str:
"""存储原始字节到 IPFS"""
if self._client is None:
await self.connect()
loop = asyncio.get_event_loop()
result = await loop.run_in_executor(
None,
lambda: self._client.add_bytes(data)
)
return result.decode() if isinstance(result, bytes) else result
async def get_json(self, cid: str) -> dict:
"""从 IPFS 获取 JSON 数据"""
if self._client is None:
await self.connect()
loop = asyncio.get_event_loop()
data = await loop.run_in_executor(
None,
lambda: self._client.cat(cid)
)
return json.loads(data.decode('utf-8'))
async def get_agent_data(self, cid: str) -> dict:
"""
获取 Agent 数据(兼容带元数据的格式)
参数:
cid: 内容标识符
返回:
Agent 数据字典
"""
data = await self.get_json(cid)
if 'msg_chain_version' in data and 'content' in data:
return data['content']
return data
async def pin_cid(self, cid: str, recursive: bool = True):
"""固定 CID 到本地节点"""
if self._client is None:
await self.connect()
loop = asyncio.get_event_loop()
await loop.run_in_executor(
None,
lambda: self._client.pin.add(cid, recursive=recursive)
)
async def unpin_cid(self, cid: str):
"""取消固定 CID"""
if self._client is None:
await self.connect()
loop = asyncio.get_event_loop()
await loop.run_in_executor(
None,
lambda: self._client.pin.rm(cid)
)
async def list_pins(self) -> List[str]:
"""列出所有固定的 CID"""
if self._client is None:
await self.connect()
loop = asyncio.get_event_loop()
result = await loop.run_in_executor(
None,
lambda: self._client.pin.ls()
)
return list(result.get('Keys', {}).keys())
async def resolve(self, cid: str) -> Optional[bytes]:
"""解析 CID 对应的原始内容"""
try:
if self._client is None:
await self.connect()
loop = asyncio.get_event_loop()
data = await loop.run_in_executor(
None,
lambda: self._client.cat(cid)
)
return data
except Exception:
return None
async def close(self):
"""关闭客户端连接"""
if self._client is not None:
self._client.close()
self._client = None
class IPFSClientTypeScript:
"""
TypeScript IPFS 客户端的 Python 参考实现
对应 JS 代码结构供前端使用
"""
@staticmethod
def get_js_imports() -> str:
return """
// TypeScript IPFS 客户端
import { create } from 'ipfs-http-client'
import { CID } from 'multiformats/cid'
import { sha256 } from 'multiformats/hashes/sha2'
const MSG_CHAIN_PREFIX = 'msg'
"""
@staticmethod
def get_js_client_code() -> str:
return """
class IPFSClient {
private node: any;
private nodeUrl: string;
constructor(nodeUrl: string = 'http://localhost:5001') {
this.nodeUrl = nodeUrl;
}
async connect(): Promise<void> {
this.node = await create({ url: this.nodeUrl });
}
async storeAgentData(data: Record<string, any>): Promise<string> {
const enriched = {
msg_chain_version: 'msg-chain-1',
bech32_prefix: MSG_CHAIN_PREFIX,
content: data,
metadata: {
stored_at: new Date().toISOString(),
},
};
const { cid } = await this.node.add(JSON.stringify(enriched));
return cid.toString();
}
async getAgentData(cid: string): Promise<Record<string, any>> {
const chunks = [];
for await (const chunk of this.node.cat(cid)) {
chunks.push(chunk);
}
const data = JSON.parse(Buffer.concat(chunks).toString());
if (data.msg_chain_version && data.content) {
return data.content;
}
return data;
}
async pin(cid: string): Promise<void> {
await this.node.pin.add(cid);
}
}
"""
2.3 IPFS 集群客户端
import asyncio
import json
from typing import List, Optional, Dict
import aiohttp
from dataclasses import dataclass
@dataclass
class ClusterStatus:
"""IPFS 集群状态"""
peer_id: str
peers: List[str]
pin_count: int
cluster_healthy: bool
class IPFSClusterClient:
"""
IPFS 集群客户端
提供高可用的 pinning 服务,跨多节点复制
"""
def __init__(self, cluster_endpoints: List[str],
replication_factor: int = 3):
self.endpoints = cluster_endpoints
self.replication_factor = replication_factor
self._session: Optional[aiohttp.ClientSession] = None
async def _ensure_session(self):
if self._session is None:
self._session = aiohttp.ClientSession()
async def replicate(self, cid: str, replicas: Optional[int] = None) -> Dict:
"""
在集群中复制 CID
参数:
cid: 要复制的内容标识符
replicas: 副本数量(默认使用集群设置)
返回:
操作状态字典
"""
await self._ensure_session()
r = self.replication_factor if replicas is None else replicas
payload = {
'cid': cid,
'replication_factor_min': r,
'replication_factor_max': r + 1,
}
results = {}
for endpoint in self.endpoints:
try:
async with self._session.post(
f'{endpoint}/pins/{cid}',
json=payload,
timeout=aiohttp.ClientTimeout(total=60),
) as resp:
results[endpoint] = await resp.json()
except Exception as e:
results[endpoint] = {'error': str(e)}
return results
async def replicate_with_policy(self, cid: str, policy: Dict) -> Dict:
"""
按策略复制 CID
策略示例:
{
'replicas': 3,
'regions': ['us-east', 'eu-west', 'ap-southeast'],
'min_providers': 2,
}
"""
await self._ensure_session()
payload = {
'cid': cid,
**policy,
}
results = {}
for endpoint in self.endpoints:
try:
async with self._session.post(
f'{endpoint}/pins/{cid}',
json=payload,
timeout=aiohttp.ClientTimeout(total=120),
) as resp:
results[endpoint] = await resp.json()
except Exception as e:
results[endpoint] = {'error': str(e)}
return results
async def unpin(self, cid: str) -> Dict:
"""从集群中取消固定 CID"""
await self._ensure_session()
results = {}
for endpoint in self.endpoints:
try:
async with self._session.delete(
f'{endpoint}/pins/{cid}',
timeout=aiohttp.ClientTimeout(total=30),
) as resp:
results[endpoint] = await resp.json()
except Exception as e:
results[endpoint] = {'error': str(e)}
return results
async def list_pins(self) -> Dict:
"""列出集群中所有 pin"""
await self._ensure_session()
endpoint = self.endpoints[0]
async with self._session.get(
f'{endpoint}/pins',
timeout=aiohttp.ClientTimeout(total=30),
) as resp:
return await resp.json()
async def status(self) -> ClusterStatus:
"""获取集群状态"""
await self._ensure_session()
endpoint = self.endpoints[0]
async with self._session.get(
f'{endpoint}/health',
timeout=aiohttp.ClientTimeout(total=10),
) as resp:
health = await resp.json()
async with self._session.get(
f'{endpoint}/peers',
timeout=aiohttp.ClientTimeout(total=10),
) as resp:
peers_data = await resp.json()
async with self._session.get(
f'{endpoint}/pins',
timeout=aiohttp.ClientTimeout(total=10),
) as resp:
pins = await resp.json()
return ClusterStatus(
peer_id=health.get('peer_id', ''),
peers=[p.get('id', '') for p in peers_data],
pin_count=len(pins) if isinstance(pins, list) else 0,
cluster_healthy=health.get('status') == 'healthy',
)
async def get_replication_status(self, cid: str) -> Dict:
"""获取特定 CID 的复制状态"""
await self._ensure_session()
endpoint = self.endpoints[0]
async with self._session.get(
f'{endpoint}/pins/{cid}',
timeout=aiohttp.ClientTimeout(total=10),
) as resp:
return await resp.json()
async def close(self):
if self._session:
await self._session.close()
self._session = None
class PinataService:
"""
Pinata 云 pinning 服务集成
提供可靠的远程 pinning 基础设施
"""
def __init__(self, api_key: str, api_secret: str):
self.api_key = api_key
self.api_secret = api_secret
self.base_url = 'https://api.pinata.cloud'
self._session: Optional[aiohttp.ClientSession] = None
async def _ensure_session(self):
if self._session is None:
self._session = aiohttp.ClientSession(
headers={
'pinata_api_key': self.api_key,
'pinata_secret_api_key': self.api_secret,
}
)
async def pin_by_cid(self, cid: str, name: str = '',
metadata: Optional[Dict] = None) -> Dict:
"""通过 CID pin 内容到 Pinata"""
await self._ensure_session()
payload = {
'pinataContent': None,
'pinataMetadata': {
'name': name or cid,
'keyvalues': {
'source': 'msg_chain_agent',
'chain': 'msg-chain-1',
**(metadata or {}),
},
},
'pinataOptions': {
'cidVersion': 1,
},
}
async with self._session.post(
f'{self.base_url}/pinning/pinByHash',
json={'hashToPin': cid, **payload['pinataMetadata'],
'pinataOptions': payload['pinataOptions']},
timeout=aiohttp.ClientTimeout(total=60),
) as resp:
return await resp.json()
async def pin_json(self, data: dict, name: str = '') -> Dict:
"""上传并 pin JSON 数据到 Pinata"""
await self._ensure_session()
payload = {
'pinataContent': data,
'pinataMetadata': {
'name': name or 'msg_agent_data',
'keyvalues': {
'source': 'msg_chain_agent',
'chain': 'msg-chain-1',
},
},
}
async with self._session.post(
f'{self.base_url}/pinning/pinJSONToIPFS',
json=payload,
timeout=aiohttp.ClientTimeout(total=60),
) as resp:
return await resp.json()
async def unpin(self, cid: str) -> Dict:
"""从 Pinata 取消 pin"""
await self._ensure_session()
async with self._session.delete(
f'{self.base_url}/pinning/unpin/{cid}',
timeout=aiohttp.ClientTimeout(total=30),
) as resp:
return await resp.json()
async def list_pins(self, status: str = 'pinned') -> List[Dict]:
"""列出 Pinata 上的 pin"""
await self._ensure_session()
params = {'status': status, 'pageLimit': 100}
async with self._session.get(
f'{self.base_url}/data/pinList',
params=params,
timeout=aiohttp.ClientTimeout(total=30),
) as resp:
data = await resp.json()
return data.get('rows', [])
async def close(self):
if self._session:
await self._session.close()
self._session = None
class InfuraIPFS:
"""
Infura IPFS 服务集成
提供免费等级的 IPFS 网关和 pinning
"""
def __init__(self, project_id: str, project_secret: str):
self.project_id = project_id
self.project_secret = project_secret
self.base_url = f'https://ipfs.infura.io:5001'
self._auth = aiohttp.BasicAuth(project_id, project_secret)
async def add(self, data: bytes) -> str:
"""上传数据到 Infura IPFS"""
async with aiohttp.ClientSession() as session:
async with session.post(
f'{self.base_url}/api/v0/add',
data={'file': data},
auth=self._auth,
timeout=aiohttp.ClientTimeout(total=120),
) as resp:
result = await resp.json()
return result.get('Hash', '')
async def pin(self, cid: str) -> bool:
"""通过 Infura pin 内容"""
async with aiohttp.ClientSession() as session:
async with session.post(
f'{self.base_url}/api/v0/pin/add?arg={cid}',
auth=self._auth,
timeout=aiohttp.ClientTimeout(total=60),
) as resp:
return resp.status == 200
async def cat(self, cid: str) -> bytes:
"""从 Infura 网关获取内容"""
async with aiohttp.ClientSession() as session:
async with session.get(
f'https://{self.project_id}.infura-ipfs.io/ipfs/{cid}',
timeout=aiohttp.ClientTimeout(total=30),
) as resp:
return await resp.read()
2.4 CID 生成与验证
import multihash
from multiformats import CID
import hashlib
from typing import Union
class CIDHelper:
"""
CID 生成与验证工具
支持 CIDv0 和 CIDv1
"""
@staticmethod
def compute_cid(data: Union[str, bytes], cid_version: int = 1) -> str:
"""计算数据的 CID"""
if isinstance(data, str):
data = data.encode('utf-8')
hash_bytes = hashlib.sha256(data).digest()
if cid_version == 0:
mh = multihash.encode(hash_bytes, 'sha2-256')
cid_str = CID('dag-pb', mh).encode('base32')
return cid_str
else:
mh = multihash.encode(hash_bytes, 'sha2-256')
cid_obj = CID('dag-pb', mh)
return cid_obj.encode('base32')
@staticmethod
def validate_cid(cid_str: str) -> bool:
"""验证 CID 格式是否有效"""
try:
cid = CID.decode(cid_str)
return cid is not None
except Exception:
return False
@staticmethod
def extract_hash(cid_str: str) -> str:
"""从 CID 中提取底层哈希值"""
cid = CID.decode(cid_str)
return cid.multihash.digest.hex()
@staticmethod
def cid_to_bytes(cid_str: str) -> bytes:
"""将 CID 转换为字节表示"""
cid = CID.decode(cid_str)
return cid.bytes
@staticmethod
def bytes_to_cid(data: bytes) -> str:
"""从字节重建 CID"""
cid = CID.decode(data)
return str(cid)
class AgentDataEncoder:
"""
Agent 数据编码器
确保数据包含 MSG Chain 命名空间
"""
MSG_CHAIN_PREFIX = 'msg'
CHAIN_ID = 'msg-chain-1'
@classmethod
def encode(cls, data: dict, agent_id: str) -> dict:
"""编码 Agent 数据,添加 MSG Chain 元数据"""
import datetime
return {
'msg_chain': {
'version': '1.0',
'chain_id': cls.CHAIN_ID,
'bech32_prefix': cls.MSG_CHAIN_PREFIX,
'agent_id': agent_id,
'encoded_at': datetime.datetime.utcnow().isoformat(),
},
'payload': data,
}
@classmethod
def decode(cls, encoded: dict) -> dict:
"""解码 Agent 数据"""
if 'msg_chain' in encoded and 'payload' in encoded:
return encoded['payload']
return encoded
2.5 MFS(Mutable File System)支持
class MFSClient:
"""
IPFS MFS(Mutable File System)客户端
支持可变内容的 IPFS 操作
"""
def __init__(self, ipfs_client):
self._client = ipfs_client._client
async def write(self, path: str, data: bytes):
"""写入文件到 MFS"""
loop = asyncio.get_event_loop()
await loop.run_in_executor(
None,
lambda: self._client.files.write(path, data, create=True, truncate=True)
)
async def read(self, path: str) -> bytes:
"""从 MFS 读取文件"""
loop = asyncio.get_event_loop()
result = await loop.run_in_executor(
None,
lambda: self._client.files.read(path)
)
return result
async def mkdir(self, path: str, parents: bool = True):
"""创建 MFS 目录"""
loop = asyncio.get_event_loop()
await loop.run_in_executor(
None,
lambda: self._client.files.mkdir(path, parents=parents)
)
async def ls(self, path: str) -> List[Dict]:
"""列出 MFS 目录内容"""
loop = asyncio.get_event_loop()
result = await loop.run_in_executor(
None,
lambda: self._client.files.ls(path)
)
return result.get('Entries', [])
async def rm(self, path: str, recursive: bool = False):
"""删除 MFS 文件"""
loop = asyncio.get_event_loop()
await loop.run_in_executor(
None,
lambda: self._client.files.rm(path, recursive=recursive)
)
2.6 IPFS 网关与公共访问
class IPFSGateway:
"""
IPFS 公共网关管理
提供通过 HTTP 网关访问 IPFS 内容的统一接口
"""
PUBLIC_GATEWAYS = [
'https://ipfs.io/ipfs/{cid}',
'https://cloudflare-ipfs.com/ipfs/{cid}',
'https://gateway.pinata.cloud/ipfs/{cid}',
'https://dweb.link/ipfs/{cid}',
]
@staticmethod
def gateway_url(cid: str, gateway_index: int = 0) -> str:
"""获取 IPFS 网关 URL"""
template = IPFSGateway.PUBLIC_GATEWAYS[gateway_index]
return template.format(cid=cid)
@staticmethod
async def fetch_via_gateway(cid: str, timeout: int = 30) -> Optional[bytes]:
"""通过公共网关获取内容(自动回退)"""
for i, template in enumerate(IPFSGateway.PUBLIC_GATEWAYS):
url = template.format(cid=cid)
try:
async with aiohttp.ClientSession() as session:
async with session.get(
url, timeout=aiohttp.ClientTimeout(total=timeout)
) as resp:
if resp.status == 200:
return await resp.read()
except Exception:
continue
return None
@staticmethod
def generate_subdomain_url(cid: str) -> str:
"""生成子域名格式的网关 URL"""
return f'https://{cid}.ipfs.dweb.link/'
2.7 IPFS 内容完整性验证
class ContentVerifier:
"""IPFS 内容完整性验证器"""
@staticmethod
def verify_content(data: bytes, expected_cid: str) -> bool:
"""验证内容是否匹配预期的 CID"""
hash_bytes = hashlib.sha256(data).digest()
mh = multihash.encode(hash_bytes, 'sha2-256')
actual_cid = str(CID('dag-pb', mh))
return actual_cid == expected_cid
@staticmethod
def verify_agent_content(data: dict, chain_record: dict) -> bool:
"""验证 Agent 内容与链上记录的一致性"""
serialized = json.dumps(data, ensure_ascii=False, default=str)
content_bytes = serialized.encode('utf-8')
actual_checksum = hashlib.sha256(content_bytes).hexdigest()
expected_checksum = chain_record.get('checksum', '')
if expected_checksum and actual_checksum != expected_checksum:
return False
actual_cid = CIDHelper.compute_cid(content_bytes)
expected_cid = chain_record.get('cid', '')
if expected_cid and actual_cid != expected_cid:
return False
return True
3. IPFS 锚定合约
3.1 合约设计
IPFS 锚定合约用于将 IPFS 内容的 CID 锚定到 MSG Chain 上,建立链下数据的永久链上引用。每个锚定记录包含:
- 所有者地址(
msg1...前缀) - 内容 CID
- 数据类型(conversation, knowledge, model 等)
- 时间戳
- 内容校验和
- 数据大小
- 可选元数据
3.2 Rust 智能合约实现
use cosmwasm_std::{
entry_point, to_binary, Binary, Deps, DepsMut, Env, MessageInfo,
Response, StdError, StdResult, Addr, Storage, Uint64,
};
use cw_storage_plus::{Item, Map, IndexedMap, MultiIndex, Index};
use schemars::JsonSchema;
use serde::{Deserialize, Serialize};
// MSG Chain bech32 前缀
const MSG_BECH32_PREFIX: &str = "msg";
const CHAIN_ID: &str = "msg-chain-1";
// ---------- 状态数据结构 ----------
/// 数据锚定记录
#[derive(Serialize, Deserialize, Clone, Debug, PartialEq, JsonSchema)]
pub struct DataAnchor {
/// 所有者地址(msg1...)
pub owner: Addr,
/// IPFS 内容标识符
pub cid: String,
/// 数据类型标识
pub data_type: String,
/// 锚定时间戳(Unix 秒)
pub timestamp: u64,
/// 内容 SHA-256 校验和
pub checksum: String,
/// 数据大小(字节)
pub size: u64,
/// 数据格式(json, binary, text)
pub format: String,
/// 是否可变(MFS 支持)
pub mutable: bool,
/// 当前版本(可变数据的版本号)
pub version: u32,
}
/// 链上查询请求
#[derive(Serialize, Deserialize, Clone, Debug, PartialEq, JsonSchema)]
#[serde(rename_all = "snake_case")]
pub enum QueryMsg {
/// 按 CID 查询锚定
GetAnchor { cid: String },
/// 查询所有者的锚定列表
GetAnchorsByOwner {
owner: String,
start_after: Option<String>,
limit: Option<u32>,
},
/// 按类型查询锚定
GetAnchorsByType {
data_type: String,
start_after: Option<String>,
limit: Option<u32>,
},
/// 查询锚定总数
GetAnchorCount {},
}
/// 合约响应消息
#[derive(Serialize, Deserialize, Clone, Debug, PartialEq, JsonSchema)]
#[serde(rename_all = "snake_case")]
pub enum ExecuteMsg {
/// 锚定新的数据引用
AnchorData {
cid: String,
data_type: String,
checksum: String,
size: u64,
format: String,
mutable: bool,
},
/// 更新可变数据的锚定
UpdateAnchor {
cid: String,
new_cid: String,
new_checksum: String,
new_size: u64,
},
/// 转移锚定所有权
TransferOwnership {
cid: String,
new_owner: String,
},
/// 删除锚定
RemoveAnchor {
cid: String,
},
}
/// 实例化消息
#[derive(Serialize, Deserialize, Clone, Debug, PartialEq, JsonSchema)]
pub struct InstantiateMsg {
pub admin: Option<String>,
pub min_anchor_size: u64,
pub max_anchor_size: u64,
}
/// 合约配置
#[derive(Serialize, Deserialize, Clone, Debug, PartialEq, JsonSchema)]
pub struct Config {
pub admin: Addr,
pub min_anchor_size: u64,
pub max_anchor_size: u64,
pub total_anchors: u64,
}
/// 锚定历史记录
#[derive(Serialize, Deserialize, Clone, Debug, PartialEq, JsonSchema)]
pub struct AnchorHistory {
pub old_cid: String,
pub new_cid: String,
pub timestamp: u64,
pub updated_by: Addr,
}
// ---------- 状态存储 ----------
pub const CONFIG: Item<Config> = Item::new("config");
pub const ANCHORS: Map<&str, DataAnchor> = Map::new("anchors");
pub const OWNER_ANCHORS: Map<(&str, &str), bool> = Map::new("owner_anchors");
pub const TYPE_ANCHORS: Map<(&str, &str), bool> = Map::new("type_anchors");
pub const ANCHOR_HISTORY: Map<&str, Vec<AnchorHistory>> = Map::new("history");
pub const ANCHOR_COUNT: Item<u64> = Item::new("anchor_count");
// ---------- 合约入口 ----------
#[entry_point]
pub fn instantiate(
deps: DepsMut,
_env: Env,
info: MessageInfo,
msg: InstantiateMsg,
) -> StdResult<Response> {
let admin = match msg.admin {
Some(admin_addr) => deps.api.addr_validate(&admin_addr)?,
None => info.sender.clone(),
};
let config = Config {
admin,
min_anchor_size: msg.min_anchor_size,
max_anchor_size: msg.max_anchor_size,
total_anchors: 0,
};
CONFIG.save(deps.storage, &config)?;
ANCHOR_COUNT.save(deps.storage, &0)?;
Ok(Response::new()
.add_attribute("method", "instantiate")
.add_attribute("admin", config.admin.to_string())
.add_attribute("chain_id", CHAIN_ID))
}
#[entry_point]
pub fn execute(
deps: DepsMut,
env: Env,
info: MessageInfo,
msg: ExecuteMsg,
) -> StdResult<Response> {
match msg {
ExecuteMsg::AnchorData {
cid, data_type, checksum, size, format, mutable,
} => execute_anchor_data(deps, env, info, cid, data_type, checksum, size, format, mutable),
ExecuteMsg::UpdateAnchor {
cid, new_cid, new_checksum, new_size,
} => execute_update_anchor(deps, env, info, cid, new_cid, new_checksum, new_size),
ExecuteMsg::TransferOwnership { cid, new_owner } => {
execute_transfer_ownership(deps, env, info, cid, new_owner)
}
ExecuteMsg::RemoveAnchor { cid } => execute_remove_anchor(deps, env, info, cid),
}
}
/// 锚定数据
#[allow(clippy::too_many_arguments)]
pub fn execute_anchor_data(
deps: DepsMut,
env: Env,
info: MessageInfo,
cid: String,
data_type: String,
checksum: String,
size: u64,
format: String,
mutable: bool,
) -> StdResult<Response> {
let config = CONFIG.load(deps.storage)?;
if cid.is_empty() {
return Err(StdError::generic_err("CID 不能为空"));
}
if size < config.min_anchor_size {
return Err(StdError::generic_err(format!(
"数据大小 {} 字节小于最小限制 {} 字节", size, config.min_anchor_size
)));
}
if size > config.max_anchor_size {
return Err(StdError::generic_err(format!(
"数据大小 {} 字节大于最大限制 {} 字节", size, config.max_anchor_size
)));
}
if ANCHORS.has(deps.storage, cid.as_str()) {
return Err(StdError::generic_err(format!("CID {} 已被锚定", cid)));
}
let anchor = DataAnchor {
owner: info.sender.clone(),
cid: cid.clone(),
data_type: data_type.clone(),
timestamp: env.block.time.seconds(),
checksum: checksum.clone(),
size,
format,
mutable,
version: 1,
};
ANCHORS.save(deps.storage, cid.as_str(), &anchor)?;
OWNER_ANCHORS.save(deps.storage, (info.sender.as_str(), cid.as_str()), &true)?;
TYPE_ANCHORS.save(deps.storage, (data_type.as_str(), cid.as_str()), &true)?;
let count = ANCHOR_COUNT.load(deps.storage)?;
ANCHOR_COUNT.save(deps.storage, &(count + 1))?;
let mut config = config;
config.total_anchors = count + 1;
CONFIG.save(deps.storage, &config)?;
Ok(Response::new()
.add_attribute("method", "anchor_data")
.add_attribute("owner", info.sender.to_string())
.add_attribute("cid", &cid)
.add_attribute("data_type", &data_type)
.add_attribute("checksum", &checksum)
.add_attribute("size", size.to_string())
.add_attribute("timestamp", env.block.time.seconds().to_string()))
}
/// 更新锚定(仅可变数据)
pub fn execute_update_anchor(
deps: DepsMut,
env: Env,
info: MessageInfo,
cid: String,
new_cid: String,
new_checksum: String,
new_size: u64,
) -> StdResult<Response> {
let mut anchor = ANCHORS
.load(deps.storage, cid.as_str())
.map_err(|_| StdError::generic_err(format!("CID {} 未找到", cid)))?;
if anchor.owner != info.sender {
return Err(StdError::generic_err("只有所有者可以更新锚定"));
}
if !anchor.mutable {
return Err(StdError::generic_err("不可变数据不能更新"));
}
let history_entry = AnchorHistory {
old_cid: cid.clone(),
new_cid: new_cid.clone(),
timestamp: env.block.time.seconds(),
updated_by: info.sender.clone(),
};
let mut history = ANCHOR_HISTORY
.load(deps.storage, cid.as_str())
.unwrap_or_default();
history.push(history_entry);
ANCHOR_HISTORY.save(deps.storage, cid.as_str(), &history)?;
OWNER_ANCHORS.remove(deps.storage, (anchor.owner.as_str(), cid.as_str()));
TYPE_ANCHORS.remove(deps.storage, (anchor.data_type.as_str(), cid.as_str()));
ANCHORS.remove(deps.storage, cid.as_str());
let new_anchor = DataAnchor {
owner: info.sender.clone(),
cid: new_cid.clone(),
data_type: anchor.data_type.clone(),
timestamp: env.block.time.seconds(),
checksum: new_checksum.clone(),
size: new_size,
format: anchor.format.clone(),
mutable: true,
version: anchor.version + 1,
};
ANCHORS.save(deps.storage, new_cid.as_str(), &new_anchor)?;
OWNER_ANCHORS.save(deps.storage, (info.sender.as_str(), new_cid.as_str()), &true)?;
TYPE_ANCHORS.save(deps.storage, (anchor.data_type.as_str(), new_cid.as_str()), &true)?;
Ok(Response::new()
.add_attribute("method", "update_anchor")
.add_attribute("from_cid", &cid)
.add_attribute("to_cid", &new_cid)
.add_attribute("new_version", (anchor.version + 1).to_string()))
}
/// 转移锚定所有权
pub fn execute_transfer_ownership(
deps: DepsMut,
env: Env,
info: MessageInfo,
cid: String,
new_owner: String,
) -> StdResult<Response> {
let new_owner_addr = deps.api.addr_validate(&new_owner)?;
let mut anchor = ANCHORS
.load(deps.storage, cid.as_str())
.map_err(|_| StdError::generic_err(format!("CID {} 未找到", cid)))?;
if anchor.owner != info.sender {
return Err(StdError::generic_err("只有所有者可以转移所有权"));
}
OWNER_ANCHORS.remove(deps.storage, (anchor.owner.as_str(), cid.as_str()));
OWNER_ANCHORS.save(deps.storage, (new_owner_addr.as_str(), cid.as_str()), &true)?;
anchor.owner = new_owner_addr.clone();
ANCHORS.save(deps.storage, cid.as_str(), &anchor)?;
Ok(Response::new()
.add_attribute("method", "transfer_ownership")
.add_attribute("cid", &cid)
.add_attribute("from", info.sender.to_string())
.add_attribute("to", new_owner_addr.to_string()))
}
/// 删除锚定
pub fn execute_remove_anchor(
deps: DepsMut,
_env: Env,
info: MessageInfo,
cid: String,
) -> StdResult<Response> {
let anchor = ANCHORS
.load(deps.storage, cid.as_str())
.map_err(|_| StdError::generic_err(format!("CID {} 未找到", cid)))?;
if anchor.owner != info.sender {
return Err(StdError::generic_err("只有所有者可以删除锚定"));
}
OWNER_ANCHORS.remove(deps.storage, (anchor.owner.as_str(), cid.as_str()));
TYPE_ANCHORS.remove(deps.storage, (anchor.data_type.as_str(), cid.as_str()));
ANCHORS.remove(deps.storage, cid.as_str());
let count = ANCHOR_COUNT.load(deps.storage)?;
ANCHOR_COUNT.save(deps.storage, &(count.saturating_sub(1)))?;
Ok(Response::new()
.add_attribute("method", "remove_anchor")
.add_attribute("cid", &cid))
}
// ---------- 查询 ----------
#[entry_point]
pub fn query(deps: Deps, _env: Env, msg: QueryMsg) -> StdResult<Binary> {
match msg {
QueryMsg::GetAnchor { cid } => to_binary(&query_anchor(deps, cid)?),
QueryMsg::GetAnchorsByOwner { owner, start_after, limit } => {
to_binary(&query_anchors_by_owner(deps, owner, start_after, limit)?)
}
QueryMsg::GetAnchorsByType { data_type, start_after, limit } => {
to_binary(&query_anchors_by_type(deps, data_type, start_after, limit)?)
}
QueryMsg::GetAnchorCount {} => to_binary(&query_anchor_count(deps)?),
}
}
fn query_anchor(deps: Deps, cid: String) -> StdResult<DataAnchor> {
ANCHORS
.load(deps.storage, cid.as_str())
.map_err(|_| StdError::not_found(format!("锚定 {} 不存在", cid)))
}
fn query_anchors_by_owner(
deps: Deps,
owner: String,
start_after: Option<String>,
limit: Option<u32>,
) -> StdResult<Vec<DataAnchor>> {
let owner_addr = deps.api.addr_validate(&owner)?;
let limit = limit.unwrap_or(10).min(30) as usize;
let anchors: StdResult<Vec<DataAnchor>> = OWNER_ANCHORS
.prefix(owner_addr.as_str())
.range(deps.storage, None, None, cosmwasm_std::Order::Ascending)
.take(limit)
.map(|item| {
let (cid, _) = item?;
ANCHORS.load(deps.storage, &cid)
})
.collect();
anchors
}
fn query_anchors_by_type(
deps: Deps,
data_type: String,
start_after: Option<String>,
limit: Option<u32>,
) -> StdResult<Vec<DataAnchor>> {
let limit = limit.unwrap_or(10).min(30) as usize;
let anchors: StdResult<Vec<DataAnchor>> = TYPE_ANCHORS
.prefix(data_type.as_str())
.range(deps.storage, None, None, cosmwasm_std::Order::Ascending)
.take(limit)
.map(|item| {
let (cid, _) = item?;
ANCHORS.load(deps.storage, &cid)
})
.collect();
anchors
}
fn query_anchor_count(deps: Deps) -> StdResult<u64> {
ANCHOR_COUNT.load(deps.storage)
}
3.3 合约部署脚本
import subprocess
import json
from pathlib import Path
class DataAnchorDeployer:
"""数据锚定合约部署工具"""
def __init__(self, chain_id: str = 'msg-chain-1',
node_url: str = 'http://localhost:26657',
key_name: str = 'default'):
self.chain_id = chain_id
self.node_url = node_url
self.key_name = key_name
self.binary = 'msgchaind'
def compile_contract(self, contract_path: str) -> Path:
"""编译 CosmWasm 合约"""
result = subprocess.run(
['cargo', 'wasm'],
cwd=contract_path,
capture_output=True,
text=True,
check=True,
)
wasm_path = Path(contract_path) / 'target/wasm32-unknown-unknown/release'
wasm_files = list(wasm_path.glob('*.wasm'))
if not wasm_files:
raise RuntimeError('编译失败:未找到 WASM 文件')
return wasm_files[0]
def store_code(self, wasm_path: Path, gas_limit: str = '5000000') -> int:
"""上传合约代码到链"""
cmd = [
self.binary, 'tx', 'wasm', 'store', str(wasm_path),
'--from', self.key_name,
'--chain-id', self.chain_id,
'--node', self.node_url,
'--gas', gas_limit,
'--gas-prices', '1000000000attoMSG',
'--gas-adjustment', '1.3',
'--yes',
'--output', 'json',
]
result = subprocess.run(cmd, capture_output=True, text=True, check=True)
tx_result = json.loads(result.stdout)
code_id = self._parse_code_id(tx_result)
self._wait_for_tx(tx_result['txhash'])
return code_id
def instantiate_contract(
self, code_id: int, admin: str, min_size: int = 1,
max_size: int = 10 * 1024 * 1024, label: str = 'msg_data_anchor',
gas_limit: str = '3000000',
) -> str:
"""实例化合约"""
init_msg = {
'admin': admin,
'min_anchor_size': min_size,
'max_anchor_size': max_size,
}
cmd = [
self.binary, 'tx', 'wasm', 'instantiate', str(code_id),
json.dumps(init_msg),
'--from', self.key_name,
'--admin', admin,
'--label', label,
'--chain-id', self.chain_id,
'--node', self.node_url,
'--gas', gas_limit,
'--gas-prices', '1000000000attoMSG',
'--gas-adjustment', '1.3',
'--yes',
'--output', 'json',
]
result = subprocess.run(cmd, capture_output=True, text=True, check=True)
tx_result = json.loads(result.stdout)
contract_addr = self._parse_contract_address(tx_result)
self._wait_for_tx(tx_result['txhash'])
return contract_addr
def _parse_code_id(self, tx_result: dict) -> int:
for event in tx_result.get('logs', []):
for evt in event.get('events', []):
if evt.get('type') == 'store_code':
for attr in evt.get('attributes', []):
if attr.get('key') == 'code_id':
return int(attr['value'])
raise ValueError('未找到 Code ID')
def _parse_contract_address(self, tx_result: dict) -> str:
for event in tx_result.get('logs', []):
for evt in event.get('events', []):
if evt.get('type') == 'instantiate':
for attr in evt.get('attributes', []):
if attr.get('key') == '_contract_address':
addr = attr['value']
assert addr.startswith('msg1'), \
f'地址必须以 msg1 开头: {addr}'
return addr
raise ValueError('未找到合约地址')
def _wait_for_tx(self, txhash: str, timeout: int = 60):
import time
start = time.time()
while time.time() - start < timeout:
cmd = [
self.binary, 'query', 'tx', txhash,
'--node', self.node_url,
'--output', 'json',
]
result = subprocess.run(cmd, capture_output=True, text=True)
try:
tx_data = json.loads(result.stdout)
if tx_data.get('code') == 0:
return
except (json.JSONDecodeError, KeyError):
pass
time.sleep(2)
raise TimeoutError('交易确认超时')
3.4 链上锚定交互
from typing import Optional
from dataclasses import dataclass
import hashlib
import json
@dataclass
class AnchorRecord:
"""链上锚定记录"""
owner: str
cid: str
data_type: str
timestamp: int
checksum: str
size: int
format: str
mutable: bool
version: int
class ChainAnchorClient:
"""链上锚定客户端"""
def __init__(self, contract_addr: str,
rpc_url: str = 'http://localhost:26657',
rest_url: str = 'http://localhost:1317'):
self.contract_addr = contract_addr
self.rpc_url = rpc_url
self.rest_url = rest_url
if not contract_addr.startswith('msg1'):
raise ValueError(f'MSG Chain 合约地址必须以 msg1 开头: {contract_addr}')
async def anchor_data(self, cid: str, data_type: str,
data: dict, key_name: str) -> str:
"""锚定数据到链上"""
serialized = json.dumps(data, ensure_ascii=False, default=str)
checksum = hashlib.sha256(serialized.encode()).hexdigest()
size = len(serialized.encode())
msg = {
'anchor_data': {
'cid': cid,
'data_type': data_type,
'checksum': checksum,
'size': size,
'format': 'json',
'mutable': False,
}
}
return await self._execute_contract(msg, key_name)
async def query_anchor(self, cid: str) -> Optional[AnchorRecord]:
"""查询链上锚定记录"""
query = {'get_anchor': {'cid': cid}}
result = await self._query_contract(query)
if not result:
return None
return AnchorRecord(
owner=result['owner'],
cid=result['cid'],
data_type=result['data_type'],
timestamp=result['timestamp'],
checksum=result['checksum'],
size=result['size'],
format=result.get('format', 'json'),
mutable=result.get('mutable', False),
version=result.get('version', 1),
)
async def query_by_owner(self, owner: str, limit: int = 10) -> list:
"""查询所有者的锚定列表"""
if not owner.startswith('msg1'):
raise ValueError(f'MSG Chain 地址必须以 msg1 开头: {owner}')
query = {'get_anchors_by_owner': {'owner': owner, 'limit': limit}}
return await self._query_contract(query)
async def query_by_type(self, data_type: str, limit: int = 10) -> list:
"""按类型查询锚定"""
query = {'get_anchors_by_type': {'data_type': data_type, 'limit': limit}}
return await self._query_contract(query)
async def verify_on_chain(self, cid: str, data: dict) -> bool:
"""验证链下数据与链上记录是否匹配"""
record = await self.query_anchor(cid)
if not record:
return False
serialized = json.dumps(data, ensure_ascii=False, default=str)
actual_checksum = hashlib.sha256(serialized.encode()).hexdigest()
return actual_checksum == record.checksum
async def _execute_contract(self, msg: dict, key_name: str) -> str:
import subprocess
cmd = [
'msgchaind', 'tx', 'wasm', 'execute',
self.contract_addr, json.dumps(msg),
'--from', key_name,
'--chain-id', 'msg-chain-1',
'--node', self.rpc_url,
'--gas', '3000000',
'--gas-prices', '1000000000attoMSG',
'--gas-adjustment', '1.3',
'--yes', '--output', 'json',
]
result = subprocess.run(cmd, capture_output=True, text=True, check=True)
tx_result = json.loads(result.stdout)
return tx_result['txhash']
async def _query_contract(self, query: dict) -> any:
import subprocess
cmd = [
'msgchaind', 'query', 'wasm', 'contract-state', 'smart',
self.contract_addr, json.dumps(query),
'--node', self.rpc_url, '--output', 'json',
]
result = subprocess.run(cmd, capture_output=True, text=True, check=True)
return json.loads(result.stdout).get('data')
4. Arweave 永久存储
4.1 Arweave 简介
Arweave 是一种永久去中心化存储网络,通过一次性付费实现数据永久保存。对于 AI Agent,Arweave 特别适合:
- 对话历史归档:Agent 的完整对话记录永久保存
- 知识库快照:Agent 知识库的不可变版本记录
- 模型权重存档:AI 模型权重和配置的永久存储
- 审计日志:Agent 操作的不可篡改审计痕迹
4.2 Python Arweave 客户端
import json
import hashlib
from typing import Optional, List, Dict, Any
from dataclasses import dataclass
from datetime import datetime
import base64
import aiohttp
@dataclass
class ArweaveTag:
"""Arweave 标签键值对"""
name: str
value: str
class ArweaveWallet:
"""Arweave 钱包管理"""
def __init__(self, jwk: dict):
self.jwk = jwk
self.address = self._derive_address()
@classmethod
def from_file(cls, path: str) -> 'ArweaveWallet':
with open(path) as f:
jwk = json.load(f)
return cls(jwk)
def _derive_address(self) -> str:
n_bytes = base64.urlsafe_b64decode(self.jwk['n'])
hash_bytes = hashlib.sha256(n_bytes).digest()
return base64.urlsafe_b64encode(hash_bytes).decode()[:43]
class ArweaveClient:
"""Arweave 客户端"""
def __init__(self, wallet: ArweaveWallet,
gateway_url: str = 'https://arweave.net',
graphql_url: str = 'https://arweave.net/graphql'):
self.wallet = wallet
self.gateway_url = gateway_url
self.graphql_url = graphql_url
self._session: Optional[aiohttp.ClientSession] = None
async def _ensure_session(self):
if self._session is None:
self._session = aiohttp.ClientSession()
async def store_permanently(self, data: bytes,
tags: Optional[List[ArweaveTag]] = None,
target: str = '', quantity: str = '0') -> str:
"""永久存储数据到 Arweave"""
await self._ensure_session()
all_tags = [
ArweaveTag('Content-Type', 'application/octet-stream'),
ArweaveTag('App-Name', 'msg-chain-agent'),
ArweaveTag('App-Version', '1.0.0'),
ArweaveTag('Chain-ID', 'msg-chain-1'),
ArweaveTag('Bech32-Prefix', 'msg'),
]
if tags:
all_tags.extend(tags)
info = await self._get_network_info()
last_tx = info.get('last_tx', '')
reward = await self._get_transaction_price(len(data))
tx = {
'format': 2,
'data': base64.urlsafe_b64encode(data).decode(),
'tags': [
{
'name': base64.urlsafe_b64encode(t.name.encode()).decode(),
'value': base64.urlsafe_b64encode(t.value.encode()).decode(),
}
for t in all_tags
],
'owner': self.wallet.jwk['n'],
'target': target,
'quantity': quantity,
'reward': reward,
'last_tx': last_tx,
}
data_root = self._compute_data_root(data)
tx['data_root'] = data_root
tx['data_size'] = str(len(data))
tx_id = await self._post_transaction(tx)
return tx_id
async def store_json_permanently(self, data: dict,
tags: Optional[List[ArweaveTag]] = None,
description: str = '') -> str:
"""永久存储 JSON 数据"""
serialized = json.dumps(data, ensure_ascii=False, default=str)
content_bytes = serialized.encode('utf-8')
all_tags = [
ArweaveTag('Content-Type', 'application/json'),
ArweaveTag('Data-Description', description or 'msg-agent-data'),
ArweaveTag('Data-Checksum', hashlib.sha256(content_bytes).hexdigest()),
]
if tags:
all_tags.extend(tags)
return await self.store_permanently(content_bytes, all_tags)
async def store_agent_state(self, agent_id: str, state: dict,
state_type: str = 'state', version: str = '1') -> str:
"""存储 Agent 状态到 Arweave"""
tags = [
ArweaveTag('Agent-ID', agent_id),
ArweaveTag('State-Type', state_type),
ArweaveTag('State-Version', version),
ArweaveTag('Chain-ID', 'msg-chain-1'),
]
return await self.store_json_permanently(state, tags, f'msg-agent-{state_type}')
async def get_transaction(self, tx_id: str) -> Optional[Dict]:
"""获取交易数据"""
await self._ensure_session()
url = f'{self.gateway_url}/tx/{tx_id}'
async with self._session.get(url) as resp:
if resp.status == 200:
return await resp.json()
return None
async def get_transaction_data(self, tx_id: str) -> Optional[bytes]:
"""获取交易数据内容"""
await self._ensure_session()
url = f'{self.gateway_url}/{tx_id}'
async with self._session.get(url) as resp:
if resp.status == 200:
return await resp.read()
return None
async def get_json_data(self, tx_id: str) -> Optional[dict]:
"""获取并解析 JSON 交易数据"""
data = await self.get_transaction_data(tx_id)
if data:
try:
return json.loads(data.decode('utf-8'))
except json.JSONDecodeError:
return None
return None
async def _get_network_info(self) -> Dict:
url = f'{self.gateway_url}/info'
async with self._session.get(url) as resp:
return await resp.json()
async def _get_transaction_price(self, data_size: int) -> str:
url = f'{self.gateway_url}/price/{data_size}'
async with self._session.get(url) as resp:
return await resp.text()
def _compute_data_root(self, data: bytes) -> str:
hash_bytes = hashlib.sha256(data).digest()
return base64.urlsafe_b64encode(hash_bytes).decode()
async def _post_transaction(self, tx: Dict) -> str:
url = f'{self.gateway_url}/tx'
async with self._session.post(url, json=tx) as resp:
if resp.status in (200, 202):
result = await resp.json()
return result.get('id', '')
raise RuntimeError(f'提交交易失败: {resp.status}')
async def wait_for_confirmation(self, tx_id: str, timeout: int = 300,
poll_interval: int = 5) -> bool:
import asyncio
start = datetime.now()
while (datetime.now() - start).seconds < timeout:
tx = await self.get_transaction(tx_id)
if tx and tx.get('status') == 'confirmed':
return True
await asyncio.sleep(poll_interval)
return False
async def close(self):
if self._session:
await self._session.close()
self._session = None
4.3 GraphQL 查询
class ArweaveGraphQL:
"""Arweave GraphQL 查询接口"""
def __init__(self, graphql_url: str = 'https://arweave.net/graphql'):
self.graphql_url = graphql_url
async def query_transactions(self, tags: Dict[str, str],
first: int = 10,
after: Optional[str] = None) -> Dict:
"""按标签查询交易"""
query = """
query($tags: [TagFilter!]!, $first: Int!, $after: String) {
transactions(tags: $tags, first: $first, after: $after) {
pageInfo { hasNextPage }
edges {
cursor
node {
id
tags { name value }
owner { address }
data { size type }
block { height timestamp }
}
}
}
}
"""
tag_filters = [{'name': k, 'values': [v]} for k, v in tags.items()]
return await self._execute(query, {
'tags': tag_filters, 'first': first, 'after': after,
})
async def query_agent_data(self, agent_id: str,
data_type: str = '', limit: int = 20) -> List[Dict]:
"""查询 Agent 存储在 Arweave 上的数据"""
tags = {'Agent-ID': agent_id}
if data_type:
tags['State-Type'] = data_type
result = await self.query_transactions(tags, limit)
edges = result.get('data', {}).get('transactions', {}).get('edges', [])
return [
{
'id': edge['node']['id'],
'tags': {t['name']: t['value'] for t in edge['node']['tags']},
'owner': edge['node']['owner']['address'],
'size': edge['node']['data']['size'],
'timestamp': edge['node']['block']['timestamp'],
}
for edge in edges
]
async def _execute(self, query: str, variables: Dict) -> Dict:
async with aiohttp.ClientSession() as session:
async with session.post(
self.graphql_url,
json={'query': query, 'variables': variables},
) as resp:
return await resp.json()
4.4 Bundled Transactions(捆绑交易)
class ArweaveBundleClient:
"""Arweave 捆绑交易客户端"""
def __init__(self, client: ArweaveClient,
bundle_service_url: str = 'https://up.arweave.net'):
self.client = client
self.bundle_service_url = bundle_service_url
self._pending_items: List[Dict] = []
def add_item(self, data: bytes, tags: Optional[List[ArweaveTag]] = None):
"""添加数据项到待提交批次"""
self._pending_items.append({
'data': data,
'tags': tags or [],
})
async def flush(self, max_items: int = 100) -> List[str]:
"""提交所有待处理项"""
if not self._pending_items:
return []
batch = self._pending_items[:max_items]
self._pending_items = self._pending_items[max_items:]
bundle_data = self._build_bundle(batch)
bundle_tags = [
ArweaveTag('Bundle-Format', 'json'),
ArweaveTag('Bundle-Version', '1.0.0'),
ArweaveTag('App-Name', 'msg-chain-agent-bundler'),
]
tx_id = await self.client.store_permanently(bundle_data, bundle_tags)
return [tx_id]
async def flush_all(self) -> List[str]:
tx_ids = []
while self._pending_items:
ids = await self.flush()
tx_ids.extend(ids)
return tx_ids
def _build_bundle(self, items: List[Dict]) -> bytes:
bundle = {
'msg_chain_bundle': {
'version': '1.0',
'created_at': datetime.utcnow().isoformat(),
'item_count': len(items),
},
'items': [
{
'data': base64.urlsafe_b64encode(item['data']).decode(),
'tags': [{'name': t.name, 'value': t.value} for t in item['tags']],
'checksum': hashlib.sha256(item['data']).hexdigest(),
}
for item in items
],
}
return json.dumps(bundle, ensure_ascii=False).encode('utf-8')
@property
def pending_count(self) -> int:
return len(self._pending_items)
5. 存储选型策略
5.1 分层存储策略
from enum import Enum
from dataclasses import dataclass, field
from typing import Optional, List, Callable, Awaitable
import time
import json
import math
class StorageTier(Enum):
"""存储层级枚举"""
HOT = 'hot' # 热数据 - 内存/本地缓存
WARM = 'warm' # 温数据 - IPFS
COLD = 'cold' # 冷数据 - Arweave
ARCHIVE = 'archive' # 归档 - Filecoin
@dataclass
class StorageTierConfig:
"""存储层级配置"""
tier: StorageTier
description: str
cost_per_gb_month: float
retrieval_cost_per_gb: float
avg_latency_ms: int
durability: float
persistence: str
max_size_gb: float
min_ttl_seconds: int
@dataclass
class DataProfile:
"""数据访问特征描述"""
data_id: str
size_bytes: int
access_frequency: float
access_pattern: str
required_durability: float
required_latency_ms: int
ttl_seconds: int
budget_per_gb: float
immutable: bool
sensitivity: str
@dataclass
class StorageDecision:
"""存储决策结果"""
tier: StorageTier
provider: str
estimated_cost: float
estimated_latency_ms: int
cid: Optional[str] = None
tx_id: Optional[str] = None
reasons: List[str] = field(default_factory=list)
class StorageStrategy:
"""
自动存储层级选择策略
根据数据特征自动选择最合适的存储方案
"""
TIER_CONFIGS = {
StorageTier.HOT: StorageTierConfig(
tier=StorageTier.HOT,
description='内存/本地缓存',
cost_per_gb_month=0.0,
retrieval_cost_per_gb=0.0,
avg_latency_ms=1,
durability=0.99,
persistence='temporary',
max_size_gb=10.0,
min_ttl_seconds=60,
),
StorageTier.WARM: StorageTierConfig(
tier=StorageTier.WARM,
description='IPFS + 多点 Pinning',
cost_per_gb_month=0.05,
retrieval_cost_per_gb=0.01,
avg_latency_ms=500,
durability=0.999,
persistence='conditional',
max_size_gb=100.0,
min_ttl_seconds=3600,
),
StorageTier.COLD: StorageTierConfig(
tier=StorageTier.COLD,
description='Arweave 永久存储',
cost_per_gb_month=0.01,
retrieval_cost_per_gb=0.02,
avg_latency_ms=2000,
durability=0.99999,
persistence='permanent',
max_size_gb=1000.0,
min_ttl_seconds=86400 * 365,
),
StorageTier.ARCHIVE: StorageTierConfig(
tier=StorageTier.ARCHIVE,
description='Filecoin 存储交易',
cost_per_gb_month=0.003,
retrieval_cost_per_gb=0.05,
avg_latency_ms=5000,
durability=0.999999,
persistence='permanent',
max_size_gb=10000.0,
min_ttl_seconds=86400 * 365,
),
}
MSG_CHAIN_PREFIX = 'msg'
CHAIN_ID = 'msg-chain-1'
def __init__(self, budget_monthly: float = 10.0):
self.budget_monthly = budget_monthly
def select_tier(self, profile: DataProfile) -> StorageTier:
"""根据数据特征选择存储层级"""
for tier, config in self.TIER_CONFIGS.items():
if config.avg_latency_ms <= profile.required_latency_ms:
break
else:
tier = StorageTier.ARCHIVE
for t in [StorageTier.ARCHIVE, StorageTier.COLD,
StorageTier.WARM, StorageTier.HOT]:
cfg = self.TIER_CONFIGS[t]
if cfg.durability >= profile.required_durability:
tier = t
break
for t in [StorageTier.HOT, StorageTier.WARM,
StorageTier.COLD, StorageTier.ARCHIVE]:
cfg = self.TIER_CONFIGS[t]
monthly_cost = (profile.size_bytes / (1024**3)) * cfg.cost_per_gb_month
if monthly_cost <= profile.budget_per_gb:
tier = t
break
return tier
def estimate_cost(self, size_bytes: int, tier: StorageTier,
duration_days: int = 30) -> float:
"""估算存储成本"""
size_gb = size_bytes / (1024**3)
config = self.TIER_CONFIGS[tier]
if tier == StorageTier.HOT:
return 0.0
elif tier == StorageTier.ARCHIVE:
return size_gb * config.cost_per_gb_month * (duration_days / 30)
elif tier == StorageTier.COLD:
return size_gb * config.cost_per_gb_month * (duration_days / 30)
else:
return size_gb * config.cost_per_gb_month * (duration_days / 30)
def decide(self, profile: DataProfile) -> StorageDecision:
"""做出完整存储决策"""
tier = self.select_tier(profile)
config = self.TIER_CONFIGS[tier]
reasons = []
if profile.access_frequency > 10:
reasons.append('高访问频率')
if profile.immutable:
reasons.append('不可变数据')
if profile.required_durability > 0.9999:
reasons.append('高持久性要求')
estimated_cost = self.estimate_cost(profile.size_bytes, tier)
provider_map = {
StorageTier.HOT: 'local_cache',
StorageTier.WARM: 'ipfs_cluster',
StorageTier.COLD: 'arweave',
StorageTier.ARCHIVE: 'filecoin',
}
return StorageDecision(
tier=tier,
provider=provider_map[tier],
estimated_cost=estimated_cost,
estimated_latency_ms=config.avg_latency_ms,
reasons=reasons,
)
class DataLifecycleManager:
"""数据生命周期管理器"""
def __init__(self, strategy: StorageStrategy):
self.strategy = strategy
self._access_log: Dict[str, List[float]] = {}
def record_access(self, data_id: str):
"""记录数据访问时间"""
if data_id not in self._access_log:
self._access_log[data_id] = []
self._access_log[data_id].append(time.time())
cutoff = time.time() - 86400
self._access_log[data_id] = [
t for t in self._access_log[data_id] if t > cutoff
]
def get_frequency(self, data_id: str) -> float:
"""获取数据访问频率(次/小时)"""
logs = self._access_log.get(data_id, [])
if not logs:
return 0.0
window = min(time.time() - min(logs), 86400)
if window < 60:
return len(logs)
return len(logs) / (window / 3600)
def should_promote(self, data_id: str, current_tier: StorageTier,
threshold: float = 5.0) -> bool:
"""判断是否需要升级存储层级"""
freq = self.get_frequency(data_id)
if current_tier == StorageTier.ARCHIVE and freq > threshold:
return True
if current_tier == StorageTier.COLD and freq > threshold * 2:
return True
if current_tier == StorageTier.WARM and freq > threshold * 10:
return True
return False
def should_demote(self, data_id: str, current_tier: StorageTier,
threshold: float = 0.1) -> bool:
"""判断是否需要降级存储层级"""
freq = self.get_frequency(data_id)
if current_tier == StorageTier.HOT and freq < threshold:
return True
if current_tier == StorageTier.WARM and freq < threshold * 0.5:
return True
if current_tier == StorageTier.COLD and freq < threshold * 0.1:
return True
return False
class CostOptimizer:
"""存储成本优化器"""
def __init__(self, monthly_budget: float):
self.monthly_budget = monthly_budget
self.allocations: Dict[str, float] = {}
def optimize_allocation(self, profiles: List[DataProfile]) -> Dict[str, StorageTier]:
"""优化多数据项的存储分配"""
sorted_profiles = sorted(
profiles,
key=lambda p: (p.required_durability, p.access_frequency, -p.budget_per_gb),
reverse=True,
)
allocations = {}
remaining_budget = self.monthly_budget
strategy = StorageStrategy(budget_monthly=remaining_budget)
for profile in sorted_profiles:
best_tier = None
best_cost = float('inf')
for tier in StorageTier:
cost = strategy.estimate_cost(profile.size_bytes, tier)
if cost <= remaining_budget and cost < best_cost:
best_tier = tier
best_cost = cost
if best_tier is None:
cheapest = min(StorageTier, key=lambda t: strategy.estimate_cost(profile.size_bytes, t))
best_tier = cheapest
best_cost = strategy.estimate_cost(profile.size_bytes, cheapest)
allocations[profile.data_id] = best_tier
remaining_budget -= best_cost
return allocations
def suggest_compression(self, data: dict, target_tier: StorageTier) -> dict:
"""建议数据压缩策略"""
import zlib
serialized = json.dumps(data, ensure_ascii=False, default=str)
original_size = len(serialized.encode())
compressed = zlib.compress(serialized.encode())
compressed_size = len(compressed)
savings = original_size - compressed_size
return {
'original_size': original_size,
'compressed_size': compressed_size,
'savings_bytes': savings,
'savings_percent': round(savings / original_size * 100, 2),
'recommended': savings > original_size * 0.3,
}
5.2 自动存储迁移
class StorageMigrator:
"""数据存储迁移器"""
def __init__(self, ipfs_client, arweave_client,
filecoin_client, chain_client):
self.ipfs = ipfs_client
self.arweave = arweave_client
self.filecoin = filecoin_client
self.chain = chain_client
async def migrate_hot_to_warm(self, data_id: str, data: dict) -> str:
"""热数据迁移到 IPFS"""
enriched = {
'msg_chain_origin': {
'data_id': data_id,
'migrated_at': datetime.utcnow().isoformat(),
'source_tier': 'hot',
'target_tier': 'warm',
},
'data': data,
}
cid = await self.ipfs.store_agent_data(enriched)
cluster = IPFSClusterClient(['http://localhost:5001'], replication_factor=3)
await cluster.replicate(cid, replicas=3)
await self.chain.anchor_data(cid, 'migrated_data', data)
return cid
async def migrate_warm_to_cold(self, cid: str) -> str:
"""IPFS 数据迁移到 Arweave"""
data = await self.ipfs.get_agent_data(cid)
anchor = await self.chain.query_anchor(cid)
tags = [
ArweaveTag('Source-CID', cid),
ArweaveTag('Migration-Source', 'ipfs'),
ArweaveTag('Migration-Target', 'arweave'),
ArweaveTag('Chain-ID', 'msg-chain-1'),
]
if anchor:
tags.append(ArweaveTag('Anchor-Owner', anchor.owner))
tx_id = await self.arweave.store_json_permanently(
data, tags, description=f'migrated-from-ipfs-{cid}',
)
await self.chain.anchor_data(tx_id, 'arweave_ref', data)
return tx_id
async def migrate_cold_to_archive(self, tx_id: str, duration_days: int = 365) -> str:
"""Arweave 数据迁移到 Filecoin"""
data = await self.arweave.get_transaction_data(tx_id)
cid = CIDHelper.compute_cid(data)
deal_cid = await self.filecoin.make_deal(
cid=cid, miner='t01000',
duration=duration_days * 86400, price=100000,
)
await self.arweave.store_json_permanently(
{
'migration': {
'source': 'arweave', 'source_tx': tx_id,
'target': 'filecoin', 'target_deal': deal_cid,
'cid': cid, 'timestamp': datetime.utcnow().isoformat(),
}
},
[ArweaveTag('Type', 'migration-log')],
)
return deal_cid
async def auto_migrate(self, data_id: str, current_tier: StorageTier,
target_tier: StorageTier) -> Dict:
"""自动执行层级间迁移"""
migration_map = {
(StorageTier.HOT, StorageTier.WARM): self.migrate_hot_to_warm,
(StorageTier.WARM, StorageTier.COLD): self.migrate_warm_to_cold,
(StorageTier.COLD, StorageTier.ARCHIVE): self.migrate_cold_to_archive,
}
key = (current_tier, target_tier)
if key not in migration_map:
raise ValueError(f'不支持的迁移路径: {current_tier.value} -> {target_tier.value}')
migrator = migration_map[key]
reference = await migrator(data_id, {}) if current_tier == StorageTier.HOT \
else await migrator(data_id)
return {
'data_id': data_id,
'from_tier': current_tier.value,
'to_tier': target_tier.value,
'reference': reference,
'timestamp': datetime.utcnow().isoformat(),
}
6. Filecoin 存储交易
6.1 Filecoin 简介
Filecoin 是一个去中心化存储网络,通过存储合约(Deal)和加密证明确保数据的安全存储。对于 AI Agent,Filecoin 的优势在于:
- 可证明存储:通过时空证明(PoSt)验证数据持续存储
- 大规模归档:适合 GB 级以上数据的长期存储
- 检索市场:支持高效的数据检索
- 与 IPFS 原生兼容:Filecoin 使用相同的 CID 格式
6.2 Python Filecoin 客户端
from typing import Optional, List, Dict, Any
from dataclasses import dataclass
import asyncio
import json
import hashlib
from datetime import datetime, timedelta
import aiohttp
@dataclass
class DealProposal:
"""Filecoin 存储合约提案"""
cid: str
miner: str
duration: int
price: int
collateral: int
verified: bool
start_epoch: int
end_epoch: int
client_addr: str
label: str
@dataclass
class DealStatus:
"""存储合约状态"""
deal_cid: str
state: str
miner: str
client: str
start_epoch: int
end_epoch: int
piece_cid: str
size: int
verified: bool
price: int
collateral: int
sector_start: Optional[int]
class LotusClient:
"""Lotus (Filecoin 节点) RPC 客户端"""
def __init__(self, api_url: str = 'http://localhost:1234/rpc/v0',
auth_token: str = ''):
self.api_url = api_url
self.auth_token = auth_token
self._session: Optional[aiohttp.ClientSession] = None
async def _ensure_session(self):
if self._session is None:
headers = {'Authorization': f'Bearer {self.auth_token}'}
self._session = aiohttp.ClientSession(headers=headers)
async def _call(self, method: str, params: List = None) -> Any:
"""调用 Lotus JSON-RPC 方法"""
await self._ensure_session()
payload = {
'jsonrpc': '2.0', 'id': 1,
'method': f'Filecoin.{method}',
'params': params or [],
}
async with self._session.post(
self.api_url, json=payload,
timeout=aiohttp.ClientTimeout(total=120),
) as resp:
result = await resp.json()
if 'error' in result:
raise RuntimeError(f'Lotus RPC 错误: {result["error"]}')
return result.get('result')
async def deal_add(self, proposal: DealProposal) -> str:
"""提交存储合约提案"""
params = [{
'Data': {
'TransferType': 'graphsync',
'Root': {'/': proposal.cid},
'PieceCid': None,
'PieceSize': None,
},
'Wallet': proposal.client_addr,
'Miner': proposal.miner,
'EpochPrice': str(proposal.price),
'MinBlocksDuration': proposal.duration,
'ProviderCollateral': str(proposal.collateral),
'DealStartEpoch': proposal.start_epoch,
'FastRetrieval': True,
'VerifiedDeal': proposal.verified,
'Label': proposal.label,
}]
return await self._call('ClientStartDeal', params)
async def deal_get_status(self, deal_cid: str) -> DealStatus:
"""获取合约状态"""
result = await self._call('StateMarketStorageDeal', [
int(deal_cid, 16), None,
])
state = result.get('State', {})
proposal = result.get('Proposal', {})
return DealStatus(
deal_cid=deal_cid,
state=self._parse_deal_state(state.get('SectorStartEpoch')),
miner=proposal.get('Provider', ''),
client=proposal.get('Client', ''),
start_epoch=proposal.get('StartEpoch', 0),
end_epoch=proposal.get('EndEpoch', 0),
piece_cid=proposal.get('PieceCID', {}).get('/', ''),
size=proposal.get('PieceSize', 0),
verified=proposal.get('VerifiedDeal', False),
price=int(proposal.get('StoragePricePerEpoch', '0')),
collateral=int(proposal.get('ProviderCollateral', '0')),
sector_start=state.get('SectorStartEpoch'),
)
def _parse_deal_state(self, sector_start_epoch: Optional[int]) -> str:
if sector_start_epoch is None or sector_start_epoch < 0:
return 'proposed'
elif sector_start_epoch > 0:
return 'active'
return 'unknown'
async def deal_list(self) -> List[Dict]:
return await self._call('ClientListDeals')
async def deal_retrieve(self, cid: str, miner: str, out_file: str) -> bool:
params = [{'Root': {'/': cid}}, out_file, miner]
result = await self._call('ClientRetrieve', params)
return result is not None
async def close(self):
if self._session:
await self._session.close()
self._session = None
class FilecoinClient:
"""Filecoin 存储客户端"""
def __init__(self, lotus_client: LotusClient,
default_duration_days: int = 365,
default_replicas: int = 3):
self.lotus = lotus_client
self.default_duration = default_duration_days * 86400
self.default_replicas = default_replicas
async def make_deal(self, cid: str, miner: str,
duration: Optional[int] = None,
price: int = 100000) -> str:
"""创建存储合约"""
duration = duration or self.default_duration
proposal = DealProposal(
cid=cid, miner=miner, duration=duration, price=price,
collateral=0, verified=False, start_epoch=0,
end_epoch=duration, client_addr='',
label=f'msg_chain_agent_{cid[:16]}',
)
return await self.lotus.deal_add(proposal)
async def replicate_deals(self, cid: str, miners: List[str],
duration: Optional[int] = None,
price: int = 100000) -> Dict[str, str]:
"""向多个矿工创建冗余存储合约"""
tasks = [
self.make_deal(cid, miner, duration, price)
for miner in miners[:self.default_replicas]
]
deal_cids = await asyncio.gather(*tasks, return_exceptions=True)
results = {}
for miner, result in zip(miners, deal_cids):
if isinstance(result, Exception):
results[miner] = f'error: {result}'
else:
results[miner] = result
return results
async def check_deal_health(self, deal_cid: str) -> Dict:
"""检查合约健康状态"""
try:
status = await self.lotus.deal_get_status(deal_cid)
return {
'deal_cid': deal_cid,
'healthy': status.state == 'active',
'state': status.state,
'miner': status.miner,
'remaining_epochs': max(0, status.end_epoch - status.start_epoch),
}
except Exception as e:
return {'deal_cid': deal_cid, 'healthy': False, 'error': str(e)}
async def verify_storage_proof(self, deal_cid: str) -> Dict:
"""验证存储证明"""
status = await self.lotus.deal_get_status(deal_cid)
if status.state != 'active':
return {'verified': False, 'reason': '合约未激活'}
miner_info = await self.lotus._call('StateMinerInfo', [status.miner, None])
return {
'verified': True,
'miner': status.miner,
'sector_start': status.sector_start,
}
async def retrieve_data(self, cid: str, miner: Optional[str] = None) -> Optional[bytes]:
"""从 Filecoin 检索数据"""
import tempfile, os
if miner:
miners = [miner]
else:
miners = ['t01000']
for m in miners:
with tempfile.NamedTemporaryFile(delete=False) as tmp:
out_path = tmp.name
try:
success = await self.lotus.deal_retrieve(cid, m, out_path)
if success:
with open(out_path, 'rb') as f:
data = f.read()
os.unlink(out_path)
return data
except Exception:
continue
finally:
if os.path.exists(out_path):
os.unlink(out_path)
return None
async def estimate_deal_cost(self, size_bytes: int, duration_days: int = 365) -> Dict:
"""估算存储合约成本"""
size_gb = size_bytes / (1024**3)
epochs = duration_days * 86400 / 30
estimated_fil = size_gb * 0.001 * epochs
return {
'size_gb': round(size_gb, 4),
'duration_days': duration_days,
'estimated_fil': round(estimated_fil, 6),
'estimated_usd': round(estimated_fil * 5, 2),
'recommended_miners': ['t01000', 't01001', 't01002'],
}
async def renew_deal(self, deal_cid: str, additional_days: int = 365) -> str:
"""续约存储合约"""
status = await self.lotus.deal_get_status(deal_cid)
return await self.make_deal(status.piece_cid, status.miner,
additional_days * 86400)
6.3 Filecoin 与 IPFS 联合使用
class FilecoinIPFSBridge:
"""Filecoin 与 IPFS 桥接"""
def __init__(self, ipfs_client: IPFSClient, filecoin_client: FilecoinClient):
self.ipfs = ipfs_client
self.filecoin = filecoin_client
async def archive_to_filecoin(self, cid: str, miners: List[str] = None) -> Dict:
"""将 IPFS 内容归档到 Filecoin"""
if miners is None:
miners = ['t01000', 't01001']
data = await self.ipfs.resolve(cid)
if data is None:
raise ValueError(f'CID {cid} 在 IPFS 上不可用')
deal_results = await self.filecoin.replicate_deals(cid, miners)
await self.ipfs.pin_cid(cid)
return {
'cid': cid,
'deals': deal_results,
'retrieval_mode': 'hybrid',
'hot_storage': 'ipfs',
'cold_storage': 'filecoin',
}
async def retrieve_from_archive(self, cid: str) -> Optional[bytes]:
"""从归档中检索数据(优先从 IPFS)"""
data = await self.ipfs.resolve(cid)
if data is not None:
return data
data = await self.filecoin.retrieve_data(cid)
if data is not None:
await self.ipfs.store_bytes(data)
await self.ipfs.pin_cid(cid)
return data
6.4 存储验证器
class StorageVerifier:
"""去中心化存储验证器"""
def __init__(self, ipfs_client: IPFSClient, arweave_client: ArweaveClient,
filecoin_client: FilecoinClient, chain_client: ChainAnchorClient):
self.ipfs = ipfs_client
self.arweave = arweave_client
self.filecoin = filecoin_client
self.chain = chain_client
async def verify_all(self, cid: str, expected_data: dict) -> Dict:
"""在所有存储层验证数据"""
results = {
'chain': await self._verify_chain(cid, expected_data),
'ipfs': await self._verify_ipfs(cid, expected_data),
}
try:
arw_tx = await self.chain.query_anchor(f'arweave:{cid}')
if arw_tx:
results['arweave'] = await self._verify_arweave(cid, expected_data)
except Exception:
pass
return {
'cid': cid,
'all_verified': all(r.get('verified', False) for r in results.values()),
'results': results,
}
async def _verify_chain(self, cid: str, data: dict) -> Dict:
try:
anchor = await self.chain.query_anchor(cid)
if not anchor:
return {'verified': False, 'error': '链上未找到锚定'}
serialized = json.dumps(data, ensure_ascii=False, default=str)
checksum = hashlib.sha256(serialized.encode()).hexdigest()
return {
'verified': checksum == anchor.checksum,
'owner': anchor.owner,
'timestamp': anchor.timestamp,
'checksum_match': checksum == anchor.checksum,
}
except Exception as e:
return {'verified': False, 'error': str(e)}
async def _verify_ipfs(self, cid: str, data: dict) -> Dict:
try:
stored_data = await self.ipfs.get_agent_data(cid)
return {'verified': True, 'content_match': stored_data == data, 'available': True}
except Exception as e:
return {'verified': False, 'available': False, 'error': str(e)}
async def _verify_arweave(self, cid: str, data: dict) -> Dict:
try:
tx_data = await self.arweave.get_json_data(cid)
return {'verified': True, 'content_match': tx_data == data, 'permanent': True}
except Exception as e:
return {'verified': False, 'error': str(e)}
7. Agent 数据管理示例
7.1 完整的数据管理器
import asyncio
import json
import logging
from typing import Optional, List, Dict, Any, Callable
from dataclasses import dataclass, field
from datetime import datetime
from enum import Enum
import hashlib
# MSG Chain 常量
MSG_BECH32_PREFIX = 'msg'
MSG_CHAIN_ID = 'msg-chain-1'
class DataCategory(Enum):
"""Agent 数据分类"""
CONVERSATION = 'conversation'
KNOWLEDGE_BASE = 'knowledge_base'
MODEL_STATE = 'model_state'
USER_PREFERENCE = 'user_preference'
AUDIT_LOG = 'audit_log'
TASK_RESULT = 'task_result'
AGENT_CONFIG = 'agent_config'
TRAINING_DATA = 'training_data'
@dataclass
class AgentDataRecord:
"""Agent 数据记录"""
record_id: str
agent_id: str
category: DataCategory
data: dict
created_at: str = field(default_factory=lambda: datetime.utcnow().isoformat())
storage_tier: Optional[str] = None
storage_ref: Optional[str] = None
chain_anchor: Optional[str] = None
size: int = 0
checksum: str = ''
class AgentDataManager:
"""
Agent 数据管理器
统一管理 Agent 在去中心化存储中的生命周期
"""
def __init__(
self,
agent_id: str,
chain_contract_addr: str,
ipfs_node_url: str = 'http://localhost:5001',
arweave_wallet_path: str = '',
filecoin_api_url: str = 'http://localhost:1234/rpc/v0',
monthly_budget: float = 10.0,
enable_auto_tiering: bool = True,
):
self.agent_id = agent_id
self.chain_contract = chain_contract_addr
self.enable_auto_tiering = enable_auto_tiering
self.ipfs = self._init_ipfs(ipfs_node_url)
self.strategy = StorageStrategy(budget_monthly=monthly_budget)
self.lifecycle = DataLifecycleManager(self.strategy)
self.chain_client = ChainAnchorClient(chain_contract_addr)
if arweave_wallet_path:
wallet = ArweaveWallet.from_file(arweave_wallet_path)
self.arweave = ArweaveClient(wallet)
else:
self.arweave = None
if filecoin_api_url:
lotus = LotusClient(filecoin_api_url)
self.filecoin = FilecoinClient(lotus)
else:
self.filecoin = None
self._local_index: Dict[str, AgentDataRecord] = {}
assert isinstance(agent_id, str) and len(agent_id) > 0
def _init_ipfs(self, node_url: str) -> IPFSClient:
return IPFSClient(node_url)
async def start(self):
"""启动数据管理器"""
await self.ipfs.connect()
logging.info(
f'Agent 数据管理器已启动: {self.agent_id} '
f'[chain: {MSG_CHAIN_ID}, prefix: {MSG_BECH32_PREFIX}]'
)
async def store_conversation_log(
self, session_id: str, messages: list,
metadata: Optional[Dict] = None, tier: Optional[str] = None,
) -> str:
"""存储对话日志"""
record_id = f'{self.agent_id}/conversation/{session_id}'
data = {
'agent_id': self.agent_id,
'session_id': session_id,
'messages': messages,
'metadata': metadata or {},
'chain_context': {
'chain_id': MSG_CHAIN_ID,
'bech32_prefix': MSG_BECH32_PREFIX,
},
}
record = AgentDataRecord(
record_id=record_id, agent_id=self.agent_id,
category=DataCategory.CONVERSATION, data=data,
)
storage_ref = await self._store_with_strategy(
data, record_id, DataCategory.CONVERSATION, tier
)
record.storage_ref = storage_ref
record.storage_tier = tier or await self._auto_detect_tier(data)
self._local_index[record_id] = record
return storage_ref
async def store_knowledge_entry(
self, entry_id: str, content: dict, tags: List[str] = None,
) -> str:
"""存储知识库条目"""
record_id = f'{self.agent_id}/knowledge/{entry_id}'
data = {
'agent_id': self.agent_id, 'entry_id': entry_id,
'content': content, 'tags': tags or [], 'version': 1,
}
record = AgentDataRecord(
record_id=record_id, agent_id=self.agent_id,
category=DataCategory.KNOWLEDGE_BASE, data=data,
)
if self.arweave:
tags_ar = [
ArweaveTag('Agent-ID', self.agent_id),
ArweaveTag('Data-Type', 'knowledge'),
ArweaveTag('Entry-ID', entry_id),
ArweaveTag('Chain-ID', MSG_CHAIN_ID),
]
if tags:
tags_ar.extend(ArweaveTag('Tag', t) for t in tags)
tx_id = await self.arweave.store_json_permanently(
data, tags_ar, f'knowledge-{entry_id}'
)
storage_ref = tx_id
cid = await self.ipfs.store_agent_data(data)
await self.ipfs.pin_cid(cid)
else:
cid = await self.ipfs.store_agent_data(data)
await self.ipfs.pin_cid(cid)
storage_ref = cid
anchor_cid = await self.chain_client.anchor_data(
storage_ref, 'knowledge', data, self.agent_id
)
record.storage_ref = storage_ref
record.chain_anchor = anchor_cid
self._local_index[record_id] = record
return storage_ref
async def store_audit_log(self, action: str, details: dict,
timestamp: Optional[str] = None) -> str:
"""存储审计日志"""
record_id = f'{self.agent_id}/audit/{action}/{timestamp or datetime.utcnow().isoformat()}'
log_entry = {
'agent_id': self.agent_id, 'action': action, 'details': details,
'timestamp': timestamp or datetime.utcnow().isoformat(),
'chain_id': MSG_CHAIN_ID,
}
record = AgentDataRecord(
record_id=record_id, agent_id=self.agent_id,
category=DataCategory.AUDIT_LOG, data=log_entry,
)
if not self.arweave:
raise RuntimeError('审计日志需要 Arweave 支持')
tags = [
ArweaveTag('Agent-ID', self.agent_id),
ArweaveTag('Action', action),
ArweaveTag('Category', 'audit'),
]
tx_id = await self.arweave.store_json_permanently(
log_entry, tags, f'audit-{action}'
)
record.storage_ref = tx_id
self._local_index[record_id] = record
return tx_id
async def get_data(self, storage_ref: str, source: str = 'auto') -> Optional[dict]:
"""从存储中获取数据"""
if source in ('auto', 'ipfs'):
try:
data = await self.ipfs.get_agent_data(storage_ref)
if data:
return data
except Exception:
pass
if source in ('auto', 'arweave') and self.arweave:
try:
data = await self.arweave.get_json_data(storage_ref)
if data:
return data
except Exception:
pass
if source in ('auto', 'filecoin') and self.filecoin:
try:
bytes_data = await self.filecoin.retrieve_data(storage_ref)
if bytes_data:
return json.loads(bytes_data.decode('utf-8'))
except Exception:
pass
return None
async def list_records(self, category: Optional[DataCategory] = None,
limit: int = 50) -> List[AgentDataRecord]:
"""列出本地索引中的数据记录"""
records = list(self._local_index.values())
if category:
records = [r for r in records if r.category == category]
return records[:limit]
async def verify_data_integrity(self, storage_ref: str) -> Dict:
"""验证数据完整性"""
anchor = await self.chain_client.query_anchor(storage_ref)
if not anchor:
return {'verified': False, 'error': '链上未找到锚定记录'}
data = await self.get_data(storage_ref)
if not data:
return {'verified': False, 'error': '无法从存储获取数据'}
serialized = json.dumps(data, ensure_ascii=False, default=str)
actual_checksum = hashlib.sha256(serialized.encode()).hexdigest()
match = actual_checksum == anchor.checksum
return {
'verified': match,
'storage_ref': storage_ref,
'chain_checksum': anchor.checksum,
'actual_checksum': actual_checksum,
'owner': anchor.owner,
'timestamp': anchor.timestamp,
}
async def _store_with_strategy(self, data: dict, record_id: str,
category: DataCategory,
forced_tier: Optional[str] = None) -> str:
"""根据策略存储数据"""
serialized = json.dumps(data, ensure_ascii=False, default=str)
size = len(serialized.encode())
profile = DataProfile(
data_id=record_id, size_bytes=size,
access_frequency=self.lifecycle.get_frequency(record_id),
access_pattern='random', required_durability=0.9999,
required_latency_ms=5000, ttl_seconds=86400 * 30,
budget_per_gb=self.strategy.budget_monthly,
immutable=category in (DataCategory.AUDIT_LOG, DataCategory.TRAINING_DATA),
sensitivity='public' if category == DataCategory.KNOWLEDGE_BASE else 'private',
)
if forced_tier:
tier = StorageTier(forced_tier)
else:
decision = self.strategy.decide(profile)
tier = decision.tier
if tier == StorageTier.HOT:
return f'cache:{record_id}'
elif tier == StorageTier.WARM:
cid = await self.ipfs.store_agent_data(data)
await self.ipfs.pin_cid(cid)
return cid
elif tier == StorageTier.COLD:
if not self.arweave:
raise RuntimeError('需要 Arweave 客户端')
tags = [
ArweaveTag('Agent-ID', self.agent_id),
ArweaveTag('Record-ID', record_id),
ArweaveTag('Category', category.value),
]
return await self.arweave.store_json_permanently(data, tags)
elif tier == StorageTier.ARCHIVE:
if not self.filecoin:
raise RuntimeError('需要 Filecoin 客户端')
cid = await self.ipfs.store_agent_data(data)
deal_cid = await self.filecoin.make_deal(cid, 't01000', 365 * 86400)
return deal_cid
raise ValueError(f'未知存储层级: {tier}')
async def _auto_detect_tier(self, data: dict) -> str:
size = len(json.dumps(data, default=str).encode())
if size < 1024:
return 'hot'
elif size < 1024 * 1024:
return 'warm'
else:
return 'cold'
async def cleanup_old_data(self, max_age_days: int = 90):
"""清理旧数据"""
cutoff = datetime.utcnow().timestamp() - (max_age_days * 86400)
for record_id, record in list(self._local_index.items()):
created = datetime.fromisoformat(record.created_at).timestamp()
if created < cutoff:
if record.storage_tier in ('hot', 'warm') and record.storage_ref:
pass
del self._local_index[record_id]
async def close(self):
"""关闭数据管理器"""
await self.ipfs.close()
if self.arweave:
await self.arweave.close()
logging.info(f'Agent 数据管理器已关闭: {self.agent_id}')
class AgentDataFactory:
"""Agent 数据工厂"""
def __init__(self):
self._managers: Dict[str, AgentDataManager] = {}
def create_manager(self, agent_id: str, **kwargs) -> AgentDataManager:
if agent_id in self._managers:
raise ValueError(f'Agent {agent_id} 已存在')
manager = AgentDataManager(agent_id=agent_id, **kwargs)
self._managers[agent_id] = manager
return manager
def get_manager(self, agent_id: str) -> Optional[AgentDataManager]:
return self._managers.get(agent_id)
async def close_all(self):
for manager in self._managers.values():
await manager.close()
self._managers.clear()
7.2 端到端使用示例
async def main_example():
"""
端到端示例:演示 Agent 如何使用去中心化存储
"""
# 1. 创建数据管理器
factory = AgentDataFactory()
manager = factory.create_manager(
agent_id='msg_agent_demo_001',
chain_contract_addr='msg1contractaddress...',
arweave_wallet_path='./wallet.json',
monthly_budget=15.0,
)
await manager.start()
try:
# 2. 存储对话日志
print('=== 存储对话日志到 IPFS ===')
session_id = 'session_通用维护记录_001'
messages = [
{'role': 'user', 'content': '什么是 MSG Chain?', 'timestamp': '2026-12-01T10:00:00Z'},
{'role': 'assistant', 'content': 'MSG Chain 是专为 AI Agent 设计的 L1 区块链。', 'timestamp': '2026-12-01T10:00:05Z'},
{'role': 'user', 'content': '如何集成去中心化存储?', 'timestamp': '2026-12-01T10:01:00Z'},
{'role': 'assistant', 'content': 'Agent 可通过 IPFS/Arweave/Filecoin 存储数据。', 'timestamp': '2026-12-01T10:01:08Z'},
]
cid = await manager.store_conversation_log(
session_id=session_id,
messages=messages,
metadata={'topic': 'decentralized_storage'},
)
print(f'对话已存储,CID: {cid}')
# 3. 存储知识库条目
print()
print('=== 存储知识库条目到 Arweave ===')
knowledge = {
'title': 'MSG Chain 存储架构',
'content': 'MSG Chain 支持 IPFS 内容寻址和链上锚定。Agent 可以通过 DataAnchor 合约将 CID 永久记录在链上。',
'references': ['https://docs.msgchain.org/storage'],
}
tx_id = await manager.store_knowledge_entry(
entry_id='kb_storage_arch_001',
content=knowledge,
tags=['architecture', 'storage', 'ipfs'],
)
print(f'知识已永久存储,Arweave TX: {tx_id}')
# 4. 记录审计日志
print()
print('=== 记录审计日志到 Arweave ===')
audit_tx = await manager.store_audit_log(
action='knowledge_update',
details={
'entry_id': 'kb_storage_arch_001',
'change': 'initial_create',
},
)
print(f'审计日志已存储,Arweave TX: {audit_tx}')
# 5. 验证数据完整性
print()
print('=== 验证数据完整性 ===')
verification = await manager.verify_data_integrity(cid)
print(f'数据完整性验证: {"通过" if verification["verified"] else "失败"}')
print(f' 链上校验和: {verification["chain_checksum"][:20]}...')
print(f' 实际校验和: {verification["actual_checksum"][:20]}...')
# 6. 检索数据
print()
print('=== 从去中心化存储检索数据 ===')
retrieved = await manager.get_data(cid, source='ipfs')
if retrieved:
msg_cnt = len(retrieved.get('messages', []))
print(f'检索成功: {msg_cnt} 条消息')
print(f' 会话: {retrieved.get("session_id")}')
# 7. 列出所有记录
print()
print('=== Agent 数据记录列表 ===')
records = await manager.list_records()
for rec in records:
print(f' [{rec.category.value}] {rec.record_id}')
print(f' 存储: {rec.storage_tier} -> {rec.storage_ref[:32]}...')
finally:
await manager.close()
await factory.close_all()
if __name__ == '__main__':
asyncio.run(main_example())
7.3 最佳实践总结
"""
MSG Chain AI Agent 去中心化存储最佳实践
1. 数据分类分级
- 对话日志: IPFS (温存储),7 天后迁移至 Arweave
- 知识库: Arweave (永久存储),保留 IPFS 缓存用于快速检索
- 审计日志: Arweave (永久存储,不可篡改)
- 用户偏好: 本地缓存 + IPFS 备份
- 模型状态: IPFS + Filecoin (大规模归档)
- 训练数据: Filecoin (长期归档,可证明存储)
2. 地址命名规范
- MSG Chain 地址: msg1 前缀 (bech32)
- Agent ID 格式: msg_agent_{network}_{id}
- 合约命名: msg_data_anchor, msg_storage_index
3. 成本控制策略
- 小数据 (< 1KB): 直接链上存储或本地缓存
- 中等数据 (1KB ~ 1MB): IPFS 优先
- 大数据 (1MB ~ 1GB): Arweave 一次性付费
- 超大数据 (> 1GB): Filecoin 存储合约
- 月度存储预算: 通过 CostOptimizer 自动分配
4. 数据完整性保证
- 每次存储计算 SHA-256 checksum
- 将 checksum 和 CID 锚定到 MSG Chain
- 定期通过 StorageVerifier 验证所有层数据
- 链下数据 + 链上锚定的双重验证
5. 安全注意事项
- 密钥管理: 使用独立 Arweave 钱包存储审计日志
- 访问控制: MSG Chain 合约控制锚定写入权限
- 加密: 敏感数据在存储前进行客户端加密
- 备份: 关键数据至少存储在 2 个不同的层
6. MSG Chain 集成要点
- 所有存储引用最终锚定到 msg-chain-1
- 使用 msg1 地址格式进行身份标识
- 通过 DataAnchor 合约统一管理数据索引
- bech32 前缀 msg 贯穿所有链上交互
"""
附录
A. 依赖安装
# Python 依赖
pip install ipfshttpclient aiohttp multiformats multihash cryptography
# Rust 依赖(Cargo.toml)
# cosmwasm-std = "1.5"
# cw-storage-plus = "1.2"
# serde = { version = "1.0", features = ["derive"] }
# schemars = "0.8"
# MSG Chain CLI
# 从 https://github.com/msgchain/msgchaind 下载
# 或使用预编译二进制
B. 快速启动脚本
#!/usr/bin/env python3
"""
快速启动脚本:初始化 Agent 去中心化存储环境
"""
import asyncio
import sys
async def quickstart():
print('MSG Chain AI Agent 去中心化存储快速启动')
print('=' * 50)
print(f'链 ID: msg-chain-1')
print(f'地址前缀: msg')
print()
# 检查依赖
try:
import ipfshttpclient
print('[OK] ipfshttpclient 已安装')
except ImportError:
print('[!] 请安装: pip install ipfshttpclient')
try:
import aiohttp
print('[OK] aiohttp 已安装')
except ImportError:
print('[!] 请安装: pip install aiohttp')
print()
# 验证环境
print('环境准备完成!')
print()
print('下一步:')
print(' 1. 启动 IPFS 节点: ipfs daemon')
print(' 2. 启动 MSG Chain 节点: msgchaind start')
print(' 3. 部署 DataAnchor 合约')
print(' 4. 运行 Agent 数据管理示例')
print()
print('详细指南请参考本文档各章节。')
if __name__ == '__main__':
asyncio.run(quickstart())
