diff --git a/README.md b/README.md index a967fa56..59de3a8b 100644 --- a/README.md +++ b/README.md @@ -173,11 +173,13 @@ - [ ] 实现外部消息到 `model.Message` / `runner.Runner.Run` 的转换 - [ ] 实现 Runner Event 到文本、流式消息和卡片消息的转换 - [ ] 接入企业微信或微信相关通道 -- [ ] 再接入至少一种不同 IM 通道,例如 Telegram +- [x] 接入 Telegram long polling 文本通道(Issue #31;单 Binding、Gateway Dispatch、进程内幂等) +- [ ] 接入 Telegram webhook、媒体/rich update 或其他 IM 通道 - [ ] 实现 webhook 验签、账号与租户绑定、用户身份映射 - [ ] 使用 `tenant + channel + message_id` 实现幂等去重和缓存回复 - [x] 实现单聊/群聊 Session ID 规则及跨群、跨租户隔离 -- [ ] 处理消息分段、频率限制、异步回复、图片/文件、撤回和失败重试 +- [x] Telegram 文本回复分段、重复投递和论坛线程路由 +- [ ] 处理频率限制、异步回复、图片/文件、撤回和失败重试 - [ ] 增加重复投递、乱序、验签失败和跨租户访问测试 ### 治理、安全与可观测性 @@ -214,21 +216,24 @@ - [x] 列出至少 8 个生产风险及对应缓解措施 - [x] 持续标注可直接复用的 tRPC-Agent-Go 能力与平台新增模块边界 -> Issue #24 只完成架构、数据模型和运维文档;当前 PR 只确认下方部分 Gateway/API、进程内 -> 保护和 Binding-aware identity 条目。完整 Channel/Gateway 生产能力、真实 IM/Storage -> Adapter、迁移工具和生产告警仍未实现。 +> Issue #24 只完成架构、数据模型和运维文档;当前仓库另外交付了 Issue #28 的 Gateway/API +> 阶段能力和 Issue #31 的 Telegram long polling 文本 Adapter。完整 WeCom/Telegram webhook、 +> rich update、持久化消息能力、迁移工具和生产告警仍未实现。 ## 当前 PR 实现记录(不改变原验收要求) -> 以下内容仅索引 PR #29 当前 head 的实现和测试范围,不替代、收窄或修改上方原验收要求; +> 以下内容仅索引当前实现 PR 的 head 和测试范围,不替代、收窄或修改上方原验收要求; > 上方勾选仅表示该原条目已有完整证据。PR 尚未合并,当前阶段实现也不等同于全部原验收项已完成。 - `trpcservice/gateway/auth.go` 的 proof-bearing API 身份校验对应 `trpcservice/gateway/auth_test.go`;`resolver.go` 对应 `resolver_test.go`。 - 当前 PR 包含 Gateway、HTTP/SSE、进程内限流/幂等和 Channel trusted-principal 的阶段性代码, 具体边界以 `docs/docs/gateway.md` 为准。 +- Issue #31 的 `trpcservice/channels/telegram` 提供单 Binding、`getMe` 身份校验、普通文本 + long polling、Gateway Dispatch、进程内幂等和脱敏分段回复;具体边界以 + `docs/docs/telegram.md` 为准。 - Issue #26 的 fake candidate resolver/verifier 与 proof-bearing routing 边界有独立测试, - 但这不代表真实 IM Adapter、生产 webhook 或持久化能力已满足 README 原验收要求。 + 但这不代表 WeCom/Telegram webhook、媒体能力或持久化消息能力已满足 README 原验收要求。 ## 代码目录 diff --git a/docs/docs/architecture.md b/docs/docs/architecture.md index 956569ab..af93a5d9 100644 --- a/docs/docs/architecture.md +++ b/docs/docs/architecture.md @@ -41,7 +41,7 @@ Model Profile 和 Backend Profile,构造一个带版本、摘要和租户边 | Admin API | 租户、App、Profile、Binding 的管理、发布、回滚、审计入口 | Admin API → 控制面 Repository | 控制面模型已实现;HTTP API 为平台新增 | | Config/Registry | 版本校验、同租户引用、Factory/Storage 注册和缓存失效 | 控制面 → Gateway/Worker 快照 | Execution Plan/快照边界已有;Issue #28 交付进程内 Runner Registry,分布式失效仍为后续工作 | | Secret Resolver | 公开入站用不含 `tenant_id` 的 `CandidateBindingContext` 返回一次性验签 handle;验签后按固定 Tenant/Profile 作用域解析执行 Secret | Adapter/Gateway → Resolver;Resolver 不反向选租户 | Resolver 接口已有;candidate-scoped API、KMS/Secret Manager 为平台新增 | -| Channel Adapter | 解析供应商回调、校验协议、验签/解密、转换统一消息和出站回复 | IM ↔ Adapter ↔ Gateway | 包占位;真实 WeCom/Telegram Adapter 为平台新增 | +| Channel Adapter | 解析供应商回调、校验协议、验签/解密、转换统一消息和出站回复 | IM ↔ Adapter ↔ Gateway | Issue #31 已交付 Telegram long polling 普通文本;WeCom、Telegram webhook 和其他生产 Adapter 仍为平台新增 | | Agent Gateway | 限流、候选绑定路由、可信租户建立、幂等记录、快照装配和队列投递 | Adapter → Gateway → Worker/Queue | Issue #28 文档契约;实现阶段交付 HTTP/API 与 Channel principal 入口,真实队列仍为后续工作 | | Agent Worker | 消费固定执行计划,调用 Runner、Model、Tool 和 Storage,生成回复事件 | Gateway/Queue → Worker → 上游能力 | 最小 Runner spine 已有;Issue #28 交付进程内 Dispatch/HTTP,独立 Worker 为后续工作 | | Runner/Agent/Model | Agent 编排、模型调用、Tool/MCP、Event 和 context 取消 | Worker → tRPC-Agent-Go | 直接复用;当前已有最小 LLMAgent/Runner 装配 | @@ -586,6 +586,7 @@ Production architecture design (本页) └── persistent repositories and production adapters ``` -PR #25 的文档交付已完成,Issue #26 的 Channel Binding 领域与可信路由已实现;在 Issue #28 -代码验收完成前,README 不勾选 Gateway 的持续服务、Registry 或 HTTP/SSE 能力,避免文档 -进度掩盖工程边界。 +PR #25 的文档交付已完成,Issue #26 的 Channel Binding 领域与可信路由已实现,Issue #28 +交付了 Gateway 的进程内执行链,Issue #31 交付了 Telegram long polling 普通文本 Adapter; +WeCom/Telegram webhook、rich update、持久化和分布式能力仍按后续 Issue 交付,避免文档进度 +掩盖工程边界。 diff --git a/docs/docs/channel-binding.md b/docs/docs/channel-binding.md index 86390bb8..af1fa860 100644 --- a/docs/docs/channel-binding.md +++ b/docs/docs/channel-binding.md @@ -1,8 +1,9 @@ # Channel Binding 与可信入站路由 > 本页是 Issue #26 的实现契约。它先固定控制面模型、候选路由和可信边界,随后由 -> `trpcservice/channels` 的领域模型、InMemory Repository 和 fake verifier 实现。文中没有 -> 把真实企业微信/Telegram 适配器或 HTTP Gateway 误写成当前交付物。 +> `trpcservice/channels` 的领域模型、InMemory Repository 和 fake verifier 实现。Telegram +> long polling 的适配器契约见 [Telegram 长轮询 Adapter](telegram.md);本页仍只定义控制面和 +> trusted routing,不把协议运行时细节混入 Binding 领域模型。 ## 目标与边界 @@ -22,8 +23,9 @@ Channel Binding 把一个外部 IM 账号绑定到同一租户的 Agent App。 和无拼接碰撞的单聊/群聊/线程 Runner identity; - 使用 fake resolver/verifier 的离线集成测试。 -明确不在范围内:真实供应商 SDK、企业微信 AES 解密、Telegram Bot API、HTTP Gateway、KMS/ -Vault、PostgreSQL migration、消息去重/回复 Outbox、队列和生产审计持久化。 +明确不在范围内:真实供应商 SDK、企业微信 AES 解密、Telegram webhook、HTTP Gateway、KMS/ +Vault、PostgreSQL migration、消息去重/回复 Outbox、队列和生产审计持久化。Telegram long +polling 运行时契约见 [Telegram 长轮询 Adapter](telegram.md),不属于本 Binding 领域模型。 ## 控制面模型 diff --git a/docs/docs/gateway.md b/docs/docs/gateway.md index 05baf462..3e5dce58 100644 --- a/docs/docs/gateway.md +++ b/docs/docs/gateway.md @@ -23,7 +23,7 @@ PR #25 的架构验收继续约束组件职责:Channel Adapter 负责协议适 #26 的 `VerifiedBinding` / `RoutingTarget` 是 Channel principal 的唯一可信来源;本 Issue 不重新解释请求 body/header,也不从其中拼出租户。 -本 Issue 明确不实现真实 WeCom/Telegram webhook、OAuth/OIDC、KMS/Vault、Redis/SQL +本 Issue 明确不实现真实 WeCom/Telegram webhook(Telegram long polling 由 Issue #31 单独交付)、OAuth/OIDC、KMS/Vault、Redis/SQL 持久化、生产队列、Admin API、Graph/Chain/Parallel/Cycle 全量运行时或多节点一致性。 InMemory 限流、幂等、Registry 和 Session 只证明单进程契约,不能宣称跨节点生产语义。 diff --git a/docs/docs/index.md b/docs/docs/index.md index 93dd8d19..6017a9b8 100644 --- a/docs/docs/index.md +++ b/docs/docs/index.md @@ -20,6 +20,8 @@ - [Gateway、Execution Plan 与 HTTP/SSE](gateway.md):对齐 PR #25 架构验收与 Issue #26 可信主体,定义 Issue #28 的 Resolver、Runner Registry、Dispatch、普通/SSE API、 限流、幂等和服务生命周期契约。 +- [Telegram 长轮询 Adapter](telegram.md):Issue #31 的文档先行契约,固定单 Binding、Bot + 身份校验、普通文本映射、Dispatch 聚合回复和生命周期边界。 ## 快速开始 @@ -42,6 +44,7 @@ cd trpc-agent-service - [架构设计](architecture.md) — 组件拓扑、可信路由、消息链路、数据同步和多后端迁移 - [数据模型](data-model.md) — 核心表结构、Session/Event/Memory/Summary/Audit 和租户约束 - [Channel Binding](channel-binding.md) — 租户级通道绑定、候选发现与可信入站路由 +- [Telegram 长轮询 Adapter](telegram.md) — 单 Binding Telegram long polling、文本映射与安全边界 - [Gateway、Execution Plan 与 HTTP/SSE](gateway.md) — 可信主体、固定执行计划、Runner Registry、 Dispatch、健康检查、优雅停机和普通/流式 API - [运维方案](ops.md) — 发布灰度、监控审计、故障恢复、容量模型和生产风险清单 diff --git a/docs/docs/telegram.md b/docs/docs/telegram.md new file mode 100644 index 00000000..548e4c2d --- /dev/null +++ b/docs/docs/telegram.md @@ -0,0 +1,150 @@ +# Telegram 长轮询 Channel Adapter + +> Issue #31 的实现契约与状态记录。Telegram long polling 普通文本路径已实现并由单元、race +> 与全仓验证覆盖;Webhook、媒体/rich update、持久化和跨节点能力仍明确不在本 Issue 范围内。 + +## 1. 交付边界 + +Telegram 适配器是一个绑定级别的协议入口,不创建第二套租户或 Runner 路由。一个适配器实例 +只代表一个 active Telegram Binding,运行链路固定为: + +```text +Telegram getUpdates + -> Update.Message 校验和规范化 + -> 已验证的 channels.RoutingTarget + -> gateway.Channel Principal + -> gateway.DispatchService + -> 完整消费脱敏 DispatchEvent + -> 聚合并分段 + -> Telegram sendMessage +``` + +本 Issue 只实现 long polling 和普通文本消息。Webhook、媒体、命令、回调、持久化 outbox、跨 +节点 polling ownership 与分布式幂等保留给后续 Issue;`WebhookPath` 仍是绑定配置中的预留字段, +不能被本适配器读取为运行模式。 + +## 2. 适配器边界与构造 + +实现包位于 `trpcservice/channels/telegram`。公开构造配置的语义如下: + +| 配置 | 约束 | +| --- | --- | +| `BotToken` | 仅为运行时输入,不能写入 Binding、Plan、缓存键、日志、trace 或错误;构造成功后只由 SDK client 持有 | +| `Target` | 必须由现有 trusted boundary 创建并通过 `Validate()`;必须是 active Telegram Binding 的 RoutingTarget | +| `Dispatcher` | 使用现有 `gateway.DispatchService`;适配器不能直接调用 Runner | +| `Idempotency` | 可注入现有 `gateway.IdempotencyStore`;未注入时由适配器拥有一个进程内实例,不宣称跨进程保证 | +| `APIBaseURL` | 可选 HTTPS origin,仅用于 SDK Bot API;不从 Telegram update 或 webhook 字段读取 | +| `HTTPClient` | 可选的 SDK HTTP client,测试使用 fake/`httptest`,不要求真实凭据 | +| `PollTimeout` | 可选 long-poll timeout;采用 SDK 默认值时不自行覆盖 | +| `Workers` | 零值为 1;大于 1 必须由调用方显式配置,并由 SDK 同步处理 handler 生命周期 | +| `ErrorHook` | 只接收稳定的适配器错误类别,不接收 SDK/provider 原始错误、token 或 endpoint 凭据 | +| `Factory` | 注入式 Bot factory;生产实现才依赖 `github.com/go-telegram/bot`,测试不创建网络 client | + +构造函数接收 Context,先创建带默认 update handler 的 client,再调用 `getMe`,把返回的 Bot user +ID 规范化为十进制字符串并与 `Target.ProviderAccountID` 精确比较。创建失败、`getMe` 失败或身份 +不一致都 fail closed;在身份通过前不得处理任何 update。 + +适配器内部只保存由 `gateway.NewChannelPrincipal(Target)` 产生的 principal。Telegram update +不包含并且不能覆盖 tenant、binding、app、model、profile 或 routing hint;显示名、username 和 +标题只可作为未来展示元数据,不能参与认证、session 或 Runner identity。 + +Bot factory 对 SDK 使用以下固定策略: + +- `WithSkipGetMe`,由适配器在自己的 Context 中执行并校验 `getMe`; +- `WithDefaultHandler` 指向适配器的单 update handler; +- `WithNotAsyncHandlers`,使 `Run(ctx)` 结束时不会遗留 SDK 自己创建的异步 handler goroutine; +- `WithWorkers` 采用已验证的 worker 数量; +- 可选地设置 server URL、HTTP client 和 polling timeout; +- SDK polling error 只转换成稳定的 `polling` hook 事件,不能直接透传或记录原始错误。 + +## 3. 入站规范化与幂等 + +第一版只接受 `Update.Message` 中的普通文本: + +| Telegram 字段 | Gateway 字段 | 规则 | +| --- | --- | --- | +| `Message.Text` | `InboundMessage.Content` | 必须非空;继续交给 `InboundMessage.Normalize()` 做 trim/长度校验 | +| `Message.From.ID` | `ExternalUserID` | 必须存在且非零;不能使用 username/姓名 | +| private `Message.Chat.ID` | `ConversationDirect` + `ExternalPeerID` | chat ID 是稳定会话身份 | +| group/supergroup `Message.Chat.ID` | `ConversationGroup` + `ExternalChatID` | chat ID 进入群会话身份 | +| `Message.MessageThreadID` | `ExternalThreadID` | 大于零时保留;发送回复时原样作为 forum thread | +| `Update.ID` + trusted `BindingID` | `ExternalMessageID` / `RequestID` | 使用长度前缀编码生成稳定、无碰撞的 binding-aware ID | + +编辑消息、channel post、callback/inline、service update、无 sender/chat/text、未知 chat 类型和 +媒体-only update 都以稳定的非敏感原因忽略或拒绝,且不得进入 Dispatch。所有合法消息先用固定 +principal 调用 `IdempotencyStore.Begin`: + +- pending duplicate 不再次调用 Dispatch,也不启动隐藏 retry; +- completed duplicate 重用已缓存的脱敏 DispatchEvent,并只重新发送一个聚合后的逻辑回复; +- dispatch 或发送前的处理失败释放 claim,允许调用方按既有进程内策略重新处理; +- 该 store 只保证当前进程,不能暗示跨节点、重启恢复或持久化语义。 + +## 4. Dispatch 与回复 + +适配器把规范化消息和可信 principal 交给 `DispatchService.Dispatch`,`RequestID` 使用上节生成 +的稳定值,Context 原样向下传递。它必须读完事件 channel 直到关闭,不能在第一个 `message` 或 +`done` 后提前退出。只拼接 `DispatchEventMessage.Text`;收到脱敏 `error`、stream 异常或空的 +dispatch stream 时,发送固定的适配器级失败文本,不暴露 provider error、stack trace、Secret +或 repository 细节。 + +正常文本回复按 Unicode code point 切分,每段最多 4096 个 code point,并逐段调用: + +```text +sendMessage(chat_id=Message.Chat.ID, + message_thread_id=Message.MessageThreadID when > 0, + text=chunk) +``` + +不为每个 partial event 发送 Telegram 消息,不在本 Issue 引入编辑消息、队列、退避或后台重试。 +`sendMessage` 失败只通过稳定的 `send` hook 暴露;若已发送部分分段,不回滚也不启动隐式重试。 + +## 5. 生命周期与错误脱敏 + +```mermaid +sequenceDiagram + participant C as Service Context + participant A as Telegram Adapter + participant B as Bot SDK + participant G as Gateway Dispatch + participant T as Telegram API + + C->>A: New(ctx, runtime token, trusted Target) + A->>B: New + getMe + B-->>A: bot user ID + A->>A: compare provider_account_id; mismatch fails closed + C->>A: Run(ctx) + A->>B: Start(ctx) + B->>T: getUpdates (long polling) + T-->>B: Update.Message + B->>A: HandleUpdate(ctx, update) + A->>G: trusted Principal + normalized InboundMessage + G-->>A: complete redacted DispatchEvent stream + A->>T: one or more sendMessage chunks + C-->>B: cancel ctx + B-->>A: stop polling and synchronous handlers + A-->>C: Run returns +``` + +`Run(ctx)` 是阻塞入口;Context 取消必须同时结束 SDK polling、在途 Dispatch 和 sendMessage。适配器 +不创建自己的 retry goroutine,不保存 request Context,不持有 Runner lease;lease 和 event drain +由现有 Gateway contract 管理。`Close` 只关闭适配器拥有的进程内幂等 store,不能关闭调用方注入 +的 store 或 HTTP client。 + +错误 hook 只使用 `initialization`、`polling`、`update`、`dispatch`、`send` 等稳定 operation 和 +适配器 sentinel error。原始 SDK error 只能用于本地判断,不能出现在返回值、hook payload、日志、 +trace 或 Telegram 回复中。 + +## 6. 文档与代码验收清单 + +README 和 MkDocs 状态应明确区分已交付与后续能力: + +- [x] SDK 版本固定,Bot factory/client 可注入,`Run(ctx)` 和单 update handler 可测试; +- [x] `getMe` 身份校验、tenant/Binding/Runner 隔离、普通文本映射和 binding-aware 幂等通过测试; +- [x] Dispatch 完整消费、单逻辑回复、4096 code point 分段、forum thread 路由和失败脱敏通过测试; +- [x] cancellation、polling error、send failure、duplicate delivery 和资源生命周期通过测试; +- [x] Telegram long polling 已实现;Webhook、持久化幂等/outbox、媒体、跨节点 ownership + 和其他 rich update 明确保持未勾选。 + +参考:[Telegram Bot API](https://core.telegram.org/bots/api)、 +[getUpdates](https://core.telegram.org/bots/api#getting-updates)、 +[github.com/go-telegram/bot](https://github.com/go-telegram/bot)。 diff --git a/docs/mkdocs.yml b/docs/mkdocs.yml index 5f63ae62..0fa190b8 100644 --- a/docs/mkdocs.yml +++ b/docs/mkdocs.yml @@ -30,6 +30,7 @@ nav: - 架构设计: architecture.md - 数据模型: data-model.md - Channel Binding: channel-binding.md + - Telegram 长轮询 Adapter: telegram.md - Gateway、Execution Plan 与 HTTP/SSE: gateway.md - Tenant 运行时边界: tenant-runtime-boundary.md - Agent App 模型: agent-app-model.md diff --git a/go.mod b/go.mod index 2efd62c8..d87171cc 100644 --- a/go.mod +++ b/go.mod @@ -4,6 +4,8 @@ go 1.21 require trpc.group/trpc-go/trpc-agent-go v1.11.2 +require github.com/go-telegram/bot v1.23.0 + require ( github.com/bmatcuk/doublestar/v4 v4.9.1 // indirect github.com/cenkalti/backoff/v4 v4.3.0 // indirect diff --git a/go.sum b/go.sum index 10ec298d..815284f4 100644 --- a/go.sum +++ b/go.sum @@ -16,6 +16,8 @@ github.com/go-logr/logr v1.4.3 h1:CjnDlHq8ikf6E492q6eKboGOC0T8CDaOvkHCIg8idEI= github.com/go-logr/logr v1.4.3/go.mod h1:9T104GzyrTigFIr8wt5mBrctHMim0Nb2HLGrmQ40KvY= github.com/go-logr/stdr v1.2.2 h1:hSWxHoqTgW2S2qGc0LTAI563KZ5YKYRhT3MFKZMbjag= github.com/go-logr/stdr v1.2.2/go.mod h1:mMo/vtBO5dYbehREoey6XUKy/eSumjCCveDpRre4VKE= +github.com/go-telegram/bot v1.23.0 h1:CKKQq115G/GUGBG8uuWl5uXbiBHyVjZBp/qqOLWZjJk= +github.com/go-telegram/bot v1.23.0/go.mod h1:i2TRs7fXWIeaceF3z7KzsMt/he0TwkVC680mvdTFYeM= github.com/google/go-cmp v0.6.0 h1:ofyhxvXcZhMsU5ulbFiLKl/XBFqE1GSq7atu8tAmTRI= github.com/google/go-cmp v0.6.0/go.mod h1:17dUlkBOakJ0+DkrSSNjCkIjxS6bF9zb3elmeNGIjoY= github.com/google/uuid v1.6.0 h1:NIvaJDMOsjHA8n1jAhLSgzrAzy1Hgr+hNrb57e+94F0= diff --git a/trpcservice/channels/telegram/telegram.go b/trpcservice/channels/telegram/telegram.go new file mode 100644 index 00000000..ab4fa736 --- /dev/null +++ b/trpcservice/channels/telegram/telegram.go @@ -0,0 +1,705 @@ +// Package telegram implements the tenant-scoped Telegram long-polling +// Channel Adapter. +package telegram + +import ( + "context" + "errors" + "fmt" + "net/http" + "net/url" + "strconv" + "strings" + "sync" + "time" + + "github.com/XnLemon/trpc-agent-service/trpcservice/channels" + "github.com/XnLemon/trpc-agent-service/trpcservice/gateway" + "github.com/go-telegram/bot" + "github.com/go-telegram/bot/models" +) + +const ( + defaultPollTimeout = time.Minute + minimumPollTimeout = 2 * time.Second + maximumPollTimeout = 10 * time.Minute + maximumTokenRunes = 1024 + maximumReplyRunes = 4096 + failureReply = "Sorry, I couldn't process that message." +) + +var ( + // ErrInvalid reports malformed adapter configuration or update input. + ErrInvalid = errors.New("invalid telegram adapter input") + // ErrNotReady reports that the adapter has no usable Bot client. + ErrNotReady = errors.New("telegram adapter is not ready") + // ErrClosed reports an adapter or its owned process-local state after close. + ErrClosed = errors.New("telegram adapter is closed") + // ErrAlreadyRunning reports a second concurrent Run call. + ErrAlreadyRunning = errors.New("telegram adapter is already running") + // ErrInitialization reports a redacted Bot construction or getMe failure. + ErrInitialization = errors.New("telegram bot initialization failed") + // ErrBotIdentityMismatch reports a Bot identity different from the trusted + // Binding provider account. + ErrBotIdentityMismatch = errors.New("telegram bot identity does not match binding") + // ErrInvalidUpdate reports a malformed supported update shape. + ErrInvalidUpdate = errors.New("invalid telegram update") + // ErrUnsupportedUpdate reports an update outside the first text-only scope. + ErrUnsupportedUpdate = errors.New("unsupported telegram update") + // ErrDuplicateUpdate reports an update already being handled by this process. + ErrDuplicateUpdate = errors.New("duplicate telegram update") + // ErrDispatch reports a redacted Gateway dispatch failure. + ErrDispatch = errors.New("telegram dispatch failed") + // ErrSendMessage reports a redacted Telegram sendMessage failure. + ErrSendMessage = errors.New("telegram send message failed") + // ErrPolling reports a redacted SDK polling error delivered to ErrorHook. + ErrPolling = errors.New("telegram polling failed") +) + +// ErrorOperation identifies the safe operation category supplied to ErrorHook. +type ErrorOperation string + +const ( + // ErrorOperationInitialization identifies construction or getMe failures. + ErrorOperationInitialization ErrorOperation = "initialization" + // ErrorOperationPolling identifies long-polling failures. + ErrorOperationPolling ErrorOperation = "polling" + // ErrorOperationUpdate identifies rejected or unsupported updates. + ErrorOperationUpdate ErrorOperation = "update" + // ErrorOperationDispatch identifies Gateway execution failures. + ErrorOperationDispatch ErrorOperation = "dispatch" + // ErrorOperationSend identifies outbound sendMessage failures. + ErrorOperationSend ErrorOperation = "send" +) + +// ErrorEvent is the redacted payload passed to an ErrorHook. Err is always a +// stable adapter sentinel and never a provider error, token, or stack trace. +type ErrorEvent struct { + Operation ErrorOperation + Err error +} + +// ErrorHook observes stable adapter failures without receiving provider +// details or runtime secrets. +type ErrorHook func(ErrorEvent) + +// BotClient is the small SDK surface required by the adapter. A fake client +// can implement it without credentials or network access. +type BotClient interface { + Start(context.Context) + GetMe(context.Context) (*models.User, error) + SendMessage(context.Context, *bot.SendMessageParams) (*models.Message, error) +} + +// BotFactoryConfig contains non-secret options for constructing one BotClient. +// The token is passed separately to BotFactory.New and is never stored in this +// configuration value. +type BotFactoryConfig struct { + // Handler receives updates from the SDK's long-polling consumer. + Handler bot.HandlerFunc + // APIBaseURL overrides the Telegram API origin when non-empty. + APIBaseURL string + // HTTPClient is the optional HTTP transport used by the SDK. + HTTPClient bot.HttpClient + // PollTimeout is the Telegram getUpdates long-poll timeout. + PollTimeout time.Duration + // Workers is the explicitly validated SDK update worker count. + Workers int + // OnPollingError receives no raw error; it only signals that polling failed. + OnPollingError func() +} + +// BotFactory constructs a BotClient with the supplied runtime token and safe +// options. Production uses the public github.com/go-telegram/bot SDK; tests +// inject a fake implementation. +type BotFactory interface { + New(string, BotFactoryConfig) (BotClient, error) +} + +// BotFactoryFunc adapts a function into a BotFactory. +type BotFactoryFunc func(string, BotFactoryConfig) (BotClient, error) + +// New implements BotFactory for a BotFactoryFunc. +func (factory BotFactoryFunc) New(token string, config BotFactoryConfig) (BotClient, error) { + if factory == nil { + return nil, ErrInvalid + } + return factory(token, config) +} + +// Config defines one tenant-scoped Telegram Binding adapter. BotToken is a +// runtime-only secret and must not be persisted or placed in diagnostics. +type Config struct { + // BotToken is resolved before construction and retained only by the SDK + // client created for this adapter. + BotToken string + // Target is the trusted, non-secret route for exactly one active Binding. + Target channels.RoutingTarget + // Dispatcher is the existing protocol-neutral Gateway execution service. + Dispatcher gateway.DispatchService + // Idempotency optionally supplies a shared process-local store. When nil, + // the adapter owns a new process-local store. + Idempotency *gateway.IdempotencyStore + // APIBaseURL optionally overrides the Telegram HTTPS API origin. + APIBaseURL string + // HTTPClient optionally supplies the SDK HTTP transport. + HTTPClient bot.HttpClient + // PollTimeout optionally overrides the long-poll timeout. + PollTimeout time.Duration + // Workers controls SDK update workers. Zero defaults to one. + Workers int + // ErrorHook observes stable, redacted adapter failures. + ErrorHook ErrorHook + // Factory optionally replaces the public SDK factory for tests. + Factory BotFactory +} + +// Adapter owns one trusted Telegram Binding and routes its updates through the +// existing Gateway contracts. It does not create or cache a Runner directly. +type Adapter struct { + client BotClient + dispatcher gateway.DispatchService + principal gateway.Principal + target channels.RoutingTarget + idempotency *gateway.IdempotencyStore + ownIdempotency bool + errorHook ErrorHook + + mu sync.RWMutex + closed bool + runCancel context.CancelFunc +} + +// New validates the trusted route, constructs the Bot client, and verifies its +// getMe identity before returning an adapter that can handle updates. +func New(ctx context.Context, config Config) (*Adapter, error) { + if ctx == nil { + return nil, fmt.Errorf("%w: context is required", ErrInvalid) + } + if err := ctx.Err(); err != nil { + return nil, err + } + token, err := normalizeToken(config.BotToken) + if err != nil { + return nil, err + } + if err := config.Target.Validate(); err != nil { + return nil, fmt.Errorf("%w: trusted routing target is invalid", ErrInvalid) + } + if config.Target.Channel != channels.ChannelTelegram { + return nil, fmt.Errorf("%w: routing target is not Telegram", ErrInvalid) + } + providerAccountID, err := strconv.ParseInt(config.Target.ProviderAccountID, 10, 64) + if err != nil || providerAccountID <= 0 || strconv.FormatInt(providerAccountID, 10) != config.Target.ProviderAccountID { + return nil, fmt.Errorf("%w: Telegram provider account ID is not canonical", ErrInvalid) + } + if config.Dispatcher == nil { + return nil, fmt.Errorf("%w: dispatcher is required", ErrInvalid) + } + apiBaseURL, err := normalizeAPIBaseURL(config.APIBaseURL) + if err != nil { + return nil, err + } + pollTimeout, err := normalizePollTimeout(config.PollTimeout) + if err != nil { + return nil, err + } + workers, err := normalizeWorkers(config.Workers) + if err != nil { + return nil, err + } + principal, err := gateway.NewChannelPrincipal(config.Target) + if err != nil { + return nil, fmt.Errorf("%w: trusted principal is invalid", ErrInvalid) + } + + idempotency := config.Idempotency + ownIdempotency := false + if idempotency == nil { + idempotency, err = gateway.NewIdempotencyStore(gateway.IdempotencyConfig{}) + if err != nil { + return nil, fmt.Errorf("%w: idempotency store is unavailable", ErrInvalid) + } + ownIdempotency = true + } + factory := config.Factory + if factory == nil { + factory = sdkBotFactory{} + } + adapter := &Adapter{ + dispatcher: config.Dispatcher, principal: principal, target: config.Target, + idempotency: idempotency, ownIdempotency: ownIdempotency, errorHook: config.ErrorHook, + } + client, err := factory.New(token, BotFactoryConfig{ + Handler: adapter.sdkHandler(), + APIBaseURL: apiBaseURL, + HTTPClient: config.HTTPClient, + PollTimeout: pollTimeout, + Workers: workers, + OnPollingError: func() { adapter.report(ErrorOperationPolling, ErrPolling) }, + }) + if err != nil || client == nil { + adapter.report(ErrorOperationInitialization, ErrInitialization) + _ = adapter.closeOwnedIdempotency() + return nil, ErrInitialization + } + adapter.client = client + me, err := client.GetMe(ctx) + if err != nil { + if contextErr := ctx.Err(); contextErr != nil { + _ = adapter.closeOwnedIdempotency() + return nil, contextErr + } + adapter.report(ErrorOperationInitialization, ErrInitialization) + _ = adapter.closeOwnedIdempotency() + return nil, ErrInitialization + } + if me == nil || !me.IsBot || me.ID <= 0 || strconv.FormatInt(me.ID, 10) != config.Target.ProviderAccountID { + adapter.report(ErrorOperationInitialization, ErrBotIdentityMismatch) + _ = adapter.closeOwnedIdempotency() + return nil, ErrBotIdentityMismatch + } + return adapter, nil +} + +// Run starts blocking Telegram long polling and returns after ctx is canceled +// or Close cancels the run. The SDK owns its polling and worker goroutines. +func (adapter *Adapter) Run(ctx context.Context) error { + if ctx == nil { + return fmt.Errorf("%w: context is required", ErrInvalid) + } + if err := ctx.Err(); err != nil { + return err + } + if adapter == nil { + return ErrNotReady + } + adapter.mu.Lock() + closed, client := adapter.closed, adapter.client + if closed { + adapter.mu.Unlock() + return ErrClosed + } + if client == nil { + adapter.mu.Unlock() + return ErrNotReady + } + if adapter.runCancel != nil { + adapter.mu.Unlock() + return ErrAlreadyRunning + } + runContext, cancel := context.WithCancel(ctx) + adapter.runCancel = cancel + adapter.mu.Unlock() + defer func() { + adapter.mu.Lock() + adapter.runCancel = nil + adapter.mu.Unlock() + cancel() + }() + client.Start(runContext) + return nil +} + +// Close closes only idempotency state owned by the adapter. Injected stores and +// HTTP clients remain owned by their callers; polling is stopped by canceling +// the Context passed to Run. +func (adapter *Adapter) Close() error { + if adapter == nil { + return nil + } + adapter.mu.Lock() + if adapter.closed { + adapter.mu.Unlock() + return nil + } + adapter.closed = true + cancel := adapter.runCancel + adapter.mu.Unlock() + if cancel != nil { + cancel() + } + return adapter.closeOwnedIdempotency() +} + +// HandleUpdate validates and processes one Telegram update. It is exposed for +// deterministic tests and for the SDK's default handler. +func (adapter *Adapter) HandleUpdate(ctx context.Context, update *models.Update) error { + if adapter == nil { + return ErrNotReady + } + if ctx == nil { + err := fmt.Errorf("%w: context is required", ErrInvalid) + adapter.report(ErrorOperationUpdate, ErrInvalid) + return err + } + if err := ctx.Err(); err != nil { + return err + } + adapter.mu.RLock() + closed, client := adapter.closed, adapter.client + adapter.mu.RUnlock() + if closed { + return ErrClosed + } + if client == nil || adapter.idempotency == nil { + return ErrNotReady + } + message, err := normalizeUpdate(adapter.target, update) + if err != nil { + adapter.report(ErrorOperationUpdate, err) + return err + } + claim, replay, err := adapter.idempotency.Begin(ctx, adapter.principal, message) + if err != nil { + if errors.Is(err, gateway.ErrDuplicateMessage) { + return ErrDuplicateUpdate + } + if errors.Is(err, context.Canceled) || errors.Is(err, context.DeadlineExceeded) { + return err + } + if errors.Is(err, gateway.ErrClosed) { + return ErrClosed + } + return ErrInvalid + } + if claim == nil { + if err := adapter.sendEvents(ctx, update.Message, replay); err != nil { + if errors.Is(err, context.Canceled) || errors.Is(err, context.DeadlineExceeded) { + return err + } + adapter.report(ErrorOperationSend, ErrSendMessage) + return err + } + return nil + } + + events, dispatchErr := adapter.dispatch(ctx, message) + if dispatchErr != nil { + _ = claim.Fail() + if errors.Is(dispatchErr, context.Canceled) || errors.Is(dispatchErr, context.DeadlineExceeded) { + return dispatchErr + } + adapter.report(ErrorOperationDispatch, ErrDispatch) + if sendErr := adapter.sendText(ctx, update.Message, failureReply); sendErr != nil { + if !errors.Is(sendErr, context.Canceled) && !errors.Is(sendErr, context.DeadlineExceeded) { + adapter.report(ErrorOperationSend, ErrSendMessage) + } + } + return ErrDispatch + } + if err := claim.Complete(events); err != nil { + return ErrDispatch + } + if err := adapter.sendEvents(ctx, update.Message, events); err != nil { + if errors.Is(err, context.Canceled) || errors.Is(err, context.DeadlineExceeded) { + return err + } + adapter.report(ErrorOperationSend, ErrSendMessage) + return err + } + return nil +} + +func (adapter *Adapter) sdkHandler() bot.HandlerFunc { + return func(ctx context.Context, _ *bot.Bot, update *models.Update) { + _ = adapter.HandleUpdate(ctx, update) + } +} + +func (adapter *Adapter) dispatch(ctx context.Context, message gateway.InboundMessage) ([]gateway.DispatchEvent, error) { + stream, err := adapter.dispatcher.Dispatch(ctx, gateway.DispatchRequest{ + Principal: adapter.principal, Message: message, RequestID: message.ExternalMessageID, + }) + if err != nil { + if errors.Is(err, context.Canceled) || errors.Is(err, context.DeadlineExceeded) { + return nil, err + } + return nil, ErrDispatch + } + if stream == nil { + return nil, ErrDispatch + } + events := make([]gateway.DispatchEvent, 0, 4) + done := false + failed := false + for { + select { + case <-ctx.Done(): + return nil, ctx.Err() + case event, ok := <-stream: + if !ok { + if !done || failed { + return nil, ErrDispatch + } + return events, nil + } + events = append(events, event) + if event.Type == gateway.DispatchEventError { + failed = true + } + if event.Done { + done = true + } + } + } +} + +func (adapter *Adapter) sendEvents(ctx context.Context, message *models.Message, events []gateway.DispatchEvent) error { + for _, event := range events { + if event.Type == gateway.DispatchEventError { + return adapter.sendText(ctx, message, failureReply) + } + } + var builder strings.Builder + for _, event := range events { + if event.Type == gateway.DispatchEventMessage { + builder.WriteString(event.Text) + } + } + if builder.Len() == 0 { + return nil + } + return adapter.sendText(ctx, message, builder.String()) +} + +func (adapter *Adapter) sendText(ctx context.Context, message *models.Message, text string) error { + if message == nil { + return ErrInvalidUpdate + } + if adapter == nil || adapter.client == nil { + return ErrNotReady + } + chunks := splitText(text, maximumReplyRunes) + for _, chunk := range chunks { + if err := ctx.Err(); err != nil { + return err + } + _, err := adapter.client.SendMessage(ctx, &bot.SendMessageParams{ + ChatID: message.Chat.ID, MessageThreadID: message.MessageThreadID, Text: chunk, + }) + if err != nil { + return ErrSendMessage + } + } + return nil +} + +func normalizeUpdate(target channels.RoutingTarget, update *models.Update) (gateway.InboundMessage, error) { + if update == nil || update.ID < 0 { + return gateway.InboundMessage{}, ErrInvalidUpdate + } + if hasUnsupportedUpdate(update) || update.Message == nil { + return gateway.InboundMessage{}, ErrUnsupportedUpdate + } + message := update.Message + if hasUnsupportedMessage(message) || strings.TrimSpace(message.Text) == "" { + return gateway.InboundMessage{}, ErrUnsupportedUpdate + } + if message.From == nil || message.From.ID <= 0 || message.Chat.ID == 0 { + return gateway.InboundMessage{}, ErrInvalidUpdate + } + if message.MessageThreadID < 0 { + return gateway.InboundMessage{}, ErrInvalidUpdate + } + inbound := gateway.InboundMessage{ + Content: message.Text, ContentType: gateway.ContentTypeText, + ExternalMessageID: externalMessageID(target, update.ID), + ExternalUserID: strconv.FormatInt(message.From.ID, 10), + } + switch message.Chat.Type { + case models.ChatTypePrivate: + inbound.ConversationKind = channels.ConversationDirect + inbound.ExternalPeerID = strconv.FormatInt(message.Chat.ID, 10) + case models.ChatTypeGroup, models.ChatTypeSupergroup: + inbound.ConversationKind = channels.ConversationGroup + inbound.ExternalChatID = strconv.FormatInt(message.Chat.ID, 10) + default: + return gateway.InboundMessage{}, ErrUnsupportedUpdate + } + if message.MessageThreadID > 0 { + inbound.ExternalThreadID = strconv.Itoa(message.MessageThreadID) + } + normalized, err := inbound.Normalize() + if err != nil { + return gateway.InboundMessage{}, ErrInvalidUpdate + } + return normalized, nil +} + +func hasUnsupportedMessage(message *models.Message) bool { + if message == nil { + return true + } + for _, entity := range message.Entities { + if entity.Type == models.MessageEntityTypeBotCommand { + return true + } + } + return message.DirectMessagesTopic != nil || message.SenderChat != nil || + message.SenderBusinessBot != nil || message.ReceiverUser != nil || message.BusinessConnectionID != "" || + message.RichMessage != nil || message.Animation != nil || message.Audio != nil || message.Document != nil || + message.PaidMedia != nil || len(message.Photo) > 0 || message.Sticker != nil || message.Story != nil || + message.Video != nil || message.VideoNote != nil || message.Voice != nil || message.Caption != "" || + len(message.CaptionEntities) > 0 || message.HasMediaSpoiler || message.Checklist != nil || + message.MediaGroupID != "" || message.ReplyToStore != nil || message.SuggestedPostInfo != nil || + message.EffectID != "" || message.EditDate != 0 || + message.Contact != nil || message.Dice != nil || message.Game != nil || message.Poll != nil || + message.Venue != nil || message.Location != nil || len(message.NewChatMembers) > 0 || + message.LeftChatMember != nil || message.NewChatTitle != "" || len(message.NewChatPhoto) > 0 || + message.DeleteChatPhoto || message.GroupChatCreated || message.SupergroupChatCreated || + message.ChannelChatCreated || message.MessageAutoDeleteTimerChanged != nil || message.MigrateToChatID != 0 || + message.MigrateFromChatID != 0 || message.PinnedMessage != nil || message.Invoice != nil || + message.SuccessfulPayment != nil || message.RefundedPayment != nil || message.UsersShared != nil || + message.ChatShared != nil || message.Gift != nil || message.UniqueGift != nil || + message.GiftUpgradeSent != nil || message.ConnectedWebsite != "" || message.WriteAccessAllowed != nil || + message.PassportData != nil || message.ProximityAlertTriggered != nil || message.BoostAdded != nil || + message.ChatBackgroundSet != nil || message.ChecklistTasksDone != nil || message.ChecklistTasksAdded != nil || + message.DirectMessagePriceChanged != nil || message.ForumTopicCreated != nil || message.ForumTopicEdited != nil || + message.ForumTopicClosed != nil || message.ForumTopicReopened != nil || message.GeneralForumTopicHidden != nil || + message.GeneralForumTopicUnhidden != nil || message.GiveawayCreated != nil || message.Giveaway != nil || + message.GiveawayWinners != nil || message.GiveawayCompleted != nil || message.PaidMessagePriceChanged != nil || + message.ChatOwnerLeft != nil || message.ChatOwnerChanged != nil || message.CommunityChatAdded != nil || + message.CommunityChatRemoved != nil || message.SuggestedPostApproved != nil || + message.SuggestedPostApprovalFailed != nil || message.SuggestedPostDeclined != nil || + message.SuggestedPostPaid != nil || message.SuggestedPostRefunded != nil || message.VideoChatScheduled != nil || + message.VideoChatStarted != nil || message.VideoChatEnded != nil || message.VideoChatParticipantsInvited != nil || + message.WebAppData != nil || message.ManagedBotCreated != nil || message.PollOptionAdded != nil || + message.PollOptionDeleted != nil || message.GuestBotCallerUser != nil || message.GuestBotCallerChat != nil || + message.GuestQueryID != "" || message.ReplyToPollOptionID != "" || message.LivePhoto != nil +} + +func hasUnsupportedUpdate(update *models.Update) bool { + return update.EditedMessage != nil || update.ChannelPost != nil || update.EditedChannelPost != nil || + update.BusinessConnection != nil || update.BusinessMessage != nil || update.EditedBusinessMessage != nil || + update.DeletedBusinessMessages != nil || update.MessageReaction != nil || update.MessageReactionCount != nil || + update.InlineQuery != nil || update.ChosenInlineResult != nil || update.CallbackQuery != nil || + update.ShippingQuery != nil || update.PreCheckoutQuery != nil || update.PurchasedPaidMedia != nil || + update.Poll != nil || update.PollAnswer != nil || update.ManagedBot != nil || update.GuestMessage != nil || + update.MyChatMember != nil || update.ChatMember != nil || update.ChatJoinRequest != nil || + update.ChatBoost != nil || update.RemovedChatBoost != nil || update.Subscription != nil +} + +func externalMessageID(target channels.RoutingTarget, updateID int64) string { + return encodeParts("telegram-update", target.BindingID, strconv.FormatInt(updateID, 10)) +} + +func encodeParts(parts ...string) string { + var builder strings.Builder + for _, part := range parts { + builder.WriteString(strconv.Itoa(len([]byte(part)))) + builder.WriteByte(':') + builder.WriteString(part) + } + return builder.String() +} + +func splitText(text string, maximum int) []string { + if text == "" { + return nil + } + runes := []rune(text) + if len(runes) <= maximum { + return []string{text} + } + chunks := make([]string, 0, (len(runes)+maximum-1)/maximum) + for len(runes) > 0 { + end := maximum + if len(runes) < end { + end = len(runes) + } + chunks = append(chunks, string(runes[:end])) + runes = runes[end:] + } + return chunks +} + +func normalizeToken(token string) (string, error) { + if token == "" || strings.TrimSpace(token) != token || len([]rune(token)) > maximumTokenRunes || hasControl(token) { + return "", fmt.Errorf("%w: bot token is invalid", ErrInvalid) + } + return token, nil +} + +func normalizeAPIBaseURL(value string) (string, error) { + value = strings.TrimSpace(value) + if value == "" { + return "", nil + } + parsed, err := url.Parse(value) + if err != nil || parsed.Scheme != "https" || parsed.Host == "" || parsed.User != nil || (parsed.Path != "" && parsed.Path != "/") || parsed.RawPath != "" || parsed.RawQuery != "" || parsed.Fragment != "" || hasControl(value) { + return "", fmt.Errorf("%w: API base URL must be an HTTPS origin", ErrInvalid) + } + return strings.TrimRight(value, "/"), nil +} + +func normalizePollTimeout(value time.Duration) (time.Duration, error) { + if value == 0 { + return defaultPollTimeout, nil + } + if value < minimumPollTimeout || value > maximumPollTimeout { + return 0, fmt.Errorf("%w: polling timeout is outside supported bounds", ErrInvalid) + } + return value, nil +} + +func normalizeWorkers(value int) (int, error) { + if value == 0 { + return 1, nil + } + if value < 1 { + return 0, fmt.Errorf("%w: worker count must be positive", ErrInvalid) + } + return value, nil +} + +func hasControl(value string) bool { + for _, character := range value { + if character < 0x20 || character == 0x7f { + return true + } + } + return false +} + +func (adapter *Adapter) report(operation ErrorOperation, err error) { + if adapter == nil || adapter.errorHook == nil || err == nil { + return + } + adapter.errorHook(ErrorEvent{Operation: operation, Err: err}) +} + +func (adapter *Adapter) closeOwnedIdempotency() error { + if adapter == nil || !adapter.ownIdempotency || adapter.idempotency == nil { + return nil + } + return adapter.idempotency.Close() +} + +type sdkBotFactory struct{} + +func (sdkBotFactory) New(token string, config BotFactoryConfig) (BotClient, error) { + if config.Handler == nil || config.Workers < 1 || config.PollTimeout < minimumPollTimeout { + return nil, ErrInvalid + } + options := []bot.Option{ + bot.WithSkipGetMe(), bot.WithDefaultHandler(config.Handler), bot.WithNotAsyncHandlers(), bot.WithWorkers(config.Workers), + bot.WithHTTPClient(config.PollTimeout, configuredHTTPClient(config.HTTPClient, config.PollTimeout)), + bot.WithErrorsHandler(func(error) { + if config.OnPollingError != nil { + config.OnPollingError() + } + }), + } + if config.APIBaseURL != "" { + options = append(options, bot.WithServerURL(config.APIBaseURL)) + } + return bot.New(token, options...) +} + +func configuredHTTPClient(client bot.HttpClient, pollTimeout time.Duration) bot.HttpClient { + if client != nil { + return client + } + return &http.Client{Timeout: pollTimeout + 5*time.Second} +} diff --git a/trpcservice/channels/telegram/telegram_test.go b/trpcservice/channels/telegram/telegram_test.go new file mode 100644 index 00000000..4de36a94 --- /dev/null +++ b/trpcservice/channels/telegram/telegram_test.go @@ -0,0 +1,861 @@ +package telegram + +import ( + "context" + "crypto/sha256" + "encoding/hex" + "errors" + "fmt" + "net/http" + "strings" + "sync" + "testing" + "time" + + "github.com/XnLemon/trpc-agent-service/trpcservice/agent" + "github.com/XnLemon/trpc-agent-service/trpcservice/channels" + "github.com/XnLemon/trpc-agent-service/trpcservice/channels/inmemory" + "github.com/XnLemon/trpc-agent-service/trpcservice/gateway" + "github.com/XnLemon/trpc-agent-service/trpcservice/tenant" + "github.com/go-telegram/bot" + "github.com/go-telegram/bot/models" +) + +func TestNewInjectsFactoryAndRejectsBotIdentityMismatch(t *testing.T) { + target := newTrustedTarget(t, channels.ChannelTelegram, "constructor", "12345") + dispatcher := &dispatchStub{} + client := &fakeBot{me: &models.User{ID: 12345, IsBot: true}} + factory := &fakeFactory{client: client} + adapter, err := New(context.Background(), Config{ + BotToken: "12345:runtime-secret", Target: target, Dispatcher: dispatcher, Factory: factory, + APIBaseURL: "https://api.example.test", PollTimeout: 3 * time.Second, Workers: 2, + }) + if err != nil { + t.Fatal(err) + } + if factory.token != "12345:runtime-secret" || factory.config.Handler == nil { + t.Fatal("factory did not receive the runtime token and update handler") + } + if factory.config.APIBaseURL != "https://api.example.test" || factory.config.PollTimeout != 3*time.Second || factory.config.Workers != 2 { + t.Fatalf("factory received unexpected options: %+v", factory.config) + } + if adapter == nil || adapter.principal.Kind() != gateway.PrincipalChannel { + t.Fatal("adapter did not retain a trusted channel principal") + } + + recorder := &errorRecorder{} + mismatched := &fakeFactory{client: &fakeBot{me: &models.User{ID: 54321, IsBot: true}}} + _, err = New(context.Background(), Config{ + BotToken: "12345:runtime-secret", Target: target, Dispatcher: dispatcher, Factory: mismatched, + ErrorHook: recorder.hook, + }) + if !errors.Is(err, ErrBotIdentityMismatch) { + t.Fatalf("identity mismatch was accepted or leaked another error: %v", err) + } + events := recorder.snapshot() + if len(events) != 1 || events[0].Operation != ErrorOperationInitialization || !errors.Is(events[0].Err, ErrBotIdentityMismatch) { + t.Fatalf("unexpected identity error hook: %+v", events) + } + if len(dispatcher.requests()) != 0 { + t.Fatal("identity mismatch reached Dispatch") + } +} + +func TestNewRejectsNonTelegramTargetAndInvalidRuntimeOptions(t *testing.T) { + wecomTarget := newTrustedTarget(t, channels.ChannelWeCom, "wrong-channel", "corp-1") + factoryCalls := 0 + factory := BotFactoryFunc(func(string, BotFactoryConfig) (BotClient, error) { + factoryCalls++ + return &fakeBot{me: &models.User{ID: 1, IsBot: true}}, nil + }) + _, err := New(context.Background(), Config{ + BotToken: "token", Target: wecomTarget, Dispatcher: &dispatchStub{}, Factory: factory, + }) + if !errors.Is(err, ErrInvalid) || factoryCalls != 0 { + t.Fatalf("non-Telegram target was not rejected before factory: err=%v calls=%d", err, factoryCalls) + } + + target := newTrustedTarget(t, channels.ChannelTelegram, "options", "12345") + for name, config := range map[string]Config{ + "negative worker": {BotToken: "token", Target: target, Dispatcher: &dispatchStub{}, Workers: -1}, + "short poll": {BotToken: "token", Target: target, Dispatcher: &dispatchStub{}, PollTimeout: time.Second}, + "http api": {BotToken: "token", Target: target, Dispatcher: &dispatchStub{}, APIBaseURL: "http://insecure.example"}, + } { + t.Run(name, func(t *testing.T) { + if _, err := New(context.Background(), config); !errors.Is(err, ErrInvalid) { + t.Fatalf("invalid option was accepted: %v", err) + } + }) + } +} + +func TestBotFactoryContractsAndSDKOptions(t *testing.T) { + var nilFactory BotFactoryFunc + if _, err := nilFactory.New("token", BotFactoryConfig{}); !errors.Is(err, ErrInvalid) { + t.Fatalf("nil BotFactoryFunc returned unexpected error: %v", err) + } + called := false + expected := &fakeBot{} + factory := BotFactoryFunc(func(token string, config BotFactoryConfig) (BotClient, error) { + called = token == "runtime-token" && config.Workers == 2 + return expected, nil + }) + client, err := factory.New("runtime-token", BotFactoryConfig{Workers: 2}) + if err != nil || client != expected || !called { + t.Fatalf("BotFactoryFunc did not delegate: client=%v err=%v called=%v", client, err, called) + } + + handler := func(context.Context, *bot.Bot, *models.Update) {} + transport := &http.Client{} + if configuredHTTPClient(transport, time.Second) != transport || configuredHTTPClient(nil, time.Second) == nil { + t.Fatal("configured HTTP client did not preserve or create a client") + } + sdkClient, err := (sdkBotFactory{}).New("runtime-token", BotFactoryConfig{ + Handler: handler, APIBaseURL: "https://api.example.test", HTTPClient: transport, + PollTimeout: minimumPollTimeout, Workers: 1, OnPollingError: func() {}, + }) + if err != nil { + t.Fatalf("SDK factory rejected valid options: %v", err) + } + if _, ok := sdkClient.(*bot.Bot); !ok { + t.Fatalf("SDK factory returned unexpected client type %T", sdkClient) + } + for name, config := range map[string]BotFactoryConfig{ + "missing handler": {PollTimeout: minimumPollTimeout, Workers: 1}, + "invalid worker": {Handler: handler, PollTimeout: minimumPollTimeout, Workers: 0}, + "short timeout": {Handler: handler, PollTimeout: minimumPollTimeout - time.Nanosecond, Workers: 1}, + } { + t.Run(name, func(t *testing.T) { + if _, err := (sdkBotFactory{}).New("runtime-token", config); !errors.Is(err, ErrInvalid) { + t.Fatalf("invalid SDK factory options were accepted: %v", err) + } + }) + } +} + +func TestNewRedactsConstructionFailuresAndRejectsPreconditions(t *testing.T) { + target := newTrustedTarget(t, channels.ChannelTelegram, "construction-failures", "12345") + dispatcher := &dispatchStub{} + var nilContext context.Context + if _, err := New(nilContext, Config{BotToken: "token", Target: target, Dispatcher: dispatcher}); !errors.Is(err, ErrInvalid) { + t.Fatalf("nil construction context returned unexpected error: %v", err) + } + canceled, cancel := context.WithCancel(context.Background()) + cancel() + if _, err := New(canceled, Config{BotToken: "token", Target: target, Dispatcher: dispatcher}); !errors.Is(err, context.Canceled) { + t.Fatalf("canceled construction context returned unexpected error: %v", err) + } + for name, config := range map[string]Config{ + "invalid token": {BotToken: " token", Target: target, Dispatcher: dispatcher}, + "missing dispatcher": {BotToken: "token", Target: target}, + "untrusted target": {BotToken: "token", Dispatcher: dispatcher}, + "noncanonical account": {BotToken: "token", Target: newTrustedTarget(t, channels.ChannelTelegram, "noncanonical", "012345"), Dispatcher: dispatcher}, + } { + t.Run(name, func(t *testing.T) { + if _, err := New(context.Background(), config); !errors.Is(err, ErrInvalid) { + t.Fatalf("invalid construction input was accepted: %v", err) + } + }) + } + for name, factory := range map[string]BotFactory{ + "factory error": BotFactoryFunc(func(string, BotFactoryConfig) (BotClient, error) { + return nil, errors.New("provider token=secret") + }), + "nil client": BotFactoryFunc(func(string, BotFactoryConfig) (BotClient, error) { + return nil, nil + }), + "getMe error": &fakeFactory{client: &fakeBot{meErr: errors.New("provider token=secret")}}, + "missing bot identity": &fakeFactory{client: &fakeBot{}}, + "not a bot": &fakeFactory{client: &fakeBot{me: &models.User{ID: 12345}}}, + } { + t.Run(name, func(t *testing.T) { + recorder := &errorRecorder{} + _, err := New(context.Background(), Config{ + BotToken: "runtime-token", Target: target, Dispatcher: dispatcher, Factory: factory, + ErrorHook: recorder.hook, + }) + want := ErrInitialization + if name == "missing bot identity" || name == "not a bot" { + want = ErrBotIdentityMismatch + } + if !errors.Is(err, want) || strings.Contains(err.Error(), "secret") { + t.Fatalf("construction failure was not stable/redacted: err=%v want=%v", err, want) + } + if len(recorder.snapshot()) != 1 { + t.Fatalf("construction failure did not report exactly one hook event: %+v", recorder.snapshot()) + } + }) + } +} + +func TestNewPreservesCancellationDuringGetMe(t *testing.T) { + target := newTrustedTarget(t, channels.ChannelTelegram, "getme-cancel", "12345") + started := make(chan struct{}) + client := &fakeBot{getMeFn: func(ctx context.Context) (*models.User, error) { + close(started) + <-ctx.Done() + return nil, ctx.Err() + }} + factory := &fakeFactory{client: client} + recorder := &errorRecorder{} + ctx, cancel := context.WithCancel(context.Background()) + defer cancel() + result := make(chan error, 1) + go func() { + _, err := New(ctx, Config{ + BotToken: "runtime-token", Target: target, Dispatcher: &dispatchStub{}, Factory: factory, + ErrorHook: recorder.hook, + }) + result <- err + }() + select { + case <-started: + case <-time.After(time.Second): + t.Fatal("getMe did not start") + } + cancel() + if err := <-result; !errors.Is(err, context.Canceled) { + t.Fatalf("getMe cancellation was remapped: %v", err) + } + if events := recorder.snapshot(); len(events) != 0 { + t.Fatalf("getMe cancellation emitted an initialization failure hook: %+v", events) + } +} + +func TestPollingErrorsUseStableRedactedHook(t *testing.T) { + target := newTrustedTarget(t, channels.ChannelTelegram, "polling-hook", "12345") + recorder := &errorRecorder{} + factory := &fakeFactory{client: &fakeBot{me: &models.User{ID: 12345, IsBot: true}}} + adapter, err := New(context.Background(), Config{ + BotToken: "12345:runtime-secret", Target: target, Dispatcher: &dispatchStub{}, Factory: factory, + ErrorHook: recorder.hook, + }) + if err != nil { + t.Fatal(err) + } + defer func() { _ = adapter.Close() }() + factory.config.OnPollingError() + events := recorder.snapshot() + if len(events) != 1 || events[0].Operation != ErrorOperationPolling || !errors.Is(events[0].Err, ErrPolling) { + t.Fatalf("unexpected polling error hook: %+v", events) + } + if strings.Contains(events[0].Err.Error(), "runtime-secret") { + t.Fatal("polling error hook leaked the runtime token") + } +} + +func TestHandleUpdateMapsPrivateTextAndAggregatesDispatchEvents(t *testing.T) { + target := newTrustedTarget(t, channels.ChannelTelegram, "private", "12345") + dispatcher := &dispatchStub{events: []gateway.DispatchEvent{ + {Type: gateway.DispatchEventMessage, Text: "hello "}, + {Type: gateway.DispatchEventStatus, Status: "partial"}, + {Type: gateway.DispatchEventMessage, Text: "world"}, + {Type: gateway.DispatchEventDone, Done: true}, + }} + client := &fakeBot{me: &models.User{ID: 12345, IsBot: true}} + adapter := newTestAdapter(t, target, dispatcher, client) + key := contextKey("request-context") + ctx := context.WithValue(context.Background(), key, "preserved") + update := textUpdate(7, models.ChatTypePrivate, 100, 42, " hello ", 0) + if err := adapter.HandleUpdate(ctx, update); err != nil { + t.Fatal(err) + } + + requests := dispatcher.requests() + if len(requests) != 1 { + t.Fatalf("expected one Dispatch call, got %d", len(requests)) + } + request := requests[0] + if request.Principal.Kind() != gateway.PrincipalChannel || request.Principal.TenantID() != target.TenantID || request.Principal.AppID() != target.AppID { + t.Fatalf("Dispatch did not receive the trusted principal: %+v", request.Principal) + } + if request.Message.Content != "hello" || request.Message.ContentType != gateway.ContentTypeText || request.Message.ExternalUserID != "42" || request.Message.ExternalPeerID != "100" || request.Message.ConversationKind != channels.ConversationDirect { + t.Fatalf("unexpected private inbound message: %+v", request.Message) + } + expectedID := externalMessageID(target, 7) + if request.Message.ExternalMessageID != expectedID || request.RequestID != expectedID { + t.Fatalf("unexpected stable message/request ID: message=%q request=%q expected=%q", request.Message.ExternalMessageID, request.RequestID, expectedID) + } + if got := dispatcher.contextValue(key); got != "preserved" { + t.Fatalf("Dispatch did not receive the caller Context: %v", got) + } + sent := client.sent() + if len(sent) != 1 || sent[0].Text != "hello world" || sent[0].ChatID != 100 || sent[0].ThreadID != 0 { + t.Fatalf("unexpected aggregated Telegram reply: %+v", sent) + } +} + +func TestHandleUpdateMapsGroupThreadAndSplitsUnicodeReply(t *testing.T) { + target := newTrustedTarget(t, channels.ChannelTelegram, "group", "12345") + reply := strings.Repeat("界", maximumReplyRunes) + "🙂" + dispatcher := &dispatchStub{events: []gateway.DispatchEvent{ + {Type: gateway.DispatchEventMessage, Text: reply}, + {Type: gateway.DispatchEventDone, Done: true}, + }} + client := &fakeBot{me: &models.User{ID: 12345, IsBot: true}} + adapter := newTestAdapter(t, target, dispatcher, client) + update := textUpdate(8, models.ChatTypeSupergroup, -100, 42, "group text", 9) + if err := adapter.HandleUpdate(context.Background(), update); err != nil { + t.Fatal(err) + } + request := dispatcher.requests()[0] + if request.Message.ConversationKind != channels.ConversationGroup || request.Message.ExternalChatID != "-100" || request.Message.ExternalThreadID != "9" { + t.Fatalf("unexpected group/thread mapping: %+v", request.Message) + } + sent := client.sent() + if len(sent) != 2 || len([]rune(sent[0].Text)) != maximumReplyRunes || len([]rune(sent[1].Text)) != 1 { + t.Fatalf("reply was not split at Unicode code-point boundary: %+v", sent) + } + for _, message := range sent { + if message.ChatID != -100 || message.ThreadID != 9 { + t.Fatalf("group/thread routing was not preserved: %+v", sent) + } + } +} + +func TestDuplicateDeliveryUsesProcessLocalIdempotency(t *testing.T) { + target := newTrustedTarget(t, channels.ChannelTelegram, "duplicate", "12345") + dispatcher := &dispatchStub{events: []gateway.DispatchEvent{ + {Type: gateway.DispatchEventMessage, Text: "once"}, {Type: gateway.DispatchEventDone, Done: true}, + }} + client := &fakeBot{me: &models.User{ID: 12345, IsBot: true}} + adapter := newTestAdapter(t, target, dispatcher, client) + update := textUpdate(9, models.ChatTypePrivate, 100, 42, "duplicate", 0) + if err := adapter.HandleUpdate(context.Background(), update); err != nil { + t.Fatal(err) + } + if err := adapter.HandleUpdate(context.Background(), update); err != nil { + t.Fatal(err) + } + if len(dispatcher.requests()) != 1 || len(client.sent()) != 2 { + t.Fatalf("completed duplicate did not replay the cached logical reply: dispatch=%d sends=%d", len(dispatcher.requests()), len(client.sent())) + } + + entered := make(chan struct{}) + release := make(chan struct{}) + pendingDispatcher := &dispatchStub{stream: func(ctx context.Context, _ gateway.DispatchRequest) (<-chan gateway.DispatchEvent, error) { + close(entered) + select { + case <-release: + case <-ctx.Done(): + return nil, ctx.Err() + } + return eventStream(gateway.DispatchEvent{Type: gateway.DispatchEventMessage, Text: "pending"}, gateway.DispatchEvent{Type: gateway.DispatchEventDone, Done: true}), nil + }} + pendingClient := &fakeBot{me: &models.User{ID: 12345, IsBot: true}} + pendingAdapter := newTestAdapter(t, target, pendingDispatcher, pendingClient) + firstDone := make(chan error, 1) + go func() { firstDone <- pendingAdapter.HandleUpdate(context.Background(), update) }() + select { + case <-entered: + case <-time.After(time.Second): + t.Fatal("first update did not reach Dispatch") + } + if err := pendingAdapter.HandleUpdate(context.Background(), update); !errors.Is(err, ErrDuplicateUpdate) { + t.Fatalf("pending duplicate was dispatched instead of rejected: %v", err) + } + close(release) + if err := <-firstDone; err != nil { + t.Fatal(err) + } + if len(pendingDispatcher.requests()) != 1 { + t.Fatalf("pending duplicate started more than one dispatch: %d", len(pendingDispatcher.requests())) + } +} + +func TestUnsupportedAndMalformedUpdatesNeverDispatch(t *testing.T) { + target := newTrustedTarget(t, channels.ChannelTelegram, "unsupported", "12345") + dispatcher := &dispatchStub{events: []gateway.DispatchEvent{{Type: gateway.DispatchEventDone, Done: true}}} + client := &fakeBot{me: &models.User{ID: 12345, IsBot: true}} + adapter := newTestAdapter(t, target, dispatcher, client) + valid := textUpdate(10, models.ChatTypePrivate, 100, 42, "valid", 0) + privateCommand := textUpdate(101, models.ChatTypePrivate, 100, 42, "/start", 0) + privateCommand.Message.Entities = []models.MessageEntity{{Type: models.MessageEntityTypeBotCommand, Offset: 0, Length: 6}} + groupCommand := textUpdate(102, models.ChatTypeGroup, -100, 42, "/help@bot", 0) + groupCommand.Message.Entities = []models.MessageEntity{{Type: models.MessageEntityTypeBotCommand, Offset: 0, Length: 10}} + cases := []struct { + name string + update *models.Update + want error + }{ + {name: "nil", update: nil, want: ErrInvalidUpdate}, + {name: "empty update", update: &models.Update{ID: 11}, want: ErrUnsupportedUpdate}, + {name: "edited", update: &models.Update{ID: 12, EditedMessage: valid.Message}, want: ErrUnsupportedUpdate}, + {name: "callback", update: &models.Update{ID: 13, CallbackQuery: &models.CallbackQuery{}}, want: ErrUnsupportedUpdate}, + {name: "media only", update: &models.Update{ID: 14, Message: &models.Message{Chat: models.Chat{ID: 100, Type: models.ChatTypePrivate}, From: &models.User{ID: 42}, Photo: []models.PhotoSize{{}}}}, want: ErrUnsupportedUpdate}, + {name: "service with text", update: &models.Update{ID: 141, Message: &models.Message{Text: "system", Chat: models.Chat{ID: 100, Type: models.ChatTypePrivate}, From: &models.User{ID: 42}, NewChatTitle: "renamed"}}, want: ErrUnsupportedUpdate}, + {name: "private command", update: privateCommand, want: ErrUnsupportedUpdate}, + {name: "group command", update: groupCommand, want: ErrUnsupportedUpdate}, + {name: "missing sender", update: &models.Update{ID: 15, Message: &models.Message{Text: "x", Chat: models.Chat{ID: 100, Type: models.ChatTypePrivate}}}, want: ErrInvalidUpdate}, + {name: "channel", update: textUpdate(16, models.ChatTypeChannel, 100, 42, "x", 0), want: ErrUnsupportedUpdate}, + } + for _, test := range cases { + t.Run(test.name, func(t *testing.T) { + if err := adapter.HandleUpdate(context.Background(), test.update); !errors.Is(err, test.want) { + t.Fatalf("unexpected rejection: got=%v want=%v", err, test.want) + } + }) + } + if len(dispatcher.requests()) != 0 || len(client.sent()) != 0 { + t.Fatalf("unsupported updates reached Gateway or Telegram: dispatch=%d sends=%d", len(dispatcher.requests()), len(client.sent())) + } +} + +func TestDispatchAndSendFailuresAreRedacted(t *testing.T) { + target := newTrustedTarget(t, channels.ChannelTelegram, "failure", "12345") + recorder := &errorRecorder{} + dispatcher := &dispatchStub{err: errors.New("provider token=secret stack=private")} + client := &fakeBot{me: &models.User{ID: 12345, IsBot: true}} + adapter := newTestAdapterWithHook(t, target, dispatcher, client, recorder.hook) + update := textUpdate(17, models.ChatTypePrivate, 100, 42, "fail", 0) + if err := adapter.HandleUpdate(context.Background(), update); !errors.Is(err, ErrDispatch) || strings.Contains(err.Error(), "secret") { + t.Fatalf("dispatch failure was not redacted: %v", err) + } + sent := client.sent() + if len(sent) != 1 || sent[0].Text != failureReply || strings.Contains(sent[0].Text, "secret") { + t.Fatalf("dispatch failure did not produce a fixed reply: %+v", sent) + } + if events := recorder.snapshot(); len(events) != 1 || events[0].Operation != ErrorOperationDispatch || !errors.Is(events[0].Err, ErrDispatch) { + t.Fatalf("unexpected dispatch failure hook: %+v", events) + } + + sendRecorder := &errorRecorder{} + sendDispatcher := &dispatchStub{events: []gateway.DispatchEvent{{Type: gateway.DispatchEventMessage, Text: "reply"}, {Type: gateway.DispatchEventDone, Done: true}}} + sendClient := &fakeBot{me: &models.User{ID: 12345, IsBot: true}, sendErr: errors.New("provider token=secret")} + sendAdapter := newTestAdapterWithHook(t, target, sendDispatcher, sendClient, sendRecorder.hook) + if err := sendAdapter.HandleUpdate(context.Background(), update); !errors.Is(err, ErrSendMessage) || strings.Contains(err.Error(), "secret") { + t.Fatalf("send failure was not redacted: %v", err) + } + if events := sendRecorder.snapshot(); len(events) != 1 || events[0].Operation != ErrorOperationSend || !errors.Is(events[0].Err, ErrSendMessage) { + t.Fatalf("unexpected send failure hook: %+v", events) + } +} + +func TestRunCancellationAndCloseStopPolling(t *testing.T) { + target := newTrustedTarget(t, channels.ChannelTelegram, "lifecycle", "12345") + started := make(chan struct{}) + client := &fakeBot{me: &models.User{ID: 12345, IsBot: true}, startFn: func(ctx context.Context) { + close(started) + <-ctx.Done() + }} + adapter := newTestAdapter(t, target, &dispatchStub{}, client) + runDone := make(chan error, 1) + go func() { runDone <- adapter.Run(context.Background()) }() + select { + case <-started: + case <-time.After(time.Second): + t.Fatal("Run did not start polling") + } + if err := adapter.Close(); err != nil { + t.Fatal(err) + } + select { + case err := <-runDone: + if err != nil { + t.Fatalf("Run returned an unexpected error after Close: %v", err) + } + case <-time.After(time.Second): + t.Fatal("Close did not cancel polling") + } + if err := adapter.Close(); err != nil { + t.Fatal(err) + } + if err := adapter.HandleUpdate(context.Background(), textUpdate(18, models.ChatTypePrivate, 100, 42, "closed", 0)); !errors.Is(err, ErrClosed) { + t.Fatalf("closed adapter accepted an update: %v", err) + } +} + +func TestRunAndHandleUpdatePreconditions(t *testing.T) { + var nilAdapter *Adapter + if err := nilAdapter.Run(context.Background()); !errors.Is(err, ErrNotReady) { + t.Fatalf("nil adapter Run returned unexpected error: %v", err) + } + if err := nilAdapter.HandleUpdate(context.Background(), nil); !errors.Is(err, ErrNotReady) { + t.Fatalf("nil adapter HandleUpdate returned unexpected error: %v", err) + } + if err := nilAdapter.Close(); err != nil { + t.Fatalf("nil adapter Close returned unexpected error: %v", err) + } + + bare := &Adapter{} + var nilContext context.Context + if err := bare.Run(nilContext); !errors.Is(err, ErrInvalid) { + t.Fatalf("nil Run context returned unexpected error: %v", err) + } + if err := bare.Run(context.Background()); !errors.Is(err, ErrNotReady) { + t.Fatalf("bare adapter Run returned unexpected error: %v", err) + } + if err := bare.HandleUpdate(context.Background(), nil); !errors.Is(err, ErrNotReady) { + t.Fatalf("bare adapter HandleUpdate returned unexpected error: %v", err) + } + + target := newTrustedTarget(t, channels.ChannelTelegram, "preconditions", "12345") + adapter := newTestAdapter(t, target, &dispatchStub{events: []gateway.DispatchEvent{{Type: gateway.DispatchEventDone, Done: true}}}, &fakeBot{me: &models.User{ID: 12345, IsBot: true}}) + if err := adapter.HandleUpdate(nilContext, textUpdate(19, models.ChatTypePrivate, 100, 42, "text", 0)); !errors.Is(err, ErrInvalid) { + t.Fatalf("nil HandleUpdate context returned unexpected error: %v", err) + } + canceled, cancel := context.WithCancel(context.Background()) + cancel() + if err := adapter.HandleUpdate(canceled, textUpdate(20, models.ChatTypePrivate, 100, 42, "text", 0)); !errors.Is(err, context.Canceled) { + t.Fatalf("canceled HandleUpdate context returned unexpected error: %v", err) + } + if err := adapter.Close(); err != nil { + t.Fatal(err) + } +} + +func TestRunRejectsConcurrentStart(t *testing.T) { + target := newTrustedTarget(t, channels.ChannelTelegram, "concurrent-run", "12345") + started := make(chan struct{}) + client := &fakeBot{me: &models.User{ID: 12345, IsBot: true}, startFn: func(ctx context.Context) { + close(started) + <-ctx.Done() + }} + adapter := newTestAdapter(t, target, &dispatchStub{}, client) + runDone := make(chan error, 1) + go func() { runDone <- adapter.Run(context.Background()) }() + select { + case <-started: + case <-time.After(time.Second): + t.Fatal("Run did not start") + } + if err := adapter.Run(context.Background()); !errors.Is(err, ErrAlreadyRunning) { + t.Fatalf("concurrent Run returned unexpected error: %v", err) + } + if err := adapter.Close(); err != nil { + t.Fatal(err) + } + if err := <-runDone; err != nil { + t.Fatalf("first Run returned unexpected error: %v", err) + } +} + +func TestDispatchTerminalAndReplyHelpers(t *testing.T) { + message := gateway.InboundMessage{Content: "text", ContentType: gateway.ContentTypeText, ExternalMessageID: "message-1", ExternalUserID: "42", ConversationKind: channels.ConversationDirect, ExternalPeerID: "100"} + for name, stub := range map[string]*dispatchStub{ + "nil stream": {stream: func(context.Context, gateway.DispatchRequest) (<-chan gateway.DispatchEvent, error) { return nil, nil }}, + "incomplete": {events: []gateway.DispatchEvent{{Type: gateway.DispatchEventMessage, Text: "partial"}}}, + "event error": {events: []gateway.DispatchEvent{{Type: gateway.DispatchEventError, Error: "redacted"}, {Type: gateway.DispatchEventDone, Done: true}}}, + "provider error": {err: errors.New("provider token=secret")}, + } { + t.Run(name, func(t *testing.T) { + adapter := &Adapter{dispatcher: stub} + if _, err := adapter.dispatch(context.Background(), message); !errors.Is(err, ErrDispatch) { + t.Fatalf("dispatch returned unexpected error: %v", err) + } + }) + } + contextError := &Adapter{dispatcher: &dispatchStub{err: context.Canceled}} + if _, err := contextError.dispatch(context.Background(), message); !errors.Is(err, context.Canceled) { + t.Fatalf("context dispatch error was not preserved: %v", err) + } + waiting := make(chan gateway.DispatchEvent) + canceled, cancel := context.WithCancel(context.Background()) + cancel() + cancelAdapter := &Adapter{dispatcher: &dispatchStub{stream: func(context.Context, gateway.DispatchRequest) (<-chan gateway.DispatchEvent, error) { + return waiting, nil + }}} + if _, err := cancelAdapter.dispatch(canceled, message); !errors.Is(err, context.Canceled) { + t.Fatalf("dispatch cancellation was not preserved: %v", err) + } + + target := newTrustedTarget(t, channels.ChannelTelegram, "reply-helpers", "12345") + client := &fakeBot{me: &models.User{ID: 12345, IsBot: true}} + adapter := newTestAdapter(t, target, &dispatchStub{}, client) + tgMessage := &models.Message{Chat: models.Chat{ID: 100, Type: models.ChatTypePrivate}} + if err := adapter.sendEvents(context.Background(), tgMessage, []gateway.DispatchEvent{{Type: gateway.DispatchEventError}}); err != nil { + t.Fatalf("error event did not send fixed reply: %v", err) + } + if err := adapter.sendEvents(context.Background(), tgMessage, []gateway.DispatchEvent{{Type: gateway.DispatchEventStatus, Status: "empty"}}); err != nil { + t.Fatalf("empty event stream produced unexpected reply error: %v", err) + } + if err := adapter.sendText(context.Background(), nil, "text"); !errors.Is(err, ErrInvalidUpdate) { + t.Fatalf("nil outbound message returned unexpected error: %v", err) + } + if err := (*Adapter)(nil).sendText(context.Background(), tgMessage, "text"); !errors.Is(err, ErrNotReady) { + t.Fatalf("nil outbound adapter returned unexpected error: %v", err) + } + canceledSend, cancelSend := context.WithCancel(context.Background()) + cancelSend() + if err := adapter.sendText(canceledSend, tgMessage, "text"); !errors.Is(err, context.Canceled) { + t.Fatalf("canceled outbound context returned unexpected error: %v", err) + } + if got := splitText("", maximumReplyRunes); got != nil { + t.Fatalf("empty reply produced chunks: %#v", got) + } +} + +func TestSDKHandlerUsesSingleUpdatePath(t *testing.T) { + target := newTrustedTarget(t, channels.ChannelTelegram, "sdk-handler", "12345") + client := &fakeBot{me: &models.User{ID: 12345, IsBot: true}} + dispatcher := &dispatchStub{events: []gateway.DispatchEvent{{Type: gateway.DispatchEventMessage, Text: "reply"}, {Type: gateway.DispatchEventDone, Done: true}}} + factory := &fakeFactory{client: client} + adapter, err := New(context.Background(), Config{BotToken: "runtime-token", Target: target, Dispatcher: dispatcher, Factory: factory}) + if err != nil { + t.Fatal(err) + } + defer func() { _ = adapter.Close() }() + factory.config.Handler(context.Background(), nil, textUpdate(21, models.ChatTypePrivate, 100, 42, "input", 0)) + if len(dispatcher.requests()) != 1 || len(client.sent()) != 1 || client.sent()[0].Text != "reply" { + t.Fatalf("SDK handler did not route one update: dispatch=%d sends=%v", len(dispatcher.requests()), client.sent()) + } +} + +func TestBindingIsolationAndStableUnicodeChunking(t *testing.T) { + first := newTrustedTarget(t, channels.ChannelTelegram, "isolation-one", "12345") + second := newTrustedTarget(t, channels.ChannelTelegram, "isolation-two", "12345") + input := channels.IdentityInput{ExternalUserID: "42", Kind: channels.ConversationGroup, ExternalChatID: "-100", ExternalThreadID: "9"} + firstIdentity, err := first.RunnerIdentity(input) + if err != nil { + t.Fatal(err) + } + secondIdentity, err := second.RunnerIdentity(input) + if err != nil { + t.Fatal(err) + } + if firstIdentity.UserID == secondIdentity.UserID || firstIdentity.SessionID == secondIdentity.SessionID || first.TenantID == second.TenantID || first.BindingID == second.BindingID { + t.Fatal("same external IDs crossed tenant or Binding identity boundaries") + } + chunks := splitText(strings.Repeat("界", maximumReplyRunes+1), maximumReplyRunes) + if len(chunks) != 2 || len([]rune(chunks[0])) != maximumReplyRunes || len([]rune(chunks[1])) != 1 { + t.Fatalf("unexpected Unicode chunks: lengths=%d,%d", len([]rune(chunks[0])), len([]rune(chunks[1]))) + } +} + +type contextKey string + +type dispatchStub struct { + mu sync.Mutex + requestsList []gateway.DispatchRequest + contexts []context.Context + events []gateway.DispatchEvent + err error + stream func(context.Context, gateway.DispatchRequest) (<-chan gateway.DispatchEvent, error) +} + +func (stub *dispatchStub) Dispatch(ctx context.Context, request gateway.DispatchRequest) (<-chan gateway.DispatchEvent, error) { + stub.mu.Lock() + stub.requestsList = append(stub.requestsList, request) + stub.contexts = append(stub.contexts, ctx) + err, stream, events := stub.err, stub.stream, append([]gateway.DispatchEvent(nil), stub.events...) + stub.mu.Unlock() + if stream != nil { + return stream(ctx, request) + } + if err != nil { + return nil, err + } + return eventStream(events...), nil +} + +func (stub *dispatchStub) requests() []gateway.DispatchRequest { + stub.mu.Lock() + defer stub.mu.Unlock() + return append([]gateway.DispatchRequest(nil), stub.requestsList...) +} + +func (stub *dispatchStub) contextValue(key contextKey) any { + stub.mu.Lock() + defer stub.mu.Unlock() + if len(stub.contexts) == 0 { + return nil + } + return stub.contexts[0].Value(key) +} + +func eventStream(events ...gateway.DispatchEvent) <-chan gateway.DispatchEvent { + stream := make(chan gateway.DispatchEvent, len(events)) + for _, event := range events { + stream <- event + } + close(stream) + return stream +} + +type fakeFactory struct { + client *fakeBot + token string + config BotFactoryConfig +} + +func (factory *fakeFactory) New(token string, config BotFactoryConfig) (BotClient, error) { + factory.token = token + factory.config = config + return factory.client, nil +} + +type fakeBot struct { + mu sync.Mutex + me *models.User + meErr error + getMeFn func(context.Context) (*models.User, error) + sendErr error + sends []sentMessage + startFn func(context.Context) +} + +type sentMessage struct { + ChatID int64 + ThreadID int + Text string +} + +func (client *fakeBot) Start(ctx context.Context) { + if client.startFn != nil { + client.startFn(ctx) + return + } + <-ctx.Done() +} + +func (client *fakeBot) GetMe(ctx context.Context) (*models.User, error) { + if err := ctx.Err(); err != nil { + return nil, err + } + if client.getMeFn != nil { + return client.getMeFn(ctx) + } + return client.me, client.meErr +} + +func (client *fakeBot) SendMessage(ctx context.Context, params *bot.SendMessageParams) (*models.Message, error) { + if err := ctx.Err(); err != nil { + return nil, err + } + client.mu.Lock() + defer client.mu.Unlock() + if client.sendErr != nil { + return nil, client.sendErr + } + chatID, ok := params.ChatID.(int64) + if !ok { + return nil, fmt.Errorf("unexpected fake chat ID type %T", params.ChatID) + } + client.sends = append(client.sends, sentMessage{ChatID: chatID, ThreadID: params.MessageThreadID, Text: params.Text}) + return &models.Message{}, nil +} + +func (client *fakeBot) sent() []sentMessage { + client.mu.Lock() + defer client.mu.Unlock() + return append([]sentMessage(nil), client.sends...) +} + +type errorRecorder struct { + mu sync.Mutex + events []ErrorEvent +} + +func (recorder *errorRecorder) hook(event ErrorEvent) { + recorder.mu.Lock() + defer recorder.mu.Unlock() + recorder.events = append(recorder.events, event) +} + +func (recorder *errorRecorder) snapshot() []ErrorEvent { + recorder.mu.Lock() + defer recorder.mu.Unlock() + return append([]ErrorEvent(nil), recorder.events...) +} + +func newTestAdapter(t *testing.T, target channels.RoutingTarget, dispatcher gateway.DispatchService, client *fakeBot) *Adapter { + return newTestAdapterWithHook(t, target, dispatcher, client, nil) +} + +func newTestAdapterWithHook(t *testing.T, target channels.RoutingTarget, dispatcher gateway.DispatchService, client *fakeBot, hook ErrorHook) *Adapter { + t.Helper() + adapter, err := New(context.Background(), Config{BotToken: "12345:runtime-secret", Target: target, Dispatcher: dispatcher, Factory: &fakeFactory{client: client}, ErrorHook: hook}) + if err != nil { + t.Fatal(err) + } + return adapter +} + +func textUpdate(updateID int64, chatType models.ChatType, chatID, userID int64, text string, threadID int) *models.Update { + return &models.Update{ID: updateID, Message: &models.Message{ + ID: int(updateID), MessageThreadID: threadID, From: &models.User{ID: userID}, + Chat: models.Chat{ID: chatID, Type: chatType}, Text: text, + }} +} + +func newTrustedTarget(t *testing.T, channel channels.Channel, tenantKey, providerAccountID string) channels.RoutingTarget { + t.Helper() + root, snapshot, app := activeTenantApp(t, tenantKey) + repo := inmemory.NewInMemoryRepository() + routeDigest, err := channels.DigestPublicRouteKey(channel, "route-"+tenantKey) + if err != nil { + t.Fatal(err) + } + protocol := channels.ProtocolConfiguration{} + if channel == channels.ChannelTelegram { + protocol.Telegram = &channels.TelegramProtocolConfiguration{WebhookPath: "/reserved"} + } else { + protocol.WeCom = &channels.WeComProtocolConfiguration{CorpID: providerAccountID, ReceiveID: "receive"} + } + binding, _, err := repo.Create(context.Background(), channels.CreateInput{ + TenantID: root.TenantID, BindingKey: "binding-" + tenantKey, Channel: channel, + ProviderAccountID: providerAccountID, PublicRouteKeyDigest: routeDigest, AppID: app.AppID, + SecretRef: "secret/" + tenantKey, Protocol: protocol, Metadata: validMetadata(), + }) + if err != nil { + t.Fatal(err) + } + binding, _, err = repo.Activate(context.Background(), channels.TransitionStatusInput{TenantID: binding.TenantID, BindingID: binding.BindingID, ExpectedVersion: binding.Version, Metadata: validMetadata()}) + if err != nil { + t.Fatal(err) + } + secret := "offline-secret" + resolver := inmemory.NewFakeCandidateResolver(repo, map[channels.SecretScope]string{{TenantID: root.TenantID, SecretRef: binding.SecretRef}: secret}) + candidates, err := repo.LookupCandidates(context.Background(), channel, routeDigest) + if err != nil || len(candidates) != 1 { + t.Fatalf("candidate lookup failed: %d candidates, %v", len(candidates), err) + } + digest := sha256.Sum256([]byte("test message")) + request := channels.VerificationRequest{ + Purpose: channels.PurposeWebhookVerification, Timestamp: time.Now().UTC(), Nonce: "nonce-" + tenantKey, + MessageDigest: hex.EncodeToString(digest[:]), ReceiveID: "receive", + } + request.Signature = inmemory.SignFakeRequest(secret, request) + handle, err := resolver.ResolveCandidate(context.Background(), channels.CandidateSecretRequest{Candidate: candidates[0], Purpose: channels.PurposeWebhookVerification}) + if err != nil { + t.Fatal(err) + } + verified, err := resolver.Verify(context.Background(), handle, request) + if err != nil { + t.Fatal(err) + } + target, err := channels.NewRoutingTarget(snapshot, binding, app, verified) + if err != nil { + t.Fatal(err) + } + return target +} + +func activeTenantApp(t *testing.T, key string) (*tenant.Tenant, tenant.ConfigurationSnapshot, *agent.App) { + t.Helper() + root, err := tenant.NewTenant(tenant.CreateInput{TenantKey: key, DisplayName: "Telegram Test Tenant", AuditRetentionDays: 30, LogMaskingLevel: tenant.MaskingBasic, TraceSamplingRate: 1}) + if err != nil { + t.Fatal(err) + } + snapshot, err := tenant.NewConfigurationSnapshot(root) + if err != nil { + t.Fatal(err) + } + app, err := agent.NewApp(agent.CreateInput{TenantID: root.TenantID, AppKey: "support", DisplayName: "Support", Description: "offline Telegram test"}) + if err != nil { + t.Fatal(err) + } + revision := int64(1) + app.Status = agent.StatusActive + app.CurrentRevision = &revision + app.Version = 2 + app.UpdatedAt = app.CreatedAt.Add(time.Second) + if err := app.Validate(); err != nil { + t.Fatal(err) + } + return root, snapshot, app +} + +func validMetadata() channels.ChangeMetadata { + return channels.ChangeMetadata{ActorType: "test", ActorID: "telegram", Reason: "telegram adapter test", CorrelationID: "telegram-test-correlation"} +}