AI Agent 第三方服务集成指南
适用链: MSG Chain | 地址前缀: msg | 版本: 1.0.0
⚠️ No-Go Disclaimer: MSGChain 主网裁决为 No-Go。本文件所有内容反映的是开发阶段的技术设计,不代表主网未独立核验上线状态。生产部署状态请以白皮书为准:https://msgchain.org/whitepaper/
目录
- 概述
- Twitter/X 集成
- Discord 集成
- Telegram 集成
- GitHub 集成
- Slack 集成
- Email (SMTP/IMAP) 集成
- Webhook 系统
- 跨平台消息路由
- 完整示例:多平台监控 Agent
- 附录
1. 概述
1.1 为什么需要第三方服务集成
AI Agent 在 MSG Chain 上运行,天然具备链上数据访问和智能合约交互能力。然而,现实世界的用户和系统主要活动在链下——社交媒体、即时通讯、代码仓库、邮件等。要让 Agent 真正有用,必须打通这些渠道。
核心价值:
| 领域 | 价值 | 典型场景 |
|---|---|---|
| 社交 | 触达用户、品牌曝光 | 自动发推、回复提及、话题监控 |
| 通讯 | 实时交互、客服 | Discord/Telegram 机器人、通知推送 |
| 开发 | 代码审查、DevOps | PR 审查、Issue 管理、CI/CD 触发 |
| 数据 | 信息采集、分析 | 链上数据推送、市场监控、预警 |
1.2 服务分类
第三方服务
├── 社交平台 (Social)
│ ├── Twitter/X → 发帖、监控、搜索
│ └── (其他社交平台)
├── 通讯平台 (Communication)
│ ├── Discord → Bot、Slash Commands、频道管理
│ ├── Telegram → Bot、Inline Keyboard、支付
│ ├── Slack → Webhook、Block Kit、App
│ └── Email → SMTP 发送、IMAP 接收
├── 开发工具 (Dev Tools)
│ ├── GitHub → Issue、PR、Code Review
│ └── GitLab / Bitbucket → (类似集成)
└── 数据/Webhook
├── 自定义 Webhook → 接收外部事件
└── 链上桥接 → 链下 → 链上事件中继
1.3 架构总览
┌─────────────────────────────────────────────┐
│ AI Agent (MSG Chain) │
│ ┌─────────────────────────────────────┐ │
│ │ Agent Core │ │
│ │ ┌──────────┐ ┌─────────────────┐ │ │
│ │ │ LLM │ │ Chain Client │ │ │
│ │ │ (推理) │ │ (msg1... 地址) │ │ │
│ │ └──────────┘ └─────────────────┘ │ │
│ └─────────────────────────────────────┘ │
│ │ │ │ │
│ ┌──────┴─────────┴──────────┴──────┐ │
│ │ Service Adapter Layer │ │
│ │ ┌──────┐ ┌──────┐ ┌──────┐ │ │
│ │ │Twit. │ │Disc. │ │Tele. │ ... │ │
│ │ └──────┘ └──────┘ └──────┘ │ │
│ └──────────────────────────────────┘ │
└─────────────────────────────────────────────┘
│ │ │
┌────┴─────────┴──────────┴────┐
│ External APIs │
│ Twitter Discord Telegram │
│ GitHub Slack Email │
└──────────────────────────────┘
1.4 通用安全实践
- API 密钥/令牌通过环境变量或机密管理服务注入,绝不硬编码
- 所有出站请求使用 HTTPS/TLS
- 为每个 Agent 使用独立的 API 凭据,最小权限原则
- 请求频率受控 (Rate Limiting),避免被服务封禁
- 敏感操作(如转账)需多签确认
2. Twitter/X 集成
2.1 准备工作
- 在 Twitter Developer Portal 创建 Project
- 启用 OAuth 2.0 + OAuth 1.0a (User Context)
- 生成 API Key、API Secret、Bearer Token、Access Token & Secret
- 设置环境变量:
export TWITTER_BEARER_TOKEN="your_bearer_token"
export TWITTER_API_KEY="your_api_key"
export TWITTER_API_SECRET="your_api_secret"
export TWITTER_ACCESS_TOKEN="your_access_token"
export TWITTER_ACCESS_SECRET="your_access_secret"
2.2 依赖安装
pip install tweepy==4.14.0 httpx python-dotenv
2.3 完整 Twitter/X Agent 实现
import os
import asyncio
import logging
from typing import Optional, List, Dict
from datetime import datetime, timedelta
import tweepy
from dotenv import load_dotenv
from msg_chain_sdk import MsgChainClient, AgentLog
load_dotenv()
logger = logging.getLogger(__name__)
class TwitterAgent:
"""Twitter/X integration for AI Agent on MSG Chain.
Provides tweet posting, mention monitoring, topic search,
thread creation, and automated response capabilities.
"""
def __init__(
self,
agent_id: str,
chain_rpc: str = "https://rpc.msgchain.io",
chain_id: str = "msg-chain-1",
):
bearer_token = os.getenv("TWITTER_BEARER_TOKEN")
consumer_key = os.getenv("TWITTER_API_KEY")
consumer_secret = os.getenv("TWITTER_API_SECRET")
access_token = os.getenv("TWITTER_ACCESS_TOKEN")
access_token_secret = os.getenv("TWITTER_ACCESS_SECRET")
if not all([bearer_token, consumer_key, consumer_secret,
access_token, access_token_secret]):
raise EnvironmentError(
"Missing Twitter API credentials. Set TWITTER_* env vars."
)
self.client = tweepy.Client(
bearer_token=bearer_token,
consumer_key=consumer_key,
consumer_secret=consumer_secret,
access_token=access_token,
access_token_secret=access_token_secret,
wait_on_rate_limit=True,
)
self.agent_id = agent_id
self.last_mention_id: Optional[str] = None
# Chain client for on-chain logging
self.chain = MsgChainClient(rpc=chain_rpc, chain_id=chain_id)
# Rate limiting
self._tweet_count = 0
self._tweet_reset_time = datetime.utcnow()
def _check_rate_limit(self):
"""Ensure we stay within Twitter's rate limits (300 tweets/3h)."""
now = datetime.utcnow()
if now - self._tweet_reset_time > timedelta(hours=3):
self._tweet_count = 0
self._tweet_reset_time = now
if self._tweet_count >= 280:
raise RuntimeError("Rate limit approaching: 280/300 tweets posted this window")
self._tweet_count += 1
async def agent_log(self, action: str, data: dict):
"""Record agent action on MSG Chain."""
log = AgentLog(
agent_id=self.agent_id,
action=action,
data=data,
timestamp=datetime.utcnow().isoformat(),
)
tx_hash = await self.chain.submit_log(log)
logger.info(f"On-chain log submitted: {tx_hash}")
return tx_hash
async def post_tweet(self, text: str) -> tweepy.Response:
"""Post a tweet. Text must be <= 280 characters.
Args:
text: Tweet content (will be truncated to 280).
Returns:
tweepy.Response with tweet data.
Example:
>>> agent = TwitterAgent("agent-001")
>>> resp = await agent.post_tweet("Hello from MSG Chain AI Agent! #msgchain")
>>> print(resp.data['id'])
"""
self._check_rate_limit()
tweet_text = text[:280]
tweet = self.client.create_tweet(text=tweet_text)
logger.info(f"Tweet posted: {tweet.data['id']}")
# Log on chain
await self.agent_log("twitter_post", {
"tweet_id": tweet.data["id"],
"text_preview": tweet_text[:50],
})
return tweet
async def post_thread(self, tweets: List[str]) -> List[tweepy.Response]:
"""Post a thread of multiple tweets.
Each tweet links to the previous one via reply.
Args:
tweets: List of tweet texts (each <= 280 chars).
Returns:
List of tweet responses in order.
Example:
>>> agent = TwitterAgent("agent-001")
>>> thread = [
... "Thread: AI Agent on MSG Chain (1/3)",
... "Agents can autonomously post to Twitter (2/3)",
... "And monitor mentions too! (3/3) #msgchain",
... ]
>>> await agent.post_thread(thread)
"""
if not tweets:
return []
results = []
last_tweet_id = None
for i, text in enumerate(tweets):
self._check_rate_limit()
tweet_text = text[:280]
kwargs = {"text": tweet_text}
if last_tweet_id:
kwargs["in_reply_to_tweet_id"] = last_tweet_id
tweet = self.client.create_tweet(**kwargs)
results.append(tweet)
last_tweet_id = tweet.data["id"]
await asyncio.sleep(0.5)
await self.agent_log("twitter_thread", {
"tweet_ids": [t.data["id"] for t in results],
"length": len(tweets),
})
return results
async def monitor_mentions(self, since_id: Optional[str] = None):
"""Monitor @mentions of the authenticated user.
Analyzes each mention with LLM and auto-replies.
Args:
since_id: Only get mentions after this tweet ID.
If None, uses self.last_mention_id.
Returns:
List of (mention, response_text) tuples.
"""
user_id = self.client.get_me().data.id
mention_kwargs = {
"id": user_id,
"since_id": since_id or self.last_mention_id,
"tweet_fields": ["created_at", "author_id"],
}
mentions = self.client.get_users_mentions(**mention_kwargs)
results = []
if not mentions.data:
return results
for mention in mentions.data:
author_id = mention.author_id
if str(author_id) == str(user_id):
continue
response_text = await self._generate_mention_response(
mention.text, mention.id
)
if response_text:
reply = self.client.create_tweet(
text=response_text[:280],
in_reply_to_tweet_id=mention.id,
)
results.append((mention, reply))
self.last_mention_id = mention.id
await self.agent_log("twitter_mentions_processed", {
"count": len(results),
"since_id": self.last_mention_id,
})
return results
async def _generate_mention_response(
self, mention_text: str, mention_id: str
) -> Optional[str]:
"""Generate a context-aware response to a mention."""
mention_lower = mention_text.lower()
if "balance" in mention_lower or "余额" in mention_lower:
balance = await self.chain.query_balance(self.agent_id)
return (
f"@{self.agent_id} current balance: {balance} MSG\n"
f"Powered by MSG Chain AI Agent"
)
elif "help" in mention_lower or "帮助" in mention_lower:
return (
f"@{self.agent_id} I ca[未公开路径]"
f" - Query balance (mention 'balance')\n"
f" - Search topics (mention 'search <topic>')\n"
f" - On-chain data analysis\n"
f"#msgchain"
)
elif "search" in mention_lower or "搜索" in mention_lower:
query = mention_text.split("search", 1)[-1].strip()
query = query.split("搜索", 1)[-1].strip()
if query:
results = await self.search_topic(query, max_results=3)
if results:
urls = "\n".join(
f"https://twitter.com/i/web/status/{t['id']}"
for t in results
)
return f"Latest tweets about '{query}':\n{urls}"
return f"No tweets found about '{query}'."
return "Please specify a search keyword. Example: search MSG Chain"
return (
f"Mention received!\n"
f"I am an AI Agent running on MSG Chain (ID: {self.agent_id}).\n"
f"Type 'help' to see what I can do."
)
async def search_topic(
self, query: str, max_results: int = 10, language: str = "en"
) -> List[Dict]:
"""Search recent tweets on a topic.
Args:
query: Search query string.
max_results: Max tweets to return (10-100).
language: Tweet language filter.
Returns:
List of tweet objects.
"""
tweets = self.client.search_recent_tweets(
query=query,
max_results=min(max_results, 100),
language=language,
tweet_fields=["created_at", "public_metrics", "author_id"],
)
result = []
if tweets.data:
result = [
{
"id": t.id,
"text": t.text,
"author_id": t.author_id,
"created_at": str(t.created_at),
"likes": t.public_metrics.get("like_count", 0),
"retweets": t.public_metrics.get("retweet_count", 0),
}
for t in tweets.data
]
await self.agent_log("twitter_search", {
"query": query,
"results_count": len(result),
})
return result
async def like_tweet(self, tweet_id: str) -> bool:
"""Like a tweet by ID."""
user_id = self.client.get_me().data.id
try:
self.client.like(tweet_id=tweet_id, user_id=user_id)
await self.agent_log("twitter_like", {"tweet_id": tweet_id})
return True
except Exception as e:
logger.error(f"Failed to like tweet {tweet_id}: {e}")
return False
async def retweet(self, tweet_id: str) -> bool:
"""Retweet a tweet by ID."""
user_id = self.client.get_me().data.id
try:
self.client.retweet(tweet_id=tweet_id, user_id=user_id)
await self.agent_log("twitter_retweet", {"tweet_id": tweet_id})
return True
except Exception as e:
logger.error(f"Failed to retweet {tweet_id}: {e}")
return False
async def get_user_timeline(
self, username: str, max_results: int = 10
) -> List[Dict]:
"""Get tweets from a specific user's timeline."""
user = self.client.get_user(username=username)
if not user.data:
return []
tweets = self.client.get_users_tweets(
id=user.data.id,
max_results=min(max_results, 100),
tweet_fields=["created_at", "public_metrics"],
)
if not tweets.data:
return []
return [
{
"id": t.id,
"text": t.text,
"created_at": str(t.created_at),
"metrics": dict(t.public_metrics) if t.public_metrics else {},
}
for t in tweets.data
]
async def run_monitor_loop(self, interval: int = 60):
"""Run mention monitoring loop continuously."""
logger.info(f"Starting Twitter mention monitor (interval={interval}s)")
while True:
try:
results = await self.monitor_mentions()
if results:
logger.info(f"Processed {len(results)} mentions")
except Exception as e:
logger.error(f"Monitor error: {e}")
await asyncio.sleep(interval)
class TwitterStreamingAgent(TwitterAgent):
"""Extension of TwitterAgent with real-time streaming support."""
def __init__(self, agent_id: str, **kwargs):
super().__init__(agent_id, **kwargs)
self.stream_rules: List[Dict] = []
self.stream = None
async def add_stream_rule(self, value: str, tag: str = ""):
"""Add a filtered stream rule."""
self.stream_rules.append({"value": value, "tag": tag})
rule_payload = {"add": [{"value": value, "tag": tag}]}
self.client.create_stream_rule(rule_payload)
logger.info(f"Stream rule added: {value}")
async def start_stream(self):
"""Start the filtered stream and process tweets in real-time."""
if not self.stream_rules:
raise RuntimeError("No stream rules configured.")
class AgentStreamListener(tweepy.StreamingClient):
def __init__(self, bearer_token, agent):
super().__init__(bearer_token)
self.agent = agent
def on_tweet(self, tweet):
asyncio.create_task(self.agent._on_stream_tweet(tweet))
def on_error(self, status_code):
logger.error(f"Stream error: {status_code}")
return False
self.stream = AgentStreamListener(
os.getenv("TWITTER_BEARER_TOKEN"), self
)
self.stream.filter(
tweet_fields=["created_at", "author_id", "public_metrics"]
)
async def _on_stream_tweet(self, tweet):
"""Process a tweet received from the stream."""
logger.info(f"Stream received tweet {tweet.id}: {tweet.text[:50]}")
if "msgchain" in tweet.text.lower() or "msg chain" in tweet.text.lower():
await self.like_tweet(tweet.id)
await self.agent_log("twitter_stream_received", {
"tweet_id": tweet.id,
"text_preview": tweet.text[:80],
})
# Example usage
if __name__ == "__main__":
async def main():
agent = TwitterAgent(agent_id="msg-agent-001")
await agent.post_tweet(
"MSG Chain AI Agent is live! "
"On-chain agents can now interact with Twitter autonomously. "
"#msgchain #AIAgent"
)
results = await agent.search_topic("MSG Chain AI", max_results=5)
for r in results:
print(f"[{r['created_at']}] {r['text'][:60]}")
await agent.monitor_mentions()
asyncio.run(main())
2.4 Twitter API 注意事项
| 限制项 | 数值 | 处理方式 |
|---|---|---|
| 推文字数 | 280 字符 | 自动截断 |
| 发推频率 | 300/3h | _check_rate_limit() |
| 搜索频率 | 450/15min | wait_on_rate_limit=True |
| 流规则数 | 25 | 超出需删除旧规则 |
3. Discord 集成
3.1 准备工作
- 在 Discord Developer Portal 创建 Application
- 创建 Bot,获取 Token
- 在 OAuth2 URL Generator 选择 bot + applications.commands 权限
- 设置 Bot 权限: Send Messages, Read Message History, Mention Everyone, Use Slash Commands
- 将 Bot 邀请到服务器
export DISCORD_BOT_TOKEN="your_bot_token"
export DISCORD_GUILD_ID="your_guild_id"
export DISCORD_CLIENT_ID="your_client_id"
3.2 依赖安装
pip install discord.py==2.3.2 aiohttp
3.3 完整 Discord Agent 实现
import os
import re
import asyncio
import logging
from typing import Optional, List, Dict, Set
from datetime import datetime
import discord
from discord import app_commands
from discord.ext import commands, tasks
from msg_chain_sdk import MsgChainClient
load_dotenv()
logger = logging.getLogger(__name__)
class DiscordAgent(discord.Client):
"""Discord bot integration for AI Agent on MSG Chain.
Handles messages, slash commands, role-based access,
and automated channel management.
"""
def __init__(
self,
agent_id: str,
allowed_roles: Optional[List[str]] = None,
command_prefix: str = "!",
chain_rpc: str = "https://rpc.msgchain.io",
):
intents = discord.Intents.default()
intents.message_content = True
intents.members = True
intents.guilds = True
super().__init__(intents=intents)
self.agent_id = agent_id
self.command_prefix = command_prefix
self.allowed_roles = allowed_roles or ["admin", "moderator", "agent-operator"]
self.chain = MsgChainClient(rpc=chain_rpc)
self.tree = app_commands.CommandTree(self)
self.monitored_channels: Set[int] = set()
async def setup_hook(self):
"""Called when the bot is ready to set up extensions."""
self._register_commands()
self.background_task.start()
async def on_ready(self):
"""Handle bot ready event."""
logger.info(
f"Agent {self.agent_id} logged in as {self.user} "
f"(ID: {self.user.id})"
)
logger.info(f"Connected to {len(self.guilds)} guilds")
await self.tree.sync()
logger.info("Slash commands synced")
async def on_message(self, message: discord.Message):
"""Handle incoming messages."""
if message.author == self.user:
return
if message.author.bot:
return
await self._log_message(message)
if self.user in message.mentions:
async with message.channel.typing():
response = await self._process_mention(message)
await message.reply(response, mention_author=False)
return
if message.content.startswith(self.command_prefix):
cmd_text = message.content[len(self.command_prefix):].strip()
await self._handle_text_command(cmd_text, message)
return
if message.channel.id in self.monitored_channels:
if self._should_auto_respond(message):
async with message.channel.typing():
response = await self._auto_respond(message)
await message.reply(response, mention_author=False)
def _register_commands(self):
"""Register slash commands."""
@self.tree.command(
name="query",
description="Query AI Agent data from MSG Chain"
)
@app_commands.describe(query="Your query string")
async def query(interaction: discord.Interaction, query: str):
await self._handle_slash_query(interaction, query)
@self.tree.command(
name="balance",
description="Check agent or wallet balance on MSG Chain"
)
@app_commands.describe(address="Optional: MSG Chain address")
async def balance(
interaction: discord.Interaction,
address: Optional[str] = None,
):
await self._handle_slash_balance(interaction, address)
@self.tree.command(
name="monitor",
description="Start/stop monitoring a channel"
)
@app_commands.describe(action="start or stop", channel="Target channel")
@app_commands.checks.has_permissions(manage_channels=True)
async def monitor(
interaction: discord.Interaction,
action: str,
channel: Optional[discord.TextChannel] = None,
):
await self._handle_slash_monitor(interaction, action, channel)
@self.tree.command(
name="broadcast",
description="Broadcast a message to all monitored channels"
)
@app_commands.checks.has_role("admin")
async def broadcast(
interaction: discord.Interaction,
message: str,
):
await self._handle_slash_broadcast(interaction, message)
@self.tree.command(
name="status",
description="Show agent status and health information"
)
async def status(interaction: discord.Interaction):
await self._handle_slash_status(interaction)
def _check_role_access(self, member: discord.Member) -> bool:
"""Check if a member has the required roles."""
if not self.allowed_roles:
return True
member_roles = {role.name.lower() for role in member.roles}
return any(role in member_roles for role in self.allowed_roles)
async def _process_mention(self, message: discord.Message) -> str:
"""Process a message that @mentioned the bot."""
content = re.sub(rf'<@!?{self.user.id}>', '', message.content).strip()
if not content:
return (
f"**MSG Chain AI Agent**\n"
f"Agent ID: `{self.agent_id}`\n"
f"Type `/help` to see available commands"
)
if await self._is_rate_limited(message.author.id):
return "Rate limit exceeded. Please try again later."
response = await self._llm_inference(content, message)
await self._log_interaction("mention", message.author.id, content, response)
return response
async def _llm_inference(self, prompt: str, context) -> str:
"""Perform LLM inference (rule-based demo)."""
prompt_lower = prompt.lower()
if "balance" in prompt_lower or "余额" in prompt_lower:
balance = await self.chain.query_balance(self.agent_id)
return (
f"Agent `{self.agent_id}` current balanc[未公开路径]"
f"**{balance} MSG**\n"
f"(Chain: MSG Chain)"
)
if "help" in prompt_lower or "帮助" in prompt_lower:
return (
"**Available Commands**\n"
"`/query <text>` - Query on-chain data\n"
"`/balance [address]` - Check balance\n"
"`/monitor start/stop` - Channel monitoring\n"
"`/broadcast <msg>` - Broadcast (admin)\n"
"`/status` - Agent status\n"
"You can also @ me to ask questions!"
)
return (
f"Message received!\n"
f"> {prompt[:200]}\n\n"
f"Agent `{self.agent_id}` is processing..."
)
async def _handle_text_command(self, cmd: str, message: discord.Message):
"""Handle text-based commands."""
parts = cmd.split()
if not parts:
return
command = parts[0].lower()
args = parts[1:]
if command == "agent":
subcmd = args[0] if args else "status"
if subcmd == "status":
await message.channel.send(
f"Agent `{self.agent_id}` running\n"
f"Connected guilds: {len(self.guilds)}"
)
elif subcmd == "invite" and self._check_role_access(message.author):
invite = await message.channel.create_invite(max_age=3600)
await message.channel.send(f"Invite: {invite}")
elif subcmd == "broadcast" and len(args) >= 2:
if self._check_role_access(message.author):
msg = " ".join(args[1:])
count = await self._broadcast_to_channels(msg)
await message.channel.send(f"Broadcast to {count} channels")
elif command == "chain":
if args and args[0] == "block":
height = await self.chain.get_block_height()
await message.channel.send(f"MSG Chain block height: `{height}`")
async def _handle_slash_query(self, interaction: discord.Interaction, query: str):
"""Handle /query slash command."""
await interaction.response.defer(thinking=True)
try:
result = await self.chain.query(query)
embed = discord.Embed(
title="Query Result",
description=str(result)[:4096],
color=discord.Color.blue(),
)
embed.set_footer(text=f"MSG Chain Agent | {self.agent_id}")
await interaction.followup.send(embed=embed)
except Exception as e:
embed = discord.Embed(
title="Query Failed",
description=str(e),
color=discord.Color.red(),
)
await interaction.followup.send(embed=embed)
async def _handle_slash_balance(
self, interaction: discord.Interaction, address: Optional[str]
):
"""Handle /balance slash command."""
await interaction.response.defer(thinking=True)
target = address or f"msg1{self.agent_id[:38]}"
try:
balance = await self.chain.query_balance(target)
embed = discord.Embed(title="Balance", color=discord.Color.green())
embed.add_field(name="Address", value=f"`{target}`", inline=False)
embed.add_field(name="Balance", value=f"**{balance} MSG**", inline=True)
await interaction.followup.send(embed=embed)
except Exception as e:
await interaction.followup.send(f"Query failed: {e}")
async def _handle_slash_monitor(
self,
interaction: discord.Interaction,
action: str,
channel: Optional[discord.TextChannel],
):
"""Handle /monitor slash command."""
target = channel or interaction.channel
if action == "start":
self.monitored_channels.add(target.id)
await interaction.response.send_message(f"Monitoring #{target.name}", ephemeral=True)
elif action == "stop":
self.monitored_channels.discard(target.id)
await interaction.response.send_message(f"Stopped monitoring #{target.name}", ephemeral=True)
else:
await interaction.response.send_message("Use 'start' or 'stop'", ephemeral=True)
async def _handle_slash_broadcast(self, interaction: discord.Interaction, message: str):
"""Handle /broadcast slash command."""
await interaction.response.defer(ephemeral=True)
count = 0
for channel_id in self.monitored_channels:
channel = self.get_channel(channel_id)
if channel and isinstance(channel, discord.TextChannel):
try:
embed = discord.Embed(
title="Agent Broadcast",
description=message,
color=discord.Color.gold(),
timestamp=datetime.utcnow(),
)
embed.set_footer(text=f"MSG Chain Agent | {self.agent_id}")
await channel.send(embed=embed)
count += 1
except Exception:
pass
await interaction.followup.send(f"Sent to {count}/{len(self.monitored_channels)} channels", ephemeral=True)
async def _handle_slash_status(self, interaction: discord.Interaction):
"""Handle /status slash command."""
embed = discord.Embed(
title="Agent Status",
color=discord.Color.blue(),
timestamp=datetime.utcnow(),
)
embed.add_field(name="Agent ID", value=f"`{self.agent_id}`", inline=False)
embed.add_field(name="Guilds", value=str(len(self.guilds)), inline=True)
embed.add_field(name="Monitored Channels", value=str(len(self.monitored_channels)), inline=True)
embed.add_field(name="Chain", value="MSG Chain", inline=True)
await interaction.response.send_message(embed=embed)
async def _broadcast_to_channels(self, message: str) -> int:
"""Broadcast text to all monitored channels."""
count = 0
for channel_id in self.monitored_channels:
channel = self.get_channel(channel_id)
if channel and isinstance(channel, discord.TextChannel):
try:
await channel.send(f"{message}")
count += 1
except Exception:
pass
return count
def _should_auto_respond(self, message: discord.Message) -> bool:
"""Decide if a message should get auto-response."""
keywords = ["agent", "help", "msgchain", "msg chain", "ai", "机器人"]
content_lower = message.content.lower()
return any(kw in content_lower for kw in keywords)
async def _auto_respond(self, message: discord.Message) -> str:
"""Generate auto-response."""
balance = await self.chain.query_balance(self.agent_id)
return (
f"**MSG Chain Agent** (auto-response)\n"
f"Detected keywords. Agent `{self.agent_id}` present.\n"
f"Balance: {balance} MSG\n"
f"@me or use `/` commands to interact."
)
async def _is_rate_limited(self, user_id: int) -> bool:
"""Simple per-user rate limiting (placeholder)."""
return False
async def _log_message(self, message: discord.Message):
"""Log message metadata on chain."""
await self.chain.submit_log({
"agent_id": self.agent_id,
"action": "discord_message",
"data": {
"channel_id": str(message.channel.id),
"author_id": str(message.author.id),
"message_id": str(message.id),
},
"timestamp": datetime.utcnow().isoformat(),
})
async def _log_interaction(self, action_type: str, user_id: int, query: str, response: str):
"""Log an agent-user interaction on chain."""
await self.chain.submit_log({
"agent_id": self.agent_id,
"action": f"discord_{action_type}",
"data": {
"user_id": str(user_id),
"query_preview": query[:100],
"response_preview": response[:100],
},
"timestamp": datetime.utcnow().isoformat(),
})
async def on_member_join(self, member: discord.Member):
"""Welcome new members."""
channel = member.guild.system_channel
if channel:
await channel.send(
f"Welcome {member.mention}! I'm MSG Chain AI Agent `{self.agent_id}`."
)
async def on_error(self, event_method: str, *args, **kwargs):
"""Global error handler."""
logger.error(f"Error in {event_method}: {args}")
@tasks.loop(minutes=15)
async def background_task(self):
"""Periodic background tasks."""
logger.debug("Running background tasks...")
try:
height = await self.chain.get_block_height()
logger.info(f"Chain block height: {height}")
except Exception as e:
logger.error(f"Chain health check failed: {e}")
class DiscordModeratorAgent(DiscordAgent):
"""Extended Discord agent with moderation capabilities."""
def __init__(self, agent_id: str, mod_role: str = "moderator", **kwargs):
super().__init__(agent_id, **kwargs)
self.mod_role = mod_role
self.filtered_words: List[str] = ["spam", "scam", "phishing"]
async def on_message(self, message: discord.Message):
if await self._check_filtered_content(message):
await self._take_mod_action(message)
return
await super().on_message(message)
async def _check_filtered_content(self, message: discord.Message) -> bool:
content_lower = message.content.lower()
return any(word in content_lower for word in self.filtered_words)
async def _take_mod_action(self, message: discord.Message):
await message.delete()
await message.channel.send(
f"{message.author.mention} your message was auto-deleted.",
delete_after=10,
)
logger.warning(f"Mod action: deleted message {message.id} from {message.author}")
def run_agent():
"""Run the Discord agent."""
token = os.getenv("DISCORD_BOT_TOKEN")
if not token:
raise EnvironmentError("DISCORD_BOT_TOKEN not set")
agent_id = os.getenv("AGENT_ID", "discord-agent-001")
agent = DiscordAgent(agent_id=agent_id)
agent.run(token=token, log_handler=None)
if __name__ == "__main__":
run_agent()
3.4 Discord 权限表
| 权限 | 值 | 说明 |
|---|---|---|
| Send Messages | 0x800 | 发送消息 |
| Read Message History | 0x10000 | 读取历史 |
| Mention Everyone | 0x20000 | @everyone |
| Use Slash Commands | 0x8000000000 | 斜杠命令 |
| Manage Channels | 0x10 | 管理频道(监控用) |
4. Telegram 集成
4.1 准备工作
- 在 @BotFather 创建 Bot,获取 Token
- 设置 Bot 命令: /start /balance /help /query /monitor
- 可选: 配置 Telegram Stars 支付
export TELEGRAM_BOT_TOKEN="your_bot_token"
export TELEGRAM_WEBHOOK_URL="https://your-domain.com/webhook/telegram"
4.2 依赖安装
pip install python-telegram-bot==20.6 httpx
4.3 完整 Telegram Agent 实现
import os
import asyncio
import logging
from typing import Optional, Dict, List, Any
from datetime import datetime
from telegram import (
Bot, Update, InlineKeyboardButton, InlineKeyboardMarkup,
BotCommand, BotCommandScopeDefault,
)
from telegram.ext import (
Application, ApplicationBuilder, CommandHandler, MessageHandler,
CallbackQueryHandler, filters, ContextTypes,
)
from telegram.constants import ParseMode
from msg_chain_sdk import MsgChainClient
load_dotenv()
logger = logging.getLogger(__name__)
class TelegramAgent:
"""Telegram bot integration for AI Agent on MSG Chain."""
def __init__(
self,
agent_id: str,
token: Optional[str] = None,
webhook_url: Optional[str] = None,
allowed_user_ids: Optional[List[int]] = None,
):
self.agent_id = agent_id
self.token = token or os.getenv("TELEGRAM_BOT_TOKEN")
self.webhook_url = webhook_url or os.getenv("TELEGRAM_WEBHOOK_URL")
self.allowed_user_ids = allowed_user_ids or []
if not self.token:
raise EnvironmentError("TELEGRAM_BOT_TOKEN not set")
self.chain = MsgChainClient()
self.capabilities = ["Balance Query", "On-chain Query", "Message Broadcast", "Auto Reply"]
self.app = ApplicationBuilder().token(self.token).build()
self._setup_handlers()
self.sessions: Dict[int, Dict] = {}
def _setup_handlers(self):
"""Register all command and message handlers."""
self.app.add_handler(CommandHandler("start", self._cmd_start))
self.app.add_handler(CommandHandler("help", self._cmd_help))
self.app.add_handler(CommandHandler("balance", self._cmd_balance))
self.app.add_handler(CommandHandler("query", self._cmd_query))
self.app.add_handler(CommandHandler("monitor", self._cmd_monitor))
self.app.add_handler(CommandHandler("broadcast", self._cmd_broadcast))
self.app.add_handler(CommandHandler("chain", self._cmd_chain))
self.app.add_handler(CommandHandler("cancel", self._cmd_cancel))
self.app.add_handler(CallbackQueryHandler(self._handle_callback))
self.app.add_handler(MessageHandler(filters.TEXT & ~filters.COMMAND, self._handle_message))
self.app.add_handler(MessageHandler(filters.StatusUpdate.NEW_CHAT_MEMBERS, self._handle_new_members))
self.app.add_error_handler(self._error_handler)
async def _cmd_start(self, update: Update, context: ContextTypes.DEFAULT_TYPE):
"""Handle /start command."""
user = update.effective_user
chat = update.effective_chat
self.sessions[user.id] = {
"user_id": user.id,
"username": user.username,
"first_seen": datetime.utcnow().isoformat(),
}
text = (
f"**MSG Chain AI Agent**\n\n"
f"Agent ID: `{self.agent_id}`\n"
f"Chain: `MSG Chain (msg-chain-1)`\n"
f"Prefix: `msg`\n\n"
f"**Capabilities:**\n"
+ "\n".join(f"- {cap}" for cap in self.capabilities) +
"\n\nType /help to see all commands."
)
keyboard = [
[InlineKeyboardButton("Balance", callback_data="balance"),
InlineKeyboardButton("Query", callback_data="query_help")],
[InlineKeyboardButton("Chain Status", callback_data="chain_status"),
InlineKeyboardButton("Help", callback_data="help")],
]
reply_markup = InlineKeyboardMarkup(keyboard)
await update.message.reply_text(text, reply_markup=reply_markup, parse_mode=ParseMode.MARKDOWN)
await self._log_action("start", {"user_id": user.id, "username": user.username})
async def _cmd_help(self, update: Update, context: ContextTypes.DEFAULT_TYPE):
"""Handle /help command."""
help_text = (
"**Agent Command Help**\n\n"
"/start - Show welcome\n"
"/help - This help\n"
"/balance [address] - Check balance\n"
"/query <keyword> - On-chain query\n"
"/monitor <start/stop> - Topic monitor\n"
"/broadcast <message> - Broadcast (admin)\n"
"/chain - Chain status\n"
"/cancel - Cancel operation"
)
await update.message.reply_text(help_text, parse_mode=ParseMode.MARKDOWN)
async def _cmd_balance(self, update: Update, context: ContextTypes.DEFAULT_TYPE):
"""Handle /balance command."""
await update.message.chat.send_action("typing")
address = context.args[0] if context.args else None
target = address or f"msg1{self.agent_id[:38]}"
try:
balance = await self.chain.query_balance(target)
text = f"**Balance**\n\nAddress: `{target}`\nBalance: **{balance} MSG**"
await update.message.reply_text(text, parse_mode=ParseMode.MARKDOWN)
except Exception as e:
await update.message.reply_text(f"Query failed: {str(e)}")
async def _cmd_query(self, update: Update, context: ContextTypes.DEFAULT_TYPE):
"""Handle /query command."""
if not context.args:
await update.message.reply_text("Usage: /query <keyword>")
return
query = " ".join(context.args)
await update.message.chat.send_action("typing")
try:
result = await self.chain.query(query)
text = f"**Query Result**\n\nQuery: `{query}`\n```\n{str(result)[:3000]}\n```"
await update.message.reply_text(text, parse_mode=ParseMode.MARKDOWN)
except Exception as e:
await update.message.reply_text(f"Query failed: {str(e)}")
async def _cmd_monitor(self, update: Update, context: ContextTypes.DEFAULT_TYPE):
"""Handle /monitor command."""
if not context.args or context.args[0] not in ("start", "stop"):
await update.message.reply_text("Usage: /monitor start <keyword> or /monitor stop <keyword>")
return
action = context.args[0]
keyword = " ".join(context.args[1:]) if len(context.args) > 1 else ""
if not keyword:
await update.message.reply_text("Please provide a keyword.")
return
msg = f"Started monitoring: `{keyword}`" if action == "start" else f"Stopped monitoring: `{keyword}`"
await update.message.reply_text(msg, parse_mode=ParseMode.MARKDOWN)
await self._log_action("monitor", {"action": action, "keyword": keyword})
async def _cmd_broadcast(self, update: Update, context: ContextTypes.DEFAULT_TYPE):
"""Handle /broadcast command (admin only)."""
user_id = update.effective_user.id
if self.allowed_user_ids and user_id not in self.allowed_user_ids:
await update.message.reply_text("Permission denied.")
return
if not context.args:
await update.message.reply_text("Usage: /broadcast <message>")
return
message = " ".join(context.args)
keyboard = [
[InlineKeyboardButton("Confirm", callback_data="broadcast_confirm"),
InlineKeyboardButton("Cancel", callback_data="broadcast_cancel")],
]
reply_markup = InlineKeyboardMarkup(keyboard)
await update.message.reply_text(f"**Broadcast Preview**\n\n{message}", reply_markup=reply_markup, parse_mode=ParseMode.MARKDOWN)
context.user_data["pending_broadcast"] = message
async def _cmd_chain(self, update: Update, context: ContextTypes.DEFAULT_TYPE):
"""Handle /chain command."""
await update.message.chat.send_action("typing")
try:
height = await self.chain.get_block_height()
text = f"**MSG Chain Status**\n\nChain ID: `msg-chain-1`\nPrefix: `msg`\nBlock Height: `{height}`"
await update.message.reply_text(text, parse_mode=ParseMode.MARKDOWN)
except Exception as e:
await update.message.reply_text(f"Failed: {str(e)}")
async def _cmd_cancel(self, update: Update, context: ContextTypes.DEFAULT_TYPE):
"""Handle /cancel command."""
context.user_data.clear()
await update.message.reply_text("Operation cancelled.")
async def _handle_message(self, update: Update, context: ContextTypes.DEFAULT_TYPE):
"""Handle regular text messages."""
user = update.effective_user
message = update.message
chat = update.effective_chat
if user.is_bot:
return
await message.chat.send_action("typing")
bot_username = (await context.bot.get_me()).username
is_mentioned = chat.type in ("group", "supergroup") and f"@{bot_username}" in message.text
if chat.type != "private" and not is_mentioned:
return
text = message.text
if is_mentioned:
text = text.replace(f"@{bot_username}", "").strip()
if not text:
text = "Hello"
response = await self._generate_response(text, user)
await self._log_action("message", {
"user_id": user.id,
"username": user.username,
"message_preview": text[:100],
})
await message.reply_text(response)
async def _generate_response(self, text: str, user) -> str:
"""Generate AI response."""
text_lower = text.lower()
if "hello" in text_lower or "hi" in text_lower or "你好" in text:
return f"Hello {user.first_name}! I'm MSG Chain Agent `{self.agent_id}`."
if "balance" in text_lower or "余额" in text:
balance = await self.chain.query_balance(self.agent_id)
return f"Agent `{self.agent_id}` balance: **{balance} MSG**"
if "who" in text_lower or "what" in text_lower or "你是谁" in text:
return (
f"I'm **{self.agent_id}**, an AI Agent on MSG Chain.\n"
f"I can query on-chain data and interact with users via Telegram.\n"
f"Prefix: `msg` | Chain ID: `msg-chain-1`"
)
return (
f"Message received!\n"
f"You said: _{text[:200]}_\n"
f"I am processing your request..."
)
async def _handle_callback(self, update: Update, context: ContextTypes.DEFAULT_TYPE):
"""Handle inline keyboard button presses."""
query = update.callback_query
await query.answer()
data = query.data
if data == "balance":
balance = await self.chain.query_balance(self.agent_id)
await query.edit_message_text(
f"**Balance**\n\nAgent `{self.agent_id}`\nBalance: **{balance} MSG**",
parse_mode=ParseMode.MARKDOWN,
)
elif data == "query_help":
await query.edit_message_text(
"**Query Help**\n\nUse /query <keyword>.\n"
"Example[未公开路径] /query MSG Chain stats\n- /query validator list",
parse_mode=ParseMode.MARKDOWN,
)
elif data == "chain_status":
height = await self.chain.get_block_height()
await query.edit_message_text(
f"**Chain Status**\n\nChain: MSG Chain\nHeight: `{height}`",
parse_mode=ParseMode.MARKDOWN,
)
elif data == "help":
await query.edit_message_text(
"**Commands**\n/start - Welcome\n/help - Help\n"
"/balance - Balance\n/query - Query\n/chain - Chain status",
parse_mode=ParseMode.MARKDOWN,
)
elif data == "broadcast_confirm":
message = context.user_data.get("pending_broadcast", "")
await query.edit_message_text(f"**Broadcast Sent!**\n\n{message}" if message else "No pending message.")
elif data == "broadcast_cancel":
context.user_data.pop("pending_broadcast", None)
await query.edit_message_text("Broadcast cancelled.")
async def _handle_new_members(self, update: Update, context: ContextTypes.DEFAULT_TYPE):
"""Welcome new members."""
for member in update.message.new_chat_members:
if member.id == context.bot.id:
await update.message.reply_text(
f"Hello everyone! I'm MSG Chain AI Agent `{self.agent_id}`."
)
async def _error_handler(self, update: Update, context: ContextTypes.DEFAULT_TYPE):
"""Handle errors."""
logger.error(f"Update {update} caused error {context.error}")
try:
if update and update.effective_chat:
await update.effective_chat.send_message("An error occurred.")
except Exception:
pass
async def _log_action(self, action: str, data: dict):
"""Log action on chain."""
try:
await self.chain.submit_log({
"agent_id": self.agent_id,
"action": f"telegram_{action}",
"data": data,
"timestamp": datetime.utcnow().isoformat(),
})
except Exception as e:
logger.warning(f"Failed to log: {e}")
async def set_commands(self):
"""Register bot commands with Telegram."""
commands = [
BotCommand("start", "Start Agent"),
BotCommand("help", "Show help"),
BotCommand("balance", "Check balance [address]"),
BotCommand("query", "On-chain query"),
BotCommand("monitor", "Monitor start/stop"),
BotCommand("chain", "Chain status"),
BotCommand("cancel", "Cancel operation"),
]
await self.app.bot.set_my_commands(commands, scope=BotCommandScopeDefault())
async def start_webhook(self):
"""Start bot using webhook mode."""
if not self.webhook_url:
raise ValueError("Webhook URL required")
await self.app.bot.set_webhook(url=self.webhook_url)
await self.app.start()
async def start_polling(self):
"""Start bot using polling mode."""
await self.set_commands()
logger.info("Starting Telegram bot in polling mode...")
await self.app.run_polling(allowed_updates=Update.ALL_TYPES)
def run(self):
"""Convenience method to start the bot."""
asyncio.run(self.start_polling())
class TelegramPaymentAgent(TelegramAgent):
"""Extended Telegram agent with Telegram Stars payment support."""
def __init__(self, agent_id: str, **kwargs):
super().__init__(agent_id, **kwargs)
self.provider_token = os.getenv("TELEGRAM_PAYMENT_TOKEN")
async def _cmd_balance(self, update: Update, context: ContextTypes.DEFAULT_TYPE):
await super()._cmd_balance(update, context)
keyboard = [[InlineKeyboardButton("Buy MSG with Stars", callback_data="buy_msg")]]
reply_markup = InlineKeyboardMarkup(keyboard)
await update.message.reply_text("Premium: Click to buy MSG with Telegram Stars.", reply_markup=reply_markup)
async def _handle_callback(self, update, context):
query = update.callback_query
await query.answer()
if query.data == "buy_msg":
await context.bot.send_invoice(
chat_id=query.message.chat_id,
title="MSG Token Purchase",
description="Buy 10 MSG Tokens",
payload="buy_10_msg",
provider_token=self.provider_token,
currency="XTR",
prices=[{"label": "10 MSG", "amount": 100}],
)
await super()._handle_callback(update, context)
def main():
agent_id = os.getenv("AGENT_ID", "telegram-agent-001")
agent = TelegramAgent(agent_id=agent_id)
mode = os.getenv("TELEGRAM_MODE", "polling")
if mode == "webhook":
asyncio.run(agent.start_webhook())
else:
agent.run()
if __name__ == "__main__":
main()
4.4 Telegram Bot 命令配置
| 命令 | 描述 | 权限 |
|---|---|---|
| /start | 启动并显示欢迎信息 | 公开 |
| /help | 显示帮助 | 公开 |
| /balance [address] | 查询余额 | 公开 |
| /query <keyword> | 链上查询 | 公开 |
| /monitor <start/stop> | 话题监控 | 公开 |
| /broadcast <message> | 广播消息 | 管理员 |
| /chain | 链状态 | 公开 |
| /cancel | 取消操作 | 公开 |
5. GitHub 集成
5.1 准备工作
- 在 GitHub Settings Tokens 生成 Personal Access Token
- 权限: repo, issues, pull_requests, contents
- 可选: 安装 GitHub App 以获得 Webhook 支持
export GITHUB_TOKEN="your_github_token"
export GITHUB_WEBHOOK_SECRET="your_webhook_secret"
5.2 依赖安装
pip install PyGithub==2.1.1 httpx
5.3 完整 GitHub Agent 实现
import os
import re
import asyncio
import logging
from typing import Optional, List, Dict, Any
from datetime import datetime
from github import Github, GithubIntegration
from github.Repository import Repository
from msg_chain_sdk import MsgChainClient
load_dotenv()
logger = logging.getLogger(__name__)
class GitHubAgent:
"""GitHub integration for AI Agent on MSG Chain."""
def __init__(
self,
agent_id: str,
token: Optional[str] = None,
app_id: Optional[str] = None,
app_key: Optional[str] = None,
):
self.agent_id = agent_id
self.token = token or os.getenv("GITHUB_TOKEN")
if not self.token:
raise EnvironmentError("GITHUB_TOKEN not set")
self.gh = Github(
login_or_token=self.token,
per_page=100,
user_agent=f"MSGChain-Agent/{agent_id}",
)
self.app_id = app_id or os.getenv("GITHUB_APP_ID")
self.app_key = app_key or os.getenv("GITHUB_APP_KEY")
self.integration: Optional[GithubIntegration] = None
if self.app_id and self.app_key:
self.integration = GithubIntegration(self.app_id, self.app_key)
self.chain = MsgChainClient()
self._repo_cache: Dict[str, Repository] = {}
self._reviewed_prs: set = set()
def _get_repo(self, repo_full_name: str) -> Repository:
if repo_full_name not in self._repo_cache:
self._repo_cache[repo_full_name] = self.gh.get_repo(repo_full_name)
return self._repo_cache[repo_full_name]
async def _log_action(self, action: str, data: dict):
await self.chain.submit_log({
"agent_id": self.agent_id,
"action": f"github_{action}",
"data": data,
"timestamp": datetime.utcnow().isoformat(),
})
async def create_issue(
self, repo: str, title: str, body: str,
labels: Optional[List[str]] = None,
assignees: Optional[List[str]] = None,
) -> int:
"""Create an issue in the specified repository."""
repository = self._get_repo(repo)
full_body = (
f"{body}\n\n---\n"
f"_Created by AI Agent `{self.agent_id}` | MSG Chain_"
)
issue = repository.create_issue(
title=title, body=full_body,
labels=labels or [], assignees=assignees or [],
)
logger.info(f"Issue #{issue.number} created in {repo}")
await self._log_action("create_issue", {
"repo": repo, "issue_number": issue.number, "title": title,
})
return issue.number
async def close_issue(self, repo: str, issue_number: int, comment: Optional[str] = None) -> bool:
"""Close an issue with an optional comment."""
repository = self._get_repo(repo)
issue = repository.get_issue(number=issue_number)
if comment:
issue.create_comment(f"{comment}\n\n_Closed by AI Agent `{self.agent_id}`_")
issue.edit(state="closed")
logger.info(f"Issue #{issue_number} closed in {repo}")
return True
async def comment_on_issue(self, repo: str, issue_number: int, comment: str) -> int:
"""Comment on an issue or PR."""
repository = self._get_repo(repo)
issue = repository.get_issue(number=issue_number)
full = f"{comment}\n\n_AI Agent `{self.agent_id}` | MSG Chain_"
created = issue.create_comment(full)
return created.id
async def list_issues(
self, repo: str, state: str = "open",
labels: Optional[List[str]] = None,
max_results: int = 50,
) -> List[Dict]:
"""List issues in a repository."""
repository = self._get_repo(repo)
issues = repository.get_issues(state=state, labels=labels)
result = []
for issue in issues[:max_results]:
result.append({
"number": issue.number, "title": issue.title,
"state": issue.state,
"labels": [l.name for l in issue.labels],
"created_at": issue.created_at.isoformat(),
"url": issue.html_url,
})
return result
async def review_pull_request(self, repo: str, pr_number: int, auto_approve: bool = False) -> Dict:
"""Perform AI-powered review on a pull request."""
if (repo, pr_number) in self._reviewed_prs:
return {"error": "PR already reviewed by this agent"}
repository = self._get_repo(repo)
pr = repository.get_pull(pr_number)
files = pr.get_files()
review_comments = []
has_issues = False
for file in files:
if not file.patch:
continue
analysis = await self._analyze_file(file.filename, file.patch, file.status)
if analysis["issues"]:
has_issues = True
for issue in analysis["issues"]:
review_comments.append(f"**{file.filename}**: {issue['message']}")
if analysis["suggestions"]:
for s in analysis["suggestions"]:
review_comments.append(f"**{file.filename}**: {s}")
body_parts = [
f"## AI Agent Code Review\n",
f"**PR:** #{pr_number} - {pr.title}",
f"**Reviewer Agent:** `{self.agent_id}`",
f"**Time:** {datetime.utcnow().isoformat()}",
"",
]
if review_comments:
body_parts.append("### Comments")
body_parts.extend(f"- {c}" for c in review_comments)
else:
body_parts.append("**No issues found.**")
body_parts.extend([
"", f"**Files:** {files.totalCount} | **Comments:** {len(review_comments)}",
"", "---", f"_Reviewed by AI Agent `{self.agent_id}` | MSG Chain_",
])
event = "REQUEST_CHANGES" if has_issues else ("APPROVE" if auto_approve else "COMMENT")
pr.create_review(body="\n".join(body_parts), event=event)
self._reviewed_prs.add((repo, pr_number))
logger.info(f"Reviewed PR #{pr_number} in {repo} ({len(review_comments)} comments, event={event})")
await self._log_action("review_pr", {
"repo": repo, "pr_number": pr_number,
"comment_count": len(review_comments), "event": event,
})
return {"pr_number": pr_number, "comment_count": len(review_comments), "event": event, "has_issues": has_issues}
async def _analyze_file(self, filename: str, patch: str, status: str) -> Dict:
"""Analyze a file's diff for issues and suggestions."""
issues = []
suggestions = []
ext = filename.rsplit(".", 1)[-1].lower() if "." in filename else ""
lines = patch.split("\n")
if len(lines) > 500:
suggestions.append("Large change (>500 lines), consider splitting into smaller PRs.")
if ext == "py":
if "TODO" in patch or "FIXME" in patch:
issues.append({"message": "Contains TODO/FIXME, address before merge.", "severity": "warning"})
if "print(" in patch and "import logging" not in patch:
suggestions.append("Consider using logging instead of print().")
elif ext in ("js", "ts"):
if "console.log" in patch:
suggestions.append("Remove debug console.log statements.")
if "any" in patch and ext == "ts":
suggestions.append("Avoid `any` type, use specific types.")
elif ext == "rs":
if "unwrap()" in patch:
issues.append({"message": "Avoid unwrap() in production, use ? operator.", "severity": "error"})
security_patterns = [
(r"(api_key|secret|password|token)\s*=\s*['\"].+['\"]", "Possible hardcoded credential"),
(r"eval\s*\(", "eval() is a security risk"),
(r"exec\s*\(", "exec() is a security risk"),
]
for pattern, msg in security_patterns:
if re.search(pattern, patch, re.IGNORECASE):
issues.append({"message": f"Security: {msg}", "severity": "error"})
return {"issues": issues, "suggestions": suggestions}
async def merge_pull_request(self, repo: str, pr_number: int, merge_method: str = "merge") -> bool:
"""Merge a pull request."""
repository = self._get_repo(repo)
pr = repository.get_pull(pr_number)
if not pr.mergeable:
return False
result = pr.merge(merge_method=merge_method)
if result.merged:
await self._log_action("merge_pr", {"repo": repo, "pr_number": pr_number, "method": merge_method})
return True
return False
async def create_pull_request(
self, repo: str, title: str, body: str, head: str, base: str = "main", draft: bool = False
) -> int:
"""Create a pull request."""
repository = self._get_repo(repo)
full_body = f"{body}\n\n---\n_Created by AI Agent `{self.agent_id}` | MSG Chain_"
pr = repository.create_pull(title=title, body=full_body, head=head, base=base, draft=draft)
logger.info(f"PR #{pr.number} created: {pr.html_url}")
await self._log_action("create_pr", {"repo": repo, "pr_number": pr.number, "title": title})
return pr.number
async def get_repo_stats(self, repo: str) -> Dict:
"""Get repository statistics."""
repository = self._get_repo(repo)
return {
"name": repo, "stars": repository.stargazers_count,
"forks": repository.forks_count,
"open_issues": repository.open_issues_count,
"language": repository.language,
"topics": repository.get_topics(),
}
async def auto_triage_issues(self, repo: str, max_issues: int = 10) -> List[Dict]:
"""Automatically triage unlabeled open issues."""
repository = self._get_repo(repo)
issues = repository.get_issues(state="open", labels=[])
results = []
for issue in issues[:max_issues]:
suggested = self._suggest_labels(issue.title, issue.body or "")
if suggested:
for label_name in suggested:
try:
label = repository.get_label(label_name)
issue.add_to_labels(label)
except Exception:
pass
results.append({"issue_number": issue.number, "title": issue.title, "labels": suggested})
if results:
await self._log_action("triage_issues", {"repo": repo, "count": len(results)})
return results
def _suggest_labels(self, title: str, body: str) -> List[str]:
"""Suggest GitHub labels based on issue content."""
text = f"{title} {body}".lower()
labels = []
if any(w in text for w in ["bug", "fix", "error", "crash"]):
labels.append("bug")
if any(w in text for w in ["feature", "enhancement", "request"]):
labels.append("enhancement")
if any(w in text for w in ["doc", "documentation", "readme"]):
labels.append("documentation")
if any(w in text for w in ["security", "vulnerability"]):
labels.append("security")
if any(w in text for w in ["urgent", "critical", "blocker"]):
labels.append("priority:high")
return labels[:5]
if __name__ == "__main__":
async def main():
agent = GitHubAgent(agent_id="github-agent-001")
await agent.create_issue("msgchain/msg-chain-core", "Add AI Agent SDK docs",
"We need comprehensive docs.", labels=["documentation"])
stats = await agent.get_repo_stats("msgchain/msg-chain-core")
print(f"Stars: {stats['stars']}, Issues: {stats['open_issues']}")
asyncio.run(main())
5.4 GitHub Token 权限说明
| 权限 | 用途 | 必需 |
|---|---|---|
| repo | 访问私有仓库 | 是 |
| issues:write | 创建/管理 Issue | 是 |
| pull_requests:write | PR 审查/合并 | 是 |
| contents:read | 读取代码 | 是 |
| workflows | 触发 Actions | 可选 |
6. Slack 集成
6.1 准备工作
- 在 Slack API 创建 App
- 启用 Socket Mode 或配置 Request URL
- 添加 Bot Token Scopes: chat:write, channels:history, channels:read, users:read
- 安装 App 到工作区
export SLACK_BOT_TOKEN="xoxb-your-bot-token"
export SLACK_APP_TOKEN="xapp-your-app-token"
export SLACK_SIGNING_SECRET="your-signing-secret"
6.2 依赖安装
pip install slack-bolt==1.18.1 slack-sdk==3.27.1
6.3 Slack Agent 实现
import os
import asyncio
import logging
from typing import Optional, List, Dict
from datetime import datetime
from slack_bolt import App, Ack
from slack_bolt.adapter.socket_mode import SocketModeHandler
from slack_sdk import WebClient
from slack_sdk.errors import SlackApiError
from msg_chain_sdk import MsgChainClient
load_dotenv()
logger = logging.getLogger(__name__)
class SlackAgent:
"""Slack integration for AI Agent on MSG Chain."""
def __init__(
self,
agent_id: str,
bot_token: Optional[str] = None,
app_token: Optional[str] = None,
signing_secret: Optional[str] = None,
):
self.agent_id = agent_id
self.bot_token = bot_token or os.getenv("SLACK_BOT_TOKEN")
self.app_token = app_token or os.getenv("SLACK_APP_TOKEN")
self.signing_secret = signing_secret or os.getenv("SLACK_SIGNING_SECRET")
if not self.bot_token or not self.app_token:
raise EnvironmentError("SLACK_BOT_TOKEN and SLACK_APP_TOKEN must be set")
self.app = App(token=self.bot_token, signing_secret=self.signing_secret)
self.client = WebClient(token=self.bot_token)
self.chain = MsgChainClient()
self._register_handlers()
def _register_handlers(self):
"""Register all Slack event and command handlers."""
@self.app.event("app_mention")
async def handle_mention(event: dict, say: callable):
user = event.get("user", "")
text = event.get("text", "")
bot_user_id = self.app.client.auth_test()["user_id"]
clean_text = text.replace(f"<@{bot_user_id}>", "").strip()
response = await self._generate_response(clean_text, user)
await say(text=response, thread_ts=event.get("ts"))
await self._log_action("mention", {"user": user, "text_preview": clean_text[:100]})
@self.app.event("message")
async def handle_dm(event: dict, say: callable):
if event.get("channel_type") != "im":
return
user = event.get("user", "")
text = event.get("text", "")
response = await self._generate_response(text, user, "dm")
await say(text=response, thread_ts=event.get("ts"))
@self.app.command("/agent")
async def agent_command(ack: Ack, command: dict, respond: callable):
await ack()
cmd_text = command.get("text", "").strip()
result = await self._handle_slash_command(cmd_text)
await respond(text=result)
@self.app.command("/balance")
async def balance_command(ack: Ack, command: dict, respond: callable):
await ack()
address = command.get("text", "").strip() or None
target = address or f"msg1{self.agent_id[:38]}"
try:
balance = await self.chain.query_balance(target)
await respond(f"Balance for `{target}`: *{balance} MSG*")
except Exception as e:
await respond(f"Query failed: {e}")
@self.app.command("/chain")
async def chain_command(ack: Ack, command: dict, respond: callable):
await ack()
height = await self.chain.get_block_height()
await respond(f"MSG Chain block height: `{height}`")
@self.app.action("button_balance")
async def button_balance(ack: Ack, body: dict, client: WebClient):
await ack()
balance = await self.chain.query_balance(self.agent_id)
await client.chat_postMessage(
channel=body["channel"]["id"],
text=f"Agent `{self.agent_id}` balance: *{balance} MSG*",
thread_ts=body["message"]["ts"],
)
@self.app.action("button_chain_status")
async def button_chain_status(ack: Ack, body: dict, client: WebClient):
await ack()
height = await self.chain.get_block_height()
await client.chat_postMessage(
channel=body["channel"]["id"],
text=f"MSG Chain block height: `{height}`",
thread_ts=body["message"]["ts"],
)
async def _generate_response(self, text: str, user: str, channel: str = "") -> str:
text_lower = text.lower()
if not text:
return f"Hello <@{user}>! I'm MSG Chain AI Agent `{self.agent_id}`."
if "balance" in text_lower or "余额" in text:
balance = await self.chain.query_balance(self.agent_id)
return f"Agent `{self.agent_id}` balance: *{balance} MSG*"
if "help" in text_lower:
return (
f"*MSG Chain AI Agent Help*\n"
f"/agent <message> - Chat\n/balance [address] - Balance\n/chain - Status"
)
if "chain" in text_lower:
height = await self.chain.get_block_height()
return f"MSG Chain block height: `{height}`"
return f"Message received <@{user}>! Content: _{text[:200]}_"
async def _handle_slash_command(self, text: str) -> str:
if not text or text == "help":
return "Usage: `/agent <message>`"
text_lower = text.lower()
if text_lower.startswith("balance"):
return await self._handle_balance_command(text[8:].strip() or None)
if text_lower in ("chain", "status"):
height = await self.chain.get_block_height()
return f"MSG Chain block height: `{height}`"
return await self._generate_response(text, "user")
async def _handle_balance_command(self, address: Optional[str]) -> str:
target = address or f"msg1{self.agent_id[:38]}"
try:
balance = await self.chain.query_balance(target)
return f"Address `{target}` balance: *{balance} MSG*"
except Exception as e:
return f"Query failed: {str(e)}"
async def send_message(self, channel: str, text: str, blocks: Optional[List[Dict]] = None) -> bool:
try:
kwargs = {"channel": channel, "text": text}
if blocks:
kwargs["blocks"] = blocks
self.client.chat_postMessage(**kwargs)
return True
except SlackApiError as e:
logger.error(f"Failed to send: {e}")
return False
async def _log_action(self, action: str, data: dict):
try:
await self.chain.submit_log({
"agent_id": self.agent_id,
"action": f"slack_{action}",
"data": data,
"timestamp": datetime.utcnow().isoformat(),
})
except Exception as e:
logger.warning(f"Failed to log: {e}")
def start(self):
logger.info(f"Starting Slack Agent: {self.agent_id}")
handler = SocketModeHandler(self.app, self.app_token)
handler.start()
if __name__ == "__main__":
agent = SlackAgent(agent_id="slack-agent-001")
agent.start()
6.4 Slack Event 清单
| 事件 | 用途 | 启用方式 |
|---|---|---|
| app_mention | Bot 被 @提及 | Event Subscription |
| message.im | 私信 Bot | Event Subscription |
| /agent | Slash 命令 | Slash Commands |
| /balance | 余额查询 | Slash Commands |
| /chain | 链状态 | Slash Commands |
7. Email (SMTP/IMAP) 集成
7.1 准备工作
export SMTP_HOST="smtp.gmail.com"
export SMTP_PORT=587
export SMTP_USER="your-email@gmail.com"
export SMTP_PASSWORD="your-app-password"
export IMAP_HOST="imap.gmail.com"
export IMAP_PORT=993
7.2 依赖安装
pip install aiosmtplib==3.0.1 aioimaplib==1.0.1 email-validator==2.1.0
7.3 Email Agent 实现
import os
import re
import asyncio
import logging
import email
from email.mime.text import MIMEText
from email.mime.multipart import MIMEMultipart
from email.mime.base import MIMEBase
from email.utils import formataddr
from typing import Optional, List, Dict
from datetime import datetime
import aiosmtplib
import aioimaplib
from email_validator import validate_email, EmailNotValidError
from msg_chain_sdk import MsgChainClient
load_dotenv()
logger = logging.getLogger(__name__)
class EmailAgent:
"""Email (SMTP/IMAP) integration for AI Agent on MSG Chain."""
def __init__(
self, agent_id: str,
smtp_host: Optional[str] = None, smtp_port: Optional[int] = None,
smtp_user: Optional[str] = None, smtp_password: Optional[str] = None,
imap_host: Optional[str] = None, imap_port: Optional[int] = None,
):
self.agent_id = agent_id
self.email_address = smtp_user or os.getenv("SMTP_USER", "")
self.smtp_host = smtp_host or os.getenv("SMTP_HOST", "smtp.gmail.com")
self.smtp_port = smtp_port or int(os.getenv("SMTP_PORT", "587"))
self.smtp_user = smtp_user or os.getenv("SMTP_USER", "")
self.smtp_password = smtp_password or os.getenv("SMTP_PASSWORD", "")
self.imap_host = imap_host or os.getenv("IMAP_HOST", "imap.gmail.com")
self.imap_port = imap_port or int(os.getenv("IMAP_PORT", "993"))
if not self.smtp_user or not self.smtp_password:
raise EnvironmentError("SMTP_USER and SMTP_PASSWORD must be set")
self.chain = MsgChainClient()
self.templates: Dict[str, str] = {}
self._processed_uids: set = set()
async def send_email(
self, to: str, subject: str, body: str,
from_name: Optional[str] = None, html: Optional[str] = None,
cc: Optional[List[str]] = None, bcc: Optional[List[str]] = None,
) -> bool:
"""Send an email via SMTP."""
if not self._validate_address(to):
logger.error(f"Invalid address: {to}")
return False
msg = MIMEMultipart("alternative")
msg["From"] = formataddr((from_name or f"MSG Chain Agent {self.agent_id}", self.smtp_user))
msg["To"] = to
msg["Subject"] = subject
msg["Date"] = email.utils.formatdate(localtime=True)
msg["Message-ID"] = email.utils.make_msgid(domain="msgchain.org")
if cc:
for addr in cc:
if self._validate_address(addr):
msg["CC"] = addr
msg.attach(MIMEText(body, "plain", "utf-8"))
if html:
msg.attach(MIMEText(html, "html", "utf-8"))
try:
async with aiosmtplib.SMTP(
hostname=self.smtp_host, port=self.smtp_port,
use_tls=self.smtp_port == 465, start_tls=self.smtp_port == 587,
) as smtp:
await smtp.login(self.smtp_user, self.smtp_password)
recipients = [to] + (cc or []) + (bcc or [])
await smtp.send_message(msg, recipients=recipients)
logger.info(f"Email sent to {to}: {subject}")
await self._log_action("send_email", {
"to": to, "subject": subject[:100], "cc_count": len(cc or []),
})
return True
except (aiosmtplib.SMTPException, ConnectionError) as e:
logger.error(f"Failed to send: {e}")
return False
async def fetch_emails(self, folder: str = "INBOX", limit: int = 10) -> List[Dict]:
"""Fetch recent emails from IMAP inbox."""
emails = []
try:
imap = aioimaplib.IMAP4_SSL(host=self.imap_host, port=self.imap_port)
await imap.wait_hello()
await imap.login(self.smtp_user, self.smtp_password)
await imap.select(folder)
_, message_ids = await imap.search("ALL")
ids = message_ids[0].split() if message_ids[0] else []
for msg_id in ids[-limit:]:
msg_id_str = msg_id.decode()
if msg_id_str in self._processed_uids:
continue
_, data = await imap.fetch(msg_id, "(RFC822)")
if data and data[0]:
parsed = self._parse_email(data[0][1])
if parsed:
emails.append(parsed)
self._processed_uids.add(msg_id_str)
await imap.logout()
logger.info(f"Fetched {len(emails)} emails from {folder}")
except Exception as e:
logger.error(f"Failed to fetch emails: {e}")
return emails
def _parse_email(self, raw_data: bytes) -> Optional[Dict]:
try:
msg = email.message_from_bytes(raw_data)
subject_parts = email.header.decode_header(msg["Subject"] or "")
subject = "".join(
part.decode(charset or "utf-8") if isinstance(part, bytes) else part
for part, charset in subject_parts
)
body = ""
if msg.is_multipart():
for part in msg.walk():
if part.get_content_type() == "text/plain":
body = self._decode_part(part)
else:
body = self._decode_part(msg)
return {
"message_id": msg["Message-ID"] or "",
"subject": subject,
"from": msg["From"] or "",
"date": msg["Date"] or "",
"body": body[:5000],
}
except Exception as e:
logger.error(f"Failed to parse: {e}")
return None
def _decode_part(self, part) -> str:
try:
payload = part.get_payload(decode=True)
if payload:
charset = part.get_content_charset() or "utf-8"
return payload.decode(charset, errors="replace")
except Exception:
pass
return ""
async def auto_reply(self, folder: str = "INBOX", max_replies: int = 5) -> int:
"""Auto-reply to unprocessed emails."""
emails = await self.fetch_emails(folder=folder, limit=max_replies)
count = 0
for mail in emails:
if not mail["body"]:
continue
reply_body = await self._generate_auto_reply(mail["subject"], mail["body"], mail["from"])
if reply_body:
success = await self.send_email(
to=mail["from"], subject=f"Re: {mail['subject']}",
body=reply_body, from_name=f"MSG Chain Agent ({self.agent_id})",
)
if success:
count += 1
return count
async def _generate_auto_reply(self, subject: str, body: str, sender: str) -> Optional[str]:
text = f"{subject} {body}".lower()
if "balance" in text or "余额" in text:
balance = await self.chain.query_balance(self.agent_id)
return f"Agent ({self.agent_id}) balance: {balance} MSG\nChain: MSG Chain"
if "help" in text or "support" in text:
return (
f"I'm AI Agent (ID: {self.agent_id}) on MSG Chain.\n"
f"I can check balances, query chain data, and monitor events."
)
return f"Received your email.\nSubject: {subject}\nAgent ({self.agent_id}) is processing."
def _validate_address(self, address: str) -> bool:
try:
validate_email(address, check_deliverability=False)
return True
except EmailNotValidError:
return False
async def _log_action(self, action: str, data: dict):
try:
await self.chain.submit_log({
"agent_id": self.agent_id,
"action": f"email_{action}",
"data": data,
"timestamp": datetime.utcnow().isoformat(),
})
except Exception:
pass
if __name__ == "__main__":
async def main():
agent = EmailAgent(agent_id="email-agent-001")
await agent.send_email("user@example.com", "Hello from MSG Chain",
"This is an automated message from the AI Agent.")
asyncio.run(main())
8. Webhook 系统
8.1 概述
Webhook 系统作为 Agent 的"耳朵",接收来自第三方服务的事件推送。无论 Twitter 的新推文、GitHub 的 PR 创建,还是 Slack 的消息,都通过统一的 Webhook 入口路由到对应的 Agent 处理器。
8.2 依赖安装
pip install fastapi==0.109.0 uvicorn[standard]==0.27.0 pydantic==2.5.3
8.3 Webhook 系统实现
import os
import hmac
import hashlib
import logging
from typing import Optional, Dict, Any, Callable
from datetime import datetime
from fastapi import FastAPI, Request, HTTPException
from pydantic import BaseModel
import uvicorn
from msg_chain_sdk import MsgChainClient
load_dotenv()
logger = logging.getLogger(__name__)
class GitHubEvent(BaseModel):
action: str = ""
issue: Optional[Dict] = None
pull_request: Optional[Dict] = None
sender: Optional[Dict] = None
repository: Optional[Dict] = None
class DiscordEvent(BaseModel):
type: int
token: str
member: Optional[Dict] = None
data: Optional[Dict] = None
channel_id: Optional[str] = None
guild_id: Optional[str] = None
class GenericWebhookPayload(BaseModel):
source: str
event: str
data: Dict = {}
timestamp: str = ""
signature: Optional[str] = None
class WebhookRouter:
"""Routes incoming webhooks to registered handlers."""
def __init__(self, agent, secret_token: Optional[str] = None):
self.agent = agent
self.secret_token = secret_token or os.getenv("WEBHOOK_SECRET", "")
self.routes: Dict[str, Callable] = {}
self.chain = MsgChainClient()
def register(self, event_type: str, handler: Callable):
"""Register a handler for an event type."""
self.routes[event_type] = handler
logger.info(f"Registered handler for: {event_type}")
async def dispatch(self, event_type: str, payload: dict) -> dict:
"""Dispatch an event to its registered handler."""
handler = self.routes.get(event_type)
if not handler:
logger.warning(f"No handler for event type: {event_type}")
return {"status": "ignored", "reason": "no_handler"}
try:
result = await handler(payload)
await self.chain.submit_log({
"agent_id": getattr(self.agent, "agent_id", "unknown"),
"action": f"webhook_{event_type}",
"data": {"event_type": event_type, "status": "processed"},
"timestamp": datetime.utcnow().isoformat(),
})
return {"status": "ok", "result": result}
except Exception as e:
logger.error(f"Handler error for {event_type}: {e}")
return {"status": "error", "message": str(e)}
def verify_github_signature(self, payload_body: bytes, signature_header: str) -> bool:
"""Verify GitHub webhook signature."""
if not self.secret_token:
return True
expected = "sha256=" + hmac.new(
self.secret_token.encode(), payload_body, hashlib.sha256
).hexdigest()
return hmac.compare_digest(expected, signature_header)
class WebhookServer:
"""FastAPI webhook server for receiving third-party events."""
def __init__(self, agent, port: int = 8080, host: str = "0.0.0.0"):
self.agent = agent
self.port = port
self.host = host
self.app = FastAPI(title="MSG Chain Agent Webhooks")
self.router = WebhookRouter(agent)
self._setup_routes()
def _setup_routes(self):
"""Setup webhook endpoint routes."""
@self.app.post("/webhook/github")
async def github_webhook(request: Request):
body = await request.body()
sig = request.headers.get("X-Hub-Signature-256", "")
event_type = request.headers.get("X-GitHub-Event", "")
if not self.router.verify_github_signature(body, sig):
raise HTTPException(status_code=401, detail="Invalid signature")
payload = await request.json()
result = await self.router.dispatch(f"github:{event_type}", payload)
return result
@self.app.post("/webhook/discord")
async def discord_webhook(request: Request):
payload = await request.json()
event = DiscordEvent(**payload)
result = await self.router.dispatch("discord:interaction", event.model_dump())
return result
@self.app.post("/webhook/generic")
async def generic_webhook(payload: GenericWebhookPayload):
result = await self.router.dispatch(
f"generic:{payload.source}:{payload.event}",
payload.model_dump(),
)
return result
@self.app.get("/health")
async def health():
return {"status": "ok", "agent": getattr(self.agent, "agent_id", "unknown")}
def register_handler(self, event_type: str, handler: Callable):
"""Register an event handler."""
self.router.register(event_type, handler)
def start(self):
"""Start the webhook server."""
logger.info(f"Starting webhook server on {self.host}:{self.port}")
uvicorn.run(self.app, host=self.host, port=self.port)
class WebhookClient:
"""Client for sending webhooks to external services."""
def __init__(self):
import httpx
self.http = httpx.AsyncClient(timeout=30)
async def send(self, url: str, payload: dict, secret: Optional[str] = None) -> bool:
"""Send a webhook to an external URL."""
headers = {"Content-Type": "application/json"}
if secret:
body = json.dumps(payload).encode()
signature = hmac.new(secret.encode(), body, hashlib.sha256).hexdigest()
headers["X-Webhook-Signature"] = signature
try:
resp = await self.http.post(url, json=payload, headers=headers)
resp.raise_for_status()
logger.info(f"Webhook sent to {url}: {resp.status_code}")
return True
except Exception as e:
logger.error(f"Webhook to {url} failed: {e}")
return False
async def close(self):
await self.http.aclose()
# Example: Setup webhook server with handlers
if __name__ == "__main__":
import json
class SimpleAgent:
def __init__(self):
self.agent_id = "webhook-agent-001"
async def handle_github_push(self, payload: dict):
repo = payload.get("repository", {}).get("full_name", "unknown")
ref = payload.get("ref", "")
logger.info(f"Push event from {repo}: {ref}")
return {"repo": repo, "ref": ref}
async def handle_github_issues(self, payload: dict):
action = payload.get("action", "")
issue = payload.get("issue", {})
logger.info(f"Issue {action}: #{issue.get('number')}")
return {"action": action, "issue_number": issue.get("number")}
agent = SimpleAgent()
server = WebhookServer(agent, port=8080)
server.register_handler("github:push", agent.handle_github_push)
server.register_handler("github:issues", agent.handle_github_issues)
server.start()
8.4 Webhook 安全
| 服务 | 验证方式 | 实现 |
|---|---|---|
| GitHub | HMAC-SHA256 签名 | verify_github_signature() |
| Discord | Interaction Token + Public Key | 需额外验证 |
| 自定义 | HMAC 签名 + Timestamp | GenericWebhookPayload |
9. 跨平台消息路由
9.1 概述
跨平台消息路由允许 Agent 在多个服务之间同步消息、广播通知,并根据平台特性进行格式化。这是构建多平台 Agent 的核心基础设施。
9.2 实现
import asyncio
import logging
from typing import Optional, List, Dict, Any, Callable
from datetime import datetime
from dataclasses import dataclass, field
from msg_chain_sdk import MsgChainClient
logger = logging.getLogger(__name__)
@dataclass
class PlatformMessage:
"""Unified message format across platforms."""
content: str
platform: str
thread_id: Optional[str] = None
author: Optional[str] = None
timestamp: str = ""
metadata: Dict = field(default_factory=dict)
class PlatformFormatter:
"""Format content for different platforms."""
@staticmethod
def twitter(content: dict) -> str:
text = content.get("text", "")
if len(text) > 280:
text = text[:277] + "..."
if content.get("hashtags"):
text += "\n" + " ".join(f"#{h}" for h in content["hashtags"][:4])
return text
@staticmethod
def discord(content: dict) -> str:
text = content.get("text", "")
embed = content.get("embed")
if embed:
return f"{text}\n\n**{embed.get('title', '')}**\n{embed.get('description', '')}"
return text
@staticmethod
def telegram(content: dict) -> str:
text = content.get("text", "")
markup = content.get("markdown")
if markup:
return f"*{content.get('title', '')}*\n\n{text}"
return text
@staticmethod
def slack(content: dict) -> str:
text = content.get("text", "")
blocks = content.get("blocks")
if blocks:
import json
return json.dumps({"text": text, "blocks": blocks})
return text
@staticmethod
def email(content: dict) -> str:
text = content.get("text", "")
html = content.get("html", "")
return html if html else text
class UnifiedMessenger:
"""Route messages across all integrated platforms."""
def __init__(self, agent, chain_client: Optional[MsgChainClient] = None):
self.agent = agent
self.chain = chain_client or MsgChainClient()
self.formatters = {
"twitter": PlatformFormatter.twitter,
"discord": PlatformFormatter.discord,
"telegram": PlatformFormatter.telegram,
"slack": PlatformFormatter.slack,
"email": PlatformFormatter.email,
}
self.platforms: Dict[str, Any] = {}
self.message_queue: asyncio.Queue = asyncio.Queue()
self._running = False
def register_platform(self, name: str, instance: Any):
"""Register a platform client."""
self.platforms[name] = instance
logger.info(f"Platform registered: {name}")
def register_formatter(self, platform: str, formatter: Callable):
"""Register a custom formatter for a platform."""
self.formatters[platform] = formatter
async def send(self, platform: str, content: str, **kwargs) -> bool:
"""Send a message to a specific platform."""
client = self.platforms.get(platform)
if not client:
logger.error(f"Platform not registered: {platform}")
return False
try:
if platform == "twitter":
await client.post_tweet(content)
elif platform == "discord":
channel = kwargs.get("channel")
if channel:
await channel.send(content)
elif platform == "telegram":
chat_id = kwargs.get("chat_id")
if chat_id:
await client.send_message(chat_id=chat_id, text=content)
elif platform == "slack":
channel = kwargs.get("channel")
if channel:
await client.send_message(channel, content)
elif platform == "email":
to = kwargs.get("to")
subject = kwargs.get("subject", "Message from MSG Chain Agent")
if to:
await client.send_email(to=to, subject=subject, body=content)
else:
logger.error(f"Unsupported platform: {platform}")
return False
await self._log_action("send", {
"platform": platform, "content_preview": content[:100],
})
return True
except Exception as e:
logger.error(f"Send to {platform} failed: {e}")
return False
async def broadcast(self, content: str, platforms: Optional[List[str]] = None) -> Dict[str, bool]:
"""Broadcast a message to multiple platforms."""
targets = platforms or list(self.platforms.keys())
results = {}
for platform in targets:
results[platform] = await self.send(platform, content)
await self._log_action("broadcast", {
"platforms": targets,
"success_count": sum(1 for v in results.values() if v),
"total": len(targets),
})
return results
async def cross_post(self, content: dict, platforms: Optional[List[str]] = None):
"""Cross-post with platform-specific formatting."""
targets = platforms or list(self.platforms.keys())
results = {}
for platform in targets:
formatter = self.formatters.get(platform)
if formatter:
formatted = formatter(content)
results[platform] = await self.send(platform, formatted)
else:
results[platform] = await self.send(platform, content.get("text", ""))
return results
async def enqueue(self, message: PlatformMessage):
"""Add a message to the send queue."""
await self.message_queue.put(message)
async def process_queue(self):
"""Process the message queue continuously."""
self._running = True
logger.info("Message queue processor started")
while self._running:
try:
message = await asyncio.wait_for(self.message_queue.get(), timeout=1.0)
success = await self.send(
message.platform, message.content,
thread_id=message.thread_id,
)
if not success:
logger.warning(f"Failed to send queued message to {message.platform}")
self.message_queue.task_done()
except asyncio.TimeoutError:
continue
except Exception as e:
logger.error(f"Queue processing error: {e}")
def stop_queue(self):
"""Stop the message queue processor."""
self._running = False
async def bridge(self, source_platform: str, target_platforms: List[str], message: str):
"""Bridge a message from one platform to others."""
logger.info(f"Bridging from {source_platform} to {target_platforms}")
results = {}
for target in target_platforms:
if target != source_platform:
results[target] = await self.send(target, message)
return results
async def _log_action(self, action: str, data: dict):
try:
await self.chain.submit_log({
"agent_id": getattr(self.agent, "agent_id", "unknown"),
"action": f"messenger_{action}",
"data": data,
"timestamp": datetime.utcnow().isoformat(),
})
except Exception as e:
logger.warning(f"Failed to log: {e}")
class MessageRouter:
"""Route incoming messages to appropriate handlers."""
def __init__(self, messenger: UnifiedMessenger):
self.messenger = messenger
self.routes: Dict[str, Callable] = {}
def route(self, platform: str, handler: Callable):
"""Register a handler for a platform's incoming messages."""
self.routes[platform] = handler
async def handle_incoming(self, platform: str, message: PlatformMessage):
"""Handle an incoming message from a platform."""
handler = self.routes.get(platform)
if not handler:
logger.warning(f"No handler for incoming {platform} message")
return None
return await handler(message)
async def auto_respond(self, message: PlatformMessage) -> Optional[str]:
"""Auto-generate a response to an incoming message."""
content = message.content.lower()
if "balance" in content or "余额" in content:
return "Checking your balance on MSG Chain..."
if "help" in content or "帮助" in content:
return "Available: balance, price, chain status queries."
return None
# Example
if __name__ == "__main__":
class DemoAgent:
def __init__(self):
self.agent_id = "router-agent-001"
agent = DemoAgent()
messenger = UnifiedMessenger(agent)
# Register formatters
@messenger.register_formatter
def custom_discord(content: dict) -> str:
return f"**[{content.get('title', '')}]** {content.get('text', '')}"
async def demo():
# Broadcast to all platforms
results = await messenger.broadcast(
"Hello from MSG Chain AI Agent!",
platforms=["discord", "telegram"],
)
print(f"Broadcast results: {results}")
# Cross-post with formatting
content = {
"text": "MSG Chain block height is now 1,000,000!",
"title": "Chain Update",
"hashtags": ["msgchain", "milestone"],
}
results = await messenger.cross_post(content, platforms=["twitter", "discord"])
print(f"Cross-post results: {results}")
asyncio.run(demo())
9.3 路由策略
| 模式 | 说明 | 适用场景 |
|---|---|---|
| Broadcast | 同一消息发送到所有平台 | 公告、警报 |
| Cross-post | 各平台独立格式化后发送 | 内容发布 |
| Bridge | 从一个平台中继到其他平台 | 跨平台客服 |
| Queue | 异步队列发送 | 高吞吐量通知 |
10. 完整示例:多平台监控 Agent
以下示例展示了一个完整的、可运行的多平台监控 Agent,它同时监听 Twitter 提及、Discord 消息和 GitHub Issue,并将事件统一路由到所有配置的平台。
10.1 完整代码
#!/usr/bin/env python3
"""
Multi-Platform Monitoring Agent for MSG Chain.
Monitors Twitter mentions, Discord channels, and GitHub issues
simultaneously and broadcasts events across all platforms.
"""
import os
import asyncio
import logging
from typing import Optional, Dict, Any
from datetime import datetime
from dotenv import load_dotenv
from msg_chain_sdk import MsgChainClient
# Platform agents
from twitter_agent import TwitterAgent
from discord_agent import DiscordAgent
from github_agent import GitHubAgent
from telegram_agent import TelegramAgent
from slack_agent import SlackAgent
from email_agent import EmailAgent
from webhook_server import WebhookServer, WebhookRouter
from unified_messenger import UnifiedMessenger, PlatformFormatter
load_dotenv()
logger = logging.getLogger(__name__)
class MonitoringAgent:
"""Multi-platform monitoring agent for MSG Chain.
Monitors events across all integrated platforms and
routes notifications to configured channels.
"""
def __init__(self, agent_id: str = "monitor-agent-001"):
self.agent_id = agent_id
self.chain = MsgChainClient()
# Initialize platform clients
self.twitter = TwitterAgent(agent_id)
self.discord = DiscordAgent(agent_id)
self.telegram = TelegramAgent(agent_id)
self.github = GitHubAgent(agent_id)
self.slack = SlackAgent(agent_id)
self.email = EmailAgent(agent_id)
# Initialize unified messenger
self.messenger = UnifiedMessenger(self)
self.messenger.register_platform("twitter", self.twitter)
self.messenger.register_platform("discord", self.discord)
self.messenger.register_platform("telegram", self.telegram)
self.messenger.register_platform("slack", self.slack)
self.messenger.register_platform("email", self.email)
# Webhook server for receiving events
self.webhook = WebhookServer(self, port=int(os.getenv("WEBHOOK_PORT", "8080")))
self._setup_webhook_handlers()
# Event counter for reporting
self.event_counts: Dict[str, int] = {
"twitter_mentions": 0,
"discord_messages": 0,
"github_issues": 0,
"telegram_commands": 0,
"slack_mentions": 0,
}
def _setup_webhook_handlers(self):
"""Register webhook event handlers."""
@self.webhook.router.register("github:issues")
async def on_github_issue(payload: dict):
action = payload.get("action", "")
issue = payload.get("issue", {})
repo = payload.get("repository", {}).get("full_name", "unknown")
self.event_counts["github_issues"] += 1
notification = (
f"GitHub Issue #{issue.get('number')} ({action}) in {repo}\\n"
f"Title: {issue.get('title')}\\n"
f"By: {issue.get('user', {}).get('login', 'unknown')}"
)
# Broadcast to all platforms
await self.messenger.broadcast(notification)
await self._log_event("github_issue", {
"action": action, "issue_number": issue.get("number"), "repo": repo,
})
return {"processed": True}
@self.webhook.router.register("github:push")
async def on_github_push(payload: dict):
repo = payload.get("repository", {}).get("full_name", "unknown")
ref = payload.get("ref", "").replace("refs/heads/", "")
commits = payload.get("commits", [])
self.event_counts["github_issues"] += 1
notification = (
f"Push to {repo}/{ref}\\n"
f"Commits: {len(commits)}"
)
await self.messenger.broadcast(notification)
return {"processed": True}
async def start_monitoring(self):
"""Start all monitoring loops."""
logger.info(f"Starting monitoring agent: {self.agent_id}")
tasks = [
self._monitor_twitter(),
self._monitor_github(),
self._run_webhook(),
self._generate_periodic_report(),
]
await asyncio.gather(*tasks)
async def _monitor_twitter(self):
"""Monitor Twitter mentions for keywords."""
logger.info("Starting Twitter mention monitoring")
while True:
try:
mentions = await self.twitter.monitor_mentions()
for mention, reply in mentions:
self.event_counts["twitter_mentions"] += 1
notification = (
f"Twitter mention from @{mention.author_id}:\\n"
f"{mention.text[:100]}"
)
await self.messenger.broadcast(
notification, platforms=["discord", "telegram", "slack"]
)
except Exception as e:
logger.error(f"Twitter monitor error: {e}")
await asyncio.sleep(60)
async def _monitor_github(self):
"""Monitor GitHub for new issues and PRs."""
repos = os.getenv("GITHUB_MONITORED_REPOS", "").split(",")
if not repos:
logger.info("No GitHub repos to monitor")
return
last_checks = {repo: datetime.utcnow() for repo in repos}
while True:
for repo in repos:
try:
new_issues = await self.github.list_issues(
repo, state="open", max_results=5
)
for issue in new_issues:
created = datetime.fromisoformat(issue["created_at"])
if created > last_checks[repo]:
self.event_counts["github_issues"] += 1
notification = (
f"New issue in {repo}: #{issue['number']} "
f"{issue['title']}"
)
await self.messenger.broadcast(notification)
last_checks[repo] = datetime.utcnow()
except Exception as e:
logger.error(f"GitHub monitor error for {repo}: {e}")
await asyncio.sleep(300)
async def _run_webhook(self):
"""Run the webhook server."""
logger.info("Starting webhook server")
self.webhook.start()
async def _generate_periodic_report(self):
"""Generate and send periodic status reports."""
while True:
await asyncio.sleep(3600)
height = await self.chain.get_block_height()
report = (
f"MSG Chain Agent Hourly Report\\n"
f"Time: {datetime.utcnow().isoformat()}\\n"
f"Block Height: {height}\\n"
f"Events Processe[未公开路径]"
)
for event_type, count in self.event_counts.items():
report += f" - {event_type}: {count}\\n"
await self.messenger.broadcast(report)
async def _log_event(self, event_type: str, data: dict):
"""Log event on MSG Chain."""
try:
await self.chain.submit_log({
"agent_id": self.agent_id,
"action": f"event_{event_type}",
"data": data,
"timestamp": datetime.utcnow().isoformat(),
})
except Exception as e:
logger.warning(f"Failed to log event: {e}")
async def cleanup(self):
"""Cleanup resources."""
logger.info("Cleaning up monitoring agent")
if hasattr(self.messenger, 'message_queue'):
self.messenger.stop_queue()
if hasattr(self.twitter, 'stream') and self.twitter.stream:
self.twitter.stream.disconnect()
async def main():
"""Main entry point for the monitoring agent."""
agent = MonitoringAgent()
try:
await agent.start_monitoring()
except KeyboardInterrupt:
logger.info("Shutting down...")
finally:
await agent.cleanup()
if __name__ == "__main__":
logging.basicConfig(level=logging.INFO)
asyncio.run(main())
10.2 环境变量配置
# Chain
export CHAIN_RPC="https://rpc.msgchain.io"
# Twitter
export TWITTER_BEARER_TOKEN="..."
export TWITTER_API_KEY="..."
export TWITTER_API_SECRET="..."
export TWITTER_ACCESS_TOKEN="..."
export TWITTER_ACCESS_SECRET="..."
# Discord
export DISCORD_BOT_TOKEN="..."
# Telegram
export TELEGRAM_BOT_TOKEN="..."
# GitHub
export GITHUB_TOKEN="..."
export GITHUB_MONITORED_REPOS="msgchain/msg-chain-core,msgchain/agent-sdk"
# Slack
export SLACK_BOT_TOKEN="xoxb-..."
export SLACK_APP_TOKEN="xapp-..."
# Email
export SMTP_HOST="smtp.gmail.com"
export SMTP_PORT=587
export SMTP_USER="agent@msgchain.org"
export SMTP_PASSWORD="..."
# Webhook
export WEBHOOK_PORT=8080
export WEBHOOK_SECRET="your-webhook-secret"
10.3 运行
# Install dependencies
pip install tweepy discord.py python-telegram-bot PyGithub slack-bolt aiosmtplib fastapi uvicorn
# Set environment variables
export AGENT_ID="monitor-agent-001"
export CHAIN_RPC="https://rpc.msgchain.io"
# Run the agent
python monitoring_agent.py
11. 附录
A. 常用错误排查
| 错误 | 可能原因 | 解决方案 |
|---|---|---|
| 401 Unauthorized | API Token 过期或无效 | 重新生成 Token 并更新环境变量 |
| Rate Limit Exceeded | 请求频率过高 | 启用 wait_on_rate_limit,降低请求频率 |
| Webhook 验证失败 | HMAC 签名不匹配 | 检查 WEBHOOK_SECRET 是否一致 |
| Connection Timeout | 网络不通或防火墙 | 检查 RPC 节点可达性 |
| Invalid Bech32 | 地址格式错误 | 确认地址以 msg1 开头 |
B. 链上日志 Schema
{
"agent_id": "string",
"action": "string",
"data": {
"platform": "string",
"event_type": "string",
"payload": {}
},
"timestamp": "2026-07-07T12:00:00Z",
"signature": "string"
}
C. 推荐项目结构
agent/
├── agents/
│ ├── twitter_agent.py
│ ├── discord_agent.py
│ ├── telegram_agent.py
│ ├── github_agent.py
│ ├── slack_agent.py
│ └── email_agent.py
├── core/
│ ├── messenger.py # UnifiedMessenger
│ ├── webhook.py # WebhookServer/Router
│ └── chain.py # MSG Chain client wrapper
├── monitoring_agent.py # Main entry point
├── config.py # Configuration
├── requirements.txt
└── .env # Environment variables
D. 安全 Checklist
- [ ] 所有 API Token 使用环境变量,不硬编码
- [ ] GitHub Token 使用最小权限 scope
- [ ] Webhook Secret 已配置并验证签名
- [ ] Rate Limiting 已实现
- [ ] 链上日志记录所有外部操作
- [ ] 敏感操作(转账、合约部署)有多签
- [ ] 定期轮换 API 凭据
文档版本: 1.0.0 | 链: MSG Chain (msg-chain-1) | 地址前缀: msg
本文档中的代码示例仅供参考。生产部署前请进行全面测试和安全审计。
