diff --git a/backend/app/channels/AGENTS.md b/backend/app/channels/AGENTS.md index 9f79bae71..357e9092c 100644 --- a/backend/app/channels/AGENTS.md +++ b/backend/app/channels/AGENTS.md @@ -12,7 +12,7 @@ Bridges external messaging platforms (Feishu, Slack, Telegram, Discord, DingTalk **What may be published from the stream is an allowlist, not a denylist** (`_accumulate_stream_text`): only assistant message types — LangChain serializes `AIMessage.type` as `"ai"` and `AIMessageChunk.type` as `"AIMessageChunk"`, plus the OpenAI-style `"assistant"` spelling for foreign runtimes — become displayable text. The previous rule rejected only payloads whose `type` contained `"tool"` and therefore published everything else, which leaked DeerFlow's hidden model context to every streaming IM channel: `DynamicContextMiddleware` injects recalled memory as a hidden `HumanMessage` (`type == "human"`) and rewrites the user's own turn into a new `HumanMessage`, `DurableContextMiddleware` injects a hidden `` `HumanMessage`, and LangGraph fans those state writes out on the `messages-tuple` stream. Proved live on a Buzz relay, which published a `` fact block and, in another run, a verbatim echo of the user's own message as the assistant's reply. Matching is by prefix (`ai` / `assistant`), never substring, because ordinary words contain `"ai"` (`chain`, `domain`). The message type is resolved by `_stream_payload_type`, which handles both the `model_dump()` shape DeerFlow's own gateway emits and LangChain's `to_json()` constructor shape (whose top-level `type` is the literal `"constructor"`, with the class name at the tail of the `id` path). A bare `str` payload is no longer accepted at all: it carries no type information, so it cannot be attributed to the assistant, and nothing in DeerFlow produces one (`runtime/serialization.py::serialize_messages_tuple` always emits `[message_dict, metadata]`). - `base.py` - Abstract `Channel` base class (start/stop/send lifecycle). Provider callbacks that submit coroutines from SDK threads must use `_submit_threadsafe_coroutine()`: it creates and retains the real `asyncio.Task` on the owner loop instead of treating `run_coroutine_threadsafe()`'s proxy Future as a completion signal. Submission is closed atomically with shutdown, and `stop()` must call `_close_and_drain_threadsafe_futures()` before tearing down SDK resources. - `service.py` - Manages lifecycle of all configured channels from `config.yaml`. Shutdown closes manager admission first and keeps transports alive until every manager worker/follow-up watcher has exited. A successful manager stop therefore owns no live handler; if the Gateway's outer timeout cancels shutdown, the service retains its channel objects and global singleton so unfinished resources are not detached and cleanup can be retried. -- `slack.py` / `feishu.py` / `telegram.py` / `discord.py` / `dingtalk.py` - Platform-specific implementations (`feishu.py` tracks the running card `message_id` in memory and patches the same card in place; `telegram.py` accepts inbound text/photos/documents, preserves media captions, hands token-free attachment bytes to the shared upload pipeline, edits the "Working on it..." stream target in place via `editMessageText`, and can optionally send final Markdown replies as Rich Messages through `channels.telegram.rich_messages`; `discord.py` registers typing loops before inbound handling yields and `_start_typing()` refuses work once `_running` is false; because `stop()` runs on the main loop while typing tasks belong to `_discord_loop`, normal cross-thread cancellation, awaiting, and map cleanup are scheduled there with a bounded wait, while `_run_client()` drains tasks in its `finally` block before an exception/disconnect can make that loop unusable; an already-stopped foreign loop must never have its tasks awaited from the main loop; to that end every outbound cross-loop call (`send` / `send_file` / `_get_channel_or_thread`) goes through `_run_on_discord_loop`, which bounds the await (`DISCORD_OUTBOUND_TIMEOUT_SECONDS` 30s; `DISCORD_UPLOAD_TIMEOUT_SECONDS` 120s for file uploads, whose unbounded-size payload needs room for a slow uplink plus 429 retry-after) and fails fast with a `RuntimeError` when the client loop is missing or not running — a dead client becomes a logged send failure instead of a permanently hung `ChannelManager` worker — and `is_running` reports client-thread aliveness (like `feishu.py`) so `ensure_channel_ready` can restart the channel after `_run_client()` exits on a fatal error; `dingtalk.py` optionally uses AI Card streaming for in-place updates when `card_template_id` is configured, and overrides `receive_file` to download inbound images (`picture`/`richText`) and documents (`file`) by `downloadCode` into the thread uploads bucket, mirroring `feishu.py`) +- `slack.py` / `feishu.py` / `telegram.py` / `discord.py` / `dingtalk.py` - Platform-specific implementations (`feishu.py` tracks the running card `message_id` in memory and patches the same card in place; `telegram.py` accepts inbound text/photos/documents, preserves media captions, hands token-free attachment bytes to the shared upload pipeline, edits the "Working on it..." stream target in place via `editMessageText`, and can optionally send final Markdown replies as Rich Messages through `channels.telegram.rich_messages` — the Rich path only fires when `rich_messages` is on **and** the text contains a rich construct (fenced code, emphasis, a table, a task list, `
`, block math, or a link), so structured command/error replies (none of these) stay plain text and their newlines and `` tokens survive; `discord.py` registers typing loops before inbound handling yields and `_start_typing()` refuses work once `_running` is false; because `stop()` runs on the main loop while typing tasks belong to `_discord_loop`, normal cross-thread cancellation, awaiting, and map cleanup are scheduled there with a bounded wait, while `_run_client()` drains tasks in its `finally` block before an exception/disconnect can make that loop unusable; an already-stopped foreign loop must never have its tasks awaited from the main loop; to that end every outbound cross-loop call (`send` / `send_file` / `_get_channel_or_thread`) goes through `_run_on_discord_loop`, which bounds the await (`DISCORD_OUTBOUND_TIMEOUT_SECONDS` 30s; `DISCORD_UPLOAD_TIMEOUT_SECONDS` 120s for file uploads, whose unbounded-size payload needs room for a slow uplink plus 429 retry-after) and fails fast with a `RuntimeError` when the client loop is missing or not running — a dead client becomes a logged send failure instead of a permanently hung `ChannelManager` worker — and `is_running` reports client-thread aliveness (like `feishu.py`) so `ensure_channel_ready` can restart the channel after `_run_client()` exits on a fatal error; `dingtalk.py` optionally uses AI Card streaming for in-place updates when `card_template_id` is configured, and overrides `receive_file` to download inbound images (`picture`/`richText`) and documents (`file`) by `downloadCode` into the thread uploads bucket, mirroring `feishu.py`) - `buzz.py` - Buzz (Nostr relay) implementation: one NIP-42-authenticated websocket, pubkey-allowlist + mention/DM/thread-follow gating, streaming replies via in-place kind-40003 edits; requires the `buzz` dependency extra. Its durable seen-event replay guard uses async `aseen()` / `arecord()` / `aflush()` boundaries: initial JSON reads and coalesced atomic writes run off the Gateway event loop, and `stop()` awaits the final flush before returning. **Subscription model** (operator-facing version in [IM_CHANNEL_CONNECTIONS.md](../../docs/IM_CHANNEL_CONNECTIONS.md#buzz-subscription-model)): the relay fans kind-9 chat events out **only** to channel-scoped subscriptions, proved against a live relay — `REQ {"kinds":[9]}` is accepted and answered with `EOSE` but never receives an event (the connector authenticated and then sat silent forever), `REQ {"kinds":[9],"#h":[uuid]}` works, and a multi-value `#h` matches nothing, so it is strictly one REQ per channel (the same shape as Buzz's own `buzz-acp` harness). Every connection therefore rebuilds three kinds of subscription after NIP-42 auth: `buzz-discovery` (`{"kinds":[39000]}`, a historical query returning exactly the channels this identity belongs to, one stored event each then `EOSE` — adding `#p` returns zero, do not "narrow" it), `buzz-membership` (`{"kinds":[44100,44101],"#p":[us],"since":}`, the relay-signed member-added/member-removed notifications whose `p` tag names the affected member and `h` tag the channel), and one `buzz-chat-` per discovered channel. Chat subscriptions open as each kind-39000 arrives; the discovery `EOSE` is the completeness barrier that retries any that failed and warns when discovery found nothing. A kind-44100 for our pubkey subscribes to the new channel **live** (then re-issues discovery so its name/type reach the DM-detection cache, but only when that channel's metadata is actually missing — an unconditional refresh is how a burst of 44100s multiplied into one discovery pass each); a kind-44101 issues `buzz_nostr.close_frame` for exactly that channel's subscription and drops its metadata. **The membership subscription is scoped to LIVE events** and this is load-bearing: buzz-relay *stores* 44100/44101 and serves history newest-first (default limit 2000), so an unscoped filter replayed the whole membership history on every connect — every stored add read as live, re-running discovery once each (M+1 discovery passes × N stored kind-39000 events per connect, observed live as two `channel discovery complete` lines and one channel logged `` because the 44100 path subscribed before its metadata arrived), re-subscribing channels we have since been removed from, and letting a stored 44101 transiently unsubscribe a channel we are still in. `since` is anchored at the moment the socket opened (`_session_started_at`) minus `MEMBERSHIP_LOOKBACK_SECONDS` (60) of slack, which covers both relay clock skew and a membership change published *during* the connect/auth handshake; the slack can only cost an idempotent replay of the last minute. **A relay `CLOSED` frame is recovered, not merely forgotten** — every subscription on the socket fails silently when dropped, so `_handle_closed` re-issues it, bounded by `MAX_RESUBSCRIBE_ATTEMPTS` (3) per subscription id per connection *and per auth epoch* (`_resubscribe_attempts` is reset on session start, on re-auth, and by `stop()`, so pre-auth rejections — which the auth branch already recovers wholesale — never spend the authenticated session's budget). The first retry is immediate (the common case is a one-off hiccup); later ones back off 1s then 2s, awaited inline in the read loop rather than in a background task that could outlive its own socket. **`auth-required:` before this socket has completed its NIP-42 handshake is the expected bootstrap sequence, not a refusal**, and `_handle_closed` short-circuits it ahead of every recovery path: the connector opens its control REQs immediately in case the relay serves unauthenticated reads, a closed relay answers `auth-required:` plus an `AUTH` challenge, and the auth branch re-opens everything. That case is logged at DEBUG and consumes neither the permanent-refusal branch nor the retry budget (a chat subscription is still dropped from `_chat_subscriptions`, since it genuinely is not subscribed and discovery is what re-opens it). Treating it as permanent produced an operator-facing warning claiming discovery/membership tracking was DOWN in the same run where discovery then completed, every channel was subscribed, and a brand-new channel's kind-44100 was picked up one second later. The boundary is the per-socket `_auth_completed` flag, set once the signed AUTH event has been sent and cleared on session entry, session exit, and `stop()`; the same reason *after* that stays loud, and says the subscription is down until the relay's next AUTH challenge or the next reconnect rather than borrowing the non-auth wording. Then `_is_transient_close` decides whether to retry at all: NIP-01/NIP-42 `auth-required:`/`restricted:`/`blocked:`/`mute:`/`invalid:`/`pow:` prefixes and buzz-relay's own removal/revocation prose are permanent (do not fight the relay over a channel that is no longer ours; a post-auth `auth-required:` is recovered by the AUTH branch, not by re-issuing), `rate-limited:`/`error:`/no-reason-at-all and anything unrecognized are transient — the default resolves toward "keep listening" because going silently deaf is the failure this exists to remove, and the attempt budget bounds a wrong guess. A chat `CLOSED` is only ever recovered for a channel already in `_chat_subscriptions`: a `CLOSED` is relay-supplied, so acting on an unknown one would let a relay induce a subscription just by naming a channel. Every subscription that goes unlistened is logged at WARNING, never INFO. Subscription ids are deterministic per channel precisely so one can be replaced or closed without disturbing the others on the socket, and `_chat_subscriptions` is per-socket state cleared on session end, on re-auth (a pre-auth REQ may have been rejected), and by `stop()`. Three bounds on remote-fed state: `MAX_CACHED_CHANNELS` (512) caps the kind-39000 metadata cache, the watermark map, and the resubscribe-attempt map; `MAX_CHANNEL_SUBSCRIPTIONS` (256, well under buzz-relay's own 1024-per-connection ceiling) caps live chat subscriptions — at the cap new channels are refused and named in a warning rather than evicting a working subscription. **Known bound (documented, not fixed):** the relay caps historical delivery at 2000 events per subscription, newest-first, even with a `since`, so >2000 unread messages in a *single* channel across a disconnect loses the oldest — the relay never sends them and the watermark advances past them. That is the one remaining path that can skip; everything else fails toward replay. **Trust model** (operator-facing version in [IM_CHANNEL_CONNECTIONS.md](../../docs/IM_CHANNEL_CONNECTIONS.md#buzz-trust-model)): every inbound `EVENT` is authenticated at the single `handle_relay_frame` choke point — the NIP-01 id is recomputed from the delivered payload and the BIP-340 Schnorr signature verified against the claimed `pubkey` (`buzz_nostr.verify_event`, pure and total: malformed input returns `False`, never raises) — so `ev["pubkey"]`, the authorization principal for both the allowlist and the `/connect` bind, cannot be forged by a relay the DeerFlow operator does not run. What remains trusted is the *authorship* of kind-39000 channel metadata: any member can sign one, and because per-channel subscriptions are now driven by discovery, a forged kind-39000 has two effects rather than one — it can mark a channel `type: "dm"` (relaxing `require_mention` for that channel) **and** it can induce a chat subscription for a channel of the forger's choosing, since the channels we listen to are exactly the channels we hold metadata for. Neither makes anything be *acted on*: `allowed_users` and per-event signature verification are independent gates, so an induced subscription only means the relay reads its own traffic back to a subscriber that drops it, bounded by `MAX_CHANNEL_SUBSCRIPTIONS` (which refuses rather than evicts, so it cannot displace a real channel). Same for a forged kind-44100, except its `p` tag is re-checked locally so it must at least name us. Closing this needs a configured trusted relay pubkey, which `relay_url` is not. `allowed_users` is deny-by-default (empty = nobody, unlike siblings' empty = everyone), so `start()` logs a WARNING when it is empty and each drop logs at DEBUG. The resubscribe cursor (`since`) is **per channel**, advances only for events that were actually processed, and never past `now + MAX_FUTURE_SKEW_SECONDS`, because it is peer-supplied (`created_at`) and a single future-dated event otherwise made the connector permanently deaf. Per channel rather than global is the safety-critical half: subscriptions are per channel, so one shared cursor is the newest event seen in *any* channel and a busy channel would drag it past a quiet channel's unread messages, skipping them on the next reconnect — measured on a live relay, three channels of one identity sat ~28h apart. Per-channel cursors can only ever cost duplicate delivery (absorbed by the manager's `event_id` dedupe), and an evicted cursor degrades to "no `since`", i.e. the relay's default backlog — both fail toward replay, never toward a miss. Streaming tracks every oversize chunk index (`_stream_targets` for chunk 0, `_stream_tails` for the rest), since the manager republishes cumulative text and reposting `chunks[1:]` per update flooded the channel; all of it is per-connection state cleared by `stop()`, and the remote-fed kind-39000 cache is capped at `MAX_CACHED_CHANNELS`. **`send()` refuses outright to publish text carrying a hidden model-context wrapper** (``, ``, `` — `_HIDDEN_CONTEXT_MARKERS`), logging at ERROR and clearing the stream bookkeeping on a blocked `is_final`. This is defense in depth behind the manager's allowlist, and it lives here rather than in a sibling connector because on Buzz a leak is permanent: every streaming update is an immutable public Nostr event, so a corrective edit only changes what clients render while the original leaked event stays on the relay. Matching is on the literal opening tag, so a reply that merely talks about memory is still published Buzz seen-event shutdown stays bounded and retryable: `aflush()` awaits any in-flight write, attempts at most one final snapshot, leaves a still-changing generation dirty for fail-open replay, and returns without a live persistence timer. A Gateway cancellation does not cancel the underlying worker-thread write; `BuzzChannel.stop()` tracks cleanup completion separately from transport admission so ChannelService can retry the retained channel without racing a second snapshot against that write. The store is quiesced before stop (including the already-stopped guard), so a timed-out relay task that records after stop only marks data dirty and cannot schedule detached file work; a repeated `stop()` still drains that dirty state, while `BuzzChannel.start()` explicitly resumes scheduling and flushes it automatically. - `github.py` - Webhook-driven GitHub channel. Inbound messages come from `POST /api/webhooks/github`; outbound is log-only because GitHub agents post explicitly with `gh` from their sandbox when they choose to comment or create a PR diff --git a/backend/app/channels/telegram.py b/backend/app/channels/telegram.py index 08972b663..0d323ecfa 100644 --- a/backend/app/channels/telegram.py +++ b/backend/app/channels/telegram.py @@ -4,6 +4,7 @@ from __future__ import annotations import asyncio import logging +import re import threading import time from collections.abc import Coroutine @@ -42,6 +43,35 @@ MAX_TRACKED_STREAM_MESSAGES = 256 # Indirection so tests can patch the clock without touching the global time module. _monotonic = time.monotonic +# Rich Messages only earn their keep when the text carries a construct that +# renders natively: a fenced code block, emphasis, a table, a task list, +#
, block math, or a link. Structured command/error replies are plain +# text with none of these, so they stay plain and their newlines and +# tokens survive verbatim instead of being collapsed by the +# rich parser. +# +# The patterns are deliberately *well-formed*, not "contains this character": +# a table must be a line that leads with a pipe, a link must be [text](url), +# so a command line like "/goal [condition|clear]" (brackets + a mid-line pipe) +# never trips the detector. +_TELEGRAM_RICH_CONSTRUCT_RE = re.compile( + r"(?m)" + r"^\s*(`{3,}|~{3,})" # fenced code block + r"|^\s*[-*+]\s+\[[ xX]\]" # task list item + r"|^\s*\|.*\|" # table row (line leads with a pipe) + r"|^\s*\|?\s*:?-{2,}\s*\|" # table separator row (needs a pipe, so "--flag" stays plain) + r"|\[[^\]]*\]\(" # markdown link [text](url) + r"|\*\*[^*\n]+\*\*" # bold + r"|\*(?!\s)[^*\n]+?(? bool: + """Whether *text* contains a construct that needs native rich rendering.""" + return _TELEGRAM_RICH_CONSTRUCT_RE.search(text) is not None + def _load_telegram_input_file(path, filename: str): from telegram import InputFile @@ -310,7 +340,11 @@ class TelegramChannel(Channel): return False def _can_send_rich(self, text: str) -> bool: - return bool(self.config.get("rich_messages")) and 0 < len(text) <= TELEGRAM_MAX_RICH_MESSAGE_LENGTH + # Rich Messages are used only when rich_messages is on and the text + # actually contains a rich construct. Structured command/error replies + # are plain text with none, so they stay plain and their newlines and + # tokens are not collapsed into one line. + return bool(self.config.get("rich_messages")) and 0 < len(text) <= TELEGRAM_MAX_RICH_MESSAGE_LENGTH and _has_rich_constructs(text) async def _edit_rich_message(self, chat_id: int, message_id: int, text: str) -> bool: """Replace a streamed preview with a persistent Telegram Rich Message.""" diff --git a/backend/tests/test_channels.py b/backend/tests/test_channels.py index 7611a6fc5..4d01b7e44 100644 --- a/backend/tests/test_channels.py +++ b/backend/tests/test_channels.py @@ -10317,6 +10317,76 @@ class TestTelegramStreaming: _run(go()) + def test_plain_command_reply_stays_plain_when_enabled(self): + """A plain command/error reply (no rich construct) must stay plain text + even when rich_messages is on, so newlines and tokens + survive. The text deliberately carries a bracketed-pipe token + ([condition|clear]) to prove the construct detector stays well-formed.""" + + async def go(): + ch, bot = self._make_channel_with_bot() + ch.config["rich_messages"] = True + help_text = "Available commands:\n/goal [condition|clear] — Set or clear a goal\n/agent use — Start with an agent" + + await ch.send(OutboundMessage(channel_name="telegram", chat_id="12345", thread_id="t1", text=help_text, is_final=True)) + + assert bot.rich == [] + assert [message["text"] for message in bot.sent] == [help_text] + + _run(go()) + + def test_plain_flag_list_stays_plain_when_enabled(self): + """A reply whose lines merely *start* with ``--`` (CLI flag lists, + signature separators) must stay plain: a table-separator row needs a + pipe, so a bare ``--verbose`` line is not a GFM delimiter row.""" + + async def go(): + ch, bot = self._make_channel_with_bot() + ch.config["rich_messages"] = True + flag_list = "Options:\n--verbose\n--help" + + await ch.send(OutboundMessage(channel_name="telegram", chat_id="12345", thread_id="t1", text=flag_list, is_final=True)) + + assert bot.rich == [] + assert [message["text"] for message in bot.sent] == [flag_list] + + _run(go()) + + def test_plain_arithmetic_stays_plain_when_enabled(self): + """A reply with spaced asterisks used as multiplication (``2 * 3 * 4``) + must stay plain: italic emphasis needs tight, non-space delimiters, so + the spans between the asterisks are not handed to the rich parser.""" + + async def go(): + ch, bot = self._make_channel_with_bot() + ch.config["rich_messages"] = True + arithmetic = "Compute: 2 * 3 * 4 = 24" + + await ch.send(OutboundMessage(channel_name="telegram", chat_id="12345", thread_id="t1", text=arithmetic, is_final=True)) + + assert bot.rich == [] + assert [message["text"] for message in bot.sent] == [arithmetic] + + _run(go()) + + def test_plain_command_reply_stays_plain_on_stream_edit(self, monkeypatch): + """A streamed-then-final plain reply (no rich construct) must not be + replaced by a rich edit of the in-flight placeholder.""" + + async def go(): + ch, bot = self._make_channel_with_bot() + ch.config["rich_messages"] = True + monkeypatch.setattr("app.channels.telegram._monotonic", lambda: 1000.0) + + await ch._send_running_reply("12345", 42) + await ch.send(OutboundMessage(channel_name="telegram", chat_id="12345", thread_id="t1", text="Available commands:\n/new — new", is_final=True, thread_ts="42")) + + assert bot.rich == [] + # Final text is applied as a plain edit of the streamed placeholder. + assert [message["text"] for message in bot.edited] == ["Available commands:\n/new — new"] + + _run(go()) + def test_final_replaces_plain_stream_with_rich_message(self, monkeypatch): async def go(): ch, bot = self._make_channel_with_bot()