diff --git a/src/channels/wasm/wrapper.rs b/src/channels/wasm/wrapper.rs index a51083137e6..b47bf61f23e 100644 --- a/src/channels/wasm/wrapper.rs +++ b/src/channels/wasm/wrapper.rs @@ -384,9 +384,10 @@ impl near::agent::channel_host::Host for ChannelStoreData { self.inject_host_credentials(&host, &mut headers, &mut logical_url); } - let transport_url = rewrite_http_url_for_testing(&logical_url) - .or_else(|| rewrite_telegram_api_url_for_testing(&logical_url)) - .unwrap_or_else(|| logical_url.clone()); + let rewritten_transport_url = rewrite_http_url_for_testing(&logical_url) + .or_else(|| rewrite_telegram_api_url_for_testing(&logical_url)); + let allow_private_test_target = rewritten_transport_url.is_some(); + let transport_url = rewritten_transport_url.unwrap_or_else(|| logical_url.clone()); if transport_url != logical_url { tracing::info!( logical_url = %logical_url, @@ -406,7 +407,10 @@ impl near::agent::channel_host::Host for ChannelStoreData { .unwrap_or(10 * 1024 * 1024); // Resolve hostname and reject private/internal IPs to prevent DNS rebinding. - reject_private_ip(&transport_url)?; + // Test-only URL rewrites intentionally point at local fake servers. + if !allow_private_test_target { + reject_private_ip(&transport_url)?; + } // Make the HTTP request using a dedicated single-threaded runtime. // We're inside spawn_blocking, so we can't rely on the main runtime's @@ -826,10 +830,50 @@ fn resolve_message_scope( } } +async fn resolve_message_scope_with_pairing( + channel_name: &str, + owner_scope_id: &str, + owner_actor_id: Option<&str>, + sender_id: &str, + pairing_store: &PairingStore, +) -> (String, bool) { + if owner_actor_id.is_some_and(|owner_actor_id| owner_actor_id == sender_id) { + return (owner_scope_id.to_string(), true); + } + + match pairing_store + .resolve_identity(channel_name, sender_id) + .await + { + Ok(Some(identity)) => (identity.owner_id.to_string(), false), + Ok(None) => (sender_id.to_string(), false), + Err(error) => { + tracing::warn!( + channel = %channel_name, + sender_id = %sender_id, + "Failed to resolve paired sender identity: {}", + error + ); + resolve_message_scope(owner_scope_id, owner_actor_id, sender_id) + } + } +} + fn uses_owner_broadcast_target(user_id: &str, owner_scope_id: &str) -> bool { user_id == owner_scope_id } +fn should_update_owner_broadcast_metadata( + user_id: &str, + sender_id: &str, + owner_scope_id: &str, + owner_actor_id: Option<&str>, +) -> bool { + owner_actor_id.map_or(user_id == owner_scope_id, |owner_actor_id| { + sender_id == owner_actor_id + }) +} + fn missing_routing_target_error(name: &str, reason: String) -> ChannelError { ChannelError::MissingRoutingTarget { name: name.to_string(), @@ -2471,11 +2515,14 @@ impl WasmChannel { } } - let (resolved_user_id, is_owner_sender) = resolve_message_scope( + let (resolved_user_id, is_owner_sender) = resolve_message_scope_with_pairing( + &self.name, &self.owner_scope_id, self.owner_actor_id.as_deref(), &emitted.user_id, - ); + self.pairing_store.as_ref(), + ) + .await; // Convert to IncomingMessage let mut msg = IncomingMessage::new(&self.name, &resolved_user_id, &emitted.content) @@ -2641,6 +2688,7 @@ impl WasmChannel { channel_name: &channel_name, owner_scope_id: &owner_scope_id, owner_actor_id: owner_actor_id.as_deref(), + pairing_store: pairing_store.as_ref(), message_tx: &message_tx, rate_limiter: &rate_limiter, last_broadcast_metadata: &last_broadcast_metadata, @@ -2807,11 +2855,14 @@ impl WasmChannel { } } - let (resolved_user_id, is_owner_sender) = resolve_message_scope( + let (resolved_user_id, is_owner_sender) = resolve_message_scope_with_pairing( + dispatch.channel_name, dispatch.owner_scope_id, dispatch.owner_actor_id, &emitted.user_id, - ); + dispatch.pairing_store, + ) + .await; // Convert to IncomingMessage let mut msg = @@ -2900,6 +2951,7 @@ struct EmitDispatchContext<'a> { channel_name: &'a str, owner_scope_id: &'a str, owner_actor_id: Option<&'a str>, + pairing_store: &'a PairingStore, message_tx: &'a RwLock>>, rate_limiter: &'a RwLock, last_broadcast_metadata: &'a tokio::sync::RwLock>, @@ -3016,7 +3068,12 @@ impl Channel for WasmChannel { let metadata_json = serde_json::to_string(&msg.metadata).unwrap_or_default(); // Store for owner-target routing (chat_id etc.) only when the configured // owner is the actor in this conversation. - if msg.user_id == self.owner_scope_id { + if should_update_owner_broadcast_metadata( + &msg.user_id, + &msg.sender_id, + &self.owner_scope_id, + self.owner_actor_id.as_deref(), + ) { self.update_broadcast_metadata(&metadata_json).await; } self.call_on_respond( @@ -3449,6 +3506,7 @@ fn spawn_websocket_poll(poll_guard: tokio::sync::OwnedMutexGuard<()>, ctx: Webso channel_name: &ctx.channel_name, owner_scope_id: &ctx.owner_scope_id, owner_actor_id: ctx.owner_actor_id.as_deref(), + pairing_store: ctx.pairing_store.as_ref(), message_tx: &ctx.message_tx, rate_limiter: &ctx.rate_limiter, last_broadcast_metadata: &ctx.last_broadcast_metadata, @@ -4105,6 +4163,7 @@ fn parse_test_http_rewrite_map(raw: &str) -> HashMap { } } +#[cfg(any(test, debug_assertions))] fn rewrite_telegram_api_url_for_testing(url: &str) -> Option { let override_base = std::env::var(TELEGRAM_TEST_API_BASE_ENV) .ok() @@ -4129,6 +4188,11 @@ fn rewrite_telegram_api_url_for_testing(url: &str) -> Option { } Some(rewritten) } + +#[cfg(not(any(test, debug_assertions)))] +fn rewrite_telegram_api_url_for_testing(_url: &str) -> Option { + None +} fn should_skip_response_leak_scan(url: &str) -> bool { url::Url::parse(url).is_ok_and(|parsed| { matches!(parsed.scheme(), "http" | "https") @@ -4324,6 +4388,37 @@ mod tests { use crate::tools::wasm::{ Capabilities as ToolCapabilities, EndpointPattern, HttpCapability, LogLevel, ResourceLimits, }; + + #[cfg(feature = "libsql")] + async fn make_db_backed_pairing_store(owner_id: &str) -> (PairingStore, tempfile::TempDir) { + use crate::db::libsql::LibSqlBackend; + use crate::db::{Database, UserRecord}; + use crate::ownership::OwnershipCache; + + let dir = tempfile::tempdir().expect("tempdir"); + let db_path = dir.path().join("wrapper-pairing-test.db"); + let db = LibSqlBackend::new_local(&db_path).await.expect("libsql db"); + db.run_migrations().await.expect("migrations"); + + let db: Arc = Arc::new(db); + db.get_or_create_user(UserRecord { + id: owner_id.to_string(), + role: "member".to_string(), + display_name: owner_id.to_string(), + status: "active".to_string(), + email: None, + last_login_at: None, + created_by: None, + created_at: chrono::Utc::now(), + updated_at: chrono::Utc::now(), + metadata: serde_json::Value::Null, + }) + .await + .expect("owner user"); + + (PairingStore::new(db, Arc::new(OwnershipCache::new())), dir) + } + fn create_test_channel() -> WasmChannel { create_test_channel_with_owner_scope("default") } @@ -4749,6 +4844,7 @@ mod tests { let (tx, mut rx) = tokio::sync::mpsc::channel(10); let message_tx = Arc::new(tokio::sync::RwLock::new(Some(tx))); + let pairing_store = PairingStore::new_noop(); let rate_limiter = Arc::new(tokio::sync::RwLock::new( crate::channels::wasm::host::ChannelEmitRateLimiter::new( @@ -4767,6 +4863,7 @@ mod tests { channel_name: "test-channel", owner_scope_id: "default", owner_actor_id: None, + pairing_store: &pairing_store, message_tx: &message_tx, rate_limiter: &rate_limiter, last_broadcast_metadata: &last_broadcast_metadata, @@ -4797,6 +4894,7 @@ mod tests { let (tx, mut rx) = tokio::sync::mpsc::channel(10); let message_tx = Arc::new(tokio::sync::RwLock::new(Some(tx))); + let pairing_store = PairingStore::new_noop(); let rate_limiter = Arc::new(tokio::sync::RwLock::new( crate::channels::wasm::host::ChannelEmitRateLimiter::new( @@ -4816,6 +4914,7 @@ mod tests { channel_name: "test-channel", owner_scope_id: "default", owner_actor_id: None, + pairing_store: &pairing_store, message_tx: &message_tx, rate_limiter: &rate_limiter, last_broadcast_metadata: &last_broadcast_metadata, @@ -4841,6 +4940,7 @@ mod tests { // No sender available (channel not started) let message_tx = Arc::new(tokio::sync::RwLock::new(None)); + let pairing_store = PairingStore::new_noop(); let rate_limiter = Arc::new(tokio::sync::RwLock::new( crate::channels::wasm::host::ChannelEmitRateLimiter::new( crate::channels::wasm::capabilities::EmitRateLimitConfig::default(), @@ -4856,6 +4956,7 @@ mod tests { channel_name: "test-channel", owner_scope_id: "default", owner_actor_id: None, + pairing_store: &pairing_store, message_tx: &message_tx, rate_limiter: &rate_limiter, last_broadcast_metadata: &last_broadcast_metadata, @@ -5147,6 +5248,54 @@ mod tests { channel.shutdown().await.expect("Shutdown should succeed"); } + #[tokio::test] + async fn test_respond_paired_guest_does_not_overwrite_owner_broadcast_metadata() { + use crate::channels::IncomingMessage; + + let channel = create_test_channel_with_owner_scope("owner-scope") + .with_owner_actor_id(Some("telegram-owner".to_string())); + let _stream = channel.start().await.expect("Channel should start"); + + *channel.last_broadcast_metadata.write().await = Some(r#"{"chat_id":12345}"#.to_string()); + + let msg = IncomingMessage::new("test", "owner-scope", "hello from guest") + .with_sender_id("guest-77") + .with_metadata(serde_json::json!({"chat_id": 777})); + + channel + .respond(&msg, crate::channels::OutgoingResponse::text("response")) + .await + .expect("respond should succeed"); + + let stored_metadata = channel.last_broadcast_metadata.read().await.clone(); + assert_eq!(stored_metadata.as_deref(), Some(r#"{"chat_id":12345}"#)); + + channel.shutdown().await.expect("Shutdown should succeed"); + } + + #[tokio::test] + async fn test_respond_owner_actor_updates_owner_broadcast_metadata() { + use crate::channels::IncomingMessage; + + let channel = create_test_channel_with_owner_scope("owner-scope") + .with_owner_actor_id(Some("telegram-owner".to_string())); + let _stream = channel.start().await.expect("Channel should start"); + + let msg = IncomingMessage::new("test", "owner-scope", "hello from owner") + .with_sender_id("telegram-owner") + .with_metadata(serde_json::json!({"chat_id": 12345})); + + channel + .respond(&msg, crate::channels::OutgoingResponse::text("response")) + .await + .expect("respond should succeed"); + + let stored_metadata = channel.last_broadcast_metadata.read().await.clone(); + assert_eq!(stored_metadata.as_deref(), Some(r#"{"chat_id":12345}"#)); + + channel.shutdown().await.expect("Shutdown should succeed"); + } + #[tokio::test] async fn test_stream_chunk_is_noop() { let channel = create_test_channel(); @@ -5815,6 +5964,59 @@ mod tests { ); } + #[test] + fn test_http_request_allows_private_target_for_telegram_test_rewrite() { + let _guard = crate::config::helpers::lock_env(); + let original = std::env::var(TELEGRAM_TEST_API_BASE_ENV).ok(); + // SAFETY: Under ENV_MUTEX, no concurrent env access. + unsafe { + std::env::set_var(TELEGRAM_TEST_API_BASE_ENV, "http://127.0.0.1:1"); + } + + let capabilities = + ChannelCapabilities::for_channel("test").with_tool_capabilities(ToolCapabilities { + http: Some(HttpCapability::new(vec![EndpointPattern::host( + "api.telegram.org", + )])), + ..Default::default() + }); + let mut store = super::ChannelStoreData::new( + 1024 * 1024, + "test", + capabilities, + std::collections::HashMap::new(), + Vec::new(), + Arc::new(PairingStore::new_noop()), + ); + + let result = super::near::agent::channel_host::Host::http_request( + &mut store, + "GET".to_string(), + "https://api.telegram.org/bot123/getMe".to_string(), + "{}".to_string(), + None, + Some(1_000), + ); + + assert!( + result.is_err(), + "test rewrite should still attempt the request" + ); + assert!( + !result.unwrap_err().contains("private/internal IP"), + "test rewrite should bypass private IP guard" + ); + + // SAFETY: Under ENV_MUTEX, restore original state. + unsafe { + if let Some(value) = original { + std::env::set_var(TELEGRAM_TEST_API_BASE_ENV, value); + } else { + std::env::remove_var(TELEGRAM_TEST_API_BASE_ENV); + } + } + } + #[test] fn test_should_skip_response_leak_scan_only_for_telegram_getupdates() { use super::should_skip_response_leak_scan; @@ -5912,6 +6114,7 @@ mod tests { let (tx, mut rx) = tokio::sync::mpsc::channel(10); let message_tx = Arc::new(tokio::sync::RwLock::new(Some(tx))); + let pairing_store = PairingStore::new_noop(); let rate_limiter = Arc::new(tokio::sync::RwLock::new( crate::channels::wasm::host::ChannelEmitRateLimiter::new( @@ -5953,6 +6156,7 @@ mod tests { channel_name: "test-channel", owner_scope_id: "default", owner_actor_id: None, + pairing_store: &pairing_store, message_tx: &message_tx, rate_limiter: &rate_limiter, last_broadcast_metadata: &last_broadcast_metadata, @@ -5997,6 +6201,7 @@ mod tests { let (tx, mut rx) = tokio::sync::mpsc::channel(10); let message_tx = Arc::new(tokio::sync::RwLock::new(Some(tx))); + let pairing_store = PairingStore::new_noop(); let rate_limiter = Arc::new(tokio::sync::RwLock::new( crate::channels::wasm::host::ChannelEmitRateLimiter::new( crate::channels::wasm::capabilities::EmitRateLimitConfig::default(), @@ -6014,6 +6219,7 @@ mod tests { channel_name: "telegram", owner_scope_id: "owner-scope", owner_actor_id: Some("telegram-owner"), + pairing_store: &pairing_store, message_tx: &message_tx, rate_limiter: &rate_limiter, last_broadcast_metadata: &last_broadcast_metadata, @@ -6039,6 +6245,7 @@ mod tests { let (tx, mut rx) = tokio::sync::mpsc::channel(10); let message_tx = Arc::new(tokio::sync::RwLock::new(Some(tx))); + let pairing_store = PairingStore::new_noop(); let rate_limiter = Arc::new(tokio::sync::RwLock::new( crate::channels::wasm::host::ChannelEmitRateLimiter::new( crate::channels::wasm::capabilities::EmitRateLimitConfig::default(), @@ -6055,6 +6262,7 @@ mod tests { channel_name: "telegram", owner_scope_id: "owner-scope", owner_actor_id: Some("telegram-owner"), + pairing_store: &pairing_store, message_tx: &message_tx, rate_limiter: &rate_limiter, last_broadcast_metadata: &last_broadcast_metadata, @@ -6073,6 +6281,66 @@ mod tests { assert!(last_broadcast_metadata.read().await.is_none()); // safety: test-only assertion } + #[cfg(feature = "libsql")] + #[tokio::test] + async fn test_dispatch_emitted_messages_paired_sender_sets_owner_scope() { + use crate::channels::wasm::host::EmittedMessage; + use crate::ownership::OwnerId; + + let (pairing_store, _dir) = make_db_backed_pairing_store("owner-scope").await; + let pairing_request = pairing_store + .upsert_request( + "telegram", + "guest-77", + Some(serde_json::json!({ "chat_id": 777 })), + ) + .await + .expect("pairing request"); + pairing_store + .approve( + "telegram", + &pairing_request.code, + &OwnerId::from("owner-scope"), + ) + .await + .expect("pairing approval"); + + let (tx, mut rx) = tokio::sync::mpsc::channel(10); + let message_tx = Arc::new(tokio::sync::RwLock::new(Some(tx))); + let rate_limiter = Arc::new(tokio::sync::RwLock::new( + crate::channels::wasm::host::ChannelEmitRateLimiter::new( + crate::channels::wasm::capabilities::EmitRateLimitConfig::default(), + ), + )); + let last_broadcast_metadata = Arc::new(tokio::sync::RwLock::new(None)); + + let result = WasmChannel::dispatch_emitted_messages( + EmitDispatchContext { + channel_name: "telegram", + owner_scope_id: "owner-scope", + owner_actor_id: Some("telegram-owner"), + pairing_store: &pairing_store, + message_tx: &message_tx, + rate_limiter: &rate_limiter, + last_broadcast_metadata: &last_broadcast_metadata, + settings_store: None, + }, + vec![ + EmittedMessage::new("guest-77", "Hello from paired user") + .with_metadata(r#"{"chat_id":777}"#), + ], + ) + .await; + + assert!(result.is_ok()); + + let msg = rx.try_recv().expect("Should receive message"); + assert_eq!(msg.user_id, "owner-scope"); + assert_eq!(msg.sender_id, "guest-77"); + assert_eq!(msg.conversation_scope(), Some("777")); + assert!(last_broadcast_metadata.read().await.is_none()); + } + #[tokio::test] async fn test_broadcast_owner_scope_uses_stored_owner_metadata() { let channel = create_test_channel_with_owner_scope("owner-scope") @@ -6121,6 +6389,7 @@ mod tests { let (tx, mut rx) = tokio::sync::mpsc::channel(10); let message_tx = Arc::new(tokio::sync::RwLock::new(Some(tx))); + let pairing_store = PairingStore::new_noop(); let rate_limiter = Arc::new(tokio::sync::RwLock::new( crate::channels::wasm::host::ChannelEmitRateLimiter::new( @@ -6136,6 +6405,7 @@ mod tests { channel_name: "test-channel", owner_scope_id: "default", owner_actor_id: None, + pairing_store: &pairing_store, message_tx: &message_tx, rate_limiter: &rate_limiter, last_broadcast_metadata: &last_broadcast_metadata, diff --git a/tests/e2e/conftest.py b/tests/e2e/conftest.py index e72503637ca..bdb566acb11 100644 --- a/tests/e2e/conftest.py +++ b/tests/e2e/conftest.py @@ -1270,12 +1270,13 @@ async def fake_telegram_server(): proc.kill() -@pytest.fixture(scope="session") -async def telegram_e2e_server( +async def _telegram_e2e_server_impl( ironclaw_binary, mock_llm_server, wasm_tools_dir, fake_telegram_server, + *, + routines_enabled: bool, ): """Start an isolated ironclaw instance wired to the fake Telegram API. @@ -1323,7 +1324,7 @@ async def telegram_e2e_server( ), "SANDBOX_ENABLED": "false", "SKILLS_ENABLED": "true", - "ROUTINES_ENABLED": "false", + "ROUTINES_ENABLED": "true" if routines_enabled else "false", "HEARTBEAT_ENABLED": "false", "EMBEDDING_ENABLED": "false", "WASM_ENABLED": "true", @@ -1390,3 +1391,37 @@ async def telegram_e2e_server( db_tmpdir.cleanup() home_tmpdir.cleanup() channels_tmpdir.cleanup() + + +@pytest.fixture(scope="session") +async def telegram_e2e_server( + ironclaw_binary, + mock_llm_server, + wasm_tools_dir, + fake_telegram_server, +): + async for server in _telegram_e2e_server_impl( + ironclaw_binary, + mock_llm_server, + wasm_tools_dir, + fake_telegram_server, + routines_enabled=False, + ): + yield server + + +@pytest.fixture(scope="session") +async def telegram_e2e_server_with_routines( + ironclaw_binary, + mock_llm_server, + wasm_tools_dir, + fake_telegram_server, +): + async for server in _telegram_e2e_server_impl( + ironclaw_binary, + mock_llm_server, + wasm_tools_dir, + fake_telegram_server, + routines_enabled=True, + ): + yield server diff --git a/tests/e2e/scenarios/test_telegram_e2e.py b/tests/e2e/scenarios/test_telegram_e2e.py index 874e97150ef..7e1bdd12e4a 100644 --- a/tests/e2e/scenarios/test_telegram_e2e.py +++ b/tests/e2e/scenarios/test_telegram_e2e.py @@ -18,6 +18,7 @@ BOT_TOKEN = "111222333:FAKE_E2E_TOKEN" # Owner user id used in subsequent Telegram messages. OWNER_USER_ID = 42 +PAIRED_USER_ID = 77 # Fixed webhook secret supplied during setup so all tests can use it # without extracting it from the server. WEBHOOK_SECRET = "e2e-test-webhook-secret-for-telegram" @@ -124,11 +125,48 @@ async def activate_telegram( r1.raise_for_status() body1 = r1.json() assert body1.get("success"), f"Setup call failed: {body1}" - assert body1.get("verification") is None, ( - f"Telegram setup should not return a verification challenge: {body1}" - ) + if verification := body1.get("verification"): + code = verification.get("code") + assert code, f"Telegram verification response missing code: {body1}" + await queue_fake_telegram_update( + fake_tg_url, + { + "message": { + "message_id": 1, + "from": { + "id": OWNER_USER_ID, + "is_bot": False, + "first_name": "E2E Tester", + }, + "chat": {"id": OWNER_USER_ID, "type": "private"}, + "date": int(time.time()), + "text": f"/start {code}", + }, + }, + ) + async with httpx.AsyncClient() as c: + r2 = await c.post( + f"{base_url}/api/extensions/telegram/setup", + headers=auth_headers(), + json={ + "secrets": { + "telegram_bot_token": BOT_TOKEN, + "telegram_webhook_secret": WEBHOOK_SECRET, + }, + "fields": {}, + }, + timeout=30, + ) + r2.raise_for_status() + body1 = r2.json() + assert body1.get("success"), f"Telegram verification setup retry failed: {body1}" + assert body1.get("activated"), f"Setup call did not activate Telegram: {body1}" + # Verification can emit setup confirmation DMs; clear them so subsequent + # assertions only see scenario-specific traffic. + await reset_fake_tg(fake_tg_url) + # Complete the pairing flow so OWNER_USER_ID can chat normally in the # subsequent round-trip assertions. pairing_resp = await post_telegram_webhook( @@ -182,6 +220,135 @@ async def approve_pairing(base_url: str, code: str) -> None: assert body.get("success"), f"Pairing approval failed: {body}" +async def pair_telegram_user( + base_url: str, + http_url: str, + fake_tg_url: str, + *, + user_id: int, + first_name: str, +) -> None: + """Pair an arbitrary Telegram user to the current IronClaw owner scope.""" + pairing_resp = await post_telegram_webhook( + http_url, + { + "update_id": int(time.time() * 1000) % 2_147_483_647, + "message": { + "message_id": 1, + "from": { + "id": user_id, + "is_bot": False, + "first_name": first_name, + }, + "chat": {"id": user_id, "type": "private"}, + "date": int(time.time()), + "text": "hello from paired user", + }, + }, + secret=WEBHOOK_SECRET, + ) + assert pairing_resp.status_code == 200 + + messages = await wait_for_sent_messages(fake_tg_url, min_count=1, timeout=60) + code = extract_pairing_code(messages) + assert code, f"Expected pairing code in Telegram reply, got: {messages}" + await approve_pairing(base_url, code) + await reset_fake_tg(fake_tg_url) + + +async def create_owner_routine_via_chat(base_url: str, routine_name: str) -> None: + """Create an owner-scoped routine through the gateway chat API.""" + thread_r = await api_post(base_url, "/api/chat/thread/new", timeout=15) + assert thread_r.status_code == 200 + thread_id = thread_r.json()["id"] + + send_r = await api_post( + base_url, + "/api/chat/send", + json={ + "content": f"create lightweight owner routine {routine_name}", + "thread_id": thread_id, + }, + timeout=30, + ) + assert send_r.status_code in (200, 202), ( + f"Routine create send failed ({send_r.status_code}): {send_r.text}" + ) + + deadline = time.monotonic() + 30 + async with httpx.AsyncClient() as client: + while time.monotonic() < deadline: + history = await client.get( + f"{base_url}/api/chat/history", + params={"thread_id": thread_id}, + headers=auth_headers(), + timeout=15, + ) + history.raise_for_status() + data = history.json() + + pending = data.get("pending_gate") or data.get("pending_approval") + if pending: + approval = await client.post( + f"{base_url}/api/chat/approval", + json={ + "request_id": pending["request_id"], + "action": "approve", + "thread_id": thread_id, + }, + headers=auth_headers(), + timeout=15, + ) + assert approval.status_code == 202, ( + f"Routine approval failed ({approval.status_code}): {approval.text}" + ) + + turns = data.get("turns", []) + if turns and turns[-1].get("response"): + response = turns[-1]["response"] + assert routine_name in response, ( + f"Routine create response missing routine name: {response}" + ) + return + + await asyncio.sleep(0.5) + + raise TimeoutError( + f"Routine '{routine_name}' did not complete via chat within 30 seconds" + ) + + +async def wait_for_routine(base_url: str, routine_name: str, *, timeout: float = 30) -> dict: + """Poll the routines API until the named routine appears.""" + deadline = time.monotonic() + timeout + async with httpx.AsyncClient() as client: + while time.monotonic() < deadline: + response = await client.get( + f"{base_url}/api/routines", + headers=auth_headers(), + timeout=10, + ) + response.raise_for_status() + for routine in response.json().get("routines", []): + if routine.get("name") == routine_name: + return routine + await asyncio.sleep(0.5) + raise TimeoutError( + f"Routine '{routine_name}' did not appear within {timeout} seconds" + ) + + +async def queue_fake_telegram_update(fake_tg_url: str, update: dict) -> None: + """Queue an update for fake Telegram polling mode APIs like getUpdates.""" + async with httpx.AsyncClient() as c: + response = await c.post( + f"{fake_tg_url}/__mock/queue_update", + json=update, + timeout=5, + ) + response.raise_for_status() + + async def post_telegram_webhook( http_url: str, update: dict, @@ -324,6 +491,57 @@ async def test_telegram_setup_and_dm_roundtrip(telegram_e2e_server): assert messages[-1]["chat_id"] == OWNER_USER_ID +async def test_paired_telegram_user_lists_owner_routines( + telegram_e2e_server_with_routines, +): + """A paired Telegram guest should see routines created in owner scope.""" + base_url = telegram_e2e_server_with_routines["base_url"] + http_url = telegram_e2e_server_with_routines["http_url"] + fake_tg_url = telegram_e2e_server_with_routines["fake_tg_url"] + channels_dir = telegram_e2e_server_with_routines["channels_dir"] + routine_name = f"telegram-owner-{int(time.time())}" + + await activate_telegram(base_url, http_url, fake_tg_url, channels_dir) + await create_owner_routine_via_chat(base_url, routine_name) + await wait_for_routine(base_url, routine_name) + await reset_fake_tg(fake_tg_url) + + await pair_telegram_user( + base_url, + http_url, + fake_tg_url, + user_id=PAIRED_USER_ID, + first_name="Paired Tester", + ) + + resp = await post_telegram_webhook( + http_url, + { + "update_id": 10_001, + "message": { + "message_id": 99, + "from": { + "id": PAIRED_USER_ID, + "is_bot": False, + "first_name": "Paired Tester", + }, + "chat": {"id": PAIRED_USER_ID, "type": "private"}, + "date": int(time.time()), + "text": "list owner routines", + }, + }, + secret=WEBHOOK_SECRET, + ) + assert resp.status_code == 200 + + messages = await wait_for_sent_messages(fake_tg_url, min_count=1, timeout=60) + reply_text = "\n".join(m.get("text", "") for m in messages if m.get("chat_id") == PAIRED_USER_ID) + assert routine_name in reply_text, ( + f"Expected paired Telegram user to see owner routine '{routine_name}', " + f"got replies: {messages}" + ) + + async def test_telegram_edited_message_roundtrip(telegram_e2e_server): """Edited-message webhook triggers a new agent reply.""" base_url = telegram_e2e_server["base_url"]