From a91a05db41741321505d6b7eb4e9c566a00a5cd0 Mon Sep 17 00:00:00 2001 From: Anonymous Date: Sun, 27 Sep 2026 12:08:03 +0300 Subject: [PATCH 01/10] Trust the bundled certificates for every HTTPS connection; release 3.0.1 Frozen macOS builds pointed OpenSSL at a folder of the build machine, so every HTTPS source failed with a connection error. One shared context now loads the system store and the certifi bundle. --- CHANGELOG.md | 6 ++++++ packaging/windows-installer.iss | 2 +- proxy_workbench/branding.py | 2 +- proxy_workbench/gateway.py | 4 ++-- proxy_workbench/probes.py | 8 +++----- proxy_workbench/proxytool.py | 3 ++- proxy_workbench/tls.py | 29 +++++++++++++++++++++++++++ pyproject.toml | 2 +- tests/test_tls_context.py | 35 +++++++++++++++++++++++++++++++++ 9 files changed, 80 insertions(+), 11 deletions(-) create mode 100644 proxy_workbench/tls.py create mode 100644 tests/test_tls_context.py diff --git a/CHANGELOG.md b/CHANGELOG.md index a8b8d2e..8ce1935 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -2,6 +2,12 @@ The format follows Keep a Changelog and semantic versioning. +## [3.0.1] — 2026-09-27 + +### Fixed + +- **macOS app: collecting and HTTPS checks work again.** The 3.0.0 macOS builds looked for trusted certificates in a folder of the build machine, so every HTTPS source failed with a connection error and no proxies could be collected. Every HTTPS connection (sources, HTTPS checks, the gateway's HTTPS upstreams) now also trusts the certificate bundle shipped inside the app, while the system store is still used. + ## [3.0.0] — 2026-09-27 The biggest release so far: a new interface, a source catalog of 150 lists, a desktop app that lives in the menu bar, persistent proxy pools with schedules, API keys, diagnostics that explain an empty result, and backups you can preview before restoring. Existing data folders are upgraded in place, with an automatic copy taken first. diff --git a/packaging/windows-installer.iss b/packaging/windows-installer.iss index f366fdc..f91ea3c 100644 --- a/packaging/windows-installer.iss +++ b/packaging/windows-installer.iss @@ -1,6 +1,6 @@ ; Per-user installer for the Windows desktop build. ; -; iscc /DProductVersion=3.0.0 /DOutDir=C:\path\to\dist /DSourceDir=C:\path\to\dist packaging\windows-installer.iss +; iscc /DProductVersion=3.0.1 /DOutDir=C:\path\to\dist /DSourceDir=C:\path\to\dist packaging\windows-installer.iss ; ; PrivilegesRequired=lowest is the whole point: the app writes to per-user ; folders, so it never needs an administrator, and it never installs anything diff --git a/proxy_workbench/branding.py b/proxy_workbench/branding.py index 4f58e6b..34e8e4e 100644 --- a/proxy_workbench/branding.py +++ b/proxy_workbench/branding.py @@ -11,7 +11,7 @@ PRODUCT_NAME = "Proxy Workbench" PRODUCT_ID = "ProxyWorkbench" -PRODUCT_VERSION = "3.0.0" +PRODUCT_VERSION = "3.0.1" DEFAULT_REQUEST_PROFILE = "workbench" PROJECT_URL = "https://github.com/DavidVoitenko/proxy-workbench" # Newest built-in source list, fetched only when the user asks for it. This is diff --git a/proxy_workbench/gateway.py b/proxy_workbench/gateway.py index c93b7ac..b878d2f 100644 --- a/proxy_workbench/gateway.py +++ b/proxy_workbench/gateway.py @@ -49,7 +49,7 @@ from collections import OrderedDict from pathlib import Path -from . import geoip, reputation, secrets as secretstore, socks4 +from . import geoip, reputation, secrets as secretstore, socks4, tls from .api import Exports, is_loopback, select from .i18n import tr @@ -1019,7 +1019,7 @@ async def open_tunnel(proxy, host, port, forward=False, ssl_context=None, creden proxy_host, _, proxy_port = address.rpartition(':') proxy_host = proxy_host.strip('[]') if scheme == 'https': - context = ssl_context or ssl.create_default_context() + context = ssl_context or tls.default_context() reader, writer = await asyncio.open_connection(proxy_host, int(proxy_port), limit=MAX_HEAD, ssl=context, server_hostname=proxy_host) else: diff --git a/proxy_workbench/probes.py b/proxy_workbench/probes.py index b46359e..781b63e 100644 --- a/proxy_workbench/probes.py +++ b/proxy_workbench/probes.py @@ -27,11 +27,12 @@ import json import math import re -import ssl import time from dataclasses import dataclass, replace from urllib.parse import parse_qs, urljoin, urlsplit +from . import tls + __all__ = [ # errors 'ProbeError', 'E_VALIDATION_SCHEMA', 'E_VALIDATION_FIELD', 'E_VALIDATION_UNKNOWN_FIELD', @@ -632,10 +633,7 @@ def build_ssl_context(ca_bundle=None): a custom CA, a header, a redirect or an API key never disables TLS verification. """ - context = ssl.create_default_context(cafile=ca_bundle) if ca_bundle else ssl.create_default_context() - context.check_hostname = True - context.verify_mode = ssl.CERT_REQUIRED - return context + return tls.default_context(ca_bundle) @dataclass(frozen=True) diff --git a/proxy_workbench/proxytool.py b/proxy_workbench/proxytool.py index 9a7dbd1..e1d93d4 100644 --- a/proxy_workbench/proxytool.py +++ b/proxy_workbench/proxytool.py @@ -37,6 +37,7 @@ from . import diagnostics from . import geoip from . import socks4 +from . import tls from . import formats from .i18n import tr, utf8_output from . import paths @@ -48,7 +49,7 @@ # generation directory a second time. ROOT = paths.PACKAGE -TLS = ssl.create_default_context() +TLS = tls.default_context() SCHEMES = {'http', 'https', 'socks4', 'socks5', 'socks5h'} SAFE_TARGET_HEADERS = {'accept', 'accept-encoding', 'accept-language', 'cache-control', 'pragma', 'user-agent', 'x-client-version', 'x-request-id'} diff --git a/proxy_workbench/tls.py b/proxy_workbench/tls.py new file mode 100644 index 0000000..27268bd --- /dev/null +++ b/proxy_workbench/tls.py @@ -0,0 +1,29 @@ +"""One verifying TLS context for every outgoing HTTPS connection. + +``ssl.create_default_context()`` trusts whatever the interpreter's OpenSSL was +built to look at. In a frozen build that is a path on the build machine +(``/opt/homebrew/etc/openssl@3``, a CI toolcache), which does not exist on the +user's computer, so every certificate failed and every HTTPS source looked +unreachable. The certifi bundle ships with httpx and inside every build, so it +is always added; the system store is still loaded first, which keeps the +Windows store and corporate roots working. +""" +from __future__ import annotations + +import ssl + +import certifi + +__all__ = ['default_context'] + + +def default_context(cafile=None): + """A verifying context: a given CA file replaces the defaults, never disables checks.""" + if cafile: + context = ssl.create_default_context(cafile=cafile) + else: + context = ssl.create_default_context() + context.load_verify_locations(cafile=certifi.where()) + context.check_hostname = True + context.verify_mode = ssl.CERT_REQUIRED + return context diff --git a/pyproject.toml b/pyproject.toml index 84fa8c1..198618b 100644 --- a/pyproject.toml +++ b/pyproject.toml @@ -6,7 +6,7 @@ build-backend = "setuptools.build_meta" name = "proxy-workbench" # Must equal proxy_workbench.branding.PRODUCT_VERSION: the binary advertised one # number and the wheel another. -version = "3.0.0" +version = "3.0.1" description = "Collect free public proxies and keep only the ones that work for your services: local GUI + CLI checker for HTTP, HTTPS and SOCKS5" readme = "README.md" requires-python = ">=3.11" diff --git a/tests/test_tls_context.py b/tests/test_tls_context.py new file mode 100644 index 0000000..3d5158f --- /dev/null +++ b/tests/test_tls_context.py @@ -0,0 +1,35 @@ +"""HTTPS must verify in a frozen build, where OpenSSL's default paths are gone.""" +from __future__ import annotations + +import ssl +import unittest +from unittest import mock + +import certifi + +from proxy_workbench import tls + + +class DefaultContextTests(unittest.TestCase): + def test_the_certifi_roots_are_loaded_even_when_the_system_paths_are_missing(self): + # A frozen app points OpenSSL at a folder of the build machine. Loading + # nothing from there must still leave a context that trusts public roots. + with mock.patch.object(ssl.SSLContext, 'load_default_certs', lambda self, purpose=None: None): + context = tls.default_context() + self.assertGreater(context.cert_store_stats()['x509_ca'], 50) + self.assertEqual(context.verify_mode, ssl.CERT_REQUIRED) + self.assertTrue(context.check_hostname) + + def test_a_custom_bundle_replaces_the_defaults_and_keeps_verification(self): + context = tls.default_context(certifi.where()) + self.assertEqual(context.verify_mode, ssl.CERT_REQUIRED) + self.assertTrue(context.check_hostname) + + def test_every_outgoing_path_uses_the_shared_context(self): + from proxy_workbench import probes, proxytool + self.assertGreater(proxytool.TLS.cert_store_stats()['x509_ca'], 50) + self.assertGreater(probes.build_ssl_context().cert_store_stats()['x509_ca'], 50) + + +if __name__ == '__main__': + unittest.main() From 2d5c89276ed01551d9ab82bc26f27bc3a623c543 Mon Sep 17 00:00:00 2001 From: Anonymous Date: Sun, 27 Sep 2026 12:08:35 +0300 Subject: [PATCH 02/10] CI: the built macOS and Windows apps must collect real sources over HTTPS --- .github/workflows/ci.yml | 13 +++++++++++++ .github/workflows/windows.yml | 12 ++++++++++++ 2 files changed, 25 insertions(+) diff --git a/.github/workflows/ci.yml b/.github/workflows/ci.yml index 2c53740..e476e41 100644 --- a/.github/workflows/ci.yml +++ b/.github/workflows/ci.yml @@ -203,6 +203,19 @@ jobs: - name: Smoke test the built application run: python packaging/smoke.py "dist/Proxy Workbench.app/Contents/MacOS/Proxy Workbench" + # The smoke test only talks to local mocks. A frozen build can still fail + # every real HTTPS download (3.0.0 did: no trusted certificates), so the + # built app has to collect the quick source set from the internet. + - name: Collect real sources over HTTPS with the built application + run: | + app="dist/Proxy Workbench.app/Contents/MacOS/Proxy Workbench" + data="$RUNNER_TEMP/pw-https-check" + "$app" source set quick --data "$data" + "$app" collect --data "$data" | tee "$RUNNER_TEMP/collect.log" + count=$(sed -n 's/^Unique candidates in the database: \([0-9]*\).*/\1/p' "$RUNNER_TEMP/collect.log") + echo "collected: ${count:-0}" + test "${count:-0}" -gt 0 + - name: Report what the build recorded run: | python - <<'PY' diff --git a/.github/workflows/windows.yml b/.github/workflows/windows.yml index 137cf58..eba1cfc 100644 --- a/.github/workflows/windows.yml +++ b/.github/workflows/windows.yml @@ -68,6 +68,18 @@ jobs: - name: Smoke test the console CLI executable run: python packaging/smoke.py dist/proxy-workbench-cli.exe + # Local mocks cannot show whether real HTTPS downloads verify in the + # frozen build, so collect the quick source set from the internet. + - name: Collect real sources over HTTPS with the built CLI + shell: bash + run: | + data="$RUNNER_TEMP/pw-https-check" + dist/proxy-workbench-cli.exe source set quick --data "$data" + dist/proxy-workbench-cli.exe collect --data "$data" | tee "$RUNNER_TEMP/collect.log" + count=$(sed -n 's/^Unique candidates in the database: \([0-9]*\).*/\1/p' "$RUNNER_TEMP/collect.log") + echo "collected: ${count:-0}" + test "${count:-0}" -gt 0 + # The GUI executable is windowed, so it has nowhere to print: what it # prints is the log the build already produced. - name: Report what the build recorded From f1d8fd80f38573c2519406f69bea6d9bf9b8425c Mon Sep 17 00:00:00 2001 From: Anonymous Date: Sun, 27 Sep 2026 12:27:33 +0300 Subject: [PATCH 03/10] Wait for a short hold on the data folder; read refused request bodies A check started from the desktop app failed with "data folder is already used" whenever the background job runner or scheduler held the lock for its one-second poll. A run now waits up to ten seconds for such a hold. The interface server answered refused POST requests without reading their body, which Windows turned into a reset connection. --- CHANGELOG.md | 2 ++ proxy_workbench/gui.py | 19 ++++++++++++ proxy_workbench/proxytool.py | 34 +++++++++++++-------- tests/test_data_lock_wait.py | 58 ++++++++++++++++++++++++++++++++++++ 4 files changed, 101 insertions(+), 12 deletions(-) create mode 100644 tests/test_data_lock_wait.py diff --git a/CHANGELOG.md b/CHANGELOG.md index 8ce1935..6cab0f5 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -7,6 +7,8 @@ The format follows Keep a Changelog and semantic versioning. ### Fixed - **macOS app: collecting and HTTPS checks work again.** The 3.0.0 macOS builds looked for trusted certificates in a folder of the build machine, so every HTTPS source failed with a connection error and no proxies could be collected. Every HTTPS connection (sources, HTTPS checks, the gateway's HTTPS upstreams) now also trusts the certificate bundle shipped inside the app, while the system store is still used. +- **“This data folder is already used by another run” when starting a check.** The desktop app's background job runner and scheduler hold the data lock for a moment every second, and a check started from the interface gave up on the first try. A run now waits up to 10 seconds for such a short hold; a folder that stays busy is still refused. +- The interface server reads or closes the body of a request it refuses, so a refused request no longer shows up as a reset connection on Windows or as a garbled next request. ## [3.0.0] — 2026-09-27 diff --git a/proxy_workbench/gui.py b/proxy_workbench/gui.py index 9828757..7e6ea06 100644 --- a/proxy_workbench/gui.py +++ b/proxy_workbench/gui.py @@ -5129,12 +5129,31 @@ def query_body(self, query): body[name] = text return body + def discard_body(self): + """Read an unused request body, or close the connection if it cannot be read. + + Answering before the body is read leaves it in the socket: on a kept-alive + connection it would be parsed as the next request, and on Windows closing + a socket with unread data resets it, so the client sees a connection + abort instead of the refusal. + """ + try: + length = int(self.headers.get('Content-Length', '0')) + except ValueError: + length = -1 + if 0 < length <= MAX_BODY: + self.rfile.read(length) + elif length != 0: + self.close_connection = True + def do_POST(self): if not self.allowed(): + self.discard_body() return try: length = int(self.headers.get('Content-Length', '0')) if not 0 < length <= MAX_BODY: + self.close_connection = True return self.respond(413, dict(error='Слишком большой запрос.')) payload = json.loads(self.rfile.read(length)) path = urlsplit(self.path).path diff --git a/proxy_workbench/proxytool.py b/proxy_workbench/proxytool.py index e1d93d4..6a90691 100644 --- a/proxy_workbench/proxytool.py +++ b/proxy_workbench/proxytool.py @@ -50,6 +50,8 @@ ROOT = paths.PACKAGE TLS = tls.default_context() +# How long a run waits for the data folder lock held by a short background task. +DATA_LOCK_WAIT_S = 10.0 SCHEMES = {'http', 'https', 'socks4', 'socks5', 'socks5h'} SAFE_TARGET_HEADERS = {'accept', 'accept-encoding', 'accept-language', 'cache-control', 'pragma', 'user-agent', 'x-client-version', 'x-request-id'} @@ -8139,19 +8141,27 @@ def main(argv=None): denylist = Denylist.from_file(denylist_path, normalizer=normalize) collect_denylist = Denylist.empty() if args.local_denylist is False else denylist # Exclusive OS lock is released even after a crash; read-only exports also lock. + # The desktop app's job runner and scheduler hold it for a moment every + # second, so a worker started from the interface waits for that instead of + # failing on the first try; only a run that keeps the folder busy is refused. lock = (args.data / 'workbench.lock').open('a+b') - try: - if os.name == 'nt': - import msvcrt - lock.write(b'0'); lock.flush(); lock.seek(0) - msvcrt.locking(lock.fileno(), msvcrt.LK_NBLCK, 1) - else: - import fcntl - fcntl.flock(lock, fcntl.LOCK_EX | fcntl.LOCK_NB) - except OSError: - print(tr('Эта папка data уже используется другим запуском.', 'This data folder is already used by another run.'), file=sys.stderr) - lock.close() - return 2 + deadline = time.monotonic() + DATA_LOCK_WAIT_S + while True: + try: + if os.name == 'nt': + import msvcrt + lock.seek(0); lock.write(b'0'); lock.flush(); lock.seek(0) + msvcrt.locking(lock.fileno(), msvcrt.LK_NBLCK, 1) + else: + import fcntl + fcntl.flock(lock, fcntl.LOCK_EX | fcntl.LOCK_NB) + break + except OSError: + if time.monotonic() >= deadline: + print(tr('Эта папка data уже используется другим запуском.', 'This data folder is already used by another run.'), file=sys.stderr) + lock.close() + return 2 + time.sleep(0.1) if args.command == 'clear-data': if not args.yes: p.error(tr('clear-data требует явного --yes', 'clear-data needs an explicit --yes')) diff --git a/tests/test_data_lock_wait.py b/tests/test_data_lock_wait.py new file mode 100644 index 0000000..b83db37 --- /dev/null +++ b/tests/test_data_lock_wait.py @@ -0,0 +1,58 @@ +"""A run started while a background task briefly holds the data folder must wait, not fail.""" +from __future__ import annotations + +import contextlib +import io +import tempfile +import threading +import time +import unittest +from pathlib import Path +from unittest import mock + +from proxy_workbench import proxytool +from proxy_workbench.maintenance import exclusive_lock + + +def hold(path, seconds, ready): + with exclusive_lock(path): + ready.set() + time.sleep(seconds) + + +class DataLockWaitTests(unittest.TestCase): + def run_export(self, data): + err = io.StringIO() + with contextlib.redirect_stdout(io.StringIO()), contextlib.redirect_stderr(err): + code = proxytool.main(['export', '--data', str(data)]) + return code, err.getvalue() + + def test_a_short_background_hold_is_waited_out(self): + # The desktop app's job runner and scheduler take the lock every second; + # a scan launched from the interface used to die with "folder busy". + with tempfile.TemporaryDirectory() as tmp: + ready = threading.Event() + holder = threading.Thread(target=hold, args=(Path(tmp) / 'workbench.lock', 0.6, ready)) + holder.start() + ready.wait(5) + started = time.monotonic() + _, err = self.run_export(tmp) + holder.join() + self.assertNotIn('already used by another run', err) + self.assertNotIn('уже используется', err) + self.assertGreaterEqual(time.monotonic() - started, 0.4) + + def test_a_folder_that_stays_busy_is_still_refused(self): + with tempfile.TemporaryDirectory() as tmp, mock.patch.object(proxytool, 'DATA_LOCK_WAIT_S', 0.3): + ready = threading.Event() + holder = threading.Thread(target=hold, args=(Path(tmp) / 'workbench.lock', 1.5, ready)) + holder.start() + ready.wait(5) + code, err = self.run_export(tmp) + holder.join() + self.assertEqual(code, 2) + self.assertTrue('already used by another run' in err or 'уже используется' in err, err) + + +if __name__ == '__main__': + unittest.main() From 57bf17be3068e6b7acff5b5174d9616c6372a454 Mon Sep 17 00:00:00 2001 From: Anonymous Date: Sun, 27 Sep 2026 12:56:50 +0300 Subject: [PATCH 04/10] Let scans use all their workers and spread ports of one host The body memory budget was a fixed 64 MiB, which with the default 1 MiB body cap held every scan at 64 checks in flight. It now grows with the worker count up to 1 GiB. Candidates are read in address order, so the thousands of ports some lists publish for one IP queued every worker behind that host; the first port of each host now goes first and further ports follow round by round. 663k collected addresses: 25 -> ~390 checks/s. --- CHANGELOG.md | 1 + proxy_workbench/proxytool.py | 62 ++++++++++++++++++++++++++- tests/test_scan_engine_regressions.py | 31 ++++++++++++++ 3 files changed, 92 insertions(+), 2 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index 6cab0f5..dfdaabd 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -6,6 +6,7 @@ The format follows Keep a Changelog and semantic versioning. ### Fixed +- **Checks run up to 16× faster on large lists.** Every scan was held at 64 checks in flight whatever the worker setting said (the memory reserved for response bodies had a fixed 64 MiB budget); the budget now grows with the workers up to 1 GiB. Ports of one IP no longer queue up together either: the first port of every host is checked first and further ports follow host by host. A full check of 663,000 collected addresses went from 25 to about 390 addresses per second with 512 workers. - **macOS app: collecting and HTTPS checks work again.** The 3.0.0 macOS builds looked for trusted certificates in a folder of the build machine, so every HTTPS source failed with a connection error and no proxies could be collected. Every HTTPS connection (sources, HTTPS checks, the gateway's HTTPS upstreams) now also trusts the certificate bundle shipped inside the app, while the system store is still used. - **“This data folder is already used by another run” when starting a check.** The desktop app's background job runner and scheduler hold the data lock for a moment every second, and a check started from the interface gave up on the first try. A run now waits up to 10 seconds for such a short hold; a folder that stays busy is still refused. - The interface server reads or closes the body of a request it refuses, so a refused request no longer shows up as a reset connection on Windows or as a garbled next request. diff --git a/proxy_workbench/proxytool.py b/proxy_workbench/proxytool.py index 6a90691..8c1cccd 100644 --- a/proxy_workbench/proxytool.py +++ b/proxy_workbench/proxytool.py @@ -3843,6 +3843,41 @@ def measurement_bytes(row): #: cursor of the corpus. CANDIDATE_PAGE = 4096 + +def _candidate_host(proxy): + """The host part of ``scheme://host:port``: the unit the per-host limit counts.""" + rest = proxy.partition('://')[2] or proxy + return rest.rpartition(':')[0] or rest + + +class HostSpread: + """First port of every host at once, further ports of a host in later rounds.""" + + def __init__(self): + self._seen: set[str] = set() + self._later: dict[str, list[str]] = {} + + def admit(self, proxy): + """True when ``proxy`` may go out now; otherwise it is kept for a round.""" + host = _candidate_host(proxy) + if host in self._seen: + self._later.setdefault(host, []).append(proxy) + return False + self._seen.add(host) + return True + + def rounds(self): + """The kept addresses, one per host per round, in first-seen host order.""" + later, index = self._later, 0 + while later: + batch = [] + for host in list(later): + batch.append(later[host][index]) + if index + 1 >= len(later[host]): + del later[host] + index += 1 + yield batch + #: One page of the scope, as an ordered range over the address index with #: membership as a test — see :func:`candidate_pages` for why it is not a join. _CANDIDATE_PAGE_SQL = ( @@ -3864,6 +3899,18 @@ def measurement_bytes(row): #: silence. DEFAULT_RAM_PER_INFLIGHT = 256 * 1024 DEFAULT_MAX_RAM_BYTES = 64 * 1024 * 1024 +#: Upper bound of the reserved body memory of one scan, whatever --workers says. +MAX_SCAN_RAM_BYTES = 1024 * 1024 * 1024 + + +def scan_ram_budget(workers, ram_per_inflight): + """Body memory a scan may reserve: one full body per worker, up to a ceiling. + + A fixed 64 MiB with the default 1 MiB body cap held every scan at 64 + checks in flight whatever --workers said. + """ + wanted = max(DEFAULT_MAX_RAM_BYTES, ram_per_inflight * max(1, int(workers))) + return min(wanted, max(MAX_SCAN_RAM_BYTES, ram_per_inflight)) def _number(value): @@ -4319,7 +4366,7 @@ def counted(row): max_inflight=max(1, int(workers)), max_open_fds=open_fd_budget(workers), fds_per_request=FDS_PER_REQUEST, - max_ram_bytes=max(DEFAULT_MAX_RAM_BYTES, ram_per_inflight * 8), + max_ram_bytes=scan_ram_budget(workers, ram_per_inflight), ram_per_inflight_bytes=ram_per_inflight, max_queue_items=max(2 * max(1, int(workers)), 64), max_results_pending=64, @@ -4428,6 +4475,13 @@ async def candidates(): phases = (True, False) if proven else (None,) for phase in phases: last = '' + # Pages come in address order, so the thousands of ports some lists + # publish for one IP arrive together and every worker queued behind + # that one host (a 663k corpus ran at 28 checks/s). The first port + # of a host goes out at once; further ports of the same host wait and + # are sent afterwards one per host per round, so hosts stay in + # parallel and the per-host limit still holds. + spread = HostSpread() while True: page = [proxy for (proxy,) in db.execute( _CANDIDATE_PAGE_SQL, (last, collection_id, CANDIDATE_PAGE))] @@ -4450,10 +4504,14 @@ async def candidates(): continue if phase is not None and (proxy in proven) is not phase: continue - body.append(proxy) + if spread.admit(proxy): + body.append(proxy) if not body: continue yield ''.join(f'{proxy}\n' for proxy in body).encode('utf-8') + for batch in spread.rounds(): + for start in range(0, len(batch), CANDIDATE_PAGE): + yield ''.join(f'{proxy}\n' for proxy in batch[start:start + CANDIDATE_PAGE]).encode('utf-8') rows: dict[str, dict] = {} store_failures: list[str] = [] diff --git a/tests/test_scan_engine_regressions.py b/tests/test_scan_engine_regressions.py index 04f23e6..a4e59d9 100644 --- a/tests/test_scan_engine_regressions.py +++ b/tests/test_scan_engine_regressions.py @@ -377,3 +377,34 @@ def test_the_wrapped_local_error_is_found(self): if __name__ == '__main__': unittest.main() + + +class RamBudgetTests(unittest.TestCase): + def test_the_worker_count_is_not_capped_at_64_by_the_body_reservation(self): + from proxy_workbench import pipeline + mib = 1024 * 1024 + for workers in (128, 512): + budgets = pipeline.Budgets(max_inflight=workers, max_open_fds=10 ** 6, fds_per_request=3, + max_ram_bytes=proxytool.scan_ram_budget(workers, mib), + ram_per_inflight_bytes=mib) + self.assertEqual(budgets.ram_ceiling >= workers, True, workers) + + def test_the_reservation_has_a_ceiling(self): + mib = 1024 * 1024 + self.assertEqual(proxytool.scan_ram_budget(100_000, mib), proxytool.MAX_SCAN_RAM_BYTES) + self.assertEqual(proxytool.scan_ram_budget(1, mib), proxytool.DEFAULT_MAX_RAM_BYTES) + + +class HostSpreadTests(unittest.TestCase): + def test_many_ports_of_one_host_do_not_come_in_a_row(self): + # Address order put thousands of ports of one IP next to each other, so + # every worker waited for that host under the per-host limit. + stream = [f'http://203.0.113.1:{port}' for port in range(1, 6)] + ['http://203.0.113.2:80', + 'http://198.51.100.7:8080', 'http://198.51.100.7:3128'] + spread = proxytool.HostSpread() + now = [proxy for proxy in stream if spread.admit(proxy)] + self.assertEqual(now, ['http://203.0.113.1:1', 'http://203.0.113.2:80', 'http://198.51.100.7:8080']) + rounds = list(spread.rounds()) + self.assertEqual(rounds[0], ['http://203.0.113.1:2', 'http://198.51.100.7:3128']) + self.assertEqual(rounds[1:], [['http://203.0.113.1:3'], ['http://203.0.113.1:4'], ['http://203.0.113.1:5']]) + self.assertEqual(sorted(now + [p for r in rounds for p in r]), sorted(stream)) From 69885d5feeb076a5d138eeaa2a10a244b4aa444c Mon Sep 17 00:00:00 2001 From: Anonymous Date: Sun, 27 Sep 2026 13:16:26 +0300 Subject: [PATCH 05/10] Gateway: race slow upstreams and never go dark while every proxy rests A connect attempt that has not answered within 2 s now has the next proxy dialled alongside it and the first tunnel wins (4 attempts), instead of waiting out each connect timeout in turn. When every matching proxy rests after failures, the one due back first is used instead of answering 502 for the whole cooldown. --- CHANGELOG.md | 1 + proxy_workbench/gateway.py | 113 ++++++++++++++++++++++++++--------- tests/test_gateway_health.py | 15 +++++ 3 files changed, 100 insertions(+), 29 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index dfdaabd..5c14a6c 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -8,6 +8,7 @@ The format follows Keep a Changelog and semantic versioning. - **Checks run up to 16× faster on large lists.** Every scan was held at 64 checks in flight whatever the worker setting said (the memory reserved for response bodies had a fixed 64 MiB budget); the budget now grows with the workers up to 1 GiB. Ports of one IP no longer queue up together either: the first port of every host is checked first and further ports follow host by host. A full check of 663,000 collected addresses went from 25 to about 390 addresses per second with 512 workers. - **macOS app: collecting and HTTPS checks work again.** The 3.0.0 macOS builds looked for trusted certificates in a folder of the build machine, so every HTTPS source failed with a connection error and no proxies could be collected. Every HTTPS connection (sources, HTTPS checks, the gateway's HTTPS upstreams) now also trusts the certificate bundle shipped inside the app, while the system store is still used. +- **The rotating proxy is far more reliable with flaky free proxies.** A proxy that does not answer within 2 seconds no longer blocks the request: the next proxy is tried alongside it and the first working tunnel wins (up to 4 proxies; a failed request now gives up after about 14 s instead of 24 s). When every proxy in a small pool was resting after failures, the gateway answered 502 instantly for five minutes; it now falls back to the proxy that is due back first. On 11 real public proxies right after a restart: 18 of 20 requests succeeded instead of 6 of 10. - **“This data folder is already used by another run” when starting a check.** The desktop app's background job runner and scheduler hold the data lock for a moment every second, and a check started from the interface gave up on the first try. A run now waits up to 10 seconds for such a short hold; a folder that stays busy is still refused. - The interface server reads or closes the body of a request it refuses, so a refused request no longer shows up as a reset connection on Windows or as a garbled next request. diff --git a/proxy_workbench/gateway.py b/proxy_workbench/gateway.py index b878d2f..f2c4f3b 100644 --- a/proxy_workbench/gateway.py +++ b/proxy_workbench/gateway.py @@ -779,8 +779,16 @@ def _pick(self, exclude, request, session, sticky, binding, reserve): # free slot for the proxy it is bound to means a refusal. return None if not candidates: - return None - choice = self._choose(candidates, self._strategy(binding), now) + # Every matching proxy is resting. Refusing outright turned a + # small pool of flaky public proxies into a gateway that answered + # 502 instantly for five minutes; the one whose rest ends first + # is still a better answer than none. + candidates = self._resting_fallback(request, now, binding, exclude) + if not candidates: + return None + choice = candidates[0] + else: + choice = self._choose(candidates, self._strategy(binding), now) if session: self.sessions[session] = (choice, now + self.session_ttl) if not reserve: @@ -788,6 +796,14 @@ def _pick(self, exclude, request, session, sticky, binding, reserve): self._take(choice) return Lease(self, choice) + def _resting_fallback(self, request, now, binding, exclude): + """Resting proxies that could serve this request, soonest back first.""" + limit = self._limit(binding) + resting = [proxy for proxy in self.matching(request, binding) + if proxy not in exclude and self.resting.get(proxy, 0) > now and self.allowed(proxy) + and (not limit or self.active.get(proxy, 0) < limit)] + return sorted(resting, key=lambda proxy: self.resting[proxy]) + def _take(self, proxy): self.active[proxy] = self.active.get(proxy, 0) + 1 @@ -1190,7 +1206,7 @@ def credential_state(source): class Gateway: - def __init__(self, pool, token=None, attempts=3, connect_timeout=8, idle_timeout=300, + def __init__(self, pool, token=None, attempts=4, connect_timeout=8, idle_timeout=300, allow_local_without_auth=False, *, handshake_timeout=30, response_timeout=15, max_clients=MAX_CLIENTS, max_session=0, bind=None, token_origin=None, drain_timeout=2.0, ssl_context=None, refresh_interval=2.0, @@ -1200,6 +1216,8 @@ def __init__(self, pool, token=None, attempts=3, connect_timeout=8, idle_timeout self.token_origin = token_origin or ('explicit' if token else 'none') self.attempts = max(1, int(attempts)) self.connect_timeout = connect_timeout + # Seconds before the next proxy joins a connect attempt that has not answered. + self.stagger = min(2.0, float(connect_timeout)) self.idle_timeout = idle_timeout self.allow_local_without_auth = allow_local_without_auth self.handshake_timeout = handshake_timeout @@ -1307,38 +1325,75 @@ async def connect(self, host, port, forward=False, request=None, session=None, b The returned lease owns the concurrency slot; the caller must release it (or use ``with``) whatever happens to the stream. + + Attempts are staggered rather than strictly one after another: when a + proxy has not answered within ``stagger`` seconds, the next one starts + alongside it and the first tunnel wins. One after another, three dead + proxies cost a client three full connect timeouts before a 502. A + strict sticky session keeps the sequential order, because it may only + fall through when its own proxy failed. """ tried = set() - for attempt in range(self.attempts): + pending = {} + failures = 0 + strict = bool(session and sticky == 'strict') + + def start_next(): lease = self.pool.reserve(request=request, exclude=tried, session=session, sticky=sticky, binding=binding) if lease is None: - break - proxy = lease.proxy - tried.add(proxy) - self.pool.stats['retries'] += attempt > 0 - started = time.monotonic() - try: - stream, credentials = await asyncio.wait_for( - self._tunnel(proxy, host, port, forward), self.connect_timeout) - except (OSError, UpstreamError, ValueError, UnicodeError) as exc: - lease.release() - self.pool.outcome(proxy, 'handshake_failed', detail=UNUSABLE) - if session and sticky == 'strict': + return False + tried.add(lease.proxy) + self.pool.stats['retries'] += len(tried) > 1 + task = asyncio.ensure_future(asyncio.wait_for( + self._tunnel(lease.proxy, host, port, forward), self.connect_timeout)) + pending[task] = (lease, time.monotonic()) + return True + + try: + if not start_next(): + self.pool.stats['failed'] += 1 + raise UpstreamError('NO_PROXIES') + while pending: + can_add = not strict and len(tried) < self.attempts + done, _ = await asyncio.wait(list(pending), timeout=self.stagger if can_add else None, + return_when=asyncio.FIRST_COMPLETED) + for task in done: + lease, started = pending.pop(task) + try: + stream, credentials = task.result() + except (OSError, UpstreamError, ValueError, UnicodeError): + lease.release() + self.pool.outcome(lease.proxy, 'handshake_failed', detail=UNUSABLE) + failures += 1 + continue + # The plain-HTTP path still has to write the credential onto the + # wire, so it rides on the lease and is dropped the moment the + # request is written. + lease.credentials = credentials + self.pool.connected(lease.proxy, (time.monotonic() - started) * 1000) + return lease, stream + if strict and failures: break - continue - except BaseException: - # Cancellation and timeout both belong here: the slot must not leak. + # Either the stagger ran out or everything that finished failed. + if not strict and len(tried) < self.attempts: + start_next() + self.pool.stats['failed'] += 1 + raise UpstreamError('NO_WORKING_PROXY') + finally: + # Losers of the race, and everything on cancellation: no slot leaks + # and no half-open tunnel stays behind. + for task, (lease, _) in pending.items(): + task.cancel() lease.release() - raise - # The plain-HTTP path still has to write the credential onto the - # wire, so it rides on the lease and is dropped the moment the - # request is written. - lease.credentials = credentials - self.pool.connected(proxy, (time.monotonic() - started) * 1000) - return lease, stream - self.pool.stats['failed'] += 1 - raise UpstreamError('NO_WORKING_PROXY' if tried else 'NO_PROXIES') + for task in pending: + with contextlib.suppress(BaseException): + await task + for task in pending: + if task.done() and not task.cancelled() and task.exception() is None: + stream = task.result()[0] + with contextlib.suppress(Exception): + stream[1].close() if isinstance(stream, tuple) else stream.close() # --- authentication ------------------------------------------------------ @@ -1866,7 +1921,7 @@ async def start(data, host=DEFAULT_HOST, port=DEFAULT_PORT, token=None, filters= on_deny='keep', bindings=None, default_binding=None, binding=None, credentials=None, handshake_timeout=30, max_clients=MAX_CLIENTS, - max_session=0, cache_limit=CACHE_LIMIT, attempts=3, connect_timeout=8, + max_session=0, cache_limit=CACHE_LIMIT, attempts=4, connect_timeout=8, idle_timeout=300, refresh_interval=2.0, response_timeout=15, max_failures=2, cooldown=300): """Listen for proxy clients and spread their connections over the export. diff --git a/tests/test_gateway_health.py b/tests/test_gateway_health.py index 380711e..10a69c2 100644 --- a/tests/test_gateway_health.py +++ b/tests/test_gateway_health.py @@ -225,3 +225,18 @@ async def test_an_open_tunnel_is_never_moved_to_another_upstream(self): if __name__ == '__main__': unittest.main() + + +class RestingFallbackTests(GatewayCase): + async def test_a_pool_where_every_proxy_rests_still_serves(self): + # With a small pool of flaky public proxies every proxy reached two + # failures within a minute, and the gateway then answered 502 instantly + # for five minutes although some of them worked most of the time. + up = await self.http_upstream('relay') + self.publish([up.url]) + server, address = await self.start(max_failures=1, cooldown=300) + pool = server.gateway.pool + pool.outcome(up.url, 'handshake_failed') + self.assertIn(up.url, pool.resting) + self.assertEqual(pool.available(), []) + self.assertEqual((await self.get(address, '/x')).status_code, 200) From 1070a54c207cab1d9bacb7a0a708a222b5b79900 Mon Sep 17 00:00:00 2001 From: Anonymous Date: Sun, 27 Sep 2026 14:34:56 +0300 Subject: [PATCH 06/10] Fix raced gateway health and sticky sessions; scope source lookup --- proxy_workbench/gateway.py | 38 ++++++++++++---- proxy_workbench/gui.py | 6 +-- proxy_workbench/proxytool.py | 20 +++++++-- tests/test_freshness.py | 4 ++ tests/test_gateway_health.py | 86 ++++++++++++++++++++++++++++++++++-- 5 files changed, 136 insertions(+), 18 deletions(-) diff --git a/proxy_workbench/gateway.py b/proxy_workbench/gateway.py index f2c4f3b..d1d01aa 100644 --- a/proxy_workbench/gateway.py +++ b/proxy_workbench/gateway.py @@ -750,7 +750,8 @@ def pick(self, exclude=(), request=None, session=None, sticky=None, binding=None """Choose a proxy without taking a slot. Prefer :meth:`reserve` in a gateway.""" return self._pick(exclude, request, session, sticky, binding, reserve=False) - def reserve(self, request=None, exclude=(), session=None, sticky=None, binding=None): + def reserve(self, request=None, exclude=(), session=None, sticky=None, binding=None, + commit_session=True): """Pick a proxy and take its concurrency slot in one indivisible step. Doing both together is what keeps parallel connects under @@ -758,9 +759,10 @@ def reserve(self, request=None, exclude=(), session=None, sticky=None, binding=N reservation, so two clients can never observe the same free slot. The caller owns the returned :class:`Lease`. """ - return self._pick(exclude, request, session, sticky, binding, reserve=True) + return self._pick(exclude, request, session, sticky, binding, reserve=True, + commit_session=commit_session) - def _pick(self, exclude, request, session, sticky, binding, reserve): + def _pick(self, exclude, request, session, sticky, binding, reserve, commit_session=True): now = time.monotonic() sticky = self._sticky(binding, sticky) with self.lock: @@ -771,7 +773,8 @@ def _pick(self, exclude, request, session, sticky, binding, reserve): if pinned and pinned in candidates: if reserve: self._take(pinned) - self.sessions[session] = (pinned, now + self.session_ttl) + if commit_session: + self.sessions[session] = (pinned, now + self.session_ttl) return Lease(self, pinned) return pinned if pinned and sticky == 'strict': @@ -789,13 +792,18 @@ def _pick(self, exclude, request, session, sticky, binding, reserve): choice = candidates[0] else: choice = self._choose(candidates, self._strategy(binding), now) - if session: + if session and commit_session: self.sessions[session] = (choice, now + self.session_ttl) if not reserve: return choice self._take(choice) return Lease(self, choice) + def bind_session(self, session, proxy): + """Pin a session only after its upstream tunnel has won the race.""" + with self.lock: + self.sessions[session] = (proxy, time.monotonic() + self.session_ttl) + def _resting_fallback(self, request, now, binding, exclude): """Resting proxies that could serve this request, soonest back first.""" limit = self._limit(binding) @@ -1336,11 +1344,11 @@ async def connect(self, host, port, forward=False, request=None, session=None, b tried = set() pending = {} failures = 0 - strict = bool(session and sticky == 'strict') + strict = bool(session and self.pool._sticky(binding, sticky) == 'strict') def start_next(): lease = self.pool.reserve(request=request, exclude=tried, session=session, - sticky=sticky, binding=binding) + sticky=sticky, binding=binding, commit_session=strict) if lease is None: return False tried.add(lease.proxy) @@ -1358,7 +1366,10 @@ def start_next(): can_add = not strict and len(tried) < self.attempts done, _ = await asyncio.wait(list(pending), timeout=self.stagger if can_add else None, return_when=asyncio.FIRST_COMPLETED) - for task in done: + winner = None + for task in list(pending): + if task not in done: + continue lease, started = pending.pop(task) try: stream, credentials = task.result() @@ -1367,11 +1378,22 @@ def start_next(): self.pool.outcome(lease.proxy, 'handshake_failed', detail=UNUSABLE) failures += 1 continue + if winner is None: + winner = (lease, stream, credentials, started) + else: + # Two tunnels can finish in one event-loop turn. Only + # one may keep its reservation and open stream. + stream[1].close() + lease.release() + if winner is not None: + lease, stream, credentials, started = winner # The plain-HTTP path still has to write the credential onto the # wire, so it rides on the lease and is dropped the moment the # request is written. lease.credentials = credentials self.pool.connected(lease.proxy, (time.monotonic() - started) * 1000) + if session: + self.pool.bind_session(session, lease.proxy) return lease, stream if strict and failures: break diff --git a/proxy_workbench/gui.py b/proxy_workbench/gui.py index 7e6ea06..ed7c0e3 100644 --- a/proxy_workbench/gui.py +++ b/proxy_workbench/gui.py @@ -2772,7 +2772,7 @@ def scoped_rows(self, plan, *, limit=None): if denylist.error: raise ValueError('Не удалось прочитать локальный denylist; обновите список.') policy = self.result_policy(plan, cfg, denylist=denylist) - source_keys = self.source_keys(conn) + source_keys = self.source_keys(conn, plan['profile']) provider_of = self.provider_resolver() now = time.time() entries = self.annotations.read().get('entries') or {} @@ -2863,9 +2863,9 @@ def matches_filters(self, plan, row, verdict): return False return True - def source_keys(self, conn): + def source_keys(self, conn, profile=None): try: - return core.source_map(conn) + return core.source_map(conn, profile) except sqlite3.Error: return {} diff --git a/proxy_workbench/proxytool.py b/proxy_workbench/proxytool.py index 8c1cccd..585dc8d 100644 --- a/proxy_workbench/proxytool.py +++ b/proxy_workbench/proxytool.py @@ -5071,23 +5071,37 @@ def listed_counts(db): return {} -def source_map(db): +def source_map(db, profile=None): """Return every source key associated with a candidate. ``candidate_meta.source`` is retained as a first-seen compatibility field, but it is not authoritative: the many-to-many ``candidate_seen`` table is what prevents source health and recommendations from depending on collection order. + + With ``profile`` only addresses that have a result in that profile are + mapped. A reader of one page of results does not need the sources of every + collected candidate: on 663k candidates the full map cost 1.3 s a request. """ result = {} + if profile is None: + meta_sql = 'SELECT proxy, source FROM candidate_meta WHERE source IS NOT NULL' + seen_sql = 'SELECT proxy, source FROM candidate_seen' + params = () + else: + meta_sql = ('SELECT m.proxy, m.source FROM results r JOIN candidate_meta m ON m.proxy = r.proxy ' + 'WHERE r.profile = ? AND m.source IS NOT NULL') + seen_sql = ('SELECT s.proxy, s.source FROM results r JOIN candidate_seen s ON s.proxy = r.proxy ' + 'WHERE r.profile = ?') + params = (profile,) try: - for proxy, source in db.execute('SELECT proxy, source FROM candidate_meta WHERE source IS NOT NULL'): + for proxy, source in db.execute(meta_sql, params): if source: result.setdefault(proxy, []).append(source) except sqlite3.Error: pass try: - for proxy, source in db.execute('SELECT proxy, source FROM candidate_seen'): + for proxy, source in db.execute(seen_sql, params): if not source: continue values = result.setdefault(proxy, []) diff --git a/tests/test_freshness.py b/tests/test_freshness.py index b919c89..d12aca1 100644 --- a/tests/test_freshness.py +++ b/tests/test_freshness.py @@ -327,8 +327,12 @@ def test_overlap_health_and_provenance_use_every_source(self): db.execute('INSERT INTO candidate_meta(proxy, source) VALUES (?,?)', (proxy, 'aaa')) mark_seen(db, (proxy, 'bbb')) store_result(db, ('fx', proxy, json.dumps(row))) + other = 'http://11.0.0.2:80' + mark_seen(db, (other, 'other')) + store_result(db, ('other-profile', other, json.dumps(measured(other, cfg)))) db.commit() self.assertEqual(p.source_map(db)[proxy], ('aaa', 'bbb')) + self.assertEqual(p.source_map(db, 'fx'), {proxy: ('aaa', 'bbb')}) report = p.export(db, 'fx', Path(temp) / 'out', min_success=1) self.assertEqual(report['source_quality'], { 'aaa': {'checked': 1, 'passed': 1}, 'bbb': {'checked': 1, 'passed': 1}}) diff --git a/tests/test_gateway_health.py b/tests/test_gateway_health.py index 10a69c2..b4be605 100644 --- a/tests/test_gateway_health.py +++ b/tests/test_gateway_health.py @@ -223,10 +223,6 @@ async def test_an_open_tunnel_is_never_moved_to_another_upstream(self): await shutdown(writer) -if __name__ == '__main__': - unittest.main() - - class RestingFallbackTests(GatewayCase): async def test_a_pool_where_every_proxy_rests_still_serves(self): # With a small pool of flaky public proxies every proxy reached two @@ -240,3 +236,85 @@ async def test_a_pool_where_every_proxy_rests_still_serves(self): self.assertIn(up.url, pool.resting) self.assertEqual(pool.available(), []) self.assertEqual((await self.get(address, '/x')).status_code, 200) + + +class RacingConnectTests(GatewayCase): + async def test_simultaneous_failure_is_recorded_before_the_winner_returns(self): + alive = 'http://11.0.0.1:80' + dead = 'http://11.0.0.2:80' + self.publish([alive, dead]) + server, _ = await self.start(max_failures=1) + running = server.gateway + running.stagger = 0 + ready = asyncio.Event() + started = set() + + class Writer: + def close(self): + pass + + async def tunnel(proxy, *_args): + started.add(proxy) + if len(started) == 2: + ready.set() + await ready.wait() + if proxy == dead: + raise gateway.UpstreamError('closed') + return (None, Writer()), None + + running._tunnel = tunnel + lease, _stream = await running.connect('127.0.0.1', self.target) + try: + self.assertEqual(lease.proxy, alive) + self.assertIn(dead, running.pool.resting) + self.assertEqual(running.pool.report(dead)['failed'], 1) + finally: + lease.release() + + + async def test_failover_session_is_pinned_to_the_winning_tunnel(self): + fast = 'http://11.0.0.1:80' + slow = 'http://11.0.0.2:80' + self.publish([fast, slow]) + server, _ = await self.start(sticky='failover') + running = server.gateway + running.stagger = 0 + + async def tunnel(proxy, *_args): + if proxy == slow: + await asyncio.Future() + return (None, None), None + + running._tunnel = tunnel + lease, _stream = await running.connect('127.0.0.1', self.target, session='one') + try: + self.assertEqual(lease.proxy, fast) + self.assertEqual(running.pool.sessions['one'][0], fast) + finally: + lease.release() + self.assertFalse(running.pool.active) + + async def test_configured_strict_session_dials_only_one_proxy(self): + proxies = ['http://11.0.0.1:80', 'http://11.0.0.2:80'] + self.publish(proxies) + server, _ = await self.start(sticky='strict') + running = server.gateway + running.stagger = 0 + dialed = [] + + async def tunnel(proxy, *_args): + dialed.append(proxy) + await asyncio.sleep(0.01) + return (None, None), None + + running._tunnel = tunnel + lease, _stream = await running.connect('127.0.0.1', self.target, session='strict') + try: + self.assertEqual(dialed, [lease.proxy]) + self.assertEqual(running.pool.sessions['strict'][0], lease.proxy) + finally: + lease.release() + + +if __name__ == '__main__': + unittest.main() From 9012eaab64876ef3079603b310cb65cc6ebc2f46 Mon Sep 17 00:00:00 2001 From: Anonymous Date: Sun, 27 Sep 2026 14:35:57 +0300 Subject: [PATCH 07/10] Document gateway and result-view fixes for 3.0.1 --- CHANGELOG.md | 2 ++ proxy_workbench/proxytool.py | 4 ++-- 2 files changed, 4 insertions(+), 2 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index 5c14a6c..13d5449 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -9,6 +9,8 @@ The format follows Keep a Changelog and semantic versioning. - **Checks run up to 16× faster on large lists.** Every scan was held at 64 checks in flight whatever the worker setting said (the memory reserved for response bodies had a fixed 64 MiB budget); the budget now grows with the workers up to 1 GiB. Ports of one IP no longer queue up together either: the first port of every host is checked first and further ports follow host by host. A full check of 663,000 collected addresses went from 25 to about 390 addresses per second with 512 workers. - **macOS app: collecting and HTTPS checks work again.** The 3.0.0 macOS builds looked for trusted certificates in a folder of the build machine, so every HTTPS source failed with a connection error and no proxies could be collected. Every HTTPS connection (sources, HTTPS checks, the gateway's HTTPS upstreams) now also trusts the certificate bundle shipped inside the app, while the system store is still used. - **The rotating proxy is far more reliable with flaky free proxies.** A proxy that does not answer within 2 seconds no longer blocks the request: the next proxy is tried alongside it and the first working tunnel wins (up to 4 proxies; a failed request now gives up after about 14 s instead of 24 s). When every proxy in a small pool was resting after failures, the gateway answered 502 instantly for five minutes; it now falls back to the proxy that is due back first. On 11 real public proxies right after a restart: 18 of 20 requests succeeded instead of 6 of 10. +- Simultaneous gateway connect outcomes now record failed proxies before serving the winner. A failover session is pinned to the tunnel that actually won, and strict sessions keep their single-proxy rule. +- The results view looks up source provenance only for its selected profile, avoiding a scan of every collected candidate on each request. - **“This data folder is already used by another run” when starting a check.** The desktop app's background job runner and scheduler hold the data lock for a moment every second, and a check started from the interface gave up on the first try. A run now waits up to 10 seconds for such a short hold; a folder that stays busy is still refused. - The interface server reads or closes the body of a request it refuses, so a refused request no longer shows up as a reset connection on Windows or as a garbled next request. diff --git a/proxy_workbench/proxytool.py b/proxy_workbench/proxytool.py index 585dc8d..11aa15b 100644 --- a/proxy_workbench/proxytool.py +++ b/proxy_workbench/proxytool.py @@ -5080,8 +5080,8 @@ def source_map(db, profile=None): on collection order. With ``profile`` only addresses that have a result in that profile are - mapped. A reader of one page of results does not need the sources of every - collected candidate: on 663k candidates the full map cost 1.3 s a request. + mapped. A GUI read for one profile does not need the sources of every + collected candidate. """ result = {} if profile is None: From d220a861063e95c071b4217dbe86566baf492a8f Mon Sep 17 00:00:00 2001 From: Anonymous Date: Sun, 27 Sep 2026 14:43:46 +0300 Subject: [PATCH 08/10] Make gateway race regression deterministic across event loops --- tests/test_gateway_health.py | 13 ++++++++++++- 1 file changed, 12 insertions(+), 1 deletion(-) diff --git a/tests/test_gateway_health.py b/tests/test_gateway_health.py index b4be605..89e9cef 100644 --- a/tests/test_gateway_health.py +++ b/tests/test_gateway_health.py @@ -9,6 +9,7 @@ """ import asyncio import unittest +from unittest import mock from proxy_workbench import gateway from tests.gateway_support import GatewayCase, shutdown @@ -263,7 +264,17 @@ async def tunnel(proxy, *_args): return (None, Writer()), None running._tunnel = tunnel - lease, _stream = await running.connect('127.0.0.1', self.target) + real_wait = asyncio.wait + + async def wait_for_both(tasks, **options): + # Deliver both finished tasks in one batch, the interleaving that + # used to leave a failed proxy unrecorded on Windows. + if len(tasks) == 2: + options.update(timeout=None, return_when=asyncio.ALL_COMPLETED) + return await real_wait(tasks, **options) + + with mock.patch.object(gateway.asyncio, 'wait', side_effect=wait_for_both): + lease, _stream = await running.connect('127.0.0.1', self.target) try: self.assertEqual(lease.proxy, alive) self.assertIn(dead, running.pool.resting) From c36358e76ad0e0f45cba8572babe966ea9be4091 Mon Sep 17 00:00:00 2001 From: Anonymous Date: Sun, 27 Sep 2026 14:48:13 +0300 Subject: [PATCH 09/10] Rest slow upstreams after a raced connection --- CHANGELOG.md | 2 +- proxy_workbench/gateway.py | 31 +++++++++++++++++++++++-------- tests/test_gateway_health.py | 23 +++++++++++++++++++++++ 3 files changed, 47 insertions(+), 9 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index 13d5449..e00df45 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -9,7 +9,7 @@ The format follows Keep a Changelog and semantic versioning. - **Checks run up to 16× faster on large lists.** Every scan was held at 64 checks in flight whatever the worker setting said (the memory reserved for response bodies had a fixed 64 MiB budget); the budget now grows with the workers up to 1 GiB. Ports of one IP no longer queue up together either: the first port of every host is checked first and further ports follow host by host. A full check of 663,000 collected addresses went from 25 to about 390 addresses per second with 512 workers. - **macOS app: collecting and HTTPS checks work again.** The 3.0.0 macOS builds looked for trusted certificates in a folder of the build machine, so every HTTPS source failed with a connection error and no proxies could be collected. Every HTTPS connection (sources, HTTPS checks, the gateway's HTTPS upstreams) now also trusts the certificate bundle shipped inside the app, while the system store is still used. - **The rotating proxy is far more reliable with flaky free proxies.** A proxy that does not answer within 2 seconds no longer blocks the request: the next proxy is tried alongside it and the first working tunnel wins (up to 4 proxies; a failed request now gives up after about 14 s instead of 24 s). When every proxy in a small pool was resting after failures, the gateway answered 502 instantly for five minutes; it now falls back to the proxy that is due back first. On 11 real public proxies right after a restart: 18 of 20 requests succeeded instead of 6 of 10. -- Simultaneous gateway connect outcomes now record failed proxies before serving the winner. A failover session is pinned to the tunnel that actually won, and strict sessions keep their single-proxy rule. +- Simultaneous gateway connect outcomes now record failed proxies before serving the winner. Proxies still connecting after another wins are marked slow, so they stop delaying each request. A failover session is pinned to the tunnel that actually won, and strict sessions keep their single-proxy rule. - The results view looks up source provenance only for its selected profile, avoiding a scan of every collected candidate on each request. - **“This data folder is already used by another run” when starting a check.** The desktop app's background job runner and scheduler hold the data lock for a moment every second, and a check started from the interface gave up on the first try. A run now waits up to 10 seconds for such a short hold; a folder that stays busy is still refused. - The interface server reads or closes the body of a request it refuses, so a refused request no longer shows up as a reset connection on Windows or as a garbled next request. diff --git a/proxy_workbench/gateway.py b/proxy_workbench/gateway.py index d1d01aa..a15f9ea 100644 --- a/proxy_workbench/gateway.py +++ b/proxy_workbench/gateway.py @@ -76,7 +76,7 @@ # itself; a target that refused must never rest a working upstream. HEALTH_GOOD = frozenset({'response', 'tunnel_bytes'}) HEALTH_TARGET = frozenset({'upstream_unavailable', 'no_response'}) -HEALTH_FAULT = frozenset({'handshake_failed', 'upstream_refused', 'closed_empty'}) +HEALTH_FAULT = frozenset({'handshake_failed', 'slow_connect', 'upstream_refused', 'closed_empty'}) HEALTH_OUTCOMES = HEALTH_GOOD | HEALTH_TARGET | HEALTH_FAULT #: Half-life of a gateway health observation, seconds. HEALTH_DECAY = 300.0 @@ -1344,6 +1344,7 @@ async def connect(self, host, port, forward=False, request=None, session=None, b tried = set() pending = {} failures = 0 + won = False strict = bool(session and self.pool._sticky(binding, sticky) == 'strict') def start_next(): @@ -1394,6 +1395,7 @@ def start_next(): self.pool.connected(lease.proxy, (time.monotonic() - started) * 1000) if session: self.pool.bind_session(session, lease.proxy) + won = True return lease, stream if strict and failures: break @@ -1403,19 +1405,32 @@ def start_next(): self.pool.stats['failed'] += 1 raise UpstreamError('NO_WORKING_PROXY') finally: - # Losers of the race, and everything on cancellation: no slot leaks - # and no half-open tunnel stays behind. - for task, (lease, _) in pending.items(): - task.cancel() + # A loser that was still connecting after the stagger has proved + # too slow for this gateway. Count it, or it gets picked again on + # every request and delays every other connection forever. + for task, (lease, started) in pending.items(): + unfinished = not task.done() + if unfinished: + task.cancel() lease.release() + if won and unfinished and time.monotonic() - started >= self.stagger: + self.pool.outcome(lease.proxy, 'slow_connect', detail='slow_connect') for task in pending: with contextlib.suppress(BaseException): await task - for task in pending: - if task.done() and not task.cancelled() and task.exception() is None: + for task, (lease, _) in pending.items(): + if task.cancelled(): + continue + try: stream = task.result()[0] + except (OSError, UpstreamError, ValueError, UnicodeError): + if won: + self.pool.outcome(lease.proxy, 'handshake_failed', detail=UNUSABLE) + except BaseException: + pass + else: with contextlib.suppress(Exception): - stream[1].close() if isinstance(stream, tuple) else stream.close() + stream[1].close() # --- authentication ------------------------------------------------------ diff --git a/tests/test_gateway_health.py b/tests/test_gateway_health.py index 89e9cef..1f48293 100644 --- a/tests/test_gateway_health.py +++ b/tests/test_gateway_health.py @@ -240,6 +240,29 @@ async def test_a_pool_where_every_proxy_rests_still_serves(self): class RacingConnectTests(GatewayCase): + async def test_a_slow_loser_rests_instead_of_delaying_every_request(self): + fast = 'http://11.0.0.1:80' + slow = 'http://11.0.0.2:80' + self.publish([fast, slow]) + server, _ = await self.start(max_failures=1) + running = server.gateway + running.stagger = 0.01 + + async def tunnel(proxy, *_args): + if proxy == slow: + await asyncio.Future() + return (None, None), None + + running._tunnel = tunnel + lease, _stream = await running.connect('127.0.0.1', self.target) + try: + self.assertEqual(lease.proxy, fast) + self.assertIn(slow, running.pool.resting) + self.assertEqual(running.pool.report(slow)['failed'], 1) + finally: + lease.release() + self.assertFalse(running.pool.active) + async def test_simultaneous_failure_is_recorded_before_the_winner_returns(self): alive = 'http://11.0.0.1:80' dead = 'http://11.0.0.2:80' From cbc3b8fd4bf93ebb86625b29d79fed8dd724239e Mon Sep 17 00:00:00 2001 From: Anonymous Date: Sun, 27 Sep 2026 14:53:41 +0300 Subject: [PATCH 10/10] Track staggered attempts without clock rounding --- proxy_workbench/gateway.py | 10 ++++++++-- 1 file changed, 8 insertions(+), 2 deletions(-) diff --git a/proxy_workbench/gateway.py b/proxy_workbench/gateway.py index a15f9ea..26318d6 100644 --- a/proxy_workbench/gateway.py +++ b/proxy_workbench/gateway.py @@ -1343,6 +1343,7 @@ async def connect(self, host, port, forward=False, request=None, session=None, b """ tried = set() pending = {} + staggered = set() failures = 0 won = False strict = bool(session and self.pool._sticky(binding, sticky) == 'strict') @@ -1367,6 +1368,11 @@ def start_next(): can_add = not strict and len(tried) < self.attempts done, _ = await asyncio.wait(list(pending), timeout=self.stagger if can_add else None, return_when=asyncio.FIRST_COMPLETED) + if not done and can_add: + # The wait itself establishes that these attempts have + # missed the policy deadline. Comparing clocks again in + # cleanup can lose a few milliseconds on Windows. + staggered.update(pending) winner = None for task in list(pending): if task not in done: @@ -1408,12 +1414,12 @@ def start_next(): # A loser that was still connecting after the stagger has proved # too slow for this gateway. Count it, or it gets picked again on # every request and delays every other connection forever. - for task, (lease, started) in pending.items(): + for task, (lease, _) in pending.items(): unfinished = not task.done() if unfinished: task.cancel() lease.release() - if won and unfinished and time.monotonic() - started >= self.stagger: + if won and unfinished and task in staggered: self.pool.outcome(lease.proxy, 'slow_connect', detail='slow_connect') for task in pending: with contextlib.suppress(BaseException):