diff --git a/.dockerignore b/.dockerignore index 0a175e7..a3ecfbc 100644 --- a/.dockerignore +++ b/.dockerignore @@ -2,6 +2,8 @@ .github .venv data +!proxy_workbench/data/ +!proxy_workbench/data/*.json docs tests packaging diff --git a/.github/repo-settings.json b/.github/repo-settings.json index ed2437d..91df619 100644 --- a/.github/repo-settings.json +++ b/.github/repo-settings.json @@ -1,5 +1,5 @@ { - "description": "Free proxies that actually work for YOUR sites: collects 55 public lists, checks HTTP/HTTPS/SOCKS4/SOCKS5 against your services, rates anonymity, filters by country and blacklists, then serves them as a rotating proxy, API, PAC and Clash config. Local GUI + CLI + Docker.", + "description": "Finds public proxies that work on the sites you need. Browse a 150-source catalog, collect supported feeds, check against your services, and use ready lists, a rotating gateway and an API. macOS/Windows app, CLI, Docker, 12 languages.", "homepage_from_pages": true, "has_discussions": true, "has_issues": true, diff --git a/.github/workflows/ci.yml b/.github/workflows/ci.yml index e476e41..b574a6e 100644 --- a/.github/workflows/ci.yml +++ b/.github/workflows/ci.yml @@ -157,6 +157,7 @@ jobs: docker run --rm proxy-workbench:ci --help > /dev/null docker run --rm proxy-workbench:ci clear-data --yes docker run --rm --entrypoint python proxy-workbench:ci -c "from proxy_workbench import source_catalog; assert source_catalog.load_bundled()['sources']" + docker run --rm proxy-workbench:ci preset list --json > /dev/null # The .app users install. This job builds it on a real macOS runner, # starts it, and checks the page it serves, its per-user data folder, a second diff --git a/CHANGELOG.md b/CHANGELOG.md index e00df45..33a4b20 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -2,6 +2,15 @@ The format follows Keep a Changelog and semantic versioning. +## [3.0.2] — 2026-09-27 + +### Fixed + +- Source retries now publish candidates only from a completed response. Byte and line limits keep valid following addresses, and the scan enforces each target's request cap even when one probe makes several requests. +- The API rejects invalid numbers and cursors cleanly, executes concurrent retries with the same idempotency key once, pages scoped jobs and events correctly, and releases SSE slots when clients leave. +- API key forms submit the selected permissions. The import wizard offers every CSV/JSON column, accepts the first or duplicate column, and discards previews after the input changes. Switched tabs remain visible when a browser pauses animations. +- Docker includes the service preset catalog. Its Compose gateway can start on the configured network bind, and the CLI shows a generated password or usable authenticated examples without printing a supplied secret. + ## [3.0.1] — 2026-09-27 ### Fixed diff --git a/Dockerfile b/Dockerfile index 0a29c13..ecb9d96 100644 --- a/Dockerfile +++ b/Dockerfile @@ -13,6 +13,7 @@ COPY requirements.txt ./ RUN python -m pip install --requirement requirements.txt COPY proxy_workbench/*.py proxy_workbench/*.json ./proxy_workbench/ +COPY proxy_workbench/data/ ./proxy_workbench/data/ COPY proxytool.py service.example.json ./ RUN useradd --create-home --uid 10001 workbench \ && mkdir -p /app/data \ diff --git a/README.md b/README.md index 0907c0c..f492d4b 100644 --- a/README.md +++ b/README.md @@ -103,7 +103,7 @@ Pick one way to install: | **Windows app** | Download `proxy-workbench-…-windows-x64-setup.exe` from the [latest release](https://github.com/DavidVoitenko/proxy-workbench/releases/latest) and run it; there is a portable `.zip` too, and a separate `proxy-workbench-cli.exe` for the command line | nothing else | | **pipx** (Windows, macOS, Linux) | `pipx install git+https://github.com/DavidVoitenko/proxy-workbench` then `proxy-workbench` | Python 3.11+ and [pipx](https://pypa.io/pipx/) | | **Source folder** | Download the code (**Code → Download ZIP** or `git clone`), then double-click `Start.bat` (Windows) / `Start.command` (macOS) or run `./run.sh` (Linux) | Python 3.11+ | -| **Docker** | `docker compose up -d` with the bundled [`compose.yml`](compose.yml) (checker + API + rotating proxy) | Docker | +| **Docker** | Set `PROXY_WORKBENCH_API_TOKEN`, then run `docker compose up -d` with the bundled [`compose.yml`](compose.yml) (checker + API + rotating proxy) | Docker | `proxy-workbench` without arguments starts the application: the interface opens in your browser, and on macOS a menu bar item shows what is happening and offers pause, start, “start at login” and quit. Starting it a second time reaches the one that is already running instead of opening a rival. `proxy-workbench run …` and the other commands below work the same way as `./run.sh …`. @@ -288,7 +288,7 @@ curl -x socks5h://127.0.0.1:8899 https://example.org/ - If a proxy fails, the same connection is retried through another one (up to 3). A proxy that fails twice rests for 5 minutes. - Plain `http://` requests reach HTTP proxies directly, because many of them allow CONNECT only to port 443. -On a server, start it with `./run.sh gateway` and narrow the pool with the usual filters, for example `gateway --protocol socks5 --country DE --max-latency 1500`. Binding to a network address (`--host 0.0.0.0`) requires a password: `--gateway-token `, or the `PROXY_WORKBENCH_GATEWAY_TOKEN` variable; when it is not given the gateway makes one and prints it. Clients then log in with any user name and that password, over HTTP Basic or SOCKS5 user/password. +On a server, start it with `./run.sh gateway` and narrow the pool with the usual filters, for example `gateway --protocol socks5 --country DE --max-latency 1500`. Binding to a network address requires `--host 0.0.0.0 --lan` and a password: `--gateway-token `, or the `PROXY_WORKBENCH_GATEWAY_TOKEN` variable; when it is not given the gateway makes one and prints it. Clients then log in with any user name and that password, over HTTP Basic or SOCKS5 user/password. **The gateway password is not the API token.** They are separate identities on purpose: whoever knows the password you handed to a phone must not be able to read the published snapshot, and a leaked API token must not be a working proxy. Pass `--api-token` to `serve` and `--gateway-token` to `gateway`; `compose.yml` shows both variables. @@ -380,11 +380,11 @@ docker run --rm --user "$(id -u):$(id -g)" -v "$PWD/data:/app/data" \ Serve fresh proxies to other containers: one container re-checks, the other answers API requests from the same data folder. ```sh -docker run -d --name pw-check -v "$PWD/data:/app/data" proxy-workbench run --want 50 --watch 30 -docker run -d --name pw-api -p 127.0.0.1:8765:8765 -e PROXY_WORKBENCH_API_TOKEN=change-me \ +docker run -d --name pw-check --user "$(id -u):$(id -g)" -v "$PWD/data:/app/data" proxy-workbench run --want 50 --watch 30 +docker run -d --name pw-api --user "$(id -u):$(id -g)" -p 127.0.0.1:8765:8765 -e PROXY_WORKBENCH_API_TOKEN=change-me \ -v "$PWD/data:/app/data" proxy-workbench serve --host 0.0.0.0 -docker run -d --name pw-gateway -p 127.0.0.1:8899:8899 -e PROXY_WORKBENCH_API_TOKEN=change-me \ - -v "$PWD/data:/app/data" proxy-workbench gateway --host 0.0.0.0 +docker run -d --name pw-gateway --user "$(id -u):$(id -g)" -p 127.0.0.1:8899:8899 -e PROXY_WORKBENCH_GATEWAY_TOKEN=another-secret \ + -v "$PWD/data:/app/data" proxy-workbench gateway --host 0.0.0.0 --lan ``` Every release also publishes a ready image to the GitHub Container Registry. It appears under **Packages** in the repository sidebar as `ghcr.io//proxy-workbench:` and `:latest`. Results land in the mounted `data/` folder exactly as with a local install. diff --git a/README.ru.md b/README.ru.md index a5c3a87..a993049 100644 --- a/README.ru.md +++ b/README.ru.md @@ -106,7 +106,7 @@ | **Программа для Windows** | Скачайте `proxy-workbench-…-windows-x64-setup.exe` из [последнего релиза](https://github.com/DavidVoitenko/proxy-workbench/releases/latest) и запустите установщик; есть portable-архив `.zip` и отдельный `proxy-workbench-cli.exe` для командной строки | больше ничего | | **pipx** (Windows, macOS, Linux) | `pipx install git+https://github.com/DavidVoitenko/proxy-workbench`, затем `proxy-workbench` | Python 3.11+ и [pipx](https://pypa.io/pipx/) | | **Папка с кодом** | Скачайте код (**Code → Download ZIP** или `git clone`) и запустите, как в таблице ниже | Python 3.11+ | -| **Docker** | `docker compose up -d` с готовым [`compose.yml`](compose.yml): проверка + API + ротирующий прокси | Docker | +| **Docker** | Задайте `PROXY_WORKBENCH_API_TOKEN`, затем запустите `docker compose up -d` с готовым [`compose.yml`](compose.yml): проверка + API + ротирующий прокси | Docker | `proxy-workbench` без аргументов запускает приложение: интерфейс открывается в браузере, а на macOS в меню-баре появляется значок с состоянием и пунктами «пауза», «запуск проверки», «запускать при входе» и «выход». Повторный запуск обращается к уже работающему экземпляру, а не поднимает вторую копию. `proxy-workbench run …` и остальные команды работают так же, как `./run.sh …`. @@ -379,7 +379,7 @@ curl -x socks5h://127.0.0.1:8899 https://example.org/ - Если прокси не ответил, то же соединение повторяется через другой (до 3 раз). Прокси, который ошибся дважды, отдыхает 5 минут. - Обычные `http://`-запросы идут в HTTP-прокси напрямую, потому что многие из них разрешают CONNECT только на порт 443. -На сервере шлюз запускает `./run.sh gateway`; пул сужается обычными фильтрами, например `gateway --protocol socks5 --country DE --max-latency 1500`. Для сетевого адреса (`--host 0.0.0.0`) нужен пароль: `--gateway-token <секрет>` или переменная `PROXY_WORKBENCH_GATEWAY_TOKEN`; если он не задан, шлюз придумывает свой и печатает его. Клиенты входят с любым именем и этим паролем через HTTP Basic или логин/пароль SOCKS5. +На сервере шлюз запускает `./run.sh gateway`; пул сужается обычными фильтрами, например `gateway --protocol socks5 --country DE --max-latency 1500`. Для сетевого адреса нужны `--host 0.0.0.0 --lan` и пароль: `--gateway-token <секрет>` или переменная `PROXY_WORKBENCH_GATEWAY_TOKEN`; если он не задан, шлюз придумывает свой и печатает его. Клиенты входят с любым именем и этим паролем через HTTP Basic или логин/пароль SOCKS5. **Пароль шлюза — не токен API.** Это разные идентичности намеренно: тот, кому вы дали пароль для телефона в сети, не должен читать опубликованный снапшот, а утёкший токен API не должен работать как прокси. `serve` получает `--api-token`, `gateway` — `--gateway-token`; в `compose.yml` показаны обе переменные. @@ -471,11 +471,11 @@ docker run --rm --user "$(id -u):$(id -g)" -v "$PWD/data:/app/data" \ Раздавать свежие прокси другим контейнерам: один контейнер перепроверяет, второй отвечает на запросы API из той же папки данных. ```sh -docker run -d --name pw-check -v "$PWD/data:/app/data" proxy-workbench run --want 50 --watch 30 -docker run -d --name pw-api -p 127.0.0.1:8765:8765 -e PROXY_WORKBENCH_API_TOKEN=change-me \ +docker run -d --name pw-check --user "$(id -u):$(id -g)" -v "$PWD/data:/app/data" proxy-workbench run --want 50 --watch 30 +docker run -d --name pw-api --user "$(id -u):$(id -g)" -p 127.0.0.1:8765:8765 -e PROXY_WORKBENCH_API_TOKEN=change-me \ -v "$PWD/data:/app/data" proxy-workbench serve --host 0.0.0.0 -docker run -d --name pw-gateway -p 127.0.0.1:8899:8899 -e PROXY_WORKBENCH_API_TOKEN=change-me \ - -v "$PWD/data:/app/data" proxy-workbench gateway --host 0.0.0.0 +docker run -d --name pw-gateway --user "$(id -u):$(id -g)" -p 127.0.0.1:8899:8899 -e PROXY_WORKBENCH_GATEWAY_TOKEN=another-secret \ + -v "$PWD/data:/app/data" proxy-workbench gateway --host 0.0.0.0 --lan ``` Каждый релиз также публикует готовый образ в GitHub Container Registry. Он появляется в разделе **Packages** на странице репозитория как `ghcr.io//proxy-workbench:<версия>` и `:latest`. Результаты сохраняются в подключённую папку `data/`, как при обычной установке. diff --git a/compose.yml b/compose.yml index 11beccd..fc46279 100644 --- a/compose.yml +++ b/compose.yml @@ -36,7 +36,7 @@ services: gateway: <<: *workbench - command: ["gateway", "--host", "0.0.0.0"] + command: ["gateway", "--host", "0.0.0.0", "--lan"] environment: # Not the API token: the gateway password is a separate identity. Unset is # fine -- the gateway then generates one and prints it. diff --git a/packaging/windows-installer.iss b/packaging/windows-installer.iss index f91ea3c..0ea289f 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.1 /DOutDir=C:\path\to\dist /DSourceDir=C:\path\to\dist packaging\windows-installer.iss +; iscc /DProductVersion=3.0.2 /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/api.py b/proxy_workbench/api.py index 7640248..a54fe81 100644 --- a/proxy_workbench/api.py +++ b/proxy_workbench/api.py @@ -1545,14 +1545,23 @@ def _job_store(self): def _op_jobs_list(self, call): store, conn = self._job_store() + limit = min(int(call.query.get('limit') or 50), 500) + offset = self._offset(call) try: - items = [_job_dict(job) - for job in store.jobs(limit=min(int(call.query.get('limit') or 50), 500))] \ - if store is not None else [] + allowed = _scope_values(call.principal, 'collections') + scoped = sorted(allowed) if allowed is not None else None + if store is None: + items, total = [], 0 + else: + jobs = store.jobs(collection_ids=scoped, limit=limit, offset=offset) + items = [_job_dict(job) for job in jobs] + total = store.count_jobs(collection_ids=scoped) finally: _close(conn) - return {'items': self._guard_objects(items, call, 'id', 'jobs'), - 'stream_id': 'jobs', 'next_seq': None} + items = self._guard_objects(items, call, 'id', 'jobs') + return {'items': items, 'stream_id': 'jobs', 'total': total, + 'cursor_seq': offset + len(items) if items else None, + 'next_seq': offset + len(items) + 1 if offset + len(items) < total else None} def _op_jobs_get(self, call): store, conn = self._job_store() @@ -2061,19 +2070,21 @@ def _op_events_system(self, call): conn = self.connection() items = [] try: - rows = conn.execute( - 'SELECT e.rowid AS rid, e.job_id, e.seq, e.at, e.type, e.code, e.data_json, ' - 'j.collection_id FROM job_event e LEFT JOIN job j ON j.id = e.job_id ' - 'WHERE e.rowid > ? ORDER BY e.rowid LIMIT ?', - (after, limit if allowed is None else limit * 20)).fetchall() \ - if conn is not None else [] + sql = ('SELECT e.rowid AS rid, e.job_id, e.seq, e.at, e.type, e.code, ' + 'e.data_json, j.collection_id FROM job_event e ' + 'LEFT JOIN job j ON j.id = e.job_id WHERE e.rowid > ?') + params = [after] + if allowed is not None: + sql += f' AND j.collection_id IN ({",".join("?" * len(allowed))})' + params.extend(sorted(allowed)) + sql += ' ORDER BY e.rowid LIMIT ?' + params.append(limit) + rows = conn.execute(sql, params).fetchall() if conn is not None else [] except sqlite3.Error: rows = [] finally: _close(conn) for row in rows: - if allowed is not None and row['collection_id'] not in allowed: - continue try: data = json.loads(row['data_json'] or '{}') except (TypeError, ValueError): @@ -4503,13 +4514,16 @@ def _control(self): self.send_header('Connection', 'close') self.end_headers() if self.command == 'HEAD': + response.stream.close() return try: for frame in response.stream: self.wfile.write(frame.encode('utf-8') if isinstance(frame, str) else frame) self.wfile.flush() except OSError: - pass # the subscriber went away; the stream releases its slot itself + pass # the subscriber went away + finally: + response.stream.close() def do_GET(self): url = urlsplit(self.path) diff --git a/proxy_workbench/apiv1.py b/proxy_workbench/apiv1.py index 763a1f1..a8c821e 100644 --- a/proxy_workbench/apiv1.py +++ b/proxy_workbench/apiv1.py @@ -34,6 +34,7 @@ from __future__ import annotations import base64 +import binascii import hashlib import json import math @@ -287,13 +288,21 @@ def _validate_scalar(field, value): if field.kind in ('int', 'float'): if isinstance(value, str): try: - value = float(value) + if field.kind == 'int': + if re.fullmatch(r'[+-]?\d+', value) is None: + raise ValueError('not an integer') + value = int(value) + else: + value = float(value) except ValueError: raise field_error(field.name, tr('ожидается число', 'a number is expected')) from None if isinstance(value, bool) or not isinstance(value, (int, float)): raise field_error(field.name, tr('ожидается число', 'a number is expected')) + if isinstance(value, float) and not math.isfinite(value): + raise field_error(field.name, tr('ожидается конечное число', + 'a finite number is expected')) if field.kind == 'int': - if float(value) != int(value): + if isinstance(value, float) and not value.is_integer(): raise field_error(field.name, tr('ожидается целое', 'an integer is expected')) value = int(value) if field.minimum is not None and value < field.minimum: @@ -817,6 +826,12 @@ def __init__(self, ttl_s, maximum, clock): self._clock = clock self._records = OrderedDict() self._lock = threading.Lock() + # Serialize identical mutations from check through execution and cache + # insertion. A fixed set of locks keeps memory bounded for random keys. + self._singleflight = tuple(threading.RLock() for _ in range(64)) + + def serialized(self, bucket, key): + return self._singleflight[hash((bucket, key)) % len(self._singleflight)] def get(self, bucket, key, digest): with self._lock: @@ -897,9 +912,10 @@ def encode_cursor(stream_id, seq, digest=None): def decode_cursor(cursor): try: padded = cursor + '=' * (-len(cursor) % 4) - data = json.loads(base64.urlsafe_b64decode(padded.encode('ascii')).decode('utf-8')) + data = json.loads(base64.b64decode(padded.encode('ascii'), altchars=b'-_', + validate=True).decode('utf-8')) return str(data['s']), int(data['n']), str(data.get('d') or '') - except (ValueError, TypeError, KeyError) as exc: + except (ValueError, TypeError, KeyError, UnicodeError, binascii.Error) as exc: raise field_error('cursor', tr('нечитаемый курсор', 'unreadable cursor'), action=tr('начните выдачу заново', 'restart the listing')) from exc @@ -963,42 +979,78 @@ class Call: class EventStream: """Server-sent events of one stream, normalized to the documented frame.""" - def __init__(self, stream_id, events, clock, history=None, cursor_seq=0): + def __init__(self, stream_id, events, clock, history=None, cursor_seq=0, + on_close=None): self.stream_id = stream_id self.events = events self.clock = clock self.history = history self.cursor_seq = cursor_seq + self._on_close = on_close + self._close_lock = threading.Lock() + self._closed = False + + def close(self): + with self._close_lock: + if self._closed: + return + self._closed = True + try: + close = getattr(self.events, 'close', None) + if close is not None: + close() + finally: + if self._on_close is not None: + self._on_close() def __iter__(self): last = self.cursor_seq - for raw in self.events: - if not isinstance(raw, dict) or 'seq' not in raw or 'type' not in raw: - yield self._frame({'seq': last + 1, 'type': 'error', - 'code': 'E_VALIDATION_FIELD', - 'data': {'reason': 'event needs seq and type'}}) - return - seq = int(raw['seq']) - if seq <= last: - yield self._frame({'seq': last + 1, 'type': 'error', - 'code': 'E_CONFLICT_REVISION', - 'data': {'reason': 'seq must increase', 'seq': seq, - 'last_seq': last}}) - return - if self.history is not None: - self.history.record(self.stream_id, seq) - last = seq - yield self._frame(raw) + try: + for raw in self.events: + if not isinstance(raw, dict) or 'seq' not in raw or 'type' not in raw: + yield self._frame({'seq': last + 1, 'type': 'error', + 'code': 'E_VALIDATION_FIELD', + '_no_cursor': True, + 'data': {'reason': 'event needs seq and type'}}) + return + try: + seq = int(raw['seq']) + except (TypeError, ValueError, OverflowError): + yield self._frame({'seq': last + 1, 'type': 'error', + 'code': 'E_VALIDATION_FIELD', + '_no_cursor': True, + 'data': {'reason': 'event seq must be an integer'}}) + return + if seq <= last: + yield self._frame({'seq': last + 1, 'type': 'error', + 'code': 'E_CONFLICT_REVISION', + '_no_cursor': True, + 'data': {'reason': 'seq must increase', 'seq': seq, + 'last_seq': last}}) + return + if self.history is not None and not raw.get('_no_cursor'): + self.history.record(self.stream_id, seq) + last = seq + yield self._frame(raw) + finally: + self.close() def _frame(self, raw): + try: + at = float(raw.get('at') or self.clock()) + if not math.isfinite(at): + raise ValueError('nonfinite timestamp') + except (TypeError, ValueError, OverflowError): + at = self.clock() event = {'stream': self.stream_id, 'seq': int(raw.get('seq') or 0), - 'at': float(raw.get('at') or self.clock()), - 'type': str(raw.get('type') or 'unknown'), + 'at': at, + 'type': re.sub(r'[\x00-\x1f\x7f]', '', str(raw.get('type') or 'unknown')), 'job_id': raw.get('job_id'), 'item_id': raw.get('item_id'), 'code': raw.get('code'), 'data': raw.get('data') or {}} payload = json.dumps(event, ensure_ascii=False, sort_keys=True, default=str) - return (f'id: {encode_cursor(self.stream_id, event["seq"])}\n' - f'event: {event["type"]}\n' + cursor = '' if raw.get('_no_cursor') else \ + f'id: {encode_cursor(self.stream_id, event["seq"])}\n' + return (cursor + f'event: {event["type"]}\n' f'data: {payload}\n\n') @@ -1846,6 +1898,10 @@ def handle(self, request): def _check_config(self): """A network bind the Host check cannot serve is refused at construction.""" + if self.config.loopback_only and not is_loopback_host(self.config.host): + raise ValueError(tr( + 'Для сетевого адреса включите allow_remote_bind и задайте allowed_hosts.', + 'A network bind requires allow_remote_bind and allowed_hosts.')) if not self.config.loopback_only and not self.config.allowed_hosts: raise ValueError(tr( f'API на {self.config.host} доступно из сети: задайте allowed_hosts, ' @@ -1926,14 +1982,17 @@ def _send(self, response, head_only=False): self.send_header('Transfer-Encoding', 'chunked') self.end_headers() if head_only: - self.wfile.write(b'0\r\n\r\n') + response.stream.close() return - for frame in response.stream: - payload = frame.encode('utf-8') - self.wfile.write(b'%x\r\n' % len(payload) + payload + b'\r\n') + try: + for frame in response.stream: + payload = frame.encode('utf-8') + self.wfile.write(b'%x\r\n' % len(payload) + payload + b'\r\n') + self.wfile.flush() + self.wfile.write(b'0\r\n\r\n') self.wfile.flush() - self.wfile.write(b'0\r\n\r\n') - self.wfile.flush() + finally: + response.stream.close() server = ThreadingHTTPServer((self.config.host, self.config.port), Handler) server.daemon_threads = True @@ -1979,7 +2038,10 @@ def _handle(self, request): require_permission(principal, route.permission, self.clock()) for kind, source in route.scope: require_scope(principal, kind, params.get(source)) - idem_key = request.header('Idempotency-Key') + idem_key = request.header('Idempotency-Key') if route.mutating else None + if idem_key is not None and (not idem_key.strip() or len(idem_key) > 128): + raise field_error('Idempotency-Key', tr('ожидается непустой ключ до 128 символов', + 'a nonempty key of at most 128 characters is expected')) if route.mutating and self.config.require_idempotency_key and not idem_key: raise field_error('Idempotency-Key', tr('обязателен для изменяющих запросов', 'required for mutations'), @@ -1987,6 +2049,11 @@ def _handle(self, request): 'repeat the request with the same key after a failure')) if route.sse: return self._stream(request, route, params, principal) + if idem_key: + bucket = (principal.key_id if principal else 'anonymous', request.method, + request.path) + with self.idempotency.serialized(bucket, idem_key): + return self._invoke(request, route, params, principal, idem_key) return self._invoke(request, route, params, principal, idem_key) def _method_error(self, method, path): @@ -2359,6 +2426,7 @@ def _response(self, route, request, result, principal, extra_headers, query=None headers.append(('Content-Disposition', 'attachment; filename=' + f'"{FILE_CHARS.sub("_", str(result["filename"]))}"')) + headers.extend(self._cors_headers(request)) return Response(200, bytes(data), str(result.get('content_type') or 'application/octet-stream'), tuple(headers)) @@ -2453,53 +2521,65 @@ def _stream(self, request, route, params, principal): # "stream_id": ...}``. The page used to be iterated as it stood, which # walks the *keys* of a dict, so every event stream answered one # ``error`` frame ("event needs seq and type") and nothing else. - if isinstance(source, tuple): - declared, events = source - elif isinstance(source, Mapping) and isinstance(source.get('items'), list): - declared, events = source.get('stream_id') or stream_id, source['items'] - else: - declared, events = stream_id, source - limit = int(query.get('limit') or self.config.max_limit) - stream = EventStream(str(declared or stream_id), - self._guarded(events, limit, request, key_slot), self.clock, - self.history, cursor_seq=seq) + try: + if isinstance(source, tuple) and len(source) == 2: + declared, events = source + elif isinstance(source, Mapping) and isinstance(source.get('items'), list): + declared, events = source.get('stream_id') or stream_id, source['items'] + else: + declared, events = stream_id, source + if not isinstance(events, Iterable) or isinstance(events, (str, bytes, Mapping)): + raise ApiError('E_SERVICE_UNAVAILABLE', status=503, + details={'reason': 'events are not iterable'}) + limit = int(query.get('limit') or self.config.max_limit) + + def close_stream(): + try: + close = getattr(events, 'close', None) + if close is not None: + close() + finally: + self.key_quota.release(key_slot) + + stream = EventStream(str(declared or stream_id), + self._guarded(events, limit, request), self.clock, + self.history, cursor_seq=seq, + on_close=close_stream) + except Exception: + self.key_quota.release(key_slot) + raise return Response(200, b'', 'text/event-stream; charset=utf-8', (('Cache-Control', 'no-store'), ('X-Accel-Buffering', 'no'), - ('X-Workbench-Stream', stream.stream_id)), stream=stream) + ('X-Workbench-Stream', stream.stream_id), + *self._cors_headers(request)), stream=stream) - def _guarded(self, events, limit, request, key_slot=None): + def _guarded(self, events, limit, request): """Bound the stream and re-check the key on a timer while it is open. - The key's own concurrency slot belongs to the open stream: it is given - back in the ``finally``, whether the client read every event, stopped - half way or vanished, so a closed subscription never locks its key out. + The stream owns and releases the key's concurrency slot when it closes. """ - if not isinstance(events, Iterable): - self.key_quota.release(key_slot) - raise ApiError('E_SERVICE_UNAVAILABLE', status=503, - details={'reason': 'events are not iterable'}) last = 0 checked_at = self.clock() - try: - for index, event in enumerate(events): - if index >= limit: - yield {'seq': last + 1, 'type': 'stream.closed', 'code': 'E_LIMIT_BUDGET', - 'data': {'reason': 'limit reached'}} + for index, event in enumerate(events): + if index >= limit: + yield {'seq': last + 1, 'type': 'stream.closed', 'code': 'E_LIMIT_BUDGET', + '_no_cursor': True, 'data': {'reason': 'limit reached'}} + return + now = self.clock() + if now - checked_at >= max(self.config.sse_reverify_s, 0.0): + try: + self._recheck(request) + except ApiError as exc: + yield {'seq': last + 1, 'type': 'session.closed', 'code': exc.code, + '_no_cursor': True, 'data': {'reason': exc.message}} return - now = self.clock() - if now - checked_at >= max(self.config.sse_reverify_s, 0.0): - try: - self._recheck(request) - except ApiError as exc: - yield {'seq': last + 1, 'type': 'session.closed', 'code': exc.code, - 'data': {'reason': exc.message}} - return - checked_at = now - if isinstance(event, dict): + checked_at = now + if isinstance(event, dict): + try: last = max(last, int(event.get('seq') or 0)) - yield event - finally: - self.key_quota.release(key_slot) + except (TypeError, ValueError, OverflowError): + pass + yield event def _recheck(self, request): """A key that expired or was revoked closes the open stream, per policy.""" diff --git a/proxy_workbench/branding.py b/proxy_workbench/branding.py index 34e8e4e..87d50f4 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.1" +PRODUCT_VERSION = "3.0.2" 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/gui.py b/proxy_workbench/gui.py index ed7c0e3..bca67ce 100644 --- a/proxy_workbench/gui.py +++ b/proxy_workbench/gui.py @@ -3853,7 +3853,7 @@ def import_plan_view(plan): # module's own `MappingSuggestion`; nothing here re-guesses a role. suggestion = plan.mapping_suggestion body['mapping_suggestion'] = suggestion.to_dict() if suggestion is not None else None - body['columns'] = list(body.get('mapping') or {}) + body['columns'] = list(suggestion.columns) if suggestion is not None else list(plan.header) body['can_commit'] = not plan.needs_mapping body['modes'] = list(importer.MODES) body['formats'] = list(importer.FORMATS) diff --git a/proxy_workbench/importer.py b/proxy_workbench/importer.py index 4d61800..37655a9 100644 --- a/proxy_workbench/importer.py +++ b/proxy_workbench/importer.py @@ -353,6 +353,7 @@ class MappingSuggestion: found: dict missing: tuple ambiguous: tuple + columns: tuple @property def usable(self) -> bool: @@ -360,7 +361,8 @@ def usable(self) -> bool: def to_dict(self) -> dict: return {'mapping': self.mapping.to_dict(), 'found': self.found, - 'missing': list(self.missing), 'ambiguous': list(self.ambiguous)} + 'missing': list(self.missing), 'ambiguous': list(self.ambiguous), + 'columns': list(self.columns)} def _fold(name) -> str: @@ -380,7 +382,7 @@ def suggest_mapping(columns: Sequence) -> MappingSuggestion: missing = tuple(role for role in ROLES[:2] if role not in found and role not in ambiguous) mapping = ColumnMapping(host=found.get('host'), port=found.get('port'), scheme=found.get('scheme'), country=found.get('country')) - return MappingSuggestion(mapping, found, missing, tuple(sorted(ambiguous))) + return MappingSuggestion(mapping, found, missing, tuple(sorted(ambiguous)), tuple(columns)) def _resolve(mapping: ColumnMapping, columns: Sequence) -> dict: diff --git a/proxy_workbench/jobs.py b/proxy_workbench/jobs.py index 00f149e..a37d9e1 100644 --- a/proxy_workbench/jobs.py +++ b/proxy_workbench/jobs.py @@ -904,7 +904,14 @@ def retry(self, job_id: str, *, items: Iterable[str] | None = None, scope: Scope (idempotency_key,)).fetchone() if row is not None: existing = self._job_from_row(row) - if existing.id == job.id or existing.kind != new_kind: + origin = conn.execute( + 'SELECT data_json FROM job_event WHERE job_id=? AND seq=1', + (existing.id,)).fetchone() + try: + retry_of = json.loads(origin[0]).get('retry_of') if origin else None + except (TypeError, ValueError): + retry_of = None + if existing.id == job.id or existing.kind != new_kind or retry_of != job.id: raise IdempotencyConflict( f'key {idempotency_key!r} is already bound to job {existing.id}') return existing @@ -1289,8 +1296,14 @@ def job(self, job_id: str) -> Job: return self._job(self._conn, job_id) def jobs(self, *, kind: str | None = None, state: str | None = None, - limit: int | None = None) -> list[Job]: + collection_ids: Sequence[str] | None = None, + limit: int | None = None, offset: int = 0) -> list[Job]: clauses, params = [], [] + if collection_ids is not None: + if not collection_ids: + return [] + clauses.append(f'collection_id IN ({",".join("?" * len(collection_ids))})') + params.extend(collection_ids) if kind is not None: clauses.append('kind=?') params.append(kind) @@ -1306,8 +1319,23 @@ def jobs(self, *, kind: str | None = None, state: str | None = None, if limit is not None: sql += ' LIMIT ?' params.append(int(limit)) + if offset: + if limit is None: + sql += ' LIMIT -1' + sql += ' OFFSET ?' + params.append(max(0, int(offset))) return [self._job_from_row(row) for row in self._conn.execute(sql, params).fetchall()] + def count_jobs(self, *, collection_ids: Sequence[str] | None = None) -> int: + if collection_ids is None: + return int(self._conn.execute('SELECT count(*) FROM job').fetchone()[0]) + if not collection_ids: + return 0 + placeholders = ','.join('?' * len(collection_ids)) + return int(self._conn.execute( + f'SELECT count(*) FROM job WHERE collection_id IN ({placeholders})', + tuple(collection_ids)).fetchone()[0]) + def events(self, job_id: str, *, after_seq: int = 0, limit: int = 500, types: Sequence[str] | None = None) -> list[JobEvent]: """Events of a job with ``seq > after_seq`` (§5.7 resume rule).""" diff --git a/proxy_workbench/pipeline.py b/proxy_workbench/pipeline.py index 7a8e848..7441887 100644 --- a/proxy_workbench/pipeline.py +++ b/proxy_workbench/pipeline.py @@ -59,7 +59,7 @@ 'ITEM_STATES', 'EMITTED_ITEM_STATES', 'DEADLINE_EXCEEDED', 'E_LIMIT_BUDGET', 'E_LIMIT_BODY', 'E_VALIDATION_FIELD', 'E_VALIDATION_SCHEMA', # errors - 'PipelineError', 'ValidationError', 'BudgetExhausted', + 'PipelineError', 'ValidationError', 'BudgetExhausted', 'TargetBudgetExhausted', # budgets and governance 'Budgets', 'Limits', 'TargetPolicy', 'ResourceGate', 'ResourceSnapshot', 'Reservation', 'AdaptiveConcurrency', 'ConcurrencySample', 'HostLimiter', @@ -175,6 +175,10 @@ class BudgetExhausted(PipelineError): code = E_LIMIT_BUDGET +class TargetBudgetExhausted(BudgetExhausted): + """A single target is out of requests; the other targets may still run.""" + + # -------------------------------------------------------------------------- # clock # -------------------------------------------------------------------------- @@ -560,9 +564,12 @@ class RequestBudget: budget until the bytes have actually been consumed. """ - def __init__(self, gate: ResourceGate, reservation: Reservation): + def __init__(self, gate: ResourceGate, reservation: Reservation, *, + target_policy: TargetPolicy | None = None, target_used: dict[str, int] | None = None): self.gate = gate self.reservation = reservation + self.target_policy = target_policy + self.target_used = target_used if target_used is not None else {} self.active = False self.requests = 0 self.bytes = 0 @@ -601,6 +608,15 @@ async def reserve(self, max_bytes: int) -> Reservation: cap = max(0, int(max_bytes)) if available is not None: cap = min(cap, available) + policy = self.target_policy + if policy is not None and self.requests: + used = self.target_used.get(policy.target_id, 0) + if policy.max_requests is not None and used >= policy.max_requests: + exc = TargetBudgetExhausted(E_LIMIT_BUDGET, + f'Лимит запросов цели {policy.target_id!r} исчерпан.') + self.exhausted = exc + raise exc + self.target_used[policy.target_id] = used + 1 gate._requests += 1 gate._reserved_bytes += cap self.requests += 1 @@ -1940,57 +1956,76 @@ async def _feed_source(self, source: SourceSpec, emit) -> bool: stream = await source.open() if self.config.parse_streaming: buffer = bytearray() + discarding_line = False async for chunk in stream: if self._stopped.is_set() or not self._check_control(): return False + if not chunk: + continue if cap is not None and used >= cap: self._truncate() return False + truncated = False if cap is not None and used + len(chunk) > cap: chunk = chunk[:max(0, cap - used)] - self._truncate() + truncated = True used += len(chunk) self._counters.source_bytes += len(chunk) - # A line may straddle two chunks, so the buffer offset of the - # current chunk's first byte is tracked explicitly: the carry - # from the previous chunk must not shift the newline positions. - offset = len(buffer) - buffer += chunk start = 0 while True: index = chunk.find(b'\n', start) if index < 0: break - end = offset + (index - start) + 1 - await self._feed_line(bytes(buffer[:end - 1]), source, emit) - del buffer[:end] - offset = 0 + if not discarding_line: + part = chunk[start:index] + if len(buffer) + len(part) <= budgets.max_line_bytes: + buffer.extend(part) + await self._feed_line(bytes(buffer), source, emit) + else: + self._reject_long_line() + buffer.clear() + discarding_line = False start = index + 1 - if cap is not None and used >= cap: - self._truncate() - return False - if len(buffer) > budgets.max_line_bytes: - await self._feed_line(bytes(buffer), source, emit) - buffer.clear() - if buffer: + if not discarding_line: + tail = chunk[start:] + if len(buffer) + len(tail) <= budgets.max_line_bytes: + buffer.extend(tail) + else: + buffer.clear() + discarding_line = True + self._reject_long_line() + if truncated: + self._truncate() + return False + if buffer and not discarding_line: await self._feed_line(bytes(buffer), source, emit) else: document = bytearray() async for chunk in stream: if self._stopped.is_set() or not self._check_control(): return False + if not chunk: + continue if cap is not None and used >= cap: self._truncate() return False if cap is not None and used + len(chunk) > cap: - chunk = chunk[:max(0, cap - used)] + accepted = max(0, cap - used) + used += accepted + self._counters.source_bytes += accepted self._truncate() + return False used += len(chunk) self._counters.source_bytes += len(chunk) document += chunk await self._feed_tokens(self._parse(bytes(document)), source, emit) return True + def _reject_long_line(self) -> None: + self._counters.parsed += 1 + self._counters.rejected += 1 + self._counters.stage_failures.setdefault(STAGE_PARSE, E_LIMIT_BODY) + def _truncate(self) -> None: """A source hit its byte budget. That is a budget stop, not an interruption, and it carries the BODY code of the shared contract.""" @@ -2254,25 +2289,34 @@ async def _stage(self, item: Item, stage: str) -> StageOutcome: gate = self._gate if runner is None or gate is None: # pragma: no cover - validated in PipelineConfig raise ValidationError(E_VALIDATION_FIELD, f'Для этапа {stage!r} нет runner или gate.') - target_id = self._pick_target(stage) - # The per-target slot and its share of the request budget are taken - # together, under one lock: a check followed by a request would let N - # concurrent stages all pass the check and then overshoot the budget. - taken = False - if target_id is not None: - if not await self._target_acquire(target_id): + # Select the target and reserve its slot under the same lock. Selecting + # first lets concurrent workers all choose a target with one request + # left while another target still has room. + taken = stage == STAGE_EXPENSIVE and bool(self._target_policies) + target_id = await self._target_acquire_any() if taken else None + if taken: + if target_id is None: return StageOutcome(stage, False, code=E_LIMIT_BUDGET, failed_stage='target', - detail=f'Лимит запросов цели {target_id!r} исчерпан.') - taken = True + detail='Лимит запросов целей исчерпан.') + target_meter = {'requests': 0} try: - return await self._measure(item, stage, runner, gate, target_id, taken) + return await self._measure(item, stage, runner, gate, target_id, target_meter) finally: if taken: - await self._target_release(target_id) - - async def _measure(self, item: Item, stage: str, runner, gate, target_id, taken) -> StageOutcome: + await self._target_release(target_id, target_meter) + + async def _measure(self, item: Item, stage: str, runner, gate, target_id, target_meter) -> StageOutcome: + target_policy = self._target_policies.get(target_id) + target_remaining = None + if target_policy is not None and target_policy.max_requests is not None: + # The first request was reserved together with the target slot. + target_remaining = 1 + max(0, target_policy.max_requests - self.target_used(target_id)) + global_remaining = gate.remaining_requests() + max_requests = global_remaining if global_remaining else (target_remaining or 1) + if target_remaining is not None: + max_requests = min(max_requests, target_remaining) limit = StageLimit(stage=stage, target_id=target_id, - max_requests=gate.remaining_requests() if gate.remaining_requests() else 1, + max_requests=max_requests, max_bytes=gate.remaining_bytes() or 0, remaining_requests=gate.remaining_requests(), remaining_bytes=gate.remaining_bytes(), remaining_s=gate.remaining_s(), @@ -2291,7 +2335,7 @@ async def _measure(self, item: Item, stage: str, runner, gate, target_id, taken) # rather than as a measurement that failed. self._stop('budget_exhausted', exc.code) return StageOutcome(stage, False, code=exc.code, detail=exc.message) - budget = RequestBudget(gate, reservation) + budget = RequestBudget(gate, reservation, target_policy=target_policy, target_used=self._target_used) limit = replace(limit, budget=budget) try: try: @@ -2299,6 +2343,8 @@ async def _measure(self, item: Item, stage: str, runner, gate, target_id, taken) budget.check() except asyncio.CancelledError: raise + except TargetBudgetExhausted as exc: + outcome = StageOutcome(stage, False, code=exc.code, failed_stage='target', detail=exc.message) except BudgetExhausted as exc: self._stop('budget_exhausted', exc.code) outcome = StageOutcome(stage, False, code=exc.code, detail=exc.message) @@ -2311,6 +2357,7 @@ async def _measure(self, item: Item, stage: str, runner, gate, target_id, taken) outcome = StageOutcome(stage, False, code=code, detail=str(exc)[:200]) self.concurrency.observe_outcome(outcome) finally: + target_meter['requests'] = budget.requests if budget.active else (outcome.requests if outcome else 0) if budget.active: reservation = replace(reservation, requests=0) if budget.requests or budget.exhausted: @@ -2323,18 +2370,20 @@ async def _measure(self, item: Item, stage: str, runner, gate, target_id, taken) bytes=0 if outcome is None else outcome.bytes) if outcome is None: # pragma: no cover - only reachable via cancellation raise ValidationError(E_VALIDATION_FIELD, f'Этап {stage!r} не дал результата.') + if target_policy is not None and not budget.active and outcome.requests > 1: + # A custom runner may only report its cost instead of using the + # request meter. Record its real cost and refuse an over-budget + # verdict; product transports reserve every request before I/O. + additional = outcome.requests - 1 + self._target_used[target_id] += additional + if target_policy.max_requests is not None and self._target_used[target_id] > target_policy.max_requests: + outcome = StageOutcome(stage, False, code=E_LIMIT_BUDGET, failed_stage='target', + requests=outcome.requests, bytes=outcome.bytes, + detail=f'Лимит запросов цели {target_id!r} превышен.') return outcome # -- per-target limits -------------------------------------------------- - def _pick_target(self, stage: str) -> str | None: - """The expensive stage's target, preferring the one that has been used - least. Cheap and basic are the pipeline's own probes and have no - configured target.""" - if stage != STAGE_EXPENSIVE or not self._target_policies: - return None - return min(self._target_policies, key=lambda name: (self._target_used.get(name, 0), name)) - def target_used(self, target_id: str) -> int: """How many requests the target has already been given.""" return self._target_used.get(target_id, 0) @@ -2345,29 +2394,38 @@ def _target_exhausted(self, target_id: str) -> bool: return False return self._target_used.get(target_id, 0) >= policy.max_requests - async def _target_acquire(self, target_id: str) -> bool: - """Take one in-flight slot and one unit of the target's request budget. + async def _target_acquire_any(self) -> str | None: + """Choose a target with budget and take its in-flight slot atomically. - Both happen under the same condition lock, so a stage either gets both - or gets neither: a budget that could be overshot by concurrency is not - a budget. False means the target is out of requests, and the caller - must not run a probe for it. + A full target makes the caller wait while any target can still free a + slot. ``None`` means every configured target has spent its budget. """ - policy = self._target_policies.get(target_id, TargetPolicy(target_id)) async with self._target_cond: while True: - if policy.max_requests is not None and self._target_used.get(target_id, 0) >= policy.max_requests: - return False - if self._target_inflight.get(target_id, 0) < policy.max_inflight: - break + available = [policy for policy in self._target_policies.values() + if policy.max_requests is None + or self._target_used.get(policy.target_id, 0) < policy.max_requests] + if not available: + if any(self._target_inflight.values()): + await self._target_cond.wait() + continue + return None + free = [policy for policy in available + if self._target_inflight.get(policy.target_id, 0) < policy.max_inflight] + if free: + policy = min(free, key=lambda item: (self._target_used.get(item.target_id, 0), item.target_id)) + target_id = policy.target_id + self._target_inflight[target_id] = self._target_inflight.get(target_id, 0) + 1 + self._target_used[target_id] = self._target_used.get(target_id, 0) + 1 + return target_id await self._target_cond.wait() - self._target_inflight[target_id] = self._target_inflight.get(target_id, 0) + 1 - if policy.max_requests is not None: - self._target_used[target_id] = self._target_used.get(target_id, 0) + 1 - return True - async def _target_release(self, target_id: str) -> None: + async def _target_release(self, target_id: str, meter) -> None: async with self._target_cond: + if not meter['requests']: + # A stage that never opened a request returns its reservation + # to the target, so waiting peers can still use the cap. + self._target_used[target_id] = max(0, self._target_used.get(target_id, 0) - 1) self._target_inflight[target_id] = max(0, self._target_inflight.get(target_id, 0) - 1) self._target_cond.notify_all() diff --git a/proxy_workbench/proxytool.py b/proxy_workbench/proxytool.py index 11aa15b..5dd235f 100644 --- a/proxy_workbench/proxytool.py +++ b/proxy_workbench/proxytool.py @@ -1851,20 +1851,8 @@ def publish(): blocked=sum(r.get('blocked', 0) for r in reports), candidates=db.execute("SELECT count(*) FROM candidates").fetchone()[0])) - def add(value, protocol='http', country=None, source=None, public_only=True, seen=None): + def queue_candidate(proxy, country=None, source=None, seen=None): nonlocal pending_writes - value = value.strip() - raw = value if '://' in value else protocol+'://'+value - # A remote source never gets to name a hostname or a private address, - # whatever the flag says: only a list the user handed over locally does. - proxy = normalize(raw) if public_only else normalize_custom(raw) - if not proxy: - return 'invalid' - if denylist.match(proxy): - return 'blocked' - if seen is not None and proxy in seen['values']: - seen['duplicate'] += 1 - return 'accepted' if max_items is not None and collection_budget['items'] >= max_items: member = db.execute('SELECT 1 FROM membership WHERE collection_id=? AND endpoint_id=?', (collection_id, schema.endpoint_id(proxy))).fetchone() @@ -1872,8 +1860,6 @@ def add(value, protocol='http', country=None, source=None, public_only=True, see _collection_budget_error(collection_budget, 'SOURCE_ITEM_BUDGET') if not (isinstance(country, str) and geoip.COUNTRY_CODE.fullmatch(country.upper())): country = None - if source and seen is not None: - seen['values'].add(proxy) pending.append((proxy, schema.endpoint_id(proxy), country and country.upper(), source, seen if source else None)) pending_writes += 1 @@ -1884,6 +1870,34 @@ def add(value, protocol='http', country=None, source=None, public_only=True, see commit() return 'accepted' + def add(value, protocol='http', country=None, source=None, public_only=True, seen=None): + value = value.strip() + raw = value if '://' in value else protocol+'://'+value + # A remote source never gets to name a hostname or a private address, + # whatever the flag says: only a list the user handed over locally does. + proxy = normalize(raw) if public_only else normalize_custom(raw) + if not proxy: + return 'invalid' + if denylist.match(proxy): + return 'blocked' + if seen is not None and proxy in seen['values']: + seen['duplicate'] += 1 + return 'accepted' + if seen is not None and seen.get('defer'): + # A streamed response can break after thousands of complete lines. + # Keep its candidates aside until the HTTP body finishes; a retry + # may return a different list and must not leave the broken prefix + # in the collection or its source provenance. + seen['values'].add(proxy) + seen['staged'].append(proxy) + if country is not None: + seen['staged_country'][proxy] = country + return 'accepted' + outcome = queue_candidate(proxy, country, source, seen) + if source and seen is not None: + seen['values'].add(proxy) + return outcome + pending = [] country_columns = [name for name in ('country', 'country_at', 'country_source') if name in endpoint_columns(db)] @@ -2050,7 +2064,8 @@ async def fetch(index, kind, url, plan=None): signatures = set() started_at = clock() state = _source_state_row(db, key) if record_provenance else None - seen = {'values': set(), 'new': 0, 'duplicate': 0, 'metadata': {}} + seen = {'values': set(), 'new': 0, 'duplicate': 0, 'metadata': {}, + 'defer': False, 'staged': [], 'staged_country': {}} received = recognized = status_code = 0 reject_reasons = {} body_digest = None @@ -2129,11 +2144,31 @@ def restore_attempt(checkpoint): endpoint_count, seen['duplicate'], values) = checkpoint budget.clear() budget.update(saved_budget) + seen['staged'].clear() + seen['staged_country'].clear() if values is not None: seen['values'] = set(values) if parse_state == 'empty': parse_state = 'pending' + def publish_staged(): + """Publish only a fully read streamed attempt's candidates.""" + staged = seen['staged'] + if not staged: + return + index = 0 + try: + for index, proxy in enumerate(staged): + queue_candidate(proxy, seen['staged_country'].get(proxy), key, seen) + except SourceFetchError: + # An item budget may stop publication halfway through. + # The accepted counter must describe the stored prefix. + seen['values'].difference_update(staged[index:]) + raise + finally: + seen['staged'] = [] + seen['staged_country'].clear() + def consume_candidate(): nonlocal candidate_count if candidate_count >= max_source_candidates: @@ -2420,8 +2455,13 @@ async def request_source(request_url, expected_page=None, headers=None): else: checkpoint = attempt_checkpoint() try: - async with asyncio.timeout(timeout): - data = await request_source(candidate, page, headers) + seen['defer'] = kind not in CATALOG_ADAPTER_KINDS and kind != 'geonode' + try: + async with asyncio.timeout(timeout): + data = await request_source(candidate, page, headers) + finally: + seen['defer'] = False + publish_staged() page_url = candidate succeeded = True error = None @@ -2443,7 +2483,15 @@ async def request_source(request_url, expected_page=None, headers=None): if error == 'SOURCE_CANDIDATE_LIMIT' and seen['values']: # A streamed page cut off by the candidate limit # was read, and what it delivered is kept. + try: + publish_staged() + except SourceFetchError as exc: + error = exc.code pages += 1 + elif error not in COLLECT_BUDGET_ERRORS: + # A failed retry or mirror is not the source's + # answer. Earlier completed pages remain intact. + restore_attempt(checkpoint) break pages += 1 if not_modified: @@ -6240,6 +6288,9 @@ def parser(): help=tr('CSV-база DB-IP Country Lite; по умолчанию data/geoip/' + geoip.DB_NAME, 'DB-IP Country Lite CSV; default data/geoip/' + geoip.DB_NAME)) p.add_argument('--min-success', type=float, default=2/3, help=tr('минимальная доля успехов КАЖДОГО target, 0..1', 'minimum success share for EACH target, 0..1')) p.add_argument('--host', default='127.0.0.1', help=tr('serve: адрес локального API; по умолчанию только этот компьютер', 'serve: API address; default is this computer only')) + p.add_argument('--lan', action='store_true', + help=tr('gateway: явно разрешить прослушивание сетевого адреса', + 'gateway: explicitly allow listening on a network address')) p.add_argument('--port', type=int, default=None, help=tr('serve/gateway: порт; по умолчанию 8765 для API и 8899 для шлюза', 'serve/gateway: port; default 8765 for the API and 8899 for the gateway')) @@ -6530,16 +6581,27 @@ def run_gateway(args, countries): async def run(): server = await gateway.start(args.data, args.host, port, gateway_token, filters, args.rotate, - max(0, args.max_per_proxy), max(0.0, args.session_ttl) * 60) + max(0, args.max_per_proxy), max(0.0, args.session_ttl) * 60, + lan=args.lan) pool = server.gateway.pool - shown = f'[{args.host}]' if ':' in args.host else args.host + shown_host = server.bind.published_host if args.lan and gateway.is_loopback(args.host) else args.host + shown = f'[{shown_host}]' if ':' in shown_host else shown_host address = f'{shown}:{server.sockets[0].getsockname()[1]}' - print(tr(f'Ротирующий прокси: {address} (HTTP и SOCKS5 TCP), в пуле {len(pool.refresh())} прокси. Ctrl+C — остановить.', - f'Rotating proxy: {address} (HTTP and SOCKS5 TCP), {len(pool.refresh())} proxies in the pool. Ctrl+C to stop.'), + available = len(pool.matching(binding=pool.default_binding)) + print(tr(f'Ротирующий прокси: {address} (HTTP и SOCKS5 TCP), в пуле {available} прокси. Ctrl+C — остановить.', + f'Rotating proxy: {address} (HTTP and SOCKS5 TCP), {available} proxies in the pool. Ctrl+C to stop.'), flush=True) - print(f' curl -x http://{address} https://example.org/', flush=True) - print(f' curl -x http://country-de-session-1:x@{address} https://example.org/', flush=True) - print(f' curl http://{address}/status', flush=True) + if server.gateway.token_origin == 'generated': + print(tr(f'Пароль шлюза (сохраните): {server.gateway.token}', + f'Gateway password (save it): {server.gateway.token}'), flush=True) + if server.gateway.token: + print(f' curl -x http://workbench:PASSWORD@{address} https://example.org/ # replace PASSWORD', + flush=True) + print(f' curl -x http://country-de-session-1:PASSWORD@{address} https://example.org/', + flush=True) + else: + print(f' curl -x http://{address} https://example.org/', flush=True) + print(f' curl http://{address}/status', flush=True) async with server: await server.serve_forever() try: @@ -8153,6 +8215,8 @@ def main(argv=None): utf8_output() p = parser() args = p.parse_intermixed_args(argv) + if args.lan and args.command != 'gateway': + p.error('--lan is available only for gateway') if (min(args.attempts, args.workers, args.max_bytes) < 1 or args.top < 0 or not math.isfinite(args.rate) or args.rate < 0 or not math.isfinite(args.timeout) or args.timeout <= 0 diff --git a/proxy_workbench/ui/app.js b/proxy_workbench/ui/app.js index a58192c..be644a3 100644 --- a/proxy_workbench/ui/app.js +++ b/proxy_workbench/ui/app.js @@ -6307,7 +6307,11 @@ async function sourceAction(path, button, message) { } catch (error) { toast(error.message, true); } finally { - if (button) button.disabled = false; + if (button) { + button.disabled = false; + button.textContent = button.dataset.busyLabel; + delete button.dataset.busyLabel; + } } } @@ -8566,7 +8570,7 @@ function setupServiceCatalogListeners() { // click on Apply replays instead of duplicating. // --------------------------------------------------------------------------- -let importState = {plan: null, source: null, mapping: null, busy: false}; +let importState = {plan: null, source: null, mapping: null, busy: false, revision: 0}; // One place decides whether "Apply" is clickable, so a finishing request can // never re-enable a button the preview had just disabled for "nothing to do". @@ -8588,6 +8592,17 @@ function setImportBusy(busy) { updateImportButtons(); } +function invalidateImportPreview() { + importState.revision += 1; + importState.plan = null; + importState.busy = false; + const preview = $('import-preview-body'); + if (preview) preview.innerHTML = ''; + const count = $('imp-count'); + if (count) count.textContent = '0'; + updateImportButtons(); +} + function importNote(text) { const node = $('import-note'); if (node) node.textContent = text || ''; @@ -8625,14 +8640,22 @@ function renderImportMapping(suggestion, problem) { const missing = (problem && problem.missing) || []; const ambiguous = (problem && problem.ambiguous) || []; if (!suggestion) { box.hidden = true; return; } - const columns = Object.values(suggestionMap).filter(value => value !== null && value !== undefined); + const columns = Array.isArray(suggestion.columns) ? suggestion.columns + : Object.values(suggestionMap).filter(value => value !== null && value !== undefined); if (!columns.length) { box.hidden = true; return; } const roles = ['host', 'port', 'scheme', 'country']; grid.innerHTML = roles.map(role => { const chosen = suggestionMap[role]; const warn = missing.includes(role) || ambiguous.includes(role); + const chosenIndex = chosen === null || chosen === undefined ? -1 + : typeof chosen === 'number' ? chosen + : columns.findIndex(name => String(name) === String(chosen)); const options = [``] - .concat(columns.map(name => ``)) + .concat(columns.map((name, index) => { + const repeated = columns.filter(other => String(other) === String(name)).length > 1; + const label = repeated || !name ? `${name || '—'} (#${index + 1})` : name; + return ``; + })) .join(''); return ` ${esc(t('imp.mapping' + role.charAt(0).toUpperCase() + role.slice(1)))} @@ -8710,14 +8733,16 @@ function mappingFromForm() { if (!grid || grid.closest('#import-mapping').hidden) return null; const mapping = {}; grid.querySelectorAll('[data-mapping-role]').forEach(node => { - mapping[node.dataset.mappingRole] = node.value || null; + mapping[node.dataset.mappingRole] = node.value === '' ? null : Number(node.value); }); - return (mapping.host && mapping.port) ? mapping : null; + return (mapping.host !== null && mapping.port !== null) ? mapping : null; } async function runImportPreview() { const source = currentImportSource(); + invalidateImportPreview(); if (!source) { importNote(t('imp.needFile')); return; } + const revision = importState.revision; setImportBusy(true); try { const plan = await api('/api/import/preview', { @@ -8729,6 +8754,7 @@ async function runImportPreview() { name: source.name, channel: source.channel }); + if (revision !== importState.revision) return; importState.plan = plan; importState.source = source; renderImportMapping(plan.mapping_suggestion, plan.mapping_problem); @@ -8737,11 +8763,13 @@ async function runImportPreview() { if (count) count.textContent = fmt((plan.counts || {}).total || 0); if (plan.needs_mapping) importNote(t('imp.needMapping')); } catch (error) { - importState.plan = null; - importNote(''); - toast(error.message, true); + if (revision === importState.revision) { + importState.plan = null; + importNote(''); + toast(error.message, true); + } } finally { - setImportBusy(false); + if (revision === importState.revision) setImportBusy(false); } } @@ -8815,6 +8843,8 @@ function setupImportListeners() { file.onchange = async () => { const chosen = file.files && file.files[0]; if (!chosen) return; + invalidateImportPreview(); + importState.source = null; if (chosen.size > 32 * 1024 * 1024) { importNote(t('error.fileTooLarge')); return; @@ -8836,15 +8866,28 @@ function setupImportListeners() { } const paste = $('import-text'); if (paste) { - paste.oninput = () => { importState.source = null; }; + paste.oninput = () => { importState.source = null; invalidateImportPreview(); }; } + const collection = $('import-collection'); + if (collection) collection.onchange = () => { + invalidateImportPreview(); + if (currentImportSource()) runImportPreview(); + }; + const mapping = $('import-mapping-grid'); + if (mapping) mapping.onchange = () => invalidateImportPreview(); if ($('import-preview')) $('import-preview').onclick = runImportPreview; if ($('import-commit')) $('import-commit').onclick = runImportCommit; if ($('imp-refresh')) $('imp-refresh').onclick = loadImportHistory; const mode = $('import-mode'); - if (mode) mode.onchange = () => { if (importState.source) runImportPreview(); }; + if (mode) mode.onchange = () => { + invalidateImportPreview(); + if (currentImportSource()) runImportPreview(); + }; const format = $('import-format'); - if (format) format.onchange = () => { if (importState.source) runImportPreview(); }; + if (format) format.onchange = () => { + invalidateImportPreview(); + if (currentImportSource()) runImportPreview(); + }; } // --------------------------------------------------------------------------- @@ -9454,9 +9497,7 @@ function renderPermissions(view) { keysState.permissions = view.permissions || []; const grid = $('keys-permissions'); if (!grid) return; - const groups = [ - ['read.*', 'read.permissions' in view ? '' : null] - ]; + const selected = new Set(selectedPermissions()); const buckets = [ {label: 'read.*', items: view.read_permissions || []}, {label: 'write.*', items: view.write_permissions || []}, @@ -9472,21 +9513,32 @@ function renderPermissions(view) {
${bucket.items.map(name => ` -
`).join(''); grid.querySelectorAll('[data-perm-group]').forEach(box => { box.onchange = () => { const target = box.checked; - grid.querySelectorAll('[data-permission]').forEach(node => { - if (node.value.startsWith(box.dataset.permGroup) || - (box.dataset.permGroup === 'export.secret' && node.value === 'export.secret')) { - node.checked = target; - } + box.closest('fieldset').querySelectorAll('[data-permission]').forEach(node => { + node.checked = target; }); + syncPermissionGroups(); }; }); - void groups; + grid.querySelectorAll('[data-permission]').forEach(node => { + node.onchange = syncPermissionGroups; + }); + syncPermissionGroups(); +} + +function syncPermissionGroups() { + const grid = $('keys-permissions'); + if (!grid) return; + grid.querySelectorAll('[data-perm-group]').forEach(box => { + const items = Array.from(box.closest('fieldset').querySelectorAll('[data-permission]')); + box.checked = items.length > 0 && items.every(node => node.checked); + box.indeterminate = items.some(node => node.checked) && !box.checked; + }); } function selectedPermissions() { @@ -9501,11 +9553,7 @@ function setPermissions(names) { if (!grid) return; const wanted = new Set(names || []); grid.querySelectorAll('[data-permission]').forEach(node => { node.checked = wanted.has(node.value); }); - grid.querySelectorAll('[data-perm-group]').forEach(box => { - const items = Array.from(grid.querySelectorAll('[data-permission]')) - .filter(node => node.value.startsWith(box.dataset.permGroup)); - box.checked = items.length > 0 && items.every(node => node.checked); - }); + syncPermissionGroups(); } function showKeySecret(issued, message) { @@ -9554,8 +9602,8 @@ function renderKeys(view) {
${keyStateBadge(row.state)}
- - ${row.state === 'disabled' + ${row.state !== 'revoked' ? `` : ''} + ${row.state === 'revoked' ? '' : row.state === 'disabled' ? `` : ``} ${row.state !== 'revoked' diff --git a/proxy_workbench/ui/style.css b/proxy_workbench/ui/style.css index 7b0d3e4..7738937 100644 --- a/proxy_workbench/ui/style.css +++ b/proxy_workbench/ui/style.css @@ -653,7 +653,7 @@ header { .page.active { display: block !important; - animation: pageFadeSlide 180ms ease-out forwards; + opacity: 1; will-change: opacity, transform; } @@ -11701,14 +11701,12 @@ aside.sidebar nav button.nav.active { animation: pw-tabIndicatorIn 0.22s var(--ease-spring, cubic-bezier(0.16,1,0.3,1)) forwards; } -/* --- 4b. Panel content entrance (class applied by app.js on tab switch) --- - app.js should: add .panel-enter to the newly activated .page.active, then - remove it on animationend. */ +/* Keep switched panels visible even when a browser suspends animation clocks. */ .panel-enter { - animation: pw-scaleIn 0.22s var(--ease-spring, cubic-bezier(0.16,1,0.3,1)) both !important; + animation: none !important; } -/* Page active baseline — keep existing pageFadeSlide, just ensure will-change */ +/* Page active baseline */ .page.active { will-change: opacity, transform; } @@ -21791,27 +21789,8 @@ tbody tr:nth-child(even) td, .chip, .filter-chip, .segmented-btn, .recheck-btn ========================================================================== */ -/* ═══════════════════════════════════════════════════════════════════════════ - 1. OFFSCREEN TAB OPTIMIZATION (content-visibility & intrinsic sizing) - Verified against index.html (8 pages: #page-scan..#page-help) & style.css:650 - Offscreen tabs cost 0 rendering/layout work until activated. - ═══════════════════════════════════════════════════════════════════════════ */ - -.page { - content-visibility: auto; - contain-intrinsic-size: 1px 800px; - contain-intrinsic-size: auto 800px; -} - -.page:not(.active) { - content-visibility: auto; - contain-intrinsic-size: 1px 800px; - contain-intrinsic-size: auto 800px; -} - -.page.active { - content-visibility: visible !important; -} +/* Inactive pages already use display:none; extra content containment can + leave the newly selected page unpainted in WebKit. */ /* ═══════════════════════════════════════════════════════════════════════════ diff --git a/pyproject.toml b/pyproject.toml index 198618b..8650237 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.1" +version = "3.0.2" 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_apiv1_resilience.py b/tests/test_apiv1_resilience.py new file mode 100644 index 0000000..1f89292 --- /dev/null +++ b/tests/test_apiv1_resilience.py @@ -0,0 +1,270 @@ +"""Regression tests for API validation, concurrency, streams and scoped reads.""" +from concurrent.futures import ThreadPoolExecutor +import http.client +import json +from pathlib import Path +import sqlite3 +import tempfile +import threading +import time +import unittest + +from proxy_workbench import api, apiv1, db, jobs + + +class Keys(apiv1.KeyStore): + def __init__(self, *, collections=(), concurrency=None): + self.principal = apiv1.Principal( + key_id='reader', kind='api_key', collections=collections, + concurrency=concurrency, permissions=frozenset(apiv1.PERMISSIONS)) + + def verify(self, secret): + return self.principal if secret == 'test-key' else None + + def audit(self, record): + pass + + +class Service(apiv1.Service): + def __init__(self): + self.calls = 0 + self.version = 1 + self.events = [{'seq': 1, 'type': 'job.note', 'data': {}}] + self.started = threading.Event() + self.release = threading.Event() + + def queue_state(self): + return {'depth': 0, 'capacity': 10} + + def invoke(self, operation, call): + if operation == 'results.list': + return {'items': [{'version': self.version}], 'stream_id': 'results', + 'next_seq': None} + if operation == 'collections.create': + self.calls += 1 + self.started.set() + self.release.wait(2) + return {'id': 'one', 'revision': 1} + if operation == 'events.system': + return {'items': self.events, 'stream_id': 'system'} + if operation == 'exports.download': + return {'data': b'proxy-list\n', 'filename': 'proxies.txt', + 'content_type': 'text/plain; charset=utf-8'} + raise NotImplementedError(operation) + + +def request(method, path, *, query='', body=None, idem=None): + headers = {'Host': '127.0.0.1', 'Authorization': 'Bearer test-key'} + if idem is not None: + headers['Idempotency-Key'] = idem + if body is not None: + headers['Content-Type'] = 'application/json' + return apiv1.Request(method, path, query=query, headers=headers, + body=json.dumps(body).encode() if body is not None else b'') + + +class ResilienceTests(unittest.TestCase): + def setUp(self): + self.service = Service() + self.control = apiv1.ApiV1(self.service, Keys()) + + def test_nonfinite_numbers_and_invalid_cursor_return_validation_errors(self): + for query in ('profile_revision=nan', 'max_latency_ms=inf', + 'min_mbps=nan', 'cursor=____', 'cursor=%21%21%21%21'): + with self.subTest(query=query): + answer = self.control.handle(request('GET', '/v1/results', query=query)) + self.assertEqual(answer.status_code, 400, answer.body) + answer = self.control.handle(request( + 'POST', '/v1/checks/quick-test', body={'endpoint': 'http://192.0.2.1:80', + 'max_seconds': float('nan')}, idem='nan')) + self.assertEqual(answer.status_code, 400, answer.body) + + def test_remote_bind_needs_explicit_network_configuration(self): + with self.assertRaises(ValueError): + apiv1.ApiV1(self.service, Keys(), apiv1.ApiConfig( + host='0.0.0.0', allowed_hosts=('service.example',))) + + def test_read_is_not_cached_by_a_mutation_idempotency_header(self): + first = self.control.handle(request('GET', '/v1/results', idem='read')) + self.service.version = 2 + second = self.control.handle(request('GET', '/v1/results', idem='read')) + self.assertEqual(first.json()['items'][0]['version'], 1) + self.assertEqual(second.json()['items'][0]['version'], 2) + + def test_allowed_browser_origin_can_read_streams_and_artifacts(self): + origin = 'https://app.example' + self.control = apiv1.ApiV1( + self.service, Keys(), apiv1.ApiConfig(allowed_origins=(origin,), + cors_origins=(origin,))) + for path in ('/v1/events', '/v1/exports/example/download/proxies.txt'): + with self.subTest(path=path): + call = request('GET', path) + call = apiv1.Request(call.method, call.path, call.query, + {**call.headers, 'Origin': origin}, call.body) + answer = self.control.handle(call) + self.assertEqual(answer.status_code, 200, answer.body) + self.assertEqual(dict(answer.headers)['Access-Control-Allow-Origin'], origin) + if answer.stream is not None: + answer.stream.close() + + def test_malformed_service_event_cannot_break_or_inject_stream_frames(self): + stream = apiv1.EventStream('system', [ + {'seq': 1, 'type': 'job.note\ndata: forged', 'at': float('nan')}, + {'seq': 'invalid', 'type': 'job.note'}, + ], lambda: 123.0) + frames = list(stream) + self.assertEqual(len(frames), 2) + self.assertIn('event: job.notedata: forged\n', frames[0]) + self.assertNotIn('\ndata: forged\n', frames[0]) + self.assertIn('"at": 123.0', frames[0]) + self.assertIn('event: error\n', frames[1]) + + def test_synthetic_stream_close_does_not_advance_resume_cursor(self): + self.service.events = [ + {'seq': 1, 'type': 'job.note', 'data': {}}, + {'seq': 2, 'type': 'job.note', 'data': {}}, + ] + response = self.control.handle(request('GET', '/v1/events', query='limit=1')) + frames = list(response.stream) + self.assertEqual(len(frames), 2) + self.assertIn('id: ', frames[0]) + self.assertIn('event: stream.closed\n', frames[1]) + self.assertNotIn('id: ', frames[1]) + + def test_concurrent_replay_executes_a_mutation_once(self): + call = request('POST', '/v1/collections', body={'name': 'one'}, idem='same') + with ThreadPoolExecutor(max_workers=2) as pool: + first = pool.submit(self.control.handle, call) + self.assertTrue(self.service.started.wait(2)) + second = pool.submit(self.control.handle, call) + deadline = time.monotonic() + 2 + while self.control.stats['requests'] < 2 and time.monotonic() < deadline: + time.sleep(0.001) + self.assertEqual(self.control.stats['requests'], 2) + time.sleep(0.05) + self.service.release.set() + answers = (first.result(3), second.result(3)) + self.assertEqual([answer.status_code for answer in answers], [200, 200]) + self.assertEqual(self.service.calls, 1) + + def test_head_event_stream_releases_key_slot_and_sends_no_body(self): + control = apiv1.ApiV1(self.service, Keys(concurrency=1), + apiv1.ApiConfig(port=0)) + server = control.make_server() + thread = threading.Thread(target=server.serve_forever, daemon=True) + thread.start() + try: + conn = http.client.HTTPConnection('127.0.0.1', server.server_port, timeout=3) + headers = {'Authorization': 'Bearer test-key'} + conn.request('HEAD', '/v1/events', headers=headers) + head = conn.getresponse() + self.assertEqual(head.status, 200) + self.assertEqual(head.read(), b'') + deadline = time.monotonic() + 2 + while control.key_quota.active('reader') and time.monotonic() < deadline: + time.sleep(0.01) + self.assertEqual(control.key_quota.active('reader'), 0) + conn.request('GET', '/v1/events', headers=headers) + response = conn.getresponse() + self.assertEqual(response.status, 200) + self.assertIn(b'event: job.note', response.read()) + conn.close() + finally: + server.shutdown() + server.server_close() + thread.join(3) + + +class ScopedJobReadsTests(unittest.TestCase): + def test_retry_key_cannot_return_an_unrelated_job(self): + conn = sqlite3.connect(':memory:') + try: + jobs.install_schema(conn) + store = jobs.JobStore(conn) + scope = jobs.Scope('mine', 'profile', 1) + first = store.submit('scan', scope, [jobs.QueueItem('endpoint-one')]) + second = store.submit('scan', scope, [jobs.QueueItem('endpoint-two')]) + retried = store.retry(first.id, idempotency_key='retry-key') + self.assertEqual(store.retry(first.id, idempotency_key='retry-key').id, + retried.id) + with self.assertRaises(jobs.IdempotencyConflict): + store.retry(second.id, idempotency_key='retry-key') + finally: + conn.close() + + def test_combined_server_head_event_stream_releases_key_slot(self): + with tempfile.TemporaryDirectory() as directory: + home = Path(directory) + manager = api.key_manager(home) + try: + secret = manager.bootstrap_admin( + local_trusted=True, permissions=sorted(apiv1.PERMISSIONS), + concurrency={'max_active': 1}).secret + finally: + manager.conn.close() + server = api.make_api_server(home, '127.0.0.1', 0) + thread = threading.Thread(target=server.serve_forever, daemon=True) + thread.start() + try: + conn = http.client.HTTPConnection('127.0.0.1', server.server_port, timeout=3) + conn.request('HEAD', '/v1/events', headers={'Authorization': f'Bearer {secret}'}) + head = conn.getresponse() + self.assertEqual(head.status, 200) + self.assertEqual(head.read(), b'') + conn.close() + conn = http.client.HTTPConnection('127.0.0.1', server.server_port, timeout=3) + conn.request('GET', '/v1/events', headers={'Authorization': f'Bearer {secret}'}) + response = conn.getresponse() + self.assertEqual(response.status, 200) + response.read() + conn.close() + finally: + server.shutdown() + server.server_close() + thread.join(3) + + def test_scope_applies_before_list_and_event_limits(self): + with tempfile.TemporaryDirectory() as directory: + home = Path(directory) + db.migrate(home) + conn = db.connect(home) + try: + store = jobs.JobStore(conn) + foreign = store.submit('scan', jobs.Scope('foreign', 'profile', 1)) + for number in range(25): + store.emit(foreign.id, 'job.note', data={'number': number}) + own = store.submit('scan', jobs.Scope('mine', 'profile', 1)) + another = store.submit('scan', jobs.Scope('mine', 'profile', 1)) + finally: + conn.close() + control = apiv1.ApiV1(api.WorkbenchService(home), Keys(collections=('mine',))) + listed = control.handle(request('GET', '/v1/jobs', query='limit=1')) + self.assertEqual(listed.status_code, 200, listed.body) + first_ids = [item['id'] for item in listed.json()['items']] + self.assertEqual(len(first_ids), 1) + self.assertEqual(listed.json()['total'], 2) + cursor = listed.json()['next_cursor'] + self.assertTrue(cursor) + second_page = control.handle(request( + 'GET', '/v1/jobs', query=f'limit=1&cursor={cursor}')) + self.assertEqual(second_page.status_code, 200, second_page.body) + second_ids = [item['id'] for item in second_page.json()['items']] + self.assertEqual(set(first_ids + second_ids), {own.id, another.id}) + self.assertIsNone(second_page.json()['next_cursor']) + streamed = control.handle(request('GET', '/v1/events', query='limit=1')) + self.assertEqual(streamed.status_code, 200, streamed.body) + frames = list(streamed.stream) + self.assertIn(own.id, ''.join(frames)) + self.assertNotIn(foreign.id, ''.join(frames)) + first_seq = json.loads(frames[0].split('data: ', 1)[1])['seq'] + event_cursor = frames[0].split('id: ', 1)[1].split('\n', 1)[0] + resumed = control.handle(request( + 'GET', '/v1/events', query=f'limit=1&cursor={event_cursor}')) + self.assertEqual(resumed.status_code, 200, resumed.body) + resumed_frames = list(resumed.stream) + next_seq = json.loads(resumed_frames[0].split('data: ', 1)[1])['seq'] + self.assertEqual(next_seq, first_seq + 1) + + +if __name__ == '__main__': + unittest.main() diff --git a/tests/test_collect_pipeline.py b/tests/test_collect_pipeline.py index 1f1541c..a7eca42 100644 --- a/tests/test_collect_pipeline.py +++ b/tests/test_collect_pipeline.py @@ -375,6 +375,45 @@ async def no_sleep(_seconds): self.assertEqual(entry['bytes'], len(body)) self.assertEqual(entry['pages'], 1) + async def test_broken_attempt_does_not_leave_candidates_from_its_prefix(self): + first = ''.join(f'11.43.0.{index}:80\n' for index in range(1, 21)).encode() + replacement = b'11.44.0.1:80\n11.44.0.2:80\n' + calls = [] + + async def handler(reader, writer): + try: + await reader.readuntil(b'\r\n\r\n') + calls.append(1) + body = first if len(calls) == 1 else replacement + head = (f'HTTP/1.1 200 OK\r\nContent-Length: {len(body)}\r\n' + 'Connection: close\r\n\r\n').encode() + writer.write(head + (first[:len(first) // 2] if len(calls) == 1 else body)) + await writer.drain() + finally: + writer.close() + with contextlib.suppress(ConnectionError): + await writer.wait_closed() + + async def no_sleep(_seconds): + return None + + server = await asyncio.start_server(handler, '127.0.0.1', 0) + original = p.COLLECT_WRITE_BATCH + p.COLLECT_WRITE_BATCH = 5 + try: + url = f'http://127.0.0.1:{server.sockets[0].getsockname()[1]}/l' + report = await p.collect(self.db, [f'http {url}'], [], allow_private_sources=True, + quiet=True, sleep=no_sleep) + finally: + p.COLLECT_WRITE_BATCH = original + server.close() + await server.wait_closed() + self.assertEqual(len(calls), 2) + self.assertIsNone(report['sources'][0]['error']) + self.assertEqual(report['sources'][0]['accepted'], 2) + self.assertEqual(self.candidates(), ['http://11.44.0.1:80', 'http://11.44.0.2:80']) + self.assertEqual(self.db.execute('SELECT count(*) FROM membership_source').fetchone()[0], 2) + async def test_a_source_cut_off_by_the_candidate_limit_reports_its_page(self): body = ''.join(f'11.41.0.{index}:80\n' for index in range(1, 21)).encode() async with self.serve({'/l': body}) as base: diff --git a/tests/test_gateway_bind.py b/tests/test_gateway_bind.py index d0678b8..63bfdb3 100644 --- a/tests/test_gateway_bind.py +++ b/tests/test_gateway_bind.py @@ -6,14 +6,17 @@ apart, keep loopback as the default, and check that the LAN choice is visible. """ import asyncio +import contextlib +import io import json import tempfile import unittest from pathlib import Path +from unittest import mock import httpx -from proxy_workbench import gateway +from proxy_workbench import gateway, proxytool from tests.gateway_support import GatewayCase, write_export #: Values that stand in for the other two identities. They are literal test @@ -67,6 +70,22 @@ async def run(): server.close() asyncio.run(run()) + def test_cli_passes_explicit_lan_opt_in_to_the_gateway(self): + with tempfile.TemporaryDirectory() as temp: + args = proxytool.parser().parse_intermixed_args( + ['gateway', '--data', temp, '--host', '0.0.0.0', '--port', '0', '--lan']) + with mock.patch('proxy_workbench.gateway.start', side_effect=ValueError('stopped')) as start: + with contextlib.redirect_stderr(io.StringIO()): + self.assertEqual(proxytool.run_gateway(args, ()), 2) + self.assertTrue(start.await_args.kwargs['lan']) + self.assertEqual(start.await_args.args[1], '0.0.0.0') + + def test_cli_refuses_lan_flag_for_other_commands(self): + with self.assertRaises(SystemExit) as caught: + with contextlib.redirect_stderr(io.StringIO()): + proxytool.main(['serve', '--lan']) + self.assertEqual(caught.exception.code, 2) + class LanBackgroundTests(unittest.TestCase): """One short, authenticated LAN bind, isolated from every other test. diff --git a/tests/test_pipeline_chain.py b/tests/test_pipeline_chain.py index 02644a8..f7887f3 100644 --- a/tests/test_pipeline_chain.py +++ b/tests/test_pipeline_chain.py @@ -51,6 +51,40 @@ def test_source_byte_budget_truncates_instead_of_buffering(self): self.assertEqual(result.counters.stage_failures.get('stop'), pl.E_LIMIT_BODY) self.assertFalse(result.feed_complete) + def test_source_exactly_at_byte_cap_is_complete(self): + endpoints = support.addresses(2) + body = ('\n'.join(endpoints) + '\n').encode('utf-8') + + async def fetch(): + yield body + yield b'' + + result = run(pl.run_pipeline(support.config( + [pl.SourceSpec('s', fetch)], + pl.Runners(cheap=support.recording_runner(), basic=None, expensive=None), + run_basic=False, budgets=pl.Budgets(max_source_bytes=len(body))))) + self.assertEqual(result.state, 'complete') + self.assertTrue(result.feed_complete) + self.assertEqual(result.counters.sources_truncated, 0) + self.assertEqual(result.metrics.measured, 2) + + def test_overlong_line_is_discarded_without_hiding_the_next_address(self): + endpoint = support.addresses(1)[0] + oversized = b'x' * (len(endpoint) + 1) + + async def fetch(): + yield oversized[:9] + yield oversized[9:] + b'\n' + endpoint.encode('utf-8') + b'\n' + + result = run(pl.run_pipeline(support.config( + [pl.SourceSpec('s', fetch)], + pl.Runners(cheap=support.recording_runner(), basic=None, expensive=None), + run_basic=False, budgets=pl.Budgets(max_line_bytes=len(endpoint))))) + self.assertEqual(result.state, 'complete') + self.assertEqual(result.counters.rejected, 1) + self.assertEqual(result.counters.parsed, 2) + self.assertEqual(result.metrics.measured, 1) + def test_per_source_byte_cap_wins_over_the_global_one(self): endpoints = support.addresses(20) source = support.chunk_source('s', endpoints, chunk=16) diff --git a/tests/test_pipeline_limits.py b/tests/test_pipeline_limits.py index e725042..75b0913 100644 --- a/tests/test_pipeline_limits.py +++ b/tests/test_pipeline_limits.py @@ -394,6 +394,59 @@ def test_a_target_request_budget_stops_only_that_targets_stage(self): codes = [item.reason_code for item in result.results if not item.ok] self.assertEqual(codes, [pl.E_LIMIT_BUDGET] * 9) + def test_expensive_stage_uses_other_targets_after_one_is_exhausted(self): + seen = [] + result = run(pl.run_pipeline(support.config( + [support.chunk_source('s', corpus(12))], + pl.Runners(cheap=None, basic=support.recording_runner(ok=True), + expensive=self.expensive_wide(seen)), + limits=pl.Limits(targets=(pl.TargetPolicy('a', max_inflight=1, max_requests=1), + pl.TargetPolicy('b', max_inflight=2, max_requests=20))), + budgets=pl.Budgets(max_inflight=4, max_open_fds=12), + expensive_policy=pl.EXPENSIVE_ALL_PASSING))) + self.assertEqual(result.state, 'complete') + self.assertEqual(result.counters.passed_endpoints, 12) + self.assertEqual(seen.count('a'), 1) + self.assertEqual(seen.count('b'), 11) + + def test_target_request_cap_counts_each_metered_request(self): + opened = [] + + async def runner(item, *, stage, limit): + self.assertEqual(limit.max_requests, 2) + with limit.budget: + for _ in range(3): + reservation = await limit.budget.reserve(8) + opened.append(limit.target_id) + await limit.budget.settle(reservation, 1) + return pl.StageOutcome(stage, True) + + chain = pl.Pipeline(support.config( + [support.chunk_source('s', corpus(1))], + pl.Runners(cheap=None, basic=support.recording_runner(ok=True), expensive=runner), + limits=pl.Limits(targets=(pl.TargetPolicy('judge', max_requests=2),)), + expensive_policy=pl.EXPENSIVE_ALL_PASSING)) + result = run(chain.run()) + self.assertEqual(opened, ['judge', 'judge']) + self.assertEqual(chain.target_used('judge'), 2) + self.assertEqual(result.state, 'complete', 'a target cap is local to the stage') + self.assertEqual(result.results[0].reason_code, pl.E_LIMIT_BUDGET) + self.assertEqual(result.results[0].stages[-1].requests, 2) + + def test_target_reservation_is_returned_when_stage_does_no_io(self): + async def runner(item, *, stage, limit): + await asyncio.sleep(0.001) + return pl.StageOutcome(stage, True, requests=0) + + chain = pl.Pipeline(support.config( + [support.chunk_source('s', corpus(3))], + pl.Runners(cheap=None, basic=support.recording_runner(ok=True), expensive=runner), + limits=pl.Limits(targets=(pl.TargetPolicy('judge', max_requests=1),)), + expensive_policy=pl.EXPENSIVE_ALL_PASSING)) + result = run(chain.run()) + self.assertEqual(result.counters.passed_endpoints, 3) + self.assertEqual(chain.target_used('judge'), 0) + def test_a_target_concurrency_limit_is_honoured(self): concurrent = {'now': 0, 'peak': 0} diff --git a/tests/test_scan_budgets.py b/tests/test_scan_budgets.py index 98fc7f3..3e84c47 100644 --- a/tests/test_scan_budgets.py +++ b/tests/test_scan_budgets.py @@ -103,6 +103,23 @@ async def refused(proxy, config, limiter): self.assertEqual(state['checked'], 20, 'a refusal spent no request and must not stop the run') self.assertEqual(state['stop_reason'], 'complete') + def test_expensive_target_cap_blocks_its_third_request(self): + self.seed(1) + opened = [] + + async def expensive(proxy, config, limiter, row): + for index in range(3): + async with proxytool._probe_io(1): + opened.append(index) + return row + + state = self.scan(self._ok_probe, expensive_probe=expensive, max_requests=4) + # The scan reserves half the run's request allowance for its expensive + # target. One basic probe and two expensive requests were spent. + self.assertEqual(opened, [0, 1]) + self.assertEqual(state['requests'], 3) + self.assertEqual(state['checked'], 0) + class CountWhatTests(unittest.TestCase): """N endpoints, N unique proxy IPs and N confirmed exit IPs are three numbers.""" diff --git a/tests/test_web_permissions_ui.py b/tests/test_web_permissions_ui.py new file mode 100644 index 0000000..1295230 --- /dev/null +++ b/tests/test_web_permissions_ui.py @@ -0,0 +1,186 @@ +"""Browser form regressions that do not require a live proxy scan.""" + +import json +from pathlib import Path +import shutil +import sys +import unittest + +sys.path.insert(0, str(Path(__file__).resolve().parents[1])) +from tests import web_support as ws +from proxy_workbench import gui + + +PERMISSIONS_DOM = r""" +const esc = value => String(value); +const t = key => key; +const fieldset = (items) => ({querySelectorAll: () => items}); +const grid = { + groups: [], permissions: [], + set innerHTML(html) { + this.groups = []; + this.permissions = []; + for (const block of html.matchAll(/
([\s\S]*?)<\/fieldset>/g)) { + const group = {checked: false, indeterminate: false, dataset: { + permGroup: block[1].match(/data-perm-group="([^"]+)"/)[1] + }}; + const items = []; + for (const match of block[1].matchAll(/]*)>/g)) { + const valueMatch = match[2].match(/\bvalue="([^"]+)"/); + const node = {value: valueMatch ? valueMatch[1] : 'on', + checked: /\bchecked\b/.test(match[2])}; + node.closest = () => fieldset(items); + items.push(node); + this.permissions.push(node); + } + group.closest = () => fieldset(items); + this.groups.push(group); + } + }, + querySelectorAll(selector) { + if (selector === '[data-perm-group]') return this.groups; + if (selector === '[data-permission]') return this.permissions; + return []; + } +}; +const $ = id => id === 'keys-permissions' ? grid : null; +let keysState = {view: null, permissions: [], groups: {}}; +""" + + +class PermissionFormTests(unittest.TestCase): + def test_named_permissions_survive_refresh_and_group_toggles(self): + source = ws.js_slice('function renderPermissions(view)', 'function showKeySecret(issued, message)') + script = PERMISSIONS_DOM + source + r""" +const view = {permissions: ['read.results', 'write.views', 'export.secret'], + read_permissions: ['read.results'], write_permissions: ['write.views'], admin_permissions: []}; +renderPermissions(view); +setPermissions(['read.results', 'export.secret']); +const first = selectedPermissions(); +renderPermissions(view); +const refreshed = selectedPermissions(); +const write = grid.groups.find(group => group.dataset.permGroup === 'write.*'); +write.checked = true; +write.onchange(); +const afterGroup = selectedPermissions(); +const read = grid.permissions.find(item => item.value === 'read.results'); +read.checked = false; +read.onchange(); +process.stdout.write(JSON.stringify({first, refreshed, afterGroup, + afterSingle: selectedPermissions(), + readGroup: grid.groups.find(group => group.dataset.permGroup === 'read.*').checked})); +""" + value = json.loads(ws.node_ok(script)) + self.assertEqual(value['first'], ['read.results', 'export.secret']) + self.assertEqual(value['refreshed'], value['first']) + self.assertEqual(value['afterGroup'], ['read.results', 'write.views', 'export.secret']) + self.assertEqual(value['afterSingle'], ['write.views', 'export.secret']) + self.assertFalse(value['readGroup']) + + def test_a_source_action_restores_the_button(self): + source = ws.js_slice('async function sourceAction(path, button, message)', "if ($('prune-sources'))") + script = r""" +const t = key => key; +const button = {disabled: false, textContent: 'Update', dataset: {}}; +const sources = {value: ''}; +const $ = id => id === 'sources' ? sources : null; +const getSettings = () => ({}); +const updateSourceCount = () => {}; +const updateCodeEditors = () => {}; +const toast = () => {}; +let settings = {}; +let fail = false; +const api = async () => {if (fail) throw Error('network'); return {settings: {sources: ['https://example.org/list']}};}; +""" + source + r""" +(async () => { + await sourceAction('/api/sources/update', button, () => 'done'); + const success = [button.textContent, button.disabled, sources.value]; + fail = true; + await sourceAction('/api/sources/update', button, () => 'done'); + process.stdout.write(JSON.stringify({success, failure: [button.textContent, button.disabled]})); +})().catch(error => {console.error(error); process.exitCode = 1;}); +""" + value = json.loads(ws.node_ok(script)) + self.assertEqual(value['success'], ['Update', False, 'https://example.org/list']) + self.assertEqual(value['failure'], ['Update', False]) + + def test_editing_import_text_invalidates_inflight_preview(self): + source = ws.js_slice('let importState =', '// Pools and schedules') + script = r""" +const t = key => key; +const toast = () => {}; +const fields = { + 'import-preview': {disabled: false}, + 'import-commit': {disabled: true, dataset: {}}, + 'import-text': {value: '198.51.100.7:8080'}, + 'import-collection': {value: 'first'}, + 'import-format': {value: ''}, + 'import-mode': {value: 'merge'}, + 'import-mapping-grid': {closest: () => ({hidden: true})}, + 'import-preview-body': {innerHTML: ''}, + 'imp-count': {textContent: '0'}, + 'import-note': {textContent: ''}, +}; +const $ = id => fields[id] || null; +let finish; +const api = async () => new Promise(resolve => {finish = resolve;}); +""" + source + r""" +(async () => { + setupImportListeners(); + const pending = runImportPreview(); + fields['import-text'].value = '198.51.100.8:8080'; + fields['import-text'].oninput(); + finish({counts: {added: 1}, batch_id: 'old-preview'}); + await pending; + process.stdout.write(JSON.stringify({plan: importState.plan, busy: importState.busy, + applyDisabled: fields['import-commit'].disabled, previewDisabled: fields['import-preview'].disabled, + count: fields['imp-count'].textContent})); +})().catch(error => {console.error(error); process.exitCode = 1;}); +""" + value = json.loads(ws.node_ok(script)) + self.assertIsNone(value['plan']) + self.assertFalse(value['busy']) + self.assertTrue(value['applyDisabled']) + self.assertFalse(value['previewDisabled']) + self.assertEqual(value['count'], '0') + + def test_import_mapping_offers_every_column_and_accepts_index_zero(self): + home = ws.build_data([], publish=False) + app = gui.App(home) + try: + plan = app.import_preview({ + 'collection': 'public-base', 'format': 'csv', 'name': 'ambiguous.csv', + 'text': 'host,ip,port\n11.0.0.1,11.0.0.2,8080\n' + }) + finally: + app.close() + shutil.rmtree(home) + self.assertTrue(plan['needs_mapping']) + self.assertEqual(plan['columns'], ['host', 'ip', 'port']) + self.assertEqual(plan['mapping_suggestion']['columns'], plan['columns']) + + render = ws.js_slice('function renderImportMapping(suggestion, problem)', 'function renderImportPreview(plan)') + mapping = ws.js_slice('function mappingFromForm()', 'async function runImportPreview()') + script = r""" +const esc = value => String(value); +const t = key => key; +const box = {hidden: true}; +const grid = {innerHTML: '', closest: () => box, querySelectorAll: () => [ + {dataset: {mappingRole: 'host'}, value: '0'}, + {dataset: {mappingRole: 'port'}, value: '2'}, + {dataset: {mappingRole: 'scheme'}, value: ''}, + {dataset: {mappingRole: 'country'}, value: ''}, +]}; +const $ = id => id === 'import-mapping' ? box : id === 'import-mapping-grid' ? grid : null; +""" + render + mapping + "\n" + "const plan = " + json.dumps(plan) + ";\n" + r""" +renderImportMapping(plan.mapping_suggestion, plan.mapping_problem); +process.stdout.write(JSON.stringify({html: grid.innerHTML, hidden: box.hidden, mapping: mappingFromForm()})); +""" + result = json.loads(ws.node_ok(script)) + self.assertFalse(result['hidden']) + self.assertIn('', result['html']) + self.assertEqual(result['mapping'], {'host': 0, 'port': 2, 'scheme': None, 'country': None}) + + +if __name__ == '__main__': + unittest.main() diff --git a/tests/test_web_tabs.py b/tests/test_web_tabs.py new file mode 100644 index 0000000..735ab11 --- /dev/null +++ b/tests/test_web_tabs.py @@ -0,0 +1,81 @@ +"""Navigation smoke for every shipped page and its first data load.""" + +import json +from pathlib import Path +import re +import sys +import unittest + +sys.path.insert(0, str(Path(__file__).resolve().parents[1])) +from tests import web_support as ws + + +class TabNavigationTests(unittest.TestCase): + def test_each_nav_opens_its_page_and_starts_its_data_load(self): + html = ws.INDEX_HTML.read_text(encoding='utf-8') + navs = re.findall(r']*\bdata-tab="([^"]+)"', html) + pages = re.findall(r']*\bclass="page(?: [^"]*)?"[^>]*\bid="page-([^"]+)"', html) + expected = ['scan', 'results', 'gateway', 'mobile', 'sources', 'pools', 'keys', 'help'] + self.assertEqual(navs, expected) + self.assertEqual(pages, expected) + + source = ws.js_slice('function showTab(name)', '// Final URLs only') + script = r""" +let currentTab = 'scan'; +const t = key => key; +const calls = []; +const pageNames = """ + json.dumps(expected) + r"""; +const mkClassList = () => { + const names = new Set(); + return {toggle(name, value) {if (value) names.add(name); else names.delete(name);}, + add(name) {names.add(name);}, remove(name) {names.delete(name);}, contains(name) {return names.has(name);}}; +}; +const pages = pageNames.map(name => ({id: 'page-' + name, classList: mkClassList(), offsetWidth: 1})); +const navs = pageNames.map(name => ({dataset: {tab: name}, classList: mkClassList()})); +const label = {textContent: ''}; +const document = { + querySelectorAll(selector) { + if (selector === '.page') return pages; + if (selector === '.nav' || selector === '[data-tab]') return navs; + if (selector === '[data-go]') return []; + return []; + }, + getElementById(id) {return id === 'page-label' ? label : pages.find(page => page.id === id);} +}; +const $ = id => document.getElementById(id); +const window = {location: {hash: ''}, scrollTo() {}}; +const history = {replaceState(_state, _title, hash) {window.location.hash = hash;}}; +function loadResults() {calls.push('results');} +function reloadCatalog() {calls.push('catalog');} +function loadPools() {calls.push('pools');} +function loadSchedules() {calls.push('schedules');} +function loadJobs() {calls.push('jobs');} +function loadGatewayOptions() {calls.push('gatewayOptions');} +function loadKeys() {calls.push('keys');} +function loadDiagnostics() {calls.push('diagnostics');} +function loadDesktop() {calls.push('desktop');} +""" + source + r""" +const results = []; +for (const nav of navs) { + calls.length = 0; + nav.onclick(); + results.push({tab: nav.dataset.tab, active: pages.filter(page => page.classList.contains('active')).map(page => page.id), + label: label.textContent, hash: window.location.hash, loaders: [...calls]}); +} +process.stdout.write(JSON.stringify(results)); +""" + results = json.loads(ws.node_ok(script)) + loaders = {'scan': [], 'results': ['results'], 'gateway': [], 'mobile': [], + 'sources': ['catalog'], + 'pools': ['pools', 'schedules', 'jobs', 'gatewayOptions'], + 'keys': ['keys'], 'help': ['diagnostics', 'desktop']} + for row in results: + name = row['tab'] + self.assertEqual(row['active'], ['page-' + name]) + self.assertEqual(row['label'], 'nav.' + name) + self.assertEqual(row['hash'], '#' + name) + self.assertEqual(row['loaders'], loaders[name]) + + +if __name__ == '__main__': + unittest.main()