mirror of
https://github.com/Flowseal/tg-ws-proxy.git
synced 2026-08-01 16:40:33 +03:00
Compare commits
10 Commits
8ac52f69b3
..
main
| Author | SHA1 | Date | |
|---|---|---|---|
| 2a699bb91b | |||
| a5b2b73dfd | |||
| aee473c9f9 | |||
| 41e97c6a05 | |||
| ecf3d6f3a1 | |||
| 21aaeb3aba | |||
| e0230ebda7 | |||
| 47b8db18d9 | |||
| a0545cca64 | |||
| 129a760996 |
@@ -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 \
|
||||
|
||||
Binary file not shown.
|
Before Width: | Height: | Size: 473 B After Width: | Height: | Size: 22 KiB |
@@ -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
@@ -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
@@ -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
@@ -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()
|
||||
|
||||
@@ -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
@@ -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:
|
||||
|
||||
Reference in New Issue
Block a user