dApp Docs/AI Agent 第三方服务集成指南
Development reference. Not independently verified for production.

AI Agent 第三方服务集成指南

适用链: MSG Chain | 地址前缀: msg | 版本: 1.0.0

⚠️ No-Go Disclaimer: MSGChain 主网裁决为 No-Go。本文件所有内容反映的是开发阶段的技术设计,不代表主网未独立核验上线状态。生产部署状态请以白皮书为准:https://msgchain.org/whitepaper/


目录

  1. 概述
  2. Twitter/X 集成
  3. Discord 集成
  4. Telegram 集成
  5. GitHub 集成
  6. Slack 集成
  7. Email (SMTP/IMAP) 集成
  8. Webhook 系统
  9. 跨平台消息路由
  10. 完整示例:多平台监控 Agent
  11. 附录

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 通用安全实践

2. Twitter/X 集成

2.1 准备工作

  1. 在 Twitter Developer Portal 创建 Project
  2. 启用 OAuth 2.0 + OAuth 1.0a (User Context)
  3. 生成 API Key、API Secret、Bearer Token、Access Token & Secret
  4. 设置环境变量:
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 准备工作

  1. 在 Discord Developer Portal 创建 Application
  2. 创建 Bot,获取 Token
  3. 在 OAuth2 URL Generator 选择 bot + applications.commands 权限
  4. 设置 Bot 权限: Send Messages, Read Message History, Mention Everyone, Use Slash Commands
  5. 将 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 准备工作

  1. 在 @BotFather 创建 Bot,获取 Token
  2. 设置 Bot 命令: /start /balance /help /query /monitor
  3. 可选: 配置 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 准备工作

  1. 在 GitHub Settings Tokens 生成 Personal Access Token
  2. 权限: repo, issues, pull_requests, contents
  3. 可选: 安装 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 准备工作

  1. 在 Slack API 创建 App
  2. 启用 Socket Mode 或配置 Request URL
  3. 添加 Bot Token Scopes: chat:write, channels:history, channels:read, users:read
  4. 安装 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


文档版本: 1.0.0 | 链: MSG Chain (msg-chain-1) | 地址前缀: msg

本文档中的代码示例仅供参考。生产部署前请进行全面测试和安全审计。