mirror of
https://github.com/by-sonic/tglock.git
synced 2026-09-05 18:16:09 +03:00
Compare commits
6 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| 935121b991 | |||
| fe9aad5ee9 | |||
| 28085f9127 | |||
| 923c22f9b4 | |||
| 8ce3c368a3 | |||
| 0dc6b7bf3e |
Generated
+1
-1
@@ -3553,7 +3553,7 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "tglock"
|
||||
version = "2.0.0-beta.10"
|
||||
version = "2.0.0-beta.13"
|
||||
dependencies = [
|
||||
"aes",
|
||||
"cipher",
|
||||
|
||||
+1
-1
@@ -1,6 +1,6 @@
|
||||
[package]
|
||||
name = "tglock"
|
||||
version = "2.0.0-beta.10"
|
||||
version = "2.0.0-beta.13"
|
||||
edition = "2021"
|
||||
rust-version = "1.88"
|
||||
description = "Telegram unblock via local WebSocket tunnel"
|
||||
|
||||
@@ -116,6 +116,64 @@ LAN-режим превращал бы машину в открытый прок
|
||||
Счётчик `ws_failures` от них отличается тем, что растёт после успешного
|
||||
рукопожатия с клиентом: там договорились с клиентом, но не смогли с Telegram.
|
||||
|
||||
### Почему не поднялся туннель
|
||||
|
||||
`ws_failures` говорит, что каскад маршрутов упал целиком, и молчит о причине.
|
||||
Текст с перечислением попыток собирался в `TransportEngine::connect` и там же
|
||||
пропадал: наверх уходил `Err`, который выбрасывался в `serve`. При `туннелей 0`
|
||||
и растущих сбоях отличить «провайдер режет закреплённые адреса» от «воркер
|
||||
отвечает отказом» было нечем — ровно та стена, в которую упёрся репортёр #50.
|
||||
|
||||
Теперь причина попадает в журнал одной строкой на каждый набор отказов:
|
||||
|
||||
```
|
||||
Не поднялся туннель до DC2: 149.154.167.51 — не отвечает (таймаут TCP);
|
||||
kws2.web.telegram.org — таймаут TLS/WebSocket; my.workers.dev — рукопожатие
|
||||
WebSocket: HTTP error: 403 Forbidden
|
||||
```
|
||||
|
||||
Дедупликация журнала делает эту строку разовой: маршруты у DC стабильны, и
|
||||
повтор той же комбинации отказов не пишется.
|
||||
|
||||
Домены Cloudflare Worker отчитываются так же. Строка, не похожая на имя хоста,
|
||||
раньше отбрасывалась молча — `https://name.workers.dev/` со схемой или слэшем не
|
||||
проходит `valid_domain`, маршрут не появлялся, и «воркер настроен» ничем не
|
||||
отличалось от «воркера нет». Теперь отвергнутая строка называется вместе с
|
||||
причиной, а принятая подтверждается: `Cloudflare Worker в списке маршрутов:
|
||||
name.workers.dev`.
|
||||
|
||||
## Туннель: два независимых направления
|
||||
|
||||
Каждое клиентское соединение получает свой 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` считает **установленные** туннели: счётчик поднимается после
|
||||
@@ -124,6 +182,14 @@ LAN-режим превращал бы машину в открытый прок
|
||||
несколько секунд каждый. Состояния «порт открыт», «идёт перебор маршрутов» и
|
||||
«туннель установлен» различимы и в GUI, и в выводе CLI.
|
||||
|
||||
Дата-центр и маршрут пишутся **одним значением**, в момент, когда туннель
|
||||
поднялся. Пока это были два независимых поля, номер писало соединение при
|
||||
разборе init, а маршрут — другое соединение после рукопожатия, и при десятках
|
||||
одновременных соединений в строку статуса попадала пара из разных из них.
|
||||
Читалась она как «до этого DC шли этим маршрутом», хотя означала другое: в
|
||||
диагностике #42 встречались строки `DC5 · Запасной Telegram IP`, а у DC5
|
||||
закреплённый адрес всего один и запасного у него не бывает вовсе.
|
||||
|
||||
Секрет прокси — половина ссылки `tg://proxy`. Для сервиса его нужно закрепить
|
||||
файлом (`--secret-file`): под `DynamicUser` и `ProtectHome` домашней папки нет,
|
||||
путь по умолчанию не определяется, и секрет генерировался бы заново при каждом
|
||||
|
||||
+1
-1
@@ -1,7 +1,7 @@
|
||||
{
|
||||
"name": "tglock-ui",
|
||||
"private": true,
|
||||
"version": "2.0.0-beta.10",
|
||||
"version": "2.0.0-beta.13",
|
||||
"type": "module",
|
||||
"scripts": {
|
||||
"dev": "vite --port 1420",
|
||||
|
||||
+2
-2
@@ -218,8 +218,8 @@ async fn watch_status(stats: Arc<proxy::Stats>) {
|
||||
let current = (
|
||||
stats.active.load(Ordering::Relaxed),
|
||||
stats.ws.load(Ordering::Relaxed),
|
||||
stats.last_dc.load(Ordering::Relaxed),
|
||||
stats.last_route.load(Ordering::Relaxed),
|
||||
stats.last_dc(),
|
||||
stats.last_route(),
|
||||
stats.ws_failures.load(Ordering::Relaxed),
|
||||
stats.route_failures(),
|
||||
stats.blocked.load(Ordering::Relaxed),
|
||||
|
||||
+2
-2
@@ -116,8 +116,8 @@ impl AppState {
|
||||
for event in self.stats.drain_events() {
|
||||
self.log(event, false);
|
||||
}
|
||||
let data_center = self.stats.last_dc.load(Ordering::Relaxed);
|
||||
let route = transport::route_label(self.stats.last_route.load(Ordering::Relaxed));
|
||||
let data_center = self.stats.last_dc();
|
||||
let route = transport::route_label(self.stats.last_route());
|
||||
StatusSnapshot {
|
||||
running: self.stats.running.load(Ordering::SeqCst),
|
||||
active_connections: self.stats.active.load(Ordering::Relaxed),
|
||||
|
||||
@@ -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] {
|
||||
|
||||
+417
-52
@@ -1,7 +1,7 @@
|
||||
use crate::config::ListenConfig;
|
||||
use std::collections::{HashSet, VecDeque};
|
||||
use std::net::{IpAddr, Ipv4Addr, SocketAddr};
|
||||
use std::sync::atomic::{AtomicBool, AtomicU16, AtomicU32, AtomicU8, Ordering};
|
||||
use std::sync::atomic::{AtomicBool, AtomicU16, AtomicU32, Ordering};
|
||||
use std::sync::{Arc, Mutex};
|
||||
use std::time::Duration;
|
||||
use tokio::io::{AsyncRead, AsyncReadExt, AsyncWrite, AsyncWriteExt};
|
||||
@@ -16,13 +16,17 @@ 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,
|
||||
pub active: AtomicU32,
|
||||
pub total: AtomicU32,
|
||||
pub ws: AtomicU32,
|
||||
pub last_dc: AtomicU16,
|
||||
/// DC последнего разобранного соединения. Показывается, пока туннеля ещё
|
||||
/// нет: клиент уже понят, маршрут ещё не выбран.
|
||||
seen_dc: AtomicU16,
|
||||
pub ws_failures: AtomicU32,
|
||||
/// Сколько запросов отклонено политикой «только Telegram».
|
||||
///
|
||||
@@ -37,8 +41,19 @@ pub struct Stats {
|
||||
/// прошлого запуска. Такое соединение закрывалось молча, и по диагностике
|
||||
/// отличить его от рабочего было нельзя.
|
||||
pub unknown_clients: AtomicU32,
|
||||
/// See `transport::RouteKind::ui_code`.
|
||||
pub last_route: AtomicU8,
|
||||
/// DC и маршрут последнего поднятого туннеля, упакованные в одно значение.
|
||||
///
|
||||
/// Раньше это были два независимых поля: номер писало соединение при
|
||||
/// разборе init, маршрут — другое соединение после рукопожатия. При
|
||||
/// нескольких десятках одновременных соединений пара в строке статуса
|
||||
/// складывалась из разных из них, и читалась она как «до этого DC шли
|
||||
/// этим маршрутом», хотя означала совсем не это. У DC1, DC3, DC5 и DC203
|
||||
/// закреплённый адрес всего один, и «запасного» у них не бывает вовсе —
|
||||
/// а строки `DC5 · Запасной Telegram IP` в диагностике встречались
|
||||
/// (by-sonic/tglock#42).
|
||||
///
|
||||
/// Формат: `dc << 8 | route`, где route — `transport::RouteKind::ui_code`.
|
||||
last_tunnel: AtomicU32,
|
||||
transport: crate::transport::TransportEngine,
|
||||
secret: [u8; 16],
|
||||
/// Почему секрет не удалось сохранить, если не удалось.
|
||||
@@ -94,11 +109,11 @@ impl Stats {
|
||||
active: AtomicU32::new(0),
|
||||
total: AtomicU32::new(0),
|
||||
ws: AtomicU32::new(0),
|
||||
last_dc: AtomicU16::new(0),
|
||||
seen_dc: AtomicU16::new(0),
|
||||
ws_failures: AtomicU32::new(0),
|
||||
blocked: AtomicU32::new(0),
|
||||
unknown_clients: AtomicU32::new(0),
|
||||
last_route: AtomicU8::new(0),
|
||||
last_tunnel: AtomicU32::new(0),
|
||||
transport: crate::transport::TransportEngine::new(),
|
||||
secret,
|
||||
secret_write_error,
|
||||
@@ -154,6 +169,34 @@ impl Stats {
|
||||
self.note(format!("Подключилось устройство из сети: {}", peer.ip()));
|
||||
}
|
||||
|
||||
/// Запомнить DC, с которым пришёл клиент. Туннеля может ещё не быть.
|
||||
fn note_dc(&self, dc: u16) {
|
||||
self.seen_dc.store(dc, Ordering::Relaxed);
|
||||
}
|
||||
|
||||
/// Запомнить, каким маршрутом поднялся туннель и до какого DC.
|
||||
///
|
||||
/// Пишется одним значением, чтобы пара в диагностике всегда была из
|
||||
/// одного соединения.
|
||||
fn note_tunnel(&self, dc: u16, route: u8) {
|
||||
self.last_tunnel
|
||||
.store(u32::from(dc) << 8 | u32::from(route), Ordering::Relaxed);
|
||||
}
|
||||
|
||||
/// Номер дата-центра для показа: из последнего поднятого туннеля, а пока
|
||||
/// туннеля не было — из последнего разобранного соединения.
|
||||
pub fn last_dc(&self) -> u16 {
|
||||
match self.last_tunnel.load(Ordering::Relaxed) {
|
||||
0 => self.seen_dc.load(Ordering::Relaxed),
|
||||
packed => (packed >> 8) as u16,
|
||||
}
|
||||
}
|
||||
|
||||
/// Маршрут последнего поднятого туннеля. См. `transport::RouteKind::ui_code`.
|
||||
pub fn last_route(&self) -> u8 {
|
||||
(self.last_tunnel.load(Ordering::Relaxed) & 0xff) as u8
|
||||
}
|
||||
|
||||
pub fn telegram_secret(&self) -> String {
|
||||
crate::mtproto::telegram_secret(&self.secret)
|
||||
}
|
||||
@@ -174,7 +217,18 @@ impl Stats {
|
||||
.filter(|value| !value.trim().is_empty())
|
||||
.map(str::to_owned)
|
||||
.collect::<Vec<_>>();
|
||||
self.transport.set_worker_domains(&domains);
|
||||
let result = self.transport.set_worker_domains(&domains);
|
||||
// Молчание здесь неотличимо от «воркер работает»: пока строка не
|
||||
// попадала в список маршрутов, об этом не сообщалось ничем, и человек
|
||||
// считал резервный маршрут настроенным (by-sonic/tglock#50).
|
||||
for rejected in &result.rejected {
|
||||
self.note(format!(
|
||||
"Cloudflare Worker «{rejected}» не похож на имя хоста — маршрут не добавлен. Нужно только имя, без https:// и без косой черты: example.workers.dev"
|
||||
));
|
||||
}
|
||||
for accepted in &result.accepted {
|
||||
self.note(format!("Cloudflare Worker в списке маршрутов: {accepted}"));
|
||||
}
|
||||
}
|
||||
|
||||
pub fn stop(&self) {
|
||||
@@ -388,7 +442,7 @@ async fn handle_socks5(
|
||||
let (dc, media) = dc_from_init(&init)
|
||||
.unwrap_or_else(|| (crate::telegram_net::dc_from_ip(ip).unwrap_or(2), false));
|
||||
|
||||
stats.last_dc.store(dc, Ordering::Relaxed);
|
||||
stats.note_dc(dc);
|
||||
|
||||
let r = ws_tunnel(s, dc, media, &init, None, stats).await;
|
||||
|
||||
@@ -486,7 +540,7 @@ async fn handle_mtproto(
|
||||
}
|
||||
};
|
||||
|
||||
stats.last_dc.store(parsed.dc, Ordering::Relaxed);
|
||||
stats.note_dc(parsed.dc);
|
||||
let result = ws_tunnel(
|
||||
stream,
|
||||
parsed.dc,
|
||||
@@ -617,57 +671,102 @@ async fn ws_tunnel(
|
||||
dc: u16,
|
||||
media: bool,
|
||||
init: &[u8; 64],
|
||||
mut crypto: Option<crate::mtproto::CryptoContext>,
|
||||
crypto: Option<crate::mtproto::CryptoContext>,
|
||||
stats: &Stats,
|
||||
) -> Result<(), Box<dyn std::error::Error + Send + Sync>> {
|
||||
use futures_util::{SinkExt, StreamExt};
|
||||
|
||||
let (mut ws, connected) = stats.transport.connect(dc, media).await?;
|
||||
let (ws, connected) = match stats.transport.connect(dc, media).await {
|
||||
Ok(connected) => connected,
|
||||
Err(error) => {
|
||||
// Единственное место, где известно, ПОЧЕМУ туннеля нет. Раньше
|
||||
// текст уходил в `Err` и там пропадал: оставался счётчик сбоев без
|
||||
// причины, и отличить «провайдер режет адреса» от «воркер отвечает
|
||||
// отказом» было нечем (by-sonic/tglock#50).
|
||||
stats.note(error.clone());
|
||||
return Err(error.into());
|
||||
}
|
||||
};
|
||||
let _tunnel = EstablishedTunnel::new(stats);
|
||||
stats
|
||||
.last_route
|
||||
.store(connected.route.kind.ui_code(), Ordering::Relaxed);
|
||||
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::<Vec<u8>>(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<dyn std::error::Error + Send + Sync>>(())
|
||||
};
|
||||
|
||||
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<dyn std::error::Error + Send + Sync>>(())
|
||||
};
|
||||
|
||||
tokio::pin!(downstream, upstream);
|
||||
tokio::select! {
|
||||
result = &mut downstream => result?,
|
||||
result = &mut upstream => result?,
|
||||
}
|
||||
Ok(())
|
||||
}
|
||||
|
||||
@@ -902,6 +1001,35 @@ mod tests {
|
||||
);
|
||||
}
|
||||
|
||||
/// Диагностика обязана показывать пару из одного соединения.
|
||||
///
|
||||
/// Пока это были два независимых поля, при десятках одновременных
|
||||
/// соединений в строку статуса попадали номер от одного и маршрут от
|
||||
/// другого. Читалось это как «до DC5 шли запасным адресом», хотя у DC5
|
||||
/// закреплённый адрес всего один и запасного не бывает вовсе
|
||||
/// (by-sonic/tglock#42).
|
||||
#[test]
|
||||
fn the_reported_data_centre_and_route_come_from_the_same_tunnel() {
|
||||
let stats = Stats::new();
|
||||
assert_eq!(stats.last_dc(), 0, "до соединений показывать нечего");
|
||||
assert_eq!(stats.last_route(), 0);
|
||||
|
||||
stats.note_dc(2);
|
||||
assert_eq!(stats.last_dc(), 2, "клиент разобран, номер известен");
|
||||
assert_eq!(stats.last_route(), 0, "а маршрут ещё не выбран");
|
||||
|
||||
stats.note_tunnel(4, 2);
|
||||
assert_eq!((stats.last_dc(), stats.last_route()), (4, 2));
|
||||
|
||||
// Ещё одно соединение до другого DC, туннеля у него пока нет.
|
||||
stats.note_dc(203);
|
||||
assert_eq!(
|
||||
(stats.last_dc(), stats.last_route()),
|
||||
(4, 2),
|
||||
"пара обязана остаться от соединения, у которого туннель был"
|
||||
);
|
||||
}
|
||||
|
||||
/// Клиент с сохранённой ссылкой от прошлого запуска. Раньше его соединение
|
||||
/// закрывалось молча, и по диагностике это было неотличимо от рабочего.
|
||||
#[tokio::test]
|
||||
@@ -1080,7 +1208,7 @@ mod tests {
|
||||
// Detection must land on MTProto, which records the data centre. The
|
||||
// SOCKS5 path would instead answer with a handshake reply.
|
||||
wait_until("the MTProto data centre to be recorded", || {
|
||||
stats.last_dc.load(Ordering::Relaxed) == 2
|
||||
stats.last_dc() == 2
|
||||
})
|
||||
.await;
|
||||
|
||||
@@ -1233,9 +1361,9 @@ mod tests {
|
||||
"Telegram must receive exactly the client's plaintext"
|
||||
);
|
||||
|
||||
assert_eq!(stats.last_dc.load(Ordering::Relaxed), 4);
|
||||
assert_eq!(stats.last_dc(), 4);
|
||||
assert_eq!(
|
||||
stats.last_route.load(Ordering::Relaxed),
|
||||
stats.last_route(),
|
||||
crate::transport::RouteKind::TelegramIp.ui_code()
|
||||
);
|
||||
assert_eq!(stats.ws_failures.load(Ordering::Relaxed), 0);
|
||||
@@ -1290,7 +1418,7 @@ mod tests {
|
||||
);
|
||||
assert_eq!(relayed, request);
|
||||
assert_eq!(
|
||||
stats.last_route.load(Ordering::Relaxed),
|
||||
stats.last_route(),
|
||||
crate::transport::RouteKind::CloudflareWorker.ui_code()
|
||||
);
|
||||
|
||||
@@ -1318,10 +1446,7 @@ mod tests {
|
||||
let mut client = TcpStream::connect(("127.0.0.1", port)).await.unwrap();
|
||||
client.write_all(&init).await.unwrap();
|
||||
|
||||
wait_until("the init to be parsed", || {
|
||||
stats.last_dc.load(Ordering::Relaxed) == 2
|
||||
})
|
||||
.await;
|
||||
wait_until("the init to be parsed", || stats.last_dc() == 2).await;
|
||||
tokio::time::sleep(Duration::from_millis(300)).await;
|
||||
|
||||
assert_eq!(
|
||||
@@ -1330,7 +1455,7 @@ mod tests {
|
||||
"a handshake still in flight must not be reported as a working tunnel"
|
||||
);
|
||||
assert_eq!(
|
||||
stats.last_route.load(Ordering::Relaxed),
|
||||
stats.last_route(),
|
||||
0,
|
||||
"no route may be announced before a tunnel is established"
|
||||
);
|
||||
@@ -1359,7 +1484,7 @@ mod tests {
|
||||
})
|
||||
.await;
|
||||
assert_eq!(
|
||||
stats.last_route.load(Ordering::Relaxed),
|
||||
stats.last_route(),
|
||||
0,
|
||||
"a route must not be reported as working when every attempt failed"
|
||||
);
|
||||
@@ -1369,6 +1494,68 @@ mod tests {
|
||||
let _ = server.await.unwrap();
|
||||
}
|
||||
|
||||
/// Диагностика обязана называть причину, а не только считать сбои.
|
||||
///
|
||||
/// При `туннелей 0` счётчик сбоев говорит, что не получилось, и молчит о
|
||||
/// том, почему. Текст ошибки собирался и выбрасывался, и разобрать
|
||||
/// «провайдер режет адреса» против «воркер отвечает отказом» было нечем
|
||||
/// (by-sonic/tglock#50).
|
||||
#[tokio::test]
|
||||
async fn a_cascade_that_failed_says_why_in_the_log() {
|
||||
let dead = TcpListener::bind("127.0.0.1:0").await.unwrap();
|
||||
let dead_port = dead.local_addr().unwrap().port();
|
||||
drop(dead);
|
||||
|
||||
let stats = Stats::new();
|
||||
stats.transport.force_local_route(dead_port);
|
||||
let (port, server) = start_proxy(stats.clone(), false).await;
|
||||
|
||||
let init = unambiguous_client_init(&stats.secret, 2);
|
||||
let mut client = TcpStream::connect(("127.0.0.1", port)).await.unwrap();
|
||||
client.write_all(&init).await.unwrap();
|
||||
|
||||
wait_until("сбой засчитан", || {
|
||||
stats.ws_failures.load(Ordering::Relaxed) > 0
|
||||
})
|
||||
.await;
|
||||
|
||||
let events = stats.drain_events();
|
||||
let named = events
|
||||
.iter()
|
||||
.find(|event| event.contains("Не поднялся туннель до DC2"))
|
||||
.unwrap_or_else(|| panic!("причина отказа не попала в журнал: {events:?}"));
|
||||
assert!(
|
||||
named.contains("127.0.0.1"),
|
||||
"в журнале должен быть назван адрес, до которого не дошли: {named}"
|
||||
);
|
||||
|
||||
stats.stop();
|
||||
let _ = server.await.unwrap();
|
||||
}
|
||||
|
||||
/// Строка, не похожая на имя хоста, отбрасывалась молча, и «воркер
|
||||
/// настроен» ничем не отличалось от «воркера нет» (by-sonic/tglock#50).
|
||||
#[test]
|
||||
fn a_worker_domain_is_confirmed_or_named_as_rejected() {
|
||||
let stats = Stats::new();
|
||||
stats.set_worker_domain("https://mine.workers.dev/, spare.workers.dev");
|
||||
let events = stats.drain_events();
|
||||
|
||||
assert!(
|
||||
events
|
||||
.iter()
|
||||
.any(|event| event.contains("https://mine.workers.dev/")
|
||||
&& event.contains("не похож на имя хоста")),
|
||||
"отвергнутый домен должен быть назван вместе с причиной: {events:?}"
|
||||
);
|
||||
assert!(
|
||||
events
|
||||
.iter()
|
||||
.any(|event| event.contains("в списке маршрутов: spare.workers.dev")),
|
||||
"принятый домен нужно подтвердить, иначе проверить нечем: {events:?}"
|
||||
);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
#[ignore = "requires live Telegram network access"]
|
||||
async fn accepts_mtproto_and_builds_live_media_tunnel() {
|
||||
@@ -1389,16 +1576,194 @@ mod tests {
|
||||
client.write_all(&init).await.unwrap();
|
||||
|
||||
tokio::time::timeout(Duration::from_secs(10), async {
|
||||
while stats.last_route.load(Ordering::Relaxed) == 0 {
|
||||
while stats.last_route() == 0 {
|
||||
tokio::task::yield_now().await;
|
||||
}
|
||||
})
|
||||
.await
|
||||
.unwrap();
|
||||
assert_eq!(stats.last_dc.load(Ordering::Relaxed), 4);
|
||||
assert_eq!(stats.last_dc(), 4);
|
||||
assert_eq!(stats.ws_failures.load(Ordering::Relaxed), 0);
|
||||
|
||||
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("клиент, который не успевает читать, всё равно должен отправлять");
|
||||
}
|
||||
}
|
||||
|
||||
+56
-18
@@ -87,6 +87,15 @@ impl Route {
|
||||
}
|
||||
}
|
||||
|
||||
/// Что случилось с доменами Worker'а, которые задал пользователь.
|
||||
#[derive(Clone, Debug, Default, Eq, PartialEq)]
|
||||
pub struct WorkerDomains {
|
||||
/// Домены, попавшие в список маршрутов.
|
||||
pub accepted: Vec<String>,
|
||||
/// Строки, не похожие на имя хоста, — маршрута из них не вышло.
|
||||
pub rejected: Vec<String>,
|
||||
}
|
||||
|
||||
#[derive(Clone, Debug)]
|
||||
pub struct ConnectedRoute {
|
||||
pub route: Route,
|
||||
@@ -151,15 +160,29 @@ impl TransportEngine {
|
||||
Self::default()
|
||||
}
|
||||
|
||||
pub fn set_worker_domains(&self, domains: &[String]) {
|
||||
let mut normalized = Vec::new();
|
||||
/// Задать домены Worker'ов, вернув принятые и отвергнутые по отдельности.
|
||||
///
|
||||
/// Отвергнутые возвращаются, потому что раньше они отбрасывались молча:
|
||||
/// вписанный со схемой или слэшем `https://name.workers.dev/` не проходил
|
||||
/// проверку, маршрут не появлялся, и «воркер настроен» ничем не отличалось
|
||||
/// от «воркера нет» (by-sonic/tglock#50).
|
||||
pub fn set_worker_domains(&self, domains: &[String]) -> WorkerDomains {
|
||||
let mut accepted = Vec::new();
|
||||
let mut rejected = Vec::new();
|
||||
for domain in domains {
|
||||
let domain = domain.trim().to_ascii_lowercase();
|
||||
if valid_domain(&domain) && !normalized.contains(&domain) {
|
||||
normalized.push(domain);
|
||||
let trimmed = domain.trim();
|
||||
if trimmed.is_empty() {
|
||||
continue;
|
||||
}
|
||||
let normalized = trimmed.to_ascii_lowercase();
|
||||
if !valid_domain(&normalized) {
|
||||
rejected.push(trimmed.to_owned());
|
||||
} else if !accepted.contains(&normalized) {
|
||||
accepted.push(normalized);
|
||||
}
|
||||
}
|
||||
*self.worker_domains.lock().unwrap() = normalized;
|
||||
*self.worker_domains.lock().unwrap() = accepted.clone();
|
||||
WorkerDomains { accepted, rejected }
|
||||
}
|
||||
|
||||
pub async fn connect(
|
||||
@@ -179,18 +202,22 @@ impl TransportEngine {
|
||||
}
|
||||
Err(error) => {
|
||||
self.record_failure(&route);
|
||||
errors.push(format!(
|
||||
"{} via {}: {}",
|
||||
route.websocket_host, route.connect_host, error
|
||||
));
|
||||
let attempt = format!("{} — {}", route.connect_host, error);
|
||||
if !errors.contains(&attempt) {
|
||||
errors.push(attempt);
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// Текст читает человек: он попадает в журнал событий, и по нему
|
||||
// отличают «провайдер режет закреплённые адреса» от «воркер отвечает
|
||||
// отказом». Раньше причина отказа не доходила никуда, и при
|
||||
// `туннелей 0` узнать, почему их ноль, было нечем (by-sonic/tglock#50).
|
||||
Err(format!(
|
||||
"all Telegram routes for DC{}{} failed: {}",
|
||||
"Не поднялся туннель до DC{}{}: {}",
|
||||
dc,
|
||||
if media { " media" } else { "" },
|
||||
if media { " (медиа)" } else { "" },
|
||||
errors.join("; ")
|
||||
))
|
||||
}
|
||||
@@ -376,8 +403,8 @@ async fn connect_route(route: &Route) -> Result<TelegramWebSocket, String> {
|
||||
TcpStream::connect((route.connect_host.as_str(), route.port)),
|
||||
)
|
||||
.await
|
||||
.map_err(|_| "TCP connect timeout".to_owned())?
|
||||
.map_err(|error| format!("TCP connect: {}", error))?;
|
||||
.map_err(|_| "не отвечает (таймаут TCP)".to_owned())?
|
||||
.map_err(|error| format!("соединение не открылось: {}", error))?;
|
||||
tcp.set_nodelay(true)
|
||||
.map_err(|error| format!("TCP_NODELAY: {}", error))?;
|
||||
|
||||
@@ -402,9 +429,9 @@ async fn connect_route(route: &Route) -> Result<TelegramWebSocket, String> {
|
||||
tokio_tungstenite::client_async(request, MaybeTlsStream::Plain(tcp)),
|
||||
)
|
||||
.await
|
||||
.map_err(|_| "WebSocket timeout".to_owned())?
|
||||
.map_err(|_| "таймаут WebSocket".to_owned())?
|
||||
.map(|(websocket, _)| websocket)
|
||||
.map_err(|error| format!("WebSocket handshake: {}", error));
|
||||
.map_err(|error| format!("рукопожатие WebSocket: {}", error));
|
||||
}
|
||||
|
||||
// The URI host remains the real Telegram hostname even when the TCP socket
|
||||
@@ -417,7 +444,7 @@ async fn connect_route(route: &Route) -> Result<TelegramWebSocket, String> {
|
||||
tokio_tungstenite::client_async_tls_with_config(request, tcp, None, Some(connector)),
|
||||
)
|
||||
.await
|
||||
.map_err(|_| "TLS/WebSocket timeout".to_owned())?
|
||||
.map_err(|_| "таймаут TLS/WebSocket".to_owned())?
|
||||
.map(|(websocket, _)| websocket)
|
||||
.map_err(|error| format!("TLS/WebSocket handshake: {}", error))
|
||||
}
|
||||
@@ -612,7 +639,7 @@ mod tests {
|
||||
#[test]
|
||||
fn worker_domains_are_rejected_unless_they_are_plain_hostnames() {
|
||||
let engine = TransportEngine::new();
|
||||
engine.set_worker_domains(&[
|
||||
let result = engine.set_worker_domains(&[
|
||||
"https://scheme.workers.dev".to_owned(),
|
||||
"with.a/path".to_owned(),
|
||||
"no-dot".to_owned(),
|
||||
@@ -633,6 +660,17 @@ mod tests {
|
||||
.filter(|route| route.kind == RouteKind::CloudflareWorker)
|
||||
.collect();
|
||||
assert_eq!(workers.len(), 1, "only the valid hostname may survive");
|
||||
assert_eq!(result.accepted, vec!["good.workers.dev".to_owned()]);
|
||||
assert!(
|
||||
result.rejected.contains(&"https://scheme.workers.dev".to_owned()),
|
||||
"отвергнутая строка обязана вернуться названной, иначе о ней некому сообщить: {:?}",
|
||||
result.rejected
|
||||
);
|
||||
assert!(
|
||||
!result.rejected.iter().any(String::is_empty),
|
||||
"пустая строка — не то, о чём стоит предупреждать: {:?}",
|
||||
result.rejected
|
||||
);
|
||||
assert_eq!(workers[0].websocket_host, "good.workers.dev");
|
||||
}
|
||||
|
||||
|
||||
+1
-1
@@ -1,7 +1,7 @@
|
||||
{
|
||||
"$schema": "https://schema.tauri.app/config/2",
|
||||
"productName": "TGLock",
|
||||
"version": "2.0.0-beta.10",
|
||||
"version": "2.0.0-beta.13",
|
||||
"identifier": "com.bysonic.tglock",
|
||||
"mainBinaryName": "tglock",
|
||||
"build": {
|
||||
|
||||
Reference in New Issue
Block a user