Merge pull request #17 from EvoScientist/bug-fix

Bug fix
This commit is contained in:
Xi Zhang
2026-02-22 14:32:07 +00:00
committed by GitHub
7 changed files with 587 additions and 163 deletions
+431 -86
View File
@@ -1,76 +1,290 @@
# Channels
EvoScientist provides unified integration with 11 messaging platforms. This document covers the architecture overview, capability matrix, and detailed deployment guide for each channel.
EvoScientist provides unified integration with 10 messaging platforms. This document covers the architecture overview, message processing pipeline, capability matrix, security model, deployment guides, and troubleshooting.
Configuration file: `~/.config/evoscientist/config.yaml` (or use environment variables with the `EVOSCIENTIST_` prefix).
## Table of Contents
- [Architecture](#architecture)
- [Message Processing Pipeline](#message-processing-pipeline)
- [Middleware Pipeline](#middleware-pipeline)
- [Capability Matrix](#capability-matrix)
- [Security and Access Control](#security-and-access-control)
- [Quick Start](#quick-start)
- [Channel Deployment Guides](#channel-deployment-guides)
- [Telegram](#telegram) | [Discord](#discord) | [Slack](#slack) | [Feishu (Lark)](#feishu-lark) | [WeChat](#wechat)
- [DingTalk](#dingtalk) | [QQ](#qq) | [Signal](#signal) | [Email](#email) | [iMessage](#imessage)
- [Running Multiple Channels](#running-multiple-channels)
- [Docker Deployment](#docker-deployment)
- [Troubleshooting](#troubleshooting)
## Architecture
```
┌──────────┐ ┌──────────┐ ┌──────────┐
│ Telegram │ │ Discord │ │ Slack │ ... (×11)
└────┬─────┘ └────┬─────┘ └────┬─────┘
│ │ │
└─────────────┼─────────────┘
▼
┌──────────────┐
│ MessageBus │ async queue, 5000 cap
└──────┬───────┘
▼
┌──────────────┐
│InboundConsumer│ → Agent → OutboundMessage
└──────┬───────┘
▼
┌──────────────┐
│ Dispatcher │ routes replies to origin channel
└──────────────┘
┌─────────────────────────────────────────────┐
│ Messaging Platforms │
│ │
│ ┌────────┐ ┌────────┐ ┌────────┐ │
│ │Telegram│ │Discord │ │ Slack │ ...x10 │
│ └───┬────┘ └───┬────┘ └───┬────┘ │
└──────┼──────────┼──────────┼───────────────┘
│ │ │
┌──────┴──────────┴──────────┴───────────────┐
│ Inbound Middleware │
│ │
│ Dedup → AllowList → Pairing → GroupHist │
│ → Mention │
└──────────────────┬─────────────────────────┘
│
▼
┌──────────────────────────────────────┐
│ MessageBus │
│ │
│ inbound queue ──► outbound queue │
│ (asyncio.Queue, capacity 5000) │
└──────────┬───────────────┬───────────┘
│ │
▼ ▼
┌──────────────────┐ ┌─────────────────┐
│ InboundConsumer │ │ Dispatcher │
│ │ │ │
│ Worker pool (8) │ │ Routes replies │
│ Per-chat locks │ │ to origin │
│ Session dedup │ │ channel │
│ Timeout handling │ │ │
│ │ │ └─────────────────┘
│ ▼ │
│ Agent Core │
└──────────────────┘
```
**Core modules:**
### Core Modules
| Module | Responsibility |
|--------|---------------|
| `base.py` | Abstract `Channel` base class — declarative readiness checks, retry strategy, mention stripping, send fallback, media handling |
| `capabilities.py` | `ChannelCapabilities` frozen dataclass — each channel declares its capabilities, framework adapts automatically |
| `mixins.py` | Reusable patterns: `WebhookMixin` (aiohttp + httpx), `WebSocketMixin` (connect/reconnect/heartbeat), `PollingMixin` (async polling), `TokenMixin` (OAuth token refresh) |
| `config.py` | `BaseChannelConfig` — shared config fields (allowed_senders, proxy, text_chunk_limit, etc.) |
| `bus/` | `MessageBus` async message queue + `InboundMessage`/`OutboundMessage` dataclasses |
| `channel_manager.py` | Lifecycle management (start/stop), health checks, channel registry |
| `consumer.py` | `InboundConsumer` — dequeue messages, invoke Agent, publish replies |
| `retry.py` | Configurable exponential backoff retry (`RetryConfig`: attempts, min/max delay, jitter) |
| `markdown_utils.py` | Universal Markdown converter with per-platform formatting plugins |
| `base.py` | Abstract `Channel` base class — readiness checks, retry strategy, mention stripping, format fallback, media handling, debounce, send locks |
| `capabilities.py` | `ChannelCapabilities` frozen dataclass — each channel declares features, framework adapts automatically |
| `plugin.py` | `ChannelPlugin` base with adapter slots — `ConfigAdapter`, `SecurityAdapter`, `GroupAdapter`, `MentionAdapter`, `OutboundAdapter`, `ThreadingAdapter`, etc. |
| `mixins.py` | Reusable async patterns: `WebhookMixin` (aiohttp server + httpx client), `WebSocketMixin` (connect/reconnect/heartbeat), `PollingMixin` (async polling loop), `TokenMixin` (OAuth token auto-refresh) |
| `config.py` | `BaseChannelConfig` — shared config fields (allowed_senders, proxy, text_chunk_limit, etc.) + `SingleAccountConfigAdapter` / `MultiAccountConfigAdapter` |
| `bus/` | `MessageBus` async event bus with `InboundMessage` / `OutboundMessage` dataclasses, decoupling channels from agent core |
| `channel_manager.py` | `ChannelManager` — lifecycle management (start/stop), health monitoring, channel registry, account management, outbound dispatch |
| `consumer.py` | `InboundConsumer` — worker pool, per-chat serial locks, session deduplication, timeout handling |
| `retry.py` | `RetryConfig` — exponential backoff retry with per-channel presets (attempts, min/max delay, jitter) |
| `formatter.py` | `UnifiedFormatter` — Markdown to platform-specific format conversion (HTML, Slack mrkdwn, Discord, plain text) |
| `standalone.py` | Headless channel runner (`run_standalone`) for running channels without the CLI |
## Message Processing Pipeline
### Inbound (User Message → Agent)
```
1. Platform SDK/Webhook receives raw message
│
2. Channel._on_message() parses into RawIncoming
│
3. Channel._enqueue_raw() runs middleware pipeline:
├── DedupMiddleware — drop duplicates (LRU cache, 60s TTL)
├── AllowListMiddleware — enforce sender/channel restrictions
├── PairingMiddleware — handle DM pairing flow (if dm_policy="pairing")
├── GroupHistoryMiddleware — buffer group context for injection
└── MentionGatingMiddleware — filter by @mention policy in groups
│
4. InboundMessage queued on Channel._queue
│
5. Channel.run() → receive() → queue_message() with debounce
│
│ (500ms debounce window: rapid messages from same sender merged)
│
6. MessageBus.publish_inbound()
│
7. InboundConsumer acquires per-chat lock → invokes Agent
│
8. Agent response → OutboundMessage → MessageBus.publish_outbound()
```
### Outbound (Agent Response → User)
```
1. OutboundMessage arrives on MessageBus outbound queue
│
2. Dispatcher routes to origin channel by name
│
3. Channel.send() processes the response:
├── Stop typing indicator
├── Format text (Markdown → platform format)
├── Chunk text to platform limit (code-block-aware splitting)
├── Send each chunk via _send_chunk() with format fallback
├── Send media attachments via _send_media_impl()
└── Retry on transient errors (exponential backoff)
```
### Text Chunking
Long responses are split intelligently with this priority:
1. Markdown code block fence boundaries
2. Double newlines (paragraph breaks)
3. Single newlines
4. Space characters
5. Hard cut at limit (last resort)
Code blocks are never split mid-block when possible. Each chunk is sent as a separate message.
## Middleware Pipeline
Middleware runs sequentially on each inbound message. Each middleware can pass, modify, or drop the message.
### 1. DedupMiddleware
Prevents duplicate message processing using a bounded LRU cache with TTL.
- Cache size: 1000 entries (configurable)
- TTL: 60 seconds
- Key: `message_id` from the platform
- Messages with the same ID within the TTL window are silently dropped
### 2. AllowListMiddleware
Enforces sender and channel restrictions based on the `dm_policy` config.
| Policy | Behavior |
|--------|----------|
| `"open"` | Accept messages from anyone |
| `"allowlist"` | Only accept from `allowed_senders` / `allowed_channels` |
| `"pairing"` | Require DM pairing before accepting (see PairingMiddleware) |
When `allowed_senders` is set (non-empty), only messages from listed sender IDs pass through. Same for `allowed_channels`.
### 3. PairingMiddleware
Handles an interactive DM pairing flow for the `"pairing"` dm_policy.
- First message from an unknown sender triggers a pairing request
- The sender must provide a valid pairing code
- Once paired, the sender is added to the allowlist for future messages
### 4. GroupHistoryMiddleware
Buffers recent group chat messages to provide conversation context.
- Only active when `capabilities.groups = True`
- Maintains a per-chat rolling buffer (default: 50 messages, 5-minute max age)
- When the bot is mentioned in a group, recent history is injected into the message metadata so the agent can see prior context
- Non-mentioned group messages are buffered but not forwarded (see MentionGating)
### 5. MentionGatingMiddleware
Controls whether the bot responds in group chats.
| `require_mention` | Behavior |
|-------------------|----------|
| `True` / `"group"` | Only respond to @mentions in groups; always respond in DMs |
| `False` / `"none"` | Respond to all messages in all contexts |
| `"always"` | Require @mention even in DMs |
Default: `"group"` — the bot ignores group messages unless explicitly @mentioned.
Mention detection is platform-specific:
- **Telegram**: checks for `@bot_username` in text
- **Discord**: checks `message.mentions` for bot user
- **Slack**: handled via separate `app_mention` event type
- **Feishu**: checks `mentions` array in event payload
- **DingTalk**: checks `isInAtList` flag or `atUsers` array
- **WeChat (WeCom)**: checks `AtUserList` XML field
## Capability Matrix
| Channel | Format | Max Len | Media | Voice | Sticker | Location | Video | Typing | Reaction | Thread | Group | @Mention | No Public IP | Token Refresh | Proxy | Allowlist |
|:--------|:------:|:-------:|:-----:|:-----:|:-------:|:--------:|:-----:|:------:|:--------:|:------:|:-----:|:--------:|:------------:|:-------------:|:-----:|:---------:|
| Telegram | HTML | 4000 | ✓ | ✓ | ✓ | ✓ | | ✓ | ✓ | | ✓ | ✓ | ✓ | | ✓ | ✓ |
| Discord | Discord | 2000 | ✓ | | | | | ✓ | ✓ | ✓ | ✓ | ✓ | ✓ | | ✓ | ✓ |
| Slack | Mrkdwn | 4000 | ✓ | | | | | ✓ | ✓ | ✓ | ✓ | ✓ | ✓ | | ✓ | ✓ |
| Feishu | MD | 4096 | ✓ | ✓ | ✓ | | | | ✓ | | ✓ | ✓ | | ✓ | ✓ | ✓ |
| WeChat | MD | 4096 | ✓ | ✓ | | ✓ | | | | | ✓ | ✓ | | ✓ | ✓ | ✓ |
| DingTalk | MD | 4096 | ✓ | ✓ | | | | | | | ✓ | ✓ | ✓ | ✓ | ✓ | ✓ |
| QQ | Plain | 4096 | ✓ | | | | | | | | ✓ | ✓ | ✓ | | | ✓ |
| Signal | Plain | 4096 | ✓ | ✓ | | | | ✓ | ✓ | | ✓ | ✓ | ✓ | | | ✓ |
| iMessage | Plain | ∞ | ✓ | ✓ | | | | | | | ✓ | | ✓ | | | ✓ |
| Email | HTML | ∞ | ✓ | | | | | | | | | | ✓ | | | ✓ |
| Telegram | HTML | 4000 | S/R | R | R | R | | 4s | emoji | | G | @ | yes | | yes | yes |
| Discord | Discord | 2000 | S/R | | | | | 8s | emoji | yes | G | @ | yes | | yes | yes |
| Slack | Mrkdwn | 4000 | S/R | | | | | post | emoji | yes | G | @ | yes | | yes | yes |
| Feishu | Post | 4096 | S/R | R | R | | | | emoji | | G | @ | no | 2h | yes | yes |
| WeChat | MD | 4096 | S/R | R | | R | R | recall | | | G | @ | no | 2h | yes | yes |
| DingTalk | MD | 4096 | S/R | R | | | R | | | | G | @ | yes | 2h | yes | yes |
| QQ | Plain | 4096 | S/R | | | | | | | | G | @ | yes | | | yes |
| Signal | Plain | 4096 | S/R | R | | | | api | emoji | | G | UUID | yes | | | yes |
| iMessage | Plain | - | S/R | R | | | | | | | G | | yes | | | yes |
| Email | HTML | - | S/R | | | | | | | | | | yes | | | yes |
Legend: **S** = send, **R** = receive, **G** = group chat, **@** = @mention detection, **-** = no practical limit
### Connection Types
| Channel | Transport | Connection Mode | Default Port |
|-------------|-----------|----------------------------------------|:------------:|
| Telegram | HTTPS | Long polling (`getUpdates`) | — |
| Discord | WebSocket | Gateway events (`discord.py`) | — |
| Slack | WebSocket | Socket Mode (`slack-sdk`) | — |
| Telegram | HTTPS | Long polling (`getUpdates`) | -- |
| Discord | WebSocket | Gateway events (`discord.py`) | -- |
| Slack | WebSocket | Socket Mode (`slack-sdk`) | -- |
| Feishu | HTTP | Webhook `POST /webhook/event` | 9000 |
| WeChat | HTTP | Webhook `POST /wechat/callback` | 9001 |
| DingTalk | WebSocket | Stream Mode (DingTalk gateway) | — |
| QQ | WebSocket | Bot Gateway (`qq-botpy`) | — |
| DingTalk | WebSocket | Stream Mode (DingTalk gateway) | -- |
| QQ | WebSocket | Bot Gateway (`qq-botpy`) | -- |
| Signal | TCP | JSON-RPC (`signal-cli` daemon) | 7583 |
| iMessage | stdio | JSON-RPC (`imsg` CLI) | — |
| iMessage | stdio | JSON-RPC (`imsg` CLI) | -- |
| Email | TCP | IMAP polling + SMTP send | 993/587 |
> **"—"** means no listening port is required — no public IP or port forwarding needed.
> **"--"** means no listening port is required -- no public IP or port forwarding needed.
### Format Conversion
The `UnifiedFormatter` converts Markdown output from the agent into platform-native formats:
| Target Format | Conversion |
|:-------------|:-----------|
| HTML (Telegram, Email) | `**bold**` → `<b>bold</b>`, `` `code` `` → `<code>code</code>`, code blocks → `<pre>`, special chars escaped |
| Slack mrkdwn | `**bold**` → `*bold*`, `_italic_` → `_italic_`, code blocks preserved, `<>&` escaped |
| Discord Markdown | Mostly passthrough, minor adjustments for Discord-specific rendering |
| Feishu Post | Markdown → Feishu rich text JSON (code blocks, bold, italic, strikethrough, links, headings, quotes, lists) |
| Plain text | All formatting stripped, structure preserved via indentation |
## Security and Access Control
### Sender Allowlist
Every channel supports `allowed_senders` to restrict who can interact with the bot:
```yaml
telegram_allowed_senders: "123456789,987654321" # Telegram user IDs
discord_allowed_senders: "111222333444555666" # Discord user IDs
slack_allowed_senders: "U0123ABCDEF" # Slack Member IDs
feishu_allowed_senders: "ou_xxxxxxxxxxxx" # Feishu open_ids
signal_allowed_senders: "+1234567890" # Phone numbers
email_allowed_senders: "alice@example.com" # Email addresses
imessage_allowed_senders: "+1234567890,user@icloud.com" # Phone or email
```
When `allowed_senders` is empty, the channel accepts messages from anyone. **For production deployments, always set an allowlist.**
### Channel Allowlist
For platforms with multiple channels/groups (Discord, Slack), restrict which channels the bot operates in:
```yaml
discord_allowed_channels: "111222333444555666,777888999000111222"
slack_allowed_channels: "C0123ABCDEF,C0456GHIJKL"
```
### Token and Secret Handling
- API tokens are stored in the config file or environment variables, never logged at INFO level
- Discord logs only the first 8 and last 4 characters of the bot token for debugging
- WeChat/Feishu tokens are auto-refreshed before expiry (5-minute margin on 2-hour TTL)
- Webhook signature verification is enforced when `token`/`encoding_aes_key` is configured (WeChat, Feishu)
### Group Chat Behavior
By default, the bot only responds in group chats when explicitly @mentioned. This prevents the bot from responding to every message in a busy group. Configure via:
```yaml
# Default: only respond when mentioned in groups
channel_require_mention: "group"
# Respond to all messages (including groups)
channel_require_mention: "none"
```
## Quick Start
@@ -119,16 +333,6 @@ curl http://localhost:8080/healthz
}
```
### Running multiple channels
Comma-separate channel names in the config to enable multiple channels simultaneously:
```yaml
channel_enabled: "telegram,discord,imessage"
```
All enabled channels run concurrently via the internal message bus.
---
## Channel Deployment Guides
@@ -142,9 +346,9 @@ All enabled channels run concurrently via the internal message bus.
**Prerequisites:**
1. Search for [@BotFather](https://t.me/BotFather) in Telegram, send `/newbot`, and follow the prompts to create a bot.
2. BotFather will return a Bot Token (format: `123456789:ABCdefGHI...`) — save it securely.
2. BotFather will return a Bot Token (format: `123456789:ABCdefGHI...`) -- save it securely.
3. Get your user ID: send any message to [@userinfobot](https://t.me/userinfobot), it will reply with your numeric ID.
4. (Optional) For group use: add the bot to a group, then in BotFather send `/setprivacy` → `Disable` so the bot can read group messages.
4. (Optional) For group use: add the bot to a group, then in BotFather send `/setprivacy` -> `Disable` so the bot can read group messages.
**Configuration:**
@@ -163,7 +367,7 @@ telegram_proxy: "" # Optional HTTPS proxy (e.g. http://proxy:808
**Env vars:** `EVOSCIENTIST_TELEGRAM_BOT_TOKEN`, `EVOSCIENTIST_TELEGRAM_ALLOWED_SENDERS`, `EVOSCIENTIST_TELEGRAM_PROXY`
**Technical details:** Long polling mode, `drop_pending_updates=True` on startup to skip backlog. Markdown→Telegram HTML auto-conversion (bold, italic, strikethrough, links, code blocks, headings, lists). Falls back to plain text on HTML parse failure. Media routed by extension to `send_photo`/`send_video`/`send_audio`/`send_document`. In groups, only responds when @mentioned; auto-strips @mention. Typing indicator refreshes every 4s. Retry: 3 attempts, min delay 0.4s, parse errors not retried. Text chunk limit: 4000 chars.
**Technical details:** Long polling mode, `drop_pending_updates=True` on startup to skip backlog. Markdown to Telegram HTML auto-conversion (bold, italic, strikethrough, links, code blocks, headings, lists). Falls back to plain text on HTML parse failure. Media routed by extension to `send_photo`/`send_video`/`send_audio`/`send_document`. In groups, only responds when @mentioned; auto-strips @mention. Typing indicator refreshes every 4s. ACK reaction (eyes emoji) on message receipt, removed after reply. Retry: 3 attempts, min delay 0.4s, parse errors not retried. Text chunk limit: 4000 chars.
---
@@ -173,14 +377,14 @@ telegram_proxy: "" # Optional HTTPS proxy (e.g. http://proxy:808
**Prerequisites:**
1. Go to [Discord Developer Portal](https://discord.com/developers/applications) → New Application → enter a name.
2. Left menu **Bot** → Reset Token → copy the Bot Token.
1. Go to [Discord Developer Portal](https://discord.com/developers/applications) -> New Application -> enter a name.
2. Left menu **Bot** -> Reset Token -> copy the Bot Token.
3. Under **Privileged Gateway Intents**, enable **Message Content Intent** (required to read message content).
4. Left menu **OAuth2** → URL Generator:
4. Left menu **OAuth2** -> URL Generator:
- Scopes: check `bot`
- Bot Permissions: check `Send Messages`, `Read Message History`, `Attach Files`, `Add Reactions`
- Copy the generated URL, open in browser, select a server to invite the bot.
5. Get user ID: Discord Settings → Advanced → enable Developer Mode → right-click username → Copy User ID.
5. Get user ID: Discord Settings -> Advanced -> enable Developer Mode -> right-click username -> Copy User ID.
**Configuration:**
@@ -201,7 +405,7 @@ discord_proxy: ""
**Env vars:** `EVOSCIENTIST_DISCORD_BOT_TOKEN`, `EVOSCIENTIST_DISCORD_ALLOWED_SENDERS`, `EVOSCIENTIST_DISCORD_ALLOWED_CHANNELS`, `EVOSCIENTIST_DISCORD_PROXY`
**Technical details:** WebSocket Gateway (`discord.py`). In server channels, only responds when @mentioned; DMs respond directly. Replies via `MessageReference`. Attachment download (max 20 MB) with safe filename sanitization. Media sent via `discord.File`. Typing indicator refreshes every 8s. Retry: 3 attempts, parses `Retry-After` header for 429s. Text chunk limit: 2000 chars.
**Technical details:** WebSocket Gateway (`discord.py`). In server channels, only responds when @mentioned; DMs respond directly. Thread-aware: messages in threads are tracked with `parent_channel_id` and `thread_id`. Replies via `MessageReference`. Message cache (200 entries) for ACK emoji reactions. Attachment download (max 20 MB) with safe filename sanitization. Media sent via `discord.File`. Typing indicator refreshes every 8s. Retry: 3 attempts, parses `Retry-After` header for 429s. Text chunk limit: 2000 chars.
---
@@ -211,13 +415,13 @@ discord_proxy: ""
**Prerequisites:**
1. Go to [Slack API](https://api.slack.com/apps) → Create New App → From scratch → select workspace.
2. Left menu **Socket Mode** → enable → Generate App-Level Token, scope `connections:write` → copy App Token (`xapp-...`).
3. Left menu **OAuth & Permissions** → add Bot Token Scopes:
1. Go to [Slack API](https://api.slack.com/apps) -> Create New App -> From scratch -> select workspace.
2. Left menu **Socket Mode** -> enable -> Generate App-Level Token, scope `connections:write` -> copy App Token (`xapp-...`).
3. Left menu **OAuth & Permissions** -> add Bot Token Scopes:
- `chat:write`, `channels:history`, `groups:history`, `im:history`, `files:read`, `files:write`, `reactions:write`
4. Click **Install to Workspace** → copy Bot User OAuth Token (`xoxb-...`).
5. Left menu **Event Subscriptions** → enable → Subscribe to bot events: `message.channels`, `message.groups`, `message.im`, `app_mention`.
6. Get Member ID: click user avatar → profile → **⋮** → Copy member ID.
4. Click **Install to Workspace** -> copy Bot User OAuth Token (`xoxb-...`).
5. Left menu **Event Subscriptions** -> enable -> Subscribe to bot events: `message.channels`, `message.groups`, `message.im`, `app_mention`.
6. Get Member ID: click user avatar -> profile -> **...** -> Copy member ID.
**Configuration:**
@@ -240,7 +444,7 @@ slack_proxy: ""
**Env vars:** `EVOSCIENTIST_SLACK_BOT_TOKEN`, `EVOSCIENTIST_SLACK_APP_TOKEN`, `EVOSCIENTIST_SLACK_ALLOWED_SENDERS`, `EVOSCIENTIST_SLACK_ALLOWED_CHANNELS`, `EVOSCIENTIST_SLACK_PROXY`
**Technical details:** Socket Mode (no public URL needed). Markdown→mrkdwn conversion. DMs respond directly; channels only respond to `app_mention` events. Thread replies via `thread_ts`. Attachments downloaded with Bearer auth. Media sent via `files_upload_v2`. Runs `auth_test()` on startup to verify credentials. Retry: 3 attempts, exponential backoff + jitter. Text chunk limit: 4000 chars.
**Technical details:** Socket Mode (no public URL needed). Markdown to mrkdwn conversion. DMs respond directly; channels respond to `app_mention` events. Thread replies via `thread_ts` -- all replies are threaded to the original message. Typing indicator approximated by posting/deleting a "..." message (Slack has no bot typing API). ACK reaction (eyes emoji) on message receipt. Attachments downloaded with Bearer auth. Media sent via `files_upload_v2`. Runs `auth_test()` on startup to verify credentials and cache bot user ID. Retry: 3 attempts, exponential backoff + jitter. Text chunk limit: 4000 chars.
---
@@ -250,11 +454,11 @@ slack_proxy: ""
**Prerequisites:**
1. Go to [Feishu Open Platform](https://open.feishu.cn/app) (international: [Lark Developer](https://open.larksuite.com/app)) → create a custom app.
1. Go to [Feishu Open Platform](https://open.feishu.cn/app) (international: [Lark Developer](https://open.larksuite.com/app)) -> create a custom app.
2. Copy the **App ID** and **App Secret**.
3. Left menu **Event Subscriptions** → set request URL to `http://your-host:9000/webhook/event` → copy **Verification Token** and **Encrypt Key**.
3. Left menu **Event Subscriptions** -> set request URL to `http://your-host:9000/webhook/event` -> copy **Verification Token** and **Encrypt Key**.
4. Add event: `im.message.receive_v1` (receive messages).
5. Left menu **Permissions** → enable `im:message:send_as_bot`.
5. Left menu **Permissions** -> enable `im:message:send_as_bot`.
6. Create a version and publish.
> Webhook must be publicly reachable. For local dev, use `ngrok http 9000`.
@@ -282,7 +486,7 @@ feishu_proxy: ""
**Env vars:** `EVOSCIENTIST_FEISHU_APP_ID`, `EVOSCIENTIST_FEISHU_APP_SECRET`, `EVOSCIENTIST_FEISHU_WEBHOOK_PORT`, `EVOSCIENTIST_FEISHU_DOMAIN`
**Technical details:** Webhook on `POST /webhook/event` with URL verification challenge-response. `tenant_access_token` auto-refresh (2h TTL, refreshes 5 min before expiry). Markdown→Post rich text conversion (code blocks, bold, italic, strikethrough, links, headings, quotes, lists). Plain text fallback. Group @mention filtering. Media: images via `/im/v1/images`, files via `/im/v1/files`. Replies via `/messages/{id}/reply`. Retry: 3 attempts, rate limit delay 2.0s, matches `99991400`/`rate limit`. Text chunk limit: 4096 chars.
**Technical details:** Webhook on `POST /webhook/event` with URL verification challenge-response. Supports both v1 (legacy) and v2 event schemas. Optional AES-256-CBC event decryption (when `encrypt_key` configured). `tenant_access_token` auto-refresh (2h TTL, refreshes 5 min before expiry). Markdown to Feishu Post rich text conversion (code blocks, bold, italic, strikethrough, links, headings, quotes, ordered/unordered lists). Plain text fallback. Group @mention filtering with mention key caching. Media: images via `/im/v1/images`, files via `/im/v1/files`. Replies via `/messages/{id}/reply` API. ACK reaction via `/messages/{id}/reactions`. Retry: 3 attempts, rate limit delay 2.0s, matches `99991400`/`rate limit`. Non-retryable: permission denied (`99991401`), invalid credentials. Text chunk limit: 4096 chars.
---
@@ -296,10 +500,10 @@ Two backends supported: **WeCom** (recommended, free, no certification needed) a
**Prerequisites:**
1. Log in to [WeCom Admin Console](https://work.weixin.qq.com) → App Management → create a custom app.
1. Log in to [WeCom Admin Console](https://work.weixin.qq.com) -> App Management -> create a custom app.
2. Copy the **Corp ID**, **AgentId**, and **Secret**.
3. In app details → Receive Messages → Set API Receive → URL: `http://your-host:9001/wechat/callback` → copy **Token** and **EncodingAESKey**.
4. In app details → **Trusted IP** → add your server's public IP address. Without this, all API calls will fail with error `60020`.
3. In app details -> Receive Messages -> Set API Receive -> URL: `http://your-host:9001/wechat/callback` -> copy **Token** and **EncodingAESKey**.
4. In app details -> **Trusted IP** -> add your server's public IP address. Without this, all API calls will fail with error `60020`.
```yaml
channel_enabled: "wechat"
@@ -328,9 +532,9 @@ wechat_proxy: ""
**Prerequisites:**
1. Log in to [WeChat Official Account Platform](https://mp.weixin.qq.com) → Settings & Development → Basic Configuration.
1. Log in to [WeChat Official Account Platform](https://mp.weixin.qq.com) -> Settings & Development -> Basic Configuration.
2. Copy the **AppID** and **AppSecret**.
3. Server Configuration → URL: `http://your-host:9001/wechat/callback` → set **Token** and **EncodingAESKey**.
3. Server Configuration -> URL: `http://your-host:9001/wechat/callback` -> set **Token** and **EncodingAESKey**.
```yaml
wechat_backend: "wechatmp"
@@ -347,7 +551,7 @@ wechat_mp_encoding_aes_key: "xxxxxxxxxxxxxxxxxx"
| `wechat_mp_token` | `str` | `""` | **Required (MP).** Server Token |
| `wechat_mp_encoding_aes_key` | `str` | `""` | **Required (MP).** Server EncodingAESKey |
**Technical details:** Webhook HTTP server. XML message parsing. Signature verification. `access_token` auto-refresh. Optional AES encryption/decryption. WeCom supports Markdown message format; Official Account uses plain text. Media send/receive. Retry + backoff. Text chunk limit: 2048 chars.
**Technical details:** Webhook HTTP server for inbound (XML message parsing). GET callback for URL verification (SHA1 signature check). POST callback for message handling -- returns `"success"` within 5s and processes asynchronously. Optional AES encryption/decryption via `WeChatCrypto`. `access_token` auto-refresh (2h TTL, 5-min margin). Token-expired errors (40014, 42001) trigger automatic retry with refreshed token. WeCom supports Markdown message format with plain text fallback; Official Account uses plain text only (customer service API). WeCom group messages sent via `/appchat/send` endpoint (group IDs start with `wr`). Typing indicator approximated by posting/recalling a "..." message (WeCom only). Supports text, image, voice, video, location, file, and link message types. Media upload via `/media/upload`. Text chunk limit: 4096 chars.
---
@@ -357,9 +561,9 @@ wechat_mp_encoding_aes_key: "xxxxxxxxxxxxxxxxxx"
**Prerequisites:**
1. Go to [DingTalk Open Platform](https://open-dev.dingtalk.com) → App Development → create a bot app.
1. Go to [DingTalk Open Platform](https://open-dev.dingtalk.com) -> App Development -> create a bot app.
2. Copy the **AppKey** (Client ID) and **AppSecret** (Client Secret).
3. Enable **Stream Mode** in the app configuration — no public IP needed.
3. Enable **Stream Mode** in the app configuration -- no public IP needed.
4. Publish the app and add the bot to a group, or test via direct message.
**Configuration:**
@@ -381,7 +585,7 @@ dingtalk_proxy: ""
**Env vars:** `EVOSCIENTIST_DINGTALK_CLIENT_ID`, `EVOSCIENTIST_DINGTALK_CLIENT_SECRET`
**Technical details:** Stream Mode (WebSocket, no public IP needed). Connects via DingTalk gateway with automatic ping/pong heartbeat and message ACK. `access_token` auto-refresh. Group @mention filtering (strips first `@bot` mention). Supports image, file, video, audio attachment download. Sends in Markdown format (`sampleMarkdown`). Auth errors (`invalidauthentication`/`forbidden`/`40014`) not retried. Text chunk limit: 4096 chars.
**Technical details:** Stream Mode via WebSocket -- connects to DingTalk gateway (`/v1.0/gateway/connections/open`) with automatic ticket-based auth. Ping/pong heartbeat with system topic handling. Message ACK via JSON response. `accessToken` auto-refresh. Sends via robot `oToMessages/batchSend` API in Markdown format (`sampleMarkdown`). Image uploads via `/media/upload` API with `sampleImageMsg`. Group @mention detection via `isInAtList` flag with `atUsers` array fallback. `downloadCode` resolution via `/robot/messageFiles/download` API for file/image/video/audio attachments. Auth errors (`invalidauthentication`/`forbidden`/`40014`) not retried. Text chunk limit: 4096 chars.
---
@@ -391,7 +595,7 @@ dingtalk_proxy: ""
**Prerequisites:**
1. Go to [QQ Open Platform](https://q.qq.com) → create a bot application.
1. Go to [QQ Open Platform](https://q.qq.com) -> create a bot application.
2. Complete developer verification, create a sandbox or production bot.
3. Copy the **AppID** and **AppSecret**.
4. Search for and add the bot as a friend in QQ, or add it to a group.
@@ -461,7 +665,7 @@ signal_rpc_port: 7583
**Prerequisites:**
1. Prepare an email account with IMAP + SMTP support (Gmail, Outlook, self-hosted, etc.).
2. **Gmail:** Enable 2FA → generate an App Password. IMAP: `imap.gmail.com:993` (SSL), SMTP: `smtp.gmail.com:587` (STARTTLS).
2. **Gmail:** Enable 2FA -> generate an App Password. IMAP: `imap.gmail.com:993` (SSL), SMTP: `smtp.gmail.com:587` (STARTTLS).
3. **Outlook/Office 365:** IMAP: `outlook.office365.com:993` (SSL), SMTP: `smtp.office365.com:587` (STARTTLS).
4. Ensure IMAP access is enabled in your email settings.
@@ -508,7 +712,7 @@ email_allowed_senders: ""
**Env vars:** `EVOSCIENTIST_EMAIL_IMAP_HOST`, `EVOSCIENTIST_EMAIL_IMAP_USERNAME`, `EVOSCIENTIST_EMAIL_IMAP_PASSWORD`, `EVOSCIENTIST_EMAIL_SMTP_HOST`, `EVOSCIENTIST_EMAIL_SMTP_USERNAME`, `EVOSCIENTIST_EMAIL_SMTP_PASSWORD`
**Technical details:** IMAP polling mode, checks for UNSEEN emails periodically (max 20 per cycle). Supports SSL and STARTTLS. Auto-parses multipart emails (prefers text/plain, falls back text/html → plain text). Attachments auto-downloaded. Replies set `In-Reply-To` and `References` headers to maintain email threads. Sends HTML + plain text dual format (multipart/alternative), falls back to plain text on HTML failure. IMAP auto-reconnects on disconnect. Auth errors (auth/login/credential) not retried. No public IP needed. Text chunk limit: no limit.
**Technical details:** IMAP polling mode, checks for UNSEEN emails periodically (max 20 per cycle). Supports SSL and STARTTLS. Auto-parses multipart emails (prefers text/plain, falls back text/html -> plain text). Attachments auto-downloaded. Replies set `In-Reply-To` and `References` headers to maintain email threads. Sends HTML + plain text dual format (multipart/alternative), falls back to plain text on HTML failure. IMAP auto-reconnects on disconnect. Auth errors (auth/login/credential) not retried. No public IP needed. Text chunk limit: no limit.
---
@@ -551,3 +755,144 @@ imessage_allowed_senders: ""
**Env vars:** `EVOSCIENTIST_IMESSAGE_CLI_PATH`, `EVOSCIENTIST_IMESSAGE_SERVICE`, `EVOSCIENTIST_IMESSAGE_ALLOWED_SENDERS`
**Technical details:** JSON-RPC over stdio with imsg CLI. Creates `watch.subscribe` on startup for real-time message streaming (not polling). Supports iMessage + SMS dual channel (`service: auto`). Target resolution supports chat_id, chat_guid, chat_identifier, and phone/email. Attachments read from local paths provided by imsg. Group detection via `is_group` field. RPC errors (AppleScript/permission/not found) not retried; only connection timeouts retried. Plain text format (no Markdown). No public IP needed. Text chunk limit: 4000 chars.
---
## Running Multiple Channels
Comma-separate channel names in the config to enable multiple channels simultaneously:
```yaml
channel_enabled: "telegram,discord,slack"
```
All enabled channels run concurrently via the internal `MessageBus`. Each channel:
- Has its own connection lifecycle (connect, reconnect, health check)
- Shares the same `InboundConsumer` worker pool and agent instance
- Routes outbound replies back to the originating channel automatically
### Multi-Channel Architecture
```
ChannelManager
├── TelegramChannel ──┐
├── DiscordChannel ──┤
├── SlackChannel ──┤──► MessageBus ──► InboundConsumer ──► Agent
├── FeishuChannel ──┤ │
└── ... ──┘ ▼
OutboundMessage
│
Dispatcher routes
to origin channel
```
### Health Monitoring
The `ChannelManager` runs a background health check task that monitors all active channels. Access health status via:
```bash
# CLI
EvoSci channel status
# HTTP (if health endpoint is enabled)
curl http://localhost:8080/healthz
```
## Docker Deployment
For webhook-based channels (Feishu, WeChat), Docker simplifies port mapping and process management:
```dockerfile
FROM python:3.11-slim
WORKDIR /app
COPY . .
RUN pip install evoscientist[feishu,wechat]
# Expose webhook ports
EXPOSE 9000 9001
CMD ["EvoSci", "serve"]
```
```bash
docker build -t evoscientist .
docker run -d \
-p 9000:9000 \
-p 9001:9001 \
-e EVOSCIENTIST_CHANNEL_ENABLED="feishu,wechat" \
-e EVOSCIENTIST_FEISHU_APP_ID="cli_xxx" \
-e EVOSCIENTIST_FEISHU_APP_SECRET="xxx" \
-e EVOSCIENTIST_WECHAT_BACKEND="wecom" \
-e EVOSCIENTIST_WECHAT_WECOM_CORP_ID="ww..." \
-e EVOSCIENTIST_WECHAT_WECOM_SECRET="xxx" \
evoscientist
```
For polling/WebSocket channels (Telegram, Discord, Slack, DingTalk, QQ), no port mapping is needed:
```bash
docker run -d \
-e EVOSCIENTIST_CHANNEL_ENABLED="telegram" \
-e EVOSCIENTIST_TELEGRAM_BOT_TOKEN="123456:ABC-xxx" \
evoscientist
```
## Troubleshooting
### Common Issues
**Bot not responding to messages**
1. Check that `channel_enabled` includes your channel name
2. Verify the bot token/credentials are correct: `EvoSci config get telegram_bot_token`
3. If using `allowed_senders`, ensure your user ID is listed
4. For group chats, ensure the bot is @mentioned (default behavior)
5. Check logs for middleware drops: `DedupMiddleware`, `AllowListMiddleware`, or `MentionGatingMiddleware`
**"channel X not found" or import errors**
Install the channel-specific dependencies:
```bash
pip install evoscientist[telegram] # or discord, slack, feishu, etc.
```
**Webhook channels (Feishu, WeChat) not receiving messages**
1. Ensure the webhook URL is publicly reachable (not behind NAT without port forwarding)
2. For local development, use a tunnel: `ngrok http 9000`
3. Verify the callback URL matches exactly (including path: `/webhook/event` for Feishu, `/wechat/callback` for WeChat)
4. Check that signature verification tokens match between the platform config and your local config
**WeChat API error `60020`**
Add your server's public IP to the WeCom app's **Trusted IP** list in the admin console.
**Token refresh failures**
- Feishu/WeChat/DingTalk tokens auto-refresh with a 5-minute safety margin before expiry
- If the refresh endpoint is unreachable (network issues), messages will fail until the next successful refresh
- Check proxy settings if your server requires a proxy to reach external APIs
**Duplicate responses**
- The `DedupMiddleware` prevents most duplicates using a 60-second LRU cache
- If you see duplicates, check if the platform is sending the same message with different IDs (some platforms retry delivery)
**Messages truncated**
- Each platform has a max text length (see Capability Matrix)
- Long responses are automatically chunked at paragraph/code-block boundaries
- Adjust `text_chunk_limit` in the channel config if needed
### Debug Logging
Enable debug logs for the channel subsystem:
```bash
export EVOSCIENTIST_LOG_LEVEL=DEBUG
# or
EvoSci config set log_level debug
```
Channel-specific log output is prefixed with the module path (e.g., `EvoScientist.channels.telegram.channel`).
+35 -20
View File
@@ -428,12 +428,18 @@ class Channel(ChannelPlugin, ABC):
self._send_locks.move_to_end(chat_id)
else:
self._send_locks[chat_id] = asyncio.Lock()
# Evict oldest unlocked entries when over capacity
while len(self._send_locks) > self._send_locks_max:
oldest_key, oldest_lock = next(iter(self._send_locks.items()))
if oldest_lock.locked():
break # don't evict a lock that's in use
del self._send_locks[oldest_key]
# Evict oldest unlocked entries when over capacity.
# Skip locked entries instead of giving up entirely,
# to prevent unbounded growth.
if len(self._send_locks) > self._send_locks_max:
to_evict = [
k for k, lock in self._send_locks.items()
if not lock.locked() and k != chat_id
]
for k in to_evict:
if len(self._send_locks) <= self._send_locks_max:
break
del self._send_locks[k]
return self._send_locks[chat_id]
async def send(self, message: OutboundMessage) -> bool:
@@ -816,26 +822,30 @@ class Channel(ChannelPlugin, ABC):
def _build_inbound(self, raw: RawIncoming) -> InboundMessage | None:
"""Run *raw* through inbound middlewares and convert to InboundMessage.
Synchronous wrapper around :meth:`_build_inbound_async`. Safe to
call from both sync and async contexts.
Synchronous wrapper around :meth:`_build_inbound_async`. When an
event loop is already running, the coroutine is scheduled on that
loop via :func:`asyncio.run_coroutine_threadsafe` to avoid
thread-safety issues with middleware state (DedupCache,
GroupHistoryBuffer, etc.).
"""
import asyncio
import concurrent.futures
try:
asyncio.get_running_loop()
# Inside a running loop — run in a worker thread
with concurrent.futures.ThreadPoolExecutor(max_workers=1) as pool:
return pool.submit(
lambda: asyncio.run(self._build_inbound_async(raw))
).result()
loop = asyncio.get_running_loop()
except RuntimeError:
# No running loop — safe to create one
loop = asyncio.new_event_loop()
loop = None
if loop is not None and loop.is_running():
future = asyncio.run_coroutine_threadsafe(
self._build_inbound_async(raw), loop,
)
return future.result()
else:
new_loop = asyncio.new_event_loop()
try:
return loop.run_until_complete(self._build_inbound_async(raw))
return new_loop.run_until_complete(self._build_inbound_async(raw))
finally:
loop.close()
new_loop.close()
def _raw_to_inbound(self, raw: RawIncoming) -> InboundMessage | None:
"""Convert a RawIncoming to InboundMessage (pure transformation, no filtering).
@@ -917,7 +927,12 @@ class Channel(ChannelPlugin, ABC):
async def debounce_callback(_s=sender, _w=wait):
await asyncio.sleep(_w)
await self._process_buffered_messages(_s)
try:
await self._process_buffered_messages(_s)
except Exception as e:
_logger.error(
f"{self.name} debounce flush error for {_s}: {e}"
)
self._debounce_tasks[sender] = asyncio.create_task(
debounce_callback()
+18 -11
View File
@@ -12,6 +12,7 @@ from __future__ import annotations
import asyncio
import logging
import uuid
from collections import OrderedDict
from dataclasses import dataclass
from typing import Any, AsyncIterator, Callable, TypeVar
@@ -131,7 +132,7 @@ class InboundConsumer:
self._on_message_received = on_message_received
self._on_streaming_event = on_streaming_event
self._on_message_sent = on_message_sent
self._sessions: dict[str, str] = {} # sender_id -> thread_id
self._sessions: OrderedDict[str, str] = OrderedDict() # sender_id -> thread_id (LRU)
# Per-chat locks: same chat is processed serially (bounded)
self._chat_locks: dict[str, asyncio.Lock] = {}
@@ -152,16 +153,22 @@ class InboundConsumer:
self._metrics = ConsumerMetrics()
def _get_thread_id(self, sender_id: str) -> str:
"""Get or create a thread ID for the given sender."""
if sender_id not in self._sessions:
if len(self._sessions) >= _MAX_SESSIONS:
# Evict oldest entry
oldest = next(iter(self._sessions))
del self._sessions[oldest]
if self.thread_id:
self._sessions[sender_id] = f"{self.thread_id}:{sender_id}"
else:
self._sessions[sender_id] = str(uuid.uuid4())
"""Get or create a thread ID for the given sender.
Uses LRU ordering: recently accessed senders are moved to the
end, so eviction always removes the least-recently-active sender.
"""
if sender_id in self._sessions:
self._sessions.move_to_end(sender_id)
return self._sessions[sender_id]
if len(self._sessions) >= _MAX_SESSIONS:
# Evict the least-recently-used entry
self._sessions.popitem(last=False)
if self.thread_id:
self._sessions[sender_id] = f"{self.thread_id}:{sender_id}"
else:
self._sessions[sender_id] = str(uuid.uuid4())
return self._sessions[sender_id]
def _get_channel(self, channel_name: str) -> Channel | None:
+34 -16
View File
@@ -26,6 +26,22 @@ from .base import RawIncoming
_logger = logging.getLogger(__name__)
# ── Task cancellation helper ─────────────────────────────────────────
async def _cancel_task(task: asyncio.Task) -> None:
"""Cancel an asyncio task and await its completion.
Suppresses ``CancelledError`` so callers don't need to handle it.
"""
task.cancel()
try:
await task
except asyncio.CancelledError:
pass
except Exception:
pass # Already logged elsewhere; prevent unhandled propagation
# ═══════════════════════════════════════════════════════════════════════
# Supporting data structures
# ═══════════════════════════════════════════════════════════════════════
@@ -135,7 +151,7 @@ class GroupHistoryBuffer:
buf = self._buffers.get(chat_id)
if not buf:
return []
now = time.time()
now = time.monotonic()
recent = [e for e in buf if now - e.timestamp < self._max_age]
return recent[-limit:]
@@ -193,7 +209,7 @@ class TypingManager:
"""Cancel the typing-indicator loop for *chat_id*."""
task = self._tasks.pop(chat_id, None)
if task:
task.cancel()
await _cancel_task(task)
async def stop_all(self) -> None:
"""Cancel all active typing-indicator loops."""
@@ -236,7 +252,7 @@ class PairingManager:
# Check if already has pending request
for code, req in list(self._pending.items()):
if req.sender_id == sender_id and req.channel == channel:
if time.time() - req.created_at < self.CODE_EXPIRY:
if time.monotonic() - req.created_at < self.CODE_EXPIRY:
return code # return existing code
else:
del self._pending[code]
@@ -254,7 +270,7 @@ class PairingManager:
sender_id=sender_id,
channel=channel,
code=code,
created_at=time.time(),
created_at=time.monotonic(),
)
_logger.info(f"Pairing code {code} generated for {channel}:{sender_id}")
return code
@@ -264,7 +280,7 @@ class PairingManager:
req = self._pending.get(code)
if not req:
return False, f"Unknown code: {code}"
if time.time() - req.created_at > self.CODE_EXPIRY:
if time.monotonic() - req.created_at > self.CODE_EXPIRY:
del self._pending[code]
return False, f"Code {code} expired"
@@ -287,7 +303,7 @@ class PairingManager:
return list(self._pending.values())
def _cleanup_expired(self):
now = time.time()
now = time.monotonic()
expired = [c for c, r in self._pending.items() if now - r.created_at > self.CODE_EXPIRY]
for c in expired:
del self._pending[c]
@@ -440,11 +456,12 @@ class DebounceMiddleware:
if self.on_ready:
await self.on_ready(inbound)
def cancel_all(self) -> None:
"""Cancel all pending debounce tasks."""
for task in self._tasks.values():
task.cancel()
async def cancel_all(self) -> None:
"""Cancel all pending debounce tasks and await their completion."""
tasks = list(self._tasks.values())
self._tasks.clear()
for task in tasks:
await _cancel_task(task)
# ── Chunking ─────────────────────────────────────────────────────────
@@ -743,11 +760,8 @@ class GroupHistoryMiddleware(InboundMiddleware):
if not raw.is_group:
return raw
ts = (
raw.timestamp.timestamp()
if hasattr(raw.timestamp, "timestamp")
else time.time()
)
# Use monotonic clock for consistent expiry calculation
ts = time.monotonic()
if not raw.was_mentioned:
self._buffer.add(
@@ -792,6 +806,7 @@ class PairingMiddleware(InboundMiddleware):
self._channel_name = channel_name
self._send_response_fn = send_response_fn
self._dm_policy = dm_policy
self._background_tasks: set[asyncio.Task] = set()
async def process_inbound(
self, raw: RawIncoming, context: dict[str, Any],
@@ -809,6 +824,9 @@ class PairingMiddleware(InboundMiddleware):
code = self._manager.request_pairing(self._channel_name, raw.sender_id)
if self._send_response_fn:
text = f"\U0001f510 Pairing required. Your code: {code}\nThis code expires in 1 hour."
asyncio.ensure_future(self._send_response_fn(raw.chat_id, text))
task = asyncio.create_task(self._send_response_fn(raw.chat_id, text))
# Track the task to prevent GC and handle exceptions
self._background_tasks.add(task)
task.add_done_callback(self._background_tasks.discard)
_logger.info(f"Pairing required for {raw.sender_id}, code sent")
return None
+19 -5
View File
@@ -229,7 +229,7 @@ class WebSocketMixin:
except Exception as e:
logger.error(f"{getattr(self, 'name', '?')} WS error: {e}")
self._ws_cleanup_heartbeat()
await self._ws_cleanup_heartbeat()
self._ws_session = None
if getattr(self, "_running", False):
@@ -244,10 +244,17 @@ class WebSocketMixin:
break
await asyncio.sleep(self._ws_heartbeat_interval)
def _ws_cleanup_heartbeat(self) -> None:
async def _ws_cleanup_heartbeat(self) -> None:
if self._ws_heartbeat_task:
self._ws_heartbeat_task.cancel()
task = self._ws_heartbeat_task
self._ws_heartbeat_task = None
task.cancel()
try:
await task
except asyncio.CancelledError:
pass
except Exception:
pass
async def _ws_send_json(self, data: dict) -> None:
"""Send JSON to the active WebSocket."""
@@ -255,7 +262,7 @@ class WebSocketMixin:
await self._ws_session.send_str(json.dumps(data))
async def _stop_ws(self) -> None:
self._ws_cleanup_heartbeat()
await self._ws_cleanup_heartbeat()
if self._ws_session:
await self._ws_session.close()
self._ws_session = None
@@ -302,5 +309,12 @@ class PollingMixin:
async def _stop_polling(self) -> None:
if self._poll_task:
self._poll_task.cancel()
task = self._poll_task
self._poll_task = None
task.cancel()
try:
await task
except asyncio.CancelledError:
pass
except Exception:
pass
+38 -12
View File
@@ -33,6 +33,7 @@
- [CLI Inference](#cli-inference)
- [Script Inference](#script-inference)
- [Web Interface](#web-interface)
- [💬 Channels](#-channels)
- [🔌 MCP Integration](#-mcp-integration)
- [📊 Evaluation](#-evaluation)
- [📝 Citation](#-citation)
@@ -221,24 +222,49 @@ EvoSci config path # Show config file path
/install-skill owner/repo@skill-name
```
### iMessage Channel
### Channels
EvoScientist can be controlled remotely via iMessage. The channel shares the same agent and conversation thread as the CLI — messages from your phone go through the same session.
EvoScientist integrates with 10 messaging platforms, allowing you to control the agent remotely from any chat app. All channels share the same agent core — messages from any platform go through the same processing pipeline.
```Shell
# Enable during onboard
EvoSci onboard # Step 7 configures iMessage
| Channel | Transport | Public IP Required | Install Extra |
|:--------|:----------|:------------------:|:--------------|
| Telegram | Long Polling | No | `pip install evoscientist[telegram]` |
| Discord | WebSocket | No | `pip install evoscientist[discord]` |
| Slack | Socket Mode | No | `pip install evoscientist[slack]` |
| Feishu / Lark | HTTP Webhook | Yes | `pip install evoscientist[feishu]` |
| WeChat (WeCom / MP) | HTTP Webhook | Yes | `pip install evoscientist[wechat]` |
| DingTalk | WebSocket Stream | No | `pip install evoscientist[dingtalk]` |
| QQ | WebSocket | No | `pip install evoscientist[qq]` |
| Signal | JSON-RPC | No | `pip install evoscientist[signal]` |
| Email | IMAP + SMTP | No | `pip install evoscientist[email]` |
| iMessage | JSON-RPC (stdio) | No | macOS only, requires [imsg](https://github.com/anthropics/imsg) CLI |
# Or enable manually
EvoSci config set imessage_enabled true
EvoSci config set imessage_allowed_senders "+1234567890"
**Quick start:**
```bash
# 1. Install channel dependencies
pip install evoscientist[telegram] # or discord, slack, feishu, etc.
# 2. Configure via wizard or CLI
EvoSci onboard # interactive setup
# or
EvoSci config set channel_enabled telegram
EvoSci config set telegram_bot_token "123456:ABC-xxx"
# 3. Start
EvoSci serve # agent + all enabled channels
```
The channel can also be started manually with `/channel` in the interactive CLI.
Multiple channels can run concurrently — comma-separate names in the config:
**Features:**
- Thinking content and todo lists are forwarded to iMessage as intermediate messages
- Media files (images, PDFs) are auto-sent when the agent writes (`write_file`) or reads (`read_file`) them — no extra commands needed
```yaml
channel_enabled: "telegram,discord,slack"
```
The channel can also be started interactively with `/channel` in the CLI session.
> [!NOTE]
> For per-channel setup guides, capability matrix, architecture details, and troubleshooting, see the **[Channel Integration Guide](./EvoScientist/channels)**.
### Runtime Directories
+12 -13
View File
@@ -1251,8 +1251,8 @@ class TestInboundConsumer:
assert tid1 == "shared_thread:alice"
assert tid2 == "shared_thread:bob"
def test_session_eviction_is_fifo_not_lru(self):
"""[B-19] Sessions evict oldest by insertion, not by access."""
def test_session_eviction_is_lru(self):
"""Sessions use LRU eviction: recently accessed senders are kept."""
consumer = self._make_consumer()
consumer._sessions.clear()
@@ -1260,13 +1260,12 @@ class TestInboundConsumer:
for i in range(10):
consumer._sessions[f"user_{i}"] = f"thread_{i}"
# Access "user_0" (should make it LRU-recent, but dict doesn't)
_ = consumer._sessions["user_0"]
# Access "user_0" via _get_thread_id (triggers LRU move_to_end)
consumer._get_thread_id("user_0")
# Force eviction by exceeding limit (simulate)
# Note: actual limit is 10_000, we test the logic pattern
# "user_0" should now be at the end (most recently used)
oldest = next(iter(consumer._sessions))
assert oldest == "user_0" # Still first in insertion order
assert oldest == "user_1" # user_1 is now the least recently used
def test_metrics_initial(self):
consumer = self._make_consumer()
@@ -1537,20 +1536,20 @@ class TestIntegration:
assert flushed.content == "will be lost"
_run(_test())
def test_send_locks_unbounded_growth(self):
"""[B-03] _send_locks grows without bound for unique chat_ids."""
def test_send_locks_bounded_growth(self):
"""_send_locks stays bounded via LRU eviction of unlocked entries."""
async def _test():
ch = StubChannel()
for i in range(100):
ch._send_locks_max = 10 # Small limit for testing
for i in range(20):
msg = OutboundMessage(
channel="stub", chat_id=f"chat_{i}",
content="hi", metadata={"chat_id": f"chat_{i}"},
)
await ch.send(msg)
# All 100 unique chat_ids created a lock
assert len(ch._send_locks) == 100
# BUG: These are never cleaned up
# Should be bounded at max + 1 (the newly inserted entry)
assert len(ch._send_locks) <= ch._send_locks_max + 1
_run(_test())