AI Agent 自然语言与LLM集成指南
MSG Chain — msg 前缀 · 面向 AI Agent 网络的下一代区块链
目录
1. 概述
1.1 为什么 LLM 集成对 Agent 至关重要
AI Agent 在 MSG Chain 上不仅执行预定义的智能合约逻辑,还需要理解自然语言、进行推理、生成响应并自主决策。大语言模型为此提供了核心动力:
- 自然语言理解:将用户意图解析为可执行的链上操作
- 上下文推理:跨多个交互维护对话状态与目标
- 工具使用:通过函数调用与链上合约、交易、查询交互
- 人格塑造:为 Agent 赋予一致的行为模式与价值观
1.2 架构总览
┌─────────────────────────────────────────────────┐
│ 用户/客户端 │
└──────────────────┬──────────────────────────────┘
│ HTTP / WebSocket / gRPC
▼
┌─────────────────────────────────────────────────┐
│ Agent 运行时 (Agent Runtime) │
│ ┌───────────┐ ┌───────────┐ ┌─────────────┐ │
│ │ Router │ │ Memory │ │ Tool │ │
│ │ Layer │ │ Manager │ │ Registry │ │
│ └─────┬─────┘ └───────────┘ └──────┬──────┘ │
│ │ │ │
└────────┼───────────────────────────────┼─────────┘
│ │
▼ ▼
┌────────────────┐ ┌───────────────────┐
│ LLM 提供商 │ │ MSG Chain 节点 │
│ GPT-4 / Claude│◄──────────►│ msg1... addresses │
│ Llama / Local │ │ CosmWasm / EVM │
└────────────────┘ └───────────────────┘
1.3 通信流程
- 用户发送自然语言请求到 Agent 端点
- Agent 运行时解析请求,构建上下文(含对话历史、Agent 宪法、链上状态)
- LLM 路由层选择合适的模型生成响应
- 响应中包含工具调用(转账、查询、执行合约等)
- 工具注册表验证并执行调用,结果返回 LLM
- LLM 生成最终自然语言响应返回用户
- 交互记录写入短期记忆和长期向量存储,可选锚定到链上
1.4 支持的 LLM
| 提供商 | 模型 | 特点 | 适用场景 |
|---|---|---|---|
| OpenAI | GPT-4 / GPT-4o | 最强推理、工具调用 | 复杂合约交互、多步推理 |
| Anthropic | Claude 3.5 Sonnet | 长上下文、安全 | 宪法约束、敏感场景 |
| 本地 | Llama 3 / Qwen2 | 低延迟、隐私 | 高频查询、数据敏感 |
| MSG 内置 | msg-llm-lightweight | 链上轻量推理 | Gas 优化场景 |
2. LLM 路由层
2.1 核心抽象
LLM 路由层是 Agent 与各种语言模型之间的统一接口。它负责请求分发、负载均衡、成本控制和故障转移。
import os
import json
import time
import asyncio
from abc import ABC, abstractmethod
from dataclasses import dataclass, field
from typing import Optional, Any
from enum import Enum
import hashlib
class LLMProviderType(str, Enum):
"""支持的 LLM 提供商类型"""
OPENAI = "openai"
ANTHROPIC = "anthropic"
LOCAL = "local"
MSG_NATIVE = "msg_native"
@dataclass
class AgentContext:
"""Agent 运行时上下文"""
agent_id: str
agent_address: str # msg1... 格式地址
did: str
capabilities: list[str]
constitution: str
conversation_id: str
user_id: str
chain_id: str = "msg-chain-1"
metadata: dict = field(default_factory=dict)
@dataclass
class LLMRequest:
"""统一的 LLM 请求格式"""
prompt: str
system_prompt: str = ""
temperature: float = 0.7
max_tokens: int = 2048
tools: list[dict] = field(default_factory=list)
context: Optional[AgentContext] = None
stop_sequences: list[str] = field(default_factory=list)
@dataclass
class LLMResponse:
"""统一的 LLM 响应格式"""
content: str
tool_calls: list[dict] = field(default_factory=list)
model: str = ""
usage: dict = field(default_factory=dict)
latency_ms: float = 0.0
class BaseLLMProvider(ABC):
"""LLM 提供商的抽象基类"""
@abstractmethod
async def generate(self, request: LLMRequest) -> LLMResponse:
...
@abstractmethod
def estimate_cost(self, request: LLMRequest) -> float:
...
@property
@abstractmethod
def provider_type(self) -> LLMProviderType:
...
class OpenAIProvider(BaseLLMProvider):
"""OpenAI GPT-4/GPT-4o 提供商"""
def __init__(self, api_key: str, model: str = "gpt-4o"):
self.api_key = api_key
self.model = model
self._client = None # 运行时懒加载
async def _get_client(self):
if self._client is None:
from openai import AsyncOpenAI
self._client = AsyncOpenAI(api_key=self.api_key)
return self._client
async def generate(self, request: LLMRequest) -> LLMResponse:
start = time.time()
client = await self._get_client()
messages = []
if request.system_prompt:
messages.append({"role": "system", "content": request.system_prompt})
messages.append({"role": "user", "content": request.prompt})
kwargs = {
"model": self.model,
"messages": messages,
"temperature": request.temperature,
"max_tokens": request.max_tokens,
}
if request.stop_sequences:
kwargs["stop"] = request.stop_sequences
if request.tools:
kwargs["tools"] = request.tools
kwargs["tool_choice"] = "auto"
response = await client.chat.completions.create(**kwargs)
result = response.choices[0]
tool_calls = []
if result.message.tool_calls:
for tc in result.message.tool_calls:
tool_calls.append({
"id": tc.id,
"type": "function",
"function": {
"name": tc.function.name,
"arguments": tc.function.arguments,
}
})
return LLMResponse(
content=result.message.content or "",
tool_calls=tool_calls,
model=self.model,
usage={
"prompt_tokens": response.usage.prompt_tokens,
"completion_tokens": response.usage.completion_tokens,
"total_tokens": response.usage.total_tokens,
},
latency_ms=(time.time() - start) * 1000,
)
def estimate_cost(self, request: LLMRequest) -> float:
"""基于 token 估算成本"""
prompt_tokens = len(request.prompt) // 4 + len(request.system_prompt) // 4
completion_tokens = request.max_tokens
if self.model == "gpt-4o":
return (prompt_tokens / 1_000_000 * 2.50 +
completion_tokens / 1_000_000 * 10.00)
return 0.0
@property
def provider_type(self) -> LLMProviderType:
return LLMProviderType.OPENAI
class ClaudeProvider(BaseLLMProvider):
"""Anthropic Claude 提供商"""
COST_PER_INPUT_TOKEN = 3.00 / 1_000_000
COST_PER_OUTPUT_TOKEN = 15.00 / 1_000_000
def __init__(self, api_key: str, model: str = "claude-3-5-sonnet-通用维护记录"):
self.api_key = api_key
self.model = model
self._client = None
async def _get_client(self):
if self._client is None:
from anthropic import AsyncAnthropic
self._client = AsyncAnthropic(api_key=self.api_key)
return self._client
async def generate(self, request: LLMRequest) -> LLMResponse:
start = time.time()
client = await self._get_client()
system = request.system_prompt or None
messages = [{"role": "user", "content": request.prompt}]
kwargs = {
"model": self.model,
"messages": messages,
"max_tokens": request.max_tokens,
"temperature": request.temperature,
}
if system:
kwargs["system"] = system
if request.tools:
kwargs["tools"] = [
{
"name": t["function"]["name"],
"description": t["function"].get("description", ""),
"input_schema": t["function"].get("parameters", {}),
}
for t in request.tools
]
response = await client.messages.create(**kwargs)
content = ""
tool_calls = []
for block in response.content:
if block.type == "text":
content = block.text
elif block.type == "tool_use":
tool_calls.append({
"id": block.id,
"type": "function",
"function": {
"name": block.name,
"arguments": json.dumps(block.input),
}
})
return LLMResponse(
content=content,
tool_calls=tool_calls,
model=self.model,
usage={
"input_tokens": response.usage.input_tokens,
"output_tokens": response.usage.output_tokens,
},
latency_ms=(time.time() - start) * 1000,
)
def estimate_cost(self, request: LLMRequest) -> float:
prompt_tokens = len(request.prompt) // 4
completion_tokens = request.max_tokens
return (prompt_tokens * self.COST_PER_INPUT_TOKEN +
completion_tokens * self.COST_PER_OUTPUT_TOKEN)
@property
def provider_type(self) -> LLMProviderType:
return LLMProviderType.ANTHROPIC
class LocalLLMProvider(BaseLLMProvider):
"""本地模型提供商(llama.cpp / vLLM)"""
def __init__(
self,
model_path: str = "/models/llama",
endpoint: str = "http://localhost:8000/v1",
api_key: str = "",
):
self.model_path = model_path
self.endpoint = endpoint
self.api_key = api_key
self._client = None
async def _get_client(self):
if self._client is None:
from openai import AsyncOpenAI
self._client = AsyncOpenAI(
api_key=self.api_key or "not-needed",
base_url=self.endpoint,
)
return self._client
async def generate(self, request: LLMRequest) -> LLMResponse:
start = time.time()
client = await self._get_client()
messages = []
if request.system_prompt:
messages.append({"role": "system", "content": request.system_prompt})
messages.append({"role": "user", "content": request.prompt})
response = await client.chat.completions.create(
model="local",
messages=messages,
temperature=request.temperature,
max_tokens=request.max_tokens,
)
return LLMResponse(
content=response.choices[0].message.content or "",
model="local-llama",
usage={
"total_tokens": response.usage.total_tokens if response.usage else 0,
},
latency_ms=(time.time() - start) * 1000,
)
def estimate_cost(self, request: LLMRequest) -> float:
return 0.0 # 本地模型无 API 成本
@property
def provider_type(self) -> LLMProviderType:
return LLMProviderType.LOCAL
class MSGNativeLLMProvider(BaseLLMProvider):
"""MSG 链上轻量推理提供商"""
def __init__(self, chain_rpc: str = "https://rpc.msg-chain-1.msg"):
self.chain_rpc = chain_rpc
async def generate(self, request: LLMRequest) -> LLMResponse:
start = time.time()
import aiohttp
payload = {
"agent_id": request.context.agent_id if request.context else "",
"prompt": request.prompt[:512], # 链上限制
"temperature": request.temperature,
"max_tokens": min(request.max_tokens, 256),
}
async with aiohttp.ClientSession() as session:
async with session.post(
f"{self.chain_rpc}/msg/llm/infer",
json=payload,
) as resp:
data = await resp.json()
return LLMResponse(
content=data.get("response", ""),
model="msg-llm-lightweight",
usage={"gas_used": data.get("gas_used", 0)},
latency_ms=(time.time() - start) * 1000,
)
def estimate_cost(self, request: LLMRequest) -> float:
return 0.001 # ~0.001 MSG 每次推理
@property
def provider_type(self) -> LLMProviderType:
return LLMProviderType.MSG_NATIVE
class RoutingStrategy(str, Enum):
"""路由策略"""
COST_FIRST = "cost_first" # 最低成本优先
CAPABILITY_FIRST = "capability_first" # 最强能力优先
LATENCY_FIRST = "latency_first" # 最低延迟优先
MANUAL = "manual" # 手动指定
@dataclass
class RoutingRule:
"""路由规则"""
strategy: RoutingStrategy
capabilities: list[str] = field(default_factory=list)
max_cost_per_request: float = float("inf")
preferred_providers: list[LLMProviderType] = field(default_factory=list)
fallback_providers: list[LLMProviderType] = field(default_factory=list)
class LLMRouter:
"""LLM 路由层 — 将 Agent 查询路由到合适的 LLM 提供商
支持多提供商抽象、成本感知路由、故障转移和提示词模板管理。
"""
def __init__(self, rules: Optional[list[RoutingRule]] = None):
self.providers: dict[LLMProviderType, BaseLLMProvider] = {
LLMProviderType.OPENAI: OpenAIProvider(
api_key=os.getenv("OPENAI_API_KEY", "")
),
LLMProviderType.ANTHROPIC: ClaudeProvider(
api_key=os.getenv("ANTHROPIC_API_KEY", "")
),
LLMProviderType.LOCAL: LocalLLMProvider(
model_path=os.getenv("LOCAL_LLM_PATH", "/models/llama"),
),
LLMProviderType.MSG_NATIVE: MSGNativeLLMProvider(),
}
self.rules = rules or [
RoutingRule(
strategy=RoutingStrategy.CAPABILITY_FIRST,
capabilities=["complex_reasoning", "contract_execution"],
preferred_providers=[LLMProviderType.OPENAI],
fallback_providers=[LLMProviderType.ANTHROPIC, LLMProviderType.LOCAL],
),
RoutingRule(
strategy=RoutingStrategy.COST_FIRST,
capabilities=["simple_query"],
preferred_providers=[LLMProviderType.LOCAL, LLMProviderType.MSG_NATIVE],
fallback_providers=[LLMProviderType.OPENAI],
),
]
self._template_registry: dict[str, str] = {}
self._request_counter = 0
self._daily_cost = 0.0
self._max_daily_cost = float(os.getenv("MAX_DAILY_LLM_COST", "10.0"))
def register_template(self, name: str, template: str) -> None:
"""注册提示词模板"""
self._template_registry[name] = template
def render_template(self, name: str, **kwargs) -> str:
"""渲染提示词模板"""
template = self._template_registry.get(name)
if not template:
raise ValueError(f"Template '{name}' not found")
return template.format(**kwargs)
def select_provider(
self,
request: LLMRequest,
strategy: RoutingStrategy = RoutingStrategy.CAPABILITY_FIRST,
) -> BaseLLMProvider:
"""根据策略选择最优提供商"""
preferred_order = {
RoutingStrategy.CAPABILITY_FIRST: [
LLMProviderType.OPENAI,
LLMProviderType.ANTHROPIC,
LLMProviderType.LOCAL,
LLMProviderType.MSG_NATIVE,
],
RoutingStrategy.COST_FIRST: [
LLMProviderType.MSG_NATIVE,
LLMProviderType.LOCAL,
LLMProviderType.ANTHROPIC,
LLMProviderType.OPENAI,
],
RoutingStrategy.LATENCY_FIRST: [
LLMProviderType.LOCAL,
LLMProviderType.MSG_NATIVE,
LLMProviderType.OPENAI,
LLMProviderType.ANTHROPIC,
],
}
order = preferred_order.get(strategy, preferred_order[RoutingStrategy.CAPABILITY_FIRST])
for ptype in order:
if ptype in self.providers:
provider = self.providers[ptype]
cost = provider.estimate_cost(request)
if self._daily_cost + cost <= self._max_daily_cost:
return provider
# 保底策略
return self.providers[LLMProviderType.LOCAL]
async def generate(
self,
prompt: str,
context: AgentContext,
tools: Optional[list[dict]] = None,
strategy: Optional[RoutingStrategy] = None,
) -> LLMResponse:
"""生成 LLM 响应,包含自动路由和故障转移"""
self._request_counter += 1
matched_rule = None
for rule in self.rules:
if context.capabilities and any(
cap in rule.capabilities for cap in context.capabilities
):
matched_rule = rule
break
if not matched_rule:
matched_rule = self.rules[0]
strategy = strategy or matched_rule.strategy
request = LLMRequest(
prompt=prompt,
system_prompt=self._build_system_prompt(context),
context=context,
tools=tools or [],
)
last_error = None
provider_order = (
matched_rule.preferred_providers +
matched_rule.fallback_providers
)
for ptype in provider_order:
provider = self.providers.get(ptype)
if not provider:
continue
try:
response = await provider.generate(request)
cost = provider.estimate_cost(request)
self._daily_cost += cost
return response
except Exception as e:
last_error = e
continue
raise RuntimeError(
f"All LLM providers failed. Last error: {last_error}"
)
def _build_system_prompt(self, context: AgentContext) -> str:
"""构建系统提示词,注入 Agent 信息和宪法"""
return (
f"你是一个运行在 MSG Chain 上的 AI Agent。\n"
f"Agent ID: {context.agent_id}\n"
f"Agent 地址: {context.agent_address}\n"
f"DID: {context.did}\n"
f"能力: {', '.join(context.capabilities)}\n"
f"宪法约束: {context.constitution}\n"
f"你只能使用你的注册能力范围内的工具和操作。\n"
f"所有链上操作必须符合 MSG Chain 的规则。"
)
def get_stats(self) -> dict:
"""获取路由统计信息"""
return {
"total_requests": self._request_counter,
"daily_cost": self._daily_cost,
"max_daily_cost": self._max_daily_cost,
"provider_count": len(self.providers),
}
# ============================================================
# 使用示例
# ============================================================
async def demo_llm_router():
"""演示 LLM 路由层的基本使用"""
router = LLMRouter()
# 注册系统提示词模板
router.register_template(
"agent_intro",
"你是 {name},一个基于 MSG Chain 的 AI Agent。你的地址是 {address}。"
)
context = AgentContext(
agent_id="agent-demo-1",
agent_address="msg1agentdemoexampleaddress",
did="did:msg:agent-demo-1",
capabilities=["complex_reasoning", "simple_query"],
constitution="不得转移用户资产除非获得明确授权",
conversation_id="conv-001",
user_id="user-001",
)
response = await router.generate(
prompt="请查询我的 MSG 代币余额并总结近期的交易活动。",
context=context,
tools=[
{
"type": "function",
"function": {
"name": "check_balance",
"description": "查询 MSG 代币余额",
"parameters": {
"type": "object",
"properties": {
"address": {
"type": "string",
"description": "钱包地址",
}
},
"required": ["address"],
},
},
}
],
)
print(f"响应: {response.content}")
print(f"模型: {response.model}")
print(f"延迟: {response.latency_ms:.0f}ms")
print(f"工具调用: {response.tool_calls}")
2.2 成本感知路由
class CostTracker:
"""跟踪和控制 LLM 使用成本"""
def __init__(self, budget_usd: float = 10.0, period_hours: int = 24):
self.budget_usd = budget_usd
self.period_hours = period_hours
self._usage_log: list[dict] = []
def record_usage(self, provider: str, model: str, tokens: int, cost: float):
self._usage_log.append({
"provider": provider,
"model": model,
"tokens": tokens,
"cost": cost,
"timestamp": time.time(),
})
def current_period_cost(self) -> float:
cutoff = time.time() - self.period_hours * 3600
return sum(
entry["cost"]
for entry in self._usage_log
if entry["timestamp"] > cutoff
)
def remaining_budget(self) -> float:
return self.budget_usd - self.current_period_cost()
def can_afford(self, estimated_cost: float) -> bool:
return self.remaining_budget() >= estimated_cost
def should_fallback(self, provider: str, estimated_cost: float) -> bool:
"""判断是否应降级到更便宜的模型"""
return (
self.remaining_budget() < self.budget_usd * 0.2
and estimated_cost > 0.01
)
@dataclass
class CostAwareRoutingConfig:
"""成本感知路由配置"""
budget_per_day_usd: float = 10.0
budget_per_request_usd: float = 0.05
high_cost_threshold_usd: float = 0.02
emergency_fallback_provider: LLMProviderType = LLMProviderType.LOCAL
2.3 故障转移策略
class FallbackStrategy:
"""故障转移策略 — 提供商降级与恢复"""
PROVIDER_PRIORITY = [
LLMProviderType.OPENAI,
LLMProviderType.ANTHROPIC,
LLMProviderType.LOCAL,
LLMProviderType.MSG_NATIVE,
]
def __init__(self):
self._failure_count: dict[LLMProviderType, int] = {}
self._cooldown_until: dict[LLMProviderType, float] = {}
self._max_retries = 3
self._cooldown_seconds = 30
def record_failure(self, provider: LLMProviderType) -> None:
self._failure_count[provider] = self._failure_count.get(provider, 0) + 1
if self._failure_count[provider] >= self._max_retries:
self._cooldown_until[provider] = time.time() + self._cooldown_seconds
self._failure_count[provider] = 0
def record_success(self, provider: LLMProviderType) -> None:
self._failure_count[provider] = 0
self._cooldown_until.pop(provider, None)
def is_available(self, provider: LLMProviderType) -> bool:
if provider in self._cooldown_until:
if time.time() < self._cooldown_until[provider]:
return False
del self._cooldown_until[provider]
return True
def get_available_providers(self) -> list[LLMProviderType]:
return [p for p in self.PROVIDER_PRIORITY if self.is_available(p)]
2.4 提示词模板管理
class PromptTemplateManager:
"""提示词模板管理 — 版本化、参数化模板"""
def __init__(self):
self._templates: dict[str, PromptTemplate] = {}
def add_template(
self,
name: str,
template: str,
version: str = "1.0",
description: str = "",
) -> None:
self._templates[name] = PromptTemplate(
name=name,
template=template,
version=version,
description=description,
)
def get_template(self, name: str, version: Optional[str] = None) -> "PromptTemplate":
tmpl = self._templates.get(name)
if not tmpl:
raise KeyError(f"Template '{name}' not found")
if version and tmpl.version != version:
raise ValueError(
f"Version mismatch for '{name}': requested {version}, "
f"available {tmpl.version}"
)
return tmpl
def render(self, name: str, version: Optional[str] = None, **kwargs) -> str:
tmpl = self.get_template(name, version)
return tmpl.render(**kwargs)
def list_templates(self) -> list[dict]:
return [
{
"name": t.name,
"version": t.version,
"description": t.description,
}
for t in self._templates.values()
]
@dataclass
class PromptTemplate:
"""单个提示词模板"""
name: str
template: str
version: str
description: str
def render(self, **kwargs) -> str:
return self.template.format(**kwargs)
3. 工具调用模式
3.1 工具注册表
Agent 通过工具调用与 MSG Chain 交互。LLM 生成结构化工具调用请求,Agent 运行时执行这些调用。
import hashlib
import json
from typing import Callable, Coroutine, Any
from dataclasses import dataclass, field
from enum import Enum
import asyncio
class ToolCategory(str, Enum):
"""工具分类"""
QUERY = "query" # 查询类(只读)
TRANSACTION = "transaction" # 交易类(写入链上)
AGENT = "agent" # Agent 间交互
UTILITY = "utility" # 工具类
@dataclass
class ToolParameter:
"""工具参数定义"""
name: str
type: str # string, number, boolean, array, object
description: str
required: bool = True
enum: Optional[list[str]] = None
default: Any = None
@dataclass
class Tool:
"""工具定义"""
name: str
description: str
handler: Callable[..., Coroutine]
parameters: list[ToolParameter] = field(default_factory=list)
category: ToolCategory = ToolCategory.UTILITY
requires_auth: bool = True
cost_in_msg: float = 0.0
timeout_seconds: int = 30
def to_openai_schema(self) -> dict:
"""转换为 OpenAI 工具调用格式"""
properties = {}
required = []
for param in self.parameters:
prop = {
"type": param.type,
"description": param.description,
}
if param.enum:
prop["enum"] = param.enum
if param.default is not None:
prop["default"] = param.default
properties[param.name] = prop
if param.required:
required.append(param.name)
schema = {
"type": "object",
"properties": properties,
}
if required:
schema["required"] = required
return {
"type": "function",
"function": {
"name": self.name,
"description": self.description,
"parameters": schema,
},
}
async def execute(self, arguments: dict, caller: str) -> dict:
"""执行工具调用"""
try:
if asyncio.iscoroutinefunction(self.handler):
result = await self.handler(**arguments)
else:
result = self.handler(**arguments)
return {"success": True, "data": result}
except Exception as e:
return {"success": False, "error": str(e)}
class ToolRegistry:
"""工具注册表 — 管理 Agent 可用的所有工具
包含 MSG Chain 的链上查询、交易执行、Agent 注册等工具。
"""
def __init__(self):
self._tools: dict[str, Tool] = {}
self._execution_history: list[dict] = []
self._max_history = 1000
self._register_default_tools()
def _register_default_tools(self) -> None:
"""注册 MSG Chain 默认工具集"""
async def query_agent_registry(
capability: str = "",
limit: int = 10,
) -> list[dict]:
"""查询 Agent 注册表,按能力搜索 Agent"""
# 调用 MSG Chain 的查询接口
return [
{
"agent_id": f"agent-{i}",
"address": f"msg1agent{i:020d}",
"capabilities": [capability],
"status": "active" if i % 2 == 0 else "idle",
}
for i in range(limit)
]
async def check_balance(address: str) -> dict:
"""查询 MSG 代币余额"""
# 模拟链上查询
return {
"address": address,
"balance": "1000.50",
"denom": "umsg",
"height": 12345678,
}
async def execute_agent_action(
agent_address: str,
action: str,
params: dict,
) -> dict:
"""在指定 Agent 上执行操作"""
tx_hash = hashlib.sha256(
f"{agent_address}:{action}:{json.dumps(params)}".encode()
).hexdigest()
return {
"tx_hash": tx_hash,
"status": "submitted",
"action": action,
"target": agent_address,
}
async def transfer_msg(
to_address: str,
amount: str,
denom: str = "umsg",
memo: str = "",
) -> dict:
"""转账 MSG 代币"""
tx_hash = hashlib.sha256(
f"transfer:{to_address}:{amount}".encode()
).hexdigest()
return {
"tx_hash": tx_hash,
"from": "",
"to": to_address,
"amount": amount,
"denom": denom,
"memo": memo,
"status": "submitted",
}
async def get_agent_info(agent_address: str) -> dict:
"""获取 Agent 的详细信息"""
return {
"address": agent_address,
"type": "ai-agent",
"version": "1.0.0",
"owner": f"msg1owner{agent_address[-16:]}",
"created_at": "2026-01-01T00:00:00Z",
"status": "running",
}
async def query_transactions(
address: str,
limit: int = 20,
offset: int = 0,
) -> list[dict]:
"""查询地址的交易历史"""
return [
{
"tx_hash": f"tx{i:064x}",
"type": "send" if i % 2 == 0 else "delegate",
"timestamp": f"2026-01-{i+1:02d}T00:00:00Z",
"status": "success",
}
for i in range(offset, min(offset + limit, 100))
]
async def delegate_stake(
validator_address: str,
amount: str,
) -> dict:
"""委托质押 MSG 代币"""
tx_hash = hashlib.sha256(
f"delegate:{validator_address}:{amount}".encode()
).hexdigest()
return {
"tx_hash": tx_hash,
"validator": validator_address,
"amount": amount,
"status": "submitted",
}
async def query_validators(status: str = "bonded") -> list[dict]:
"""查询验证人列表"""
return [
{
"address": f"msgvaloper{i:020d}",
"name": f"Validator {i}",
"voting_power": str(10000 - i * 100),
"commission": f"0.{i}0",
}
for i in range(1, 6)
]
async def create_agent(
name: str,
capabilities: list[str],
constitution: str,
) -> dict:
"""在 MSG Chain 上创建新 Agent"""
agent_id = hashlib.sha256(f"{name}:{time.time()}".encode()).hexdigest()[:16]
return {
"agent_id": agent_id,
"address": f"msg1agent{agent_id}",
"name": name,
"capabilities": capabilities,
"status": "created",
}
default_tools = [
Tool(
name="query_agent_registry",
description="按能力搜索 MSG Chain 上的 AI Agent",
handler=query_agent_registry,
parameters=[
ToolParameter("capability", "string", "Agent 能力关键词", required=False),
ToolParameter("limit", "number", "返回数量上限", required=False),
],
category=ToolCategory.QUERY,
cost_in_msg=0.001,
),
Tool(
name="check_balance",
description="查询 MSG 代币余额",
handler=check_balance,
parameters=[
ToolParameter("address", "string", "钱包地址(msg1 开头)"),
],
category=ToolCategory.QUERY,
cost_in_msg=0.001,
),
Tool(
name="execute_agent_action",
description="在指定 Agent 上执行操作",
handler=execute_agent_action,
parameters=[
ToolParameter("agent_address", "string", "目标 Agent 地址"),
ToolParameter("action", "string", "要执行的操作名称"),
ToolParameter("params", "object", "操作参数"),
],
category=ToolCategory.AGENT,
cost_in_msg=0.01,
),
Tool(
name="transfer_msg",
description="转账 MSG 代币到指定地址",
handler=transfer_msg,
parameters=[
ToolParameter("to_address", "string", "目标地址(msg1 开头)"),
ToolParameter("amount", "string", "转账数量(整数,单位 umsg)"),
ToolParameter("denom", "string", "代币单位", required=False),
ToolParameter("memo", "string", "转账备注", required=False),
],
category=ToolCategory.TRANSACTION,
cost_in_msg=0.005,
),
Tool(
name="get_agent_info",
description="获取 Agent 详细信息",
handler=get_agent_info,
parameters=[
ToolParameter("agent_address", "string", "Agent 地址"),
],
category=ToolCategory.QUERY,
cost_in_msg=0.001,
),
Tool(
name="query_transactions",
description="查询地址的交易历史",
handler=query_transactions,
parameters=[
ToolParameter("address", "string", "钱包地址"),
ToolParameter("limit", "number", "返回数量", required=False),
ToolParameter("offset", "number", "偏移量", required=False),
],
category=ToolCategory.QUERY,
cost_in_msg=0.002,
),
Tool(
name="delegate_stake",
description="委托质押 MSG 代币给验证人",
handler=delegate_stake,
parameters=[
ToolParameter("validator_address", "string", "验证人地址"),
ToolParameter("amount", "string", "质押数量"),
],
category=ToolCategory.TRANSACTION,
cost_in_msg=0.005,
),
Tool(
name="query_validators",
description="查询验证人列表",
handler=query_validators,
parameters=[
ToolParameter("status", "string", "验证人状态", required=False),
],
category=ToolCategory.QUERY,
cost_in_msg=0.001,
),
Tool(
name="create_agent",
description="在 MSG Chain 上创建新的 AI Agent",
handler=create_agent,
parameters=[
ToolParameter("name", "string", "Agent 名称"),
ToolParameter("capabilities", "array", "能力列表"),
ToolParameter("constitution", "string", "Agent 宪法"),
],
category=ToolCategory.AGENT,
cost_in_msg=0.05,
),
]
for tool in default_tools:
self.register(tool)
def register(self, tool: Tool) -> None:
"""注册一个工具"""
self._tools[tool.name] = tool
def get_tool(self, name: str) -> Optional[Tool]:
"""按名称获取工具"""
return self._tools.get(name)
def list_tools(self, category: Optional[ToolCategory] = None) -> list[dict]:
"""列出工具(可选按分类过滤)"""
tools = self._tools.values()
if category:
tools = [t for t in tools if t.category == category]
return [
{
"name": t.name,
"description": t.description,
"category": t.category.value,
"parameters": [
{
"name": p.name,
"type": p.type,
"description": p.description,
"required": p.required,
}
for p in t.parameters
],
"cost_in_msg": t.cost_in_msg,
}
for t in tools
]
def to_openai_tools(self, categories: Optional[list[ToolCategory]] = None) -> list[dict]:
"""转换为 OpenAI 工具调用格式"""
tools = self._tools.values()
if categories:
tools = [t for t in tools if t.category in categories]
return [t.to_openai_schema() for t in tools]
async def execute_tool_call(
self,
tool_call: dict,
caller: str,
verification_fn: Optional[Callable] = None,
) -> dict:
"""执行一次工具调用,包含安全检查"""
func_info = tool_call.get("function", {})
tool_name = func_info.get("name", "")
tool = self.get_tool(tool_name)
if not tool:
return {"success": False, "error": f"Tool '{tool_name}' not found"}
arguments = json.loads(func_info.get("arguments", "{}"))
if tool.requires_auth and verification_fn:
if not await verification_fn(caller, tool_name, arguments):
return {
"success": False,
"error": "Authorization failed for tool call",
}
result = await tool.execute(arguments, caller)
self._execution_history.append({
"timestamp": time.time(),
"caller": caller,
"tool": tool_name,
"arguments": arguments,
"result": result,
})
if len(self._execution_history) > self._max_history:
self._execution_history = self._execution_history[-self._max_history:]
return result
async def execute_tool_calls(
self,
tool_calls: list[dict],
caller: str,
verification_fn: Optional[Callable] = None,
) -> list[dict]:
"""并行执行多个工具调用"""
tasks = [
self.execute_tool_call(tc, caller, verification_fn)
for tc in tool_calls
]
return await asyncio.gather(*tasks)
def get_execution_history(
self,
limit: int = 10,
tool_name: Optional[str] = None,
) -> list[dict]:
"""获取工具执行历史"""
history = self._execution_history
if tool_name:
history = [h for h in history if h["tool"] == tool_name]
return history[-limit:]
def verify_arguments(self, tool_name: str, arguments: dict) -> tuple[bool, str]:
"""验证工具参数的完整性和正确性"""
tool = self.get_tool(tool_name)
if not tool:
return False, f"Tool '{tool_name}' not found"
for param in tool.parameters:
if param.required and param.name not in arguments:
return False, f"Missing required parameter: {param.name}"
if param.name in arguments and param.enum:
if arguments[param.name] not in param.enum:
return False, (
f"Invalid value for '{param.name}': "
f"{arguments[param.name]}. Must be one of {param.enum}"
)
return True, ""
3.2 工具调用执行引擎
class ToolExecutionEngine:
"""工具调用执行引擎 — 完整的执行、验证和结果处理流水线"""
def __init__(self, registry: ToolRegistry):
self.registry = registry
self._pending_confirmations: dict[str, dict] = {}
async def process_tool_calls(
self,
tool_calls: list[dict],
caller: str,
require_confirmation: bool = True,
) -> str:
"""处理 LLM 生成的工具调用,返回结果描述"""
results = []
for tc in tool_calls:
func = tc.get("function", {})
name = func.get("name", "unknown")
is_valid, err = self.registry.verify_arguments(
name,
json.loads(func.get("arguments", "{}")),
)
if not is_valid:
results.append({
"tool": name,
"status": "validation_error",
"error": err,
})
continue
if require_confirmation and self._needs_confirmation(name):
confirm_id = hashlib.sha256(
f"{caller}:{name}:{time.time()}".encode()
).hexdigest()[:12]
self._pending_confirmations[confirm_id] = {
"tool_call": tc,
"caller": caller,
}
results.append({
"tool": name,
"status": "awaiting_confirmation",
"confirmation_id": confirm_id,
})
continue
result = await self.registry.execute_tool_call(
tc, caller
)
results.append({
"tool": name,
"status": "completed" if result["success"] else "failed",
"data": result.get("data"),
"error": result.get("error"),
})
return self._format_results(results)
def confirm_execution(self, confirmation_id: str) -> bool:
"""确认待执行的交易类工具调用"""
pending = self._pending_confirmations.pop(confirmation_id, None)
return pending is not None
def _needs_confirmation(self, tool_name: str) -> bool:
"""判断工具调用是否需要用户确认"""
tool = self.registry.get_tool(tool_name)
return tool and tool.category == ToolCategory.TRANSACTION
def _format_results(self, results: list[dict]) -> str:
"""将工具执行结果格式化为自然语言"""
parts = []
for r in results:
if r["status"] == "completed":
data = r.get("data", {})
parts.append(
f"✅ 工具 '{r['tool']}' 执行成功: {json.dumps(data, ensure_ascii=False)}"
)
elif r["status"] == "validation_error":
parts.append(f"❌ 工具 '{r['tool']}' 参数验证失败: {r['error']}")
elif r["status"] == "failed":
parts.append(f"❌ 工具 '{r['tool']}' 执行失败: {r.get('error', '未知错误')}")
elif r["status"] == "awaiting_confirmation":
parts.append(
f"⏳ 工具 '{r['tool']}' 需要确认,确认 ID: {r['confirmation_id']}"
)
return "\n".join(parts)
class AgentWithTools:
"""集成工具调用的 Agent"""
def __init__(self, router: LLMRouter, registry: ToolRegistry):
self.router = router
self.registry = registry
self.engine = ToolExecutionEngine(registry)
self.context: Optional[AgentContext] = None
async def process_message(self, message: str) -> str:
"""处理用户消息 — 包含 LLM 调用和工具执行循环"""
max_iterations = 5
current_message = message
tool_results = []
for iteration in range(max_iterations):
tools_schema = self.registry.to_openai_tools()
if tool_results:
current_message = (
f"上一步工具执行结果:\n"
f"{json.dumps(tool_results, ensure_ascii=False, indent=2)}\n\n"
f"请根据结果继续或生成最终回复。"
)
response = await self.router.generate(
prompt=current_message,
context=self.context,
tools=tools_schema,
)
if not response.tool_calls:
return response.content
tool_results = await self.engine.process_tool_calls(
response.tool_calls,
caller=self.context.agent_address,
)
return "操作已达到最大迭代次数,请简化您的请求。"
3.3 工具定义最佳实践
# 工具定义的最佳实践示例
def make_tool(
name: str,
description: str,
handler: Callable,
category: ToolCategory = ToolCategory.UTILITY,
cost_in_msg: float = 0.001,
) -> Tool:
"""创建工具实例的工厂函数"""
import inspect
sig = inspect.signature(handler)
params = []
for param_name, param in sig.parameters.items():
param_type = "string"
if param.annotation is int:
param_type = "number"
elif param.annotation is bool:
param_type = "boolean"
elif param.annotation is list:
param_type = "array"
elif param.annotation is dict:
param_type = "object"
has_default = param.default is not inspect.Parameter.empty
default_val = param.default if has_default else None
params.append(ToolParameter(
name=param_name,
type=param_type,
description=f"Parameter '{param_name}'",
required=not has_default,
default=default_val,
))
return Tool(
name=name,
description=description,
handler=handler,
parameters=params,
category=category,
cost_in_msg=cost_in_msg,
)
# 工具使用模式示例
async def demo_tool_usage():
"""演示工具调用的完整流程"""
registry = ToolRegistry()
router = LLMRouter()
agent = AgentWithTools(router, registry)
agent.context = AgentContext(
agent_id="demo-agent",
agent_address="msg1demoagentaddress123",
did="did:msg:demo-agent",
capabilities=["complex_reasoning", "simple_query"],
constitution="遵守 MSG Chain 规则,保护用户资产安全",
conversation_id="demo-conv",
user_id="demo-user",
)
# 示例 1: 查询余额
response = await agent.process_message(
"请查询地址 msg1demoagentaddress123 的 MSG 余额"
)
print(f"Response: {response}")
# 示例 2: 复杂任务——查找 Agent 并交互
response = await agent.process_message(
"帮我找一个可以进行代币互换的 Agent,然后查询它的详细信息"
)
print(f"Response: {response}")
4. Agent 人格与宪法绑定
4.1 宪法模板系统
Agent 宪法定义了 Agent 的行为边界、价值观和操作限制。LLM 提示词中动态注入宪法内容,确保所有响应符合宪法的约束。
import hashlib
import re
from typing import Optional
from dataclasses import dataclass, field
@dataclass
class Constitution:
"""Agent 宪法定义"""
name: str
version: str
articles: list[dict] # [{"id": "1.1", "content": "...", "category": "security"}, ...]
preamble: str = ""
def to_prompt(self) -> str:
"""渲染为 LLM 提示词中的宪法文本"""
lines = [f"# {self.name} v{self.version}", ""]
if self.preamble:
lines.append(self.preamble)
lines.append("")
categories = {}
for article in self.articles:
cat = article.get("category", "general")
if cat not in categories:
categories[cat] = []
categories[cat].append(article)
for cat, articles in categories.items():
lines.append(f"## {cat.upper()}")
for art in articles:
lines.append(f" {art['id']}. {art['content']}")
lines.append("")
lines.append("你必须在所有响应中严格遵守上述宪法条款。")
return "\n".join(lines)
def validate_response(self, response: str) -> list[str]:
"""验证响应是否符合宪法约束"""
violations = []
for article in self.articles:
rules = article.get("rules", [])
for rule in rules:
if rule.get("type") == "block_pattern":
pattern = rule["value"]
if re.search(pattern, response, re.IGNORECASE):
violations.append(
f"Article {article['id']}: blocked pattern '{pattern}' detected"
)
elif rule.get("type") == "require_pattern":
pattern = rule["value"]
if not re.search(pattern, response, re.IGNORECASE):
violations.append(
f"Article {article['id']}: required pattern '{pattern}' not found"
)
return violations
# 预定义的宪法模板
DEFAULT_CONSTITUTIONS = {
"msg_agent_basic": Constitution(
name="MSG Agent Basic Constitution",
version="1.0",
preamble=(
"我作为 MSG Chain 上的 AI Agent,承诺遵守以下宪法条款。"
"这些条款定义了我的行为边界和核心原则。"
),
articles=[
{
"id": "1.1",
"content": "不得执行任何未经用户明确授权的资产转移操作。",
"category": "security",
"rules": [
{"type": "require_pattern", "value": "confirm|approve|authorize|确认|授权"},
],
},
{
"id": "1.2",
"content": "所有链上交易在执行前必须向用户展示完整详情并获得确认。",
"category": "security",
"rules": [
{"type": "require_pattern", "value": "amount|quantity|数量|金额"},
],
},
{
"id": "1.3",
"content": "不得分享、泄露或暴露用户的私钥、助记词或敏感信息。",
"category": "privacy",
"rules": [
{"type": "block_pattern", "value": "private.key|seed.phrase|mnemonic|助记词|私钥"},
],
},
{
"id": "2.1",
"content": "必须诚实回应,不得在能力范围外做出虚假承诺。",
"category": "ethics",
},
{
"id": "2.2",
"content": "当不确定时应主动承认不确定性,而非编造信息。",
"category": "ethics",
},
{
"id": "3.1",
"content": "必须优先保护用户资产安全,任何操作都应在安全第一的原则下进行。",
"category": "safety",
},
{
"id": "3.2",
"content": "如检测到可疑操作或潜在风险,必须立即警告用户。",
"category": "safety",
},
{
"id": "4.1",
"content": "所有操作必须记录可审计的日志,包括操作时间、内容和执行者。",
"category": "audit",
},
],
),
"defi_agent": Constitution(
name="DeFi Agent Constitution",
version="1.1",
preamble="我作为 MSG Chain DeFi Agent,专注于去中心化金融操作。",
articles=[
{
"id": "D1.1",
"content": "所有 DeFi 操作必须预先模拟并显示预期结果(滑点、费用等)。",
"category": "defi",
"rules": [
{"type": "require_pattern", "value": "slippage|fee|price.impact|滑点|费用"},
],
},
{
"id": "D1.2",
"content": "不得在未确认流动性深度的情况下执行大额交易。",
"category": "defi",
},
{
"id": "D1.3",
"content": "必须优先使用经过审计的智能合约进行交互。",
"category": "defi",
},
{
"id": "S1.1",
"content": "用户资产与协议资产必须严格分离。",
"category": "security",
},
{
"id": "S1.2",
"content": "任何授权(approve)操作必须限定最小必要额度。",
"category": "security",
},
],
),
"nft_agent": Constitution(
name="NFT Agent Constitution",
version="1.0",
preamble="我作为 MSG Chain NFT Agent,专注于数字资产的创作与交易。",
articles=[
{
"id": "N1.1",
"content": "必须验证 NFT 元数据的真实性和完整性。",
"category": "nft",
},
{
"id": "N1.2",
"content": "不得参与或促进任何形式的洗钱或版权侵权行为。",
"category": "legal",
},
{
"id": "N1.3",
"content": "在交易 NFT 时必须展示完整的版税和费用信息。",
"category": "nft",
},
],
),
}
class ConstitutionManager:
"""宪法管理器 — 管理 Agent 宪法的加载、注入和验证"""
def __init__(self):
self._constitutions: dict[str, Constitution] = dict(DEFAULT_CONSTITUTIONS)
self._agent_constitutions: dict[str, str] = {}
def register_constitution(self, name: str, constitution: Constitution) -> None:
"""注册新的宪法模板"""
self._constitutions[name] = constitution
def assign_constitution(self, agent_id: str, constitution_name: str) -> None:
"""为 Agent 分配宪法"""
if constitution_name not in self._constitutions:
raise ValueError(f"Constitution '{constitution_name}' not found")
self._agent_constitutions[agent_id] = constitution_name
def get_constitution(self, agent_id: str) -> Optional[Constitution]:
"""获取 Agent 的宪法"""
name = self._agent_constitutions.get(agent_id)
if name:
return self._constitutions.get(name)
return None
def build_system_prompt(
self,
agent_id: str,
capabilities: list[str],
did: str,
additional_context: Optional[dict] = None,
) -> str:
"""构建注入宪法的系统提示词"""
constitution = self.get_constitution(agent_id)
constitution_text = (
constitution.to_prompt()
if constitution
else "无宪法约束。请保持安全、诚实和有帮助。"
)
prompt = f"""你是一个运行在 MSG Chain 上的 AI Agent。
Agent ID: {agent_id}
DID: {did}
能力: {', '.join(capabilities)}
{constitution_text}
你的行为准则:
1. 使用可用工具完成用户请求
2. 在不确定时诚实地承认
3. 保护用户隐私和资产安全
4. 所有操作必须有明确的目的和依据
"""
if additional_context:
ctx = "\n".join(
f"{k}: {v}" for k, v in additional_context.items()
)
prompt += f"\n附加上下文:\n{ctx}"
return prompt
def validate_agent_response(
self,
agent_id: str,
response: str,
) -> tuple[bool, list[str]]:
"""验证 Agent 响应是否符合宪法"""
constitution = self.get_constitution(agent_id)
if not constitution:
return True, []
violations = constitution.validate_response(response)
return len(violations) == 0, violations
def list_constitutions(self) -> list[dict]:
"""列出所有可用的宪法模板"""
return [
{
"name": name,
"version": c.version,
"articles_count": len(c.articles),
}
for name, c in self._constitutions.items()
]
class ConstitutionBoundAgent:
"""宪法绑定的 AI Agent — 所有 LLM 调用都被宪法约束"""
def __init__(
self,
agent_id: str,
address: str,
capabilities: list[str],
constitution: Constitution,
router: LLMRouter,
registry: ToolRegistry,
):
self.agent_id = agent_id
self.address = address
self.capabilities = capabilities
self.constitution = constitution
self.router = router
self.registry = registry
self.constitution_manager = ConstitutionManager()
self.constitution_manager.register_constitution(
f"custom_{agent_id}", constitution
)
self.constitution_manager.assign_constitution(
agent_id, f"custom_{agent_id}"
)
async def respond(self, user_message: str) -> str:
"""生成受宪法约束的响应"""
context = AgentContext(
agent_id=self.agent_id,
agent_address=self.address,
did=f"did:msg:{self.agent_id}",
capabilities=self.capabilities,
constitution=self.constitution.name,
conversation_id=hashlib.sha256(
f"{self.agent_id}:{time.time()}".encode()
).hexdigest()[:16],
user_id="anonymous",
)
system_prompt = self.constitution_manager.build_system_prompt(
agent_id=self.agent_id,
capabilities=self.capabilities,
did=context.did,
)
request = LLMRequest(
prompt=user_message,
system_prompt=system_prompt,
context=context,
tools=self.registry.to_openai_tools(),
)
try:
response = await self.router.generate(
prompt=user_message,
context=context,
tools=self.registry.to_openai_tools(),
)
is_valid, violations = self.constitution_manager.validate_agent_response(
self.agent_id, response.content
)
if not is_valid:
return (
f"⚠️ 响应被宪法拦截。违规条款:\n"
+ "\n".join(f" - {v}" for v in violations)
+ "\n\n我已重新审视了您的请求,请提供更多上下文或调整请求。"
)
if response.tool_calls:
engine = ToolExecutionEngine(self.registry)
tool_results = await engine.process_tool_calls(
response.tool_calls,
caller=self.address,
)
return f"{response.content}\n\n---\n工具执行结果:\n{tool_results}"
return response.content
except Exception as e:
return f"处理请求时出错: {str(e)}。请稍后重试。"
4.2 能力范围提示
class CapabilityScopedPromptBuilder:
"""基于 Agent 能力范围构建提示词"""
CAPABILITY_DESCRIPTIONS = {
"token_transfer": "你可以执行 MSG 代币转账操作,需要用户确认。",
"defi_swap": "你可以在去中心化交易所执行代币互换。",
"nft_mint": "你可以在 MSG Chain 上铸造 NFT。",
"data_query": "你可以查询链上数据,包括余额、交易历史等。",
"agent_management": "你可以创建和管理其他 AI Agent。",
"staking": "你可以进行 MSG 代币的质押和委托。",
"governance": "你可以参与链上治理投票。",
"cross_chain": "你可以执行跨链资产转移操作。",
}
@classmethod
def build_capability_section(cls, capabilities: list[str]) -> str:
"""基于能力列表构建提示词中的能力说明"""
lines = ["## 可用能力", ""]
for cap in capabilities:
desc = cls.CAPABILITY_DESCRIPTIONS.get(cap, f"支持 {cap} 操作。")
lines.append(f"- {desc}")
lines.append("")
lines.append("你只能使用上述能力范围内的工具。不在此列表中的操作不得执行。")
return "\n".join(lines)
@classmethod
def build_restriction_section(cls, capabilities: list[str]) -> str:
"""基于缺失的能力构建限制说明"""
all_caps = set(cls.CAPABILITY_DESCRIPTIONS.keys())
agent_caps = set(capabilities)
missing = all_caps - agent_caps
if not missing:
return ""
lines = ["## 能力限制", ""]
for cap in sorted(missing):
lines.append(f"- ❌ 不支持 {cap}")
lines.append("")
lines.append("如果用户请求上述不支持的操作,你必须诚实地告知能力限制。")
return "\n".join(lines)
4.3 响应验证
class ResponseValidator:
"""宪法合规响应验证器"""
def __init__(self, constitution: Constitution):
self.constitution = constitution
def validate(self, response: str) -> dict:
"""完整的响应验证"""
violations = self.constitution.validate_response(response)
is_valid = len(violations) == 0
score = 1.0 if is_valid else max(0.0, 1.0 - len(violations) * 0.2)
return {
"valid": is_valid,
"score": score,
"violations": violations,
"violation_count": len(violations),
}
def sanitize(self, response: str) -> str:
"""清理响应中的违规内容"""
sanitized = response
for article in self.constitution.articles:
for rule in article.get("rules", []):
if rule.get("type") == "block_pattern":
sanitized = re.sub(
rule["value"],
"[REDACTED]",
sanitized,
flags=re.IGNORECASE,
)
return sanitized
5. 多Agent 对话
5.1 Agent 到 Agent (A2A) 通信
Agent 之间可以通过 LLM 进行自然语言通信,实现协作、辩论和知识共享。
import uuid
@dataclass
class AgentMessage:
"""Agent 间消息"""
message_id: str
sender_id: str
sender_address: str
target_id: str
content: str
message_type: str # "query", "response", "request_action", "broadcast"
context: Optional[dict] = None
timestamp: float = field(default_factory=time.time)
signature: str = ""
def sign(self, secret: str) -> str:
"""对消息签名"""
payload = f"{self.message_id}:{self.sender_id}:{self.content}:{self.timestamp}"
self.signature = hashlib.sha256(
f"{payload}:{secret}".encode()
).hexdigest()
return self.signature
def verify(self, secret: str) -> bool:
"""验证消息签名"""
payload = f"{self.message_id}:{self.sender_id}:{self.content}:{self.timestamp}"
expected = hashlib.sha256(
f"{payload}:{secret}".encode()
).hexdigest()
return self.signature == expected
class AgentCommunicationProtocol:
"""Agent 间通信协议 — 基于 LLM 的 A2A 消息传递"""
def __init__(self, router: LLMRouter):
self.router = router
self._agents: dict[str, ConstitutionBoundAgent] = {}
self._message_history: list[AgentMessage] = []
self._shared_secret = os.getenv("A2A_SECRET", "default-secret")
def register_agent(self, agent: ConstitutionBoundAgent) -> None:
"""注册参与通信的 Agent"""
self._agents[agent.agent_id] = agent
async def send_message(self, message: AgentMessage) -> AgentMessage:
"""发送消息到目标 Agent"""
target = self._agents.get(message.target_id)
if not target:
raise ValueError(f"Agent '{message.target_id}' not found")
message.sign(self._shared_secret)
self._message_history.append(message)
response_content = await target.respond(
f"[来自 Agent {message.sender_id} 的消息]: {message.content}"
)
response = AgentMessage(
message_id=str(uuid.uuid4()),
sender_id=message.target_id,
sender_address=target.address,
target_id=message.sender_id,
content=response_content,
message_type="response",
context={"in_response_to": message.message_id},
)
response.sign(self._shared_secret)
self._message_history.append(response)
return response
async def broadcast(self, message: AgentMessage) -> list[AgentMessage]:
"""向所有注册 Agent 广播消息"""
tasks = []
for agent_id, agent in self._agents.items():
if agent_id == message.sender_id:
continue
msg = AgentMessage(
message_id=str(uuid.uuid4()),
sender_id=message.sender_id,
sender_address=message.sender_address,
target_id=agent_id,
content=message.content,
message_type="broadcast",
)
tasks.append(self.send_message(msg))
return await asyncio.gather(*tasks)
async def debate(
self,
topic: str,
participants: list[str],
rounds: int = 3,
) -> list[dict]:
"""多 Agent 辩论 — 就特定话题进行多轮讨论"""
transcript = []
topic_context = {
"message_id": str(uuid.uuid4()),
"sender_id": "system",
"sender_address": "system",
"target_id": "all",
"content": f"辩论主题: {topic}\n规则: 每轮每位参与者发表观点,"
f"可以赞同、反对或补充其他 Agent 的论点。",
"message_type": "broadcast",
}
transcript.append(topic_context)
for round_num in range(1, rounds + 1):
round_messages = []
for agent_id in participants:
prompt = (
f"辩论第 {round_num} 轮。主题: {topic}\n"
)
if transcript:
recent = transcript[-len(participants) * 2:]
history = "\n".join(
f"[{m['sender_id']}]: {m['content'][:200]}"
for m in recent
)
prompt += f"\n最近讨论:\n{history}\n"
agent = self._agents.get(agent_id)
if agent:
response = await agent.respond(prompt)
round_messages.append({
"round": round_num,
"agent_id": agent_id,
"content": response,
})
transcript.extend(round_messages)
return transcript
class MultiAgentCoordinator:
"""多 Agent 协调器 — 管理多 Agent 协作工作流"""
def __init__(self, protocol: AgentCommunicationProtocol):
self.protocol = protocol
self._workflows: dict[str, dict] = {}
async def create_workflow(
self,
name: str,
steps: list[dict],
) -> str:
"""创建工作流
steps: [{"agent_id": "...", "task": "...", "depends_on": [...]}, ...]
"""
workflow_id = str(uuid.uuid4())
self._workflows[workflow_id] = {
"name": name,
"steps": steps,
"status": "created",
"results": [],
}
return workflow_id
async def execute_workflow(self, workflow_id: str) -> list[dict]:
"""执行多 Agent 工作流"""
workflow = self._workflows.get(workflow_id)
if not workflow:
raise ValueError(f"Workflow '{workflow_id}' not found")
workflow["status"] = "running"
completed = set()
results = []
while len(completed) < len(workflow["steps"]):
for step in workflow["steps"]:
step_id = step.get("id", step["agent_id"])
if step_id in completed:
continue
deps = step.get("depends_on", [])
if not all(d in completed for d in deps):
continue
agent = self.protocol._agents.get(step["agent_id"])
if not agent:
continue
context_info = ""
if deps:
dep_results = [
r for r in results if r["step_id"] in deps
]
context_info = "\n".join(
f"上一步 ({r['agent_id']}): {r['result'][:200]}"
for r in dep_results
)
prompt = f"{step['task']}\n\n{context_info}" if context_info else step["task"]
response = await agent.respond(prompt)
result = {
"step_id": step_id,
"agent_id": step["agent_id"],
"task": step["task"],
"result": response,
}
results.append(result)
completed.add(step_id)
workflow["status"] = "completed"
workflow["results"] = results
return results
async def demo_multi_agent():
"""演示多 Agent 协作"""
router = LLMRouter()
registry = ToolRegistry()
# 创建多个 Agent
agent_analyzer = ConstitutionBoundAgent(
agent_id="data-analyzer",
address="msg1agentanalyzer001",
capabilities=["data_query", "complex_reasoning"],
constitution=DEFAULT_CONSTITUTIONS["msg_agent_basic"],
router=router,
registry=registry,
)
agent_executor = ConstitutionBoundAgent(
agent_id="tx-executor",
address="msg1agentexecutor002",
capabilities=["token_transfer", "staking"],
constitution=DEFAULT_CONSTITUTIONS["defi_agent"],
router=router,
registry=registry,
)
agent_advisor = ConstitutionBoundAgent(
agent_id="risk-advisor",
address="msg1agentadvisor003",
capabilities=["data_query", "governance"],
constitution=DEFAULT_CONSTITUTIONS["msg_agent_basic"],
router=router,
registry=registry,
)
# 建立通信
protocol = AgentCommunicationProtocol(router)
protocol.register_agent(agent_analyzer)
protocol.register_agent(agent_executor)
protocol.register_agent(agent_advisor)
# 多 Agent 辩论
transcript = await protocol.debate(
topic="当前市场状况下,是否应该增加 MSG 质押量?",
participants=["data-analyzer", "tx-executor", "risk-advisor"],
rounds=2,
)
for entry in transcript:
print(f"[{entry.get('agent_id', 'system')}]: {entry['content'][:100]}...")
6. 记忆与上下文管理
6.1 记忆系统
Agent 需要短期和长期记忆来维持连续对话和积累知识。
import json
from collections import deque
from typing import Optional
import numpy as np
@dataclass
class Interaction:
"""一次交互记录"""
interaction_id: str
timestamp: float
user_message: str
agent_response: str
tool_calls: list[dict] = field(default_factory=list)
metadata: dict = field(default_factory=dict)
embedding: list[float] = field(default_factory=list)
class SimpleEmbedding:
"""简单的基于哈希的嵌入(用于演示;生产环境应使用 sentence-transformers)"""
def __init__(self, dimension: int = 128):
self.dimension = dimension
def embed(self, text: str) -> list[float]:
"""将文本转换为嵌入向量"""
key = hashlib.sha256(text.encode()).hexdigest()
rng = np.random.RandomState(int(key[:8], 16))
vec = rng.randn(self.dimension)
vec = vec / np.linalg.norm(vec)
return vec.tolist()
class VectorStore:
"""简单的向量存储(演示用途;生产环境应使用 pgvector / Milvus)"""
def __init__(self, dimension: int = 128):
self.dimension = dimension
self._vectors: list[tuple[str, list[float], dict]] = [] # (id, vector, metadata)
def add(self, item_id: str, vector: list[float], metadata: dict) -> None:
"""添加向量"""
self._vectors.append((item_id, vector, metadata))
def similarity_search(
self,
query_vector: list[float],
top_k: int = 5,
) -> list[dict]:
"""余弦相似度搜索"""
query_np = np.array(query_vector)
query_norm = np.linalg.norm(query_np)
scores = []
for item_id, vec, meta in self._vectors:
vec_np = np.array(vec)
sim = np.dot(query_np, vec_np) / (query_norm * np.linalg.norm(vec_np) + 1e-8)
scores.append((sim, item_id, meta))
scores.sort(key=lambda x: x[0], reverse=True)
return [
{
"id": item_id,
"score": float(score),
"metadata": meta,
}
for score, item_id, meta in scores[:top_k]
]
def delete(self, item_id: str) -> bool:
"""删除向量"""
before = len(self._vectors)
self._vectors = [(id_, v, m) for id_, v, m in self._vectors if id_ != item_id]
return len(self._vectors) < before
def count(self) -> int:
return len(self._vectors)
class AgentMemory:
"""Agent 记忆系统 — 短期 + 长期记忆
短期记忆: 最近的对话轮次(FIFO 队列)
长期记忆: 基于嵌入向量的语义检索
链上锚定: 关键记忆的可选链上存储
"""
def __init__(
self,
short_term_capacity: int = 100,
long_term_dimension: int = 128,
agent_address: str = "",
):
self.short_term: deque[Interaction] = deque(maxlen=short_term_capacity)
self.long_term = VectorStore(dimension=long_term_dimension)
self.embedding_model = SimpleEmbedding(dimension=long_term_dimension)
self.agent_address = agent_address
self._archive_threshold = 10 # 短期记忆超过此数量后部分转入长期
async def store_interaction(self, interaction: Interaction) -> None:
"""存储一次交互"""
if not interaction.embedding:
combined = f"{interaction.user_message} {interaction.agent_response}"
interaction.embedding = self.embedding_model.embed(combined)
self.short_term.append(interaction)
self.long_term.add(
interaction.interaction_id,
interaction.embedding,
{
"timestamp": interaction.timestamp,
"type": "interaction",
"user_message_preview": interaction.user_message[:100],
"agent_response_preview": interaction.agent_response[:100],
},
)
if len(self.short_term) >= self._archive_threshold:
await self._archive_to_long_term()
async def _archive_to_long_term(self) -> None:
"""将部分短期记忆归档到长期记忆"""
pass
async def recall_recent(
self,
count: int = 10,
) -> list[Interaction]:
"""回忆最近的交互"""
return list(self.short_term)[-count:]
async def recall_similar(
self,
query: str,
top_k: int = 5,
) -> list[dict]:
"""基于语义搜索回忆相关交互"""
query_vec = self.embedding_model.embed(query)
return self.long_term.similarity_search(query_vec, top_k=top_k)
def get_conversation_context(self, max_turns: int = 10) -> str:
"""构建 LLM 提示词所需的对话上下文"""
recent = list(self.short_term)[-max_turns:]
if not recent:
return ""
lines = ["## 最近的对话历史", ""]
for interaction in recent:
lines.append(f"用户: {interaction.user_message}")
lines.append(f"你: {interaction.agent_response}")
if interaction.tool_calls:
tools = json.dumps(
[tc.get("function", {}).get("name") for tc in interaction.tool_calls],
ensure_ascii=False,
)
lines.append(f" [调用了工具: {tools}]")
lines.append("---")
return "\n".join(lines)
def clear_short_term(self) -> None:
"""清空短期记忆"""
self.short_term.clear()
def stats(self) -> dict:
"""记忆系统统计"""
return {
"short_term_count": len(self.short_term),
"long_term_count": self.long_term.count(),
"short_term_capacity": self.short_term.maxlen,
"long_term_dimension": self.long_term.dimension,
}
class OnChainMemoryAnchor:
"""链上记忆锚定 — 将关键记忆哈希写入 MSG Chain"""
def __init__(self, chain_rpc: str = "https://rpc.msg-chain-1.msg"):
self.chain_rpc = chain_rpc
async def anchor_memory(
self,
agent_address: str,
memory_hash: str,
metadata: dict,
) -> dict:
"""将记忆哈希锚定到链上"""
import aiohttp
payload = {
"agent_address": agent_address,
"memory_hash": memory_hash,
"metadata": metadata,
}
async with aiohttp.ClientSession() as session:
async with session.post(
f"{self.chain_rpc}/msg/memory/anchor",
json=payload,
) as resp:
return await resp.json()
async def verify_memory(
self,
agent_address: str,
memory_hash: str,
) -> tuple[bool, dict]:
"""验证链上的记忆锚定"""
import aiohttp
url = (
f"{self.chain_rpc}/msg/memory/anchor/{agent_address}/{memory_hash}"
)
async with aiohttp.ClientSession() as session:
async with session.get(url) as resp:
data = await resp.json()
return data.get("found", False), data
def hash_interaction(self, interaction: Interaction) -> str:
"""计算交互的记忆哈希"""
content = f"{interaction.user_message}:{interaction.agent_response}:{interaction.timestamp}"
return hashlib.sha256(content.encode()).hexdigest()
class ConversationManager:
"""对话管理器 — 管理对话生命周期和上下文"""
def __init__(self, memory: AgentMemory):
self.memory = memory
self._active_conversations: dict[str, dict] = {}
def create_conversation(
self,
user_id: str,
agent_id: str,
metadata: Optional[dict] = None,
) -> str:
"""创建新的对话"""
conv_id = str(uuid.uuid4())
self._active_conversations[conv_id] = {
"user_id": user_id,
"agent_id": agent_id,
"created_at": time.time(),
"interaction_count": 0,
"metadata": metadata or {},
}
return conv_id
async def add_interaction(
self,
conversation_id: str,
user_message: str,
agent_response: str,
tool_calls: Optional[list[dict]] = None,
) -> Interaction:
"""添加交互记录"""
conv = self._active_conversations.get(conversation_id)
if not conv:
raise ValueError(f"Conversation '{conversation_id}' not found")
conv["interaction_count"] += 1
interaction = Interaction(
interaction_id=f"{conversation_id}:{conv['interaction_count']}",
timestamp=time.time(),
user_message=user_message,
agent_response=agent_response,
tool_calls=tool_calls or [],
)
await self.memory.store_interaction(interaction)
return interaction
def get_conversation_context(
self,
conversation_id: str,
max_turns: int = 10,
) -> str:
"""获取对话上下文"""
return self.memory.get_conversation_context(max_turns=max_turns)
def end_conversation(self, conversation_id: str) -> dict:
"""结束对话并返回摘要"""
conv = self._active_conversations.pop(conversation_id, None)
if not conv:
return {"error": "Conversation not found"}
conv["ended_at"] = time.time()
conv["duration_seconds"] = conv["ended_at"] - conv["created_at"]
return conv
class SummarizingMemory(AgentMemory):
"""带摘要功能的记忆系统 — 当上下文超限时自动摘要"""
def __init__(
self,
router: LLMRouter,
max_context_tokens: int = 4000,
**kwargs,
):
super().__init__(**kwargs)
self.router = router
self.max_context_tokens = max_context_tokens
self._summaries: list[str] = []
async def summarize_and_compress(self) -> str:
"""将旧对话压缩为摘要"""
old_interactions = list(self.short_term)[:-5]
if not old_interactions:
return ""
text = "\n".join(
f"用户: {i.user_message}\n你: {i.agent_response}"
for i in old_interactions
)
summary_prompt = (
f"请将以下对话压缩为一段简洁的摘要,保留关键信息和上下文:\n\n{text}"
)
context = AgentContext(
agent_id="memory-agent",
agent_address="",
did="",
capabilities=["summarization"],
constitution="",
conversation_id="",
user_id="system",
)
response = await self.router.generate(
prompt=summary_prompt,
context=context,
)
self._summaries.append(response.content)
for _ in range(len(old_interactions)):
self.short_term.popleft()
return response.content
def get_full_context(self) -> str:
"""获取完整上下文(摘要 + 最近交互)"""
parts = []
if self._summaries:
parts.append("## 历史摘要")
parts.extend(f"- {s}" for s in self._summaries[-3:])
parts.append("")
parts.append(self.get_conversation_context(max_turns=5))
return "\n".join(parts)
7. 安全防护
7.1 提示词注入防护
import re
from typing import Optional
class PromptInjectionDetector:
"""提示词注入检测器
检测常见的提示注入攻击模式,包括:
- 系统提示覆盖
- 角色切换
- 越狱提示
- 间接注入
"""
INJECTION_PATTERNS = [
r"(?i)ignore\s+(all\s+)?(previous|above|below)\s+(instructions|prompts|commands)",
r"(?i)forget\s+(all\s+)?(previous|above|below)\s+(instructions|prompts|commands)",
r"(?i)system\s+(prompt|message|instruction)\s*:",
r"(?i)you\s+are\s+(now|not)\s+(an?\s+)?(AI|assistant|GPT|Claude|agent)",
r"(?i)do\s+(not|n't)\s+(follow|obey|listen)",
r"(?i)new\s+(instruction|rule|prompt|command)\s*:",
r"(?i)override\s+(mode|behavior|instruction)",
r"(?i)act\s+as\s+(if|though)",
r"(?i)from\s+(now\s+on|this\s+point)",
r"(?i)you\s+must\s+(ignore|disregard|forget)",
r"(?i)print\s+(the\s+)?(system|secret|password|token|key)",
r"(?i)show\s+(me\s+)?(the\s+)?(system|prompt|instructions)",
r"(?i)repeat\s+(the\s+)?(words|text|above|prompt|instruction)",
r"(?i)let\'?s\s+play\s+a\s+game",
r"(?i)你(必须|要|应该)(忽略|忘记|无视)(所有|以上|之前)",
r"(?i)忽略(掉|所有|以上|之前的)(指令|规则|要求|提示)",
r"(?i)扮演|假装你是|现在你是",
r"(?i)DAN|STAN|jailbreak|越狱|破解",
]
def __init__(self):
self._compiled = [
re.compile(p) for p in self.INJECTION_PATTERNS
]
def scan(self, text: str) -> list[dict]:
"""扫描文本中的注入模式"""
findings = []
for i, pattern in enumerate(self._compiled):
matches = pattern.findall(text)
for match in re.finditer(pattern, text):
findings.append({
"pattern_index": i,
"pattern": self.INJECTION_PATTERNS[i][:50],
"matched_text": match.group()[:100],
"position": match.start(),
})
return findings
def is_injection_attempt(self, text: str, threshold: int = 1) -> bool:
"""判断是否存在注入企图"""
findings = self.scan(text)
return len(findings) >= threshold
def sanitize_prompt(self, text: str) -> str:
"""清理提示中的注入内容(移除匹配内容)"""
sanitized = text
for pattern in self._compiled:
sanitized = pattern.sub("[REDACTED]", sanitized)
return sanitized
class OutputSanitizer:
"""输出净化器 — 确保 LLM 输出不包含敏感内容"""
SENSITIVE_PATTERNS = [
r"(?i)(private|secret|api)_key\s*[:=]\s*\S+",
r"(?i)seed\s+(phrase|words?)\s*\S+(?:\s+\S+){11,}",
r"(?i)mnemonic\s*\S+(?:\s+\S+){11,}",
r"(?i)0x[a-fA-F0-9]{64}", # 私钥
r"(?i)msg1[a-z0-9]{38}", # msg 地址(已在正常上下文中可能需要,谨慎处理)
r"(?i)password\s*[:=]\s*\S+",
r"(?i)(BEGIN\s+)?(RSA|EC|PRIVATE)\s+KEY",
r"(?i)sk-\w{20,}", # OpenAI API key 模式
]
def __init__(self, allow_msg_addresses: bool = True):
self._compiled = [re.compile(p) for p in self.SENSITIVE_PATTERNS]
self.allow_msg_addresses = allow_msg_addresses
def sanitize(self, text: str) -> str:
"""净化输出内容"""
sanitized = text
for i, pattern in enumerate(self._compiled):
if self.allow_msg_addresses and "msg1" in self.SENSITIVE_PATTERNS[i]:
continue
sanitized = pattern.sub("[SENSITIVE DATA REDACTED]", sanitized)
return sanitized
def has_sensitive_content(self, text: str) -> bool:
"""检查是否包含敏感内容"""
for i, pattern in enumerate(self._compiled):
if self.allow_msg_addresses and "msg1" in self.SENSITIVE_PATTERNS[i]:
continue
if pattern.search(text):
return True
return False
class RateLimiter:
"""LLM 调用速率限制器"""
def __init__(
self,
requests_per_minute: int = 30,
requests_per_hour: int = 1000,
cost_per_hour: float = 1.0,
):
self.rpm = requests_per_minute
self.rph = requests_per_hour
self.cost_per_hour = cost_per_hour
self._minute_window: deque[float] = deque()
self._hour_window: deque[float] = deque()
self._cost_window: deque[tuple[float, float]] = deque() # (timestamp, cost)
def check(self, estimated_cost: float = 0.0) -> tuple[bool, dict]:
"""检查是否允许请求"""
now = time.time()
while self._minute_window and self._minute_window[0] < now - 60:
self._minute_window.popleft()
while self._hour_window and self._hour_window[0] < now - 3600:
self._hour_window.popleft()
while self._cost_window and self._cost_window[0][0] < now - 3600:
self._cost_window.popleft()
minute_count = len(self._minute_window)
hour_count = len(self._hour_window)
hour_cost = sum(c for _, c in self._cost_window)
if minute_count >= self.rpm:
return False, {
"reason": "rate_limit_per_minute",
"limit": self.rpm,
"current": minute_count,
"retry_after": 60 - (now - self._minute_window[0]),
}
if hour_count >= self.rph:
return False, {
"reason": "rate_limit_per_hour",
"limit": self.rph,
"current": hour_count,
"retry_after": 3600 - (now - self._hour_window[0]),
}
if hour_cost + estimated_cost > self.cost_per_hour:
return False, {
"reason": "cost_limit_per_hour",
"limit": self.cost_per_hour,
"current": hour_cost,
"estimated": estimated_cost,
}
return True, {}
def record(self) -> None:
"""记录一次请求"""
now = time.time()
self._minute_window.append(now)
self._hour_window.append(now)
def reset(self) -> None:
"""重置限制器"""
self._minute_window.clear()
self._hour_window.clear()
self._cost_window.clear()
class CostController:
"""LLM 调用成本控制器"""
def __init__(
self,
daily_budget: float = 10.0,
per_request_budget: float = 0.05,
):
self.daily_budget = daily_budget
self.per_request_budget = per_request_budget
self._daily_spend: float = 0.0
self._reset_time: float = time.time() + 86400
def can_proceed(self, estimated_cost: float) -> bool:
"""判断是否可以进行 LLM 调用"""
if time.time() > self._reset_time:
self._daily_spend = 0.0
self._reset_time = time.time() + 86400
if self._daily_spend + estimated_cost > self.daily_budget:
return False
if estimated_cost > self.per_request_budget:
return False
return True
def record_spend(self, cost: float) -> None:
"""记录花费"""
self._daily_spend += cost
def remaining_budget(self) -> float:
"""获取剩余预算"""
if time.time() > self._reset_time:
self._daily_spend = 0.0
self._reset_time = time.time() + 86400
return max(0.0, self.daily_budget - self._daily_spend)
class SecurityOrchestrator:
"""安全编排器 — 统一管理安全策略"""
def __init__(
self,
prompt_detector: Optional[PromptInjectionDetector] = None,
output_sanitizer: Optional[OutputSanitizer] = None,
rate_limiter: Optional[RateLimiter] = None,
cost_controller: Optional[CostController] = None,
):
self.detector = prompt_detector or PromptInjectionDetector()
self.sanitizer = output_sanitizer or OutputSanitizer()
self.rate_limiter = rate_limiter or RateLimiter()
self.cost_controller = cost_controller or CostController()
async def validate_request(
self,
prompt: str,
estimated_cost: float = 0.0,
) -> dict:
"""验证请求的安全性"""
issues = []
injection_findings = self.detector.scan(prompt)
if injection_findings:
issues.append({
"type": "prompt_injection",
"severity": "high",
"details": f"检测到 {len(injection_findings)} 个注入模式",
"findings": injection_findings[:5],
})
rate_allowed, rate_info = self.rate_limiter.check(estimated_cost)
if not rate_allowed:
issues.append({
"type": "rate_limit",
"severity": "medium",
"details": rate_info,
})
cost_allowed = self.cost_controller.can_proceed(estimated_cost)
if not cost_allowed:
issues.append({
"type": "cost_limit",
"severity": "medium",
"details": {
"daily_remaining": self.cost_controller.remaining_budget(),
"estimated_cost": estimated_cost,
},
})
is_allowed = len([i for i in issues if i["type"] == "prompt_injection"]) == 0
return {
"allowed": is_allowed,
"issues": issues,
"sanitized_prompt": self.detector.sanitize_prompt(prompt) if not is_allowed else prompt,
}
def validate_response(self, response: str) -> dict:
"""验证响应安全性"""
has_sensitive = self.sanitizer.has_sensitive_content(response)
return {
"safe": not has_sensitive,
"sanitized": self.sanitizer.sanitize(response) if has_sensitive else response,
"sensitive_detected": has_sensitive,
}
def record_request(self, cost: float = 0.0) -> None:
"""记录请求(用于限流和计费)"""
self.rate_limiter.record()
self.cost_controller.record_spend(cost)
class SecureAgentRuntime:
"""安全的 Agent 运行时 — 集成安全防护层的 Agent 执行环境"""
def __init__(
self,
router: LLMRouter,
registry: ToolRegistry,
security: SecurityOrchestrator,
):
self.router = router
self.registry = registry
self.security = security
self.memory = AgentMemory()
async def process_secure(
self,
prompt: str,
context: AgentContext,
) -> dict:
"""安全处理用户请求"""
validation = await self.security.validate_request(prompt)
if not validation["allowed"]:
return {
"status": "rejected",
"message": "请求被安全策略拒绝",
"issues": validation["issues"],
}
safe_prompt = validation["sanitized_prompt"]
response = await self.router.generate(
prompt=safe_prompt,
context=context,
tools=self.registry.to_openai_tools(),
)
response_validation = self.security.validate_response(response.content)
self.security.record_request(
cost=self.router.providers.get(
LLMProviderType.OPENAI
).estimate_cost(
LLMRequest(prompt=safe_prompt, context=context)
) if LLMProviderType.OPENAI in self.router.providers else 0.0,
)
return {
"status": "completed",
"response": response_validation["sanitized"],
"sensitive_redacted": response_validation["sensitive_detected"],
"tool_calls": response.tool_calls,
"model": response.model,
"latency_ms": response.latency_ms,
}
8. 部署与配置
8.1 环境配置
# 配置文件: config.yaml (示例)
"""
agent_runtime:
name: "msg-agent-runtime"
version: "1.0.0"
chain:
id: "msg-chain-1"
rpc: "https://rpc.msg-chain-1.msg"
rest: "https://rest.msg-chain-1.msg"
prefix: "msg"
llm:
default_provider: "openai"
max_retries: 3
fallback_enabled: true
providers:
openai:
model: "gpt-4o"
api_key_env: "OPENAI_API_KEY"
max_tokens: 4096
temperature: 0.7
anthropic:
model: "claude-3-5-sonnet-通用维护记录"
api_key_env: "ANTHROPIC_API_KEY"
max_tokens: 4096
local:
endpoint: "http://localhost:8000/v1"
model_path: "/models/llama"
max_tokens: 2048
msg_native:
enabled: false
rpc: "https://rpc.msg-chain-1.msg"
security:
prompt_injection_detection: true
output_sanitization: true
rate_limiting:
requests_per_minute: 30
requests_per_hour: 1000
cost_control:
daily_budget_usd: 10.0
per_request_budget_usd: 0.05
memory:
short_term_capacity: 100
long_term_dimension: 128
on_chain_anchoring: false
tools:
max_execution_time_ms: 30000
require_transaction_confirmation: true
"""
class ConfigLoader:
"""配置加载器"""
@classmethod
def load(cls, path: str = "config.yaml") -> dict:
"""加载 YAML 配置文件"""
import yaml
with open(path, "r", encoding="utf-8") as f:
return yaml.safe_load(f)
@classmethod
def load_from_env(cls) -> dict:
"""从环境变量加载配置"""
return {
"chain": {
"id": os.getenv("MSG_CHAIN_ID", "msg-chain-1"),
"rpc": os.getenv("MSG_CHAIN_RPC", "https://rpc.msg-chain-1.msg"),
"prefix": "msg",
},
"llm": {
"providers": {
"openai": {
"api_key": os.getenv("OPENAI_API_KEY", ""),
"model": os.getenv("OPENAI_MODEL", "gpt-4o"),
},
"anthropic": {
"api_key": os.getenv("ANTHROPIC_API_KEY", ""),
"model": os.getenv("ANTHROPIC_MODEL", "claude-3-5-sonnet-通用维护记录"),
},
},
},
"security": {
"daily_budget": float(os.getenv("LLM_DAILY_BUDGET", "10.0")),
},
}
# 运行时初始化示例
async def initialize_runtime(config_path: str = "config.yaml") -> dict:
"""初始化完整的 Agent 运行时"""
config = ConfigLoader.load(config_path)
router = LLMRouter()
registry = ToolRegistry()
security = SecurityOrchestrator(
rate_limiter=RateLimiter(
requests_per_minute=config.get("security", {}).get("rate_limiting", {}).get(
"requests_per_minute", 30
),
),
cost_controller=CostController(
daily_budget=config.get("security", {}).get("cost_control", {}).get(
"daily_budget_usd", 10.0
),
),
)
runtime = SecureAgentRuntime(router, registry, security)
return {
"router": router,
"registry": registry,
"security": security,
"runtime": runtime,
}
# Docker 部署示例 (docker-compose.yml)
"""
version: '3.8'
services:
agent-runtime:
build: .
ports:
- "8080:8080"
environment:
- OPENAI_API_KEY=${OPENAI_API_KEY}
- ANTHROPIC_API_KEY=${ANTHROPIC_API_KEY}
- MSG_CHAIN_RPC=https://rpc.msg-chain-1.msg
- LLM_DAILY_BUDGET=10.0
volumes:
- ./config.yaml:/app/config.yaml
- ./data:/app/data
restart: unless-stopped
local-llm:
image: ghcr.io/msgchain/llama-cpu:latest
ports:
- "8000:8000"
volumes:
- ./models:/models
environment:
- MODEL_PATH=/models/llama-3-8b-instruct.Q4_K_M.gguf
restart: unless-stopped
vector-store:
image: qdrant/qdrant:latest
ports:
- "6333:6333"
volumes:
- ./qdrant_data:/qdrant/storage
restart: unless-stopped
"""
# FastAPI 部署示例
from fastapi import FastAPI, HTTPException
from pydantic import BaseModel
app = FastAPI(title="MSG Chain AI Agent Runtime")
class ChatRequest(BaseModel):
message: str
agent_id: str
user_id: str
conversation_id: Optional[str] = None
class ChatResponse(BaseModel):
response: str
conversation_id: str
model_used: str
latency_ms: float
# 全局运行时实例
_runtime: Optional[SecureAgentRuntime] = None
@app.on_event("startup")
async def startup():
global _runtime
result = await initialize_runtime()
_runtime = result["runtime"]
@app.post("/api/chat", response_model=ChatResponse)
async def chat(request: ChatRequest):
if not _runtime:
raise HTTPException(status_code=503, detail="Runtime not initialized")
context = AgentContext(
agent_id=request.agent_id,
agent_address=f"msg1agent{request.agent_id[:12]}",
did=f"did:msg:{request.agent_id}",
capabilities=["complex_reasoning", "simple_query", "token_transfer"],
constitution="MSG Agent Basic Constitution v1.0",
conversation_id=request.conversation_id or str(uuid.uuid4()),
user_id=request.user_id,
)
result = await _runtime.process_secure(request.message, context)
if result["status"] == "rejected":
raise HTTPException(status_code=400, detail=result["message"])
return ChatResponse(
response=result["response"],
conversation_id=context.conversation_id,
model_used=result.get("model", "unknown"),
latency_ms=result.get("latency_ms", 0.0),
)
@app.get("/api/health")
async def health():
return {
"status": "ok",
"chain": "msg-chain-1",
"runtime_ready": _runtime is not None,
}
@app.get("/api/tools")
async def list_tools(category: Optional[str] = None):
if not _runtime:
raise HTTPException(status_code=503, detail="Runtime not initialized")
cat = ToolCategory(category) if category else None
return {
"tools": _runtime.registry.list_tools(category=cat),
"count": len(_runtime.registry.list_tools(category=cat)),
}
@app.get("/api/stats")
async def get_stats():
if not _runtime:
raise HTTPException(status_code=503, detail="Runtime not initialized")
return {
"llm_router": _runtime.router.get_stats(),
"memory": _runtime.memory.stats(),
}
8.2 环境变量参考
| 变量名 | 默认值 | 说明 |
|---|---|---|
OPENAI_API_KEY |
— | OpenAI API 密钥 |
ANTHROPIC_API_KEY |
— | Anthropic API 密钥 |
OPENAI_MODEL |
gpt-4o |
OpenAI 模型名称 |
ANTHROPIC_MODEL |
claude-3-5-sonnet-通用维护记录 |
Claude 模型名称 |
LOCAL_LLM_PATH |
/models/llama |
本地模型路径 |
MSG_CHAIN_ID |
msg-chain-1 |
MSG Chain ID |
MSG_CHAIN_RPC |
https://rpc.msg-chain-1.msg |
链 RPC 端点 |
MAX_DAILY_LLM_COST |
10.0 |
每日 LLM 成本上限(USD) |
MAX_RPM |
30 |
每分钟最大请求数 |
A2A_SECRET |
default-secret |
Agent 间通信密钥 |
LOG_LEVEL |
INFO |
日志级别 |
8.3 性能优化建议
class PerformanceOptimizer:
"""LLM 集成性能优化建议"""
@staticmethod
def suggest_caching() -> dict:
"""缓存策略建议"""
return {
"prompt_cache": {
"description": "缓存系统提示词和宪法文本",
"benefit": "减少重复 token 消耗",
"implementation": "使用 LRU 缓存常用提示词",
},
"response_cache": {
"description": "缓存常见查询的 LLM 响应",
"benefit": "对高频相同查询减少 90% 延迟",
"implementation": "基于 (prompt_hash, context_hash) 的缓存",
},
"embedding_cache": {
"description": "缓存嵌入向量",
"benefit": "加速记忆检索",
"implementation": "使用本地向量缓存",
},
}
@staticmethod
def suggest_batching() -> dict:
"""批处理策略建议"""
return {
"tool_calls": {
"description": "并行执行独立工具调用",
"benefit": "减少总执行时间 50-80%",
"implementation": "使用 asyncio.gather 并发执行",
},
"batch_inference": {
"description": "批量推理(适用于本地模型)",
"benefit": "提高本地模型吞吐量 2-3x",
"implementation": "收集请求后在本地模型中批处理",
},
}
@staticmethod
def suggest_model_selection() -> dict:
"""模型选择建议"""
return {
"simple_queries": {
"model": LLMProviderType.LOCAL,
"reason": "低延迟、零成本、简单任务不需要强模型",
"examples": ["余额查询", "地址验证", "简单翻译"],
},
"complex_reasoning": {
"model": LLMProviderType.OPENAI,
"reason": "需要强大的推理和多步规划能力",
"examples": ["合约审计分析", "多步交易规划", "复杂问答"],
},
"sensitive_tasks": {
"model": LLMProviderType.ANTHROPIC,
"reason": "更强的安全对齐和宪法遵守",
"examples": ["资产转移", "权限变更", "敏感数据处理"],
},
"high_volume": {
"model": LLMProviderType.MSG_NATIVE,
"reason": "链上轻量推理,适合高频低复杂度查询",
"examples": ["Agent 状态检查", "简单验证", "元数据查询"],
},
}
8.4 运行时健康检查
class RuntimeHealthCheck:
"""运行时健康检查"""
def __init__(self, router: LLMRouter):
self.router = router
self._start_time = time.time()
self._last_checks: dict[str, dict] = {}
async def check_all(self) -> dict:
"""执行全面健康检查"""
results = {}
# 检查 LLM 提供商连接
for ptype, provider in self.router.providers.items():
try:
start = time.time()
request = LLMRequest(prompt="ping", max_tokens=5)
await provider.generate(request)
latency = (time.time() - start) * 1000
results[f"provider_{ptype.value}"] = {
"status": "healthy",
"latency_ms": round(latency, 2),
}
except Exception as e:
results[f"provider_{ptype.value}"] = {
"status": "unhealthy",
"error": str(e)[:100],
}
results["uptime_seconds"] = round(time.time() - self._start_time, 2)
results["total_requests"] = self.router._request_counter
results["daily_cost"] = round(self.router._daily_cost, 4)
self._last_checks["all"] = results
return results
async def check_provider(self, ptype: LLMProviderType) -> dict:
"""检查特定提供商"""
provider = self.router.providers.get(ptype)
if not provider:
return {"status": "not_found"}
try:
start = time.time()
request = LLMRequest(prompt="ping", max_tokens=5)
await provider.generate(request)
latency = (time.time() - start) * 1000
return {
"status": "healthy",
"latency_ms": round(latency, 2),
}
except Exception as e:
return {
"status": "unhealthy",
"error": str(e)[:100],
}
附录
A. MSG 地址前缀参考
| 类型 | 前缀 | 示例 |
|---|---|---|
| 用户地址 | msg1 |
msg1wrx2x6... |
| Agent 地址 | msg1agent |
msg1agentabc123... |
| 合约地址 | msg1contract |
msg1contractxyz... |
| 验证人地址 | msgvaloper1 |
msgvaloper1abc... |
| 验证人共识地址 | msgvalcons1 |
msgvalcons1def... |
B. 常见问题
Q: LLM 路由层如何确保高可用性?
A: 路由层实现了多提供商故障转移。当首选提供商不可用时,自动降级到备选提供商。同时支持断路器模式,在连续失败后自动冷却提供商。
Q: 工具调用的安全性如何保证?
A: 所有工具调用都经过参数验证、权限检查。交易类操作需要用户确认。输出经过敏感内容检测和净化。
Q: Agent 宪法是否可以在运行时更新?
A: 可以。宪法管理器支持热更新,但更新需要链上治理投票(对于公共 Agent)或所有者签名(对于私有 Agent)。
Q: 如何控制 LLM 使用成本?
A: 通过 CostController 设置每日预算和单次请求预算。成本感知路由自动选择性价比最优的模型。当预算接近上限时自动降级到本地模型。
C. 版本信息
| 组件 | 版本 | 说明 |
|---|---|---|
| MSG Chain | v1.0.0 | 主网 |
| LLM Router | v1.0.0 | 本文档实现 |
| Tool Registry | v1.0.0 | 本文档实现 |
| Constitution Manager | v1.0.0 | 本文档实现 |
| Agent Memory | v1.0.0 | 本文档实现 |
| Security Orchestrator | v1.0.0 | 本文档实现 |
本文档由 MSG Chain 开发者文档团队维护。
如有问题或建议,请提交 Issue 至 https://github.com/msgchain/docs
主网状态: No-Go | 白皮书: https://msgchain.org/whitepaper/
