diff --git a/src/crates/assembly/core/src/service/remote_connect/bot/weixin.rs b/src/crates/assembly/core/src/service/remote_connect/bot/weixin.rs index ffcf3926c5..80ded71f63 100644 --- a/src/crates/assembly/core/src/service/remote_connect/bot/weixin.rs +++ b/src/crates/assembly/core/src/service/remote_connect/bot/weixin.rs @@ -372,9 +372,7 @@ impl WeixinBot { let ret = resp["ret"].as_i64().unwrap_or(0); let errcode = resp["errcode"].as_i64().unwrap_or(0); - if (ret != 0 && ret != WEIXIN_SESSION_EXPIRED_ERRCODE) - || (errcode != 0 && errcode != WEIXIN_SESSION_EXPIRED_ERRCODE) - { + if weixin_provider::updates_failed(&resp) { if errcode == WEIXIN_SESSION_EXPIRED_ERRCODE || ret == WEIXIN_SESSION_EXPIRED_ERRCODE { @@ -386,7 +384,10 @@ impl WeixinBot { continue; } - if let Some(new_buf) = resp["get_updates_buf"].as_str() { + if let Some(new_buf) = resp["get_updates_buf"] + .as_str() + .filter(|buf| !buf.is_empty()) + { buf = new_buf.to_string(); weixin_provider::save_sync_buf(&self.api.config().bot_account_id, &buf); } @@ -489,9 +490,7 @@ impl WeixinBot { let ret = resp["ret"].as_i64().unwrap_or(0); let errcode = resp["errcode"].as_i64().unwrap_or(0); - if (ret != 0 && ret != WEIXIN_SESSION_EXPIRED_ERRCODE) - || (errcode != 0 && errcode != WEIXIN_SESSION_EXPIRED_ERRCODE) - { + if weixin_provider::updates_failed(&resp) { if errcode == WEIXIN_SESSION_EXPIRED_ERRCODE || ret == WEIXIN_SESSION_EXPIRED_ERRCODE { @@ -502,7 +501,10 @@ impl WeixinBot { continue; } - if let Some(new_buf) = resp["get_updates_buf"].as_str() { + if let Some(new_buf) = resp["get_updates_buf"] + .as_str() + .filter(|buf| !buf.is_empty()) + { buf = new_buf.to_string(); weixin_provider::save_sync_buf(&self.api.config().bot_account_id, &buf); } diff --git a/src/crates/services/services-integrations/src/remote_connect/bot/weixin.rs b/src/crates/services/services-integrations/src/remote_connect/bot/weixin.rs index 5c306b475e..f86d2416d6 100644 --- a/src/crates/services/services-integrations/src/remote_connect/bot/weixin.rs +++ b/src/crates/services/services-integrations/src/remote_connect/bot/weixin.rs @@ -6,7 +6,7 @@ use aes::cipher::{BlockDecrypt, BlockEncrypt, KeyInit}; use aes::Aes128; -use anyhow::{anyhow, Result}; +use anyhow::{anyhow, Context, Result}; use base64::{engine::general_purpose::STANDARD as B64, Engine as _}; use log::{debug, warn}; use rand::{Rng, RngCore}; @@ -15,18 +15,19 @@ use serde_json::{json, Value}; use std::collections::HashMap; use std::path::{Path, PathBuf}; use std::sync::{Arc, Mutex, OnceLock}; -use std::time::{Duration, SystemTime, UNIX_EPOCH}; +use std::time::{Duration, Instant, SystemTime, UNIX_EPOCH}; use tokio::sync::RwLock; const DEFAULT_BASE_URL: &str = "https://ilinkai.weixin.qq.com"; const DEFAULT_ILINK_BOT_TYPE: &str = "3"; -// Wire compatibility baseline audited against @tencent-weixin/openclaw-weixin 2.4.6. -const CHANNEL_VERSION: &str = "2.4.6"; +// Wire compatibility baseline audited against @tencent-weixin/openclaw-weixin 2.4.8. +const CHANNEL_VERSION: &str = "2.4.8"; const ILINK_APP_ID: &str = "bot"; -const ILINK_APP_CLIENT_VERSION: &str = "132102"; // 0x00020406 +const ILINK_APP_CLIENT_VERSION: &str = "132104"; // 0x00020408 const BOT_AGENT: &str = concat!("OpenBitFun/", env!("CARGO_PKG_VERSION")); const API_TIMEOUT_SECS: u64 = 20; -const QR_POLL_TIMEOUT_SECS: u64 = 36; +const QR_POLL_TIMEOUT_SECS: u64 = 35; +const QR_SESSION_TTL_MS: i64 = 5 * 60_000; pub const WEIXIN_SESSION_EXPIRED_ERRCODE: i64 = -14; const SESSION_PAUSE_SECS: u64 = 3600; const MAX_TEXT_CHUNK: usize = 4000; @@ -51,7 +52,7 @@ pub struct WeixinQrStartResponse { pub message: String, } -#[derive(Debug, Serialize)] +#[derive(Debug, Clone, Serialize)] #[serde(rename_all = "snake_case")] pub enum WeixinQrPollStatus { Wait, @@ -62,7 +63,7 @@ pub enum WeixinQrPollStatus { Error, } -#[derive(Debug, Serialize)] +#[derive(Debug, Clone, Serialize)] pub struct WeixinQrPollResponse { pub status: WeixinQrPollStatus, pub message: String, @@ -104,6 +105,10 @@ struct QrLoginSession { existing_ilink_token: Option, existing_bot_account_id: Option, existing_base_url: Option, + // Serialize polls and retain a confirmed result until TTL expiry so a lost + // controller response can be retried without losing the one-time login. + poll_lock: Arc>, + completed: Option, } enum QrSessionLookup { @@ -120,6 +125,8 @@ struct QrCodeApiResponse { #[derive(Debug, Deserialize)] struct QrStatusApiResponse { + ret: Option, + errcode: Option, status: Option, bot_token: Option, ilink_bot_id: Option, @@ -129,6 +136,7 @@ struct QrStatusApiResponse { pub struct WeixinProviderClient { config: WeixinConfig, + http_client: reqwest::Client, typing_tickets: Arc>>, session_pause_until_ms: Arc>>, } @@ -188,6 +196,7 @@ impl WeixinProviderClient { pub fn new(config: WeixinConfig) -> Self { Self { config, + http_client: crate::reqwest_client(), typing_tickets: Arc::new(RwLock::new(HashMap::new())), session_pause_until_ms: Arc::new(RwLock::new(HashMap::new())), } @@ -254,9 +263,10 @@ impl WeixinProviderClient { async fn post_ilink(&self, endpoint: &str, body: Value, timeout: Duration) -> Result { let url = format!("{}{}", self.base_url(), endpoint.trim_start_matches('/')); let body_str = serde_json::to_string(&body)?; - let client = crate::reqwest_client_builder().timeout(timeout).build()?; - let resp = client + let resp = self + .http_client .post(&url) + .timeout(timeout) .headers(self.build_auth_headers()) .body(body_str) .send() @@ -266,25 +276,22 @@ impl WeixinProviderClient { if !status.is_success() { return Err(anyhow!("ilink {endpoint} HTTP {status}: {text}")); } - if endpoint.contains("sendmessage") - || endpoint.contains("sendtyping") - || endpoint.contains("getconfig") - || endpoint.contains("notifystart") - || endpoint.contains("notifystop") - { - if let Ok(value) = serde_json::from_str::(&text) { - let ret = value["ret"].as_i64().unwrap_or(0); - let errcode = value["errcode"].as_i64().unwrap_or(0); - if ret != 0 || errcode != 0 { - let errmsg = value["errmsg"] - .as_str() - .or_else(|| value["msg"].as_str()) - .unwrap_or("") - .to_string(); - return Err(anyhow!( - "ilink {endpoint} application error ret={ret} errcode={errcode} errmsg={errmsg}" - )); - } + let value: Value = serde_json::from_str(&text) + .with_context(|| format!("ilink {endpoint}: invalid JSON response"))?; + if !value.is_object() { + return Err(anyhow!("ilink {endpoint}: expected a JSON object")); + } + if endpoint != "ilink/bot/getupdates" { + let ret = value["ret"].as_i64().unwrap_or(0); + let errcode = value["errcode"].as_i64().unwrap_or(0); + if ret != 0 || errcode != 0 { + let errmsg = value["errmsg"] + .as_str() + .or_else(|| value["msg"].as_str()) + .unwrap_or(""); + return Err(anyhow!( + "ilink {endpoint} application error ret={ret} errcode={errcode} errmsg={errmsg}" + )); } } Ok(text) @@ -294,9 +301,10 @@ impl WeixinProviderClient { if self.is_session_paused().await { tokio::time::sleep(Duration::from_secs(2)).await; return Ok(json!({ - "ret": 0, - "msgs": [], - "get_updates_buf": buf + "ret": WEIXIN_SESSION_EXPIRED_ERRCODE, + "errcode": WEIXIN_SESSION_EXPIRED_ERRCODE, + "errmsg": "WeChat login token is stale; sign in again", + "msgs": [] })); } @@ -309,13 +317,29 @@ impl WeixinProviderClient { }), timeout, ) - .await?; + .await; + let raw = match raw { + Ok(raw) => raw, + Err(err) + if err + .downcast_ref::() + .is_some_and(|err| err.is_timeout()) => + { + // An empty long poll is normal, not a reason to delay the next poll. + return Ok(json!({ "ret": 0, "msgs": [], "get_updates_buf": buf })); + } + Err(err) => return Err(err), + }; let value: Value = serde_json::from_str(&raw)?; let ret = value["ret"].as_i64().unwrap_or(0); let errcode = value["errcode"].as_i64().unwrap_or(0); if errcode == WEIXIN_SESSION_EXPIRED_ERRCODE || ret == WEIXIN_SESSION_EXPIRED_ERRCODE { self.pause_session().await; } + debug!( + "weixin: getupdates ret={ret} errcode={errcode} messages={}", + value["msgs"].as_array().map_or(0, Vec::len) + ); Ok(value) } @@ -874,6 +898,10 @@ pub async fn weixin_qr_start( .ok_or_else(|| anyhow!("get_bot_qrcode: missing qrcode_img_content"))?; let session_key = uuid::Uuid::new_v4().to_string(); + qr_sessions() + .lock() + .map_err(|err| anyhow!("qr session lock: {err}"))? + .retain(|_, session| now_ms() - session.started_at_ms <= QR_SESSION_TTL_MS); qr_sessions() .lock() .map_err(|err| anyhow!("qr session lock: {err}"))? @@ -888,6 +916,8 @@ pub async fn weixin_qr_start( existing_ilink_token, existing_bot_account_id, existing_base_url, + poll_lock: Arc::new(tokio::sync::Mutex::new(())), + completed: None, }, ); @@ -912,7 +942,7 @@ pub async fn weixin_qr_poll( match sessions.get(session_key) { None => QrSessionLookup::Missing, Some(session) => { - if now_ms() - session.started_at_ms > 5 * 60_000 { + if now_ms() - session.started_at_ms > QR_SESSION_TTL_MS { sessions.remove(session_key); QrSessionLookup::TimedOut } else { @@ -926,6 +956,27 @@ pub async fn weixin_qr_poll( QrSessionLookup::Missing => Ok(qr_error("No active QR session. Start login again.")), QrSessionLookup::TimedOut => Ok(qr_error("QR session expired. Start again.")), QrSessionLookup::Found(session) => { + let poll_lock = session.poll_lock.clone(); + let _poll = poll_lock.lock().await; + // Re-read after waiting for any earlier request to finish. + let current = qr_sessions() + .lock() + .map_err(|err| anyhow!("qr session lock: {err}"))? + .get(session_key) + .cloned(); + let Some(session) = current else { + return Ok(qr_error("No active QR session. Start login again.")); + }; + if now_ms() - session.started_at_ms > QR_SESSION_TTL_MS { + qr_sessions() + .lock() + .map_err(|err| anyhow!("qr session lock: {err}"))? + .remove(session_key); + return Ok(qr_error("QR session expired. Start again.")); + } + if let Some(completed) = session.completed { + return Ok(completed); + } poll_found_qr_session(session_key, session, &base, verify_code.as_deref()).await } } @@ -945,12 +996,10 @@ async fn poll_found_qr_session( url.push_str(&urlencoding::encode(code)); } - let client = crate::reqwest_client_builder() - .timeout(Duration::from_secs(QR_POLL_TIMEOUT_SECS)) - .build()?; - - let resp = client + let started = Instant::now(); + let resp = qr_http_client() .get(&url) + .timeout(Duration::from_secs(QR_POLL_TIMEOUT_SECS)) .headers(build_common_headers()) .send() .await; @@ -959,22 +1008,43 @@ async fn poll_found_qr_session( Ok(resp) => resp, Err(err) => { if err.is_timeout() { + debug!("weixin: QR status long poll timed out; continuing"); return Ok(qr_wait("waiting")); } - warn!("weixin: transient QR status request failed; retrying: {err}"); + warn!( + "weixin: transient QR status request failed; retrying: {}", + err.without_url() + ); return Ok(qr_wait("waiting")); } }; let status = resp.status(); if !status.is_success() { - let body = resp.text().await.unwrap_or_default(); - warn!("weixin: transient QR status HTTP {status}; retrying: {body}"); + warn!("weixin: QR status HTTP {status}"); + if status.is_client_error() { + return Ok(qr_error(&format!( + "WeChat rejected the login request (HTTP {status}). Start login again." + ))); + } return Ok(qr_wait("waiting")); } let status_json: QrStatusApiResponse = resp.json().await?; - match status_json.status.as_deref().unwrap_or("wait") { + let ret = status_json.ret.unwrap_or(0); + let errcode = status_json.errcode.unwrap_or(0); + if ret != 0 || errcode != 0 { + warn!("weixin: QR status application error ret={ret} errcode={errcode}"); + return Ok(qr_error(&format!( + "WeChat login failed (ret={ret}, errcode={errcode}). Start login again." + ))); + } + debug!( + "weixin: QR status={} elapsed_ms={}", + status_json.status.as_deref().unwrap_or("missing"), + started.elapsed().as_millis() + ); + match status_json.status.as_deref().unwrap_or("") { "wait" => Ok(qr_wait("waiting")), "scaned" => Ok(WeixinQrPollResponse { status: WeixinQrPollStatus::Scanned, @@ -990,7 +1060,9 @@ async fn poll_found_qr_session( "binded_redirect" => confirm_existing_qr_session(session_key, session), "confirmed" => confirm_qr_session(session_key, status_json, &base), "expired" => refresh_qr_session(session_key, original_base).await, - other => Ok(qr_wait(other)), + _ => Ok(qr_error( + "WeChat returned an unsupported login state. Start login again.", + )), } } @@ -1037,24 +1109,23 @@ fn confirm_existing_qr_session( .existing_bot_account_id .filter(|value| !value.is_empty()) .ok_or_else(|| anyhow!("binded_redirect received without existing bot account id"))?; - qr_sessions() - .lock() - .map_err(|err| anyhow!("qr session lock: {err}"))? - .remove(session_key); - Ok(WeixinQrPollResponse { - status: WeixinQrPollStatus::Confirmed, - message: "WeChat was already linked; reused the existing login.".to_string(), - qr_image_url: None, - ilink_token: Some(token), - bot_account_id: Some(account_id), - base_url: Some( - session - .existing_base_url - .unwrap_or(session.current_base_url) - .trim_end_matches('/') - .to_string(), - ), - }) + cache_qr_completion( + session_key, + WeixinQrPollResponse { + status: WeixinQrPollStatus::Confirmed, + message: "WeChat was already linked; reused the existing login.".to_string(), + qr_image_url: None, + ilink_token: Some(token), + bot_account_id: Some(account_id), + base_url: Some( + session + .existing_base_url + .unwrap_or(session.current_base_url) + .trim_end_matches('/') + .to_string(), + ), + }, + ) } fn confirm_qr_session( @@ -1076,19 +1147,31 @@ fn confirm_qr_session( .filter(|s| !s.is_empty()) .unwrap_or_else(|| base.trim_end_matches('/').to_string()); - qr_sessions() - .lock() - .map_err(|err| anyhow!("qr session lock: {err}"))? - .remove(session_key); + cache_qr_completion( + session_key, + WeixinQrPollResponse { + status: WeixinQrPollStatus::Confirmed, + message: "WeChat linked.".to_string(), + qr_image_url: None, + ilink_token: Some(token), + bot_account_id: Some(normalized), + base_url: Some(baseurl), + }, + ) +} - Ok(WeixinQrPollResponse { - status: WeixinQrPollStatus::Confirmed, - message: "WeChat linked.".to_string(), - qr_image_url: None, - ilink_token: Some(token), - bot_account_id: Some(normalized), - base_url: Some(baseurl), - }) +fn cache_qr_completion( + session_key: &str, + response: WeixinQrPollResponse, +) -> Result { + let mut sessions = qr_sessions() + .lock() + .map_err(|err| anyhow!("qr session lock: {err}"))?; + let session = sessions + .get_mut(session_key) + .ok_or_else(|| anyhow!("QR session lost before confirmation"))?; + session.completed = Some(response.clone()); + Ok(response) } async fn refresh_qr_session(session_key: &str, base: &str) -> Result { @@ -1231,6 +1314,10 @@ pub fn context_token(msg: &Value) -> Option { .filter(|s| !s.is_empty()) } +pub fn updates_failed(response: &Value) -> bool { + response["ret"].as_i64().unwrap_or(0) != 0 || response["errcode"].as_i64().unwrap_or(0) != 0 +} + pub fn suggested_long_poll_timeout(response: &Value, current: Duration) -> Duration { response["longpolling_timeout_ms"] .as_u64() @@ -1429,6 +1516,11 @@ fn sniff_image_mime(bytes: &[u8]) -> &'static str { "image/jpeg" } +fn qr_http_client() -> &'static reqwest::Client { + static CLIENT: OnceLock = OnceLock::new(); + CLIENT.get_or_init(crate::reqwest_client) +} + fn qr_sessions() -> &'static Mutex> { static CELL: OnceLock>> = OnceLock::new(); CELL.get_or_init(|| Mutex::new(HashMap::new())) @@ -1508,11 +1600,9 @@ async fn fetch_qr_code(base: &str, local_token_list: &[String]) -> Result, + ) -> ( + String, + tokio::sync::mpsc::UnboundedReceiver, + tokio::task::JoinHandle<()>, + ) { + use tokio::io::{AsyncReadExt, AsyncWriteExt}; + let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap(); + let base = format!("http://{}/", listener.local_addr().unwrap()); + let (tx, rx) = tokio::sync::mpsc::unbounded_channel(); + let task = tokio::spawn(async move { + for (delay_ms, body) in replies { + let (mut stream, _) = listener.accept().await.unwrap(); + let mut request = Vec::new(); + let mut chunk = [0; 4096]; + loop { + let read = stream.read(&mut chunk).await.unwrap(); + if read == 0 { + break; + } + request.extend_from_slice(&chunk[..read]); + if let Some(end) = request.windows(4).position(|w| w == b"\r\n\r\n") { + let headers = String::from_utf8_lossy(&request[..end]); + let length = headers + .lines() + .find_map(|line| { + let (key, value) = line.split_once(':')?; + key.eq_ignore_ascii_case("content-length") + .then(|| value.trim().parse::().ok()) + .flatten() + }) + .unwrap_or(0); + if request.len() >= end + 4 + length { + break; + } + } + } + tx.send(String::from_utf8(request).unwrap()).unwrap(); + tokio::time::sleep(Duration::from_millis(delay_ms)).await; + let response = format!( + "HTTP/1.1 200 OK\r\nContent-Type: application/json\r\nContent-Length: {}\r\nConnection: close\r\n\r\n{}", + body.len(), body + ); + let _ = stream.write_all(response.as_bytes()).await; + } + }); + (base, rx, task) + } + + fn test_client(base_url: String) -> WeixinProviderClient { + WeixinProviderClient::new(WeixinConfig { + ilink_token: "fixture-token".into(), + base_url, + bot_account_id: "fixture-bot".into(), + }) + } + + fn test_qr_session(base: &str) -> String { + let key = uuid::Uuid::new_v4().to_string(); + qr_sessions().lock().unwrap().insert( + key.clone(), + QrLoginSession { + qrcode: "fixture-qr".into(), + started_at_ms: now_ms(), + refresh_count: 0, + current_base_url: base.into(), + local_token_list: vec!["fixture-token".into()], + existing_ilink_token: Some("fixture-token".into()), + existing_bot_account_id: Some("fixture-bot".into()), + existing_base_url: Some(base.into()), + poll_lock: Arc::new(tokio::sync::Mutex::new(())), + completed: None, + }, + ); + key + } + + #[tokio::test] + async fn qr_scan_confirmation_is_replayable_after_response_loss() { + let (base, mut requests, task) = mock_ilink(vec![ + (0, r#"{"status":"scaned"}"#), + (0, r#"{"status":"confirmed","bot_token":"new-token","ilink_bot_id":"fixture@bot","baseurl":"https://assigned.example"}"#), + ]).await; + let key = test_qr_session(&base); + let scanned = weixin_qr_poll(&key, None, None).await.unwrap(); + assert!(matches!(scanned.status, WeixinQrPollStatus::Scanned)); + let (confirmed, overlapping) = tokio::join!( + weixin_qr_poll(&key, None, None), + weixin_qr_poll(&key, None, None), + ); + let confirmed = confirmed.unwrap(); + assert_eq!( + serde_json::to_value(&confirmed).unwrap(), + serde_json::to_value(overlapping.unwrap()).unwrap() + ); + assert!(matches!(confirmed.status, WeixinQrPollStatus::Confirmed)); + assert_eq!( + confirmed.base_url.as_deref(), + Some("https://assigned.example") + ); + let replay = weixin_qr_poll(&key, None, None).await.unwrap(); + assert_eq!( + serde_json::to_value(&confirmed).unwrap(), + serde_json::to_value(replay).unwrap() + ); + task.await.unwrap(); + assert!(requests + .recv() + .await + .unwrap() + .starts_with("GET /ilink/bot/get_qrcode_status?qrcode=fixture-qr ")); + assert!(requests.recv().await.is_some()); + assert!(requests.recv().await.is_none()); + qr_sessions().lock().unwrap().remove(&key); + } + + #[tokio::test] + async fn verification_and_reused_login_follow_reference_states() { + let (base, mut requests, task) = mock_ilink(vec![ + (0, r#"{"status":"need_verifycode"}"#), + (0, r#"{"status":"binded_redirect"}"#), + ]) + .await; + let key = test_qr_session(&base); + assert!(matches!( + weixin_qr_poll(&key, None, None).await.unwrap().status, + WeixinQrPollStatus::NeedVerifyCode + )); + let reused = weixin_qr_poll(&key, None, Some("12 34&56".into())) + .await + .unwrap(); + assert!(matches!(reused.status, WeixinQrPollStatus::Confirmed)); + assert_eq!(reused.ilink_token.as_deref(), Some("fixture-token")); + task.await.unwrap(); + requests.recv().await.unwrap(); + assert!(requests + .recv() + .await + .unwrap() + .contains("verify_code=12%2034%2656")); + qr_sessions().lock().unwrap().remove(&key); + } + + #[tokio::test] + async fn expired_qr_refresh_preserves_local_tokens_and_resets_poll_host() { + let (base, mut requests, task) = mock_ilink(vec![( + 0, + r#"{"qrcode":"fresh-code","qrcode_img_content":"https://example.test/scan"}"#, + )]) + .await; + let key = test_qr_session("https://old.example/"); + let refreshed = refresh_qr_session(&key, &base).await.unwrap(); + assert!(matches!(refreshed.status, WeixinQrPollStatus::Expired)); + assert_eq!( + refreshed.qr_image_url.as_deref(), + Some("https://example.test/scan") + ); + let session = qr_sessions().lock().unwrap().remove(&key).unwrap(); + assert_eq!(session.current_base_url, base); + assert_eq!(session.qrcode, "fresh-code"); + task.await.unwrap(); + let request = requests.recv().await.unwrap(); + assert!(request.starts_with("POST /ilink/bot/get_bot_qrcode?bot_type=3 ")); + assert!(request.contains(r#""local_token_list":["fixture-token"]"#)); + } + + #[tokio::test] + async fn missing_or_failed_qr_status_is_not_reported_as_waiting() { + for response in [r#"{"ret":-14}"#, r#"{}"#, r#"{"status":"unsupported"}"#] { + let (base, _requests, task) = mock_ilink(vec![(0, response)]).await; + let key = test_qr_session(&base); + assert!(matches!( + weixin_qr_poll(&key, None, None).await.unwrap().status, + WeixinQrPollStatus::Error + )); + task.await.unwrap(); + qr_sessions().lock().unwrap().remove(&key); + } + } + + #[tokio::test] + async fn long_poll_timeout_preserves_cursor_without_backoff_error() { + let (base, _requests, task) = mock_ilink(vec![(100, r#"{"ret":0,"msgs":[]}"#)]).await; + let response = test_client(base) + .get_updates_once("previous-cursor", Duration::from_millis(20)) + .await + .unwrap(); + assert_eq!(response["get_updates_buf"], "previous-cursor"); + assert!(!updates_failed(&response)); + task.await.unwrap(); + } + + #[tokio::test] + async fn stale_token_remains_an_error_while_requests_are_paused() { + let (base, mut requests, task) = mock_ilink(vec![(0, r#"{"ret":-14}"#)]).await; + let client = test_client(base); + let response = client + .get_updates_once("cursor", Duration::from_secs(1)) + .await + .unwrap(); + assert!(updates_failed(&response)); + let paused = client + .get_updates_once("cursor", Duration::from_secs(1)) + .await + .unwrap(); + assert!(updates_failed(&paused)); + assert!(paused.get("get_updates_buf").is_none()); + task.await.unwrap(); + assert!(requests.recv().await.is_some()); + assert!(requests.recv().await.is_none()); + } + + #[tokio::test] + async fn send_response_must_be_valid_json_and_successful() { + for response in [ + "not-json", + "null", + r#"{"ret":-1,"errmsg":"denied"}"#, + r#"{"errcode":42}"#, + ] { + let (base, _requests, task) = mock_ilink(vec![(0, response)]).await; + assert!(test_client(base) + .send_text_chunks("peer", "context", "hello") + .await + .is_err()); + task.await.unwrap(); + } + } + + #[tokio::test] + async fn outbound_text_uses_inbound_context_and_reference_envelope() { + let (base, mut requests, task) = mock_ilink(vec![(0, "{}")]).await; + test_client(base) + .send_text_chunks("peer", "incoming-context", "hello") + .await + .unwrap(); + task.await.unwrap(); + let request = requests.recv().await.unwrap(); + assert!(request + .to_ascii_lowercase() + .contains("authorization: bearer fixture-token")); + let body: Value = serde_json::from_str(request.split_once("\r\n\r\n").unwrap().1).unwrap(); + assert_eq!(body["msg"]["context_token"], "incoming-context"); + assert_eq!(body["msg"]["message_type"], 2); + assert_eq!(body["msg"]["message_state"], 2); + assert_eq!(body["msg"]["item_list"][0]["text_item"]["text"], "hello"); + assert_eq!(body["base_info"]["channel_version"], "2.4.8"); + } + #[test] fn context_token_error_heuristic() { let app_err = anyhow!( diff --git a/src/web-ui/src/app/components/RemoteConnectDialog/RemoteConnectDialog.scss b/src/web-ui/src/app/components/RemoteConnectDialog/RemoteConnectDialog.scss index be311a859c..c2794c2f4a 100644 --- a/src/web-ui/src/app/components/RemoteConnectDialog/RemoteConnectDialog.scss +++ b/src/web-ui/src/app/components/RemoteConnectDialog/RemoteConnectDialog.scss @@ -960,6 +960,18 @@ text-align: center; } +.openbitfun-remote-connect__weixin-progress { + display: flex; + flex-direction: column; + align-items: center; + gap: var(--openbitfun-space-2); + max-width: 100%; + + .openbitfun-remote-connect__hint { + margin: 0; + } +} + .openbitfun-remote-connect__weixin-verify { display: flex; width: 100%; diff --git a/src/web-ui/src/app/components/RemoteConnectDialog/RemoteConnectDialog.tsx b/src/web-ui/src/app/components/RemoteConnectDialog/RemoteConnectDialog.tsx index 2639bcb0d8..ce877ba1da 100644 --- a/src/web-ui/src/app/components/RemoteConnectDialog/RemoteConnectDialog.tsx +++ b/src/web-ui/src/app/components/RemoteConnectDialog/RemoteConnectDialog.tsx @@ -65,6 +65,7 @@ import { } from './remoteConnectOperationCleanup'; import { ChatAppBrandIcon } from './ChatAppBrandIcon'; import { RemotePairingCard } from './RemotePairingCard'; +import { WeixinLoginProgress } from './WeixinLoginProgress'; import './RemoteConnectDialog.scss'; // ── Types ──────────────────────────────────────────────────────────── @@ -1367,7 +1368,7 @@ export const RemoteConnectDialog: React.FC = ({ /> )} -

{t('remoteConnect.botWeixinPolling')}

+ @@ -1375,12 +1376,20 @@ export const RemoteConnectDialog: React.FC = ({ )} {weixinQrSessionKey && !weixinQrImageUrl && weixinAwaitingPhoneConfirm && (
-

{t('remoteConnect.botWeixinAwaitingPhoneConfirm')}

+
)} + {weixinQrSessionKey && !weixinQrImageUrl && !weixinAwaitingPhoneConfirm && !weixinNeedsVerifyCode && ( +
+ + +
+ )} {weixinQrSessionKey && !weixinQrImageUrl && weixinNeedsVerifyCode && (
({ useI18n: () => ({ t: (key: string) => key }) })); + +describe('WeChat login progress', () => { + let root: Root; + let container: HTMLDivElement; + beforeEach(() => { + vi.useFakeTimers(); + container = document.createElement('div'); + root = createRoot(container); + }); + afterEach(async () => { + await act(async () => root.unmount()); + vi.useRealTimers(); + }); + + it('explains a slow long poll without reporting a failed connection', async () => { + await act(async () => root.render()); + expect(container.textContent).toContain('botWeixinPolling'); + expect(container.querySelector('[role="status"]')).not.toBeNull(); + await act(async () => vi.advanceTimersByTime(10_000)); + expect(container.textContent).toContain('botWeixinSyncSlow'); + expect(container.textContent).not.toContain('stateConnected'); + }); + + it('updates the phase after phone confirmation and resets stale delay copy', async () => { + await act(async () => root.render()); + await act(async () => vi.advanceTimersByTime(10_000)); + expect(container.textContent).toContain('botWeixinSyncSlow'); + await act(async () => root.render()); + expect(container.textContent).toContain('botWeixinStarting'); + expect(container.textContent).not.toContain('botWeixinSyncSlow'); + await act(async () => vi.advanceTimersByTime(10_000)); + expect(container.textContent).toContain('botWeixinStartingSlow'); + expect(container.querySelectorAll('[role="status"]')).toHaveLength(1); + }); +}); diff --git a/src/web-ui/src/app/components/RemoteConnectDialog/WeixinLoginProgress.tsx b/src/web-ui/src/app/components/RemoteConnectDialog/WeixinLoginProgress.tsx new file mode 100644 index 0000000000..7cf8635bea --- /dev/null +++ b/src/web-ui/src/app/components/RemoteConnectDialog/WeixinLoginProgress.tsx @@ -0,0 +1,41 @@ +import { Spinner, StatusPill } from '@openbitfun/ui'; +import { useEffect, useState } from 'react'; +import { useI18n } from '@/infrastructure/i18n'; + +export type WeixinLoginPhase = 'scan' | 'confirm' | 'starting'; + +/** iLink holds a QR status request for up to 35 seconds; waiting is not a failed login. */ +export function WeixinLoginProgress({ phase }: { phase: WeixinLoginPhase }) { + const { t } = useI18n('common'); + const [slowPhase, setSlowPhase] = useState(null); + useEffect(() => { + setSlowPhase(null); + const timer = window.setTimeout(() => setSlowPhase(phase), 10_000); + return () => window.clearTimeout(timer); + }, [phase]); + + const label = phase === 'scan' + ? t('remoteConnect.botWeixinPolling') + : phase === 'confirm' + ? t('remoteConnect.botWeixinConfirming') + : t('remoteConnect.botWeixinStarting'); + const description = phase === 'scan' + ? t('remoteConnect.botWeixinScanProgressHint') + : phase === 'confirm' + ? t('remoteConnect.botWeixinAwaitingPhoneConfirm') + : t('remoteConnect.botWeixinStartingHint'); + + return ( +
+ }>{label} +

{description}

+ {slowPhase === phase && ( +

+ {phase === 'starting' + ? t('remoteConnect.botWeixinStartingSlow') + : t('remoteConnect.botWeixinSyncSlow')} +

+ )} +
+ ); +} diff --git a/src/web-ui/src/locales/en-US/common.json b/src/web-ui/src/locales/en-US/common.json index e47c817077..807ef5b804 100644 --- a/src/web-ui/src/locales/en-US/common.json +++ b/src/web-ui/src/locales/en-US/common.json @@ -1012,6 +1012,12 @@ "botWeixinVerifyCodePlaceholder": "Enter the number from WeChat", "botWeixinVerifyCodeHint": "Enter the number shown in WeChat on your phone to continue connecting.", "botWeixinVerifyCodeSubmit": "Submit number", + "botWeixinConfirming": "Syncing phone confirmation", + "botWeixinStarting": "Login confirmed, preparing pairing code", + "botWeixinScanProgressHint": "Confirm login on your phone after scanning. The pairing code appears once confirmation is synced.", + "botWeixinStartingHint": "Saving login details and starting WeChat message reception. Please wait.", + "botWeixinSyncSlow": "WeChat status updates can take about 35 seconds. If you have confirmed on your phone, wait for this page to update; no need to scan again.", + "botWeixinStartingSlow": "Setup is still in progress and the network is responding slowly. You can keep waiting or cancel and retry.", "botWeixinPolling": "Waiting for scan…", "botVerboseMode": "Verbose Mode", "botConciseMode": "Concise Mode", diff --git a/src/web-ui/src/locales/zh-CN/common.json b/src/web-ui/src/locales/zh-CN/common.json index a78ab9bbec..ca00db291b 100644 --- a/src/web-ui/src/locales/zh-CN/common.json +++ b/src/web-ui/src/locales/zh-CN/common.json @@ -1012,6 +1012,12 @@ "botWeixinVerifyCodePlaceholder": "输入微信显示的数字", "botWeixinVerifyCodeHint": "请填写手机微信中显示的数字,以继续连接", "botWeixinVerifyCodeSubmit": "提交数字", + "botWeixinConfirming": "正在同步手机确认状态", + "botWeixinStarting": "登录已确认,正在准备配对码", + "botWeixinScanProgressHint": "扫码后请在手机上确认登录,确认结果同步后会显示配对码", + "botWeixinStartingHint": "正在保存登录信息并启动微信消息接收,请稍候", + "botWeixinSyncSlow": "微信状态同步可能需要约 35 秒;若手机已确认,请等待本页更新,无需重复扫码", + "botWeixinStartingSlow": "连接仍在准备中,网络响应较慢;你可以继续等待或取消后重试", "botWeixinPolling": "等待扫码确认…", "botVerboseMode": "详细模式", "botConciseMode": "简洁模式", diff --git a/src/web-ui/src/locales/zh-TW/common.json b/src/web-ui/src/locales/zh-TW/common.json index 22e9bc00ae..9d26a2038c 100644 --- a/src/web-ui/src/locales/zh-TW/common.json +++ b/src/web-ui/src/locales/zh-TW/common.json @@ -1012,6 +1012,12 @@ "botWeixinVerifyCodePlaceholder": "輸入微信顯示的數字", "botWeixinVerifyCodeHint": "請填寫手機微信中顯示的數字,以繼續連接", "botWeixinVerifyCodeSubmit": "提交數字", + "botWeixinConfirming": "正在同步手機確認狀態", + "botWeixinStarting": "登入已確認,正在準備配對碼", + "botWeixinScanProgressHint": "掃碼後請在手機上確認登入,確認結果同步後會顯示配對碼", + "botWeixinStartingHint": "正在儲存登入資訊並啟動微信訊息接收,請稍候", + "botWeixinSyncSlow": "微信狀態同步可能需要約 35 秒;若手機已確認,請等待本頁更新,無需重複掃碼", + "botWeixinStartingSlow": "連線仍在準備中,網路回應較慢;你可以繼續等待或取消後重試", "botWeixinPolling": "等待掃碼確認…", "botVerboseMode": "詳細模式", "botConciseMode": "簡潔模式",