1pub mod approvals;
26pub mod bridge;
27pub mod platform;
28pub mod platforms;
29pub mod render;
30
31pub use bridge::{ConnectBridge, ConnectContext, SessionKey};
32pub use platform::{
33 Button, CallbackQuery, Capabilities, Inbound, InboundMessage, MessageRef, OutboundMessage,
34 Platform, PlatformError, PlatformResult, ReplyCtx,
35};
36
37use std::path::PathBuf;
38use std::sync::Arc;
39
40use tokio::sync::mpsc;
41
42use bamboo_config::ConnectPlatformConfig;
43use bamboo_llm::Config;
44
45pub struct ConnectManager {
51 tasks: Vec<tokio::task::JoinHandle<()>>,
52}
53
54impl ConnectManager {
55 pub async fn start(
60 ctx: ConnectContext,
61 config_snapshot: &Config,
62 data_dir: Option<PathBuf>,
63 ) -> Self {
64 let map_path = data_dir.map(|dir| dir.join("connect_sessions.json"));
65 let bridge = Arc::new(ConnectBridge::new(ctx, map_path));
66 bridge.load_session_map().await;
67
68 let mut tasks = Vec::new();
69 let start_ok = multi_bot_guard(&config_snapshot.connect.platforms);
70 for (index, platform_cfg) in config_snapshot.connect.platforms.iter().enumerate() {
71 match platform_cfg.platform_type.as_str() {
72 "telegram" => {
73 let token = platform_cfg.token.clone().unwrap_or_default();
74 if token.trim().is_empty() {
75 tracing::warn!(
76 "connect: telegram platform configured without a token; skipping"
77 );
78 continue;
79 }
80 if !start_ok[index] {
91 tracing::warn!(
92 "connect: multiple telegram platform entries are configured; only \
93 the FIRST is started. A second telegram bot on this instance would \
94 collide with the first on the same session-routing key \
95 (`telegram:<chat_id>:<user_id>`), silently mixing sessions for any \
96 user who messages both bots. Remove the extra entry, or track \
97 issue #454 for per-bot session keys."
98 );
99 continue;
100 }
101 if platform_cfg.allow_from.is_empty() {
102 tracing::warn!(
103 "connect: telegram platform has an EMPTY allow_from list — every \
104 inbound message will be denied until you add allowed user ids to \
105 connect.platforms[].allow_from"
106 );
107 }
108
109 let platform: Arc<dyn Platform> =
110 Arc::new(platforms::telegram::TelegramPlatform::new(token));
111 spawn_platform_tasks(
112 &mut tasks,
113 &bridge,
114 platform,
115 platform_cfg.allow_from.clone(),
116 );
117 }
118 "feishu" => {
119 let app_id = platform_cfg.app_id.clone().unwrap_or_default();
120 let app_secret = platform_cfg.app_secret.clone().unwrap_or_default();
121 if app_id.trim().is_empty() || app_secret.trim().is_empty() {
122 tracing::warn!(
123 "connect: feishu platform configured without app_id/app_secret; \
124 skipping"
125 );
126 continue;
127 }
128 if !start_ok[index] {
133 tracing::warn!(
134 "connect: multiple feishu platform entries are configured; only the \
135 FIRST is started. A second feishu app on this instance would \
136 collide with the first on the same session-routing key \
137 (`feishu:<chat_id>:<open_id>`) and on inbound message dedup. \
138 Remove the extra entry, or track issue #454 for per-bot session \
139 keys."
140 );
141 continue;
142 }
143 let Some(base_url) = resolve_feishu_base_url(platform_cfg.domain.as_deref())
144 else {
145 tracing::warn!(
146 domain = platform_cfg.domain.as_deref().unwrap_or_default(),
147 "connect: feishu platform has an invalid domain (expected \"feishu\", \
148 \"lark\", or an https:// base URL); skipping"
149 );
150 continue;
151 };
152 if platform_cfg.allow_from.is_empty() {
153 tracing::warn!(
154 "connect: feishu platform has an EMPTY allow_from list — every \
155 inbound message will be denied until you add allowed open_ids to \
156 connect.platforms[].allow_from"
157 );
158 }
159
160 let platform: Arc<dyn Platform> = Arc::new(
161 platforms::feishu::FeishuPlatform::new(app_id, app_secret, base_url),
162 );
163 spawn_platform_tasks(
164 &mut tasks,
165 &bridge,
166 platform,
167 platform_cfg.allow_from.clone(),
168 );
169 }
170 other => {
171 tracing::warn!("connect: unknown platform type '{other}'; skipping");
172 }
173 }
174 }
175
176 Self { tasks }
177 }
178}
179
180fn spawn_platform_tasks(
185 tasks: &mut Vec<tokio::task::JoinHandle<()>>,
186 bridge: &Arc<ConnectBridge>,
187 platform: Arc<dyn Platform>,
188 allow_from: Vec<String>,
189) {
190 let name = platform.name().to_string();
191 let (tx, rx) = mpsc::channel(64);
192
193 let platform_for_start = platform.clone();
194 let name_for_start = name.clone();
195 tasks.push(tokio::spawn(async move {
196 if let Err(error) = platform_for_start.start(tx).await {
197 tracing::warn!("connect: {name_for_start} platform loop exited: {error}");
198 }
199 }));
200
201 tasks.push(tokio::spawn(dispatch_loop(
202 bridge.clone(),
203 platform,
204 allow_from,
205 rx,
206 )));
207
208 tracing::info!("connect: started {name} platform");
209}
210
211fn resolve_feishu_base_url(domain: Option<&str>) -> Option<String> {
217 match domain.map(str::trim).filter(|d| !d.is_empty()) {
218 None | Some("feishu") => Some("https://open.feishu.cn".to_string()),
219 Some("lark") => Some("https://open.larksuite.com".to_string()),
220 Some(custom) if custom.starts_with("https://") => {
221 Some(custom.trim_end_matches('/').to_string())
222 }
223 Some(_) => None,
224 }
225}
226
227impl Drop for ConnectManager {
228 fn drop(&mut self) {
229 for task in &self.tasks {
230 task.abort();
231 }
232 }
233}
234
235pub(crate) fn multi_bot_guard(platforms: &[ConnectPlatformConfig]) -> Vec<bool> {
260 let mut seen_valid: std::collections::HashSet<&str> = std::collections::HashSet::new();
261 platforms
262 .iter()
263 .map(|platform_cfg| {
264 let non_empty =
265 |field: &Option<String>| field.as_deref().is_some_and(|v| !v.trim().is_empty());
266 let valid = match platform_cfg.platform_type.as_str() {
267 "telegram" => non_empty(&platform_cfg.token),
268 "feishu" => non_empty(&platform_cfg.app_id) && non_empty(&platform_cfg.app_secret),
269 _ => return true,
270 };
271 if !valid {
272 return true;
273 }
274 seen_valid.insert(platform_cfg.platform_type.as_str())
275 })
276 .collect()
277}
278
279pub(crate) fn platform_config_will_start(
285 platform: &ConnectPlatformConfig,
286 guard_allows: bool,
287) -> bool {
288 let non_empty = |field: &Option<String>| {
289 field
290 .as_deref()
291 .is_some_and(|value| !value.trim().is_empty())
292 };
293 match platform.platform_type.as_str() {
294 "telegram" => guard_allows && non_empty(&platform.token),
295 "feishu" => {
296 guard_allows
297 && non_empty(&platform.app_id)
298 && non_empty(&platform.app_secret)
299 && resolve_feishu_base_url(platform.domain.as_deref()).is_some()
300 }
301 _ => false,
302 }
303}
304
305async fn dispatch_loop(
313 bridge: Arc<ConnectBridge>,
314 platform: Arc<dyn Platform>,
315 allow_from: Vec<String>,
316 mut rx: mpsc::Receiver<Inbound>,
317) {
318 while let Some(event) = rx.recv().await {
319 match event {
320 Inbound::Message(msg) => {
321 ConnectBridge::handle_inbound(
322 bridge.clone(),
323 platform.clone(),
324 allow_from.clone(),
325 msg,
326 )
327 .await;
328 }
329 Inbound::Callback(callback) => {
330 ConnectBridge::handle_callback(
331 bridge.clone(),
332 platform.clone(),
333 allow_from.clone(),
334 callback,
335 )
336 .await;
337 }
338 }
339 }
340}
341
342#[cfg(test)]
343mod tests {
344 use super::*;
345
346 fn platform(platform_type: &str, token: Option<&str>) -> ConnectPlatformConfig {
347 ConnectPlatformConfig {
348 id: None,
349 project_id: None,
350 platform_type: platform_type.to_string(),
351 token: token.map(str::to_string),
352 token_encrypted: None,
353 token_credential_ref: None,
354 token_configured: false,
355 app_id: None,
356 app_secret: None,
357 app_secret_encrypted: None,
358 app_secret_credential_ref: None,
359 app_secret_configured: false,
360 domain: None,
361 allow_from: Vec::new(),
362 admin_from: Vec::new(),
363 }
364 }
365
366 fn feishu_platform(app_id: Option<&str>, app_secret: Option<&str>) -> ConnectPlatformConfig {
367 ConnectPlatformConfig {
368 app_id: app_id.map(str::to_string),
369 app_secret: app_secret.map(str::to_string),
370 ..platform("feishu", None)
371 }
372 }
373
374 #[test]
375 fn multi_bot_guard_allows_a_single_telegram_entry() {
376 let platforms = vec![platform("telegram", Some("tok-1"))];
377 assert_eq!(multi_bot_guard(&platforms), vec![true]);
378 }
379
380 #[test]
381 fn multi_bot_guard_rejects_every_telegram_entry_after_the_first() {
382 let platforms = vec![
383 platform("telegram", Some("tok-1")),
384 platform("telegram", Some("tok-2")),
385 platform("telegram", Some("tok-3")),
386 ];
387 assert_eq!(multi_bot_guard(&platforms), vec![true, false, false]);
388 }
389
390 #[test]
393 fn multi_bot_guard_is_scoped_per_platform_type() {
394 let platforms = vec![
395 platform("telegram", Some("tok-1")),
396 feishu_platform(Some("cli_a"), Some("secret-a")),
397 platform("telegram", Some("tok-2")),
398 feishu_platform(Some("cli_b"), Some("secret-b")),
399 ];
400 assert_eq!(multi_bot_guard(&platforms), vec![true, true, false, false]);
401 }
402
403 #[test]
408 fn multi_bot_guard_does_not_count_a_credentialless_entry_against_the_budget() {
409 let platforms = vec![
410 platform("telegram", None),
411 platform("telegram", Some("")),
412 platform("telegram", Some("tok-real")),
413 feishu_platform(Some("cli_a"), None),
414 feishu_platform(Some("cli_b"), Some("secret-real")),
415 ];
416 assert_eq!(
417 multi_bot_guard(&platforms),
418 vec![true, true, true, true, true]
419 );
420 }
421
422 #[test]
423 fn multi_bot_guard_leaves_unknown_platform_types_alone() {
424 let platforms = vec![
425 platform("dingtalk", Some("x")),
426 platform("dingtalk", Some("y")),
427 ];
428 assert_eq!(multi_bot_guard(&platforms), vec![true, true]);
429 }
430
431 #[test]
432 fn multi_bot_guard_handles_an_empty_platform_list() {
433 assert_eq!(multi_bot_guard(&[]), Vec::<bool>::new());
434 }
435
436 #[test]
437 fn resolve_feishu_base_url_covers_the_three_domain_forms() {
438 assert_eq!(
439 resolve_feishu_base_url(None).as_deref(),
440 Some("https://open.feishu.cn")
441 );
442 assert_eq!(
443 resolve_feishu_base_url(Some("feishu")).as_deref(),
444 Some("https://open.feishu.cn")
445 );
446 assert_eq!(
447 resolve_feishu_base_url(Some("")).as_deref(),
448 Some("https://open.feishu.cn")
449 );
450 assert_eq!(
451 resolve_feishu_base_url(Some("lark")).as_deref(),
452 Some("https://open.larksuite.com")
453 );
454 assert_eq!(
455 resolve_feishu_base_url(Some("https://feishu.example.corp/")).as_deref(),
456 Some("https://feishu.example.corp")
457 );
458 assert_eq!(
459 resolve_feishu_base_url(Some("http://insecure.example")),
460 None
461 );
462 assert_eq!(resolve_feishu_base_url(Some("dingtalk")), None);
463 }
464}