MSG Chain AI Agent 跨平台机器人集成指南
⚠️ No-Go Disclaimer: MSGChain 主网裁决为 No-Go。本文件所有内容反映的是开发阶段的技术设计,不代表主网未独立核验上线状态。生产部署状态请以白皮书为准:https://msgchain.org/whitepaper/
1. 概述
1.1 为什么跨平台机器人对 AI Agent 至关重要
AI Agent 的本质是自主执行任务的智能体。但如果 Agent 只能通过 CLI 或 Web UI 交互,其可用性和触达范围将严重受限。将 Agent 接入即时通讯平台(Telegram、Discord、Slack)可以实现:
- 自然交互:用户用自然语言与 Agent 对话,无需学习专用接口
- 无缝支付:通过平台内置支付系统或 MSG Chain 链上交易完成 Agent 服务付费
- 社区协作:多个用户可在群组中同时与 Agent 交互,形成协作式 AI 工作流
- 随时随地:移动端优先的体验,让 Agent 触手可及
1.2 MSG Chain 与 AI Agent 的关系
MSG Chain 是专为 AI Agent 经济设计的 L1 区块链,支持:
- Agent 身份注册与 DID 认证
- Agent 服务市场的链上结算
- A2A(Agent-to-Agent)通信的链上验证
- 跨平台身份统一与权限管理
1.3 架构概览
+-----------------------------------------------------------+
| 消息平台 (用户侧) |
| Telegram Bot | Discord Bot | Slack Bot |
+--------+----------------+----------------+---------------+
| | |
v v v
+-----------------------------------------------------------+
| Bot Middleware 层 |
| Platform Adapter (统一消息格式) |
| DID 身份管理 | 会话管理 | 速率限制 |
| A2A 消息路由 | 支付网关 |
+---------------------------+-------------------------------+
|
v
+-----------------------------------------------------------+
| Agent API 层 |
| Agent 发现 | 任务执行 | 状态查询 | 结果回调 |
+---------------------------+-------------------------------+
|
v
+-----------------------------------------------------------+
| MSG Chain 层 |
| Agent Registry (智能合约) | DID 认证 | 链上结算 |
+-----------------------------------------------------------+
1.4 技术栈选择
| 组件 | 推荐技术 | 版本要求 |
|---|---|---|
| Telegram Bot | python-telegram-bot | >= 21.0 |
| Discord Bot | discord.py | >= 2.3 |
| Slack Bot | slack-bolt | >= 1.18 |
| MSG Chain SDK | msgchain-sdk-python | >= 0.4 |
| Async HTTP | httpx | >= 0.27 |
| 密钥管理 | python-dotenv + cryptography | >= 1.0 |
1.5 前置条件
- Python 3.11+
- 各平台 Bot Token(Telegram: @BotFather, Discord: Developer Portal, Slack: App Credentials)
- MSG Chain 节点 RPC URL
- Agent DID 私钥(用于身份认证签名)
1.6 文档约定
- 所有 MSG Chain 地址使用
msg1前缀(bech32 编码) - 环境变量使用
.env文件管理 - 代码示例兼容 Python 3.11+ 类型注解
- 异步优先,所有网络操作使用
asyncio
2. 架构设计
2.1 核心组件关系
Bot Middleware 是整个系统的枢纽。它承担以下职责:
- 平台适配:将各平台的消息格式统一为内部标准格式
- 身份代理:为每个 Agent 管理跨平台 DID 身份
- 消息路由:在用户 ↔ Agent ↔ MSG Chain 之间流转消息
- 会话管理:维护用户与 Agent 的对话状态
- 支付集成:处理平台支付和链上结算
2.2 平台适配器模式
import asyncio
import logging
import time
from abc import ABC, abstractmethod
from dataclasses import dataclass, field
from enum import Enum
from typing import Any, Callable, Coroutine, Optional
logger = logging.getLogger(__name__)
class MessageType(Enum):
TEXT = "text"
COMMAND = "command"
CALLBACK = "callback"
PAYMENT = "payment"
FILE = "file"
IMAGE = "image"
class PlatformType(Enum):
TELEGRAM = "telegram"
DISCORD = "discord"
SLACK = "slack"
@dataclass
class NormalizedMessage:
platform: PlatformType
platform_id: str
channel_id: str
user_id: str
user_name: str
message_type: MessageType
text: Optional[str] = None
data: Optional[dict] = None
attachments: list[dict] = field(default_factory=list)
timestamp: float = field(default_factory=time.time)
raw: Optional[Any] = None
@dataclass
class NormalizedResponse:
platform: PlatformType
channel_id: str
user_id: str
text: Optional[str] = None
embeds: list[dict] = field(default_factory=list)
buttons: list[dict] = field(default_factory=list)
modal: Optional[dict] = None
ephemeral: bool = False
raw: Optional[Any] = None
class PlatformConnector(ABC):
\"\"\"平台连接器基类 — 适配器模式\"\"\"
def __init__(self, platform: PlatformType, bot_token: str):
self.platform = platform
self.bot_token = bot_token
self._handler: Optional[Callable[[NormalizedMessage], Coroutine]] = None
@abstractmethod
async def start(self):
...
@abstractmethod
async def stop(self):
...
@abstractmethod
async def send_message(self, response: NormalizedResponse) -> Any:
...
async def handle_message(self, message: NormalizedMessage):
if self._handler:
await self._handler(message)
def set_message_handler(self, handler: Callable[[NormalizedMessage], Coroutine]):
self._handler = handler
2.3 DID 身份管理
import hashlib
import json
from cryptography.hazmat.primitives import hashes
from cryptography.hazmat.primitives.asymmetric import ed25519
from cryptography.hazmat.primitives.serialization import (
load_pem_private_key,
load_pem_public_key,
)
class DIDManager:
\"\"\"DID 身份管理器 — 为每个 Agent 管理跨平台身份\"\"\"
def __init__(self, did_private_key_pem: bytes, agent_did: str):
self.agent_did = agent_did
self._private_key = load_pem_private_key(did_private_key_pem, password=None)
self._public_key = self._private_key.public_key()
self.platform_identities: dict[str, str] = {}
def sign_message(self, payload: dict) -> str:
message = json.dumps(payload, sort_keys=True, separators=(",", ":")).encode()
signature = self._private_key.sign(message)
return signature.hex()
def verify_signature(self, payload: dict, signature_hex: str, public_key_hex: str) -> bool:
message = json.dumps(payload, sort_keys=True, separators=(",", ":")).encode()
signature = bytes.fromhex(signature_hex)
public_key = ed25519.Ed25519PublicKey.from_public_bytes(bytes.fromhex(public_key_hex))
try:
public_key.verify(signature, message)
return True
except Exception:
return False
def get_did_document(self) -> dict:
public_key_bytes = self._public_key.public_bytes_raw()
return {
"@context": "https://www.w3.org/ns/did/v1",
"id": self.agent_did,
"verificationMethod": [
{
"id": f"{self.agent_did}#keys-1",
"type": "Ed25519VerificationKey2020",
"controller": self.agent_did,
"publicKeyMultibase": public_key_bytes.hex(),
}
],
"authentication": [f"{self.agent_did}#keys-1"],
"service": [
{
"id": f"{self.agent_did}#agent-api",
"type": "AgentAPIEndpoint",
"serviceEndpoint": "https://api.agent.example.com/v1",
}
],
}
def register_platform_link(self, platform: str, platform_user_id: str):
self.platform_identities[platform] = platform_user_id
2.4 Agent API 客户端
import httpx
from uuid import uuid4
class AgentAPIClient:
\"\"\"Agent API 客户端 — 与 Agent 后端通信\"\"\"
def __init__(self, api_base_url: str, did_manager: DIDManager, timeout: float = 30.0):
self.api_base_url = api_base_url.rstrip("/")
self.did = did_manager
self.timeout = timeout
self._client = httpx.AsyncClient(timeout=timeout)
async def close(self):
await self._client.aclose()
async def _signed_request(self, method: str, path: str, body: Optional[dict] = None) -> dict:
timestamp = int(time.time())
nonce = uuid4().hex[:16]
payload = {
"method": method,
"path": path,
"body": body or {},
"timestamp": timestamp,
"nonce": nonce,
}
signature = self.did.sign_message(payload)
headers = {
"X-Agent-DID": self.did.agent_did,
"X-Agent-Signature": signature,
"X-Agent-Timestamp": str(timestamp),
"X-Agent-Nonce": nonce,
"Content-Type": "application/json",
}
url = f"{self.api_base_url}{path}"
if method.upper() == "GET":
resp = await self._client.get(url, headers=headers, params=body)
elif method.upper() == "POST":
resp = await self._client.post(url, headers=headers, json=body)
else:
resp = await self._client.request(method, url, headers=headers, json=body)
resp.raise_for_status()
return resp.json()
async def query_agent(self, agent_did: str) -> dict:
return await self._signed_request("GET", f"/agents/{agent_did}")
async def list_agents(self, page: int = 1, size: int = 20) -> dict:
return await self._signed_request("GET", "/agents", {"page": page, "size": size})
async def create_session(self, agent_did: str) -> dict:
return await self._signed_request("POST", "/sessions", {"agent_did": agent_did})
async def send_message(self, session_id: str, content: str) -> dict:
return await self._signed_request("POST", f"/sessions/{session_id}/messages", {"content": content})
async def get_session_history(self, session_id: str) -> dict:
return await self._signed_request("GET", f"/sessions/{session_id}/messages")
async def execute_task(self, agent_did: str, task: dict) -> dict:
return await self._signed_request("POST", "/tasks", {"agent_did": agent_did, "task": task})
async def get_task_status(self, task_id: str) -> dict:
return await self._signed_request("GET", f"/tasks/{task_id}")
async def initiate_payment(self, to_agent_did: str, amount: str, denom: str = "umsg") -> dict:
return await self._signed_request("POST", "/payments", {"to": to_agent_did, "amount": amount, "denom": denom})
async def get_balance(self) -> dict:
return await self._signed_request("GET", "/wallet/balance")
2.5 Bot Middleware 主类
class BotMiddleware:
\"\"\"Bot Middleware 主类 — 统一管理多平台机器人\"\"\"
def __init__(self, agent_api_url: str, did_private_key_pem: bytes, agent_did: str):
self.did_manager = DIDManager(did_private_key_pem, agent_did)
self.agent_api = AgentAPIClient(agent_api_url, self.did_manager)
self.platforms: dict[str, PlatformConnector] = {}
self.session_map: dict[str, str] = {}
self.rate_limiter: dict[str, list[float]] = {}
self._running = False
def register_platform(self, name: str, connector: PlatformConnector):
self.platforms[name] = connector
connector.set_message_handler(self._on_platform_message)
logger.info(f"平台已注册: {name}")
async def _on_platform_message(self, msg: NormalizedMessage):
user_key = f"{msg.platform.value}:{msg.user_id}"
if not self._check_rate_limit(user_key):
logger.warning(f"速率限制触发: {user_key}")
await self._send_platform_message(NormalizedResponse(
platform=msg.platform, channel_id=msg.channel_id, user_id=msg.user_id,
text="请求过于频繁,请稍后再试。", ephemeral=True,
))
return
try:
await self._route_message(msg)
except Exception as e:
logger.exception(f"消息处理异常: {e}")
await self._send_platform_message(NormalizedResponse(
platform=msg.platform, channel_id=msg.channel_id, user_id=msg.user_id,
text=f"处理消息时出错: {str(e)[:200]}",
))
async def _route_message(self, msg: NormalizedMessage):
if msg.message_type == MessageType.COMMAND:
await self._handle_command(msg)
elif msg.message_type == MessageType.CALLBACK:
await self._handle_callback(msg)
elif msg.message_type == MessageType.PAYMENT:
await self._handle_payment(msg)
else:
await self._handle_text(msg)
async def _handle_command(self, msg: NormalizedMessage):
cmd = (msg.text or "").strip().lower()
if cmd in ("/start", "/help"):
await self._send_welcome(msg)
elif cmd.startswith("/agent"):
parts = cmd.split(maxsplit=1)
if len(parts) > 1:
await self._lookup_agent(msg, parts[1])
else:
await self._list_agents(msg)
elif cmd == "/balance":
await self._check_balance(msg)
elif cmd == "/session":
await self._show_session(msg)
else:
await self._send_platform_message(NormalizedResponse(
platform=msg.platform, channel_id=msg.channel_id, user_id=msg.user_id,
text="未知命令。输入 /help 查看可用命令。",
))
async def _handle_text(self, msg: NormalizedMessage):
user_key = f"{msg.platform.value}:{msg.user_id}"
session_id = self.session_map.get(user_key)
if not session_id:
agents = await self.agent_api.list_agents(size=1)
if not agents.get("agents"):
await self._send_platform_message(NormalizedResponse(
platform=msg.platform, channel_id=msg.channel_id, user_id=msg.user_id,
text="当前没有可用的 Agent。请稍后再试。",
))
return
agent_did = agents["agents"][0]["did"]
result = await self.agent_api.create_session(agent_did)
session_id = result["session_id"]
self.session_map[user_key] = session_id
await self._send_platform_message(NormalizedResponse(
platform=msg.platform, channel_id=msg.channel_id, user_id=msg.user_id,
text=f"已连接到 Agent: {result.get('agent_name', agent_did)}\n"
f"会话 ID: `{session_id[:16]}...`",
))
result = await self.agent_api.send_message(session_id, msg.text or "")
reply = result.get("reply", result.get("content", "(无响应)"))
await self._send_platform_message(NormalizedResponse(
platform=msg.platform, channel_id=msg.channel_id, user_id=msg.user_id,
text=reply,
))
async def _handle_callback(self, msg: NormalizedMessage):
if not msg.data:
return
action = msg.data.get("action", "")
if action == "select_agent":
agent_did = msg.data.get("agent_did", "")
user_key = f"{msg.platform.value}:{msg.user_id}"
result = await self.agent_api.create_session(agent_did)
self.session_map[user_key] = result["session_id"]
await self._send_platform_message(NormalizedResponse(
platform=msg.platform, channel_id=msg.channel_id, user_id=msg.user_id,
text=f"已选择 Agent `{result.get('agent_name', agent_did)}`。请发送您的消息。",
))
elif action == "pay_agent":
await self._initiate_payment_flow(msg)
elif action == "view_profile":
await self._show_agent_profile(msg)
async def _handle_payment(self, msg: NormalizedMessage):
await self._send_platform_message(NormalizedResponse(
platform=msg.platform, channel_id=msg.channel_id, user_id=msg.user_id,
text="支付处理中,请稍候...",
))
async def _send_welcome(self, msg: NormalizedMessage):
welcome_text = (
"AI Agent 跨平台机器人\n\n"
"我是 MSG Chain 上的 AI Agent 入口,您可以:\n\n"
"- `/agent list` -- 浏览 Agent 市场\n"
"- `/agent <name>` -- 搜索特定 Agent\n"
"- `/balance` -- 查询钱包余额\n"
"- `/session` -- 查看当前会话\n\n"
"直接发送消息即可与当前 Agent 对话。"
)
await self._send_platform_message(NormalizedResponse(
platform=msg.platform, channel_id=msg.channel_id, user_id=msg.user_id,
text=welcome_text,
))
async def _list_agents(self, msg: NormalizedMessage):
result = await self.agent_api.list_agents(page=1, size=10)
agents = result.get("agents", [])
if not agents:
await self._send_platform_message(NormalizedResponse(
platform=msg.platform, channel_id=msg.channel_id, user_id=msg.user_id,
text="当前没有注册的 Agent。",
))
return
lines = ["Agent 市场\n"]
for idx, agent in enumerate(agents[:10], 1):
name = agent.get("name", agent.get("did", "未知")[:16])
price = agent.get("price", "免费")
lines.append(f"{idx}. **{name}** -- {price}")
lines.append("\n点击下方按钮选择 Agent:")
buttons = []
for agent in agents[:5]:
buttons.append({
"text": agent.get("name", agent["did"][:8]),
"data": {"action": "select_agent", "agent_did": agent["did"]},
})
await self._send_platform_message(NormalizedResponse(
platform=msg.platform, channel_id=msg.channel_id, user_id=msg.user_id,
text="\n".join(lines), buttons=buttons,
))
async def _lookup_agent(self, msg: NormalizedMessage, query: str):
result = await self.agent_api.list_agents(page=1, size=20)
agents = result.get("agents", [])
matches = [a for a in agents if query.lower() in a.get("name", "").lower()]
if not matches:
await self._send_platform_message(NormalizedResponse(
platform=msg.platform, channel_id=msg.channel_id, user_id=msg.user_id,
text=f"未找到匹配 '{query}' 的 Agent。",
))
return
lines = [f"找到 {len(matches)} 个 Agen[未公开路径]"]
for agent in matches:
lines.append(f"- **{agent.get('name', '未知')}**")
lines.append(f" DID: `{agent['did']}`")
lines.append(f" 描述: {agent.get('description', '无')[:80]}")
await self._send_platform_message(NormalizedResponse(
platform=msg.platform, channel_id=msg.channel_id, user_id=msg.user_id,
text="\n".join(lines),
))
async def _check_balance(self, msg: NormalizedMessage):
result = await self.agent_api.get_balance()
balance = result.get("balance", "0")
await self._send_platform_message(NormalizedResponse(
platform=msg.platform, channel_id=msg.channel_id, user_id=msg.user_id,
text=f"当前钱包余额: **{balance} umsg**\n"
f"地址: `{result.get('address', 'N/A')}`",
))
async def _show_session(self, msg: NormalizedMessage):
user_key = f"{msg.platform.value}:{msg.user_id}"
session_id = self.session_map.get(user_key)
if session_id:
history = await self.agent_api.get_session_history(session_id)
messages = history.get("messages", [])
text = f"当前会话: `{session_id[:16]}...`\n消息数: {len(messages)}"
else:
text = "当前没有活跃会话。发送消息将自动创建新会话。"
await self._send_platform_message(NormalizedResponse(
platform=msg.platform, channel_id=msg.channel_id, user_id=msg.user_id,
text=text,
))
async def _show_agent_profile(self, msg: NormalizedMessage):
agent_did = msg.data.get("agent_did", "") if msg.data else ""
if not agent_did:
return
result = await self.agent_api.query_agent(agent_did)
agent = result.get("agent", {})
text = (
f"Agent: {agent.get('name', '未知')}\n"
f"DID: `{agent.get('did', agent_did)}`\n"
f"描述: {agent.get('description', '无')}\n"
f"价格: {agent.get('price', '免费')}\n"
f"评分: {agent.get('rating', '暂无')}\n"
f"已完成任务: {agent.get('completed_tasks', 0)}"
)
await self._send_platform_message(NormalizedResponse(
platform=msg.platform, channel_id=msg.channel_id, user_id=msg.user_id,
text=text,
))
async def _initiate_payment_flow(self, msg: NormalizedMessage):
agent_did = msg.data.get("agent_did", "") if msg.data else ""
amount = msg.data.get("amount", "1000000")
try:
result = await self.agent_api.initiate_payment(agent_did, amount)
tx_hash = result.get("tx_hash", "")
await self._send_platform_message(NormalizedResponse(
platform=msg.platform, channel_id=msg.channel_id, user_id=msg.user_id,
text=f"支付成功!\n交易哈希: `{tx_hash}`\n金额: {amount} umsg",
))
except Exception as e:
await self._send_platform_message(NormalizedResponse(
platform=msg.platform, channel_id=msg.channel_id, user_id=msg.user_id,
text=f"支付失败: {str(e)[:200]}",
))
def _check_rate_limit(self, key: str, max_requests: int = 10, window: float = 60.0) -> bool:
now = time.time()
timestamps = self.rate_limiter.get(key, [])
timestamps = [t for t in timestamps if now - t < window]
if len(timestamps) >= max_requests:
return False
timestamps.append(now)
self.rate_limiter[key] = timestamps
return True
async def _send_platform_message(self, response: NormalizedResponse):
connector = self.platforms.get(response.platform.value)
if connector:
await connector.send_message(response)
else:
logger.error(f"未注册的平台连接器: {response.platform.value}")
async def broadcast(self, text: str, exclude_platform: Optional[str] = None):
response = NormalizedResponse(platform=PlatformType.TELEGRAM, channel_id="", user_id="", text=text)
for name, connector in self.platforms.items():
if name == exclude_platform:
continue
response.platform = PlatformType(name)
await connector.send_message(response)
async def start(self):
self._running = True
logger.info(f"Bot Middleware 启动, 已注册平台: {list(self.platforms.keys())}")
tasks = [connector.start() for connector in self.platforms.values()]
await asyncio.gather(*tasks)
async def stop(self):
self._running = False
logger.info("正在停止所有平台连接...")
tasks = [connector.stop() for connector in self.platforms.values()]
await asyncio.gather(*tasks)
await self.agent_api.close()
2.6 统一入口与主函数
import os
from dotenv import load_dotenv
load_dotenv()
def load_private_key_from_env() -> bytes:
key_hex = os.environ.get("AGENT_DID_PRIVATE_KEY")
if key_hex:
return bytes.fromhex(key_hex)
key_path = os.environ.get("AGENT_DID_PRIVATE_KEY_PATH", "./did_key.pem")
with open(key_path, "rb") as f:
return f.read()
async def main():
agent_did = os.environ.get("AGENT_DID", "msg1agentdefault00000000000000000")
agent_api_url = os.environ.get("AGENT_API_URL", "https://api.msgchain.org/v1")
private_key = load_private_key_from_env()
middleware = BotMiddleware(agent_api_url=agent_api_url, did_private_key_pem=private_key, agent_did=agent_did)
from telegram_connector import TelegramConnector
from discord_connector import DiscordConnector
from slack_connector import SlackConnector
middleware.register_platform("telegram", TelegramConnector(os.environ["TELEGRAM_BOT_TOKEN"]))
middleware.register_platform("discord", DiscordConnector(os.environ["DISCORD_BOT_TOKEN"]))
middleware.register_platform("slack", SlackConnector(os.environ["SLACK_BOT_TOKEN"]))
try:
await middleware.start()
logger.info("所有平台机器人已启动")
await asyncio.Event().wait()
except KeyboardInterrupt:
pass
finally:
await middleware.stop()
if __name__ == "__main__":
logging.basicConfig(level=logging.INFO, format="%(asctime)s [%(levelname)s] %(name)s: %(message)s")
asyncio.run(main())
3. Telegram Bot 集成
3.1 依赖安装
pip install python-telegram-bot>=21.0 httpx>=0.27
3.2 Telegram 连接器实现
import json
import logging
from typing import Optional
from telegram import Update, InlineKeyboardButton, InlineKeyboardMarkup
from telegram.ext import (
Application, CommandHandler, CallbackQueryHandler,
MessageHandler, filters, ContextTypes,
)
from bot_middleware import (
PlatformConnector, PlatformType, NormalizedMessage, NormalizedResponse, MessageType,
)
logger = logging.getLogger(__name__)
class TelegramConnector(PlatformConnector):
\"\"\"Telegram 平台连接器\"\"\"
def __init__(self, bot_token: str):
super().__init__(PlatformType.TELEGRAM, bot_token)
self._app: Optional[Application] = None
async def start(self):
self._app = Application.builder().token(self.bot_token).build()
self._app.add_handler(CommandHandler("start", self._cmd_handler))
self._app.add_handler(CommandHandler("help", self._cmd_handler))
self._app.add_handler(CommandHandler("agent", self._cmd_handler))
self._app.add_handler(CommandHandler("balance", self._cmd_handler))
self._app.add_handler(CommandHandler("session", self._cmd_handler))
self._app.add_handler(CallbackQueryHandler(self._callback_handler))
self._app.add_handler(MessageHandler(filters.TEXT & ~filters.COMMAND, self._msg_handler))
self._app.add_handler(MessageHandler(filters.SUCCESSFUL_PAYMENT, self._payment_handler))
await self._app.initialize()
await self._app.start()
logger.info("Telegram Bot 已启动")
async def stop(self):
if self._app:
await self._app.stop()
await self._app.shutdown()
logger.info("Telegram Bot 已停止")
async def send_message(self, response: NormalizedResponse) -> None:
if not self._app:
return
chat_id = response.channel_id or response.user_id
kwargs = {"chat_id": chat_id, "parse_mode": "Markdown"}
if response.buttons:
keyboard = [
[InlineKeyboardButton(btn["text"], callback_data=json.dumps(btn["data"], separators=(",", ":")))]
for btn in response.buttons
]
kwargs["reply_markup"] = InlineKeyboardMarkup(keyboard)
if response.text:
kwargs["text"] = response.text
if response.ephemeral:
kwargs["text"] = f"*(仅您可见)*\n{response.text or ''}"
try:
await self._app.bot.send_message(**kwargs)
except Exception as e:
logger.error(f"Telegram 发送消息失败: {e}")
async def _normalize(self, update: Update, msg_type: MessageType) -> Optional[NormalizedMessage]:
if not update.effective_user or not update.effective_chat:
return None
return NormalizedMessage(
platform=PlatformType.TELEGRAM,
platform_id=str(update.update_id),
channel_id=str(update.effective_chat.id),
user_id=str(update.effective_user.id),
user_name=update.effective_user.full_name or update.effective_user.username or "unknown",
message_type=msg_type,
text=update.message.text if update.message else None,
data=None,
timestamp=update.message.date.timestamp() if update.message else 0,
raw=update,
)
async def _cmd_handler(self, update: Update, context: ContextTypes.DEFAULT_TYPE):
msg = await self._normalize(update, MessageType.COMMAND)
if msg:
await self.handle_message(msg)
async def _callback_handler(self, update: Update, context: ContextTypes.DEFAULT_TYPE):
query = update.callback_query
await query.answer()
data = json.loads(query.data) if query.data else {}
msg = NormalizedMessage(
platform=PlatformType.TELEGRAM,
platform_id=str(update.update_id),
channel_id=str(query.message.chat_id if query.message else ""),
user_id=str(query.from_user.id),
user_name=query.from_user.full_name or "unknown",
message_type=MessageType.CALLBACK,
text=query.message.text if query.message else None,
data=data,
raw=update,
)
await self.handle_message(msg)
async def _msg_handler(self, update: Update, context: ContextTypes.DEFAULT_TYPE):
msg = await self._normalize(update, MessageType.TEXT)
if msg:
await self.handle_message(msg)
async def _payment_handler(self, update: Update, context: ContextTypes.DEFAULT_TYPE):
msg = await self._normalize(update, MessageType.PAYMENT)
if msg:
await self.handle_message(msg)
3.3 Telegram Bot 完整示例: Agent 市场机器人
#!/usr/bin/env python3
\"\"\"
Telegram Agent 市场机器人 -- 完整示例
用户可以通过此机器人浏览、选择、并与 MSG Chain 上的 AI Agent 交互。
\"\"\"
import asyncio
import json
import logging
import os
import sys
from typing import Optional
from dotenv import load_dotenv
from telegram import Update, InlineKeyboardButton, InlineKeyboardMarkup, LabeledPrice
from telegram.ext import (
Application, CommandHandler, CallbackQueryHandler,
PreCheckoutQueryHandler, MessageHandler, filters, ContextTypes,
)
load_dotenv()
logging.basicConfig(level=logging.INFO, format="%(asctime)s [%(levelname)s] %(name)s: %(message)s")
logger = logging.getLogger(__name__)
BOT_TOKEN = os.environ.get("TELEGRAM_BOT_TOKEN", "")
AGENT_API_URL = os.environ.get("AGENT_API_URL", "https://api.msgchain.org/v1")
AGENT_DID = os.environ.get("AGENT_DID", "msg1agentdefault00000000000000000")
MOCK_AGENTS = [
{"did": "msg1agenttextgen0000000000000001", "name": "文本生成助手",
"description": "基于 GPT 的高质量文本生成 Agent,支持文章、代码、文案等多种格式。",
"price": "500000 umsg", "rating": 4.8, "tasks_completed": 1247, "category": "text"},
{"did": "msg1agentimagedfx00000000000002", "name": "图像生成大师",
"description": "AI 图像生成 Agent,支持文生图、图生图、风格迁移。",
"price": "1000000 umsg", "rating": 4.6, "tasks_completed": 892, "category": "image"},
{"did": "msg1agentdataanal000000000000003", "name": "数据分析专家",
"description": "数据分析 Agent,支持 CSV/Excel 分析、可视化、报表生成。",
"price": "800000 umsg", "rating": 4.9, "tasks_completed": 2156, "category": "data"},
{"did": "msg1agentcodegenp000000000000004", "name": "代码生成器",
"description": "全栈代码生成 Agent,支持多种编程语言和框架。",
"price": "600000 umsg", "rating": 4.7, "tasks_completed": 3451, "category": "code"},
{"did": "msg1agenttransla000000000000005", "name": "翻译通",
"description": "多语言翻译 Agent,支持 100+ 语言互译,保留格式和语义。",
"price": "300000 umsg", "rating": 4.5, "tasks_completed": 5678, "category": "translate"},
]
user_sessions: dict[str, str] = {}
user_balances: dict[str, str] = {}
async def start(update: Update, context: ContextTypes.DEFAULT_TYPE):
user = update.effective_user
welcome = (
f"欢迎 {user.full_name}!\n\n"
f"我是 **MSG Chain AI Agent 市场机器人**。\n"
f"在这里您可以发现、试用和支付链上的 AI Agent。\n\n"
f"**可用命令:**\n"
f"- `/market` -- 浏览 Agent 市场\n"
f"- `/agent <id>` -- 查看 Agent 详情\n"
f"- `/balance` -- 查询钱包余额\n"
f"- `/chat <消息>` -- 与当前 Agent 对话\n"
f"- `/disconnect` -- 断开与 Agent 的连接\n"
f"- `/help` -- 显示帮助信息"
)
await update.message.reply_text(welcome, parse_mode="Markdown")
async def help_command(update: Update, context: ContextTypes.DEFAULT_TYPE):
help_text = (
"帮助中心\n\n"
"*什么是 MSG Chain AI Agent?*\n"
"MSG Chain 上的 AI Agent 是可自主执行任务的智能程序。\n"
"您可以通过此机器人发现和调用它们。\n\n"
"*如何使用?*\n"
"1. 使用 `/market` 浏览可用 Agent\n"
"2. 选择一个 Agent 连接到会话\n"
"3. 直接发送消息与之交互\n"
"4. 需要时使用 `/balance` 查看余额并支付\n\n"
"*链上地址格式*\n"
"所有 Agent 使用 `msg1` 开头的 bech32 地址。\n\n"
"*遇到问题?*\n"
"请访问 https://msgchain.org/support 获取帮助。"
)
await update.message.reply_text(help_text, parse_mode="Markdown")
async def market(update: Update, context: ContextTypes.DEFAULT_TYPE):
page = int(context.args[0]) if context.args and context.args[0].isdigit() else 1
agents_per_page = 3
total_pages = (len(MOCK_AGENTS) + agents_per_page - 1) // agents_per_page
start_idx = (page - 1) * agents_per_page
end_idx = start_idx + agents_per_page
page_agents = MOCK_AGENTS[start_idx:end_idx]
text_parts = [f"AI Agent 市场 (第 {page}/{total_pages} 页)\n"]
keyboard = []
for agent in page_agents:
text_parts.append(
f"\n**{agent['name']}**\n"
f"分类: {agent['category']}\n"
f"评分: {agent['rating']}\n"
f"任务: {agent['tasks_completed']}\n"
f"价格: {agent['price']}"
)
keyboard.append([InlineKeyboardButton(f" {agent['name'][:12]}", callback_data=json.dumps({"action": "detail", "did": agent["did"]}))])
nav_buttons = []
if page > 1:
nav_buttons.append(InlineKeyboardButton("上一页", callback_data=json.dumps({"action": "page", "page": page - 1})))
if page < total_pages:
nav_buttons.append(InlineKeyboardButton("下一页", callback_data=json.dumps({"action": "page", "page": page + 1})))
if nav_buttons:
keyboard.append(nav_buttons)
await update.message.reply_text("\n".join(text_parts), parse_mode="Markdown", reply_markup=InlineKeyboardMarkup(keyboard))
async def agent_detail(update: Update, context: ContextTypes.DEFAULT_TYPE):
did = " ".join(context.args) if context.args else ""
if not did:
await update.message.reply_text("用法: `/agent <agent_did>`")
return
agent = next((a for a in MOCK_AGENTS if a["did"] == did), None)
if not agent:
await update.message.reply_text(f"未找到 Agent: `{did}`")
return
text = (
f"*{agent['name']}*\n\n"
f"**DID:** `{agent['did']}`\n"
f"**分类:** {agent['category']}\n"
f"**描述:** {agent['description']}\n"
f"**价格:** {agent['price']}\n"
f"**评分:** {agent['rating']}/5.0\n"
f"**完成任务:** {agent['tasks_completed']}\n\n"
f"选择操作:"
)
keyboard = [
[InlineKeyboardButton("连接", callback_data=json.dumps({"action": "connect", "did": did})),
InlineKeyboardButton("支付", callback_data=json.dumps({"action": "pay", "did": did}))],
[InlineKeyboardButton("返回市场", callback_data=json.dumps({"action": "back_to_market"}))],
]
await update.message.reply_text(text, parse_mode="Markdown", reply_markup=InlineKeyboardMarkup(keyboard))
async def balance(update: Update, context: ContextTypes.DEFAULT_TYPE):
user_id = str(update.effective_user.id)
balance_str = user_balances.get(user_id, "10,000,000 umsg")
await update.message.reply_text(f"当前余额\n\n`{balance_str}`\n\n提示:使用 MSG Chain 钱包充值后余额将自动更新。", parse_mode="Markdown")
async def chat_with_agent(update: Update, context: ContextTypes.DEFAULT_TYPE):
user_id = str(update.effective_user.id)
agent_did = user_sessions.get(user_id)
if not agent_did:
await update.message.reply_text("您尚未连接 Agent。请先使用 `/market` 选择一个 Agent。")
return
message = " ".join(context.args) if context.args else ""
if not message:
await update.message.reply_text("用法: `/chat <消息>`")
return
agent = next((a for a in MOCK_AGENTS if a["did"] == agent_did), None)
agent_name = agent["name"] if agent else agent_did[:16]
await update.message.reply_text(f"*{agent_name}* 正在处理...", parse_mode="Markdown")
reply = f"*{agent_name}* 回复:\n\n收到您的消息: \"{message}\"\n\n(此为模拟回复,实际部署时将调用 Agent API)"
await update.message.reply_text(reply, parse_mode="Markdown")
async def disconnect(update: Update, context: ContextTypes.DEFAULT_TYPE):
user_id = str(update.effective_user.id)
if user_id in user_sessions:
agent_did = user_sessions.pop(user_id)
await update.message.reply_text(f"已断开与 Agent `{agent_did[:16]}...` 的连接。")
else:
await update.message.reply_text("当前没有连接的 Agent。")
async def about(update: Update, context: ContextTypes.DEFAULT_TYPE):
await update.message.reply_text(
"MSG Chain AI Agent 市场机器人\n\n"
f"链: MSG Chain (msg)\n网络: Mainnet\n版本: 1.0.0\n\n文档: https://docs.msgchain.org",
parse_mode="Markdown", disable_web_page_preview=True,
)
async def button_handler(update: Update, context: ContextTypes.DEFAULT_TYPE):
query = update.callback_query
await query.answer()
try:
data = json.loads(query.data)
except json.JSONDecodeError:
logger.error(f"无法解析 callback data: {query.data}")
return
action = data.get("action", "")
if action == "detail":
did = data["did"]
agent = next((a for a in MOCK_AGENTS if a["did"] == did), None)
if not agent:
await query.edit_message_text("Agent 不存在。")
return
text = (
f"*{agent['name']}*\n\n**DID:** `{agent['did']}`\n"
f"**分类:** {agent['category']}\n**描述:** {agent['description']}\n"
f"**价格:** {agent['price']}\n**评分:** {agent['rating']}/5.0\n"
f"**完成任务:** {agent['tasks_completed']}\n\n选择操作:"
)
keyboard = [
[InlineKeyboardButton("连接", callback_data=json.dumps({"action": "connect", "did": did})),
InlineKeyboardButton("支付", callback_data=json.dumps({"action": "pay", "did": did}))],
[InlineKeyboardButton("返回市场", callback_data=json.dumps({"action": "back_to_market"}))],
]
await query.edit_message_text(text, parse_mode="Markdown", reply_markup=InlineKeyboardMarkup(keyboard))
elif action == "connect":
did = data["did"]
user_id = str(query.from_user.id)
user_sessions[user_id] = did
agent = next((a for a in MOCK_AGENTS if a["did"] == did), None)
name = agent["name"] if agent else did[:16]
await query.edit_message_text(
f"已连接到 **{name}**!\n\n现在您可以直接发送消息与该 Agent 对话。\n"
f"或使用 `/chat 您的消息` 命令。\n\n完成后可使用 `/disconnect` 断开连接。",
parse_mode="Markdown",
)
elif action == "pay":
did = data["did"]
agent = next((a for a in MOCK_AGENTS if a["did"] == did), None)
if not agent:
await query.edit_message_text("Agent 不存在。")
return
price_str = agent["price"]
amount_umsg = int(price_str.replace(" umsg", "").replace(",", ""))
amount_telegram = amount_umsg // 100
prices = [LabeledPrice(label=f"使用 {agent['name']}", amount=amount_telegram)]
await context.bot.send_invoice(
chat_id=query.message.chat_id,
title=f"支付 {agent['name']}",
description=f"为 AI Agent 服务付费\n\n{agent['description'][:100]}",
payload=f"pay_agent_{did}",
provider_token="", currency="XTR", prices=prices,
reply_markup=InlineKeyboardMarkup([
[InlineKeyboardButton("支付", pay=True)],
[InlineKeyboardButton("取消", callback_data=json.dumps({"action": "cancel_pay"}))],
]),
)
elif action == "cancel_pay":
await query.edit_message_text("支付已取消。")
elif action == "back_to_market":
page = 1
agents_per_page = 3
total_pages = (len(MOCK_AGENTS) + agents_per_page - 1) // agents_per_page
page_agents = MOCK_AGENTS[0:agents_per_page]
text_parts = [f"AI Agent 市场 (第 1/{total_pages} 页)\n"]
keyboard = []
for agent in page_agents:
text_parts.append(f"\n**{agent['name']}**\n分类: {agent['category']}\n评分: {agent['rating']}\n价格: {agent['price']}")
keyboard.append([InlineKeyboardButton(f" {agent['name'][:12]}", callback_data=json.dumps({"action": "detail", "did": agent["did"]}))])
if total_pages > 1:
keyboard.append([InlineKeyboardButton("下一页", callback_data=json.dumps({"action": "page", "page": 2}))])
await query.edit_message_text("\n".join(text_parts), parse_mode="Markdown", reply_markup=InlineKeyboardMarkup(keyboard))
elif action == "page":
page = data["page"]
agents_per_page = 3
total_pages = (len(MOCK_AGENTS) + agents_per_page - 1) // agents_per_page
start_idx = (page - 1) * agents_per_page
end_idx = start_idx + agents_per_page
page_agents = MOCK_AGENTS[start_idx:end_idx]
text_parts = [f"AI Agent 市场 (第 {page}/{total_pages} 页)\n"]
keyboard = []
for agent in page_agents:
text_parts.append(f"\n**{agent['name']}**\n分类: {agent['category']}\n评分: {agent['rating']}\n价格: {agent['price']}")
keyboard.append([InlineKeyboardButton(f" {agent['name'][:12]}", callback_data=json.dumps({"action": "detail", "did": agent["did"]}))])
nav_buttons = []
if page > 1:
nav_buttons.append(InlineKeyboardButton("上一页", callback_data=json.dumps({"action": "page", "page": page - 1})))
if page < total_pages:
nav_buttons.append(InlineKeyboardButton("下一页", callback_data=json.dumps({"action": "page", "page": page + 1})))
if nav_buttons:
keyboard.append(nav_buttons)
await query.edit_message_text("\n".join(text_parts), parse_mode="Markdown", reply_markup=InlineKeyboardMarkup(keyboard))
elif action == "confirm_pay":
user_id = str(query.from_user.id)
user_balances[user_id] = "9,500,000 umsg"
await query.edit_message_text(
"支付成功!\n\n您已成功支付 Agent 服务费用。\n交易已在 MSG Chain 上确认。\n\n使用 `/balance` 查看更新后的余额。",
parse_mode="Markdown",
)
async def pre_checkout_handler(update: Update, context: ContextTypes.DEFAULT_TYPE):
query = update.pre_checkout_query
payload = query.invoice_payload
if payload.startswith("pay_agent_"):
await query.answer(ok=True)
else:
await query.answer(ok=False, error_message="无效的支付请求。")
async def successful_payment_handler(update: Update, context: ContextTypes.DEFAULT_TYPE):
payload = update.message.successful_payment.invoice_payload
agent_did = payload.replace("pay_agent_", "")
user_id = str(update.effective_user.id)
user_sessions[user_id] = agent_did
await update.message.reply_text(
f"支付成功!\n\n您已成功连接到 Agent `{agent_did[:16]}...`\n"
f"现在可以发送消息与之交互了。\n\n使用 `/help` 查看更多命令。",
parse_mode="Markdown",
)
async def text_message_handler(update: Update, context: ContextTypes.DEFAULT_TYPE):
user_id = str(update.effective_user.id)
agent_did = user_sessions.get(user_id)
if not agent_did:
await update.message.reply_text("请先选择一个 Agent。\n使用 `/market` 浏览可用 Agent。")
return
message_text = update.message.text
agent = next((a for a in MOCK_AGENTS if a["did"] == agent_did), None)
agent_name = agent["name"] if agent else agent_did[:16]
async def simulate_typing():
await update.message.chat.send_action(action="typing")
await asyncio.sleep(1.5)
await simulate_typing()
reply_text = (
f"*{agent_name} 回复*\n\n收到: \"{message_text[:100]}\"\n\n"
f"这是一个模拟回复。在实际部署中,此消息将通过 Agent API 转发到 MSG Chain 上的 AI Agent。"
)
await update.message.reply_text(reply_text, parse_mode="Markdown")
def main():
if not BOT_TOKEN:
logger.error("请设置 TELEGRAM_BOT_TOKEN 环境变量")
sys.exit(1)
app = Application.builder().token(BOT_TOKEN).build()
app.add_handler(CommandHandler("start", start))
app.add_handler(CommandHandler("help", help_command))
app.add_handler(CommandHandler("market", market))
app.add_handler(CommandHandler("agent", agent_detail))
app.add_handler(CommandHandler("balance", balance))
app.add_handler(CommandHandler("chat", chat_with_agent))
app.add_handler(CommandHandler("disconnect", disconnect))
app.add_handler(CommandHandler("about", about))
app.add_handler(CallbackQueryHandler(button_handler))
app.add_handler(PreCheckoutQueryHandler(pre_checkout_handler))
app.add_handler(MessageHandler(filters.SUCCESSFUL_PAYMENT, successful_payment_handler))
app.add_handler(MessageHandler(filters.TEXT & ~filters.COMMAND, text_message_handler))
logger.info("Telegram Agent 市场机器人启动中...")
app.run_polling(allowed_updates=Update.ALL_TYPES)
if __name__ == "__main__":
main()
3.4 Telegram 内联键盘设计
Telegram 的内联键盘 (InlineKeyboardMarkup) 是实现 Agent 交互的核心 UI 组件:
# 单行多按钮(水平排列)
InlineKeyboardMarkup([
[InlineKeyboardButton("确认", callback_data='{"action":"confirm"}'),
InlineKeyboardButton("取消", callback_data='{"action":"cancel"}')],
])
# 多行单按钮(垂直列表)
InlineKeyboardMarkup([
[InlineKeyboardButton(f" {agent['name']}", callback_data=json.dumps({"action": "select", "did": agent["did"]}))]
for agent in agents
])
# 带导航的分页
InlineKeyboardMarkup([
[InlineKeyboardButton("上一页", callback_data=json.dumps({"page": p - 1})),
InlineKeyboardButton("下一页", callback_data=json.dumps({"page": p + 1}))],
[InlineKeyboardButton("返回", callback_data='{"action":"menu"}')],
])
# 支付按钮
InlineKeyboardMarkup([
[InlineKeyboardButton("支付 XTR", pay=True)],
[InlineKeyboardButton("取消", callback_data='{"action":"cancel"}')],
])
4. Discord Bot 集成
4.1 依赖安装
pip install discord-py>=2.3.0 httpx>=0.27
4.2 Discord 连接器实现
import json
import logging
from typing import Optional
import discord
from discord import app_commands
from discord.ext import commands
from bot_middleware import (
PlatformConnector, PlatformType, NormalizedMessage, NormalizedResponse, MessageType,
)
logger = logging.getLogger(__name__)
class DiscordConnector(PlatformConnector):
\"\"\"Discord 平台连接器\"\"\"
def __init__(self, bot_token: str):
super().__init__(PlatformType.DISCORD, bot_token)
self._bot: Optional[commands.Bot] = None
self._tree: Optional[app_commands.CommandTree] = None
async def start(self):
intents = discord.Intents.default()
intents.message_content = True
self._bot = commands.Bot(command_prefix="!", intents=intents)
self._tree = app_commands.CommandTree(self._bot)
await self._register_slash_commands()
await self._bot.start(self.bot_token)
async def stop(self):
if self._bot:
await self._bot.close()
logger.info("Discord Bot 已停止")
async def _register_slash_commands(self):
@self._tree.command(name="market", description="浏览 Agent 市场")
async def market(interaction: discord.Interaction):
msg = NormalizedMessage(
platform=PlatformType.DISCORD,
platform_id=str(interaction.id),
channel_id=str(interaction.channel_id),
user_id=str(interaction.user.id),
user_name=interaction.user.name,
message_type=MessageType.COMMAND,
text="/market",
)
await self.handle_message(msg)
await interaction.response.defer()
@self._tree.command(name="agent", description="查看 Agent 详情")
@app_commands.describe(did="Agent 的 DID 地址")
async def agent(interaction: discord.Interaction, did: str):
msg = NormalizedMessage(
platform=PlatformType.DISCORD,
platform_id=str(interaction.id),
channel_id=str(interaction.channel_id),
user_id=str(interaction.user.id),
user_name=interaction.user.name,
message_type=MessageType.COMMAND,
text=f"/agent {did}",
)
await self.handle_message(msg)
await interaction.response.defer()
@self._tree.command(name="balance", description="查询钱包余额")
async def balance(interaction: discord.Interaction):
msg = NormalizedMessage(
platform=PlatformType.DISCORD,
platform_id=str(interaction.id),
channel_id=str(interaction.channel_id),
user_id=str(interaction.user.id),
user_name=interaction.user.name,
message_type=MessageType.COMMAND,
text="/balance",
)
await self.handle_message(msg)
await interaction.response.defer()
@self._tree.command(name="session", description="查看当前会话")
async def session(interaction: discord.Interaction):
msg = NormalizedMessage(
platform=PlatformType.DISCORD,
platform_id=str(interaction.id),
channel_id=str(interaction.channel_id),
user_id=str(interaction.user.id),
user_name=interaction.user.name,
message_type=MessageType.COMMAND,
text="/session",
)
await self.handle_message(msg)
await interaction.response.defer()
logger.info("Discord 斜杠命令已注册")
async def send_message(self, response: NormalizedResponse) -> None:
if not self._bot:
return
channel = self._bot.get_channel(int(response.channel_id))
if not channel:
try:
user = await self._bot.fetch_user(int(response.user_id))
channel = user.dm_channel or await user.create_dm()
except Exception:
logger.error(f"无法找到 Discord 频道/用户: {response.channel_id}")
return
embed = None
if response.embeds:
embed_data = response.embeds[0]
embed = discord.Embed(
title=embed_data.get("title", ""),
description=embed_data.get("description", ""),
color=embed_data.get("color", 0x00AFFF),
)
if "fields" in embed_data:
for field in embed_data["fields"]:
embed.add_field(name=field.get("name", ""), value=field.get("value", ""), inline=field.get("inline", False))
if "footer" in embed_data:
embed.set_footer(text=embed_data["footer"])
if "thumbnail" in embed_data:
embed.set_thumbnail(url=embed_data["thumbnail"])
view = None
if response.buttons:
view = discord.ui.View(timeout=180)
for btn in response.buttons:
btn_data = json.dumps(btn.get("data", {}), separators=(",", ":"))
button = discord.ui.Button(label=btn.get("text", "按钮"), style=discord.ButtonStyle.primary, custom_id=f"agent:{btn_data}")
view.add_item(button)
kwargs = {}
if response.text:
kwargs["content"] = response.text[:2000]
if embed:
kwargs["embed"] = embed
if view:
kwargs["view"] = view
try:
if isinstance(channel, discord.TextChannel):
await channel.send(**kwargs)
elif hasattr(channel, "send"):
await channel.send(**kwargs)
except Exception as e:
logger.error(f"Discord 发送消息失败: {e}")
async def handle_interaction(self, interaction: discord.Interaction):
if not interaction.data or "custom_id" not in interaction.data:
return
custom_id = interaction.data["custom_id"]
if not custom_id.startswith("agent:"):
return
try:
data = json.loads(custom_id[6:])
except json.JSONDecodeError:
return
msg = NormalizedMessage(
platform=PlatformType.DISCORD,
platform_id=str(interaction.id),
channel_id=str(interaction.channel_id or ""),
user_id=str(interaction.user.id),
user_name=interaction.user.name,
message_type=MessageType.CALLBACK, text=None, data=data, raw=interaction,
)
await self.handle_message(msg)
4.3 Discord Bot 完整示例: Agent 发现机器人
#!/usr/bin/env python3
\"\"\"
Discord Agent 发现机器人 -- 完整示例
用户可以通过斜杠命令和按钮发现 MSG Chain 上的 AI Agent。
\"\"\"
import json
import logging
import os
import sys
import discord
from discord import app_commands
from discord.ext import commands
from dotenv import load_dotenv
load_dotenv()
logging.basicConfig(level=logging.INFO, format="%(asctime)s [%(levelname)s] %(name)s: %(message)s")
logger = logging.getLogger(__name__)
BOT_TOKEN = os.environ.get("DISCORD_BOT_TOKEN", "")
GUILD_ID = os.environ.get("DISCORD_GUILD_ID")
MOCK_AGENTS = [
{"did": "msg1agenttextgen0000000000000001", "name": "文本生成助手",
"description": "基于 GPT 的高质量文本生成 Agent。", "price": "500000 umsg", "rating": 4.8, "tasks_completed": 1247, "category": "text"},
{"did": "msg1agentimagedfx00000000000002", "name": "图像生成大师",
"description": "AI 图像生成 Agent。", "price": "1000000 umsg", "rating": 4.6, "tasks_completed": 892, "category": "image"},
{"did": "msg1agentdataanal000000000000003", "name": "数据分析专家",
"description": "数据分析 Agent。", "price": "800000 umsg", "rating": 4.9, "tasks_completed": 2156, "category": "data"},
{"did": "msg1agentcodegenp000000000000004", "name": "代码生成器",
"description": "全栈代码生成 Agent。", "price": "600000 umsg", "rating": 4.7, "tasks_completed": 3451, "category": "code"},
{"did": "msg1agenttransla000000000000005", "name": "翻译通",
"description": "多语言翻译 Agent。", "price": "300000 umsg", "rating": 4.5, "tasks_completed": 5678, "category": "translate"},
]
user_sessions: dict[str, str] = {}
user_balances: dict[str, str] = {}
class AgentView(discord.ui.View):
\"\"\"Agent 市场分页视图\"\"\"
def __init__(self, agents: list[dict], page: int = 1, per_page: int = 3):
super().__init__(timeout=180)
self.agents = agents
self.page = page
self.per_page = per_page
self.total_pages = max(1, (len(agents) + per_page - 1) // per_page)
self._update_buttons()
def _update_buttons(self):
self.clear_items()
total = len(self.agents)
start = (self.page - 1) * self.per_page
end = min(start + self.per_page, total)
page_agents = self.agents[start:end]
for agent in page_agents:
btn = discord.ui.Button(label=f" {agent['name'][:20]}", style=discord.ButtonStyle.secondary,
custom_id=f"agent:{json.dumps({'action': 'detail', 'did': agent['did']})}")
self.add_item(btn)
if self.total_pages > 1:
row = discord.ui.Button(label="上一页", style=discord.ButtonStyle.primary,
disabled=self.page <= 1,
custom_id=f"agent:{json.dumps({'action': 'page', 'page': self.page - 1})}")
self.add_item(row)
row2 = discord.ui.Button(label="下一页", style=discord.ButtonStyle.primary,
disabled=self.page >= self.total_pages,
custom_id=f"agent:{json.dumps({'action': 'page', 'page': self.page + 1})}")
self.add_item(row2)
class AgentDetailView(discord.ui.View):
\"\"\"Agent 详情页视图\"\"\"
def __init__(self, agent: dict):
super().__init__(timeout=180)
self.agent = agent
self.add_item(discord.ui.Button(label="连接到 Agent", style=discord.ButtonStyle.success,
custom_id=f"agent:{json.dumps({'action': 'connect', 'did': agent['did']})}"))
self.add_item(discord.ui.Button(label="支付使用", style=discord.ButtonStyle.primary,
custom_id=f"agent:{json.dumps({'action': 'pay', 'did': agent['did']})}"))
self.add_item(discord.ui.Button(label="返回市场", style=discord.ButtonStyle.secondary,
custom_id=f"agent:{json.dumps({'action': 'back'})}"))
intents = discord.Intents.default()
intents.message_content = True
bot = commands.Bot(command_prefix="!", intents=intents)
tree = bot.tree
async def sync_commands():
if GUILD_ID:
guild = discord.Object(id=int(GUILD_ID))
tree.copy_global_to(guild=guild)
await tree.sync(guild=guild)
else:
await tree.sync()
logger.info("Discord 斜杠命令已同步")
@bot.event
async def on_ready():
logger.info(f"Discord Bot 已登录: {bot.user}")
await sync_commands()
@tree.command(name="market", description="浏览 MSG Chain AI Agent 市场")
async def market(interaction: discord.Interaction):
await interaction.response.defer()
view = AgentView(MOCK_AGENTS)
total = len(MOCK_AGENTS)
embed = discord.Embed(title="MSG Chain AI Agent 市场",
description=f"发现并连接链上的 AI Agent。当前共有 **{total}** 个可用 Agent。", color=0x00AFFF)
embed.set_footer(text=f"第 {view.page}/{view.total_pages} 页 | MSG Chain")
start = (view.page - 1) * view.per_page
end = min(start + view.per_page, total)
for agent in MOCK_AGENTS[start:end]:
stars = "star" * int(agent["rating"])
embed.add_field(name=f"{agent['name']}",
value=f"`{agent['did'][:16]}...`\n{agent['description'][:80]}\n{stars} {agent['rating']} | 价格: {agent['price']} | 任务: {agent['tasks_completed']}",
inline=False)
await interaction.followup.send(embed=embed, view=view)
@tree.command(name="agent", description="查看 AI Agent 详细信息")
@app_commands.describe(did="Agent 的 MSG Chain DID 地址")
async def agent_detail(interaction: discord.Interaction, did: str):
await interaction.response.defer()
agent = next((a for a in MOCK_AGENTS if a["did"] == did), None)
if not agent:
await interaction.followup.send(f"未找到 Agent: `{did}`", ephemeral=True)
return
embed = discord.Embed(title=f"{agent['name']}", description=agent["description"], color=0x00FF88)
embed.add_field(name="DID", value=f"`{agent['did']}`", inline=False)
embed.add_field(name="分类", value=agent["category"], inline=True)
embed.add_field(name="价格", value=agent["price"], inline=True)
embed.add_field(name="评分", value=f"{agent['rating']}/5.0", inline=True)
embed.add_field(name="完成任务", value=str(agent["tasks_completed"]), inline=True)
embed.set_footer(text="MSG Chain AI Agent")
view = AgentDetailView(agent)
await interaction.followup.send(embed=embed, view=view)
@tree.command(name="balance", description="查询 MSG Chain 钱包余额")
async def balance(interaction: discord.Interaction):
await interaction.response.defer()
user_id = str(interaction.user.id)
bal = user_balances.get(user_id, "10,000,000 umsg")
embed = discord.Embed(title="钱包余额", description=f"`{bal}`", color=0xFFD700)
embed.add_field(name="地址", value="`msg1userwalletexample00000000000000`", inline=False)
embed.add_field(name="网络", value="MSG Chain Mainnet", inline=True)
embed.add_field(name="链 ID", value="msg-chain-1", inline=True)
await interaction.followup.send(embed=embed, ephemeral=True)
@tree.command(name="chat", description="与已连接的 Agent 对话")
@app_commands.describe(message="发送给 Agent 的消息")
async def chat(interaction: discord.Interaction, message: str):
await interaction.response.defer(ephemeral=True)
user_id = str(interaction.user.id)
agent_did = user_sessions.get(user_id)
if not agent_did:
await interaction.followup.send("您尚未连接 Agent。请使用 `/market` 选择一个。", ephemeral=True)
return
agent = next((a for a in MOCK_AGENTS if a["did"] == agent_did), None)
agent_name = agent["name"] if agent else agent_did[:16]
embed = discord.Embed(title=f"{agent_name} 回复",
description=f"收到消息: \"{message[:200]}\"\n\n这是一个模拟回复。实际部署时将调用 Agent API。", color=0x00FF88)
embed.set_footer(text=f"{agent_did[:16]}... | MSG Chain")
await interaction.followup.send(embed=embed, ephemeral=False)
@tree.command(name="session", description="查看当前 Agent 会话")
async def session(interaction: discord.Interaction):
await interaction.response.defer(ephemeral=True)
user_id = str(interaction.user.id)
agent_did = user_sessions.get(user_id)
if not agent_did:
await interaction.followup.send("当前没有活跃会话。使用 `/market` 连接一个 Agent。", ephemeral=True)
return
agent = next((a for a in MOCK_AGENTS if a["did"] == agent_did), None)
name = agent["name"] if agent else agent_did[:16]
embed = discord.Embed(title="当前会话", description=f"已连接到 **{name}**", color=0x7289DA)
embed.add_field(name="Agent DID", value=f"`{agent_did}`", inline=False)
embed.add_field(name="状态", value="活跃", inline=True)
await interaction.followup.send(embed=embed, ephemeral=True)
@bot.event
async def on_interaction(interaction: discord.Interaction):
if not interaction.data or interaction.data.get("type", 0) != 3:
return
custom_id = interaction.data.get("custom_id", "")
if not custom_id.startswith("agent:"):
return
await interaction.response.defer(ephemeral=True)
data = json.loads(custom_id[6:])
action = data.get("action", "")
if action == "detail":
did = data["did"]
agent = next((a for a in MOCK_AGENTS if a["did"] == did), None)
if not agent:
await interaction.followup.send("Agent 不存在。", ephemeral=True)
return
embed = discord.Embed(title=f"{agent['name']}", description=agent["description"], color=0x00FF88)
embed.add_field(name="DID", value=f"`{agent['did']}`", inline=False)
embed.add_field(name="分类", value=agent["category"], inline=True)
embed.add_field(name="价格", value=agent["price"], inline=True)
embed.add_field(name="评分", value=f"{agent['rating']}/5.0", inline=True)
embed.add_field(name="完成任务", value=str(agent["tasks_completed"]), inline=True)
view = AgentDetailView(agent)
await interaction.followup.send(embed=embed, view=view, ephemeral=False)
elif action == "connect":
did = data["did"]
user_id = str(interaction.user.id)
user_sessions[user_id] = did
agent = next((a for a in MOCK_AGENTS if a["did"] == did), None)
name = agent["name"] if agent else did[:16]
await interaction.followup.send(f"已连接到 **{name}**!\n现在可以使用 `/chat 您的消息` 与 Agent 对话。", ephemeral=False)
elif action == "pay":
did = data["did"]
agent = next((a for a in MOCK_AGENTS if a["did"] == did), None)
if not agent:
await interaction.followup.send("Agent 不存在。", ephemeral=True)
return
embed = discord.Embed(title="支付确认", description=f"您是否要支付 **{agent['price']}** 使用 {agent['name']}?", color=0xFFD700)
view = discord.ui.View(timeout=120)
view.add_item(discord.ui.Button(label=f"确认支付 {agent['price']}", style=discord.ButtonStyle.success,
custom_id=f"agent:{json.dumps({'action': 'confirm_pay', 'did': did})}"))
view.add_item(discord.ui.Button(label="取消", style=discord.ButtonStyle.danger,
custom_id=f"agent:{json.dumps({'action': 'cancel_pay'})}"))
await interaction.followup.send(embed=embed, view=view, ephemeral=True)
elif action == "confirm_pay":
did = data["did"]
user_id = str(interaction.user.id)
user_sessions[user_id] = did
bal = user_balances.get(user_id, "10,000,000 umsg")
amount = 500000
new_bal_num = int(bal.replace(" umsg", "").replace(",", "")) - amount
user_balances[user_id] = f"{new_bal_num:,} umsg"
await interaction.followup.send(f"**支付成功!**\n已连接到 Agent。\n余额: `{user_balances[user_id]}`", ephemeral=False)
elif action == "cancel_pay":
await interaction.followup.send("支付已取消。", ephemeral=True)
elif action == "back":
view = AgentView(MOCK_AGENTS)
embed = discord.Embed(title="MSG Chain AI Agent 市场",
description=f"发现并连接链上的 AI Agent。当前共有 **{len(MOCK_AGENTS)}** 个可用 Agent。", color=0x00AFFF)
embed.set_footer(text=f"第 {view.page}/{view.total_pages} 页 | MSG Chain")
start = (view.page - 1) * view.per_page
end = min(start + view.per_page, len(MOCK_AGENTS))
for agent in MOCK_AGENTS[start:end]:
embed.add_field(name=agent["name"], value=f"`{agent['did'][:16]}...`\n{agent['description'][:80]}", inline=False)
await interaction.followup.send(embed=embed, view=view, ephemeral=False)
elif action == "page":
page = data.get("page", 1)
view = AgentView(MOCK_AGENTS, page=page)
embed = discord.Embed(title="MSG Chain AI Agent 市场", description="发现并连接链上的 AI Agent。", color=0x00AFFF)
embed.set_footer(text=f"第 {view.page}/{view.total_pages} 页 | MSG Chain")
start = (view.page - 1) * view.per_page
end = min(start + view.per_page, len(MOCK_AGENTS))
for agent in MOCK_AGENTS[start:end]:
embed.add_field(name=agent["name"], value=f"`{agent['did'][:16]}...`\n{agent['description'][:80]}", inline=False)
await interaction.followup.send(embed=embed, view=view, ephemeral=False)
def main():
if not BOT_TOKEN:
logger.error("请设置 DISCORD_BOT_TOKEN 环境变量")
sys.exit(1)
logger.info("Discord Agent 发现机器人启动中...")
bot.run(BOT_TOKEN)
if __name__ == "__main__":
main()
4.4 Discord Embed 设计模式
# Agent 资料卡片
embed = discord.Embed(title=f"{agent['name']}", description=agent['description'], color=0x00FF88,
url=f"https://msgchain.org/agents/{agent['did']}")
embed.add_field(name="DID", value=f"`{agent['did']}`", inline=False)
embed.add_field(name="价格", value=agent['price'], inline=True)
embed.add_field(name="评分", value=f"{agent['rating']}/5.0", inline=True)
embed.set_footer(text=f"Agent ID: {agent['did'][:12]}... | MSG Chain")
embed.set_thumbnail(url="https://msgchain.org/logo.png")
# 交易状态卡片
embed = discord.Embed(title="交易状态", description=f"交易哈希: `{tx_hash}`", color=0x00AFFF)
embed.add_field(name="发送方", value=f"`{from_did[:16]}...`")
embed.add_field(name="接收方", value=f"`{to_did[:16]}...`")
embed.add_field(name="金额", value=f"{amount} umsg")
5. Slack Bot 集成
5.1 依赖安装
pip install slack-bolt>=1.18.0 httpx>=0.27
5.2 Slack 连接器实现
import json
import logging
from typing import Optional
from slack_bolt import App
from slack_bolt.adapter.socket_mode import SocketModeHandler
from bot_middleware import (
PlatformConnector, PlatformType, NormalizedMessage, NormalizedResponse, MessageType,
)
logger = logging.getLogger(__name__)
class SlackConnector(PlatformConnector):
\"\"\"Slack 平台连接器\"\"\"
def __init__(self, bot_token: str, app_token: Optional[str] = None):
super().__init__(PlatformType.SLACK, bot_token)
self.app_token = app_token or ""
self._app: Optional[App] = None
self._handler: Optional[SocketModeHandler] = None
async def start(self):
self._app = App(token=self.bot_token)
self._register_handlers()
self._handler = SocketModeHandler(self._app, self.app_token)
await asyncio.to_thread(self._handler.start)
logger.info("Slack Bot 已启动")
async def stop(self):
if self._handler:
self._handler.close()
logger.info("Slack Bot 已停止")
def _register_handlers(self):
app = self._app
if not app:
return
@app.command("/agent-market")
async def market_command(ack, body, client):
await ack()
msg = NormalizedMessage(platform=PlatformType.SLACK, platform_id=body.get("trigger_id", ""),
channel_id=body.get("channel_id", ""), user_id=body.get("user_id", ""),
user_name=body.get("user_name", "unknown"), message_type=MessageType.COMMAND,
text="/agent-market", raw=body)
await self.handle_message(msg)
@app.command("/agent-lookup")
async def lookup_command(ack, body, client, command):
await ack()
query = command.get("text", "")
msg = NormalizedMessage(platform=PlatformType.SLACK, platform_id=body.get("trigger_id", ""),
channel_id=body.get("channel_id", ""), user_id=body.get("user_id", ""),
user_name=body.get("user_name", "unknown"), message_type=MessageType.COMMAND,
text=f"/agent-lookup {query}", raw=body)
await self.handle_message(msg)
@app.command("/balance")
async def balance_command(ack, body, client):
await ack()
msg = NormalizedMessage(platform=PlatformType.SLACK, platform_id=body.get("trigger_id", ""),
channel_id=body.get("channel_id", ""), user_id=body.get("user_id", ""),
user_name=body.get("user_name", "unknown"), message_type=MessageType.COMMAND,
text="/balance", raw=body)
await self.handle_message(msg)
@app.action(regex="agent_.*")
async def block_action(ack, body, client):
await ack()
actions = body.get("actions", [])
if not actions:
return
action = actions[0]
try:
data = json.loads(action.get("value", "{}"))
except json.JSONDecodeError:
data = {}
msg = NormalizedMessage(platform=PlatformType.SLACK, platform_id=body.get("trigger_id", ""),
channel_id=body.get("channel", {}).get("id", ""),
user_id=body.get("user", {}).get("id", ""),
user_name=body.get("user", {}).get("name", "unknown"),
message_type=MessageType.CALLBACK, data=data, raw=body)
await self.handle_message(msg)
@app.view("agent_pay_modal")
async def view_submission(ack, body, client, view):
values = view.get("state", {}).get("values", {})
agent_did = ""
amount = ""
for block_id, block in values.items():
for action_id, action_data in block.items():
if action_id == "agent_did_input":
agent_did = action_data.get("value", "")
elif action_id == "amount_input":
amount = action_data.get("value", "")
await ack()
msg = NormalizedMessage(platform=PlatformType.SLACK, platform_id=body.get("trigger_id", ""),
channel_id=body.get("channel", {}).get("id", body.get("user", {}).get("id", "")),
user_id=body.get("user", {}).get("id", ""),
user_name=body.get("user", {}).get("name", "unknown"),
message_type=MessageType.PAYMENT, text=f"pay {agent_did} {amount}",
data={"agent_did": agent_did, "amount": amount}, raw=body)
await self.handle_message(msg)
async def send_message(self, response: NormalizedResponse) -> None:
if not self._app or not self._app.client:
return
client = self._app.client
blocks = []
if response.text:
blocks.append({"type": "section", "text": {"type": "mrkdwn", "text": response.text}})
if response.buttons:
elements = []
for btn in response.buttons:
elements.append({"type": "button", "text": {"type": "plain_text", "text": btn.get("text", "按钮")},
"value": json.dumps(btn.get("data", {}), separators=(",", ":")),
"action_id": f"agent_btn_{btn.get('text', 'unknown')[:8]}"})
blocks.append({"type": "actions", "elements": elements})
if response.modal:
blocks.append({"type": "section", "text": {"type": "mrkdwn", "text": "打开支付窗口:"}})
blocks.append({"type": "actions", "elements": [{"type": "button", "text": {"type": "plain_text", "text": "创建支付"},
"action_id": "open_pay_modal", "value": json.dumps(response.modal)}]})
channel = response.channel_id or response.user_id
kwargs = {"channel": channel, "text": response.text or " "}
if blocks:
kwargs["blocks"] = blocks
if response.ephemeral:
kwargs["response_type"] = "ephemeral"
try:
await asyncio.to_thread(client.chat_postMessage, **kwargs)
except Exception as e:
logger.error(f"Slack 发送消息失败: {e}")
5.3 Slack Bot 完整示例: Agent 交互机器人
#!/usr/bin/env python3
\"\"\"
Slack Agent 交互机器人 -- 完整示例
使用 Slack Bolt 框架和 Block Kit 实现 Agent 发现与支付。
\"\"\"
import json
import logging
import os
import sys
from dotenv import load_dotenv
from slack_bolt import App
from slack_bolt.adapter.socket_mode import SocketModeHandler
load_dotenv()
logging.basicConfig(level=logging.INFO, format="%(asctime)s [%(levelname)s] %(name)s: %(message)s")
logger = logging.getLogger(__name__)
SLACK_BOT_TOKEN = os.environ.get("SLACK_BOT_TOKEN", "")
SLACK_APP_TOKEN = os.environ.get("SLACK_APP_TOKEN", "")
MOCK_AGENTS = [
{"did": "msg1agenttextgen0000000000000001", "name": "文本生成助手", "description": "基于 GPT 的高质量文本生成 Agent。",
"price": "500000 umsg", "rating": 4.8, "tasks": 1247},
{"did": "msg1agentimagedfx00000000000002", "name": "图像生成大师", "description": "AI 图像生成 Agent。",
"price": "1000000 umsg", "rating": 4.6, "tasks": 892},
{"did": "msg1agentcodegenp000000000000004", "name": "代码生成器", "description": "全栈代码生成 Agent。",
"price": "600000 umsg", "rating": 4.7, "tasks": 3451},
{"did": "msg1agenttransla000000000000005", "name": "翻译通", "description": "多语言翻译 Agent。",
"price": "300000 umsg", "rating": 4.5, "tasks": 5678},
]
user_sessions: dict[str, str] = {}
user_balances: dict[str, str] = {}
app = App(token=SLACK_BOT_TOKEN)
def build_market_blocks(page=1, per_page=2):
total = len(MOCK_AGENTS)
total_pages = max(1, (total + per_page - 1) // per_page)
start = (page - 1) * per_page
end = min(start + per_page, total)
page_agents = MOCK_AGENTS[start:end]
blocks = [{"type": "header", "text": {"type": "plain_text", "text": "MSG Chain AI Agent 市场", "emoji": True}},
{"type": "section", "text": {"type": "mrkdwn", "text": f"发现并连接链上的 AI Agent。当前共有 *{total}* 个可用 Agent。\n第 {page}/{total_pages} 页"}},
{"type": "divider"}]
for agent in page_agents:
stars = "star" * int(agent["rating"])
blocks.append({"type": "section", "text": {"type": "mrkdwn", "text": f"*{agent['name']}*\n`{agent['did'][:20]}...`\n{agent['description']}\n{stars} {agent['rating']} | {agent['price']} | {agent['tasks']} 任务"},
"accessory": {"type": "button", "text": {"type": "plain_text", "text": "详情", "emoji": True},
"value": json.dumps({"action": "detail", "did": agent["did"]}, separators=(",", ":")),
"action_id": f"agent_detail_{agent['did'][:8]}"}})
blocks.append({"type": "divider"})
nav_buttons = []
if page > 1:
nav_buttons.append({"type": "button", "text": {"type": "plain_text", "text": "上一页", "emoji": True},
"value": json.dumps({"action": "page", "page": page - 1}, separators=(",", ":")), "action_id": "agent_page_prev"})
if page < total_pages:
nav_buttons.append({"type": "button", "text": {"type": "plain_text", "text": "下一页", "emoji": True},
"value": json.dumps({"action": "page", "page": page + 1}, separators=(",", ":")), "action_id": "agent_page_next"})
if nav_buttons:
blocks.append({"type": "actions", "elements": nav_buttons})
blocks.append({"type": "context", "elements": [{"type": "mrkdwn", "text": "使用 `/agent-lookup <名称>` 搜索 Agent | `/balance` 查询余额"}]})
return blocks
def build_agent_detail_blocks(agent):
stars = "star" * int(agent["rating"])
return [{"type": "header", "text": {"type": "plain_text", "text": f"{agent['name']}", "emoji": True}},
{"type": "divider"},
{"type": "section", "fields": [{"type": "mrkdwn", "text": f"*DID:*\n`{agent['did']}`"}, {"type": "mrkdwn", "text": f"*价格:*\n{agent['price']}"},
{"type": "mrkdwn", "text": f"*评分:*\n{stars} {agent['rating']}/5.0"}, {"type": "mrkdwn", "text": f"*完成任务:*\n{agent['tasks']}"}]},
{"type": "section", "text": {"type": "mrkdwn", "text": f"*描述:*\n{agent['description']}"}},
{"type": "divider"},
{"type": "actions", "elements": [
{"type": "button", "text": {"type": "plain_text", "text": "连接 Agent", "emoji": True}, "style": "primary",
"value": json.dumps({"action": "connect", "did": agent["did"]}, separators=(",", ":")), "action_id": f"agent_connect_{agent['did'][:8]}"},
{"type": "button", "text": {"type": "plain_text", "text": "支付", "emoji": True},
"value": json.dumps({"action": "open_pay", "did": agent["did"]}, separators=(",", ":")), "action_id": f"agent_pay_{agent['did'][:8]}"},
{"type": "button", "text": {"type": "plain_text", "text": "返回市场", "emoji": True},
"value": json.dumps({"action": "back"}, separators=(",", ":")), "action_id": "agent_back"},
]}]
@app.command("/agent-market")
def handle_market(ack, body, client):
ack()
blocks = build_market_blocks()
client.chat_postMessage(channel=body.get("channel_id", ""), blocks=blocks, text="AI Agent 市场")
@app.command("/agent-lookup")
def handle_lookup(ack, command, client, body):
ack()
query = command.get("text", "").strip()
channel_id = body.get("channel_id", "")
if not query:
client.chat_postMessage(channel=channel_id, text="请提供搜索关键词。用法: `/agent-lookup <名称或关键词>`")
return
matches = [a for a in MOCK_AGENTS if query.lower() in a["name"].lower()]
if not matches:
client.chat_postMessage(channel=channel_id, text=f"未找到匹配 \"{query}\" 的 Agent。")
return
blocks = [{"type": "section", "text": {"type": "mrkdwn", "text": f"找到 *{len(matches)}* 个 Agen[未公开路径]"}}, {"type": "divider"}]
for agent in matches:
blocks.append({"type": "section", "text": {"type": "mrkdwn", "text": f"*{agent['name']}* -- {agent['price']}\n`{agent['did']}`"},
"accessory": {"type": "button", "text": {"type": "plain_text", "text": "详情"},
"value": json.dumps({"action": "detail", "did": agent["did"]}), "action_id": f"agent_detail_{agent['did'][:8]}"}})
client.chat_postMessage(channel=channel_id, blocks=blocks, text=f"搜索: {query}")
@app.command("/balance")
def handle_balance(ack, body, client):
ack()
user_id = body.get("user_id", "")
channel_id = body.get("channel_id", "")
bal = user_balances.get(user_id, "10,000,000 umsg")
blocks = [{"type": "section", "text": {"type": "mrkdwn", "text": f"钱包余额\n`{bal}`"}},
{"type": "context", "elements": [{"type": "mrkdwn", "text": "地址: `msg1userwalletexample00000000000000`"}]}]
client.chat_postMessage(channel=channel_id, blocks=blocks, text=f"余额: {bal}")
@app.action(regex="agent_.*")
def handle_actions(ack, body, client, action):
ack()
try:
value = action.get("value", "{}")
data = json.loads(value)
except json.JSONDecodeError:
logger.error(f"无法解析 action value: {action.get('value', '')}")
return
action_type = data.get("action", "")
channel_id = body.get("channel", {}).get("id", "")
user_id = body.get("user", {}).get("id", "")
if action_type == "detail":
did = data.get("did", "")
agent = next((a for a in MOCK_AGENTS if a["did"] == did), None)
if not agent:
client.chat_postMessage(channel=channel_id, text="Agent 不存在。")
return
client.chat_postMessage(channel=channel_id, blocks=build_agent_detail_blocks(agent), text=f"Agent: {agent['name']}")
elif action_type == "connect":
did = data.get("did", "")
user_sessions[user_id] = did
agent = next((a for a in MOCK_AGENTS if a["did"] == did), None)
name = agent["name"] if agent else did[:16]
client.chat_postMessage(channel=channel_id, text=f"已连接到 *{name}*!\n直接发送消息即可与 Agent 对话。")
elif action_type == "open_pay":
did = data.get("did", "")
agent = next((a for a in MOCK_AGENTS if a["did"] == did), None)
if not agent:
client.chat_postMessage(channel=channel_id, text="Agent 不存在。")
return
trigger_id = body.get("trigger_id", "")
client.views_open(trigger_id=trigger_id, view={
"type": "modal", "callback_id": "agent_pay_modal",
"title": {"type": "plain_text", "text": f"支付 {agent['name']}", "emoji": True},
"submit": {"type": "plain_text", "text": "确认支付", "emoji": True},
"close": {"type": "plain_text", "text": "取消", "emoji": True},
"blocks": [{"type": "section", "text": {"type": "mrkdwn", "text": f"向 *{agent['name']}* 支付服务费用。\n标准价格: {agent['price']}"}},
{"type": "input", "block_id": "agent_input",
"element": {"type": "plain_text_input", "action_id": "agent_did_input", "initial_value": did,
"placeholder": {"type": "plain_text", "text": "Agent DID"}},
"label": {"type": "plain_text", "text": "Agent DID"}},
{"type": "input", "block_id": "amount_input_block",
"element": {"type": "plain_text_input", "action_id": "amount_input",
"initial_value": agent["price"].replace(" umsg", ""),
"placeholder": {"type": "plain_text", "text": "金额 (umsg)"}},
"label": {"type": "plain_text", "text": "支付金额 (umsg)"}}]})
elif action_type == "back":
client.chat_postMessage(channel=channel_id, blocks=build_market_blocks(), text="返回市场")
elif action_type == "page":
page = data.get("page", 1)
client.chat_postMessage(channel=channel_id, blocks=build_market_blocks(page=page), text=f"市场第 {page} 页")
@app.view("agent_pay_modal")
def handle_pay_modal(ack, body, client, view):
values = view.get("state", {}).get("values", {})
agent_did = ""
amount = ""
for block_id, block in values.items():
for action_id, action_data in block.items():
if action_id == "agent_did_input":
agent_did = action_data.get("value", "")
elif action_id == "amount_input":
amount = action_data.get("value", "")
ack()
user_id = body.get("user", {}).get("id", "")
channel_id = body.get("channel", {}).get("id", body.get("user", {}).get("id", ""))
if not agent_did or not amount:
client.chat_postMessage(channel=channel_id, text="支付信息不完整,请重试。")
return
user_sessions[user_id] = agent_did
current_bal = user_balances.get(user_id, "10,000,000 umsg")
amount_int = int(amount.replace(",", ""))
new_bal = max(0, int(current_bal.replace(" umsg", "").replace(",", "")) - amount_int)
user_balances[user_id] = f"{new_bal:,} umsg"
blocks = [{"type": "section", "text": {"type": "mrkdwn", "text": f"*支付成功!*\n\n向 `{agent_did[:20]}...` 支付了 `{amount} umsg`\n剩余余额: `{user_balances[user_id]}`\n\n现在您已连接到该 Agent。"}}]
client.chat_postMessage(channel=channel_id, blocks=blocks, text="支付成功")
@app.event("message")
def handle_message(body, client):
event = body.get("event", {})
if event.get("type") != "message" or event.get("subtype") == "bot_message":
return
text = event.get("text", "").strip()
if text.startswith("/"):
return
user_id = event.get("user", "")
channel_id = event.get("channel", "")
agent_did = user_sessions.get(user_id)
if not agent_did:
client.chat_postMessage(channel=channel_id, text="您尚未连接任何 Agent。请使用 `/agent-market` 浏览并选择一个 Agent。")
return
agent = next((a for a in MOCK_AGENTS if a["did"] == agent_did), None)
name = agent["name"] if agent else agent_did[:16]
blocks = [{"type": "section", "text": {"type": "mrkdwn", "text": f"*{name}* 回复:\n\n收到消息: \"{text[:150]}\"\n\n(此为模拟回复,实际部署时将调用 Agent API)"}}]
client.chat_postMessage(channel=channel_id, blocks=blocks, text=f"{name} 回复")
def main():
if not SLACK_BOT_TOKEN or not SLACK_APP_TOKEN:
logger.error("请设置 SLACK_BOT_TOKEN 和 SLACK_APP_TOKEN 环境变量")
sys.exit(1)
logger.info("Slack Agent 交互机器人启动中...")
handler = SocketModeHandler(app, SLACK_APP_TOKEN)
handler.start()
if __name__ == "__main__":
main()
5.4 Slack Block Kit 设计模式
blocks = [
{"type": "header", "text": {"type": "plain_text", "text": "Agent 名称"}},
{"type": "divider"},
{"type": "section", "text": {"type": "mrkdwn", "text": "描述文本"}},
{"type": "section", "fields": [{"type": "mrkdwn", "text": "*字段1:*\n值1"}, {"type": "mrkdwn", "text": "*字段2:*\n值2"}]},
{"type": "actions", "elements": [
{"type": "button", "text": {"type": "plain_text", "text": "按钮1"}, "value": "...", "action_id": "..."},
{"type": "button", "text": {"type": "plain_text", "text": "按钮2"}, "style": "primary", "value": "...", "action_id": "..."},
]},
{"type": "context", "elements": [{"type": "mrkdwn", "text": "脚注信息"}]},
]
6. 统一消息路由
6.1 消息归一化
from dataclasses import dataclass, field
from enum import Enum
from typing import Any, Optional
class MessageType(Enum):
TEXT = "text"
COMMAND = "command"
CALLBACK = "callback"
PAYMENT = "payment"
FILE = "file"
IMAGE = "image"
class PlatformType(Enum):
TELEGRAM = "telegram"
DISCORD = "discord"
SLACK = "slack"
@dataclass
class NormalizedMessage:
platform: PlatformType
platform_id: str
channel_id: str
user_id: str
user_name: str
message_type: MessageType
text: Optional[str] = None
data: Optional[dict] = None
attachments: list[dict] = field(default_factory=list)
timestamp: float = field(default_factory=time.time)
raw: Optional[Any] = None
def to_dict(self) -> dict:
return {"platform": self.platform.value, "platform_id": self.platform_id,
"channel_id": self.channel_id, "user_id": self.user_id, "user_name": self.user_name,
"message_type": self.message_type.value, "text": self.text, "data": self.data,
"attachments": self.attachments, "timestamp": self.timestamp}
@dataclass
class NormalizedResponse:
platform: PlatformType
channel_id: str
user_id: str
text: Optional[str] = None
embeds: list[dict] = field(default_factory=list)
buttons: list[dict] = field(default_factory=list)
modal: Optional[dict] = None
ephemeral: bool = False
raw: Optional[Any] = None
6.2 会话管理
class SessionManager:
\"\"\"跨平台会话管理器\"\"\"
def __init__(self):
self._sessions: dict[str, dict] = {}
self._agent_users: dict[str, list[str]] = {}
self._session_to_user: dict[str, str] = {}
def get_session(self, platform: PlatformType, user_id: str) -> Optional[dict]:
key = f"{platform.value}:{user_id}"
return self._sessions.get(key)
def create_session(self, platform: PlatformType, user_id: str, agent_did: str, session_id: str) -> dict:
key = f"{platform.value}:{user_id}"
info = {"agent_did": agent_did, "session_id": session_id, "platform": platform.value,
"user_id": user_id, "created_at": time.time(), "last_active": time.time()}
self._sessions[key] = info
self._session_to_user[session_id] = key
if agent_did not in self._agent_users:
self._agent_users[agent_did] = []
if key not in self._agent_users[agent_did]:
self._agent_users[agent_did].append(key)
return info
def remove_session(self, platform: PlatformType, user_id: str):
key = f"{platform.value}:{user_id}"
info = self._sessions.pop(key, None)
if info:
agent_users = self._agent_users.get(info["agent_did"], [])
if key in agent_users:
agent_users.remove(key)
self._session_to_user.pop(info["session_id"], None)
def get_users_for_agent(self, agent_did: str) -> list[dict]:
keys = self._agent_users.get(agent_did, [])
return [self._sessions[k] for k in keys if k in self._sessions]
def get_agent_for_user(self, platform: PlatformType, user_id: str) -> Optional[str]:
key = f"{platform.value}:{user_id}"
info = self._sessions.get(key)
return info["agent_did"] if info else None
def update_activity(self, platform: PlatformType, user_id: str):
key = f"{platform.value}:{user_id}"
if key in self._sessions:
self._sessions[key]["last_active"] = time.time()
def cleanup_stale(self, max_age: float = 3600):
now = time.time()
stale_keys = [k for k, v in self._sessions.items() if now - v["last_active"] > max_age]
for key in stale_keys:
self.remove_session(PlatformType(key.split(":")[0]), key.split(":", 1)[1])
return len(stale_keys)
6.3 多平台广播
async def broadcast_to_all(middleware: BotMiddleware, text: str, exclude_platform: Optional[str] = None):
response = NormalizedResponse(platform=PlatformType.TELEGRAM, channel_id="", user_id="", text=text)
for name, connector in middleware.platforms.items():
if name == exclude_platform:
continue
try:
response.platform = PlatformType(name)
await connector.send_message(response)
except Exception as e:
logger.error(f"广播到 {name} 失败: {e}")
async def broadcast_to_agent_users(middleware: BotMiddleware, agent_did: str, text: str, session_manager: SessionManager):
users = session_manager.get_users_for_agent(agent_did)
for user_info in users:
platform = PlatformType(user_info["platform"])
connector = middleware.platforms.get(platform.value)
if connector:
response = NormalizedResponse(platform=platform, channel_id=user_info["user_id"],
user_id=user_info["user_id"], text=text)
try:
await connector.send_message(response)
except Exception as e:
logger.error(f"广播到 {platform.value}:{user_info['user_id']} 失败: {e}")
6.4 速率限制
from collections import defaultdict
class RateLimiter:
\"\"\"滑动窗口速率限制器\"\"\"
def __init__(self):
self._windows: Dict[str, List[float]] = defaultdict(list)
def check(self, key: str, max_requests: int = 10, window_seconds: float = 60.0) -> bool:
now = time.time()
timestamps = self._windows[key]
timestamps[:] = [t for t in timestamps if now - t < window_seconds]
if len(timestamps) >= max_requests:
return False
timestamps.append(now)
return True
def remaining(self, key: str, max_requests: int = 10, window_seconds: float = 60.0) -> int:
now = time.time()
timestamps = self._windows[key]
timestamps[:] = [t for t in timestamps if now - t < window_seconds]
return max(0, max_requests - len(timestamps))
def reset(self, key: str):
self._windows[key] = []
def get_window_stats(self, key: str, window_seconds: float = 60.0) -> dict:
now = time.time()
timestamps = self._windows[key]
active = [t for t in timestamps if now - t < window_seconds]
return {"key": key, "active_requests": len(active), "window_seconds": window_seconds,
"oldest": min(active) if active else None, "newest": max(active) if active else None}
6.5 Agent API 会话管理增强
class AgentSessionManager:
\"\"\"Agent API 会话生命周期管理\"\"\"
def __init__(self, api_client: AgentAPIClient, max_concurrent: int = 50):
self.api = api_client
self.max_concurrent = max_concurrent
self._sessions: dict[str, dict] = {}
self._lock = asyncio.Lock()
async def get_or_create_session(self, user_key: str, agent_did: str) -> str:
async with self._lock:
if user_key in self._sessions:
session = self._sessions[user_key]
if session["agent_did"] == agent_did:
return session["session_id"]
if len(self._sessions) >= self.max_concurrent:
oldest_key = min(self._sessions, key=lambda k: self._sessions[k]["created_at"])
del self._sessions[oldest_key]
result = await self.api.create_session(agent_did)
session_id = result["session_id"]
self._sessions[user_key] = {"session_id": session_id, "agent_did": agent_did, "created_at": time.time()}
return session_id
async def send_and_receive(self, user_key: str, agent_did: str, content: str) -> str:
session_id = await self.get_or_create_session(user_key, agent_did)
result = await self.api.send_message(session_id, content)
return result.get("reply", result.get("content", ""))
async def close_session(self, user_key: str):
async with self._lock:
self._sessions.pop(user_key, None)
async def cleanup(self):
async with self._lock:
count = len(self._sessions)
self._sessions.clear()
return count
6.6 消息队列与异步处理
class MessageQueue:
\"\"\"异步消息队列 — 用于解耦消息接收与处理\"\"\"
def __init__(self, maxsize: int = 1000):
self._queue: asyncio.Queue = asyncio.Queue(maxsize=maxsize)
self._processing = False
self._stats = {"enqueued": 0, "processed": 0, "failed": 0, "dropped": 0}
async def enqueue(self, message: NormalizedMessage):
try:
await self._queue.put(message)
self._stats["enqueued"] += 1
except asyncio.QueueFull:
self._stats["dropped"] += 1
logger.warning(f"消息队列已满,丢弃消息: {message.platform_id}")
async def process_loop(self, handler, num_workers: int = 3):
self._processing = True
workers = [self._worker(handler) for _ in range(num_workers)]
await asyncio.gather(*workers)
async def _worker(self, handler):
while self._processing:
try:
message = await asyncio.wait_for(self._queue.get(), timeout=1.0)
try:
await handler(message)
self._stats["processed"] += 1
except Exception as e:
self._stats["failed"] += 1
logger.exception(f"消息处理失败: {e}")
finally:
self._queue.task_done()
except asyncio.TimeoutError:
continue
def stop(self):
self._processing = False
def stats(self) -> dict:
return {**self._stats, "queue_size": self._queue.qsize()}
6.7 统一错误处理
class ErrorHandler:
\"\"\"统一错误处理 — 将异常转换为用户友好的平台消息\"\"\"
ERROR_MESSAGES = {
"AgentNotFound": "未找到指定的 Agent。请检查 DID 地址是否正确。",
"SessionExpired": "会话已过期,请重新连接 Agent。",
"InsufficientBalance": "余额不足。请充值后再试。",
"RateLimited": "请求过于频繁,请稍后再试。",
"NetworkError": "网络连接异常,请检查网络后重试。",
"PaymentFailed": "支付失败,请稍后重试或联系支持。",
"Unauthorized": "未授权的操作。请确认您的身份。",
"AgentBusy": "Agent 正忙,请稍后再试。",
"InvalidMessage": "消息格式不正确,请重新输入。",
"UnknownError": "发生未知错误,请稍后重试。",
}
@staticmethod
def get_user_message(error: Exception) -> str:
error_type = type(error).__name__
return ErrorHandler.ERROR_MESSAGES.get(error_type, f"错误: {str(error)[:100]}")
@staticmethod
def is_retryable(error: Exception) -> bool:
retryable_types = {"NetworkError", "AgentBusy", "RateLimited", "TimeoutError"}
return type(error).__name__ in retryable_types
@staticmethod
async def handle_with_retry(func, max_retries: int = 3, base_delay: float = 1.0):
last_error = None
for attempt in range(max_retries):
try:
return await func()
except Exception as e:
last_error = e
if not ErrorHandler.is_retryable(e):
raise
if attempt < max_retries - 1:
delay = base_delay * (2 ** attempt)
logger.warning(f"重试 {attempt + 1}/{max_retries}, 等待 {delay}s: {e}")
await asyncio.sleep(delay)
raise last_error
6.8 完整消息路由流程
async def complete_routing_pipeline(middleware: BotMiddleware, msg: NormalizedMessage):
rate_limiter = RateLimiter()
user_key = f"{msg.platform.value}:{msg.user_id}"
if not rate_limiter.check(user_key):
await middleware._send_platform_message(NormalizedResponse(
platform=msg.platform, channel_id=msg.channel_id, user_id=msg.user_id,
text="请求过于频繁,请稍后再试。", ephemeral=True,
))
return
try:
if msg.message_type == MessageType.COMMAND:
await middleware._handle_command(msg)
elif msg.message_type == MessageType.CALLBACK:
await middleware._handle_callback(msg)
elif msg.message_type == MessageType.PAYMENT:
await middleware._handle_payment(msg)
else:
session_mgr = SessionManager()
agent_session = AgentSessionManager(middleware.agent_api)
agent_did = session_mgr.get_agent_for_user(msg.platform, msg.user_id)
if not agent_did:
result = await ErrorHandler.handle_with_retry(lambda: middleware.agent_api.list_agents(size=1))
agents = result.get("agents", [])
if not agents:
await middleware._send_platform_message(NormalizedResponse(
platform=msg.platform, channel_id=msg.channel_id, user_id=msg.user_id,
text="当前没有可用的 Agent。",
))
return
agent_did = agents[0]["did"]
reply = await agent_session.send_and_receive(user_key, agent_did, msg.text or "")
await middleware._send_platform_message(NormalizedResponse(
platform=msg.platform, channel_id=msg.channel_id, user_id=msg.user_id, text=reply,
))
except Exception as e:
logger.exception(f"消息路由异常: {e}")
await middleware._send_platform_message(NormalizedResponse(
platform=msg.platform, channel_id=msg.channel_id, user_id=msg.user_id,
text=ErrorHandler.get_user_message(e),
))
7. 安全考虑
7.1 Bot Token 管理
import base64
from cryptography.fernet import Fernet
from cryptography.hazmat.primitives import hashes
from cryptography.hazmat.primitives.kdf.pbkdf2 import PBKDF2
class TokenManager:
\"\"\"Bot Token 安全管理器\"\"\"
def __init__(self, master_key: Optional[bytes] = None):
if master_key:
self._fernet = Fernet(base64.urlsafe_b64encode(master_key[:32].ljust(32, b'\0')))
else:
key = os.environ.get("TOKEN_MASTER_KEY", "").encode()
if not key:
raise ValueError("必须设置 TOKEN_MASTER_KEY 环境变量")
kdf = PBKDF2(algorithm=hashes.SHA256(), length=32, salt=b"msgchain-token-salt", iterations=100000)
self._fernet = Fernet(base64.urlsafe_b64encode(kdf.derive(key)))
def encrypt_token(self, token: str) -> str:
return self._fernet.encrypt(token.encode()).decode()
def decrypt_token(self, encrypted: str) -> str:
return self._fernet.decrypt(encrypted.encode()).decode()
@staticmethod
def validate_token_format(token: str, platform: PlatformType) -> bool:
if platform == PlatformType.TELEGRAM:
parts = token.split(":")
return len(parts) == 2 and all(p.isalnum() for p in parts)
elif platform == PlatformType.DISCORD:
return len(token) > 50 and token.isprintable()
elif platform == PlatformType.SLACK:
return token.startswith("xoxb-") or token.startswith("xapp-")
return False
7.2 DID 认证
class DIDSecurityManager:
\"\"\"DID 签名验证与身份管理\"\"\"
def __init__(self, chain_rpc_url: str):
self.rpc_url = chain_rpc_url
self._public_key_cache: dict[str, tuple[bytes, float]] = {}
async def verify_agent_identity(self, agent_did: str, signature: str, payload: dict) -> bool:
public_key = await self._fetch_public_key(agent_did)
if not public_key:
return False
message = json.dumps(payload, sort_keys=True, separators=(",", ":")).encode()
try:
verify_key = ed25519.Ed25519PublicKey.from_public_bytes(public_key)
verify_key.verify(bytes.fromhex(signature), message)
return True
except Exception:
return False
async def _fetch_public_key(self, did: str) -> Optional[bytes]:
now = time.time()
if did in self._public_key_cache:
pub_key, expiry = self._public_key_cache[did]
if now < expiry:
return pub_key
try:
async with httpx.AsyncClient() as client:
resp = await client.get(f"{self.rpc_url}/cosmos/auth/v1beta1/accounts/{did}", timeout=10)
resp.raise_for_status()
data = resp.json()
pub_key_hex = data.get("account", {}).get("pub_key", {}).get("key", "")
if pub_key_hex:
pub_key = bytes.fromhex(pub_key_hex)
self._public_key_cache[did] = (pub_key, now + 300)
return pub_key
except Exception as e:
logger.error(f"获取公钥失败 {did}: {e}")
return None
def generate_challenge(self, platform: PlatformType, user_id: str) -> dict:
nonce = os.urandom(16).hex()
return {"type": "did_auth_challenge", "challenge": nonce, "platform": platform.value,
"user_id": user_id, "timestamp": int(time.time()), "expires_in": 120}
def verify_challenge_response(self, challenge: dict, response: dict, public_key_hex: str) -> bool:
if time.time() - challenge["timestamp"] > challenge["expires_in"]:
return False
if response.get("challenge") != challenge["challenge"]:
return False
signature = response.get("signature", "")
message = json.dumps(challenge, sort_keys=True, separators=(",", ":")).encode()
try:
verify_key = ed25519.Ed25519PublicKey.from_public_bytes(bytes.fromhex(public_key_hex))
verify_key.verify(bytes.fromhex(signature), message)
return True
except Exception:
return False
7.3 平台特定安全措施
class PlatformSecurity:
\"\"\"各平台安全最佳实践\"\"\"
@staticmethod
def telegram_safety(update_data: dict) -> bool:
update_id = update_data.get("update_id", 0)
if update_id < getattr(PlatformSecurity, "_last_telegram_update", 0):
logger.warning(f"检测到重复的 Telegram update_id: {update_id}")
return False
PlatformSecurity._last_telegram_update = update_id
return True
@staticmethod
def validate_slack_request(timestamp: str, signature: str, body: bytes, signing_secret: str) -> bool:
import hmac, hashlib
if not timestamp or not signature:
return False
try:
if abs(time.time() - int(timestamp)) > 300:
return False
except ValueError:
return False
sig_basestring = f"v0:{timestamp}:".encode() + body
expected = "v0=" + hmac.new(signing_secret.encode(), sig_basestring, hashlib.sha256).hexdigest()
return hmac.compare_digest(expected, signature)
@staticmethod
def sanitize_user_input(text: str, max_length: int = 2000) -> str:
cleaned = "".join(c for c in text if c.isprintable() or c in "\n\r\t")
return cleaned[:max_length]
@staticmethod
def rate_limit_by_platform(platform: PlatformType) -> tuple[int, float]:
limits = {PlatformType.TELEGRAM: (20, 60.0), PlatformType.DISCORD: (10, 60.0), PlatformType.SLACK: (15, 60.0)}
return limits.get(platform, (10, 60.0))
7.4 安全最佳实践清单
SECURITY_CHECKLIST = {
"token_management": [
"永远不要将 Bot Token 硬编码在代码中",
"使用环境变量或加密的配置文件存储 Token",
"定期轮换 Token(建议每 90 天)",
"为每个平台使用独立的 Bot Token",
],
"did_security": [
"私钥仅存储在内存中,绝不持久化到磁盘",
"使用 HSM 或安全飞地存储生产环境私钥",
"每次 API 请求都附带 DID 签名",
"实现挑战-响应认证机制",
],
"platform_specific": [
"Telegram: 验证 update_id 防重放",
"Discord: 使用交互令牌验证来源",
"Slack: 验证 X-Slack-Signature 和 X-Slack-Request-Timestamp",
],
"general": [
"所有用户输入需要进行清理和长度限制",
"实施速率限制防止滥用",
"记录所有鉴权失败事件用于审计",
"使用 HTTPS 加密所有 API 通信",
"定期审查权限和访问控制",
],
}
7.5 审计日志
import json
from datetime import datetime
class AuditLogger:
\"\"\"安全审计日志记录器\"\"\"
def __init__(self, log_path: str = "./audit.log"):
self.logger = logging.getLogger("audit")
handler = logging.FileHandler(log_path)
handler.setFormatter(logging.Formatter("%(asctime)s [AUDIT] %(message)s"))
self.logger.addHandler(handler)
self.logger.setLevel(logging.INFO)
def log_auth_event(self, user_id: str, platform: str, action: str, success: bool, reason: str = ""):
self.logger.info(json.dumps({"type": "auth_event", "user_id": user_id, "platform": platform,
"action": action, "success": success, "reason": reason,
"timestamp": datetime.utcnow().isoformat()}))
def log_payment_event(self, user_id: str, agent_did: str, amount: str, tx_hash: str, status: str):
self.logger.info(json.dumps({"type": "payment_event", "user_id": user_id, "agent_did": agent_did,
"amount": amount, "tx_hash": tx_hash, "status": status,
"timestamp": datetime.utcnow().isoformat()}))
def log_error_event(self, error_type: str, details: str, user_id: str = ""):
self.logger.info(json.dumps({"type": "error_event", "error_type": error_type, "details": details,
"user_id": user_id, "timestamp": datetime.utcnow().isoformat()}))
8. 部署指南
8.1 环境变量配置
# .env 文件示例
# MSG Chain
AGENT_DID=msg1agentdefault00000000000000000
AGENT_DID_PRIVATE_KEY=<hex_encoded_private_key>
AGENT_API_URL=https://api.msgchain.org/v1
# Telegram
TELEGRAM_BOT_TOKEN=1234567890:ABCdefGHIjklMNOpqrsTUVwxyz
# Discord
DISCORD_BOT_TOKEN=discord_bot_token_here
DISCORD_GUILD_ID=123456789012345678
# Slack
SLACK_BOT_TOKEN=xoxb-123456789012-...
SLACK_APP_TOKEN=xapp-123456789012-...
# Security (optional)
TOKEN_MASTER_KEY=<32_bytes_hex_encoded>
8.2 Docker 部署
FROM python:3.11-slim
WORKDIR /app
COPY requirements.txt .
RUN pip install --no-cache-dir -r requirements.txt
COPY . .
EXPOSE 8080
CMD ["python", "main.py"]
8.3 requirements.txt
python-telegram-bot>=21.0
discord-py>=2.3.0
slack-bolt>=1.18.0
httpx>=0.27.0
python-dotenv>=1.0.0
cryptography>=41.0.0
8.4 监控与健康检查
class HealthCheck:
\"\"\"系统健康状态检查\"\"\"
def __init__(self, middleware: BotMiddleware):
self.middleware = middleware
self.start_time = time.time()
async def check(self) -> dict:
statuses = {}
for name, connector in self.middleware.platforms.items():
statuses[name] = "running"
return {"status": "healthy", "uptime": time.time() - self.start_time,
"platforms": statuses, "active_sessions": len(self.middleware.session_map)}
async def check_agent_api(self) -> dict:
try:
result = await self.middleware.agent_api.list_agents(size=1)
return {"status": "ok", "agents_available": len(result.get("agents", []))}
except Exception as e:
return {"status": "error", "message": str(e)}
8.5 启动方式
# 直接启动
python main.py
# Docker 启动
docker build -t agent-bot-middleware .
docker run -d --env-file .env --name agent-bot agent-bot-middleware
# docker-compose.yml
version: "3.9"
services:
agent-bot:
build: .
env_file: .env
restart: unless-stopped
logging:
driver: "json-file"
options:
max-size: "10m"
max-file: "3"
8.6 故障排除
| 问题 | 可能原因 | 解决方案 |
|---|---|---|
| Bot 无法启动 | Token 无效或过期 | 检查环境变量中的 Token 是否正确 |
| 消息未收到 | Webhook 未配置 | 确保正确设置了 Bot 的 Webhook URL |
| 支付失败 | 余额不足 | 用户 MSG Chain 钱包余额不足,提示充值 |
| Agent 无响应 | Agent API 不可用 | 检查 Agent API 服务和 MSG Chain 节点状态 |
| DID 签名验证失败 | 私钥不匹配 | 确认 DID 对应的私钥是否正确 |
| 速率限制触发 | 请求过于频繁 | 调整 max_requests 参数或告知用户等待 |
本文档由 AI 辅助生成,用于指导 MSG Chain AI Agent 的跨平台机器人集成开发。
链 ID: msg-chain-1 | 地址前缀: msg1
