mirror of
https://github.com/by-sonic/tglock.git
synced 2026-09-05 18:16:09 +03:00
Compare commits
6 Commits
v2.0.0-beta.11
...
main
| Author | SHA1 | Date | |
|---|---|---|---|
| 8617d25f3a | |||
| e9af114f3c | |||
| 935121b991 | |||
| fe9aad5ee9 | |||
| 28085f9127 | |||
| 923c22f9b4 |
Generated
+1
-1
@@ -3553,7 +3553,7 @@ dependencies = [
|
|||||||
|
|
||||||
[[package]]
|
[[package]]
|
||||||
name = "tglock"
|
name = "tglock"
|
||||||
version = "2.0.0-beta.11"
|
version = "2.0.0-beta.14"
|
||||||
dependencies = [
|
dependencies = [
|
||||||
"aes",
|
"aes",
|
||||||
"cipher",
|
"cipher",
|
||||||
|
|||||||
+7
-1
@@ -1,6 +1,6 @@
|
|||||||
[package]
|
[package]
|
||||||
name = "tglock"
|
name = "tglock"
|
||||||
version = "2.0.0-beta.11"
|
version = "2.0.0-beta.14"
|
||||||
edition = "2021"
|
edition = "2021"
|
||||||
rust-version = "1.88"
|
rust-version = "1.88"
|
||||||
description = "Telegram unblock via local WebSocket tunnel"
|
description = "Telegram unblock via local WebSocket tunnel"
|
||||||
@@ -65,3 +65,9 @@ rand = "0.8"
|
|||||||
|
|
||||||
[build-dependencies]
|
[build-dependencies]
|
||||||
tauri-build = { version = "2", features = [], optional = true }
|
tauri-build = { version = "2", features = [], optional = true }
|
||||||
|
|
||||||
|
[dev-dependencies]
|
||||||
|
# `start_paused` в тестах: таймаут ожидания запроса от клиента — десять секунд,
|
||||||
|
# и ждать их по-настоящему в тесте нельзя. В сборку не попадает: dev-зависимости
|
||||||
|
# участвуют только в тестах.
|
||||||
|
tokio = { version = "1", features = ["test-util"] }
|
||||||
|
|||||||
@@ -116,6 +116,64 @@ LAN-режим превращал бы машину в открытый прок
|
|||||||
Счётчик `ws_failures` от них отличается тем, что растёт после успешного
|
Счётчик `ws_failures` от них отличается тем, что растёт после успешного
|
||||||
рукопожатия с клиентом: там договорились с клиентом, но не смогли с Telegram.
|
рукопожатия с клиентом: там договорились с клиентом, но не смогли с 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` считает **установленные** туннели: счётчик поднимается после
|
`Stats::ws` считает **установленные** туннели: счётчик поднимается после
|
||||||
@@ -175,6 +233,14 @@ Worker должен принимать WebSocket на:
|
|||||||
добавить собственную авторизацию до стабильного релиза; поэтому Worker
|
добавить собственную авторизацию до стабильного релиза; поэтому Worker
|
||||||
остаётся расширенной опцией alpha-версии.
|
остаётся расширенной опцией alpha-версии.
|
||||||
|
|
||||||
|
Запись в сокет Telegram обязана быть последовательной: следующий чанк уходит
|
||||||
|
после того, как записан предыдущий, и только когда писатель к этому готов
|
||||||
|
(`writer.ready`). В `worker/tglock-worker.js` этого не было — `write()`
|
||||||
|
вызывался поверх незавершённого, без backpressure. Пока в клиенте голодала
|
||||||
|
отправка, настоящего потока вверх через воркер не возникало и это не
|
||||||
|
проявлялось; после того как голодание починили, поток появился. Кто разворачивал
|
||||||
|
воркер раньше — обновите скрипт.
|
||||||
|
|
||||||
## Current limitations
|
## Current limitations
|
||||||
|
|
||||||
- SNI camouflage не включена: небезопасное отключение hostname verification
|
- SNI camouflage не включена: небезопасное отключение hostname verification
|
||||||
@@ -197,5 +263,9 @@ Worker должен принимать WebSocket на:
|
|||||||
CLI собирается одной командой. Подробно — в разделе README про антивирус.
|
CLI собирается одной командой. Подробно — в разделе README про антивирус.
|
||||||
- Работоспособность медиа зависит от конкретного DC аккаунта и доступности
|
- Работоспособность медиа зависит от конкретного DC аккаунта и доступности
|
||||||
Telegram/Cloudflare у провайдера.
|
Telegram/Cloudflare у провайдера.
|
||||||
|
- Соединение, открытое клиентом и молчащее дольше `IO_TIMEOUT` (10 секунд),
|
||||||
|
закрывается. Для клиента, открывающего соединения про запас, это норма; счётчик
|
||||||
|
«промолчали» показывает, как часто это происходит, — раньше такие соединения
|
||||||
|
не попадали никуда.
|
||||||
- Пулы заранее открытых WebSocket-соединений будут добавлены после измерения,
|
- Пулы заранее открытых WebSocket-соединений будут добавлены после измерения,
|
||||||
что они не создают лишнюю нагрузку и не ухудшают стабильность.
|
что они не создают лишнюю нагрузку и не ухудшают стабильность.
|
||||||
|
|||||||
@@ -41,6 +41,8 @@
|
|||||||
|
|
||||||
Если вернулось `not found` — проверь, что путь именно `/apiws`. Если ошибка про `cloudflare:sockets` — у воркера слишком старая дата совместимости, поставь в **Settings → Compatibility date** сегодняшнюю.
|
Если вернулось `not found` — проверь, что путь именно `/apiws`. Если ошибка про `cloudflare:sockets` — у воркера слишком старая дата совместимости, поставь в **Settings → Compatibility date** сегодняшнюю.
|
||||||
|
|
||||||
|
> **Разворачивал воркер до 2.0.0-beta.14 — обнови скрипт.** В прежней версии запись в сокет Telegram шла без ожидания предыдущей и без backpressure. Пока в клиенте голодала отправка, через воркер не проходило настоящего потока вверх и это не проявлялось; после того как голодание починили в beta.12, поток появился.
|
||||||
|
|
||||||
## Подключение в TGLock
|
## Подключение в TGLock
|
||||||
|
|
||||||
**В приложении:** Настройки → поле **Cloudflare Worker** → вставь `tglock.имя.workers.dev` → Сохранить. Настройки меняются только при выключенной защите.
|
**В приложении:** Настройки → поле **Cloudflare Worker** → вставь `tglock.имя.workers.dev` → Сохранить. Настройки меняются только при выключенной защите.
|
||||||
|
|||||||
+1
-1
@@ -1,7 +1,7 @@
|
|||||||
{
|
{
|
||||||
"name": "tglock-ui",
|
"name": "tglock-ui",
|
||||||
"private": true,
|
"private": true,
|
||||||
"version": "2.0.0-beta.11",
|
"version": "2.0.0-beta.14",
|
||||||
"type": "module",
|
"type": "module",
|
||||||
"scripts": {
|
"scripts": {
|
||||||
"dev": "vite --port 1420",
|
"dev": "vite --port 1420",
|
||||||
|
|||||||
+4
-2
@@ -224,14 +224,16 @@ async fn watch_status(stats: Arc<proxy::Stats>) {
|
|||||||
stats.route_failures(),
|
stats.route_failures(),
|
||||||
stats.blocked.load(Ordering::Relaxed),
|
stats.blocked.load(Ordering::Relaxed),
|
||||||
stats.unknown_clients.load(Ordering::Relaxed),
|
stats.unknown_clients.load(Ordering::Relaxed),
|
||||||
|
stats.silent_clients.load(Ordering::Relaxed),
|
||||||
);
|
);
|
||||||
if previous.as_ref() == Some(¤t) {
|
if previous.as_ref() == Some(¤t) {
|
||||||
continue;
|
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!(
|
let line = format!(
|
||||||
"соединений {active} · туннелей {tunnels} · {} · {} · сбоев {failures} · \
|
"соединений {active} · туннелей {tunnels} · {} · {} · сбоев {failures} · \
|
||||||
падений маршрутов {route_failures} · отклонено {blocked} · не опознано {unknown}",
|
падений маршрутов {route_failures} · отклонено {blocked} · не опознано {unknown} · промолчали {silent}",
|
||||||
if dc > 0 {
|
if dc > 0 {
|
||||||
format!("DC{dc}")
|
format!("DC{dc}")
|
||||||
} else {
|
} else {
|
||||||
|
|||||||
@@ -56,6 +56,12 @@ struct StatusSnapshot {
|
|||||||
/// Клиенты, которые дошли, но не сумели договориться. Почти всегда это
|
/// Клиенты, которые дошли, но не сумели договориться. Почти всегда это
|
||||||
/// ссылка `tg://proxy` от прошлого запуска, то есть другой секрет.
|
/// ссылка `tg://proxy` от прошлого запуска, то есть другой секрет.
|
||||||
unknown_clients: u32,
|
unknown_clients: u32,
|
||||||
|
/// Соединения, которые открылись и ничего не прислали до таймаута.
|
||||||
|
///
|
||||||
|
/// Растущее число в такт с переподключениями клиента означает, что он
|
||||||
|
/// открывает соединения впрок, а мы закрываем их по таймауту
|
||||||
|
/// (by-sonic/tglock#42).
|
||||||
|
silent_clients: u32,
|
||||||
uptime_seconds: u64,
|
uptime_seconds: u64,
|
||||||
port: u16,
|
port: u16,
|
||||||
/// Адрес, который нужно вписать в Telegram на другом устройстве.
|
/// Адрес, который нужно вписать в Telegram на другом устройстве.
|
||||||
@@ -128,6 +134,7 @@ impl AppState {
|
|||||||
route_failures: self.stats.route_failures(),
|
route_failures: self.stats.route_failures(),
|
||||||
blocked: self.stats.blocked.load(Ordering::Relaxed),
|
blocked: self.stats.blocked.load(Ordering::Relaxed),
|
||||||
unknown_clients: self.stats.unknown_clients.load(Ordering::Relaxed),
|
unknown_clients: self.stats.unknown_clients.load(Ordering::Relaxed),
|
||||||
|
silent_clients: self.stats.silent_clients.load(Ordering::Relaxed),
|
||||||
uptime_seconds: self
|
uptime_seconds: self
|
||||||
.started_at
|
.started_at
|
||||||
.lock()
|
.lock()
|
||||||
|
|||||||
@@ -43,6 +43,49 @@ impl CryptoContext {
|
|||||||
self.telegram_decrypt.apply_keystream(data);
|
self.telegram_decrypt.apply_keystream(data);
|
||||||
self.client_encrypt.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] {
|
pub fn generate_secret() -> [u8; 16] {
|
||||||
|
|||||||
+514
-35
@@ -16,6 +16,8 @@ const SOCKS5_VERSION: u8 = 0x05;
|
|||||||
/// byte as the start of a SOCKS5 greeting.
|
/// byte as the start of a SOCKS5 greeting.
|
||||||
const PROTOCOL_PROBE_TIMEOUT: Duration = Duration::from_millis(250);
|
const PROTOCOL_PROBE_TIMEOUT: Duration = Duration::from_millis(250);
|
||||||
const PROTOCOL_PROBE_INTERVAL: Duration = Duration::from_millis(5);
|
const PROTOCOL_PROBE_INTERVAL: Duration = Duration::from_millis(5);
|
||||||
|
/// Сколько неотвеченных Ping'ов держать, пока отправляющая половина занята.
|
||||||
|
const PONG_QUEUE: usize = 4;
|
||||||
|
|
||||||
pub struct Stats {
|
pub struct Stats {
|
||||||
pub running: AtomicBool,
|
pub running: AtomicBool,
|
||||||
@@ -39,6 +41,15 @@ pub struct Stats {
|
|||||||
/// прошлого запуска. Такое соединение закрывалось молча, и по диагностике
|
/// прошлого запуска. Такое соединение закрывалось молча, и по диагностике
|
||||||
/// отличить его от рабочего было нельзя.
|
/// отличить его от рабочего было нельзя.
|
||||||
pub unknown_clients: AtomicU32,
|
pub unknown_clients: AtomicU32,
|
||||||
|
/// Сколько соединений открылось и ничего не прислало за `IO_TIMEOUT`.
|
||||||
|
///
|
||||||
|
/// Такое соединение закрывается по таймауту и до сих пор не попадало ни в
|
||||||
|
/// один счётчик: `unknown_clients` растёт, только когда запрос пришёл и не
|
||||||
|
/// разобрался, а не когда его не дождались. По диагностике это выглядело
|
||||||
|
/// как соединение без туннеля и ничего больше — а репортёр #42 видел, что
|
||||||
|
/// его телефон переустанавливает соединение примерно раз в десять секунд,
|
||||||
|
/// то есть ровно с этим периодом.
|
||||||
|
pub silent_clients: AtomicU32,
|
||||||
/// DC и маршрут последнего поднятого туннеля, упакованные в одно значение.
|
/// DC и маршрут последнего поднятого туннеля, упакованные в одно значение.
|
||||||
///
|
///
|
||||||
/// Раньше это были два независимых поля: номер писало соединение при
|
/// Раньше это были два независимых поля: номер писало соединение при
|
||||||
@@ -111,6 +122,7 @@ impl Stats {
|
|||||||
ws_failures: AtomicU32::new(0),
|
ws_failures: AtomicU32::new(0),
|
||||||
blocked: AtomicU32::new(0),
|
blocked: AtomicU32::new(0),
|
||||||
unknown_clients: AtomicU32::new(0),
|
unknown_clients: AtomicU32::new(0),
|
||||||
|
silent_clients: AtomicU32::new(0),
|
||||||
last_tunnel: AtomicU32::new(0),
|
last_tunnel: AtomicU32::new(0),
|
||||||
transport: crate::transport::TransportEngine::new(),
|
transport: crate::transport::TransportEngine::new(),
|
||||||
secret,
|
secret,
|
||||||
@@ -156,6 +168,25 @@ impl Stats {
|
|||||||
self.note(format!("Клиент {who}: {reason}"));
|
self.note(format!("Клиент {who}: {reason}"));
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/// Клиент открыл соединение и не сказал ничего.
|
||||||
|
///
|
||||||
|
/// Отличается от `note_unknown_client` тем, что там запрос пришёл и не
|
||||||
|
/// разобрался, а здесь его не дождались. Для клиента, который открывает
|
||||||
|
/// соединения про запас, это норма; для клиента, который переоткрывает их
|
||||||
|
/// в такт с таймаутом, — нет, и разницу видно только по счётчику
|
||||||
|
/// (by-sonic/tglock#42).
|
||||||
|
fn note_silent_client(&self, peer: Option<SocketAddr>, 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-режима: если строки нет,
|
/// Это первое, что нужно знать при разборе LAN-режима: если строки нет,
|
||||||
@@ -215,7 +246,18 @@ impl Stats {
|
|||||||
.filter(|value| !value.trim().is_empty())
|
.filter(|value| !value.trim().is_empty())
|
||||||
.map(str::to_owned)
|
.map(str::to_owned)
|
||||||
.collect::<Vec<_>>();
|
.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) {
|
pub fn stop(&self) {
|
||||||
@@ -336,9 +378,13 @@ async fn detect_protocol(
|
|||||||
stats: &Stats,
|
stats: &Stats,
|
||||||
) -> Result<Protocol, Box<dyn std::error::Error + Send + Sync>> {
|
) -> Result<Protocol, Box<dyn std::error::Error + Send + Sync>> {
|
||||||
let mut probe = [0; INIT_LEN];
|
let mut probe = [0; INIT_LEN];
|
||||||
let peeked = tokio::time::timeout(IO_TIMEOUT, stream.peek(&mut probe[..1]))
|
let peeked = match tokio::time::timeout(IO_TIMEOUT, stream.peek(&mut probe[..1])).await {
|
||||||
.await
|
Ok(peeked) => peeked?,
|
||||||
.map_err(|_| "client protocol detection timeout")??;
|
Err(_) => {
|
||||||
|
stats.note_silent_client(stream.peer_addr().ok(), "не прислал ни байта");
|
||||||
|
return Err("client protocol detection timeout".into());
|
||||||
|
}
|
||||||
|
};
|
||||||
if peeked == 0 {
|
if peeked == 0 {
|
||||||
return Ok(Protocol::Empty);
|
return Ok(Protocol::Empty);
|
||||||
}
|
}
|
||||||
@@ -509,9 +555,15 @@ async fn handle_mtproto(
|
|||||||
stream.set_nodelay(true)?;
|
stream.set_nodelay(true)?;
|
||||||
let peer = stream.peer_addr().ok();
|
let peer = stream.peer_addr().ok();
|
||||||
let mut init = [0; 64];
|
let mut init = [0; 64];
|
||||||
tokio::time::timeout(IO_TIMEOUT, stream.read_exact(&mut init))
|
match tokio::time::timeout(IO_TIMEOUT, stream.read_exact(&mut init)).await {
|
||||||
.await
|
Ok(result) => {
|
||||||
.map_err(|_| "MTProto init timeout")??;
|
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) {
|
let parsed = match crate::mtproto::parse_client_init(&init, &stats.secret) {
|
||||||
Some(parsed) => parsed,
|
Some(parsed) => parsed,
|
||||||
None => {
|
None => {
|
||||||
@@ -658,55 +710,102 @@ async fn ws_tunnel(
|
|||||||
dc: u16,
|
dc: u16,
|
||||||
media: bool,
|
media: bool,
|
||||||
init: &[u8; 64],
|
init: &[u8; 64],
|
||||||
mut crypto: Option<crate::mtproto::CryptoContext>,
|
crypto: Option<crate::mtproto::CryptoContext>,
|
||||||
stats: &Stats,
|
stats: &Stats,
|
||||||
) -> Result<(), Box<dyn std::error::Error + Send + Sync>> {
|
) -> Result<(), Box<dyn std::error::Error + Send + Sync>> {
|
||||||
use futures_util::{SinkExt, StreamExt};
|
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);
|
let _tunnel = EstablishedTunnel::new(stats);
|
||||||
stats.note_tunnel(dc, connected.route.kind.ui_code());
|
stats.note_tunnel(dc, connected.route.kind.ui_code());
|
||||||
|
|
||||||
let (mut tcp_r, mut tcp_w) = tokio::io::split(tcp);
|
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
|
// 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! {
|
// `select!`, и любое ожидание внутри него останавливало вторую половину:
|
||||||
biased;
|
// непрерывная загрузка не давала опросить клиента вообще, а клиент,
|
||||||
|
// который не успевал разбирать входящий поток, замораживал заодно и свою
|
||||||
msg = ws.next() => match msg {
|
// отправку. Telegram при этом ждёт от клиента подтверждений — без них
|
||||||
Some(Ok(tungstenite::Message::Binary(mut data))) => {
|
// сессия встаёт при живом туннеле (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 {
|
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.write_all(data.as_ref()).await?;
|
||||||
tcp_w.flush().await?;
|
tcp_w.flush().await?;
|
||||||
}
|
}
|
||||||
Some(Ok(tungstenite::Message::Ping(p))) => {
|
Ok(tungstenite::Message::Ping(payload)) => {
|
||||||
let _ = ws.send(tungstenite::Message::Pong(p)).await;
|
if pong_tx.send(payload).await.is_err() {
|
||||||
}
|
break;
|
||||||
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]);
|
|
||||||
}
|
}
|
||||||
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(())
|
Ok(())
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -1434,6 +1533,208 @@ mod tests {
|
|||||||
let _ = server.await.unwrap();
|
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:?}"
|
||||||
|
);
|
||||||
|
}
|
||||||
|
|
||||||
|
/// 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::<Result<Vec<u8>, 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]
|
#[tokio::test]
|
||||||
#[ignore = "requires live Telegram network access"]
|
#[ignore = "requires live Telegram network access"]
|
||||||
async fn accepts_mtproto_and_builds_live_media_tunnel() {
|
async fn accepts_mtproto_and_builds_live_media_tunnel() {
|
||||||
@@ -1466,4 +1767,182 @@ mod tests {
|
|||||||
stats.stop();
|
stats.stop();
|
||||||
server.await.unwrap().unwrap();
|
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)]
|
#[derive(Clone, Debug)]
|
||||||
pub struct ConnectedRoute {
|
pub struct ConnectedRoute {
|
||||||
pub route: Route,
|
pub route: Route,
|
||||||
@@ -151,15 +160,29 @@ impl TransportEngine {
|
|||||||
Self::default()
|
Self::default()
|
||||||
}
|
}
|
||||||
|
|
||||||
pub fn set_worker_domains(&self, domains: &[String]) {
|
/// Задать домены Worker'ов, вернув принятые и отвергнутые по отдельности.
|
||||||
let mut normalized = Vec::new();
|
///
|
||||||
|
/// Отвергнутые возвращаются, потому что раньше они отбрасывались молча:
|
||||||
|
/// вписанный со схемой или слэшем `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 {
|
for domain in domains {
|
||||||
let domain = domain.trim().to_ascii_lowercase();
|
let trimmed = domain.trim();
|
||||||
if valid_domain(&domain) && !normalized.contains(&domain) {
|
if trimmed.is_empty() {
|
||||||
normalized.push(domain);
|
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(
|
pub async fn connect(
|
||||||
@@ -179,18 +202,22 @@ impl TransportEngine {
|
|||||||
}
|
}
|
||||||
Err(error) => {
|
Err(error) => {
|
||||||
self.record_failure(&route);
|
self.record_failure(&route);
|
||||||
errors.push(format!(
|
let attempt = format!("{} — {}", route.connect_host, error);
|
||||||
"{} via {}: {}",
|
if !errors.contains(&attempt) {
|
||||||
route.websocket_host, route.connect_host, error
|
errors.push(attempt);
|
||||||
));
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// Текст читает человек: он попадает в журнал событий, и по нему
|
||||||
|
// отличают «провайдер режет закреплённые адреса» от «воркер отвечает
|
||||||
|
// отказом». Раньше причина отказа не доходила никуда, и при
|
||||||
|
// `туннелей 0` узнать, почему их ноль, было нечем (by-sonic/tglock#50).
|
||||||
Err(format!(
|
Err(format!(
|
||||||
"all Telegram routes for DC{}{} failed: {}",
|
"Не поднялся туннель до DC{}{}: {}",
|
||||||
dc,
|
dc,
|
||||||
if media { " media" } else { "" },
|
if media { " (медиа)" } else { "" },
|
||||||
errors.join("; ")
|
errors.join("; ")
|
||||||
))
|
))
|
||||||
}
|
}
|
||||||
@@ -376,8 +403,8 @@ async fn connect_route(route: &Route) -> Result<TelegramWebSocket, String> {
|
|||||||
TcpStream::connect((route.connect_host.as_str(), route.port)),
|
TcpStream::connect((route.connect_host.as_str(), route.port)),
|
||||||
)
|
)
|
||||||
.await
|
.await
|
||||||
.map_err(|_| "TCP connect timeout".to_owned())?
|
.map_err(|_| "не отвечает (таймаут TCP)".to_owned())?
|
||||||
.map_err(|error| format!("TCP connect: {}", error))?;
|
.map_err(|error| format!("соединение не открылось: {}", error))?;
|
||||||
tcp.set_nodelay(true)
|
tcp.set_nodelay(true)
|
||||||
.map_err(|error| format!("TCP_NODELAY: {}", error))?;
|
.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)),
|
tokio_tungstenite::client_async(request, MaybeTlsStream::Plain(tcp)),
|
||||||
)
|
)
|
||||||
.await
|
.await
|
||||||
.map_err(|_| "WebSocket timeout".to_owned())?
|
.map_err(|_| "таймаут WebSocket".to_owned())?
|
||||||
.map(|(websocket, _)| websocket)
|
.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
|
// 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)),
|
tokio_tungstenite::client_async_tls_with_config(request, tcp, None, Some(connector)),
|
||||||
)
|
)
|
||||||
.await
|
.await
|
||||||
.map_err(|_| "TLS/WebSocket timeout".to_owned())?
|
.map_err(|_| "таймаут TLS/WebSocket".to_owned())?
|
||||||
.map(|(websocket, _)| websocket)
|
.map(|(websocket, _)| websocket)
|
||||||
.map_err(|error| format!("TLS/WebSocket handshake: {}", error))
|
.map_err(|error| format!("TLS/WebSocket handshake: {}", error))
|
||||||
}
|
}
|
||||||
@@ -612,7 +639,7 @@ mod tests {
|
|||||||
#[test]
|
#[test]
|
||||||
fn worker_domains_are_rejected_unless_they_are_plain_hostnames() {
|
fn worker_domains_are_rejected_unless_they_are_plain_hostnames() {
|
||||||
let engine = TransportEngine::new();
|
let engine = TransportEngine::new();
|
||||||
engine.set_worker_domains(&[
|
let result = engine.set_worker_domains(&[
|
||||||
"https://scheme.workers.dev".to_owned(),
|
"https://scheme.workers.dev".to_owned(),
|
||||||
"with.a/path".to_owned(),
|
"with.a/path".to_owned(),
|
||||||
"no-dot".to_owned(),
|
"no-dot".to_owned(),
|
||||||
@@ -633,6 +660,17 @@ mod tests {
|
|||||||
.filter(|route| route.kind == RouteKind::CloudflareWorker)
|
.filter(|route| route.kind == RouteKind::CloudflareWorker)
|
||||||
.collect();
|
.collect();
|
||||||
assert_eq!(workers.len(), 1, "only the valid hostname may survive");
|
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");
|
assert_eq!(workers[0].websocket_host, "good.workers.dev");
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
+1
-1
@@ -1,7 +1,7 @@
|
|||||||
{
|
{
|
||||||
"$schema": "https://schema.tauri.app/config/2",
|
"$schema": "https://schema.tauri.app/config/2",
|
||||||
"productName": "TGLock",
|
"productName": "TGLock",
|
||||||
"version": "2.0.0-beta.11",
|
"version": "2.0.0-beta.14",
|
||||||
"identifier": "com.bysonic.tglock",
|
"identifier": "com.bysonic.tglock",
|
||||||
"mainBinaryName": "tglock",
|
"mainBinaryName": "tglock",
|
||||||
"build": {
|
"build": {
|
||||||
|
|||||||
@@ -15,6 +15,8 @@ type Status = {
|
|||||||
blocked: number;
|
blocked: number;
|
||||||
/// Клиенты, которые дошли, но не сумели договориться о рукопожатии.
|
/// Клиенты, которые дошли, но не сумели договориться о рукопожатии.
|
||||||
unknownClients: number;
|
unknownClients: number;
|
||||||
|
/// Соединения, которые открылись и ничего не прислали до таймаута.
|
||||||
|
silentClients: number;
|
||||||
uptimeSeconds: number;
|
uptimeSeconds: number;
|
||||||
port: number;
|
port: number;
|
||||||
/// Адрес для других устройств. Приходит только в LAN-режиме.
|
/// Адрес для других устройств. Приходит только в LAN-режиме.
|
||||||
@@ -47,6 +49,7 @@ let status: Status = {
|
|||||||
routeFailures: 0,
|
routeFailures: 0,
|
||||||
blocked: 0,
|
blocked: 0,
|
||||||
unknownClients: 0,
|
unknownClients: 0,
|
||||||
|
silentClients: 0,
|
||||||
uptimeSeconds: 0,
|
uptimeSeconds: 0,
|
||||||
port: 1080,
|
port: 1080,
|
||||||
shareAddress: null,
|
shareAddress: null,
|
||||||
@@ -297,6 +300,10 @@ function renderDiagnostics(): void {
|
|||||||
<span>Не опознаны</span>
|
<span>Не опознаны</span>
|
||||||
<strong>${status.unknownClients}</strong>
|
<strong>${status.unknownClients}</strong>
|
||||||
</article>
|
</article>
|
||||||
|
<article class="metric-card">
|
||||||
|
<span>Промолчали</span>
|
||||||
|
<strong>${status.silentClients}</strong>
|
||||||
|
</article>
|
||||||
</div>
|
</div>
|
||||||
|
|
||||||
<p class="field-hint">
|
<p class="field-hint">
|
||||||
|
|||||||
+14
-1
@@ -68,12 +68,25 @@ export default {
|
|||||||
}
|
}
|
||||||
};
|
};
|
||||||
|
|
||||||
|
// Запись сериализуется: следующий чанк уходит только после того, как
|
||||||
|
// записан предыдущий, и только когда писатель к этому готов.
|
||||||
|
//
|
||||||
|
// Раньше `write()` вызывался поверх незавершённого, а `writer.ready` не
|
||||||
|
// спрашивался вовсе — backpressure не применялся. Пока в клиенте отправка
|
||||||
|
// голодала, поверх воркера настоящего потока вверх не бывало и это не
|
||||||
|
// проявлялось. Как только голодание починили, в воркер пошёл настоящий
|
||||||
|
// поток (by-sonic/tglock#42).
|
||||||
|
let pending = Promise.resolve();
|
||||||
|
|
||||||
server.addEventListener("message", (event) => {
|
server.addEventListener("message", (event) => {
|
||||||
const chunk =
|
const chunk =
|
||||||
event.data instanceof ArrayBuffer
|
event.data instanceof ArrayBuffer
|
||||||
? new Uint8Array(event.data)
|
? new Uint8Array(event.data)
|
||||||
: 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("close", shutdown);
|
||||||
server.addEventListener("error", shutdown);
|
server.addEventListener("error", shutdown);
|
||||||
|
|||||||
Reference in New Issue
Block a user