From 7dd4c1f42aaa75b237d25cf14c9296e6ef1fa4ec Mon Sep 17 00:00:00 2001 From: "marco.mengelkoch" Date: Sat, 3 Oct 2026 17:25:13 +0200 Subject: [PATCH 1/2] update mq-bridge version --- frameworks/mq-bridge-py/Dockerfile | 8 +- frameworks/mq-bridge-py/meta.json | 3 +- frameworks/mq-bridge-py/server.py | 290 ++++++----------------------- site/data/tls/mq-bridge-py.json | 5 + 4 files changed, 67 insertions(+), 239 deletions(-) create mode 100644 site/data/tls/mq-bridge-py.json diff --git a/frameworks/mq-bridge-py/Dockerfile b/frameworks/mq-bridge-py/Dockerfile index 8250b9927..f9c7330d6 100644 --- a/frameworks/mq-bridge-py/Dockerfile +++ b/frameworks/mq-bridge-py/Dockerfile @@ -1,20 +1,20 @@ # HttpArena image for the mq-bridge-py (Python) entry. # -# Installs the prebuilt mq_bridge_py wheel from PyPI and runs server.py on port +# Installs the prebuilt mq-bridge wheel from PyPI and runs server.py on port # 8080. The wheel is abi3 (cp38+) manylinux2014, so it works on this base image # without compiling from source. Pin MQB_VERSION to a released version. FROM python:3.12-slim -ARG MQB_VERSION=0.4.8 +ARG MQB_VERSION=0.4.17 RUN groupadd --system appuser \ && useradd --system --gid appuser --create-home --home-dir /home/appuser --shell /usr/sbin/nologin appuser \ && mkdir -p /app \ && chown -R appuser:appuser /app -# psycopg[binary] powers the DB-backed profiles (/async-db, /fortunes, /crud); absent +# psycopg[binary] powers the DB-backed profiles (/async-db, /fortunes); absent # DATABASE_URL they are unused. Jinja2 renders /fortunes — the fortunes profile # requires a real template engine, not string concatenation in the handler. -RUN pip install --no-cache-dir "mq-bridge-py==${MQB_VERSION}" "psycopg[binary]>=3.1" "psycopg_pool>=3.2" "redis>=5.2" "jinja2>=3.1" +RUN pip install --no-cache-dir "mq-bridge==${MQB_VERSION}" "psycopg[binary]>=3.1" "psycopg_pool>=3.2" "jinja2>=3.1" WORKDIR /app COPY --chown=appuser:appuser server.py /app/server.py # The fortunes template is a separate artifact, as the profile requires. diff --git a/frameworks/mq-bridge-py/meta.json b/frameworks/mq-bridge-py/meta.json index 0173a5539..a872b4159 100644 --- a/frameworks/mq-bridge-py/meta.json +++ b/frameworks/mq-bridge-py/meta.json @@ -10,7 +10,7 @@ "response": true }, "engine": "mq-bridge", - "description": "mq-bridge Python bindings (mq_bridge_py). A single http->response route dispatches on request metadata while HTTP framing stays in Rust; cleartext HTTP/1.1 + h2c run on 8080/8082, optional TLS listeners serve HTTP/2 on 8443 and HTTP/1.1 on 8081, and the inline-response fast path keeps responses off the GIL so the Python handler runs only the per-request dispatch. Postgres via psycopg for async-db, crud (Redis cache-aside) and fortunes, which renders per request with the Jinja2 template engine; static assets are read from /data/static per request, serving the on-disk .gz variant when the client accepts it; /delay waits with a blocking sleep, so async concurrency is the worker-process count.", + "description": "mq-bridge Python bindings (mq_bridge_py). A single http->response route dispatches on request metadata while HTTP framing stays in Rust; cleartext HTTP/1.1 + h2c run on 8080/8082, optional TLS listeners serve HTTP/2 on 8443 and HTTP/1.1 on 8081, and the inline-response fast path keeps responses off the GIL so the Python handler runs only the per-request dispatch. Postgres via psycopg for async-db and fortunes, which renders per request with the Jinja2 template engine; static assets are read from /data/static per request, serving the on-disk .gz variant when the client accepts it; /delay waits with a blocking sleep, so async concurrency is the worker-process count.", "repo": "https://github.com/marcomq/mq-bridge", "enabled": true, "tests": [ @@ -25,6 +25,7 @@ "fortunes", "json-tls", "baseline-h2", + "static-tls", "static-h2", "baseline-h2c", "json-h2c", diff --git a/frameworks/mq-bridge-py/server.py b/frameworks/mq-bridge-py/server.py index 9b0476498..081dc7a90 100644 --- a/frameworks/mq-bridge-py/server.py +++ b/frameworks/mq-bridge-py/server.py @@ -21,24 +21,19 @@ ``GET /delay/{ms}`` ``{ms}`` after waiting async ``GET /async-db?min=&max=&limit=`` ``items`` rows async-db ``GET /fortunes`` rendered HTML table fortunes -``GET /static/{file}`` file from /data/static static-h2 -``GET /crud/items?...`` paginated list crud -``GET /crud/items/{id}`` cached item crud -``POST /crud/items`` + JSON 201 + upserted item crud -``PUT /crud/items/{id}`` + JSON updated item crud +``GET /static/{file}`` file from /data/static static-tls, static-h2 =============================== ========================== ===================== Listeners: 8080 HTTP/1.1 + h2c (auto), 8082 h2c-only, 8443 h2-over-TLS, 8081 HTTP/1.1-over-TLS. The TLS ports bind only when certs are mounted. Harness inputs: ``DATASET_PATH``, ``STATIC_DIR``, ``DATABASE_URL``, -``DATABASE_MAX_CONN``, ``REDIS_URL``. A missing database is non-fatal — the +``DATABASE_MAX_CONN``. A missing database is non-fatal — the DB-backed endpoints degrade rather than blocking the cleartext profiles. ``json-comp`` is served by mq-bridge's own response compression -(``compression_enabled``): bodies over the threshold are gzip-encoded when the -client advertises ``Accept-Encoding: gzip``, identity otherwise — so one -``/json`` handler serves both ``json`` and ``json-comp``. +(``compression_enabled``): bodies over the threshold are encoded per request +from the client's ``Accept-Encoding``, identity otherwise. """ from __future__ import annotations @@ -46,7 +41,6 @@ import json as _json import os import signal -import tempfile import threading import time from pathlib import Path @@ -71,14 +65,13 @@ OCTET_META = {"content-type": "application/octet-stream", "Server": SERVER} -def _status_meta(code: int, base: dict = JSON_META) -> dict: - return dict(base, http_status_code=str(code)) +def _status(code: int, reason: bytes) -> tuple[bytes, dict]: + return reason, dict(TEXT_META, http_status_code=str(code)) -NOT_FOUND = (b"Not Found", _status_meta(404, TEXT_META)) -BAD_REQUEST = (b"Bad Request", _status_meta(400, TEXT_META)) -SERVER_ERROR = (b"Internal Server Error", _status_meta(500, TEXT_META)) -UNAVAILABLE = (b"Service Unavailable", _status_meta(503, TEXT_META)) +NOT_FOUND = _status(404, b"Not Found") +SERVER_ERROR = _status(500, b"Internal Server Error") +UNAVAILABLE = _status(503, b"Service Unavailable") # ---------- route configuration ---------- @@ -87,74 +80,49 @@ def _tls_available() -> bool: return Path(TLS_CERT).is_file() and Path(TLS_KEY).is_file() -def _http_route(name: str, listen: str, http_workers: int, extra: str = "") -> str: - return f""" - {name}: - concurrency: 1 - batch_size: 1024 - input: - http: - url: "{listen}" - workers: {http_workers} - concurrency_limit: 65536 - internal_buffer_size: 16384 - inline_response_fast_path: true - compression_enabled: true - compression_threshold_bytes: 256 -{extra} - output: - response: {{}} -""" - +def _http_route(listen: str, http_workers: int, **extra) -> dict: + http = { + "url": listen, + "workers": http_workers, + "concurrency_limit": 65536, + "internal_buffer_size": 16384, + "inline_response_fast_path": True, + "compression_enabled": True, + "compression_threshold_bytes": 256, + **extra, + } + return { + "concurrency": 1, + "batch_size": 1024, + "input": {"http": http}, + "output": {"response": {}}, + } -def _config(http_workers: int) -> tuple[str, list[str]]: - # `http_workers` is the number of accept loops (each its own SO_REUSEPORT - # listener) inside this process. When we fan out across processes we keep - # this small (the single Python worker is the per-process bottleneck); in - # single-process mode we use all cores, matching the previous default. - names = ["httparena"] - routes = [_http_route(names[0], LISTEN, http_workers)] - - # HTTP/2 cleartext (prior-knowledge) on 8082 (baseline-h2c / json-h2c), using - # the same handlers as the plaintext listener. `http2_only` makes the port - # refuse HTTP/1.1, satisfying the h2c-only anti-cheat (a dual-serving port is - # rejected). Cleartext, so no certs are needed. - names.append("httparena-h2c") - routes.append( - _http_route( - names[-1], - H2C_LISTEN, - http_workers, - " server_protocol: http2_only", - ) - ) - # TLS listeners, only when the harness has mounted certs — a local +def _routes(http_workers: int) -> dict[str, dict]: + """Route bodies by name. `http_workers` is the number of accept loops (each + its own SO_REUSEPORT listener) inside this process: small when fanned out + across processes, all cores in single-process mode.""" + routes = { + "httparena": _http_route(LISTEN, http_workers), + # h2c prior knowledge only (baseline-h2c / json-h2c): the port has to + # refuse HTTP/1.1, a dual-serving one fails the h2c-only anti-cheat. + "httparena-h2c": _http_route( + H2C_LISTEN, http_workers, server_protocol="http2_only" + ), + } + # TLS listeners only when the harness has mounted certs, so a local # plaintext-only run still works. if _tls_available(): - tls_block = ( - f' tls:\n' - f' required: true\n' - f' cert_file: "{TLS_CERT}"\n' - f' key_file: "{TLS_KEY}"' - ) - # HTTP/2 over TLS on 8443 (baseline-h2 / static-h2): ALPN advertises `h2`. - names.append("httparena-tls") - routes.append(_http_route(names[-1], TLS_LISTEN, http_workers, tls_block)) - # JSON over HTTP/1.1 + TLS on 8081 (json-tls): the same `/json` handler, - # but the port advertises ALPN `http/1.1` only so the wrk load generator - # negotiates HTTP/1.1 rather than upgrading to h2. - names.append("httparena-json-tls") - routes.append( - _http_route( - names[-1], - H1TLS_LISTEN, - http_workers, - tls_block + "\n server_protocol: http1_only", - ) + tls = {"required": True, "cert_file": TLS_CERT, "key_file": TLS_KEY} + # 8443: ALPN `h2` (baseline-h2 / static-h2). + routes["httparena-tls"] = _http_route(TLS_LISTEN, http_workers, tls=tls) + # 8081: ALPN `http/1.1` only (json-tls / static-tls), so wrk does not + # negotiate h2. + routes["httparena-h1-tls"] = _http_route( + H1TLS_LISTEN, http_workers, tls=tls, server_protocol="http1_only" ) - - return "routes:\n" + "\n".join(routes), names + return routes # ---------- static ---------- @@ -273,7 +241,6 @@ def _build_json(qs: dict[str, list[str]], count: int) -> tuple[bytes, dict]: # ---------- database ---------- _POOL = None -_REDIS = None ITEM_COLUMNS = ( "id, name, category, price, quantity, active, tags, rating_score, rating_count" @@ -297,11 +264,9 @@ def _init_pool(): except ImportError: return None budget = int(os.environ.get("DATABASE_MAX_CONN", "256")) - # One pool per forked worker, so the budget has to be divided by the worker - # count rather than handed to each. Postgres runs with max_connections=256 - # and reserves a few of those for the superuser; the crud profile hands the - # container a 62-CPU cpuset, so a per-worker max of the full budget asked - # for 62 x 256. + # One pool per forked worker, so the budget is divided by the worker count + # rather than handed to each. Postgres runs with max_connections=256 and + # reserves a few of those for the superuser. max_conn = max(1, (budget - 8) // max(1, _worker_count())) try: return ConnectionPool(url, min_size=1, max_size=max_conn, open=True) @@ -310,18 +275,6 @@ def _init_pool(): return None -def _init_redis(): - url = os.environ.get("REDIS_URL", "") - if not url: - return None - try: - import redis - - return redis.Redis.from_url(url, decode_responses=True) - except Exception: # noqa: BLE001 - non-fatal, crud degrades to DB-only - return None - - def _fetch(sql: str, params: tuple = (), one: bool = False): """Run one query on the pool. Returns the rows (or the single row, which is None when nothing matched); raises `_DbError` when the pool is missing or @@ -369,122 +322,6 @@ def _async_db(qs: dict[str, list[str]]) -> tuple[bytes, dict]: return _dumps({"count": len(items), "items": items}), JSON_META -# ---------- crud ---------- -# -# Cache-aside on Redis where the harness provides it — crud is the one profile -# that does, and the cache is shared across the forked workers as a per-process -# dict would not be. The profile reads and writes the same ids, so a long TTL -# would answer from a copy the writes have already moved past. - -CRUD_ITEM = "/crud/items/" -CRUD_TTL_MS = 200 - - -def _crud_invalidate(item_id: int) -> None: - if _REDIS is not None: - try: - _REDIS.delete("crud:%d" % item_id) - except Exception: # noqa: BLE001 - pass - - -def _crud_list(qs: dict[str, list[str]]) -> tuple[bytes, dict]: - """One query, no `SELECT COUNT(*)`: `total` reports the rows in this - response, which is what the profile asks for — the full-filter count was - dropped from the spec because it dominated Postgres CPU under writes.""" - page = max(1, _query_int(qs, "page", 1)) - limit = max(1, min(_query_int(qs, "limit", 10), 50)) - rows = _fetch( - f"SELECT {ITEM_COLUMNS} FROM items WHERE category = %s ORDER BY id " - "LIMIT %s OFFSET %s", - ((qs.get("category") or ["electronics"])[0], limit, (page - 1) * limit), - ) - items = [_item_row(r) for r in rows] - return ( - _dumps({"items": items, "total": len(items), "page": page, "limit": limit}), - JSON_META, - ) - - -def _crud_read(item_id: int) -> tuple[bytes, dict]: - key = "crud:%d" % item_id - if _REDIS is not None: - try: - hit = _REDIS.get(key) - except Exception: # noqa: BLE001 - hit = None - if hit: - return hit.encode(), dict(JSON_META, **{"X-Cache": "HIT"}) - row = _fetch(f"SELECT {ITEM_COLUMNS} FROM items WHERE id = %s", (item_id,), one=True) - if row is None: - return NOT_FOUND - body = _dumps(_item_row(row)) - if _REDIS is not None: - try: - _REDIS.set(key, body, px=CRUD_TTL_MS) - except Exception: # noqa: BLE001 - pass - return body, dict(JSON_META, **{"X-Cache": "MISS"}) - - -def _crud_create(payload: bytes) -> tuple[bytes, dict]: - """`active` / `tags` / `rating_*` are NOT NULL with no default and the create - body carries none of them, so a new row seeds them; a conflict updates only - the four fields the body actually sends.""" - try: - body = _json.loads(payload) - except ValueError: - return BAD_REQUEST - row = _fetch( - f"INSERT INTO items ({ITEM_COLUMNS}) " - "VALUES (%s, %s, %s, %s, %s, true, '[]'::jsonb, 0, 0) " - "ON CONFLICT (id) DO UPDATE SET name = EXCLUDED.name, " - "category = EXCLUDED.category, price = EXCLUDED.price, " - f"quantity = EXCLUDED.quantity RETURNING {ITEM_COLUMNS}", - ( - body.get("id"), - body.get("name", "New Product"), - body.get("category", "test"), - body.get("price", 0), - body.get("quantity", 0), - ), - one=True, - ) - if row is None: - return SERVER_ERROR - # The upsert may have replaced a row someone already read. - _crud_invalidate(row[0]) - return _dumps(_item_row(row)), _status_meta(201) - - -def _crud_update(item_id: int, payload: bytes) -> tuple[bytes, dict]: - """A partial PUT leaves the rest of the row alone: each bind is COALESCEd - against the current value, so an absent field is not an overwrite with a - default. The casts are what let psycopg send an untyped NULL.""" - try: - body = _json.loads(payload) - except ValueError: - return BAD_REQUEST - row = _fetch( - "UPDATE items SET name = COALESCE(%s::text, name), " - "category = COALESCE(%s::text, category), price = COALESCE(%s::int, price), " - "quantity = COALESCE(%s::int, quantity) " - f"WHERE id = %s RETURNING {ITEM_COLUMNS}", - ( - body.get("name"), - body.get("category"), - body.get("price"), - body.get("quantity"), - item_id, - ), - one=True, - ) - if row is None: - return NOT_FOUND - _crud_invalidate(item_id) - return _dumps(_item_row(row)), JSON_META - - # ---------- fortunes ---------- RUNTIME_FORTUNE = (0, "Additional fortune added at request time.") @@ -566,8 +403,6 @@ def _baseline(message: Message, qs: dict[str, list[str]]) -> tuple[bytes, dict]: ("POST", "/echo"): lambda m, qs: (m.payload, OCTET_META), ("GET", "/async-db"): lambda m, qs: _async_db(qs), ("GET", "/fortunes"): lambda m, qs: _fortunes(), - ("GET", "/crud/items"): lambda m, qs: _crud_list(qs), - ("POST", "/crud/items"): lambda m, qs: _crud_create(bytes(m.payload)), } # Prefix matches, checked only when no exact route matched. The prefixes are @@ -576,8 +411,6 @@ def _baseline(message: Message, qs: dict[str, list[str]]) -> tuple[bytes, dict]: ("GET", "/json/", lambda m, qs, t: _build_json(qs, _int_or(t, 0))), ("GET", "/static/", lambda m, qs, t: _serve_static(t, _accepts_gzip(m))), ("GET", "/delay/", lambda m, qs, t: _delay(t)), - ("GET", CRUD_ITEM, lambda m, qs, t: _crud_by_id(t, _crud_read)), - ("PUT", CRUD_ITEM, lambda m, qs, t: _crud_by_id(t, _crud_update, bytes(m.payload))), ) @@ -588,14 +421,6 @@ def _int_or(text: str, default: int) -> int: return default -def _crud_by_id(tail: str, fn, *args) -> tuple[bytes, dict]: - try: - item_id = int(tail) - except ValueError: - return NOT_FOUND - return fn(item_id, *args) - - _DB_ERRORS = {503: UNAVAILABLE, 500: SERVER_ERROR} @@ -619,7 +444,7 @@ def handle(message: Message) -> Message: except _DbError as exc: body, reply_meta = _DB_ERRORS[exc.code] - return message.__class__(body, reply_meta) + return Message(body, reply_meta) # ---------- process model ---------- @@ -634,15 +459,12 @@ def _run_secondary_listener(route: Route) -> None: def _run_worker(http_workers: int) -> None: # Per-process setup: the Postgres pool (background threads) and the Rust # runtime must be created AFTER any fork, never inherited across it. - global _POOL, _REDIS + global _POOL _POOL = _init_pool() - _REDIS = _init_redis() - config, names = _config(http_workers) - with tempfile.NamedTemporaryFile("w", suffix=".yaml", delete=False) as f: - f.write(config) - config_path = f.name - - routes = [Route.from_file(config_path, name).with_handler(handle) for name in names] + routes = [ + Route.from_config(body, name).with_handler(handle) + for name, body in _routes(http_workers).items() + ] # Keep every port fail-fast: if any secondary listener exits, signal this # worker so the parent supervisor restarts a clean set instead of leaving a # partially serving process behind. diff --git a/site/data/tls/mq-bridge-py.json b/site/data/tls/mq-bridge-py.json new file mode 100644 index 000000000..78db97a80 --- /dev/null +++ b/site/data/tls/mq-bridge-py.json @@ -0,0 +1,5 @@ +{ + "framework": "mq-bridge-py", + "tls": "pass", + "check": "none" +} From 8e10238129e20bb3152149959b4fc2666bd09375 Mon Sep 17 00:00:00 2001 From: "marco.mengelkoch" Date: Sat, 3 Oct 2026 17:30:35 +0200 Subject: [PATCH 2/2] rename --- frameworks/mq-bridge-py/Dockerfile | 3 ++- frameworks/mq-bridge-py/meta.json | 2 +- frameworks/mq-bridge-py/server.py | 3 ++- 3 files changed, 5 insertions(+), 3 deletions(-) diff --git a/frameworks/mq-bridge-py/Dockerfile b/frameworks/mq-bridge-py/Dockerfile index f9c7330d6..e8168738f 100644 --- a/frameworks/mq-bridge-py/Dockerfile +++ b/frameworks/mq-bridge-py/Dockerfile @@ -1,6 +1,7 @@ # HttpArena image for the mq-bridge-py (Python) entry. # -# Installs the prebuilt mq-bridge wheel from PyPI and runs server.py on port +# "mq-bridge-py" is only this entry's name: the PyPI package is `mq-bridge` +# (imported as `mq_bridge`). Installs its prebuilt wheel and runs server.py on port # 8080. The wheel is abi3 (cp38+) manylinux2014, so it works on this base image # without compiling from source. Pin MQB_VERSION to a released version. FROM python:3.12-slim diff --git a/frameworks/mq-bridge-py/meta.json b/frameworks/mq-bridge-py/meta.json index a872b4159..95f6b79e6 100644 --- a/frameworks/mq-bridge-py/meta.json +++ b/frameworks/mq-bridge-py/meta.json @@ -10,7 +10,7 @@ "response": true }, "engine": "mq-bridge", - "description": "mq-bridge Python bindings (mq_bridge_py). A single http->response route dispatches on request metadata while HTTP framing stays in Rust; cleartext HTTP/1.1 + h2c run on 8080/8082, optional TLS listeners serve HTTP/2 on 8443 and HTTP/1.1 on 8081, and the inline-response fast path keeps responses off the GIL so the Python handler runs only the per-request dispatch. Postgres via psycopg for async-db and fortunes, which renders per request with the Jinja2 template engine; static assets are read from /data/static per request, serving the on-disk .gz variant when the client accepts it; /delay waits with a blocking sleep, so async concurrency is the worker-process count.", + "description": "mq-bridge Python bindings (PyPI package mq-bridge, imported as mq_bridge). A single http->response route dispatches on request metadata while HTTP framing stays in Rust; cleartext HTTP/1.1 + h2c run on 8080/8082, optional TLS listeners serve HTTP/2 on 8443 and HTTP/1.1 on 8081, and the inline-response fast path keeps responses off the GIL so the Python handler runs only the per-request dispatch. Postgres via psycopg for async-db and fortunes, which renders per request with the Jinja2 template engine; static assets are read from /data/static per request, serving the on-disk .gz variant when the client accepts it; /delay waits with a blocking sleep, so async concurrency is the worker-process count.", "repo": "https://github.com/marcomq/mq-bridge", "enabled": true, "tests": [ diff --git a/frameworks/mq-bridge-py/server.py b/frameworks/mq-bridge-py/server.py index 081dc7a90..e18d02e5c 100644 --- a/frameworks/mq-bridge-py/server.py +++ b/frameworks/mq-bridge-py/server.py @@ -1,4 +1,5 @@ -"""HttpArena entry for mq-bridge-py (Python). +"""HttpArena entry "mq-bridge-py": the mq-bridge Python bindings (PyPI package +``mq-bridge``, imported as ``mq_bridge``). One catch-all ``http -> response`` route per listener, dispatching on the request's ``http_method`` / ``http_path`` / ``http_query`` metadata. mq-bridge