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 diff --git a/CHANGELOG.md b/CHANGELOG.md index a8b8d2e..e00df45 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -2,6 +2,18 @@ The format follows Keep a Changelog and semantic versioning. +## [3.0.1] — 2026-09-27 + +### 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. +- **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. 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. + ## [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..26318d6 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 @@ -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 @@ -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': @@ -779,15 +782,36 @@ 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) - if session: + # 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 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) + 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 @@ -1019,7 +1043,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: @@ -1190,7 +1214,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 +1224,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 +1333,110 @@ 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 = {} + staggered = set() + failures = 0 + won = False + 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: - 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) + 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: + continue + 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 + 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) + won = True + 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: + # 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, _) in pending.items(): + unfinished = not task.done() + if unfinished: + 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') + 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): + await task + 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() # --- authentication ------------------------------------------------------ @@ -1866,7 +1964,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/proxy_workbench/gui.py b/proxy_workbench/gui.py index 9828757..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 {} @@ -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/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..11aa15b 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,9 @@ # generation directory a second time. ROOT = paths.PACKAGE -TLS = ssl.create_default_context() +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'} @@ -3840,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 = ( @@ -3861,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): @@ -4316,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, @@ -4425,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))] @@ -4447,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] = [] @@ -5010,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 GUI read for one profile does not need the sources of every + collected candidate. """ 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, []) @@ -8138,19 +8213,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/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_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() 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 380711e..1f48293 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 @@ -223,5 +224,131 @@ async def test_an_open_tunnel_is_never_moved_to_another_upstream(self): await shutdown(writer) +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) + + +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' + 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 + 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) + 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() 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)) 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()