From 923c22f9b45965fa32d268409216f63404ac5da0 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?=D0=9D=D0=B8=D0=BA=D0=B8=D1=82=D0=B0=20Sonic?= Date: Wed, 26 Aug 2026 16:28:43 +0300 Subject: [PATCH] =?UTF-8?q?fix(proxy):=20=D0=B7=D0=B0=D0=B3=D1=80=D1=83?= =?UTF-8?q?=D0=B7=D0=BA=D0=B0=20=D0=BE=D1=81=D1=82=D0=B0=D0=BD=D0=B0=D0=B2?= =?UTF-8?q?=D0=BB=D0=B8=D0=B2=D0=B0=D0=BB=D0=B0=20=D0=BE=D1=82=D0=BF=D1=80?= =?UTF-8?q?=D0=B0=D0=B2=D0=BA=D1=83,=20=D1=82=D1=83=D0=BD=D0=BD=D0=B5?= =?UTF-8?q?=D0=BB=D1=8C=20=D1=88=D1=91=D0=BB=20=D0=B2=20=D0=BE=D0=B4=D0=BD?= =?UTF-8?q?=D1=83=20=D1=81=D1=82=D0=BE=D1=80=D0=BE=D0=BD=D1=83=20(#42,=20#?= =?UTF-8?q?32)=20(#51)?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Оба направления туннеля обслуживал один `select!` с пометкой `biased`. `biased` опрашивает ветки строго по порядку: пока в первой — «Telegram → клиент» — есть данные, до второй очередь не доходит вообще. При непрерывном потоке вниз, то есть при первичной синхронизации телефона или загрузке медиа, исходящие пакеты клиента не читались. Второй дефект в том же цикле: одна задача на оба направления. `tcp_w.write_all` ждёт, пока клиент разберёт присланное, и всё это время не опрашивается чтение от клиента. Телефон по Wi-Fi разбирает поток медленнее, чем Telegram Desktop на той же машине через loopback, — отсюда асимметрия «на компьютере работает, на телефоне нет». Для MTProto это фатально: клиент обязан слать подтверждения, а за каждым следующим куском файла — свой `upload.getFile`. Первый запрос уходит, дальше идёт поток вниз, и следующие запросы наверх не попадают. Снаружи это выглядит как «Подключено» при живом туннеле, нулевых сбоях и нулевых отклонениях: чаты на месте, иконки не грузятся, отправка виснет с часиками. Направления разделены на две независимые половины: `ws.split()` плюс `CryptoContext::split()`, потому что шифры направлений независимы — два потока AES-CTR со своими ключами. Ping приходит в читающую половину, а отвечает на него пишущая, через канал на четыре слота: владелец отправляющей половины должен оставаться ровно один. Два теста, падающие на beta.11: за пять секунд непрерывной загрузки наверх не уходит ни одного байта, и клиент, не успевающий читать, замораживает собственную отправку. Снятие одного `biased` чинит только первый — это и показывает, что дефекта два. Co-authored-by: by-sonic <171230345+by-sonic@users.noreply.github.com> --- docs/ARCHITECTURE_V2.md | 32 +++++ src/mtproto.rs | 43 +++++++ src/proxy.rs | 273 +++++++++++++++++++++++++++++++++++----- 3 files changed, 320 insertions(+), 28 deletions(-) diff --git a/docs/ARCHITECTURE_V2.md b/docs/ARCHITECTURE_V2.md index 5b307ad..26166ea 100644 --- a/docs/ARCHITECTURE_V2.md +++ b/docs/ARCHITECTURE_V2.md @@ -116,6 +116,38 @@ LAN-режим превращал бы машину в открытый прок Счётчик `ws_failures` от них отличается тем, что растёт после успешного рукопожатия с клиентом: там договорились с клиентом, но не смогли с Telegram. +## Туннель: два независимых направления + +Каждое клиентское соединение получает свой WebSocket-туннель, и внутри него +данные идут в обе стороны сразу. До 2.0.0-beta.12 оба направления обслуживал +один `select!` с пометкой `biased`, и это давало два дефекта, снаружи +выглядевших одинаково: «Подключено», а ничего не идёт. + +`biased` опрашивает ветки строго по порядку. Пока в первой — «Telegram → +клиент» — есть данные, до второй очередь не доходит вообще. То есть при +непрерывном потоке вниз (первичная синхронизация телефона, загрузка медиа) +исходящие пакеты клиента не читались. + +Второй дефект — одна задача на оба направления. `tcp_w.write_all` ждёт, пока +клиент разберёт присланное, и всё это время не опрашивается чтение от клиента. +Телефон по Wi-Fi разбирает поток медленнее, чем Telegram Desktop на той же +машине через loopback, — отсюда асимметрия «на компьютере работает, на телефоне +нет» из #42. + +Для MTProto это фатально: клиент обязан слать подтверждения, а за каждым +следующим куском файла — свой `upload.getFile`. Первый запрос уходит, дальше +идёт поток вниз, и следующие запросы наверх не попадают. Загрузка встаёт при +живом туннеле, нулевых сбоях и нулевых отклонениях — ровно картина из #32. + +Теперь это две независимые половины: `ws.split()` плюс `CryptoContext::split()`, +потому что шифры направлений независимы — два потока AES-CTR со своими ключами. +Ping приходит в читающую половину, а отвечает на него пишущая, через канал на +четыре слота: владелец отправляющей половины должен оставаться ровно один. + +Оба дефекта закрыты тестами, которые падают на beta.11. Первый: за пять секунд +непрерывной загрузки наверх не уходит ни одного байта. Второй: клиент, не +успевающий читать, замораживает собственную отправку. + ## Учёт состояния `Stats::ws` считает **установленные** туннели: счётчик поднимается после diff --git a/src/mtproto.rs b/src/mtproto.rs index cfb3542..26e9a99 100644 --- a/src/mtproto.rs +++ b/src/mtproto.rs @@ -43,6 +43,49 @@ impl CryptoContext { self.telegram_decrypt.apply_keystream(data); self.client_encrypt.apply_keystream(data); } + + /// Разделить шифры по направлениям, чтобы туннель шёл в обе стороны сразу. + /// + /// Направления независимы: это два потока AES-CTR со своими ключами, и ни + /// один байт одного не влияет на другой. + pub fn split(self) -> (Upstream, Downstream) { + ( + Upstream { + client_decrypt: self.client_decrypt, + telegram_encrypt: self.telegram_encrypt, + }, + Downstream { + telegram_decrypt: self.telegram_decrypt, + client_encrypt: self.client_encrypt, + }, + ) + } +} + +/// Шифры направления «клиент -> Telegram». +pub struct Upstream { + client_decrypt: AesCtr, + telegram_encrypt: AesCtr, +} + +impl Upstream { + pub fn apply(&mut self, data: &mut [u8]) { + self.client_decrypt.apply_keystream(data); + self.telegram_encrypt.apply_keystream(data); + } +} + +/// Шифры направления «Telegram -> клиент». +pub struct Downstream { + telegram_decrypt: AesCtr, + client_encrypt: AesCtr, +} + +impl Downstream { + pub fn apply(&mut self, data: &mut [u8]) { + self.telegram_decrypt.apply_keystream(data); + self.client_encrypt.apply_keystream(data); + } } pub fn generate_secret() -> [u8; 16] { diff --git a/src/proxy.rs b/src/proxy.rs index cbf56f2..c567909 100644 --- a/src/proxy.rs +++ b/src/proxy.rs @@ -16,6 +16,8 @@ const SOCKS5_VERSION: u8 = 0x05; /// byte as the start of a SOCKS5 greeting. const PROTOCOL_PROBE_TIMEOUT: Duration = Duration::from_millis(250); const PROTOCOL_PROBE_INTERVAL: Duration = Duration::from_millis(5); +/// Сколько неотвеченных Ping'ов держать, пока отправляющая половина занята. +const PONG_QUEUE: usize = 4; pub struct Stats { pub running: AtomicBool, @@ -658,55 +660,92 @@ async fn ws_tunnel( dc: u16, media: bool, init: &[u8; 64], - mut crypto: Option, + crypto: Option, stats: &Stats, ) -> Result<(), Box> { use futures_util::{SinkExt, StreamExt}; - let (mut ws, connected) = stats.transport.connect(dc, media).await?; + let (ws, connected) = stats.transport.connect(dc, media).await?; let _tunnel = EstablishedTunnel::new(stats); stats.note_tunnel(dc, connected.route.kind.ui_code()); let (mut tcp_r, mut tcp_w) = tokio::io::split(tcp); + let (mut ws_w, mut ws_r) = ws.split(); + let (upstream_crypto, downstream_crypto) = match crypto.map(|crypto| crypto.split()) { + Some((upstream, downstream)) => (Some(upstream), Some(downstream)), + None => (None, None), + }; // Send buffered init as first frame - ws.send(tungstenite::Message::Binary(init.to_vec())).await?; + ws_w.send(tungstenite::Message::Binary(init.to_vec())) + .await?; - let mut buf = vec![0u8; 65536]; + // Ping приходит в половину, которая читает, а отвечать на него должна та, + // которая пишет: владелец у отправляющей половины строго один. + let (pong_tx, mut pong_rx) = tokio::sync::mpsc::channel::>(PONG_QUEUE); - loop { - tokio::select! { - biased; - - msg = ws.next() => match msg { - Some(Ok(tungstenite::Message::Binary(mut data))) => { + // Направления работают независимо друг от друга. Раньше это был один + // `select!`, и любое ожидание внутри него останавливало вторую половину: + // непрерывная загрузка не давала опросить клиента вообще, а клиент, + // который не успевал разбирать входящий поток, замораживал заодно и свою + // отправку. Telegram при этом ждёт от клиента подтверждений — без них + // сессия встаёт при живом туннеле (by-sonic/tglock#42, #32). + let downstream = async { + let mut crypto = downstream_crypto; + while let Some(message) = ws_r.next().await { + match message { + Ok(tungstenite::Message::Binary(mut data)) => { if let Some(crypto) = &mut crypto { - crypto.telegram_to_client(data.as_mut()); + crypto.apply(data.as_mut()); } tcp_w.write_all(data.as_ref()).await?; tcp_w.flush().await?; } - Some(Ok(tungstenite::Message::Ping(p))) => { - let _ = ws.send(tungstenite::Message::Pong(p)).await; - } - Some(Ok(tungstenite::Message::Close(_))) | None => break, - Some(Err(_)) => break, - _ => {} - }, - - n = tcp_r.read(&mut buf) => match n { - Ok(0) | Err(_) => break, - Ok(n) => { - if let Some(crypto) = &mut crypto { - crypto.client_to_telegram(&mut buf[..n]); + Ok(tungstenite::Message::Ping(payload)) => { + if pong_tx.send(payload).await.is_err() { + break; } - ws.send(tungstenite::Message::Binary(buf[..n].to_vec())).await?; } - }, + Ok(tungstenite::Message::Close(_)) | Err(_) => break, + Ok(_) => {} + } } - } + Ok::<(), Box>(()) + }; - let _ = ws.close(None).await; + let upstream = async { + let mut crypto = upstream_crypto; + let mut buf = vec![0u8; 65536]; + loop { + tokio::select! { + read = tcp_r.read(&mut buf) => match read { + Ok(0) | Err(_) => break, + Ok(read) => { + if let Some(crypto) = &mut crypto { + crypto.apply(&mut buf[..read]); + } + ws_w + .send(tungstenite::Message::Binary(buf[..read].to_vec())) + .await?; + } + }, + payload = pong_rx.recv() => match payload { + Some(payload) => { + ws_w.send(tungstenite::Message::Pong(payload)).await?; + } + None => break, + }, + } + } + let _ = ws_w.close().await; + Ok::<(), Box>(()) + }; + + tokio::pin!(downstream, upstream); + tokio::select! { + result = &mut downstream => result?, + result = &mut upstream => result?, + } Ok(()) } @@ -1466,4 +1505,182 @@ mod tests { stats.stop(); server.await.unwrap().unwrap(); } + + /// Скачивание не должно затыкать отправку. + /// + /// В `ws_tunnel` цикл `select!` помечен `biased`, то есть сначала всегда + /// опрашивается ветка чтения из WebSocket. Пока Telegram присылает данные + /// непрерывно — а именно так выглядит загрузка медиа или первичная + /// синхронизация телефона — ветка чтения из клиента не опрашивается + /// вообще, и исходящие пакеты клиента наверх не уходят. + #[allow(clippy::result_large_err)] + #[tokio::test] + async fn a_download_in_flight_must_not_stop_the_client_from_sending() { + use futures_util::{SinkExt, StreamExt}; + + let relay_listener = TcpListener::bind("127.0.0.1:0").await.unwrap(); + let relay_port = relay_listener.local_addr().unwrap().port(); + let uploads = Arc::new(std::sync::atomic::AtomicUsize::new(0)); + let relay_uploads = uploads.clone(); + + tokio::spawn(async move { + let (tcp, _) = relay_listener.accept().await.unwrap(); + let websocket = + tokio_tungstenite::accept_hdr_async(tcp, |_: &Request, mut response: Response| { + response.headers_mut().insert( + "Sec-WebSocket-Protocol", + "binary".parse().expect("static header value"), + ); + Ok(response) + }) + .await + .unwrap(); + let (mut sink, mut stream) = websocket.split(); + + match stream.next().await { + Some(Ok(Message::Binary(_))) => {} + other => panic!("expected an init frame, got {other:?}"), + } + + tokio::spawn(async move { + while let Some(message) = stream.next().await { + if let Ok(Message::Binary(data)) = message { + relay_uploads.fetch_add(data.len(), Ordering::Relaxed); + } + } + }); + + // Непрерывный поток вниз — так выглядит загрузка медиа. + while sink.send(Message::Binary(vec![0; 32 * 1024])).await.is_ok() {} + }); + + let stats = Stats::new(); + stats.transport.force_local_route(relay_port); + let (port, server) = start_proxy(stats.clone(), false).await; + + let init = unambiguous_client_init(&stats.secret, -4); + let client = TcpStream::connect(("127.0.0.1", port)).await.unwrap(); + let (mut client_r, mut client_w) = client.into_split(); + client_w.write_all(&init).await.unwrap(); + + // Клиент исправно читает загрузку, иначе он затыкал бы туннель сам. + let downloaded = Arc::new(std::sync::atomic::AtomicUsize::new(0)); + let counted = downloaded.clone(); + tokio::spawn(async move { + let mut drain = vec![0; 64 * 1024]; + while let Ok(read) = client_r.read(&mut drain).await { + if read == 0 { + break; + } + counted.fetch_add(read, Ordering::Relaxed); + } + }); + + wait_until("загрузка пошла", || { + downloaded.load(Ordering::Relaxed) > 1024 * 1024 + }) + .await; + + // Telegram ждёт от клиента подтверждений и запросов. Без них сессия + // встаёт: «Подключено», а сообщения висят с часиками. + tokio::spawn(async move { + while client_w.write_all(&[0x42; 128]).await.is_ok() { + tokio::time::sleep(Duration::from_millis(20)).await; + } + }); + + tokio::time::timeout(Duration::from_secs(5), async { + while uploads.load(Ordering::Relaxed) == 0 { + tokio::time::sleep(Duration::from_millis(10)).await; + } + }) + .await + .expect("пока идёт загрузка, клиент должен доставить наверх хоть один байт"); + + stats.stop(); + let _ = server.await.unwrap(); + } + + /// Медленный клиент не должен останавливать весь туннель. + /// + /// `ws_tunnel` читает и пишет в одной задаче: пока `tcp_w.write_all` ждёт, + /// когда клиент разберёт присланное, ветка чтения из клиента не + /// опрашивается, и наверх не уходит ничего. Телефон по Wi-Fi разбирает + /// поток медленнее, чем десктоп на той же машине по loopback — отсюда + /// асимметрия «на компьютере работает, на телефоне нет». + #[allow(clippy::result_large_err)] + #[tokio::test] + async fn a_slow_client_must_not_freeze_its_own_uploads() { + use futures_util::{SinkExt, StreamExt}; + + let relay_listener = TcpListener::bind("127.0.0.1:0").await.unwrap(); + let relay_port = relay_listener.local_addr().unwrap().port(); + let uploads = Arc::new(std::sync::atomic::AtomicUsize::new(0)); + let relay_uploads = uploads.clone(); + + tokio::spawn(async move { + let (tcp, _) = relay_listener.accept().await.unwrap(); + let websocket = + tokio_tungstenite::accept_hdr_async(tcp, |_: &Request, mut response: Response| { + response.headers_mut().insert( + "Sec-WebSocket-Protocol", + "binary".parse().expect("static header value"), + ); + Ok(response) + }) + .await + .unwrap(); + let (mut sink, mut stream) = websocket.split(); + + match stream.next().await { + Some(Ok(Message::Binary(_))) => {} + other => panic!("expected an init frame, got {other:?}"), + } + + tokio::spawn(async move { + while let Some(message) = stream.next().await { + if let Ok(Message::Binary(data)) = message { + relay_uploads.fetch_add(data.len(), Ordering::Relaxed); + } + } + }); + + while sink.send(Message::Binary(vec![0; 32 * 1024])).await.is_ok() {} + }); + + let stats = Stats::new(); + stats.transport.force_local_route(relay_port); + let (port, server) = start_proxy(stats.clone(), false).await; + + let init = unambiguous_client_init(&stats.secret, -4); + let client = TcpStream::connect(("127.0.0.1", port)).await.unwrap(); + let (client_r, mut client_w) = client.into_split(); + client_w.write_all(&init).await.unwrap(); + + // Клиент занят и не разбирает входящий поток: его приёмное окно + // закрывается, и запись в него встаёт. + wait_until("туннель поднялся", || { + stats.last_route() != 0 + }) + .await; + tokio::time::sleep(Duration::from_secs(2)).await; + + tokio::spawn(async move { + while client_w.write_all(&[0x42; 128]).await.is_ok() { + tokio::time::sleep(Duration::from_millis(20)).await; + } + }); + + let result = tokio::time::timeout(Duration::from_secs(5), async { + while uploads.load(Ordering::Relaxed) == 0 { + tokio::time::sleep(Duration::from_millis(10)).await; + } + }) + .await; + + drop(client_r); + stats.stop(); + let _ = server.await.unwrap(); + result.expect("клиент, который не успевает читать, всё равно должен отправлять"); + } }