From e9af114f3ce4a9c244875cf0add7ed512f5a787d 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: Thu, 27 Aug 2026 02:34:20 +0300 Subject: [PATCH] =?UTF-8?q?fix(worker):=20=D0=B7=D0=B0=D0=BF=D0=B8=D1=81?= =?UTF-8?q?=D1=8C=20=D0=B2=20Telegram=20=D1=88=D0=BB=D0=B0=20=D0=B1=D0=B5?= =?UTF-8?q?=D0=B7=20=D0=BE=D0=B6=D0=B8=D0=B4=D0=B0=D0=BD=D0=B8=D1=8F=20?= =?UTF-8?q?=D0=B8=20backpressure=20(#42)=20(#56)?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Скрипт воркера писал в сокет Telegram так: server.addEventListener("message", (event) => { writer.write(chunk).catch(shutdown); }); `write()` вызывался поверх незавершённого, `writer.ready` не спрашивался вовсе. Пока в клиенте голодала отправка, настоящего потока вверх через воркер не возникало, и код держался. В beta.12 голодание починили — поток появился, и репортёр #42 сразу получил переподключения на обоих клиентах, которых на beta.11 с тем же воркером не было. Запись сериализована цепочкой промисов: следующий чанк уходит после того, как записан предыдущий, и только когда писатель готов. Кто разворачивал воркер раньше — нужен передеплой, о чём сказано в docs/CLOUDFLARE_WORKER.md. Причина у репортёра не подтверждена: рантайма Workers у меня нет, проверить можно только у него. Заодно счётчик «промолчали». Соединение, которое открылось и ничего не прислало за `IO_TIMEOUT`, закрывалось и не попадало ни в один счётчик: `unknown_clients` растёт, только когда запрос пришёл и не разобрался, а не когда его не дождались. Тот же репортёр сообщил, что его телефон переустанавливает соединение примерно раз в десять секунд — ровно период `IO_TIMEOUT`. Проверить это по диагностике было нечем, теперь есть чем. Тесты: Ping через туннель (путь не был покрыт вовсе, а Ping бывает только на маршруте воркера) и молчащий клиент. Второй гоняет виртуальное время, чтобы не ждать десять секунд по-настоящему, — отсюда dev-зависимость на tokio/test-util. Co-authored-by: by-sonic <171230345+by-sonic@users.noreply.github.com> --- Cargo.toml | 6 ++ docs/ARCHITECTURE_V2.md | 12 +++ docs/CLOUDFLARE_WORKER.md | 2 + src/bin/cli.rs | 6 +- src/main.rs | 7 ++ src/proxy.rs | 191 ++++++++++++++++++++++++++++++++++++-- ui/main.ts | 7 ++ worker/tglock-worker.js | 15 ++- 8 files changed, 237 insertions(+), 9 deletions(-) diff --git a/Cargo.toml b/Cargo.toml index 8c5bedd..563803c 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -65,3 +65,9 @@ rand = "0.8" [build-dependencies] tauri-build = { version = "2", features = [], optional = true } + +[dev-dependencies] +# `start_paused` в тестах: таймаут ожидания запроса от клиента — десять секунд, +# и ждать их по-настоящему в тесте нельзя. В сборку не попадает: dev-зависимости +# участвуют только в тестах. +tokio = { version = "1", features = ["test-util"] } diff --git a/docs/ARCHITECTURE_V2.md b/docs/ARCHITECTURE_V2.md index eca4204..536db80 100644 --- a/docs/ARCHITECTURE_V2.md +++ b/docs/ARCHITECTURE_V2.md @@ -233,6 +233,14 @@ Worker должен принимать WebSocket на: добавить собственную авторизацию до стабильного релиза; поэтому Worker остаётся расширенной опцией alpha-версии. +Запись в сокет Telegram обязана быть последовательной: следующий чанк уходит +после того, как записан предыдущий, и только когда писатель к этому готов +(`writer.ready`). В `worker/tglock-worker.js` этого не было — `write()` +вызывался поверх незавершённого, без backpressure. Пока в клиенте голодала +отправка, настоящего потока вверх через воркер не возникало и это не +проявлялось; после того как голодание починили, поток появился. Кто разворачивал +воркер раньше — обновите скрипт. + ## Current limitations - SNI camouflage не включена: небезопасное отключение hostname verification @@ -255,5 +263,9 @@ Worker должен принимать WebSocket на: CLI собирается одной командой. Подробно — в разделе README про антивирус. - Работоспособность медиа зависит от конкретного DC аккаунта и доступности Telegram/Cloudflare у провайдера. +- Соединение, открытое клиентом и молчащее дольше `IO_TIMEOUT` (10 секунд), + закрывается. Для клиента, открывающего соединения про запас, это норма; счётчик + «промолчали» показывает, как часто это происходит, — раньше такие соединения + не попадали никуда. - Пулы заранее открытых WebSocket-соединений будут добавлены после измерения, что они не создают лишнюю нагрузку и не ухудшают стабильность. diff --git a/docs/CLOUDFLARE_WORKER.md b/docs/CLOUDFLARE_WORKER.md index 50629b8..127d184 100644 --- a/docs/CLOUDFLARE_WORKER.md +++ b/docs/CLOUDFLARE_WORKER.md @@ -41,6 +41,8 @@ Если вернулось `not found` — проверь, что путь именно `/apiws`. Если ошибка про `cloudflare:sockets` — у воркера слишком старая дата совместимости, поставь в **Settings → Compatibility date** сегодняшнюю. +> **Разворачивал воркер до 2.0.0-beta.14 — обнови скрипт.** В прежней версии запись в сокет Telegram шла без ожидания предыдущей и без backpressure. Пока в клиенте голодала отправка, через воркер не проходило настоящего потока вверх и это не проявлялось; после того как голодание починили в beta.12, поток появился. + ## Подключение в TGLock **В приложении:** Настройки → поле **Cloudflare Worker** → вставь `tglock.имя.workers.dev` → Сохранить. Настройки меняются только при выключенной защите. diff --git a/src/bin/cli.rs b/src/bin/cli.rs index 1300a09..44be4a2 100644 --- a/src/bin/cli.rs +++ b/src/bin/cli.rs @@ -224,14 +224,16 @@ async fn watch_status(stats: Arc) { stats.route_failures(), stats.blocked.load(Ordering::Relaxed), stats.unknown_clients.load(Ordering::Relaxed), + stats.silent_clients.load(Ordering::Relaxed), ); if previous.as_ref() == Some(¤t) { continue; } - let (active, tunnels, dc, route, failures, route_failures, blocked, unknown) = current; + let (active, tunnels, dc, route, failures, route_failures, blocked, unknown, silent) = + current; let line = format!( "соединений {active} · туннелей {tunnels} · {} · {} · сбоев {failures} · \ - падений маршрутов {route_failures} · отклонено {blocked} · не опознано {unknown}", + падений маршрутов {route_failures} · отклонено {blocked} · не опознано {unknown} · промолчали {silent}", if dc > 0 { format!("DC{dc}") } else { diff --git a/src/main.rs b/src/main.rs index 49dfeaa..1cc7232 100644 --- a/src/main.rs +++ b/src/main.rs @@ -56,6 +56,12 @@ struct StatusSnapshot { /// Клиенты, которые дошли, но не сумели договориться. Почти всегда это /// ссылка `tg://proxy` от прошлого запуска, то есть другой секрет. unknown_clients: u32, + /// Соединения, которые открылись и ничего не прислали до таймаута. + /// + /// Растущее число в такт с переподключениями клиента означает, что он + /// открывает соединения впрок, а мы закрываем их по таймауту + /// (by-sonic/tglock#42). + silent_clients: u32, uptime_seconds: u64, port: u16, /// Адрес, который нужно вписать в Telegram на другом устройстве. @@ -128,6 +134,7 @@ impl AppState { route_failures: self.stats.route_failures(), blocked: self.stats.blocked.load(Ordering::Relaxed), unknown_clients: self.stats.unknown_clients.load(Ordering::Relaxed), + silent_clients: self.stats.silent_clients.load(Ordering::Relaxed), uptime_seconds: self .started_at .lock() diff --git a/src/proxy.rs b/src/proxy.rs index b072edd..6b6c3b1 100644 --- a/src/proxy.rs +++ b/src/proxy.rs @@ -41,6 +41,15 @@ pub struct Stats { /// прошлого запуска. Такое соединение закрывалось молча, и по диагностике /// отличить его от рабочего было нельзя. pub unknown_clients: AtomicU32, + /// Сколько соединений открылось и ничего не прислало за `IO_TIMEOUT`. + /// + /// Такое соединение закрывается по таймауту и до сих пор не попадало ни в + /// один счётчик: `unknown_clients` растёт, только когда запрос пришёл и не + /// разобрался, а не когда его не дождались. По диагностике это выглядело + /// как соединение без туннеля и ничего больше — а репортёр #42 видел, что + /// его телефон переустанавливает соединение примерно раз в десять секунд, + /// то есть ровно с этим периодом. + pub silent_clients: AtomicU32, /// DC и маршрут последнего поднятого туннеля, упакованные в одно значение. /// /// Раньше это были два независимых поля: номер писало соединение при @@ -113,6 +122,7 @@ impl Stats { ws_failures: AtomicU32::new(0), blocked: AtomicU32::new(0), unknown_clients: AtomicU32::new(0), + silent_clients: AtomicU32::new(0), last_tunnel: AtomicU32::new(0), transport: crate::transport::TransportEngine::new(), secret, @@ -158,6 +168,25 @@ impl Stats { self.note(format!("Клиент {who}: {reason}")); } + /// Клиент открыл соединение и не сказал ничего. + /// + /// Отличается от `note_unknown_client` тем, что там запрос пришёл и не + /// разобрался, а здесь его не дождались. Для клиента, который открывает + /// соединения про запас, это норма; для клиента, который переоткрывает их + /// в такт с таймаутом, — нет, и разницу видно только по счётчику + /// (by-sonic/tglock#42). + fn note_silent_client(&self, peer: Option, reason: &str) { + self.silent_clients.fetch_add(1, Ordering::Relaxed); + let who = match peer { + Some(peer) => peer.ip().to_string(), + None => "неизвестный адрес".to_owned(), + }; + self.note(format!( + "Клиент {who}: {reason} за {} с — соединение закрыто", + IO_TIMEOUT.as_secs() + )); + } + /// Отметить, что до прокси дотянулось устройство из сети, а не с этой машины. /// /// Это первое, что нужно знать при разборе LAN-режима: если строки нет, @@ -349,9 +378,13 @@ async fn detect_protocol( stats: &Stats, ) -> Result> { let mut probe = [0; INIT_LEN]; - let peeked = tokio::time::timeout(IO_TIMEOUT, stream.peek(&mut probe[..1])) - .await - .map_err(|_| "client protocol detection timeout")??; + let peeked = match tokio::time::timeout(IO_TIMEOUT, stream.peek(&mut probe[..1])).await { + Ok(peeked) => peeked?, + Err(_) => { + stats.note_silent_client(stream.peer_addr().ok(), "не прислал ни байта"); + return Err("client protocol detection timeout".into()); + } + }; if peeked == 0 { return Ok(Protocol::Empty); } @@ -522,9 +555,15 @@ async fn handle_mtproto( stream.set_nodelay(true)?; let peer = stream.peer_addr().ok(); let mut init = [0; 64]; - tokio::time::timeout(IO_TIMEOUT, stream.read_exact(&mut init)) - .await - .map_err(|_| "MTProto init timeout")??; + match tokio::time::timeout(IO_TIMEOUT, stream.read_exact(&mut init)).await { + Ok(result) => { + result?; + } + Err(_) => { + stats.note_silent_client(peer, "начал MTProto-init и не дослал его"); + return Err("MTProto init timeout".into()); + } + } let parsed = match crate::mtproto::parse_client_init(&init, &stats.secret) { Some(parsed) => parsed, None => { @@ -1556,6 +1595,146 @@ mod tests { ); } + /// Ping от той стороны обязан получить Pong, и туннель обязан это пережить. + /// + /// Ping приходит в читающую половину, а отвечать на него должна пишущая. + /// Единственный маршрут, где Ping вообще бывает, — Cloudflare Worker: + /// у Telegram его нет. Поэтому поломка на этом пути видна только тем, у + /// кого настроен воркер (by-sonic/tglock#42). + #[allow(clippy::result_large_err)] + #[tokio::test] + async fn a_ping_is_answered_and_the_tunnel_survives_it() { + 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 (report_tx, report_rx) = tokio::sync::oneshot::channel::, String>>(); + + tokio::spawn(async move { + let (tcp, _) = relay_listener.accept().await.unwrap(); + let mut 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 init = match websocket.next().await { + Some(Ok(Message::Binary(data))) => data, + other => { + let _ = report_tx.send(Err(format!("ожидался init, пришло {other:?}"))); + return; + } + }; + let header: [u8; INIT_LEN] = init.as_slice().try_into().unwrap(); + let mut relay = crate::mtproto::test_relay_peer(&header); + + websocket + .send(Message::Ping(b"keepalive".to_vec())) + .await + .unwrap(); + + let mut answered = false; + let mut payload = Vec::new(); + while payload.is_empty() { + match websocket.next().await { + Some(Ok(Message::Pong(pong))) => { + answered = pong == b"keepalive"; + } + Some(Ok(Message::Binary(mut data))) => { + relay.decrypt(&mut data); + payload.extend_from_slice(&data); + } + Some(Ok(_)) => {} + other => { + let _ = report_tx.send(Err(format!( + "туннель закрылся до полезных данных: {other:?}" + ))); + return; + } + } + } + let _ = report_tx.send(if answered { + Ok(payload) + } else { + Err("Pong не пришёл".to_owned()) + }); + }); + + 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 mut peer = crate::mtproto::test_client_peer(&init, &stats.secret); + let mut client = TcpStream::connect(("127.0.0.1", port)).await.unwrap(); + client.write_all(&init).await.unwrap(); + + // Клиент шлёт после Ping'а: если ответ на Ping ломает пишущую половину, + // это не дойдёт. + tokio::time::sleep(Duration::from_millis(200)).await; + let request = b"a request sent after the ping".to_vec(); + let mut wire = request.clone(); + peer.encrypt(&mut wire); + client.write_all(&wire).await.unwrap(); + + let report = tokio::time::timeout(Duration::from_secs(5), report_rx) + .await + .expect("реле должно доложить о результате") + .unwrap(); + assert_eq!(report, Ok(request), "Ping не должен ломать туннель"); + + stats.stop(); + let _ = server.await.unwrap(); + } + + /// Клиент, открывший соединение и промолчавший, обязан быть посчитан. + /// + /// До этого он не попадал никуда: `не опознано` растёт только когда запрос + /// пришёл и не разобрался. Репортёр #42 видел, что телефон переоткрывает + /// соединение примерно раз в десять секунд — ровно период `IO_TIMEOUT`, — + /// и проверить это по диагностике было нечем. + #[tokio::test(start_paused = true)] + async fn a_client_that_says_nothing_is_counted_and_named() { + let stats = Stats::new(); + let (port, server) = start_proxy(stats.clone(), false).await; + + // Соединение открыто и молчит. Время в тесте идёт само, как только + // рантайму больше нечего делать. + let _client = TcpStream::connect(("127.0.0.1", port)).await.unwrap(); + + // Бюджет ожидания должен быть больше `IO_TIMEOUT`: время в тесте + // виртуальное и прыгает к ближайшему сроку, поэтому пятисекундный + // предел `wait_until` сработал бы первым. + tokio::time::timeout(Duration::from_secs(60), async { + while stats.silent_clients.load(Ordering::Relaxed) == 0 { + tokio::time::sleep(Duration::from_millis(50)).await; + } + }) + .await + .expect("молчащий клиент должен быть посчитан"); + assert_eq!( + stats.unknown_clients.load(Ordering::Relaxed), + 0, + "молчание — не то же самое, что неразобранный запрос" + ); + + let events = stats.drain_events(); + assert!( + events + .iter() + .any(|event| event.contains("не прислал ни байта")), + "в журнале должно быть сказано, что клиент молчал: {events:?}" + ); + + stats.stop(); + let _ = server.await.unwrap(); + } + #[tokio::test] #[ignore = "requires live Telegram network access"] async fn accepts_mtproto_and_builds_live_media_tunnel() { diff --git a/ui/main.ts b/ui/main.ts index debfd08..a182025 100644 --- a/ui/main.ts +++ b/ui/main.ts @@ -15,6 +15,8 @@ type Status = { blocked: number; /// Клиенты, которые дошли, но не сумели договориться о рукопожатии. unknownClients: number; + /// Соединения, которые открылись и ничего не прислали до таймаута. + silentClients: number; uptimeSeconds: number; port: number; /// Адрес для других устройств. Приходит только в LAN-режиме. @@ -47,6 +49,7 @@ let status: Status = { routeFailures: 0, blocked: 0, unknownClients: 0, + silentClients: 0, uptimeSeconds: 0, port: 1080, shareAddress: null, @@ -297,6 +300,10 @@ function renderDiagnostics(): void { Не опознаны ${status.unknownClients} +
+ Промолчали + ${status.silentClients} +

diff --git a/worker/tglock-worker.js b/worker/tglock-worker.js index bcd6df0..43a24fb 100644 --- a/worker/tglock-worker.js +++ b/worker/tglock-worker.js @@ -68,12 +68,25 @@ export default { } }; + // Запись сериализуется: следующий чанк уходит только после того, как + // записан предыдущий, и только когда писатель к этому готов. + // + // Раньше `write()` вызывался поверх незавершённого, а `writer.ready` не + // спрашивался вовсе — backpressure не применялся. Пока в клиенте отправка + // голодала, поверх воркера настоящего потока вверх не бывало и это не + // проявлялось. Как только голодание починили, в воркер пошёл настоящий + // поток (by-sonic/tglock#42). + let pending = Promise.resolve(); + server.addEventListener("message", (event) => { const chunk = event.data instanceof ArrayBuffer ? new Uint8Array(event.data) : event.data; - writer.write(chunk).catch(shutdown); + pending = pending + .then(() => writer.ready) + .then(() => writer.write(chunk)) + .catch(shutdown); }); server.addEventListener("close", shutdown); server.addEventListener("error", shutdown);