MSG Chain 去中心化存储集成指南:IPFS 与 Arweave
数据来源:MSG Chain 代码库核实
主网状态: No-Go — 当前 MSGChain 主网裁决为 No-Go,以下内容反映代码实际状态,不代表生产可用。
目录
1. 概述
1.1 为什么在区块链 DApp 中使用去中心化存储
区块链本身不是为存储大量数据而设计的。MSG Chain 的每个区块都有大小限制(通常为 MB 级别),而将大文件(图片、视频、结构化数据集等)直接写入链上会导致:
- 存储成本急剧上升:每字节需支付 Gas 费用,大文件存储成本难以接受
- 网络吞吐量下降:区块空间被数据占用,影响交易处理速度
- 全节点存储膨胀:节点需存储所有链上数据,降低去中心化程度
去中心化存储方案解决了这个问题:链上存证明(指针/哈希),链下存数据本体。
典型架构:
+------------------+ CID / txid +--------------------+
| MSG Chain 合约 | <---------------------> | 去中心化存储网络 |
| (存储哈希+指针) | | (IPFS / Arweave) |
+------------------+ +--------------------+
^ ^
| |
v v
DApp 前端 / SDK 网关 / 节点读取
这种架构的优势:
- 低成本:链上只需存储少量字节的 CID/哈希
- 高吞吐:数据上传不占用区块空间
- 可扩展:理论上可存储任意大小的文件
- 内容寻址:数据内容本身决定地址,天然防篡改
1.2 IPFS:内容寻址存储
IPFS(InterPlanetary File System) 是一个点对点的分布式文件系统,核心概念:
| 概念 | 说明 |
|---|---|
| CID | 内容标识符(Content Identifier),基于文件内容的哈希生成。相同内容产生相同 CID。格式:Qm...(v0)或 bafy...(v1) |
| 内容寻址 | 通过 CID 而非路径定位文件。ipfs cat QmX... |
| IPNS | 星际名称系统(Inter-Planetary Name System),提供可变指针指向可变 CID。DNS 风格的名称解析 |
| DNSLink | 通过 DNS TXT 记录将域名指向 IPFS 内容,实现可变站点 |
| Pinning | 固定内容以确保节点保存该数据,不被垃圾回收 |
| MFS | 可变文件系统(Mutable File System),类似传统文件系统的目录结构 |
IPFS 的核心优点:
- 内容去重:相同文件只存储一次
- 带宽节省:P2P 传输,从最近节点获取
- 离线可用:数据缓存在多个节点
- 网络弹性:无单点故障
- 生态丰富:大量工具和 pinning 服务
IPFS 的局限性:
- 持久性依赖 pinning:如果没有任何节点 pin 数据,数据将丢失
- 可能需要 pinning 服务费用
- 公共网关可能限速或屏蔽
1.3 Arweave:永久存储
Arweave 是专为永久数据存储设计的区块链网络,核心特点:
| 概念 | 说明 |
|---|---|
| Blockweave | Arweave 的底层数据结构,每个区块链接到前一个区块和之前的一个 recall block |
| 一次付费,永久存储 | 用户预付存储费,资金通过 endowment 机制产生利息持续支付给矿工 |
| Transaction | 上传数据的基本单位,每个 tx 包含数据和标签 |
| Tags | 键值对标签,用于索引和查询数据 |
| Wallet | JWK(JSON Web Key)格式的 RSA 密钥对 |
Arweave 的核心优点:
- 真正的永久存储,无需持续支付固定费用
- 一次写入,不可篡改
- 丰富的标签系统支持高效查询
- 与 IPFS 互补:IPFS 适合频繁访问的热数据,Arweave 适合长期归档
- 通过 Bundlr/Irys 实现即时确认
Arweave 的局限性:
- 每次上传需支付 AR 代币
- 上传确认时间可能较长(原生)
- 查询灵活性不如中心化存储
1.4 IPFS vs Arweave 对比
| 维度 | IPFS | Arweave |
|---|---|---|
| 存储模型 | 内容寻址,需 pinning 保证可用 | 永久存储,一次付费 |
| 持久性 | 需手动 pin 或使用 pinning 服务 | 协议级保证永久存储 |
| 可变性 | 通过 IPNS/DNSLink | 不可变(新 tx 更新指针) |
| 费用模型 | 存储免费,pinning 服务收费 | 按数据量一次付费 |
| 上传速度 | 取决于节点连接 | 通过 Bundlr/Irys 快速确认 |
| 查询能力 | 基于 CID 精确查找 | 基于标签的键值对查询 |
| 数据大小限制 | 无硬限制(分块) | 取决于交易大小 |
| 隐私 | 公开(可加密) | 公开(可加密) |
| 适合场景 | NFT 媒体,频繁更新的元数据 | 长期归档,法规合规,数据集 |
1.5 MSG Chain 上的典型应用场景
NFT 元数据和媒体
MSG Chain 上的 NFT(CW721)配合 IPFS/Arweave 存储元数据和媒体文件:
链上: NFT 合约 (token_id, owner, cid)
链下 IPFS: metadata.json (name, description, image)
链下 IPFS: image.png (实际媒体文件)
// 合约中存储的仅为 CID 字符串
pub struct NFTMetadata {
pub cid: String, // IPFS CID 指向 metadata.json
pub updated_at: Timestamp,
}
数据市场文件
链上: DataAsset 合约 (asset_id, owner, cid, hash_commitment, price)
链下 Arweave: dataset.zip (实际数据文件,永久保存)
DAO 治理文档
链上: Proposal 合约 (proposal_id, title, description_cid, votes)
链下 IPFS: proposal_detail.md (详细的提案文档,可更新)
链下 Arweave: final_report.pdf (最终报告,不可篡改)
去中心化身份(DID)文档
链上: DID Registry (did, controller, document_cid)
链下 IPFS: did_document.json (公钥、服务端点等)
2. IPFS 集成基础
2.1 IPFS CLI 基础操作
安装 IPFS
# 下载 Kubo (Go 实现)
wget https://dist.ipfs.tech/kubo/v0.29.0/kubo_v0.29.0_linux-amd64.tar.gz
tar -xzf kubo_v0.29.0_linux-amd64.tar.gz
cd kubo
sudo bash install.sh
# 初始化
ipfs init
# 启动守护进程
ipfs daemon
# 检查版本
ipfs --version
# 输出: ipfs version 0.29.0
核心命令
# ========== 添加文件 ==========
# 添加单个文件,返回 CID
echo "Hello MSG Chain" > hello.txt
ipfs add hello.txt
# 输出: added QmZ4tDuMPZqG3qgCnfnqGvkKzBiV5z7MNM9q8Jw8jYzB6c hello.txt
# 添加目录(递归),-r 参数
mkdir nft-assets
cat > nft-assets/metadata.json << 'JSON'
{"name":"MSG NFT #1","description":"A test NFT"}
JSON
echo "fake_image_data" > nft-assets/image.png
ipfs add -r nft-assets/
# 输出:
# added QmY8... metadata.json
# added QmX9... image.png
# added QmDir... nft-assets
# 添加文件并指定 pin(自动固定)
ipfs add --pin=true large-dataset.csv
# 分块选项:指定块大小(默认 256KB)
ipfs add --chunker=size-1024 largefile.bin
# CID 版本控制
ipfs add --cid-version=1 data.txt
# 输出: added bafybei... data.txt (v1 CID)
# v1 优点: 区分大小写、支持多哈希、Base32 编码
# ========== 读取文件 ==========
# 通过 CID 读取文件
ipfs cat QmZ4tDuMPZqG3qgCnfnqGvkKzBiV5z7MNM9q8Jw8jYzB6c
# 输出: Hello MSG Chain
# 列出目录内容
ipfs ls QmDir...
# 读取到文件
ipfs get QmZ4tDuMPZqG3qgCnfnqGvkKzBiV5z7MNM9q8Jw8jYzB6c -o downloaded.txt
# ========== Pinning(固定内容)==========
# 固定 CID(确保本地保存该内容)
ipfs pin add QmZ4tDuMPZqG3qgCnfnqGvkKzBiV5z7MNM9q8Jw8jYzB6c
# 列出所有固定内容
ipfs pin ls
# 按类型列出
ipfs pin ls --type=recursive
ipfs pin ls --type=indirect
ipfs pin ls --type=direct
# 移除固定(允许垃圾回收)
ipfs pin rm QmZ4tDuMPZqG3qgCnfnqGvkKzBiV5z7MNM9q8Jw8jYzB6c
# 垃圾回收(清理未固定内容)
ipfs repo gc
# ========== IPNS(可变指针)==========
# 发布当前节点根哈希到 IPNS
ipfs name publish /ipfs/QmDir...
# 输出: Published to k51qzi5uqu5d...
# 解析 IPNS 名称
ipfs name resolve k51qzi5uqu5d...
# 输出: /ipfs/QmDir...
# ========== 文件系统(MFS)==========
# 创建目录
ipfs files mkdir /my-nft-collection
# 复制文件到 MFS
ipfs files cp /ipfs/QmX... /my-nft-collection/metadata.json
# 列出目录
ipfs files ls /my-nft-collection
# 查看 MFS 根哈希(可用于 IPNS 发布)
ipfs files stat /
2.2 IPFS HTTP API
IPFS 守护进程默认在 localhost:5001 暴露 HTTP API。所有 CLI 命令都有对应的 API 端点。
API 基础
POST /api/v0/add - 添加文件
POST /api/v0/cat - 读取文件
POST /api/v0/pin/add - 固定内容
POST /api/v0/pin/rm - 移除固定
POST /api/v0/name/publish - 发布 IPNS
POST /api/v0/dag/get - 获取 DAG 节点
Python 客户端
"""
ipfs_client.py - IPFS HTTP API Python 客户端
支持:添加文件、读取、pin、dag 操作
"""
import json
import os
import hashlib
from typing import Optional, Dict, Any
import requests
class IPFSClient:
"""IPFS HTTP API 客户端"""
def __init__(
self,
api_url: str = "http://127.0.0.1:5001",
gateway_url: str = "https://ipfs.io/ipfs",
):
self.api_url = api_url.rstrip("/")
self.gateway_url = gateway_url.rstrip("/")
self.api_version = "v0"
def _api_url(self, endpoint: str) -> str:
return f"{self.api_url}/api/{self.api_version}/{endpoint.lstrip('/')}"
# ---------- 添加文件 ----------
def add_file(
self,
file_path: str,
pin: bool = True,
cid_version: int = 1,
only_hash: bool = False,
) -> Dict[str, Any]:
"""
添加文件到 IPFS
Args:
file_path: 本地文件路径
pin: 是否自动固定
cid_version: CID 版本(0 或 1)
only_hash: 仅计算哈希不上传
Returns:
{"Name": "filename", "Hash": "Qm...", "Size": "1234"}
"""
params = {
"pin": str(pin).lower(),
"cid-version": str(cid_version),
"only-hash": str(only_hash).lower(),
}
url = self._api_url("add")
with open(file_path, "rb") as f:
files = {"file": (os.path.basename(file_path), f)}
resp = requests.post(url, params=params, files=files, timeout=60)
resp.raise_for_status()
return resp.json()
def add_bytes(
self,
data: bytes,
filename: str = "data.bin",
pin: bool = True,
cid_version: int = 1,
) -> Dict[str, Any]:
"""直接添加 bytes 数据"""
params = {
"pin": str(pin).lower(),
"cid-version": str(cid_version),
}
url = self._api_url("add")
files = {"file": (filename, data)}
resp = requests.post(url, params=params, files=files, timeout=60)
resp.raise_for_status()
return resp.json()
def add_json(
self,
obj: Any,
filename: str = "data.json",
pin: bool = True,
cid_version: int = 1,
) -> Dict[str, Any]:
"""添加 JSON 对象到 IPFS"""
data = json.dumps(obj, ensure_ascii=False).encode("utf-8")
return self.add_bytes(data, filename=filename, pin=pin, cid_version=cid_version)
def add_directory(
self,
dir_path: str,
pin: bool = True,
cid_version: int = 1,
) -> list[Dict[str, Any]]:
"""
递归添加目录到 IPFS
最后一个元素的 Hash 为目录 CID
"""
params = {
"pin": str(pin).lower(),
"cid-version": str(cid_version),
"recursive": "true",
}
url = self._api_url("add")
files = []
for root, _, filenames in os.walk(dir_path):
for filename in filenames:
full_path = os.path.join(root, filename)
arcname = os.path.relpath(full_path, dir_path)
files.append(("file", (arcname, open(full_path, "rb"))))
resp = requests.post(url, params=params, files=files, timeout=300)
resp.raise_for_status()
# 每行一个 JSON
results = [json.loads(line) for line in resp.text.strip().split("\n")]
return results
# ---------- 读取文件 ----------
def cat(self, cid: str, timeout: int = 30) -> bytes:
"""通过 CID 读取文件内容"""
url = self._api_url("cat")
params = {"arg": cid}
resp = requests.post(url, params=params, timeout=timeout)
resp.raise_for_status()
return resp.content
def cat_json(self, cid: str) -> Any:
"""读取 JSON 文件"""
content = self.cat(cid)
return json.loads(content.decode("utf-8"))
def get(self, cid: str, output_dir: str) -> bytes:
"""下载文件/目录到本地"""
url = self._api_url("get")
params = {"arg": cid, "output": output_dir}
resp = requests.post(url, params=params, timeout=300)
resp.raise_for_status()
return resp.content
# ---------- 网关访问(只读)----------
def gateway_url_for(self, cid: str) -> str:
"""生成公共网关 URL"""
return f"{self.gateway_url}/{cid}"
def fetch_via_gateway(self, cid: str, gateway: Optional[str] = None) -> bytes:
"""通过公共网关获取内容"""
base = gateway.rstrip("/") if gateway else self.gateway_url
url = f"{base}/{cid}"
resp = requests.get(url, timeout=30)
resp.raise_for_status()
return resp.content
# ---------- Pinning ----------
def pin_add(self, cid: str) -> Dict[str, Any]:
"""固定 CID"""
url = self._api_url("pin/add")
params = {"arg": cid}
resp = requests.post(url, params=params, timeout=30)
resp.raise_for_status()
return resp.json()
def pin_rm(self, cid: str) -> Dict[str, Any]:
"""移除固定"""
url = self._api_url("pin/rm")
params = {"arg": cid}
resp = requests.post(url, params=params, timeout=30)
resp.raise_for_status()
return resp.json()
def pin_ls(self, cid: Optional[str] = None) -> Dict[str, Any]:
"""列出固定"""
url = self._api_url("pin/ls")
params = {}
if cid:
params["arg"] = cid
resp = requests.post(url, params=params, timeout=30)
resp.raise_for_status()
return resp.json()
# ---------- DAG 操作 ----------
def dag_put(self, obj: Any, store_codec: str = "dag-cbor") -> Dict[str, Any]:
"""存储 DAG 节点"""
url = self._api_url("dag/put")
params = {"store-codec": store_codec, "input-codec": "dag-json"}
data = json.dumps(obj).encode("utf-8")
resp = requests.post(
url, params=params, data=data,
headers={"Content-Type": "application/json"},
timeout=30,
)
resp.raise_for_status()
return resp.json()
def dag_get(self, cid: str) -> Any:
"""获取 DAG 节点"""
url = self._api_url("dag/get")
params = {"arg": cid}
resp = requests.post(url, params=params, timeout=30)
resp.raise_for_status()
return resp.json()
# ---------- 工具方法 ----------
@staticmethod
def validate_cid(cid: str) -> bool:
"""
简单验证 CID 格式
v0: Qm 开头,base58,46 字符
v1: b 开头,Base32,59+ 字符
"""
if cid.startswith("Qm") and len(cid) == 46:
return True
if cid.startswith("b") and len(cid) >= 59:
return True
return False
@staticmethod
def compute_file_hash(file_path: str) -> str:
"""计算文件 SHA256(用于本地校验)"""
h = hashlib.sha256()
with open(file_path, "rb") as f:
for chunk in iter(lambda: f.read(65536), b""):
h.update(chunk)
return h.hexdigest()
# ---------- 使用示例 ----------
if __name__ == "__main__":
client = IPFSClient()
# 1. 添加 JSON 数据
result = client.add_json(
{"name": "MSG NFT #1", "description": "A test NFT on MSG Chain"},
filename="metadata.json",
)
cid = result["Hash"]
print(f"Added JSON, CID: {cid}")
# 2. 读取并验证
loaded = client.cat_json(cid)
print(f"Loaded data: {loaded}")
# 3. 固定确保可用
client.pin_add(cid)
# 4. 生成网关地址
gw_url = client.gateway_url_for(cid)
print(f"Gateway URL: {gw_url}")
# 5. 验证 CID
print(f"Valid CID: {IPFSClient.validate_cid(cid)}")
TypeScript 客户端
// ipfs-client.ts - IPFS HTTP API TypeScript 客户端
import axios, { AxiosInstance } from "axios";
import * as fs from "fs";
import * as path from "path";
import FormData from "form-data";
import { createHash } from "crypto";
interface IPFSAddResult {
Name: string;
Hash: string;
Size: string;
}
interface IPFSPinResult {
Pins: string[];
}
interface IPFSObject {
[key: string]: unknown;
}
export class IPFSClient {
private api: AxiosInstance;
private gatewayUrl: string;
constructor(
apiUrl: string = "http://127.0.0.1:5001",
gatewayUrl: string = "https://ipfs.io/ipfs"
) {
this.api = axios.create({
baseURL: `${apiUrl}/api/v0`,
timeout: 60000,
});
this.gatewayUrl = gatewayUrl;
}
// ---------- 添加文件 ----------
async addFile(
filePath: string,
pin: boolean = true,
cidVersion: number = 1
): Promise<IPFSAddResult> {
const form = new FormData();
form.append("file", fs.createReadStream(filePath));
const { data } = await this.api.post<IPFSAddResult>("/add", form, {
params: { pin, "cid-version": cidVersion },
headers: form.getHeaders(),
maxContentLength: Infinity,
maxBodyLength: Infinity,
});
return data;
}
async addBytes(
data: Buffer,
filename: string = "data.bin",
pin: boolean = true,
cidVersion: number = 1
): Promise<IPFSAddResult> {
const form = new FormData();
form.append("file", data, filename);
const { data: result } = await this.api.post<IPFSAddResult>("/add", form, {
params: { pin, "cid-version": cidVersion },
headers: form.getHeaders(),
});
return result;
}
async addJson(
obj: unknown,
filename: string = "data.json",
pin: boolean = true,
cidVersion: number = 1
): Promise<IPFSAddResult> {
const jsonStr = JSON.stringify(obj);
return this.addBytes(Buffer.from(jsonStr, "utf-8"), filename, pin, cidVersion);
}
async addDirectory(
dirPath: string,
pin: boolean = true,
cidVersion: number = 1
): Promise<IPFSAddResult[]> {
const form = new FormData();
const addFiles = (dir: string, base: string) => {
for (const entry of fs.readdirSync(dir, { withFileTypes: true })) {
const fullPath = path.join(dir, entry.name);
const relativePath = path.join(base, entry.name);
if (entry.isFile()) {
form.append("file", fs.createReadStream(fullPath), {
filepath: relativePath,
});
} else if (entry.isDirectory()) {
addFiles(fullPath, relativePath);
}
}
};
addFiles(dirPath, "");
const { data: raw } = await this.api.post("/add", form, {
params: { pin, "cid-version": cidVersion, recursive: true },
headers: form.getHeaders(),
maxContentLength: Infinity,
maxBodyLength: Infinity,
});
const lines = raw.trim().split("\n");
return lines.map((line: string) => JSON.parse(line));
}
// ---------- 读取文件 ----------
async cat(cid: string): Promise<Buffer> {
const { data } = await this.api.post<ArrayBuffer>("/cat", null, {
params: { arg: cid },
responseType: "arraybuffer",
});
return Buffer.from(data);
}
async catJson<T = IPFSObject>(cid: string): Promise<T> {
const buf = await this.cat(cid);
return JSON.parse(buf.toString("utf-8")) as T;
}
// ---------- Pinning ----------
async pinAdd(cid: string): Promise<IPFSPinResult> {
const { data } = await this.api.post<IPFSPinResult>("/pin/add", null, {
params: { arg: cid },
});
return data;
}
async pinRm(cid: string): Promise<IPFSPinResult> {
const { data } = await this.api.post<IPFSPinResult>("/pin/rm", null, {
params: { arg: cid },
});
return data;
}
async pinLs(cid?: string): Promise<Record<string, unknown>> {
const params: Record<string, string> = {};
if (cid) params.arg = cid;
const { data } = await this.api.post("/pin/ls", null, { params });
return data;
}
// ---------- DAG ----------
async dagPut(obj: unknown): Promise<{ Cid: string }> {
const { data } = await this.api.post(
"/dag/put",
JSON.stringify(obj),
{
params: { "store-codec": "dag-cbor", "input-codec": "dag-json" },
headers: { "Content-Type": "application/json" },
}
);
return data;
}
async dagGet<T = IPFSObject>(cid: string): Promise<T> {
const { data } = await this.api.post<T>("/dag/get", null, {
params: { arg: cid },
});
return data;
}
// ---------- 网关 ----------
gatewayUrlFor(cid: string): string {
return `${this.gatewayUrl}/${cid}`;
}
async fetchViaGateway(cid: string, gateway?: string): Promise<Buffer> {
const base = gateway ?? this.gatewayUrl;
const { data } = await axios.get<ArrayBuffer>(`${base}/${cid}`, {
responseType: "arraybuffer",
});
return Buffer.from(data);
}
// ---------- 工具 ----------
static validateCID(cid: string): boolean {
if (cid.startsWith("Qm") && cid.length === 46) return true;
if (cid.startsWith("b") && cid.length >= 59) return true;
return false;
}
static computeFileHash(filePath: string): string {
const hash = createHash("sha256");
const data = fs.readFileSync(filePath);
hash.update(data);
return hash.digest("hex");
}
}
// ---------- 使用示例 ----------
async function main() {
const client = new IPFSClient();
// 1. 添加 JSON
const result = await client.addJson({
name: "MSG NFT #1",
description: "A test NFT on MSG Chain",
});
const cid = result.Hash;
console.log(`Added JSON, CID: ${cid}`);
// 2. 读取
const loaded = await client.catJson(cid);
console.log("Loaded:", loaded);
// 3. 固定
await client.pinAdd(cid);
// 4. 网关
console.log(`Gateway URL: ${client.gatewayUrlFor(cid)}`);
// 5. 验证
console.log(`Valid CID: ${IPFSClient.validateCID(cid)}`);
}
// main().catch(console.error);
2.3 Pinning 服务
运行自己的 IPFS 节点需要持续维护。生产环境中推荐使用 pinning 服务。
Pinata
"""
pinata_client.py - Pinata IPFS Pinning 服务客户端
"""
import json
from typing import Optional, Dict, Any
import requests
class PinataClient:
"""Pinata pinning service 封装"""
def __init__(self, api_key: str, secret_key: str, jwt: Optional[str] = None):
self.base_url = "https://api.pinata.cloud"
self.api_key = api_key
self.secret_key = secret_key
self.jwt = jwt
self.session = requests.Session()
if jwt:
self.session.headers["Authorization"] = f"Bearer {jwt}"
else:
self.session.headers["pinata_api_key"] = api_key
self.session.headers["pinata_secret_api_key"] = secret_key
def pin_file(
self,
file_path: str,
name: Optional[str] = None,
keyvalues: Optional[Dict[str, str]] = None,
) -> Dict[str, Any]:
"""
上传文件到 Pinata
Returns: { IpfsHash: "Qm...", PinSize: ..., Timestamp: "..." }
"""
url = f"{self.base_url}/pinning/pinFileToIPFS"
with open(file_path, "rb") as f:
files = {"file": (file_path.split("/")[-1], f)}
data = {}
if name:
data["pinataMetadata"] = json.dumps({"name": name})
if keyvalues:
meta = json.loads(data.get("pinataMetadata", "{}"))
meta["keyvalues"] = keyvalues
data["pinataMetadata"] = json.dumps(meta)
resp = self.session.post(url, files=files, data=data, timeout=120)
resp.raise_for_status()
return resp.json()
def pin_json(
self,
obj: Any,
name: Optional[str] = None,
keyvalues: Optional[Dict[str, str]] = None,
) -> Dict[str, Any]:
"""上传 JSON 到 Pinata"""
url = f"{self.base_url}/pinning/pinJSONToIPFS"
payload = {
"pinataContent": obj,
}
if name or keyvalues:
meta = {}
if name:
meta["name"] = name
if keyvalues:
meta["keyvalues"] = keyvalues
payload["pinataMetadata"] = meta
resp = self.session.post(url, json=payload, timeout=30)
resp.raise_for_status()
return resp.json()
def unpin(self, cid: str) -> Dict[str, Any]:
"""取消固定"""
url = f"{self.base_url}/pinning/unpin/{cid}"
resp = self.session.delete(url, timeout=30)
resp.raise_for_status()
return resp.json()
def list_pins(
self,
status: str = "pinned",
page_limit: int = 10,
) -> Dict[str, Any]:
"""列出固定内容"""
url = f"{self.base_url}/data/pinList"
params = {"status": status, "pageLimit": page_limit}
resp = self.session.get(url, params=params, timeout=30)
resp.raise_for_status()
return resp.json()
# 使用示例
if __name__ == "__main__":
import os
client = PinataClient(
api_key=os.environ["PINATA_API_KEY"],
secret_key=os.environ["PINATA_SECRET_KEY"],
)
metadata = {
"name": "MSG NFT #42",
"description": "A beautiful NFT on MSG Chain",
"image": "ipfs://QmImageCID...",
"attributes": [
{"trait_type": "Background", "value": "Sunset"},
{"trait_type": "Rarity", "value": "Legendary"},
],
}
result = client.pin_json(metadata, name="MSG NFT #42 Metadata")
cid = result["IpfsHash"]
print(f"Pinned to IPFS: {cid}")
print(f"Gateway URL: https://gateway.pinata.cloud/ipfs/{cid}")
// pinata-client.ts - Pinata TypeScript 客户端
import axios, { AxiosInstance } from "axios";
import FormData from "form-data";
import * as fs from "fs";
interface PinataResult {
IpfsHash: string;
PinSize: number;
Timestamp: string;
}
interface PinataPinListResult {
count: number;
rows: Array<{
id: string;
ipfs_pin_hash: string;
size: number;
date_pinned: string;
metadata: Record<string, unknown>;
}>;
}
export class PinataClient {
private api: AxiosInstance;
constructor(
apiKey: string,
secretKey: string,
jwt?: string
) {
this.api = axios.create({
baseURL: "https://api.pinata.cloud",
timeout: 120000,
});
if (jwt) {
this.api.defaults.headers.Authorization = `Bearer ${jwt}`;
} else {
this.api.defaults.headers["pinata_api_key"] = apiKey;
this.api.defaults.headers["pinata_secret_api_key"] = secretKey;
}
}
async pinFile(
filePath: string,
name?: string,
keyvalues?: Record<string, string>
): Promise<PinataResult> {
const form = new FormData();
form.append("file", fs.createReadStream(filePath));
if (name || keyvalues) {
const meta: Record<string, unknown> = {};
if (name) meta.name = name;
if (keyvalues) meta.keyvalues = keyvalues;
form.append("pinataMetadata", JSON.stringify(meta));
}
const { data } = await this.api.post<PinataResult>(
"/pinning/pinFileToIPFS",
form,
{ headers: form.getHeaders() }
);
return data;
}
async pinJson(
obj: unknown,
name?: string,
keyvalues?: Record<string, string>
): Promise<PinataResult> {
const payload: Record<string, unknown> = {
pinataContent: obj,
pinataMetadata: {},
};
if (name) payload.pinataMetadata = { ...(payload.pinataMetadata as object), name };
if (keyvalues) payload.pinataMetadata = { ...(payload.pinataMetadata as object), keyvalues };
const { data } = await this.api.post<PinataResult>(
"/pinning/pinJSONToIPFS",
payload
);
return data;
}
async unpin(cid: string): Promise<void> {
await this.api.delete(`/pinning/unpin/${cid}`);
}
async listPins(status: string = "pinned", pageLimit: number = 10): Promise<PinataPinListResult> {
const { data } = await this.api.get<PinataPinListResult>("/data/pinList", {
params: { status, pageLimit },
});
return data;
}
}
Web3.Storage
"""
web3storage_client.py - Web3.Storage 客户端
"""
import json
from typing import Any
import requests
class Web3StorageClient:
"""Web3.Storage (基于 Filecoin) 的 Python 封装"""
def __init__(self, api_token: str):
self.base_url = "https://api.web3.storage"
self.session = requests.Session()
self.session.headers["Authorization"] = f"Bearer {api_token}"
def upload_file(self, file_path: str) -> dict:
"""
上传文件,返回包含 CID 的结果
"""
url = f"{self.base_url}/upload"
with open(file_path, "rb") as f:
files = {"file": (file_path.split("/")[-1], f)}
resp = self.session.post(url, files=files, timeout=300)
resp.raise_for_status()
return resp.json()
def upload_json(self, obj: Any, filename: str = "data.json") -> dict:
"""上传 JSON 对象"""
url = f"{self.base_url}/upload"
data = json.dumps(obj).encode("utf-8")
files = {"file": (filename, data)}
resp = self.session.post(url, files=files, timeout=60)
resp.raise_for_status()
return resp.json()
def list_uploads(self) -> list:
"""列出所有上传"""
url = f"{self.base_url}/user/uploads"
resp = self.session.get(url, timeout=30)
resp.raise_for_status()
return resp.json()
def status(self, cid: str) -> dict:
"""查询上传状态"""
url = f"{self.base_url}/status/{cid}"
resp = self.session.get(url, timeout=30)
resp.raise_for_status()
return resp.json()
2.4 内容寻址详解
CID 的结构
CID v0 (Qm...):
- 固定长度 46 字符
- Base58 编码
- 多哈希前缀: 0x12 (SHA256), 0x20 (32 字节)
- 示例: QmZ4tDuMPZqG3qgCnfnqGvkKzBiV5z7MNM9q8Jw8jYzB6c
CID v1 (bafy...):
- 可变长度,最小 59 字符
- Base32 编码(小写,URL 安全)
- 明确指定编解码器(dag-pb, dag-cbor, raw 等)
- 示例: bafybeigdyrzt5sfp7udm7hu76uh7y26nf3efuylqabf3oclgtqy55fbzdi
CID 版本选择
"""
cid_utils.py - CID 工具函数
"""
import hashlib
from typing import Optional
def cid_to_gateway_url(
cid: str,
gateway: str = "https://ipfs.io",
) -> str:
"""CID 转网关 URL"""
gateway = gateway.rstrip("/")
return f"{gateway}/ipfs/{cid}"
def resolve_ipfs_uri(uri: str) -> Optional[str]:
"""解析 ipfs:// URI 到 HTTP 网关 URL"""
if uri.startswith("ipfs://"):
cid = uri[7:]
return cid_to_gateway_url(cid)
return None
def resolve_ar_uri(uri: str) -> Optional[str]:
"""解析 ar:// URI 到 Arweave 网关 URL"""
if uri.startswith("ar://"):
txid = uri[5:]
return f"https://arweave.net/{txid}"
return None
def is_valid_cid(cid: str) -> bool:
"""检查 CID 是否为有效格式"""
import re
# v0: Qm 开头,46 字符,Base58
if re.match(r"^Qm[1-9A-HJ-NP-Za-km-z]{44}$", cid):
return True
# v1: b 开头 + Base32
if re.match(r"^b[1-9a-z]{58,}$", cid):
return True
return False
3. Arweave 集成基础
3.1 Arweave 钱包创建
使用 ArConnect 或 Arweave.app(推荐)
对于生产环境,推荐使用 ArConnect 浏览器插件或 arweave.app 网页钱包。
程序化创建钱包:
"""
arweave_wallet.py - Arweave 钱包操作
"""
import json
import os
from typing import Optional
class ArweaveWallet:
"""
Arweave 钱包管理
Arweave 钱包使用 RSA-PSS 密钥对,格式为 JWK (JSON Web Key)
"""
def __init__(self, jwk_path: Optional[str] = None):
"""
初始化钱包
如果 jwk_path 提供则加载已有钱包
否则生成新钱包
"""
if jwk_path and os.path.exists(jwk_path):
self.load_from_file(jwk_path)
else:
self.generate()
self.address = self._derive_address()
def generate(self):
"""
生成新的 Arweave 钱包(JWK 格式)
"""
from cryptography.hazmat.primitives.asymmetric import rsa
import base64
private_key = rsa.generate_private_key(
public_exponent=65537,
key_size=4096,
)
public_key = private_key.public_key()
pub_numbers = public_key.public_numbers()
pub_bytes = pub_numbers.n.to_bytes((pub_numbers.n.bit_length() + 7) // 8, "big")
self.jwk = {
"kty": "RSA",
"n": base64.urlsafe_b64encode(pub_bytes).decode().rstrip("="),
"e": "AQAB",
}
def load_from_file(self, path: str):
"""从文件加载 JWK"""
with open(path) as f:
self.jwk = json.load(f)
def save_to_file(self, path: str):
"""保存 JWK 到文件"""
with open(path, "w") as f:
json.dump(self.jwk, f, indent=2)
os.chmod(path, 0o600)
def _derive_address(self) -> str:
"""
从公钥派生 Arweave 地址
地址 = Base64URL(SHA256(公钥))
"""
import hashlib
import base64
n_bytes = base64.urlsafe_b64decode(
self.jwk["n"] + "=="
)
d = hashlib.sha256(n_bytes).digest()
return base64.urlsafe_b64encode(d).decode().rstrip("=")
使用 Irys (原 Bundlr) 客户端
Irys 提供 Arweave 的即时确认上传,无需等待 Arweave 主网确认。
"""
irys_client.py - Irys (Bundlr) 客户端封装
"""
import json
from typing import Any, Optional
import requests
class IrysClient:
"""
Irys (原 Bundlr) 客户端
使用 Irys 上传到 Arweave,即时确认
"""
def __init__(
self,
private_key: str,
node_url: str = "https://node2.irys.xyz",
currency: str = "arweave",
):
self.node_url = node_url
self.currency = currency
self.private_key = private_key
self.session = requests.Session()
def upload_file(
self,
file_path: str,
tags: Optional[list[tuple[str, str]]] = None,
) -> str:
"""
上传文件到 Irys
Args:
file_path: 本地文件路径
tags: Arweave 标签列表 [(key, value), ...]
Returns:
Arweave 交易 ID
"""
with open(file_path, "rb") as f:
data = f.read()
return self.upload_bytes(data, tags=tags)
def upload_json(
self,
obj: Any,
tags: Optional[list[tuple[str, str]]] = None,
) -> str:
"""上传 JSON 对象到 Arweave"""
if tags is None:
tags = [("Content-Type", "application/json")]
data = json.dumps(obj).encode("utf-8")
return self.upload_bytes(data, tags=tags)
def upload_bytes(
self,
data: bytes,
tags: Optional[list[tuple[str, str]]] = None,
) -> str:
"""上传 bytes 数据到 Arweave"""
if tags is None:
tags = []
# 确保 Content-Type 标签存在
has_content_type = any(t[0] == "Content-Type" for t in tags)
if not has_content_type:
tags = [("Content-Type", "application/octet-stream")] + tags
# 构造 Irys 上传请求
url = f"{self.node_url}/tx/{self.currency}"
payload = {
"data": data.hex(),
"tags": [{"name": k, "value": v} for k, v in tags],
}
resp = self.session.post(url, json=payload, timeout=120)
resp.raise_for_status()
result = resp.json()
return result.get("id", result.get("txid", ""))
def get_balance(self) -> float:
"""查询 Irys 账户余额"""
url = f"{self.node_url}/account/balance"
resp = self.session.get(url, timeout=30)
resp.raise_for_status()
return float(resp.json().get("balance", 0))
@staticmethod
def gateway_url_for(txid: str) -> str:
"""生成 Arweave 网关 URL"""
return f"https://arweave.net/{txid}"
# 使用原生 arweave-python-client
class ArweaveNativeClient:
"""
原生 Arweave 客户端(使用 arweave-python-client)
pip install arweave-python-client
"""
def __init__(self, wallet_path: str):
from arweave import Wallet, Transaction
with open(wallet_path) as f:
jwk = json.load(f)
self.wallet = Wallet(jwk)
def upload_data(self, data: bytes, tags: list) -> str:
"""
上传数据到 Arweave
Args:
data: 要上传的数据
tags: [(name, value), ...]
Returns:
交易 ID (txid)
"""
from arweave import Transaction
tx = Transaction(self.wallet, data=data, tags=tags)
tx.sign()
tx.send()
return tx.id
def upload_file(self, file_path: str, tags: list) -> str:
"""上传文件"""
with open(file_path, "rb") as f:
data = f.read()
return self.upload_data(data, tags)
def get_transaction(self, txid: str) -> dict:
"""获取交易详情"""
from arweave import Transaction
return Transaction.get(txid)
def get_data(self, txid: str) -> bytes:
"""获取交易数据"""
from arweave import Transaction
return Transaction.get_data(txid)
// arweave-client.ts - Arweave TypeScript 客户端
import Arweave from "arweave";
import Irys from "@irys/sdk";
import * as fs from "fs";
const ARWEAVE_CONFIG = {
host: "arweave.net",
port: 443,
protocol: "https" as const,
};
// ========== 原生 Arweave 客户端 ==========
export class ArweaveClient {
private arweave: Arweave;
private wallet: JsonWebKey | null = null;
constructor(jwkPath?: string) {
this.arweave = Arweave.init(ARWEAVE_CONFIG);
}
async loadWallet(jwkPath: string): Promise<void> {
const jwk = JSON.parse(fs.readFileSync(jwkPath, "utf-8"));
this.wallet = jwk;
}
async createWallet(): Promise<JsonWebKey> {
const jwk = await this.arweave.wallets.generate();
this.wallet = jwk;
return jwk;
}
async getAddress(): Promise<string> {
if (!this.wallet) throw new Error("Wallet not loaded");
return await this.arweave.wallets.jwkToAddress(this.wallet);
}
async getBalance(address?: string): Promise<string> {
const addr = address ?? (await this.getAddress());
const winston = await this.arweave.wallets.getBalance(addr);
return this.arweave.ar.winstonToAr(winston);
}
async uploadData(
data: Buffer | string,
tags: Array<{ name: string; value: string }> = [],
contentType: string = "application/octet-stream"
): Promise<string> {
if (!this.wallet) throw new Error("Wallet not loaded");
if (!tags.find((t) => t.name === "Content-Type")) {
tags.push({ name: "Content-Type", value: contentType });
}
const tx = await this.arweave.createTransaction(
{ data: typeof data === "string" ? data : data.toString("base64") },
this.wallet
);
tags.forEach((tag) => tx.addTag(tag.name, tag.value));
await this.arweave.transactions.sign(tx, this.wallet);
const response = await this.arweave.transactions.post(tx);
if (response.status !== 200) {
throw new Error(`Upload failed: ${response.statusText}`);
}
return tx.id;
}
async uploadFile(
filePath: string,
tags: Array<{ name: string; value: string }> = []
): Promise<string> {
const data = fs.readFileSync(filePath);
const ext = filePath.split(".").pop()?.toLowerCase();
const mimeMap: Record<string, string> = {
json: "application/json",
png: "image/png",
jpg: "image/jpeg",
jpeg: "image/jpeg",
gif: "image/gif",
svg: "image/svg+xml",
pdf: "application/pdf",
zip: "application/zip",
};
const contentType = mimeMap[ext ?? ""] ?? "application/octet-stream";
return this.uploadData(data, tags, contentType);
}
async uploadJson(
obj: unknown,
tags: Array<{ name: string; value: string }> = []
): Promise<string> {
const data = JSON.stringify(obj);
return this.uploadData(
data,
[{ name: "Content-Type", value: "application/json" }, ...tags]
);
}
async getData(txid: string): Promise<Buffer> {
const data = await this.arweave.transactions.getData(txid, {
decode: true,
string: false,
});
return Buffer.from(data as ArrayBuffer);
}
async getTransaction(txid: string): Promise<Record<string, unknown>> {
const tx = await this.arweave.transactions.get(txid);
const tags: Record<string, string> = {};
tx.get("tags").forEach(
(tag: { get: (name: string) => string }) => {
const key = tag.get("name");
const value = tag.get("value");
tags[key] = value;
}
);
return {
id: txid,
owner: tx.owner,
tags,
data_size: tx.data_size,
data_root: tx.data_root,
reward: tx.reward,
};
}
gatewayUrlFor(txid: string): string {
return `https://arweave.net/${txid}`;
}
}
// ========== Irys (原 Bundlr) 客户端 ==========
export class IrysClient {
private irys: Irys;
constructor(
privateKey: string,
nodeUrl: string = "https://node2.irys.xyz",
currency: string = "arweave"
) {
this.irys = new Irys({
url: nodeUrl,
token: currency,
key: privateKey,
});
}
async uploadFile(
filePath: string,
tags: Array<{ name: string; value: string }> = []
): Promise<{ id: string; size: number }> {
const result = await this.irys.uploadFile(filePath, { tags });
return { id: result.id, size: result.size };
}
async uploadJson(
obj: unknown,
tags: Array<{ name: string; value: string }> = []
): Promise<{ id: string; size: number }> {
const result = await this.irys.upload(JSON.stringify(obj), {
tags: [
{ name: "Content-Type", value: "application/json" },
...tags,
],
});
return { id: result.id, size: result.size };
}
async uploadData(
data: Buffer,
tags: Array<{ name: string; value: string }> = []
): Promise<{ id: string; size: number }> {
const result = await this.irys.upload(data, { tags });
return { id: result.id, size: result.size };
}
async getBalance(): Promise<string> {
const balance = await this.irys.getLoadedBalance();
return balance.toString();
}
async fund(amount: string): Promise<void> {
await this.irys.fund(parseFloat(amount));
}
gatewayUrlFor(txid: string): string {
return `https://arweave.net/${txid}`;
}
// 通过 GraphQL 查询
async queryByTags(
tags: Array<{ name: string; values: string[] }>,
limit: number = 10
): Promise<Array<{ id: string; tags: Array<{ name: string; value: string }> }>> {
const query = `
query ($tags: [TagFilter!]!, $limit: Int!) {
transactions(tags: $tags, first: $limit) {
edges {
node {
id
tags { name value }
}
}
}
}
`;
const result = await this.irys.query(query, { tags, limit });
return result;
}
}
// ---------- 使用示例 ----------
async function main() {
const irys = new IrysClient("your-private-key-hex");
const metadata = {
name: "MSG NFT #1",
description: "Permanently stored on Arweave",
image: "ar://image-txid",
attributes: [{ trait_type: "Color", value: "Blue" }],
};
const result = await irys.uploadJson(metadata, [
{ name: "App-Name", value: "MSG-NFT" },
{ name: "App-Version", value: "1.0.0" },
{ name: "Network", value: "msg-chain-1" },
]);
console.log(`Uploaded! Arweave TX: ${result.id}`);
console.log(`View at: https://arweave.net/${result.id}`);
}
3.2 Arweave 标签系统
Arweave 使用键值对标签进行数据索引。正确使用标签是查询数据的关键。
"""
arweave_tags.py - Arweave 标签最佳实践
"""
# 推荐标签用于 NFT 元数据
NFT_TAGS = [
("Content-Type", "application/json"),
("App-Name", "MSG-NFT"),
("Network", "msg-chain-1"),
("Type", "nft-metadata"),
("Protocol", "cw721"),
("Contract", "msg1xyz..."),
("Token-ID", "1"),
("Schema-Version", "1.0.0"),
]
# 推荐标签用于数据市场文件
DATASET_TAGS = [
("Content-Type", "application/zip"),
("App-Name", "MSG-Data-Market"),
("Network", "msg-chain-1"),
("Type", "dataset"),
("Asset-ID", "asset-001"),
("Data-Format", "csv"),
("Compression", "zip"),
("Original-Hash", "sha256:abc123..."),
("License", "CC-BY-4.0"),
]
def create_nft_tags(
contract_address: str,
token_id: str,
network: str = "msg-chain-1",
) -> list[tuple[str, str]]:
"""生成 NFT 标签"""
return [
("Content-Type", "application/json"),
("App-Name", "MSG-NFT"),
("Network", network),
("Type", "nft-metadata"),
("Protocol", "cw721"),
("Contract", contract_address),
("Token-ID", str(token_id)),
("Schema-Version", "1.0.0"),
]
def create_dataset_tags(
asset_id: str,
data_format: str = "csv",
original_hash: str = "",
) -> list[tuple[str, str]]:
"""生成数据资产标签"""
tags = [
("Content-Type", f"application/{data_format}"),
("App-Name", "MSG-Data-Market"),
("Network", "msg-chain-1"),
("Type", "dataset"),
("Asset-ID", asset_id),
]
if original_hash:
tags.append(("Original-Hash", original_hash))
return tags
// arweave-tags.ts
export interface ArweaveTag {
name: string;
value: string;
}
export const NFT_TAGS: ArweaveTag[] = [
{ name: "Content-Type", value: "application/json" },
{ name: "App-Name", value: "MSG-NFT" },
{ name: "Network", value: "msg-chain-1" },
{ name: "Type", value: "nft-metadata" },
{ name: "Protocol", value: "cw721" },
{ name: "Schema-Version", value: "1.0.0" },
];
export function createNFTTags(
contractAddress: string,
tokenId: string,
extraTags?: ArweaveTag[]
): ArweaveTag[] {
const tags: ArweaveTag[] = [
...NFT_TAGS,
{ name: "Contract", value: contractAddress },
{ name: "Token-ID", value: tokenId },
];
if (extraTags) tags.push(...extraTags);
return tags;
}
3.3 通过 GraphQL 查询 Arweave
"""
arweave_graphql.py - Arweave GraphQL 查询
"""
import json
from typing import Optional, Any
import requests
class ArweaveGraphQL:
"""Arweave GraphQL 查询接口"""
def __init__(self, endpoint: str = "https://arweave.net/graphql"):
self.endpoint = endpoint
def query(self, query_str: str, variables: Optional[dict] = None) -> dict:
"""执行 GraphQL 查询"""
payload = {"query": query_str}
if variables:
payload["variables"] = variables
resp = requests.post(self.endpoint, json=payload, timeout=30)
resp.raise_for_status()
return resp.json()
def find_by_tags(
self,
tags: list[tuple[str, str]],
limit: int = 10,
) -> list[dict]:
"""
通过标签查找交易
Args:
tags: [(name, value), ...]
limit: 返回数量
Returns:
交易列表 [{ id, tags, ... }]
"""
tag_filters = [
f'{{ name: "{name}", values: ["{value}"] }}'
for name, value in tags
]
query = f"""
{{
transactions(
tags: [{', '.join(tag_filters)}]
first: {limit}
) {{
edges {{
node {{
id
tags {{
name
value
}}
data {{ size }}
block {{
height
timestamp
}}
}}
}}
}}
}}
"""
result = self.query(query)
edges = result.get("data", {}).get("transactions", {}).get("edges", [])
return [edge["node"] for edge in edges]
def find_by_app(
self,
app_name: str,
type_filter: Optional[str] = None,
limit: int = 20,
) -> list[dict]:
"""通过 App-Name 查找交易"""
tags = [("App-Name", app_name)]
if type_filter:
tags.append(("Type", type_filter))
return self.find_by_tags(tags, limit)
def get_transaction(self, txid: str) -> dict:
"""获取单个交易详情"""
query = f"""
{{
transaction(id: "{txid}") {{
id
owner {{ address }}
tags {{ name value }}
data {{ size type }}
block {{ height timestamp }}
fee {{ winston ar }}
}}
}}
"""
result = self.query(query)
return result.get("data", {}).get("transaction", {})
# 使用示例
if __name__ == "__main__":
gql = ArweaveGraphQL()
# 查询 MSG Chain 上的所有 NFT 元数据
txs = gql.find_by_app("MSG-NFT", type_filter="nft-metadata", limit=5)
for tx in txs:
print(f"TX: {tx['id']}")
for tag in tx["tags"]:
print(f" {tag['name']}: {tag['value']}")
# 查询特定合约的元数据
contract_tags = [
("App-Name", "MSG-NFT"),
("Contract", "msg1contractaddress..."),
]
results = gql.find_by_tags(contract_tags)
print(f"Found {len(results)} transactions")
// arweave-graphql.ts
import axios from "axios";
interface GQLTag {
name: string;
value: string;
}
interface GQLTransaction {
id: string;
tags: GQLTag[];
data?: { size: string };
block?: { height: number; timestamp: number };
}
interface GQLResponse {
data: {
transactions: {
edges: Array<{ node: GQLTransaction }>;
};
transaction?: GQLTransaction;
};
}
export class ArweaveGraphQL {
private endpoint: string;
constructor(endpoint: string = "https://arweave.net/graphql") {
this.endpoint = endpoint;
}
async query(query: string, variables?: Record<string, unknown>): Promise<GQLResponse> {
const { data } = await axios.post<GQLResponse>(
this.endpoint,
{ query, variables },
{ timeout: 30000 }
);
return data;
}
async findByTags(
tags: Array<{ name: string; value: string }>,
limit: number = 10
): Promise<GQLTransaction[]> {
const tagFilters = tags
.map((t) => `{ name: "${t.name}", values: ["${t.value}"] }`)
.join(", ");
const query = `
{
transactions(tags: [${tagFilters}], first: ${limit}) {
edges {
node {
id
tags { name value }
data { size }
block { height timestamp }
}
}
}
}
`;
const result = await this.query(query);
return result.data.transactions.edges.map((e) => e.node);
}
async findByApp(
appName: string,
type?: string,
limit: number = 20
): Promise<GQLTransaction[]> {
const tags = [{ name: "App-Name", value: appName }];
if (type) tags.push({ name: "Type", value: type });
return this.findByTags(tags, limit);
}
async getTransaction(txid: string): Promise<GQLTransaction | null> {
const query = `
{
transaction(id: "${txid}") {
id
owner { address }
tags { name value }
data { size }
block { height timestamp }
}
}
`;
const result = await this.query(query);
return result.data.transaction ?? null;
}
}
4. 合约集成(IPFS)
4.1 在 CosmWasm 合约中存储 IPFS CID
// src/contract.rs - IPFS 集成的 NFT 合约
use cosmwasm_std::{
entry_point, Binary, Deps, DepsMut, Env, MessageInfo, Response, StdError, StdResult,
Timestamp, to_binary,
};
use cw_storage_plus::{Item, Map};
use serde::{Deserialize, Serialize};
use schemars::JsonSchema;
// ========== 状态定义 ==========
/// NFT 元数据结构
/// 链上仅存 CID,实际数据在 IPFS
#[derive(Serialize, Deserialize, Clone, Debug, PartialEq, JsonSchema)]
pub struct NFTMetadata {
/// IPFS CID 指向 metadata.json
pub cid: String,
/// 元数据最后更新时间
pub updated_at: Timestamp,
/// 可选的 IPNS 名称(可变指针)
pub ipns_name: Option<String>,
}
/// 数据资产结构(用于数据市场)
#[derive(Serialize, Deserialize, Clone, Debug, PartialEq, JsonSchema)]
pub struct DataAsset {
/// 数据集的 IPFS CID 或 Arweave txid
pub storage_cid: String,
/// 存储类型: "ipfs" | "arweave"
pub storage_type: StorageType,
/// 文件原始哈希(SHA256),用于完整性验证
pub content_hash: String,
/// 数据价格(在 MSG Chain 上)
pub price: Uint128,
/// 数据所有者
pub owner: Addr,
/// 是否已加密
pub encrypted: bool,
/// 加密密钥的 CID(用所有者公钥加密后存 IPFS)
pub encrypted_key_cid: Option<String>,
}
#[derive(Serialize, Deserialize, Clone, Debug, PartialEq, JsonSchema)]
pub enum StorageType {
Ipfs,
Arweave,
}
// 存储
pub const TOKEN_METADATA: Map<&str, NFTMetadata> = Map::new("tokens");
pub const DATA_ASSETS: Map<&str, DataAsset> = Map::new("assets");
pub const CONFIG: Item<Config> = Item::new("config");
#[derive(Serialize, Deserialize, Clone, Debug, PartialEq, JsonSchema)]
pub struct Config {
pub owner: Addr,
pub name: String,
pub symbol: String,
pub gateway_fallbacks: Vec<String>,
}
// ========== 消息定义 ==========
#[derive(Serialize, Deserialize, Clone, Debug, PartialEq, JsonSchema)]
pub struct InstantiateMsg {
pub name: String,
pub symbol: String,
pub owner: Option<String>,
pub gateway_fallbacks: Option<Vec<String>>,
}
#[derive(Serialize, Deserialize, Clone, Debug, PartialEq, JsonSchema)]
#[serde(rename_all = "snake_case")]
pub enum ExecuteMsg {
/// 存储 NFT 元数据 CID
StoreMetadata {
token_id: String,
cid: String,
ipns_name: Option<String>,
},
/// 批量存储元数据
BatchStoreMetadata {
token_ids: Vec<String>,
cids: Vec<String>,
},
/// 更新已有元数据
UpdateMetadata {
token_id: String,
new_cid: String,
ipns_name: Option<String>,
},
/// 注册数据资产
RegisterDataAsset {
asset_id: String,
cid: String,
storage_type: StorageType,
content_hash: String,
price: Uint128,
},
/// 验证数据完整性
VerifyDataIntegrity {
asset_id: String,
claimed_hash: String,
},
}
#[derive(Serialize, Deserialize, Clone, Debug, PartialEq, JsonSchema)]
#[serde(rename_all = "snake_case")]
pub enum QueryMsg {
/// 获取 token 元数据
GetMetadata { token_id: String },
/// 列出所有 tokens
AllTokens { start_after: Option<String>, limit: Option<u32> },
/// 获取数据资产
GetDataAsset { asset_id: String },
/// 验证数据
VerifyData { asset_id: String },
/// 获取配置
GetConfig {},
}
// ========== 消息处理 ==========
#[entry_point]
pub fn instantiate(
deps: DepsMut,
_env: Env,
info: MessageInfo,
msg: InstantiateMsg,
) -> StdResult<Response> {
let config = Config {
owner: deps.api.addr_validate(&msg.owner.unwrap_or(info.sender.to_string()))?,
name: msg.name,
symbol: msg.symbol,
gateway_fallbacks: msg.gateway_fallbacks.unwrap_or_else(|| {
vec![
"https://ipfs.io/ipfs".to_string(),
"https://gateway.pinata.cloud/ipfs".to_string(),
"https://cloudflare-ipfs.com/ipfs".to_string(),
]
}),
};
CONFIG.save(deps.storage, &config)?;
Ok(Response::new()
.add_attribute("action", "instantiate")
.add_attribute("name", &config.name))
}
#[entry_point]
pub fn execute(
deps: DepsMut,
env: Env,
info: MessageInfo,
msg: ExecuteMsg,
) -> StdResult<Response> {
match msg {
ExecuteMsg::StoreMetadata { token_id, cid, ipns_name } => {
execute_store_metadata(deps, env, info, token_id, cid, ipns_name)
}
ExecuteMsg::BatchStoreMetadata { token_ids, cids } => {
execute_batch_store_metadata(deps, env, info, token_ids, cids)
}
ExecuteMsg::UpdateMetadata { token_id, new_cid, ipns_name } => {
execute_update_metadata(deps, env, info, token_id, new_cid, ipns_name)
}
ExecuteMsg::RegisterDataAsset { asset_id, cid, storage_type, content_hash, price } => {
execute_register_asset(deps, env, info, asset_id, cid, storage_type, content_hash, price)
}
ExecuteMsg::VerifyDataIntegrity { asset_id, claimed_hash } => {
execute_verify_integrity(deps, env, info, asset_id, claimed_hash)
}
}
}
// ========== 存储元数据 ==========
/// 存储 NFT 元数据的 IPFS CID
pub fn execute_store_metadata(
deps: DepsMut,
env: Env,
_info: MessageInfo,
token_id: String,
cid: String,
ipns_name: Option<String>,
) -> StdResult<Response> {
validate_cid_format(&cid)?;
if TOKEN_METADATA.has(deps.storage, &token_id) {
return Err(StdError::generic_err(format!(
"Token ID '{}' already has metadata. Use UpdateMetadata to change.",
token_id
)));
}
let metadata = NFTMetadata {
cid: cid.clone(),
updated_at: env.block.time,
ipns_name,
};
TOKEN_METADATA.save(deps.storage, &token_id, &metadata)?;
Ok(Response::new()
.add_attribute("action", "store_metadata")
.add_attribute("token_id", &token_id)
.add_attribute("cid", &cid)
.add_attribute("timestamp", env.block.time.to_string()))
}
/// 批量存储元数据
pub fn execute_batch_store_metadata(
deps: DepsMut,
env: Env,
_info: MessageInfo,
token_ids: Vec<String>,
cids: Vec<String>,
) -> StdResult<Response> {
if token_ids.len() != cids.len() {
return Err(StdError::generic_err(
"token_ids and cids must have the same length",
));
}
if token_ids.is_empty() {
return Err(StdError::generic_err("Empty batch not allowed"));
}
if token_ids.len() > 100 {
return Err(StdError::generic_err("Batch too large (max 100)"));
}
for (token_id, cid) in token_ids.iter().zip(cids.iter()) {
validate_cid_format(cid)?;
if TOKEN_METADATA.has(deps.storage, token_id) {
return Err(StdError::generic_err(format!(
"Token ID '{}' already exists. Rollback entire batch.",
token_id
)));
}
}
for (token_id, cid) in token_ids.iter().zip(cids.iter()) {
let metadata = NFTMetadata {
cid: cid.clone(),
updated_at: env.block.time,
ipns_name: None,
};
TOKEN_METADATA.save(deps.storage, token_id, &metadata)?;
}
Ok(Response::new()
.add_attribute("action", "batch_store_metadata")
.add_attribute("count", token_ids.len().to_string()))
}
/// 更新已有元数据(仅所有者/合约 Owner)
pub fn execute_update_metadata(
deps: DepsMut,
env: Env,
info: MessageInfo,
token_id: String,
new_cid: String,
ipns_name: Option<String>,
) -> StdResult<Response> {
validate_cid_format(&new_cid)?;
let config = CONFIG.load(deps.storage)?;
if info.sender != config.owner {
return Err(StdError::generic_err("Unauthorized: only contract owner can update metadata"));
}
if !TOKEN_METADATA.has(deps.storage, &token_id) {
return Err(StdError::generic_err(format!(
"Token '{}' not found. Use StoreMetadata first.",
token_id
)));
}
let metadata = NFTMetadata {
cid: new_cid.clone(),
updated_at: env.block.time,
ipns_name,
};
TOKEN_METADATA.save(deps.storage, &token_id, &metadata)?;
Ok(Response::new()
.add_attribute("action", "update_metadata")
.add_attribute("token_id", &token_id)
.add_attribute("new_cid", &new_cid)
.add_attribute("timestamp", env.block.time.to_string()))
}
// ========== 数据市场资产 ==========
pub fn execute_register_asset(
deps: DepsMut,
env: Env,
info: MessageInfo,
asset_id: String,
cid: String,
storage_type: StorageType,
content_hash: String,
price: Uint128,
) -> StdResult<Response> {
if content_hash.len() != 64 {
return Err(StdError::generic_err(
"Invalid content_hash: must be 64 hex chars (SHA256)",
));
}
if DATA_ASSETS.has(deps.storage, &asset_id) {
return Err(StdError::generic_err(format!(
"Asset '{}' already registered",
asset_id
)));
}
let asset = DataAsset {
storage_cid: cid.clone(),
storage_type,
content_hash: content_hash.clone(),
price,
owner: info.sender.clone(),
encrypted: false,
encrypted_key_cid: None,
};
DATA_ASSETS.save(deps.storage, &asset_id, &asset)?;
Ok(Response::new()
.add_attribute("action", "register_asset")
.add_attribute("asset_id", &asset_id)
.add_attribute("cid", &cid)
.add_attribute("content_hash", &content_hash)
.add_attribute("price", &price.to_string())
.add_attribute("owner", &info.sender))
}
/// 验证数据完整性
pub fn execute_verify_integrity(
deps: DepsMut,
_env: Env,
_info: MessageInfo,
asset_id: String,
claimed_hash: String,
) -> StdResult<Response> {
let asset = DATA_ASSETS.load(deps.storage, &asset_id)?;
let matches = asset.content_hash == claimed_hash;
Ok(Response::new()
.add_attribute("action", "verify_integrity")
.add_attribute("asset_id", &asset_id)
.add_attribute("stored_hash", &asset.content_hash)
.add_attribute("claimed_hash", &claimed_hash)
.add_attribute("matches", matches.to_string()))
}
// ========== 查询 ==========
#[entry_point]
pub fn query(deps: Deps, _env: Env, msg: QueryMsg) -> StdResult<Binary> {
match msg {
QueryMsg::GetMetadata { token_id } => {
let meta = TOKEN_METADATA.load(deps.storage, &token_id)?;
to_binary(&meta)
}
QueryMsg::AllTokens { start_after, limit } => {
let limit = limit.unwrap_or(30).min(100);
let start = start_after.as_deref();
let tokens: StdResult<Vec<_>> = TOKEN_METADATA
.range(deps.storage, start, None, cosmwasm_std::Order::Ascending)
.take(limit as usize)
.map(|r| r.map(|(k, v)| (k, v)))
.collect();
to_binary(&tokens?)
}
QueryMsg::GetDataAsset { asset_id } => {
let asset = DATA_ASSETS.load(deps.storage, &asset_id)?;
to_binary(&asset)
}
QueryMsg::VerifyData { asset_id } => {
let asset = DATA_ASSETS.load(deps.storage, &asset_id)?;
to_binary(&asset)
}
QueryMsg::GetConfig {} => {
let config = CONFIG.load(deps.storage)?;
to_binary(&config)
}
}
}
// ========== 工具函数 ==========
/// 验证 CID 格式
fn validate_cid_format(cid: &str) -> StdResult<()> {
let is_v0 = cid.starts_with("Qm") && cid.len() == 46;
let is_v1 = cid.starts_with('b') && cid.len() >= 59
&& cid.chars().all(|c| c.is_ascii_lowercase() || c.is_ascii_digit());
if !is_v0 && !is_v1 {
return Err(StdError::generic_err(format!(
"Invalid CID format: '{}'. Must be v0 (Qm...) or v1 (bafy...)",
cid
)));
}
Ok(())
}
4.2 合约单元测试
// src/testing.rs
#[cfg(test)]
mod tests {
use super::*;
use cosmwasm_std::testing::{mock_dependencies, mock_env, mock_info};
use cosmwasm_std::{from_binary, Uint128};
const CONTRACT_NAME: &str = "MSG Storage Test";
const CONTRACT_SYMBOL: &str = "MST";
fn setup_contract() -> (DepsMut, Env) {
let mut deps = mock_dependencies();
let env = mock_env();
let info = mock_info("creator", &[]);
let msg = InstantiateMsg {
name: CONTRACT_NAME.to_string(),
symbol: CONTRACT_SYMBOL.to_string(),
owner: None,
gateway_fallbacks: None,
};
instantiate(deps.as_mut(), env.clone(), info, msg).unwrap();
(deps.as_mut(), env)
}
#[test]
fn test_instantiate() {
let (deps, _) = setup_contract();
let config = CONFIG.load(deps.as_ref().storage).unwrap();
assert_eq!(config.name, CONTRACT_NAME);
assert_eq!(config.symbol, CONTRACT_SYMBOL);
assert_eq!(config.gateway_fallbacks.len(), 3);
}
#[test]
fn test_store_metadata() {
let (mut deps, env) = setup_contract();
let info = mock_info("user", &[]);
let cid = "QmZ4tDuMPZqG3qgCnfnqGvkKzBiV5z7MNM9q8Jw8jYzB6c";
let resp = execute(
deps.as_mut(),
env.clone(),
info.clone(),
ExecuteMsg::StoreMetadata {
token_id: "1".to_string(),
cid: cid.to_string(),
ipns_name: None,
},
).unwrap();
assert_eq!(resp.attributes[0].value, "store_metadata");
assert_eq!(resp.attributes[1].value, "1");
assert_eq!(resp.attributes[2].value, cid);
let query_resp = query(
deps.as_ref(),
env.clone(),
QueryMsg::GetMetadata { token_id: "1".to_string() },
).unwrap();
let meta: NFTMetadata = from_binary(&query_resp).unwrap();
assert_eq!(meta.cid, cid);
}
#[test]
fn test_duplicate_metadata() {
let (mut deps, env) = setup_contract();
let info = mock_info("user", &[]);
let cid = "QmZ4tDuMPZqG3qgCnfnqGvkKzBiV5z7MNM9q8Jw8jYzB6c";
execute(
deps.as_mut(),
env.clone(),
info.clone(),
ExecuteMsg::StoreMetadata {
token_id: "1".to_string(),
cid: cid.to_string(),
ipns_name: None,
},
).unwrap();
let err = execute(
deps.as_mut(),
env.clone(),
info,
ExecuteMsg::StoreMetadata {
token_id: "1".to_string(),
cid: cid.to_string(),
ipns_name: None,
},
).unwrap_err();
assert!(err.to_string().contains("already has metadata"));
}
#[test]
fn test_invalid_cid() {
let (mut deps, env) = setup_contract();
let info = mock_info("user", &[]);
let err = execute(
deps.as_mut(),
env,
info,
ExecuteMsg::StoreMetadata {
token_id: "1".to_string(),
cid: "invalid-cid".to_string(),
ipns_name: None,
},
).unwrap_err();
assert!(err.to_string().contains("Invalid CID format"));
}
#[test]
fn test_update_metadata() {
let (mut deps, env) = setup_contract();
let creator_info = mock_info("creator", &[]);
let cid1 = "QmZ4tDuMPZqG3qgCnfnqGvkKzBiV5z7MNM9q8Jw8jYzB6c";
execute(
deps.as_mut(),
env.clone(),
creator_info.clone(),
ExecuteMsg::StoreMetadata {
token_id: "1".to_string(),
cid: cid1.to_string(),
ipns_name: None,
},
).unwrap();
let cid2 = "QmNewCID1234567890123456789012345678901234567890123";
execute(
deps.as_mut(),
env.clone(),
creator_info,
ExecuteMsg::UpdateMetadata {
token_id: "1".to_string(),
new_cid: cid2.to_string(),
ipns_name: None,
},
).unwrap();
let query_resp = query(
deps.as_ref(),
env,
QueryMsg::GetMetadata { token_id: "1".to_string() },
).unwrap();
let meta: NFTMetadata = from_binary(&query_resp).unwrap();
assert_eq!(meta.cid, cid2);
}
#[test]
fn test_unauthorized_update() {
let (mut deps, env) = setup_contract();
let user_info = mock_info("user", &[]);
let cid = "QmZ4tDuMPZqG3qgCnfnqGvkKzBiV5z7MNM9q8Jw8jYzB6c";
execute(
deps.as_mut(),
env.clone(),
user_info.clone(),
ExecuteMsg::StoreMetadata {
token_id: "1".to_string(),
cid: cid.to_string(),
ipns_name: None,
},
).unwrap();
let err = execute(
deps.as_mut(),
env,
user_info,
ExecuteMsg::UpdateMetadata {
token_id: "1".to_string(),
new_cid: "QmNewCID...".to_string(),
ipns_name: None,
},
).unwrap_err();
assert!(err.to_string().contains("Unauthorized"));
}
#[test]
fn test_register_data_asset() {
let (mut deps, env) = setup_contract();
let info = mock_info("owner", &[]);
let cid = "QmDataCID1234567890123456789012345678901234567890123";
let content_hash = "a".repeat(64);
execute(
deps.as_mut(),
env.clone(),
info.clone(),
ExecuteMsg::RegisterDataAsset {
asset_id: "dataset-001".to_string(),
cid: cid.to_string(),
storage_type: StorageType::Ipfs,
content_hash: content_hash.clone(),
price: Uint128::new(1000),
},
).unwrap();
let query_resp = query(
deps.as_ref(),
env,
QueryMsg::GetDataAsset { asset_id: "dataset-001".to_string() },
).unwrap();
let asset: DataAsset = from_binary(&query_resp).unwrap();
assert_eq!(asset.storage_cid, cid);
assert_eq!(asset.content_hash, content_hash);
}
#[test]
fn test_gateway_fallbacks() {
let mut deps = mock_dependencies();
let env = mock_env();
let info = mock_info("creator", &[]);
let custom_gateways = vec![
"https://ipfs.io/ipfs".to_string(),
"https://cloudflare-ipfs.com/ipfs".to_string(),
];
let msg = InstantiateMsg {
name: "Test".to_string(),
symbol: "TST".to_string(),
owner: None,
gateway_fallbacks: Some(custom_gateways.clone()),
};
instantiate(deps.as_mut(), env, info, msg).unwrap();
let config = CONFIG.load(deps.as_ref().storage).unwrap();
assert_eq!(config.gateway_fallbacks, custom_gateways);
}
}
4.3 网关备用列表
前端读取 IPFS 数据时应实现多个网关的自动切换:
// 合约中存储的默认网关列表
pub const DEFAULT_GATEWAYS: &[&str] = &[
"https://ipfs.io/ipfs",
"https://cloudflare-ipfs.com/ipfs",
"https://gateway.pinata.cloud/ipfs",
"https://ipfs.fleek.co/ipfs",
"https://dweb.link/ipfs",
"https://4everland.io/ipfs",
"https://w3s.link/ipfs",
];
/// 获取带 CID 的网关 URLs
pub fn get_gateway_urls(cid: &str) -> Vec<String> {
DEFAULT_GATEWAYS
.iter()
.map(|gw| format!("{}/{}", gw, cid))
.collect()
}
4.4 事件索引(Off-chain Indexer)
"""
indexer.py - 监听合约事件并索引 IPFS 元数据
"""
import json
import time
from typing import Optional
import requests
class MetadataIndexer:
"""
监听链上 StoreMetadata 事件,自动验证 IPFS 数据并缓存
"""
def __init__(
self,
rpc_url: str,
contract_address: str,
gateways: Optional[list[str]] = None,
):
self.rpc_url = rpc_url
self.contract_address = contract_address
self.gateways = gateways or [
"https://ipfs.io/ipfs",
"https://gateway.pinata.cloud/ipfs",
]
self.last_processed_height = 0
def fetch_events(self, start_height: int, end_height: int) -> list[dict]:
"""
从链上获取事件(Cosmos SDK 的 TxSearch)
实际使用应适配 MSG Chain 的 RPC 接口
"""
url = f"{self.rpc_url}/cosmos/tx/v1beta1/txs"
params = {
"events": [
f"action='store_metadata'",
f"contract='{self.contract_address}'",
],
"pagination.limit": "100",
}
resp = requests.get(url, params=params, timeout=30)
resp.raise_for_status()
txs = resp.json().get("tx_responses", [])
return txs
def verify_ipfs_data(self, cid: str) -> Optional[dict]:
"""
通过多个网关验证 IPFS 数据可达性
返回第一个成功获取的 JSON 内容
"""
for gateway in self.gateways:
url = f"{gateway}/{cid}"
try:
resp = requests.get(url, timeout=10)
if resp.status_code == 200:
return resp.json()
except requests.RequestException:
continue
return None
def run_once(self) -> int:
"""单轮索引处理"""
new_events = self.fetch_events(
self.last_processed_height,
self.last_processed_height + 1000,
)
processed = 0
for tx in new_events:
for event in tx.get("events", []):
if event.get("type") == "store_metadata":
attrs = {a["key"]: a["value"] for a in event.get("attributes", [])}
cid = attrs.get("cid")
token_id = attrs.get("token_id")
if cid:
data = self.verify_ipfs_data(cid)
if data:
self._cache_metadata(token_id, cid, data)
processed += 1
else:
print(f"WARNING: CID {cid} not reachable via any gateway")
return processed
def _cache_metadata(self, token_id: str, cid: str, data: dict):
"""缓存元数据到本地"""
cache = {
"token_id": token_id,
"cid": cid,
"metadata": data,
"indexed_at": int(time.time()),
"verified": True,
}
with open(f"cache/{token_id}.json", "w") as f:
json.dump(cache, f, indent=2)
def run_forever(self, interval: int = 6):
"""持续监听(每 interval 秒轮询)"""
while True:
try:
count = self.run_once()
if count > 0:
print(f"Indexed {count} new metadata entries")
except Exception as e:
print(f"Error: {e}")
time.sleep(interval)
5. NFT 元数据实践
5.1 CW721 元数据标准
MSG Chain 上的 NFT 使用 CW721 标准,其元数据格式兼容 OpenSea 规范:
{
"name": "MSG NFT #42",
"description": "A unique NFT on MSG Chain with decentralized storage",
"image": "ipfs://QmImageCID1234567890123456789012345678901234567890",
"image_data": "",
"external_url": "https://msg.space/nft/42",
"background_color": "000000",
"animation_url": "ipfs://QmAnimationCID...",
"youtube_url": "",
"attributes": [
{
"trait_type": "Background",
"value": "Cosmic Nebula"
},
{
"trait_type": "Rarity",
"value": "Legendary"
},
{
"trait_type": "Level",
"value": 42,
"display_type": "number"
},
{
"trait_type": "Boost",
"value": 15.5,
"display_type": "boost_percentage"
}
]
}
IPFS URI 解析
"""
nft_uri_resolver.py - NFT URI 解析器
支持:ipfs://, ipns://, ar://, https://
"""
import re
from dataclasses import dataclass
from typing import Optional
import requests
@dataclass
class ResolvedURI:
"""解析后的 URI 信息"""
original: str
scheme: str # ipfs, ipns, ar, https
identifier: str # CID, IPNS name, Arweave txid
gateway_url: Optional[str] = None
class NFTURIResolver:
"""NFT URI 解析器 - 支持多存储后端"""
def __init__(
self,
ipfs_gateways: Optional[list[str]] = None,
ipns_gateway: str = "https://ipfs.io/ipns",
arweave_gateway: str = "https://arweave.net",
):
self.ipfs_gateways = ipfs_gateways or [
"https://ipfs.io/ipfs",
"https://gateway.pinata.cloud/ipfs",
"https://cloudflare-ipfs.com/ipfs",
]
self.ipns_gateway = ipns_gateway
self.arweave_gateway = arweave_gateway
def resolve(self, uri: str) -> ResolvedURI:
"""解析 URI 为标准格式"""
if uri.startswith("ipfs://"):
cid = uri[7:]
return ResolvedURI(
original=uri,
scheme="ipfs",
identifier=cid,
gateway_url=f"{self.ipfs_gateways[0]}/{cid}",
)
if uri.startswith("ipns://"):
name = uri[7:]
return ResolvedURI(
original=uri,
scheme="ipns",
identifier=name,
gateway_url=f"{self.ipns_gateway}/{name}",
)
if uri.startswith("ar://"):
txid = uri[5:]
return ResolvedURI(
original=uri,
scheme="ar",
identifier=txid,
gateway_url=f"{self.arweave_gateway}/{txid}",
)
if uri.startswith("http"):
return ResolvedURI(
original=uri,
scheme="https",
identifier=uri,
gateway_url=uri,
)
raise ValueError(f"Unsupported URI scheme: {uri}")
def fetch_metadata(self, uri: str) -> Optional[dict]:
"""获取并解析 NFT 元数据,自动尝试多个网关"""
resolved = self.resolve(uri)
if resolved.scheme == "ipfs":
return self._fetch_ipfs(resolved.identifier)
elif resolved.scheme == "ar":
return self._fetch_arweave(resolved.identifier)
elif resolved.scheme == "https":
resp = requests.get(uri, timeout=30)
resp.raise_for_status()
return resp.json()
return None
def _fetch_ipfs(self, cid: str) -> Optional[dict]:
"""通过多个 IPFS 网关获取"""
for gateway in self.ipfs_gateways:
url = f"{gateway}/{cid}"
try:
resp = requests.get(url, timeout=15)
if resp.status_code == 200:
return resp.json()
except requests.RequestException:
continue
return None
def _fetch_arweave(self, txid: str) -> Optional[dict]:
"""从 Arweave 获取"""
url = f"{self.arweave_gateway}/{txid}"
try:
resp = requests.get(url, timeout=15)
if resp.status_code == 200:
return resp.json()
except requests.RequestException:
pass
return None
def get_direct_image_url(self, uri: str) -> str:
"""获取 NFT 图片的直接可访问 URL"""
resolved = self.resolve(uri)
if resolved.gateway_url:
return resolved.gateway_url
return uri
@staticmethod
def normalize_scheme(uri: str) -> str:
"""规范化 URI(确保有 scheme 前缀)"""
if uri.startswith("Qm") or uri.startswith("bafy"):
return f"ipfs://{uri}"
if re.match(r"^[a-zA-Z0-9_-]{43}$", uri):
return f"ar://{uri}"
return uri
// nft-resolver.ts - NFT URI TypeScript 解析器
export interface ResolvedURI {
original: string;
scheme: "ipfs" | "ipns" | "ar" | "https";
identifier: string;
gatewayUrl: string;
}
export class NFTResolver {
private ipfsGateways: string[];
private arweaveGateway: string;
constructor(
ipfsGateways?: string[],
arweaveGateway?: string
) {
this.ipfsGateways = ipfsGateways ?? [
"https://ipfs.io/ipfs",
"https://gateway.pinata.cloud/ipfs",
"https://cloudflare-ipfs.com/ipfs",
];
this.arweaveGateway = arweaveGateway ?? "https://arweave.net";
}
resolve(uri: string): ResolvedURI {
if (uri.startsWith("ipfs://")) {
const cid = uri.slice(7);
return {
original: uri,
scheme: "ipfs",
identifier: cid,
gatewayUrl: `${this.ipfsGateways[0]}/${cid}`,
};
}
if (uri.startsWith("ipns://")) {
const name = uri.slice(7);
return {
original: uri,
scheme: "ipns",
identifier: name,
gatewayUrl: `https://ipfs.io/ipns/${name}`,
};
}
if (uri.startsWith("ar://")) {
const txid = uri.slice(5);
return {
original: uri,
scheme: "ar",
identifier: txid,
gatewayUrl: `${this.arweaveGateway}/${txid}`,
};
}
if (uri.startsWith("http")) {
return {
original: uri,
scheme: "https",
identifier: uri,
gatewayUrl: uri,
};
}
throw new Error(`Unsupported URI scheme: ${uri}`);
}
async fetchMetadata(uri: string): Promise<Record<string, unknown> | null> {
const resolved = this.resolve(uri);
if (resolved.scheme === "ipfs") {
return this.fetchFromIPFS(resolved.identifier);
}
if (resolved.scheme === "ar") {
return this.fetchFromArweave(resolved.identifier);
}
if (resolved.scheme === "https") {
const resp = await fetch(uri);
if (!resp.ok) throw new Error(`HTTP ${resp.status}`);
return resp.json();
}
return null;
}
private async fetchFromIPFS(cid: string): Promise<Record<string, unknown> | null> {
for (const gateway of this.ipfsGateways) {
try {
const resp = await fetch(`${gateway}/${cid}`, {
signal: AbortSignal.timeout(10000),
});
if (resp.ok) return resp.json();
} catch {
continue;
}
}
return null;
}
private async fetchFromArweave(txid: string): Promise<Record<string, unknown> | null> {
try {
const resp = await fetch(`${this.arweaveGateway}/${txid}`, {
signal: AbortSignal.timeout(15000),
});
if (resp.ok) return resp.json();
} catch {
// ignore
}
return null;
}
}
5.2 完整的 NFT 创建与存储流程
"""
nft_mint_flow.py - 完整 NFT Mint 流程
步骤:
1. 准备 NFT 元数据和媒体
2. 上传媒体到 IPFS
3. 上传元数据 JSON 到 IPFS
4. 调用合约存储 CID
"""
import json
import os
from typing import Optional
from pinata_client import PinataClient
from ipfs_client import IPFSClient
class NFTMinter:
"""NFT 创建工具:生成元数据 -> 上传 IPFS -> 写入链上"""
def __init__(
self,
pinata_key: str,
pinata_secret: str,
):
self.pinata = PinataClient(pinata_key, pinata_secret)
self.ipfs = IPFSClient()
def create_metadata(
self,
name: str,
description: str,
image_cid: str,
animation_cid: Optional[str] = None,
attributes: Optional[list[dict]] = None,
external_url: Optional[str] = None,
) -> dict:
"""构建符合 CW721/OpenSea 的元数据"""
metadata = {
"name": name,
"description": description,
"image": f"ipfs://{image_cid}",
"external_url": external_url or "",
"background_color": "",
"attributes": attributes or [],
}
if animation_cid:
metadata["animation_url"] = f"ipfs://{animation_cid}"
return metadata
def upload_image(self, image_path: str) -> str:
"""上传图片到 IPFS,返回 CID"""
result = self.pinata.pin_file(image_path, name=os.path.basename(image_path))
return result["IpfsHash"]
def upload_metadata(self, metadata: dict) -> str:
"""上传元数据 JSON,返回 CID"""
result = self.pinata.pin_json(metadata, name=metadata.get("name", "nft-metadata"))
return result["IpfsHash"]
def mint(
self,
token_id: str,
image_path: str,
name: str,
description: str,
attributes: Optional[list[dict]] = None,
) -> dict:
"""完整 Mint 流程"""
print(f"[1/4] Uploading image: {image_path}")
image_cid = self.upload_image(image_path)
print(f" Image CID: {image_cid}")
print(f"[2/4] Building metadata...")
metadata = self.create_metadata(
name=name,
description=description,
image_cid=image_cid,
attributes=attributes,
)
print(f"[3/4] Uploading metadata to IPFS...")
metadata_cid = self.upload_metadata(metadata)
print(f" Metadata CID: {metadata_cid}")
print(f"[4/4] CID ready for chain: {metadata_cid}")
self.ipfs.pin_add(metadata_cid)
return {
"token_id": token_id,
"image_cid": image_cid,
"metadata_cid": metadata_cid,
}
def batch_mint(
self,
token_ids: list[str],
image_dir: str,
name_template: str,
description: str,
) -> list[dict]:
"""批量 Mint"""
results = []
for token_id in token_ids:
image_path = os.path.join(image_dir, f"{token_id}.png")
name = name_template.replace("{id}", token_id)
result = self.mint(
token_id=token_id,
image_path=image_path,
name=name,
description=description,
)
results.append(result)
return results
# 使用示例
if __name__ == "__main__":
import os
minter = NFTMinter(
pinata_key=os.environ["PINATA_API_KEY"],
pinata_secret=os.environ["PINATA_SECRET_KEY"],
)
result = minter.mint(
token_id="42",
image_path="./assets/nft-42.png",
name="MSG NFT #42",
description="The forty-second NFT on MSG Chain",
attributes=[
{"trait_type": "Rarity", "value": "Epic"},
{"trait_type": "Level", "value": 42, "display_type": "number"},
],
)
print(json.dumps(result, indent=2))
5.3 Lazy Minting 与 Revealed Metadata
"""
lazy_mint.py - 延迟铸造与揭示机制
工作流:
1. 部署合约时上传隐藏元数据(占位图)
2. 用户 Mint 时获取隐藏 token
3. 收集期结束后,上传真实元数据
4. 调用 Reveal 更新所有 token 的 CID
"""
import json
import hashlib
from typing import Optional
class LazyMinter:
"""Lazy Minting / Reveal 机制"""
def __init__(self, ipfs_client, contract_client):
self.ipfs = ipfs_client
self.contract = contract_client
def create_hidden_metadata(self) -> str:
"""创建隐藏元数据(揭示前使用)"""
hidden_metadata = {
"name": "MSG NFT (Hidden)",
"description": "This NFT is waiting to be revealed. Check back soon!",
"image": "ipfs://QmPlaceholderImageCID...",
"attributes": [
{"trait_type": "Status", "value": "Hidden"},
],
}
result = self.ipfs.add_json(hidden_metadata, filename="hidden.json")
return result["Hash"]
def create_revealed_metadata(
self,
token_id: str,
name: str,
description: str,
image_cid: str,
attributes: list[dict],
trait_assignments: Optional[dict] = None,
) -> dict:
"""生成揭示后的真实元数据"""
metadata = {
"name": name,
"description": description,
"image": f"ipfs://{image_cid}",
"attributes": attributes,
}
if trait_assignments:
for trait_type, value in trait_assignments.items():
metadata["attributes"].append({
"trait_type": trait_type,
"value": value,
})
return metadata
def reveal(self, token_ids: list[str], real_cids: list[str]) -> dict:
"""揭示:更新合约中 token 的 CID"""
result = self.contract.batch_store_metadata(token_ids, real_cids)
return result
class CommitRevealNFT:
"""使用 commit-reveal 机制的 NFT 铸造"""
@staticmethod
def hash_metadata(metadata: dict) -> str:
"""计算元数据的哈希(commitment)"""
serialized = json.dumps(metadata, sort_keys=True).encode("utf-8")
return hashlib.sha256(serialized).hexdigest()
def commit(self, token_id: str, metadata_hash: str) -> str:
"""提交 commitment(调用合约 commit_metadata)"""
pass
def reveal(self, token_id: str, metadata_cid: str, metadata: dict) -> str:
"""揭示并验证:合约验证 hash == previous_commitment"""
expected_hash = self.hash_metadata(metadata)
# 合约验证通过后更新 CID
pass
5.4 Arweave NFT 元数据
"""
arweave_nft.py - 使用 Arweave 存储 NFT 元数据
"""
import json
from typing import Optional
class ArweaveNFT:
"""使用 Arweave 永久存储的 NFT"""
def __init__(self, arweave_client):
self.arweave = arweave_client
def upload_metadata(
self,
name: str,
description: str,
image_txid: str,
attributes: Optional[list[dict]] = None,
contract_address: Optional[str] = None,
token_id: Optional[str] = None,
) -> str:
"""上传 NFT 元数据到 Arweave,返回交易 ID"""
metadata = {
"name": name,
"description": description,
"image": f"ar://{image_txid}",
"attributes": attributes or [],
}
tags = [
("Content-Type", "application/json"),
("App-Name", "MSG-NFT"),
("Network", "msg-chain-1"),
]
if contract_address:
tags.append(("Contract", contract_address))
if token_id:
tags.append(("Token-ID", str(token_id)))
txid = self.arweave.upload_json(metadata, tags=tags)
return txid
def upload_image(self, image_path: str) -> str:
"""上传图片到 Arweave"""
mime_map = {
"png": "image/png",
"jpg": "image/jpeg",
"jpeg": "image/jpeg",
"gif": "image/gif",
"svg": "image/svg+xml",
"webp": "image/webp",
}
ext = image_path.split(".")[-1].lower()
content_type = mime_map.get(ext, "application/octet-stream")
tags = [
("Content-Type", content_type),
("App-Name", "MSG-NFT"),
]
return self.arweave.upload_file(image_path, tags=tags)
def get_metadata(self, txid: str) -> dict:
"""从 Arweave 获取元数据"""
import json
data = self.arweave.get_data(txid)
return json.loads(data.decode("utf-8"))
6. 数据市场集成
6.1 数据资产流程
去中心化数据市场使用以下模式:
1. 提供者:
上传数据 -> 计算哈希 -> 注册资产到合约
数据: IPFS 或 Arweave
链上: asset_id, cid, hash, price, owner
2. 消费者:
支付 -> 获取 cid/hash -> 从存储下载 -> 验证哈希
3. 链上验证:
VerifyDataIntegrity(asset_id, claimed_hash) -> bool
"""
data_market.py - 数据市场集成示例
"""
import json
import hashlib
import os
from typing import Optional
from dataclasses import dataclass
@dataclass
class DataAssetInfo:
"""数据资产信息"""
asset_id: str
title: str
description: str
file_path: str
storage_type: str
cid: str
content_hash: str
price: int
encrypted: bool = False
encryption_key_cid: Optional[str] = None
class DataMarketProvider:
"""数据市场提供者"""
def __init__(self, ipfs_client, arweave_client, contract_client):
self.ipfs = ipfs_client
self.arweave = arweave_client
self.contract = contract_client
def prepare_dataset(
self,
file_path: str,
use_arweave: bool = True,
encrypt: bool = False,
encryption_key: Optional[bytes] = None,
) -> tuple[str, str]:
"""
准备数据集并上传
Returns:
(cid/txid, sha256_hash)
"""
sha256_hash = self._compute_sha256(file_path)
if encrypt:
encrypted_path = file_path + ".encrypted"
if encryption_key:
self._encrypt_file(file_path, encrypted_path, encryption_key)
file_path = encrypted_path
if use_arweave:
tags = [
("Content-Type", "application/octet-stream"),
("Original-Hash", f"sha256:{sha256_hash}"),
("App-Name", "MSG-Data-Market"),
]
cid = self.arweave.upload_file(file_path, tags=tags)
else:
result = self.ipfs.add_file(file_path, pin=True)
cid = result["Hash"]
return cid, sha256_hash
def register_asset(
self,
asset_id: str,
cid: str,
content_hash: str,
price: int,
storage_type: str = "arweave",
) -> str:
"""在链上注册数据资产"""
tx_hash = self.contract.register_asset(
asset_id=asset_id,
cid=cid,
storage_type=storage_type,
content_hash=content_hash,
price=price,
)
return tx_hash
def _compute_sha256(self, file_path: str) -> str:
"""计算文件 SHA256"""
h = hashlib.sha256()
with open(file_path, "rb") as f:
for chunk in iter(lambda: f.read(65536), b""):
h.update(chunk)
return h.hexdigest()
def _encrypt_file(self, input_path: str, output_path: str, key: bytes):
"""AES-256-GCM 加密文件"""
from cryptography.fernet import Fernet
f = Fernet(key)
with open(input_path, "rb") as fin:
data = fin.read()
encrypted = f.encrypt(data)
with open(output_path, "wb") as fout:
fout.write(encrypted)
class DataMarketConsumer:
"""数据市场消费者"""
def __init__(self, ipfs_client, arweave_client, contract_client):
self.ipfs = ipfs_client
self.arweave = arweave_client
self.contract = contract_client
def purchase_and_verify(
self,
asset_id: str,
output_path: str,
expected_hash: Optional[str] = None,
) -> bool:
"""
购买并验证数据
Steps:
1. 查询链上资产信息
2. 支付(调用合约)
3. 从存储下载
4. 验证哈希
"""
# 1. 查询
asset = self.contract.get_asset(asset_id)
print(f"Asset: {asset}")
# 2. 支付(调用合约的购买方法)
# tx_hash = self.contract.purchase(asset_id)
# 3. 下载
if asset["storage_type"] == "arweave":
data = self.arweave.get_data(asset["cid"])
else:
data = self.ipfs.cat(asset["cid"])
with open(output_path, "wb") as f:
f.write(data)
print(f"Downloaded to {output_path}")
# 4. 验证
if expected_hash is None:
expected_hash = asset["content_hash"]
actual_hash = hashlib.sha256(data).hexdigest()
matches = actual_hash == expected_hash
print(f"Hash check: {'PASS' if matches else 'FAIL'}")
return matches
def download_and_decrypt(
self,
asset_id: str,
output_path: str,
decryption_key: bytes,
) -> bool:
"""下载并解密加密的数据集"""
asset = self.contract.get_asset(asset_id)
if not asset.get("encrypted"):
return self.purchase_and_verify(asset_id, output_path)
if asset["storage_type"] == "arweave":
data = self.arweave.get_data(asset["cid"])
else:
data = self.ipfs.cat(asset["cid"])
from cryptography.fernet import Fernet
f = Fernet(decryption_key)
decrypted = f.decrypt(data)
with open(output_path, "wb") as fout:
fout.write(decrypted)
actual_hash = hashlib.sha256(decrypted).hexdigest()
matches = actual_hash == asset["content_hash"]
print(f"Decryption + hash check: {'PASS' if matches else 'FAIL'}")
return matches
6.2 加密数据集方案
"""
encrypted_dataset.py - 加密数据集 + 链上密钥管理
架构:
- 数据用对称密钥 (AES-256-GCM) 加密后上传
- 对称密钥用购买者的公钥加密后存储在 IPFS
- 链上存储加密密钥的 CID
"""
import os
import json
from typing import Optional
from cryptography.hazmat.primitives.ciphers.aead import AESGCM
from cryptography.hazmat.primitives import serialization, hashes
from cryptography.hazmat.primitives.asymmetric import padding
class EncryptedDatasetManager:
"""加密数据集管理器"""
@staticmethod
def generate_symmetric_key() -> bytes:
"""生成 AES-256 密钥(32 字节)"""
return AESGCM.generate_key(bit_length=256)
@staticmethod
def encrypt_data(data: bytes, key: bytes) -> tuple[bytes, bytes]:
"""加密数据,返回 (nonce, ciphertext)"""
aesgcm = AESGCM(key)
nonce = os.urandom(12)
ciphertext = aesgcm.encrypt(nonce, data, None)
return nonce, ciphertext
@staticmethod
def decrypt_data(ciphertext: bytes, key: bytes, nonce: bytes) -> bytes:
"""解密数据"""
aesgcm = AESGCM(key)
return aesgcm.decrypt(nonce, ciphertext, None)
@staticmethod
def encrypt_key_with_public_key(
symmetric_key: bytes,
public_key_pem: bytes,
) -> bytes:
"""用购买者的 RSA 公钥加密对称密钥"""
public_key = serialization.load_pem_public_key(public_key_pem)
encrypted_key = public_key.encrypt(
symmetric_key,
padding.OAEP(
mgf=padding.MGF1(algorithm=hashes.SHA256()),
algorithm=hashes.SHA256(),
label=None,
),
)
return encrypted_key
@staticmethod
def decrypt_key_with_private_key(
encrypted_key: bytes,
private_key_pem: bytes,
) -> bytes:
"""用私钥解密对称密钥"""
private_key = serialization.load_pem_private_key(
private_key_pem,
password=None,
)
return private_key.decrypt(
encrypted_key,
padding.OAEP(
mgf=padding.MGF1(algorithm=hashes.SHA256()),
algorithm=hashes.SHA256(),
label=None,
),
)
def prepare_encrypted_dataset(
self,
file_path: str,
buyer_public_key_pem: bytes,
ipfs_client,
) -> dict:
"""
准备加密数据集
Returns: { encrypted_data_cid, encrypted_key_cid, original_hash }
"""
with open(file_path, "rb") as f:
original_data = f.read()
original_hash = hashlib.sha256(original_data).hexdigest()
sym_key = self.generate_symmetric_key()
nonce, ciphertext = self.encrypt_data(original_data, sym_key)
encrypted_result = ipfs_client.add_bytes(
nonce + ciphertext,
filename=os.path.basename(file_path) + ".encrypted",
)
encrypted_cid = encrypted_result["Hash"]
encrypted_key = self.encrypt_key_with_public_key(
sym_key, buyer_public_key_pem
)
key_result = ipfs_client.add_bytes(
encrypted_key,
filename="encrypted_key.bin",
)
key_cid = key_result["Hash"]
return {
"encrypted_data_cid": encrypted_cid,
"encrypted_key_cid": key_cid,
"original_hash": original_hash,
}
def download_and_decrypt(
self,
encrypted_cid: str,
key_cid: str,
private_key_pem: bytes,
ipfs_client,
) -> bytes:
"""下载并解密"""
encrypted_data = ipfs_client.cat(encrypted_cid)
nonce = encrypted_data[:12]
ciphertext = encrypted_data[12:]
encrypted_key = ipfs_client.cat(key_cid)
sym_key = self.decrypt_key_with_private_key(
encrypted_key, private_key_pem
)
return self.decrypt_data(ciphertext, sym_key, nonce)
7. 前端集成
7.1 React 组件:上传到 IPFS
// components/UploadToIPFS.tsx - React 组件:上传文件到 IPFS
import React, { useState, useCallback } from "react";
import axios from "axios";
interface UploadResult {
cid: string;
gatewayUrl: string;
size: number;
fileName: string;
}
interface UploadToIPFSProps {
onUploadComplete?: (result: UploadResult) => void;
pinataApiKey?: string;
pinataSecretKey?: string;
gateway?: string;
}
export const UploadToIPFS: React.FC<UploadToIPFSProps> = ({
onUploadComplete,
pinataApiKey,
pinataSecretKey,
gateway = "https://gateway.pinata.cloud/ipfs",
}) => {
const [uploading, setUploading] = useState(false);
const [progress, setProgress] = useState(0);
const [error, setError] = useState<string | null>(null);
const [result, setResult] = useState<UploadResult | null>(null);
const uploadFile = useCallback(
async (file: File) => {
setUploading(true);
setError(null);
setProgress(0);
try {
if (pinataApiKey && pinataSecretKey) {
const formData = new FormData();
formData.append("file", file);
formData.append(
"pinataMetadata",
JSON.stringify({ name: file.name })
);
const { data } = await axios.post(
"https://api.pinata.cloud/pinning/pinFileToIPFS",
formData,
{
headers: {
pinata_api_key: pinataApiKey,
pinata_secret_api_key: pinataSecretKey,
},
onUploadProgress: (e) => {
if (e.total) setProgress(Math.round((e.loaded / e.total) * 100));
},
}
);
const uploadResult: UploadResult = {
cid: data.IpfsHash,
gatewayUrl: `${gateway}/${data.IpfsHash}`,
size: data.PinSize,
fileName: file.name,
};
setResult(uploadResult);
onUploadComplete?.(uploadResult);
} else {
const formData = new FormData();
formData.append("file", file);
const { data } = await axios.post(
"http://127.0.0.1:5001/api/v0/add",
formData,
{
params: { pin: "true", "cid-version": "1" },
onUploadProgress: (e) => {
if (e.total) setProgress(Math.round((e.loaded / e.total) * 100));
},
}
);
const uploadResult: UploadResult = {
cid: data.Hash,
gatewayUrl: `${gateway}/${data.Hash}`,
size: parseInt(data.Size),
fileName: file.name,
};
setResult(uploadResult);
onUploadComplete?.(uploadResult);
}
} catch (err: any) {
setError(err.message ?? "Upload failed");
} finally {
setUploading(false);
}
},
[pinataApiKey, pinataSecretKey, gateway, onUploadComplete]
);
const handleDrop = useCallback(
(e: React.DragEvent<HTMLDivElement>) => {
e.preventDefault();
const file = e.dataTransfer.files[0];
if (file) uploadFile(file);
},
[uploadFile]
);
const copyToClipboard = (text: string) => {
navigator.clipboard.writeText(text);
};
return (
<div className="ipfs-uploader">
<div
onDrop={handleDrop}
onDragOver={(e) => e.preventDefault()}
className="drop-zone"
style={{
border: "2px dashed #6366f1",
borderRadius: "12px",
padding: "48px",
textAlign: "center",
cursor: "pointer",
backgroundColor: "#f8fafc",
}}
>
{uploading ? (
<div>
<div
className="progress-bar"
style={{
width: `${progress}%`,
height: "4px",
backgroundColor: "#6366f1",
transition: "width 0.3s",
}}
/>
<p>Uploading... {progress}%</p>
</div>
) : (
<div>
<p style={{ fontSize: "18px", marginBottom: "8px" }}>
Drag & drop your file here
</p>
<p style={{ color: "#666" }}>or click to browse</p>
<input
type="file"
onChange={(e) => e.target.files?.[0] && uploadFile(e.target.files[0])}
style={{ display: "none" }}
id="file-input"
/>
</div>
)}
</div>
{error && (
<div className="error" style={{ color: "red", marginTop: "12px" }}>
Error: {error}
</div>
)}
{result && (
<div className="result" style={{ marginTop: "16px" }}>
<h4>Upload Complete</h4>
<table style={{ width: "100%", borderCollapse: "collapse" }}>
<tbody>
<tr>
<td style={{ fontWeight: 600, padding: "4px 8px" }}>CID</td>
<td>
<code>{result.cid}</code>
<button onClick={() => copyToClipboard(result.cid)}>
Copy
</button>
</td>
</tr>
<tr>
<td style={{ fontWeight: 600, padding: "4px 8px" }}>Gateway</td>
<td>
<a
href={result.gatewayUrl}
target="_blank"
rel="noopener noreferrer"
>
{result.gatewayUrl}
</a>
</td>
</tr>
<tr>
<td style={{ fontWeight: 600, padding: "4px 8px" }}>Size</td>
<td>{(result.size / 1024).toFixed(2)} KB</td>
</tr>
</tbody>
</table>
</div>
)}
</div>
);
};
7.2 React 组件:显示 IPFS/Arweave 图片
// components/NFTDisplay.tsx - 显示 NFT 图片(支持 IPFS + Arweave)
import React, { useState, useEffect } from "react";
interface NFTDisplayProps {
uri: string;
alt?: string;
width?: number;
height?: number;
fallbackGateways?: string[];
className?: string;
style?: React.CSSProperties;
}
const DEFAULT_IPFS_GATEWAYS = [
"https://ipfs.io/ipfs",
"https://gateway.pinata.cloud/ipfs",
"https://cloudflare-ipfs.com/ipfs",
"https://dweb.link/ipfs",
];
const ARWEAVE_GATEWAY = "https://arweave.net";
function resolveURI(
uri: string,
gateways: string[] = DEFAULT_IPFS_GATEWAYS
): string[] {
if (uri.startsWith("ipfs://")) {
const cid = uri.slice(7);
return gateways.map((g) => `${g}/${cid}`);
}
if (uri.startsWith("ipns://")) {
const name = uri.slice(7);
return [`https://ipfs.io/ipns/${name}`];
}
if (uri.startsWith("ar://")) {
const txid = uri.slice(5);
return [`${ARWEAVE_GATEWAY}/${txid}`];
}
if (uri.startsWith("data:")) {
return [uri];
}
if (uri.startsWith("http")) {
return [uri];
}
if (uri.startsWith("Qm") || uri.startsWith("bafy")) {
return gateways.map((g) => `${g}/${uri}`);
}
return [uri];
}
export const NFTDisplay: React.FC<NFTDisplayProps> = ({
uri,
alt = "NFT",
width = 400,
height = 400,
fallbackGateways,
className,
style,
}) => {
const gateways = fallbackGateways ?? DEFAULT_IPFS_GATEWAYS;
const urls = resolveURI(uri, gateways);
const [currentIndex, setCurrentIndex] = useState(0);
const [failed, setFailed] = useState(false);
const handleError = () => {
if (currentIndex < urls.length - 1) {
setCurrentIndex((prev) => prev + 1);
} else {
setFailed(true);
}
};
useEffect(() => {
setCurrentIndex(0);
setFailed(false);
}, [uri]);
if (failed) {
return (
<div
className="nft-fallback"
style={{
width,
height,
backgroundColor: "#f1f5f9",
display: "flex",
alignItems: "center",
justifyContent: "center",
borderRadius: "8px",
}}
>
<div style={{ textAlign: "center", color: "#94a3b8" }}>
<p>Failed to load</p>
<p style={{ fontSize: "12px" }}>{uri}</p>
</div>
</div>
);
}
return (
<img
src={urls[currentIndex]}
alt={alt}
width={width}
height={height}
onError={handleError}
className={className}
style={{
borderRadius: "8px",
objectFit: "cover",
...style,
}}
/>
);
};
7.3 上传到 Arweave(React)
// components/UploadToArweave.tsx
import React, { useState } from "react";
import { IrysClient } from "../lib/arweave-client";
interface ArweaveUploadProps {
irysClient: IrysClient;
onUploadComplete: (txid: string) => void;
tags?: Array<{ name: string; value: string }>;
}
export const UploadToArweave: React.FC<ArweaveUploadProps> = ({
irysClient,
onUploadComplete,
tags = [],
}) => {
const [uploading, setUploading] = useState(false);
const [txid, setTxid] = useState<string | null>(null);
const [error, setError] = useState<string | null>(null);
const handleFileChange = async (e: React.ChangeEvent<HTMLInputElement>) => {
const file = e.target.files?.[0];
if (!file) return;
setUploading(true);
setError(null);
try {
const defaultTags = [
{ name: "Content-Type", value: file.type || "application/octet-stream" },
{ name: "App-Name", value: "MSG-DApp" },
{ name: "Network", value: "msg-chain-1" },
...tags,
];
const buffer = await file.arrayBuffer();
const result = await irysClient.uploadData(
Buffer.from(buffer),
defaultTags
);
setTxid(result.id);
onUploadComplete(result.id);
} catch (err: any) {
setError(err.message ?? "Upload to Arweave failed");
} finally {
setUploading(false);
}
};
return (
<div className="arweave-uploader">
<input
type="file"
onChange={handleFileChange}
disabled={uploading}
/>
{uploading && <p>Uploading to Arweave via Irys...</p>}
{error && <p style={{ color: "red" }}>Error: {error}</p>}
{txid && (
<div>
<p>Uploaded!</p>
<p>
Transaction:{" "}
<a
href={`https://arweave.net/${txid}`}
target="_blank"
rel="noopener noreferrer"
>
{txid.slice(0, 12)}...
</a>
</p>
</div>
)}
</div>
);
};
8. 完整示例
8.1 NFT Mint + IPFS 元数据(TypeScript)
// example-nft-mint.ts - 完整 NFT Mint 流程
import { IPFSClient } from "./ipfs-client";
import { PinataClient } from "./pinata-client";
import { NFTResolver } from "./nft-resolver";
import * as fs from "fs";
interface NFTMintResult {
tokenId: string;
imageCid: string;
metadataCid: string;
imageUrl: string;
metadataUrl: string;
}
async function mintNFT(
tokenId: string,
imagePath: string,
name: string,
description: string,
attributes: Array<{ trait_type: string; value: string | number }>
): Promise<NFTMintResult> {
const pinata = new PinataClient(
process.env.PINATA_API_KEY!,
process.env.PINATA_SECRET_KEY!
);
const resolver = new NFTResolver();
const ipfs = new IPFSClient();
// 1. 上传图片
console.log(`[1/4] Uploading image: ${imagePath}`);
const imageResult = await pinata.pinFile(imagePath, name);
const imageCid = imageResult.IpfsHash;
console.log(` Image CID: ${imageCid}`);
// 2. 构建元数据
console.log(`[2/4] Building metadata...`);
const metadata = {
name,
description,
image: `ipfs://${imageCid}`,
external_url: `https://msg.space/nft/${tokenId}`,
attributes,
};
// 3. 上传元数据
console.log(`[3/4] Uploading metadata...`);
const metaResult = await pinata.pinJson(metadata, `${name} Metadata`);
const metadataCid = metaResult.IpfsHash;
console.log(` Metadata CID: ${metadataCid}`);
// 4. 固定到本地节点
await ipfs.pinAdd(metadataCid);
await ipfs.pinAdd(imageCid);
// 5. 准备写链上
console.log(`[4/4] Ready for on-chain storage`);
console.log(` Token ID: ${tokenId}`);
console.log(` Metadata CID: ${metadataCid}`);
console.log(` Now call contract: store_metadata(token_id="${tokenId}", cid="${metadataCid}")`);
return {
tokenId,
imageCid,
metadataCid,
imageUrl: resolver.resolve(metadata.image).gatewayUrl,
metadataUrl: resolver.resolve(`ipfs://${metadataCid}`).gatewayUrl,
};
}
// 完整使用示例
async function main() {
const result = await mintNFT(
"1",
"./assets/msg-nft-1.png",
"MSG Genesis #1",
"The first NFT minted on MSG Chain with IPFS storage",
[
{ trait_type: "Background", value: "Cosmic" },
{ trait_type: "Rarity", value: "Genesis" },
{ trait_type: "Power", value: 100 },
]
);
console.log("\n=== Mint Complete ===");
console.log(JSON.stringify(result, null, 2));
console.log(`\nMetadata Gateway: ${result.metadataUrl}`);
console.log(`Image Gateway: ${result.imageUrl}`);
}
// main().catch(console.error);
8.2 数据市场 + Arweave 存储(TypeScript)
// example-data-market.ts - 数据市场完整示例
import { IrysClient } from "./arweave-client";
import { ArweaveGraphQL } from "./arweave-graphql";
import { createHash } from "crypto";
import * as fs from "fs";
import * as path from "path";
interface DataAssetOnChain {
assetId: string;
storageTxid: string;
contentHash: string;
price: string;
owner: string;
}
class DataMarketExample {
private irys: IrysClient;
private gql: ArweaveGraphQL;
constructor(privateKey: string) {
this.irys = new IrysClient(privateKey);
this.gql = new ArweaveGraphQL();
}
// 计算文件 SHA256 哈希
private computeHash(filePath: string): string {
const data = fs.readFileSync(filePath);
return createHash("sha256").update(data).digest("hex");
}
// 提供者:上传数据集到 Arweave 并注册
async provideDataset(
filePath: string,
assetId: string,
price: string,
license: string = "CC-BY-4.0"
): Promise<DataAssetOnChain> {
console.log(`[1/3] Computing hash...`);
const contentHash = this.computeHash(filePath);
console.log(`[2/3] Uploading to Arweave via Irys...`);
const tags = [
{ name: "Content-Type", value: "application/octet-stream" },
{ name: "App-Name", value: "MSG-Data-Market" },
{ name: "Network", value: "msg-chain-1" },
{ name: "Asset-ID", value: assetId },
{ name: "Original-Hash", value: `sha256:${contentHash}` },
{ name: "License", value: license },
];
const result = await this.irys.uploadFile(filePath, tags);
console.log(` Arweave TX: ${result.id}`);
console.log(`[3/3] Ready for on-chain registration`);
console.log(` Call: register_asset(
asset_id="${assetId}",
cid="${result.id}",
storage_type="arweave",
content_hash="${contentHash}",
price=${price}
)`);
return {
assetId,
storageTxid: result.id,
contentHash,
price,
owner: "msg1...",
};
}
// 消费者:查询资产并验证
async verifyDataset(
storageTxid: string,
localFilePath: string,
expectedHash: string
): Promise<boolean> {
console.log(`[1/2] Fetching data from Arweave...`);
const data = await this.fetchFromGateway(storageTxid);
console.log(`[2/2] Verifying hash...`);
const actualHash = createHash("sha256").update(data).digest("hex");
const matches = actualHash === expectedHash;
console.log(` Expected: ${expectedHash}`);
console.log(` Actual: ${actualHash}`);
console.log(` Match: ${matches ? "PASS" : "FAIL"}`);
if (matches) {
fs.writeFileSync(localFilePath, data);
console.log(` Saved to: ${localFilePath}`);
}
return matches;
}
// 通过 Arweave 网关获取数据
private async fetchFromGateway(txid: string): Promise<Buffer> {
const url = `https://arweave.net/${txid}`;
const resp = await fetch(url);
if (!resp.ok) throw new Error(`HTTP ${resp.status}`);
const buffer = await resp.arrayBuffer();
return Buffer.from(buffer);
}
// 查询提供者的所有数据集
async queryProviderDatasets(contractAddress: string): Promise<void> {
const transactions = await this.gql.findByTags([
{ name: "App-Name", value: "MSG-Data-Market" },
{ name: "Contract", value: contractAddress },
]);
console.log(`Found ${transactions.length} datasets:`);
for (const tx of transactions) {
const tagMap: Record<string, string> = {};
tx.tags.forEach((t) => (tagMap[t.name] = t.value));
console.log(` TX: ${tx.id}`);
console.log(` Asset-ID: ${tagMap["Asset-ID"]}`);
console.log(` Hash: ${tagMap["Original-Hash"]}`);
console.log(` License: ${tagMap["License"]}`);
console.log("---");
}
}
}
async function main() {
const market = new DataMarketExample("your-private-key");
// 提供者流程
const asset = await market.provideDataset(
"./datasets/weather-data-2025.csv",
"weather-2025",
"1000000",
"CC-BY-4.0"
);
// 消费者下载验证
const valid = await market.verifyDataset(
asset.storageTxid,
"./downloads/verified-weather-data.csv",
asset.contentHash
);
console.log(`\nVerification: ${valid ? "SUCCESS" : "FAILED"}`);
}
// main().catch(console.error);
8.3 完整流程:IPFS + Arweave 混合(Python)
"""
example_hybrid.py - IPFS + Arweave 混合存储完整示例
场景: NFT 创建
- 元数据存 IPFS(可更新)
- 元数据和图片同时备份到 Arweave(永久保存)
"""
import json
import hashlib
import os
from typing import Optional
import requests
class HybridStorageExample:
"""
混合存储示例
IPFS: 热数据,快速访问,可更新
Arweave: 冷备份,永久保存,不可篡改
"""
def __init__(
self,
pinata_key: str,
pinata_secret: str,
arweave_key: str,
):
self.pinata_headers = {
"pinata_api_key": pinata_key,
"pinata_secret_api_key": pinata_secret,
}
self.arweave_key = arweave_key
def upload_to_ipfs(self, data: dict, name: str) -> str:
"""上传 JSON 到 IPFS"""
url = "https://api.pinata.cloud/pinning/pinJSONToIPFS"
payload = {
"pinataContent": data,
"pinataMetadata": {"name": name},
}
resp = requests.post(url, json=payload, headers=self.pinata_headers, timeout=30)
resp.raise_for_status()
return resp.json()["IpfsHash"]
def upload_to_arweave(self, data: dict, tags: list) -> str:
"""上传 JSON 到 Arweave"""
import base64
url = "https://node2.irys.xyz/tx/arweave"
payload = {
"data": json.dumps(data).encode("utf-8").hex(),
"tags": [{"name": k, "value": v} for k, v in tags],
}
resp = requests.post(url, json=payload, timeout=60)
resp.raise_for_status()
return resp.json()["id"]
def create_nft(
self,
token_id: str,
name: str,
description: str,
image_url: str,
attributes: Optional[list] = None,
) -> dict:
"""
创建 NFT,同时存储到 IPFS(可更新)和 Arweave(永久备份)
Returns:
{ token_id, ipfs_cid, arweave_txid }
"""
metadata = {
"name": name,
"description": description,
"image": image_url,
"attributes": attributes or [],
"token_id": token_id,
}
# 上传到 IPFS(快速访问)
print("[1/3] Uploading metadata to IPFS (Pinata)...")
ipfs_cid = self.upload_to_ipfs(metadata, name)
print(f" IPFS CID: {ipfs_cid}")
print(f" Gateway: https://gateway.pinata.cloud/ipfs/{ipfs_cid}")
# 上传到 Arweave(永久备份)
print("[2/3] Uploading to Arweave (Irys)...")
arweave_tags = [
("Content-Type", "application/json"),
("App-Name", "MSG-NFT-Hybrid"),
("Network", "msg-chain-1"),
("Token-ID", token_id),
("IPFS-CID", ipfs_cid),
]
arweave_txid = self.upload_to_arweave(metadata, arweave_tags)
print(f" Arweave TX: {arweave_txid}")
print(f" Gateway: https://arweave.net/{arweave_txid}")
# 记录映射关系
mapping = {
"token_id": token_id,
"name": name,
"ipfs_cid": ipfs_cid,
"arweave_txid": arweave_txid,
"ipfs_url": f"https://gateway.pinata.cloud/ipfs/{ipfs_cid}",
"arweave_url": f"https://arweave.net/{arweave_txid}",
}
with open(f"nft_{token_id}_mapping.json", "w") as f:
json.dump(mapping, f, indent=2)
print(f"[3/3] Complete! Mapping saved to nft_{token_id}_mapping.json")
return mapping
def verify_consistency(self, ipfs_cid: str, arweave_txid: str) -> bool:
"""
验证 IPFS 和 Arweave 上的数据是否一致
"""
# 从 IPFS 获取
ipfs_url = f"https://gateway.pinata.cloud/ipfs/{ipfs_cid}"
ipfs_resp = requests.get(ipfs_url, timeout=15)
ipfs_data = ipfs_resp.json()
# 从 Arweave 获取
ar_url = f"https://arweave.net/{arweave_txid}"
ar_resp = requests.get(ar_url, timeout=15)
ar_data = ar_resp.json()
# 比较哈希
ipfs_hash = hashlib.sha256(
json.dumps(ipfs_data, sort_keys=True).encode()
).hexdigest()
ar_hash = hashlib.sha256(
json.dumps(ar_data, sort_keys=True).encode()
).hexdigest()
consistent = ipfs_hash == ar_hash
print(f"Consistency check: {'PASS' if consistent else 'FAIL'}")
return consistent
# 使用示例
if __name__ == "__main__":
import os
storage = HybridStorageExample(
pinata_key=os.environ["PINATA_API_KEY"],
pinata_secret=os.environ["PINATA_SECRET_KEY"],
arweave_key=os.environ["ARWEAVE_PRIVATE_KEY"],
)
result = storage.create_nft(
token_id="1",
name="MSG Hybrid NFT #1",
description="NFT stored on both IPFS and Arweave",
image_url="ipfs://QmExampleImageCID...",
attributes=[
{"trait_type": "Storage", "value": "Hybrid"},
{"trait_type": "IPFS", "value": "Updateable"},
{"trait_type": "Arweave", "value": "Permanent"},
],
)
# 验证一致性
storage.verify_consistency(
result["ipfs_cid"],
result["arweave_txid"],
)
附录
A. MSG Chain 相关信息
| 项目 | 值 |
|---|---|
| 链 ID | msg-chain-1 |
| 地址前缀 | msg |
| 共识机制 | Tendermint / Cosmos SDK |
| 智能合约 | CosmWasm |
| 代币 | MSG |
| RPC 端点 | 请参考官方文档 |
| LCD/REST | 请参考官方文档 |
B. IPFS 网关列表
| 网关 | URL | 限制 |
|---|---|---|
| ipfs.io | https://ipfs.io/ipfs | 可能限速 |
| Pinata | https://gateway.pinata.cloud/ipfs | 需 API 密钥(上传) |
| Cloudflare | https://cloudflare-ipfs.com/ipfs | 无上传 |
| dweb.link | https://dweb.link/ipfs | 无上传 |
| Fleek | https://ipfs.fleek.co/ipfs | 需注册 |
| 4everland | https://4everland.io/ipfs | 需注册 |
| w3s.link | https://w3s.link/ipfs | Web3.Storage 网关 |
C. Arweave 相关资源
| 资源 | URL |
|---|---|
| 公共网关 | https://arweave.net |
| GraphQL 端点 | https://arweave.net/graphql |
| Irys (Bundlr) | https://node2.irys.xyz |
| ArConnect 钱包 | https://arconnect.io |
| Arweave.app | https://arweave.app |
| 区块浏览器 | https://viewblock.io/arweave |
D. 常用依赖
// package.json
{
"dependencies": {
"axios": "^1.7.0",
"form-data": "^4.0.0",
"arweave": "^1.15.0",
"@irys/sdk": "^0.2.0"
}
}
# requirements.txt
requests==2.31.0
arweave-python-client==1.0.0 # 或使用 irys-sdk
cryptography==41.0.0
E. 故障排除
IPFS 上传后内容不可达
- 检查是否已 pin:
ipfs pin ls | grep <CID> - 检查网关是否支持该 CID 版本(v1 需要较新网关)
- 使用多个网关重试
- 通过 IPFS Desktop 或本地节点验证
Arweave 交易未确认
- 检查余额:
irys get_balance - 检查节点状态:
curl https://arweave.net/status - 等待主网确认(原生客户端)
- 使用 Irys 即时确认
合约交易失败
- 检查 CID 格式(必须为有效的 v0 或 v1)
- 检查 token_id 是否重复
- 检查调用者权限
- 检查 Gas 是否充足
F. 安全建议
- 永远不要将私钥硬编码在代码中,使用环境变量或密钥管理服务
- 验证用户输入的 CID 格式后再发送交易,避免浪费 Gas
- 使用多个网关 fallback,防止单点故障
- 定期 pin 重要数据,或使用 pinning 服务确保持久性
- Arweave 标签应包含足够元信息,便于后续查询和审计
- 加密敏感数据后再上传,密钥管理在链下进行
- 验证数据完整性:下载后对比 SHA256 哈希
本文档为 MSG Chain 开发者提供 IPFS 和 Arweave 集成的完整参考。
如有问题,请查阅 MSG Chain 官方文档或提交 GitHub Issue。
