Compare commits

...

10 Commits

Author SHA1 Message Date
Flowseal 2a699bb91b Version bump 2026-07-31 16:53:31 +03:00
Flowseal a5b2b73dfd Updated icon 2026-07-31 16:52:17 +03:00
Flowseal aee473c9f9 CF Worker pool refactoring 2026-07-31 16:38:42 +03:00
Flowseal 41e97c6a05 hard limit workers pool 2026-07-31 14:20:58 +03:00
Flowseal ecf3d6f3a1 Fixes #1161 fixes #1155 2026-07-31 13:41:23 +03:00
Flowseal 21aaeb3aba macos build fix 2026-07-28 11:43:28 +03:00
Flowseal e0230ebda7 Version bump 2026-07-28 11:33:42 +03:00
Flowseal 47b8db18d9 Added pool refill backoff; don't refill pool on .get if ip is blocked 2026-07-28 11:08:24 +03:00
Flowseal a0545cca64 don't retry fronting connect 2026-07-28 10:55:36 +03:00
Flowseal 129a760996 fixes #1079 #1083 2026-07-28 10:41:39 +03:00
8 changed files with 154 additions and 75 deletions
+4 -2
View File
@@ -236,6 +236,8 @@ jobs:
build-macos:
runs-on: macos-latest
if: ${{ github.event.inputs.build_macos == 'true' }}
env:
CFFI_VERSION: "2.0.0"
steps:
- name: Checkout
uses: actions/checkout@v6
@@ -261,7 +263,7 @@ jobs:
--python-version 3.12 \
--implementation cp \
-d wheelhouse/arm64 \
'cffi>=2.0.0' \
"cffi==$CFFI_VERSION" \
Pillow==12.1.0 \
psutil==7.0.0
@@ -271,7 +273,7 @@ jobs:
--python-version 3.12 \
--implementation cp \
-d wheelhouse/x86_64 \
'cffi>=2.0.0' \
"cffi==$CFFI_VERSION" \
Pillow==12.1.0
python3.12 -m pip download \
BIN
View File
Binary file not shown.

Before

Width:  |  Height:  |  Size: 473 B

After

Width:  |  Height:  |  Size: 22 KiB

+4 -4
View File
@@ -4,8 +4,8 @@
# http://msdn.microsoft.com/en-us/library/ms646997.aspx
VSVersionInfo(
ffi=FixedFileInfo(
filevers=(1, 8, 1, 0),
prodvers=(1, 8, 1, 0),
filevers=(1, 9, 1, 0),
prodvers=(1, 9, 1, 0),
mask=0x3f,
flags=0x0,
OS=0x40004,
@@ -21,12 +21,12 @@ VSVersionInfo(
[
StringStruct(u'CompanyName', u'Flowseal'),
StringStruct(u'FileDescription', u'Telegram Desktop WebSocket Bridge Proxy'),
StringStruct(u'FileVersion', u'1.8.1.0'),
StringStruct(u'FileVersion', u'1.9.1.0'),
StringStruct(u'InternalName', u'TgWsProxy'),
StringStruct(u'LegalCopyright', u'Copyright (c) Flowseal. MIT License.'),
StringStruct(u'OriginalFilename', u'TgWsProxy.exe'),
StringStruct(u'ProductName', u'TG WS Proxy'),
StringStruct(u'ProductVersion', u'1.8.1.0'),
StringStruct(u'ProductVersion', u'1.9.1.0'),
]
)
]
+1 -1
View File
@@ -1,6 +1,6 @@
from .config import parse_dc_ip_list, proxy_config, coerce_domain_list
from .utils import get_link_host, build_github_opener
__version__ = "1.8.1"
__version__ = "1.9.1"
__all__ = ["__version__", "get_link_host", "proxy_config", "parse_dc_ip_list", "build_github_opener", "coerce_domain_list"]
+25 -22
View File
@@ -1,7 +1,6 @@
import asyncio
import logging
import struct
import random
from typing import List, Optional
from urllib.parse import urlencode
@@ -181,21 +180,22 @@ async def _cfproxy_worker_fallback(reader, writer, relay_init, label,
worker_domains = proxy_config.cfproxy_worker_domains
if not worker_domains:
return False
random.shuffle(worker_domains)
for worker_domain in worker_domains:
ws = None if is_test_dc else await cf_worker_pool.get(dc, worker_domain, fallback_dst)
if ws:
log.info("[%s] DC%d%s -> CF worker pool hit for %s",
label, dc, media_tag, fallback_dst)
else:
query = urlencode({
'dst': fallback_dst,
'dc': str(dc),
})
path = f'/apiws?{query}'
pooled = None if is_test_dc else await cf_worker_pool.get(
dc, fallback_dst, worker_domains)
if pooled:
ws, worker_domain = pooled
log.info("[%s] DC%d%s -> CF worker pool hit via %s for %s",
label, dc, media_tag, worker_domain, fallback_dst)
else:
query = urlencode({
'dst': fallback_dst,
'dc': str(dc),
})
path = f'/apiws?{query}'
ws = None
for worker_domain in cf_worker_pool.available_domains(worker_domains):
log.info("[%s] DC%d%s -> trying CF worker %s for %s",
label, dc, media_tag, worker_domain, fallback_dst)
@@ -203,17 +203,20 @@ async def _cfproxy_worker_fallback(reader, writer, relay_init, label,
ws = await RawWebSocket.connect(worker_domain, worker_domain,
timeout=10.0, path=path)
except Exception as exc:
cf_worker_pool.report_failure(worker_domain, exc)
log.warning("[%s] DC%d%s CF worker %s failed: %s",
label, dc, media_tag, worker_domain, repr(exc))
continue
stats.connections_cfproxy += 1
await ws.send(relay_init)
await bridge_ws_reencrypt(reader, writer, ws, label, ctx,
dc=dc, is_media=is_media,
splitter=splitter)
return True
return False
if ws is None:
return False
stats.connections_cfproxy += 1
await ws.send(relay_init)
await bridge_ws_reencrypt(reader, writer, ws, label, ctx,
dc=dc, is_media=is_media,
splitter=None)
return True
async def _cfproxy_fallback(reader, writer, relay_init, label,
@@ -422,4 +425,4 @@ async def _bridge_tcp_reencrypt(reader, writer, remote_reader, remote_writer,
w.close()
await w.wait_closed()
except BaseException:
pass
pass
+113 -44
View File
@@ -1,5 +1,6 @@
import asyncio
import logging
import random
import time
from collections import deque
@@ -16,15 +17,20 @@ log = logging.getLogger('tg-mtproto-proxy')
class _WsPool:
WS_POOL_MAX_AGE = 120.0
WS_POOL_CHECK_INTERVAL = 5.0
REFILL_BACKOFF_INITIAL = 60.0
REFILL_BACKOFF_MAX = 3600.0
def __init__(self):
self._idle: Dict[Tuple[int, bool], deque] = {}
self._refilling: Set[Tuple[int, bool]] = set()
self._rotating: Dict[Tuple[int, bool], asyncio.Task] = {}
self._refill_failures: Dict[Tuple[int, bool], int] = {}
self._refill_after: Dict[Tuple[int, bool], float] = {}
self.try_fronting_first = False
async def get(self, dc: int, is_media: bool,
target_ip: str, domains: List[str]
target_ip: str, domains: List[str],
*, allow_refill: bool = True
) -> Optional[RawWebSocket]:
key = (dc, is_media)
now = time.monotonic()
@@ -43,19 +49,28 @@ class _WsPool:
stats.pool_hits += 1
log.debug("WS pool hit DC%d%s (age=%.1fs, left=%d)",
dc, 'm' if is_media else '', age, len(bucket))
self._schedule_refill(key, target_ip, domains)
self.report_success(dc, is_media)
if allow_refill:
self._schedule_refill(key, target_ip, domains)
return ws
stats.pool_misses += 1
self._schedule_refill(key, target_ip, domains)
if allow_refill:
self._schedule_refill(key, target_ip, domains)
return None
def _schedule_refill(self, key, target_ip, domains):
if key in self._refilling:
if (key in self._refilling
or time.monotonic() < self._refill_after.get(key, 0)):
return
self._refilling.add(key)
asyncio.create_task(self._refill(key, target_ip, domains))
def report_success(self, dc: int, is_media: bool) -> None:
key = (dc, is_media)
self._refill_failures.pop(key, None)
self._refill_after.pop(key, None)
async def _refill(self, key, target_ip, domains):
dc, is_media = key
try:
@@ -63,6 +78,7 @@ class _WsPool:
needed = proxy_config.pool_size - len(bucket)
if needed <= 0:
return
connected = 0
tasks = [asyncio.create_task(
self._connect_one(target_ip, domains))
for _ in range(needed)]
@@ -71,9 +87,24 @@ class _WsPool:
ws = await t
if ws:
bucket.append((ws, time.monotonic()))
connected += 1
self._schedule_rotation(key, target_ip, domains)
except Exception:
pass
if connected:
self.report_success(dc, is_media)
else:
failures = self._refill_failures.get(key, 0) + 1
self._refill_failures[key] = failures
delay = min(
self.REFILL_BACKOFF_INITIAL
* (2 ** min(failures - 1, 6)),
self.REFILL_BACKOFF_MAX,
)
self._refill_after[key] = time.monotonic() + delay
log.info(
"WS pool refill failed for DC%d%s, retry in %.0fs",
dc, 'm' if is_media else '', delay)
log.debug("WS pool refilled DC%d%s: %d ready",
dc, 'm' if is_media else '', len(bucket))
finally:
@@ -136,6 +167,8 @@ class _WsPool:
self.try_fronting_first = False
return ws
except asyncio.TimeoutError:
if self.try_fronting_first:
return None
return await self._connect_fronted(target_ip, domain)
except WsHandshakeError as exc:
if exc.is_redirect:
@@ -179,80 +212,115 @@ class _WsPool:
self._idle.clear()
self._refilling.clear()
self._rotating.clear()
self._refill_failures.clear()
self._refill_after.clear()
self.try_fronting_first = False
class _CfWorkerPool:
WS_POOL_MAX_AGE = 100.0
PER_DC_LIMIT = 1
def __init__(self):
self._idle: Dict[Tuple[int, str], deque] = {}
self._refilling: Set[Tuple[int, str]] = set()
self._idle: Dict[int, deque] = {}
self._refilling: Set[int] = set()
self._exhausted_until: Dict[str, float] = {}
async def get(self, dc: int, worker_domain: str, fallback_dst: str) -> Optional[RawWebSocket]:
async def get(self, dc: int, fallback_dst: str,
worker_domains: List[str]
) -> Optional[Tuple[RawWebSocket, str]]:
now = time.monotonic()
key = (dc, worker_domain)
bucket = self._idle.get(key)
bucket = self._idle.get(dc)
if bucket is None:
bucket = deque()
self._idle[key] = bucket
self._idle[dc] = bucket
while bucket:
ws, created = bucket.popleft()
ws, created, worker_domain = bucket.popleft()
age = now - created
if (age > self.WS_POOL_MAX_AGE or ws._closed
or ws.writer.transport.is_closing()):
asyncio.create_task(self._quiet_close(ws))
continue
stats.cf_pool_hits += 1
log.debug("CF worker pool hit DC%d (age=%.1fs, left=%d)",
dc, age, len(bucket))
self._schedule_refill(key, fallback_dst)
return ws
log.debug(
"CF worker pool hit DC%d via %s (age=%.1fs, left=%d)",
dc, worker_domain, age, len(bucket))
self._schedule_refill(dc, fallback_dst, worker_domains)
return ws, worker_domain
stats.cf_pool_misses += 1
self._schedule_refill(key, fallback_dst)
return None
def _schedule_refill(self, key, fallback_dst):
if key in self._refilling:
def _schedule_refill(self, dc, fallback_dst, worker_domains):
if dc in self._refilling:
return
self._refilling.add(key)
asyncio.create_task(self._refill(key, fallback_dst))
self._refilling.add(dc)
asyncio.create_task(self._refill(
dc, fallback_dst, list(worker_domains)))
async def _refill(self, key, fallback_dst):
dc, worker_domain = key
async def _refill(self, dc, fallback_dst, worker_domains):
try:
bucket = self._idle.setdefault(key, deque())
needed = proxy_config.pool_size - len(bucket)
bucket = self._idle.setdefault(dc, deque())
target_size = min(proxy_config.pool_size, self.PER_DC_LIMIT)
needed = target_size - len(bucket)
if needed <= 0:
return
tasks = [asyncio.create_task(
self._connect_one(worker_domain, fallback_dst, dc))
for _ in range(needed)]
for t in tasks:
try:
ws = await t
if ws:
bucket.append((ws, time.monotonic()))
except Exception:
pass
for _ in range(needed):
connected = await self._connect_one(
worker_domains, fallback_dst, dc)
if connected is None:
break
ws, worker_domain = connected
bucket.append((ws, time.monotonic(), worker_domain))
log.debug("CF worker pool refilled DC%d: %d ready",
dc, len(bucket))
finally:
self._refilling.discard(key)
self._refilling.discard(dc)
async def _connect_one(self, worker_domain, fallback_dst, dc) -> Optional[RawWebSocket]:
async def _connect_one(self, worker_domains, fallback_dst, dc):
query = urlencode({
'dst': fallback_dst,
'dc': str(dc),
})
path = f'/apiws?{query}'
try:
return await RawWebSocket.connect(
worker_domain, worker_domain, timeout=8, path=path)
except Exception:
return None
for worker_domain in self.available_domains(worker_domains):
try:
ws = await RawWebSocket.connect(
worker_domain, worker_domain, timeout=8, path=path)
return ws, worker_domain
except Exception as exc:
self.report_failure(worker_domain, exc)
return None
def available_domains(self, worker_domains: List[str]) -> List[str]:
now = time.time()
domains = list()
for domain in worker_domains:
if domain in domains:
continue
exhausted_until = self._exhausted_until.get(domain, 0)
if exhausted_until > now:
continue
if exhausted_until:
self._exhausted_until.pop(domain, None)
domains.append(domain)
random.shuffle(domains)
return domains
def report_failure(self, worker_domain: str, exc: Exception) -> None:
return # TODO: check status code after daily limit reached
if not isinstance(exc, WsHandshakeError) or exc.status_code != 429:
return
now = time.time()
if self._exhausted_until.get(worker_domain, 0) > now:
return
exhausted_until = now + (86400 - (now % 86400))
self._exhausted_until[worker_domain] = exhausted_until
log.warning(
"CF worker %s reached its request limit, disabled for %d seconds", worker_domain, int(exhausted_until - now))
async def _quiet_close(self, ws):
try:
@@ -269,15 +337,16 @@ class _CfWorkerPool:
if not cf_fallbacks or not proxy_config.cfproxy_worker_domains:
return
for worker_domain in proxy_config.cfproxy_worker_domains:
for dc, fallback_dst in cf_fallbacks.items():
self._schedule_refill((dc, worker_domain), fallback_dst)
worker_domains = list(proxy_config.cfproxy_worker_domains)
for dc, fallback_dst in cf_fallbacks.items():
self._schedule_refill(dc, fallback_dst, worker_domains)
log.info("CF worker pool warmup started for %d DC(s)", len(cf_fallbacks))
def reset(self):
self._idle.clear()
self._refilling.clear()
self._exhausted_until.clear()
ws_pool = _WsPool()
+6 -1
View File
@@ -339,7 +339,11 @@ async def _handle_client(reader, writer, secret: bytes):
ws_timed_out = False
all_redirects = True
ws = await ws_pool.get(dc, is_media, target, domains) if not is_test_dc else None
allow_pool_refill = now >= ip_fail_until.get(target, 0)
ws = await ws_pool.get(
dc, is_media, target, domains,
allow_refill=allow_pool_refill,
) if not is_test_dc else None
if ws:
log.info("[%s] DC%d%s -> pool hit via %s",
label, dc, media_tag, target)
@@ -413,6 +417,7 @@ async def _handle_client(reader, writer, secret: bytes):
dc_fail_until.pop(dc_key, None)
ip_fail_until.pop(target, None)
ws_pool.report_success(dc, is_media)
stats.connections_ws += 1
splitter = None
+1 -1
View File
@@ -525,7 +525,7 @@ def _edit_config_dialog() -> None:
merged["force_test_dc"] = _config.get("force_test_dc", DEFAULT_CONFIG["force_test_dc"])
_ui_only_keys = {"appearance", "autostart", "check_updates", "language"}
config_changed = any(merged.get(k) != _config.get(k) for k in merged)
config_changed = any(merged.get(k) != cfg.get(k) for k in merged)
proxy_changed = any(merged.get(k) != _config.get(k) for k in merged if k not in _ui_only_keys)
if not config_changed: