diff --git a/.github/workflows/pr85523-materialize.yml b/.github/workflows/pr85523-materialize.yml new file mode 100644 index 000000000000..f57e997a4d97 --- /dev/null +++ b/.github/workflows/pr85523-materialize.yml @@ -0,0 +1,52 @@ +name: PR85523 materialize closure + +on: + pull_request: + types: [opened, reopened, synchronize] + paths: + - '.hermes-patches/task10/**' + - '.github/workflows/pr85523-materialize.yml' + +permissions: + contents: read + +jobs: + materialize: + runs-on: ubuntu-latest + steps: + - name: Checkout PR head + uses: actions/checkout@v4 + with: + fetch-depth: 0 + + - name: Materialize exact current-main candidate + shell: bash + env: + BASE_SHA: ${{ github.event.pull_request.base.sha }} + run: | + set -euo pipefail + cat .hermes-patches/task10/part*.patch > /tmp/pr85523.patch + echo '6db360799eec591e6b4fc0aa6350d7de5208ca62485ee0d7f5b3700ac3459b1d /tmp/pr85523.patch' | sha256sum -c - + git fetch origin "$BASE_SHA" + git checkout --detach "$BASE_SHA" + git apply --check /tmp/pr85523.patch + git apply /tmp/pr85523.patch + git diff --check + python -m py_compile gateway/platforms/webhook.py tests/gateway/test_webhook_adapter.py tests/gateway/test_webhook_task10_closure.py + python -m pytest -q tests/gateway/test_webhook_adapter.py tests/gateway/test_webhook_task10_closure.py tests/gateway/test_webhook_http_contract.py tests/gateway/test_webhook_intake_hardening.py + git hash-object gateway/platforms/webhook.py | grep '^2cbe1d560757e0bbf86d38c3e9572acfc2b4eb9c$' + git hash-object tests/gateway/test_webhook_adapter.py | grep '^21be925f865da8e992e229e446a504ffea884a47$' + git hash-object tests/gateway/test_webhook_task10_closure.py | grep '^92fb404c70bf536a8abceb81a7abf569d1593532$' + mkdir -p /tmp/pr85523-candidate/gateway/platforms /tmp/pr85523-candidate/tests/gateway + cp gateway/platforms/webhook.py /tmp/pr85523-candidate/gateway/platforms/webhook.py + cp tests/gateway/test_webhook_adapter.py /tmp/pr85523-candidate/tests/gateway/test_webhook_adapter.py + cp tests/gateway/test_webhook_task10_closure.py /tmp/pr85523-candidate/tests/gateway/test_webhook_task10_closure.py + + - name: Upload exact candidate files + uses: actions/upload-artifact@v4 + with: + name: pr85523-current-main-candidate + path: /tmp/pr85523-candidate + if-no-files-found: error + +# retrigger: 2026-08-19T21:09Z diff --git a/.hermes-patches/task10/part00.patch b/.hermes-patches/task10/part00.patch new file mode 100644 index 000000000000..7b2c80fc1aee --- /dev/null +++ b/.hermes-patches/task10/part00.patch @@ -0,0 +1,173 @@ +diff --git a/gateway/platforms/webhook.py b/gateway/platforms/webhook.py +index bc3ebe0..2cbe1d5 100644 +--- a/gateway/platforms/webhook.py ++++ b/gateway/platforms/webhook.py +@@ -40,8 +40,11 @@ import logging + import re + import subprocess + import sys ++import threading + import time ++import uuid + from collections import deque ++from enum import Enum + from typing import Any, Deque, Dict, List, Optional + + try: +@@ -131,6 +134,14 @@ DEFAULT_PORT = 8644 + _INSECURE_NO_AUTH = "INSECURE_NO_AUTH" + _DYNAMIC_ROUTES_FILENAME = "webhook_subscriptions.json" + _RATE_WINDOW_SECONDS = 60.0 ++_IDEMPOTENCY_DEFAULT_MAX_ENTRIES = 4096 ++_IDEMPOTENCY_MAX_ENTRIES_LIMIT = 1_000_000 ++_RAW_PAYLOAD_DEFAULT_CAP_BYTES = 4_000 ++_RAW_PAYLOAD_MIN_CAP_BYTES = 64 ++_RAW_PAYLOAD_MAX_CAP_BYTES = 1_000_000 ++_PROMPT_TOKEN_RE = re.compile( ++ r"\{(?P__raw__(?::(?P[^{}]*))?|[a-zA-Z0-9_.]+)\}" ++) + # Hostnames/IP literals that only serve connections originating on the same + # machine. Anything else is treated as a public bind for safety-rail purposes. + _LOOPBACK_HOSTS = frozenset({ +@@ -174,6 +185,26 @@ def check_webhook_requirements() -> bool: + return AIOHTTP_AVAILABLE + + ++class IdempotencyResult(str, Enum): ++ """Outcome of binding a stable provider delivery identity.""" ++ ++ ACCEPTED = "accepted" ++ DUPLICATE = "duplicate" ++ CONFLICT = "conflict" ++ ++ ++def _bounded_positive_int(value: Any, *, default: int, maximum: int) -> int: ++ """Parse an integer setting and keep it inside a safe positive range.""" ++ if isinstance(value, bool): ++ parsed = default ++ else: ++ try: ++ parsed = int(value) ++ except (TypeError, ValueError, OverflowError): ++ parsed = default ++ return min(max(parsed, 1), maximum) ++ ++ + class WebhookAdapter(BasePlatformAdapter): + """Generic webhook receiver that triggers agent runs from HTTP POSTs.""" + +@@ -217,14 +248,25 @@ class WebhookAdapter(BasePlatformAdapter): + # Reference to gateway runner for cross-platform delivery (set externally) + self.gateway_runner = None + +- # Idempotency: TTL cache of recently processed delivery IDs. +- # Prevents duplicate agent runs when webhook providers retry. +- self._seen_deliveries: Dict[str, float] = {} ++ # Idempotency: TTL cache of provider-native retry identities. ++ # Keys include every authority boundary that can otherwise alias: ++ # (profile, route, provider, delivery_id). The original body hash is ++ # retained so conflicting reuse can be rejected with HTTP 409. ++ self._seen_deliveries: Dict[tuple[str, str, str, str], float] = {} ++ self._seen_delivery_bodies: Dict[tuple[str, str, str, str], str] = {} + self._idempotency_ttl: int = 3600 # 1 hour ++ self._idempotency_max_entries: int = _bounded_positive_int( ++ config.extra.get( ++ "idempotency_max_entries", _IDEMPOTENCY_DEFAULT_MAX_ENTRIES ++ ), ++ default=_IDEMPOTENCY_DEFAULT_MAX_ENTRIES, ++ maximum=_IDEMPOTENCY_MAX_ENTRIES_LIMIT, ++ ) ++ self._idempotency_lock = threading.RLock() + self._seen_deliveries_next_prune_at: float = 0.0 + +- # Rate limiting: per-route timestamps in a fixed window. +- self._rate_counts: Dict[str, Deque[float]] = {} ++ # Rate limiting is isolated by profile and route. ++ self._rate_counts: Dict[tuple[str, str], Deque[float]] = {} + self._rate_limit: int = int(config.extra.get("rate_limit", 30)) # per minute + + # Body size limit (auth-before-body pattern) +@@ -423,22 +465,65 @@ class WebhookAdapter(BasePlatformAdapter): + self._delivery_info.pop(key, None) + self._delivery_info_created.pop(key, None) + +- def _prune_seen_deliveries(self, now: float) -> None: +- """Occasionally prune expired delivery IDs without scanning every POST.""" +- if now < self._seen_deliveries_next_prune_at: +- return +- cutoff = now - self._idempotency_ttl +- stale = [k for k, t in self._seen_deliveries.items() if t < cutoff] +- for k in stale: +- self._seen_deliveries.pop(k, None) +- self._seen_deliveries_next_prune_at = now + min(60.0, max(1.0, self._idempotency_ttl / 10)) +- +- def _record_rate_limit_hit(self, route_name: str, now: float) -> bool: +- """Return True if route is still within limit after recording this hit.""" +- window = self._rate_counts.get(route_name) ++ def _prune_seen_deliveries( ++ self, ++ now: float, ++ *, ++ reserve: int = 0, ++ force: bool = False, ++ ) -> None: ++ """Expire old identities and enforce the configured hard ceiling. ++ ++ ``reserve`` leaves room for an imminent insertion. This makes the ++ configured ceiling true at the insertion boundary, including values ++ below the historical implicit floor of 128. ++ """ ++ target_size = max(0, self._idempotency_max_entries - reserve) ++ with self._idempotency_lock: ++ if ( ++ not force ++ and now < self._seen_deliveries_next_prune_at ++ and len(self._seen_deliveries) <= target_size ++ ): ++ return ++ ++ cutoff = now - self._idempotency_ttl ++ stale = [ ++ key ++ for key, seen_at in self._seen_deliveries.items() ++ if seen_at < cutoff ++ ] ++ for key in stale: ++ self._seen_deliveries.pop(key, None) ++ self._seen_delivery_bodies.pop(key, None) ++ ++ overflow = len(self._seen_deliveries) - target_size ++ if overflow > 0: ++ oldest = sorted( ++ self._seen_deliveries, ++ key=lambda key: self._seen_deliveries[key], ++ )[:overflow] ++ for key in oldest: ++ self._seen_deliveries.pop(key, None) ++ self._seen_delivery_bodies.pop(key, None) ++ ++ self._seen_deliveries_next_prune_at = now + min( ++ 60.0, max(1.0, self._idempotency_ttl / 10) ++ ) ++ ++ def _record_rate_limit_hit( ++ self, ++ route_name: str, ++ now: float, ++ *, ++ profile: Optional[str] = None, ++ ) -> bool: ++ """Record one hit in a profile/route-scoped fixed window.""" ++ key = ((profile or "default"), route_name) ++ window = self._rate_counts.get(key) + if not isinstance(window, deque): + new_window: Deque[float] = deque(window or ()) +- self._rate_counts[route_name] = new_window ++ self._rate_counts[key] = new_window + window = new_window + cutoff = now - _RATE_WINDOW_SECONDS + while window and window[0] < cutoff: +@@ -448,17 +533,178 @@ class WebhookAdapter(BasePlatformAdapter): + window.append(now) + return True + diff --git a/.hermes-patches/task10/part01.patch b/.hermes-patches/task10/part01.patch new file mode 100644 index 000000000000..b61382072135 --- /dev/null +++ b/.hermes-patches/task10/part01.patch @@ -0,0 +1,178 @@ +- def _record_delivery_id(self, delivery_id: str, now: float) -> bool: +- """Return True when this delivery should be processed.""" +- seen_at = self._seen_deliveries.get(delivery_id) +- if seen_at is not None and now - seen_at < self._idempotency_ttl: +- return False +- if seen_at is not None: +- self._seen_deliveries.pop(delivery_id, None) +- self._seen_deliveries[delivery_id] = now +- if len(self._seen_deliveries) > max(self._rate_limit * 2, 128): +- self._prune_seen_deliveries(now) +- return True ++ def _record_delivery_id( ++ self, ++ delivery_id: str, ++ now: float, ++ body_hash: str, ++ *, ++ profile: str, ++ route: str, ++ provider: str, ++ ) -> IdempotencyResult: ++ """Bind a stable provider identity to one body and authority scope.""" ++ normalized_id = str(delivery_id).strip() ++ normalized_hash = str(body_hash).strip() ++ normalized_scope = tuple( ++ str(value).strip() for value in (profile, route, provider) ++ ) ++ if not normalized_id: ++ raise ValueError("delivery_id must be a non-empty stable identity") ++ if not normalized_hash: ++ raise ValueError("body_hash must be non-empty") ++ if not all(normalized_scope): ++ raise ValueError("profile, route, and provider must be non-empty") ++ key = (*normalized_scope, normalized_id) ++ ++ with self._idempotency_lock: ++ seen_at = self._seen_deliveries.get(key) ++ if seen_at is not None and now - seen_at < self._idempotency_ttl: ++ previous_hash = self._seen_delivery_bodies.get(key, "") ++ if previous_hash and normalized_hash != previous_hash: ++ return IdempotencyResult.CONFLICT ++ return IdempotencyResult.DUPLICATE ++ ++ if seen_at is not None: ++ self._seen_deliveries.pop(key, None) ++ self._seen_delivery_bodies.pop(key, None) ++ ++ self._prune_seen_deliveries(now, reserve=1) ++ self._seen_deliveries[key] = now ++ self._seen_delivery_bodies[key] = normalized_hash ++ return IdempotencyResult.ACCEPTED ++ ++ @staticmethod ++ def _nonempty_identity(value: Any) -> Optional[str]: ++ """Normalize a provider identity without collapsing blanks together.""" ++ if isinstance(value, bool) or value is None: ++ return None ++ if not isinstance(value, (str, int)): ++ return None ++ normalized = str(value).strip() ++ return normalized or None ++ ++ @staticmethod ++ def _delivery_provider( ++ request: "web.Request", route_config: dict ++ ) -> tuple[str, bool]: ++ """Return the provider namespace and whether the route declared it.""" ++ configured: Optional[str] = None ++ for candidate in ( ++ route_config.get("provider"), ++ route_config.get("signature_mode"), ++ ): ++ if isinstance(candidate, str) and candidate.strip(): ++ configured = candidate.strip().lower() ++ break ++ ++ aliases = { ++ "agentmail": "svix", ++ "generic_v1": "generic", ++ "generic_v2": "generic", ++ "github_hmac_sha256": "github", ++ "gitlab_token": "gitlab", ++ "standard_webhooks": "gitlab_standard", ++ } ++ if configured is not None: ++ return aliases.get(configured, configured), True ++ ++ headers = request.headers ++ if any( ++ headers.get(name, "").strip() ++ for name in ("svix-id", "svix-signature", "svix-timestamp") ++ ): ++ return "svix", False ++ if any( ++ headers.get(name, "").strip() ++ for name in ("X-Hub-Signature-256", "X-GitHub-Delivery") ++ ): ++ return "github", False ++ if any( ++ headers.get(name, "").strip() ++ for name in ( ++ "X-Gitlab-Token", ++ "X-Gitlab-Event-UUID", ++ "X-Gitlab-Webhook-UUID", ++ "X-Gitlab-Idempotency-Key", ++ ) ++ ): ++ return "gitlab", False ++ if any( ++ headers.get(name, "").strip() ++ for name in ("webhook-id", "webhook-signature") ++ ): ++ return "gitlab_standard", False ++ if headers.get("X-Chatwoot-Delivery", "").strip(): ++ return "chatwoot", False ++ if headers.get("linear-signature", "").strip(): ++ return "linear", False ++ if headers.get("X-Hindsight-Signature", "").strip(): ++ return "hindsight", False ++ if any( ++ headers.get(name, "").strip() ++ for name in ( ++ "X-Webhook-Signature-V2", ++ "X-Webhook-Signature", ++ "X-Request-ID", ++ ) ++ ): ++ return "generic", False ++ return "generic", False ++ ++ def _resolve_delivery_identity( ++ self, ++ request: "web.Request", ++ route_config: dict, ++ payload: dict, ++ ) -> Optional[tuple[str, str]]: ++ """Resolve one trustworthy provider-native retry identity. ++ ++ No timestamp or blank-header fallback is permitted. When this returns ++ ``None`` the caller creates a unique trace/session ID but deliberately ++ skips retry deduplication. ++ """ ++ provider, declared = self._delivery_provider(request, route_config) ++ headers = request.headers ++ ++ header_names: tuple[str, ...] ++ if provider == "github": ++ header_names = ("X-GitHub-Delivery",) ++ elif provider == "svix": ++ header_names = ("svix-id",) ++ elif provider in {"gitlab", "gitlab_standard"}: ++ header_names = ( ++ "X-Gitlab-Event-UUID", ++ "X-Gitlab-Webhook-UUID", ++ "X-Gitlab-Idempotency-Key", ++ "Idempotency-Key", ++ "webhook-id", ++ ) ++ elif provider == "chatwoot": ++ header_names = ("X-Chatwoot-Delivery",) ++ elif provider == "generic": ++ header_names = ("X-Request-ID",) ++ else: ++ header_names = () ++ ++ for header_name in header_names: ++ candidate = self._nonempty_identity(headers.get(header_name)) ++ if candidate is not None: ++ return provider, candidate ++ ++ if provider == "stripe" and declared: ++ candidate = self._nonempty_identity(payload.get("id")) ++ if candidate is not None: ++ return provider, candidate ++ ++ # Chatwoot payload IDs are ambiguous outside a route that explicitly ++ # declares the Chatwoot scheme; never infer this fallback from shape. ++ if provider == "chatwoot" and declared: diff --git a/.hermes-patches/task10/part02.patch b/.hermes-patches/task10/part02.patch new file mode 100644 index 000000000000..ca15e4f04d86 --- /dev/null +++ b/.hermes-patches/task10/part02.patch @@ -0,0 +1,176 @@ ++ candidate = self._nonempty_identity(payload.get("id")) ++ if candidate is not None: ++ return provider, candidate ++ ++ return None + + async def get_chat_info(self, chat_id: str) -> Dict[str, Any]: + return {"name": chat_id, "type": "webhook"} +@@ -669,6 +915,14 @@ class WebhookAdapter(BasePlatformAdapter): + {"error": "Payload too large"}, status=413 + ) + ++ content_encoding = request.headers.get( ++ "Content-Encoding", "" ++ ).strip().lower() ++ if content_encoding not in {"", "identity"}: ++ return web.json_response( ++ {"error": "Unsupported Content-Encoding"}, status=415 ++ ) ++ + # Read body (must be done before any validation) + try: + raw_body = await request.read() +@@ -713,26 +967,50 @@ class WebhookAdapter(BasePlatformAdapter): + + # ── Rate limiting (after auth) ─────────────────────────── + now = time.time() +- if not self._record_rate_limit_hit(route_name, now): ++ effective_profile = profile or "default" ++ if not self._record_rate_limit_hit( ++ route_name, now, profile=effective_profile ++ ): + return web.json_response( +- {"error": "Rate limit exceeded"}, status=429 ++ {"error": "Rate limit exceeded"}, ++ status=429, ++ headers={"Retry-After": str(int(_RATE_WINDOW_SECONDS))}, + ) + +- # Parse payload +- try: +- payload = json.loads(raw_body) +- except json.JSONDecodeError: +- # Try form-encoded as fallback ++ # Parse only the declared media type. Invalid JSON must never fall ++ # through into form decoding, and webhook envelopes must be objects. ++ media_type = request.headers.get("Content-Type", "").split( ++ ";", 1 ++ )[0].strip().lower() ++ if media_type == "application/x-www-form-urlencoded": + try: + import urllib.parse + + payload = dict( + urllib.parse.parse_qsl(raw_body.decode("utf-8")) + ) +- except Exception: ++ except (UnicodeDecodeError, ValueError): + return web.json_response( +- {"error": "Cannot parse body"}, status=400 ++ {"error": "Cannot parse form body"}, status=400 + ) ++ elif media_type == "application/json" or media_type.endswith( ++ "+json" ++ ): ++ try: ++ payload = json.loads(raw_body) ++ except (json.JSONDecodeError, UnicodeDecodeError): ++ return web.json_response( ++ {"error": "Cannot parse JSON body"}, status=400 ++ ) ++ else: ++ return web.json_response( ++ {"error": "Unsupported Content-Type"}, status=415 ++ ) ++ ++ if not isinstance(payload, dict): ++ return web.json_response( ++ {"error": "JSON body must be an object"}, status=400 ++ ) + + # Check event type filter + event_type = ( +@@ -770,6 +1048,46 @@ class WebhookAdapter(BasePlatformAdapter): + } + ) + ++ identity = self._resolve_delivery_identity( ++ request, route_config, payload ++ ) ++ if identity is None: ++ provider, _declared = self._delivery_provider( ++ request, route_config ++ ) ++ delivery_id = f"generated-{uuid.uuid4().hex}" ++ stable_delivery_id = None ++ else: ++ provider, delivery_id = identity ++ stable_delivery_id = delivery_id ++ idem_result = self._record_delivery_id( ++ delivery_id, ++ now, ++ hashlib.sha256(raw_body).hexdigest(), ++ profile=effective_profile, ++ route=route_name, ++ provider=provider, ++ ) ++ if idem_result is IdempotencyResult.CONFLICT: ++ return web.json_response( ++ { ++ "status": "conflict", ++ "delivery_id": delivery_id, ++ "error": ( ++ "Idempotency key reused with a different body" ++ ), ++ }, ++ status=409, ++ ) ++ if idem_result is IdempotencyResult.DUPLICATE: ++ logger.info( ++ "[webhook] Skipping duplicate delivery %s", delivery_id ++ ) ++ return web.json_response( ++ {"status": "duplicate", "delivery_id": delivery_id}, ++ status=200, ++ ) ++ + if route_config.get("script"): + # run_route_script shells out (subprocess.run, up to its timeout); + # run it in a worker thread so it can't block the gateway event loop. +@@ -795,9 +1113,20 @@ class WebhookAdapter(BasePlatformAdapter): + + # Format prompt from template + prompt_template = route_config.get("prompt", "") +- prompt = self._render_prompt( +- prompt_template, payload, event_type, route_name +- ) ++ try: ++ prompt = self._render_prompt( ++ prompt_template, payload, event_type, route_name ++ ) ++ except ValueError as exc: ++ logger.error( ++ "[webhook] Invalid prompt template for route %s: %s", ++ route_name, ++ exc, ++ ) ++ return web.json_response( ++ {"error": "Webhook route has an invalid prompt template"}, ++ status=500, ++ ) + + # Inject skill content if configured. + # We call build_skill_invocation_message() directly rather than +@@ -828,27 +1157,6 @@ class WebhookAdapter(BasePlatformAdapter): + except Exception as e: + logger.warning("[webhook] Skill loading failed: %s", e) + +- # Build a unique delivery ID +- delivery_id = request.headers.get( +- "X-GitHub-Delivery", +- request.headers.get( +- "svix-id", +- request.headers.get("X-Request-ID", str(int(time.time() * 1000))), +- ), +- ) +- +- # ── Idempotency ───────────────────────────────────────── +- # Skip duplicate deliveries (webhook retries). +- now = time.time() +- if not self._record_delivery_id(delivery_id, now): +- logger.info( +- "[webhook] Skipping duplicate delivery %s", delivery_id +- ) +- return web.json_response( +- {"status": "duplicate", "delivery_id": delivery_id}, diff --git a/.hermes-patches/task10/part03.patch b/.hermes-patches/task10/part03.patch new file mode 100644 index 000000000000..65577112f0aa --- /dev/null +++ b/.hermes-patches/task10/part03.patch @@ -0,0 +1,163 @@ +- status=200, +- ) +- + # ── Direct delivery mode (deliver_only) ───────────────── + # Skip the agent entirely — the rendered prompt IS the message we + # deliver. Use case: external services (Supabase, monitoring, +@@ -967,6 +1275,9 @@ class WebhookAdapter(BasePlatformAdapter): + "route": route_name, + "event": event_type, + "delivery_id": delivery_id, ++ "deduplication": ( ++ "stable" if stable_delivery_id is not None else "skipped" ++ ), + }, + status=202, + ) +@@ -1251,40 +1562,114 @@ class WebhookAdapter(BasePlatformAdapter): + event_type: str, + route_name: str, + ) -> str: +- """Render a prompt template with the webhook payload. ++ """Render a prompt template with bounded, parseable raw envelopes. + +- Supports dot-notation access into nested dicts: +- ``{pull_request.title}`` → ``payload["pull_request"]["title"]`` +- +- Special token ``{__raw__}`` dumps the entire payload as indented +- JSON (truncated to 4000 chars). Useful for monitoring alerts or +- any webhook where the agent needs to see the full payload. ++ ``{__raw__}`` retains its 4,000-byte default. ``{__raw__:N}`` ++ selects an explicit complete-envelope byte cap between 64 and ++ 1,000,000 bytes. The cap includes metadata and JSON escaping. + """ + if not template: +- truncated = json.dumps(payload, indent=2)[:4000] ++ raw_envelope = self._render_raw_payload( ++ payload, _RAW_PAYLOAD_DEFAULT_CAP_BYTES ++ ) + return ( + f"Webhook event '{event_type}' on route " +- f"'{route_name}':\n\n```json\n{truncated}\n```" ++ f"'{route_name}':\n\n```json\n{raw_envelope}\n```" + ) + + def _resolve(match: re.Match) -> str: +- key = match.group(1) +- # Special token: dump the entire payload as JSON +- if key == "__raw__": +- return json.dumps(payload, indent=2)[:4000] +- if key == "event_type": ++ token = match.group("token") ++ if token.startswith("__raw__"): ++ cap_text = match.group("raw_cap") ++ if cap_text is None: ++ cap = _RAW_PAYLOAD_DEFAULT_CAP_BYTES ++ else: ++ try: ++ cap = int(cap_text) ++ except (TypeError, ValueError) as exc: ++ raise ValueError( ++ f"invalid raw payload cap: {cap_text!r}" ++ ) from exc ++ return self._render_raw_payload(payload, cap) ++ if token == "event_type": + return event_type + value: Any = payload +- for part in key.split("."): ++ for part in token.split("."): + if isinstance(value, dict): +- value = value.get(part, f"{{{key}}}") ++ value = value.get(part, f"{{{token}}}") + else: +- return f"{{{key}}}" ++ return f"{{{token}}}" + if isinstance(value, (dict, list)): + return json.dumps(value, indent=2)[:2000] + return str(value) + +- return re.sub(r"\{([a-zA-Z0-9_.]+)\}", _resolve, template) ++ # One substitution pass is deliberate: payload text emitted by ++ # ``{__raw__}`` may itself contain brace-shaped strings and must never ++ # be interpreted as a second template layer. ++ return _PROMPT_TOKEN_RE.sub(_resolve, template) ++ ++ def _render_raw_payload( ++ self, payload: dict, cap: int = _RAW_PAYLOAD_DEFAULT_CAP_BYTES ++ ) -> str: ++ """Render a valid JSON envelope within the complete UTF-8 byte cap.""" ++ try: ++ requested_cap = int(cap) ++ except (TypeError, ValueError) as exc: ++ raise ValueError(f"invalid raw payload cap: {cap!r}") from exc ++ if not ( ++ _RAW_PAYLOAD_MIN_CAP_BYTES ++ <= requested_cap ++ <= _RAW_PAYLOAD_MAX_CAP_BYTES ++ ): ++ raise ValueError( ++ "raw payload cap must be between " ++ f"{_RAW_PAYLOAD_MIN_CAP_BYTES} and " ++ f"{_RAW_PAYLOAD_MAX_CAP_BYTES} bytes" ++ ) ++ ++ serialized = json.dumps(payload, indent=2, ensure_ascii=False) ++ serialized_bytes = serialized.encode("utf-8") ++ original_bytes = len(serialized_bytes) ++ ++ def _envelope(payload_text: str, *, truncated: bool) -> str: ++ return json.dumps( ++ { ++ "payload": payload_text, ++ "truncated": truncated, ++ "original_bytes": original_bytes, ++ }, ++ ensure_ascii=False, ++ separators=(",", ":"), ++ ) ++ ++ full = _envelope(serialized, truncated=False) ++ if len(full.encode("utf-8")) <= requested_cap: ++ return full ++ ++ smallest = _envelope("", truncated=True) ++ if len(smallest.encode("utf-8")) > requested_cap: ++ raise ValueError( ++ "raw payload cap is too small for the envelope metadata" ++ ) ++ ++ # JSON-string encoding size is monotonic as the prefix grows, so a ++ # binary search finds the largest code-point-safe payload prefix whose ++ # complete escaped envelope still fits the selected cap. ++ low = 0 ++ high = min(original_bytes, requested_cap) ++ best = smallest ++ while low <= high: ++ midpoint = (low + high) // 2 ++ bounded = serialized_bytes[:midpoint].decode( ++ "utf-8", errors="ignore" ++ ) ++ candidate = _envelope(bounded, truncated=True) ++ if len(candidate.encode("utf-8")) <= requested_cap: ++ best = candidate ++ low = midpoint + 1 ++ else: ++ high = midpoint - 1 ++ return best + + def _render_delivery_extra( + self, extra: dict, payload: dict +diff --git a/tests/gateway/test_webhook_adapter.py b/tests/gateway/test_webhook_adapter.py +index 4f5cdb1..21be925 100644 +--- a/tests/gateway/test_webhook_adapter.py ++++ b/tests/gateway/test_webhook_adapter.py +@@ -791,8 +791,11 @@ class TestRawTemplateToken: + "Action={action} Raw={__raw__}", payload, "push", "test" + ) + assert result.startswith("Action=closed Raw=") +- assert '"action": "closed"' in result +- assert '"number": 7' in result ++ envelope = json.loads(result.split("Raw=", 1)[1]) ++ assert envelope["truncated"] is False ++ assert envelope["original_bytes"] > 0 ++ assert '"action": "closed"' in envelope["payload"] diff --git a/.hermes-patches/task10/part04.patch b/.hermes-patches/task10/part04.patch new file mode 100644 index 000000000000..054d292fe746 --- /dev/null +++ b/.hermes-patches/task10/part04.patch @@ -0,0 +1,213 @@ ++ assert '"number": 7' in envelope["payload"] + + + # =================================================================== +diff --git a/tests/gateway/test_webhook_task10_closure.py b/tests/gateway/test_webhook_task10_closure.py +new file mode 100644 +index 0000000..92fb404 +--- /dev/null ++++ b/tests/gateway/test_webhook_task10_closure.py +@@ -0,0 +1,557 @@ ++"""Exact contracts for webhook Task 10 intake/idempotency closure. ++ ++These tests cover the authority boundaries and malformed-input cases that were ++missing from the original fan-out patch. They intentionally exercise the real ++HTTP handler where request ordering matters, and pure helpers where the ++negative matrix is clearer without transport noise. ++""" ++ ++from __future__ import annotations ++ ++import asyncio ++import hashlib ++import hmac ++import json ++from concurrent.futures import ThreadPoolExecutor ++from types import SimpleNamespace ++ ++import pytest ++from aiohttp import web ++from aiohttp.test_utils import TestClient, TestServer ++ ++from gateway.config import PlatformConfig ++from gateway.platforms.webhook import ( ++ IdempotencyResult, ++ WebhookAdapter, ++ _IDEMPOTENCY_DEFAULT_MAX_ENTRIES, ++ _IDEMPOTENCY_MAX_ENTRIES_LIMIT, ++ _INSECURE_NO_AUTH, ++ _RAW_PAYLOAD_DEFAULT_CAP_BYTES, ++ _RAW_PAYLOAD_MAX_CAP_BYTES, ++ _RAW_PAYLOAD_MIN_CAP_BYTES, ++) ++ ++ ++def _adapter( ++ routes: dict | None = None, ++ *, ++ max_entries: object = 8, ++ rate_limit: int = 30, ++) -> WebhookAdapter: ++ adapter = WebhookAdapter( ++ PlatformConfig( ++ enabled=True, ++ extra={ ++ "host": "127.0.0.1", ++ "port": 0, ++ "routes": routes or {}, ++ "idempotency_max_entries": max_entries, ++ "rate_limit": rate_limit, ++ }, ++ ) ++ ) ++ # Dynamic-route behavior is covered elsewhere. Freezing it here keeps these ++ # tests about request semantics rather than the caller's Hermes home. ++ adapter._reload_dynamic_routes = lambda: None ++ return adapter ++ ++ ++def _app(adapter: WebhookAdapter) -> web.Application: ++ app = web.Application(client_max_size=adapter._max_body_bytes) ++ app.router.add_post("/webhooks/{route_name}", adapter._handle_webhook) ++ return app ++ ++ ++def _github_headers( ++ body: bytes, ++ secret: str, ++ *, ++ delivery_id: str | None, ++ event: str = "push", ++) -> dict[str, str]: ++ headers = { ++ "Content-Type": "application/json", ++ "X-Hub-Signature-256": "sha256=" ++ + hmac.new(secret.encode(), body, hashlib.sha256).hexdigest(), ++ "X-GitHub-Event": event, ++ } ++ if delivery_id is not None: ++ headers["X-GitHub-Delivery"] = delivery_id ++ return headers ++ ++ ++def _route(secret: str = "secret", **extra: object) -> dict: ++ return { ++ "secret": secret, ++ "provider": "github", ++ "prompt": "{event_type}", ++ "deliver": "log", ++ **extra, ++ } ++ ++ ++class TestIdempotencyState: ++ def test_key_isolated_by_profile_route_and_provider(self) -> None: ++ adapter = _adapter() ++ scopes = [ ++ ("alpha", "route-a", "github"), ++ ("beta", "route-a", "github"), ++ ("alpha", "route-b", "github"), ++ ("alpha", "route-a", "gitlab"), ++ ] ++ ++ for profile, route, provider in scopes: ++ result = adapter._record_delivery_id( ++ "same-id", ++ 1000.0, ++ "same-body", ++ profile=profile, ++ route=route, ++ provider=provider, ++ ) ++ assert result is IdempotencyResult.ACCEPTED ++ ++ assert len(adapter._seen_deliveries) == len(scopes) ++ ++ def test_duplicate_and_conflict_are_explicit_results(self) -> None: ++ adapter = _adapter() ++ kwargs = { ++ "profile": "default", ++ "route": "route", ++ "provider": "github", ++ } ++ ++ assert adapter._record_delivery_id( ++ "delivery", 1000.0, "body-a", **kwargs ++ ) is IdempotencyResult.ACCEPTED ++ assert adapter._record_delivery_id( ++ "delivery", 1001.0, "body-a", **kwargs ++ ) is IdempotencyResult.DUPLICATE ++ assert adapter._record_delivery_id( ++ "delivery", 1002.0, "body-b", **kwargs ++ ) is IdempotencyResult.CONFLICT ++ ++ def test_concurrent_same_key_has_one_winner(self) -> None: ++ adapter = _adapter() ++ ++ def record() -> IdempotencyResult: ++ return adapter._record_delivery_id( ++ "concurrent", ++ 1000.0, ++ "body", ++ profile="default", ++ route="route", ++ provider="github", ++ ) ++ ++ with ThreadPoolExecutor(max_workers=16) as pool: ++ results = list(pool.map(lambda _index: record(), range(32))) ++ ++ assert results.count(IdempotencyResult.ACCEPTED) == 1 ++ assert results.count(IdempotencyResult.DUPLICATE) == 31 ++ assert len(adapter._seen_deliveries) == 1 ++ ++ def test_configured_ceiling_is_true_at_every_insertion(self) -> None: ++ adapter = _adapter(max_entries=4) ++ ++ for index in range(30): ++ assert adapter._record_delivery_id( ++ str(index), ++ float(index), ++ f"body-{index}", ++ profile="default", ++ route="route", ++ provider="github", ++ ) is IdempotencyResult.ACCEPTED ++ assert len(adapter._seen_deliveries) <= 4 ++ assert len(adapter._seen_delivery_bodies) <= 4 ++ ++ surviving_ids = {key[-1] for key in adapter._seen_deliveries} ++ assert surviving_ids == {"26", "27", "28", "29"} ++ ++ @pytest.mark.parametrize( ++ ("configured", "expected"), ++ [ ++ (0, 1), ++ (-50, 1), ++ ("not-an-int", _IDEMPOTENCY_DEFAULT_MAX_ENTRIES), ++ (True, _IDEMPOTENCY_DEFAULT_MAX_ENTRIES), ++ (False, _IDEMPOTENCY_DEFAULT_MAX_ENTRIES), ++ (float("inf"), _IDEMPOTENCY_DEFAULT_MAX_ENTRIES), ++ ( ++ _IDEMPOTENCY_MAX_ENTRIES_LIMIT + 1, ++ _IDEMPOTENCY_MAX_ENTRIES_LIMIT, ++ ), ++ ], ++ ) ++ def test_invalid_or_nonpositive_ceiling_is_clamped( ++ self, configured: object, expected: int ++ ) -> None: ++ assert _adapter(max_entries=configured)._idempotency_max_entries == expected ++ ++ @pytest.mark.parametrize( ++ ("overrides", "message"), ++ [ ++ ({"delivery_id": ""}, "delivery_id"), ++ ({"body_hash": ""}, "body_hash"), ++ ({"profile": ""}, "profile, route, and provider"), ++ ({"route": ""}, "profile, route, and provider"), ++ ({"provider": ""}, "profile, route, and provider"), ++ ], ++ ) ++ def test_invalid_binding_components_fail_closed( ++ self, overrides: dict[str, str], message: str diff --git a/.hermes-patches/task10/part05.patch b/.hermes-patches/task10/part05.patch new file mode 100644 index 000000000000..ea88de1df173 --- /dev/null +++ b/.hermes-patches/task10/part05.patch @@ -0,0 +1,189 @@ ++ ) -> None: ++ arguments = { ++ "delivery_id": "delivery", ++ "now": 1000.0, ++ "body_hash": "body", ++ "profile": "default", ++ "route": "route", ++ "provider": "github", ++ } ++ arguments.update(overrides) ++ with pytest.raises(ValueError, match=message): ++ _adapter()._record_delivery_id(**arguments) ++ ++ def test_rate_limit_isolated_by_profile_and_route(self) -> None: ++ adapter = _adapter(rate_limit=1) ++ ++ assert adapter._record_rate_limit_hit("route", 1000.0, profile="alpha") ++ assert adapter._record_rate_limit_hit("route", 1000.0, profile="beta") ++ assert adapter._record_rate_limit_hit("other", 1000.0, profile="alpha") ++ assert not adapter._record_rate_limit_hit( ++ "route", 1000.1, profile="alpha" ++ ) ++ ++ ++class TestProviderDeliveryIdentity: ++ @staticmethod ++ def _request(headers: dict[str, str]) -> SimpleNamespace: ++ return SimpleNamespace(headers=headers) ++ ++ @pytest.mark.parametrize( ++ ("route", "headers", "payload", "expected"), ++ [ ++ ( ++ {"provider": "github"}, ++ {"X-GitHub-Delivery": "gh-1"}, ++ {}, ++ ("github", "gh-1"), ++ ), ++ ({"provider": "svix"}, {"svix-id": "msg-1"}, {}, ("svix", "msg-1")), ++ ( ++ {"provider": "gitlab"}, ++ {"X-Gitlab-Event-UUID": "gl-1"}, ++ {}, ++ ("gitlab", "gl-1"), ++ ), ++ ( ++ {"signature_mode": "gitlab_standard"}, ++ {"webhook-id": "std-1"}, ++ {}, ++ ("gitlab_standard", "std-1"), ++ ), ++ ( ++ {"provider": "generic"}, ++ {"X-Request-ID": "req-1"}, ++ {}, ++ ("generic", "req-1"), ++ ), ++ ({"provider": "stripe"}, {}, {"id": "evt_1"}, ("stripe", "evt_1")), ++ ( ++ {"provider": "chatwoot"}, ++ {"X-Chatwoot-Delivery": "cw-header"}, ++ {"id": 77}, ++ ("chatwoot", "cw-header"), ++ ), ++ ({"provider": "chatwoot"}, {}, {"id": 77}, ("chatwoot", "77")), ++ ( ++ {}, ++ {"X-Chatwoot-Delivery": "cw-inferred"}, ++ {}, ++ ("chatwoot", "cw-inferred"), ++ ), ++ ], ++ ) ++ def test_stable_provider_native_matrix( ++ self, ++ route: dict, ++ headers: dict[str, str], ++ payload: dict, ++ expected: tuple[str, str], ++ ) -> None: ++ adapter = _adapter() ++ assert adapter._resolve_delivery_identity( ++ self._request(headers), route, payload ++ ) == expected ++ ++ @pytest.mark.parametrize( ++ ("route", "headers", "payload"), ++ [ ++ ({"provider": "github"}, {"X-GitHub-Delivery": " "}, {}), ++ ({"provider": "generic"}, {"X-Request-ID": ""}, {}), ++ ({}, {}, {"id": "shape-is-not-proof"}), ++ ( ++ {"provider": "generic"}, ++ {"X-Webhook-Timestamp": "1700000000"}, ++ {}, ++ ), ++ ({"provider": "stripe"}, {}, {"id": []}), ++ ({"provider": "chatwoot"}, {}, {"id": None}), ++ ], ++ ) ++ def test_blank_or_unqualified_identity_skips_dedup( ++ self, route: dict, headers: dict[str, str], payload: dict ++ ) -> None: ++ adapter = _adapter() ++ assert adapter._resolve_delivery_identity( ++ self._request(headers), route, payload ++ ) is None ++ ++ ++class TestHttpContract: ++ @pytest.mark.asyncio ++ async def test_same_delivery_runs_each_route_once(self) -> None: ++ secret = "secret" ++ adapter = _adapter({"a": _route(secret), "b": _route(secret)}) ++ events = [] ++ ++ async def capture(event) -> None: ++ events.append(event) ++ ++ adapter.handle_message = capture ++ body = json.dumps({"value": 1}).encode() ++ headers = _github_headers(body, secret, delivery_id="fanout-1") ++ ++ async with TestClient(TestServer(_app(adapter))) as client: ++ first_a = await client.post("/webhooks/a", data=body, headers=headers) ++ first_b = await client.post("/webhooks/b", data=body, headers=headers) ++ retry_a = await client.post("/webhooks/a", data=body, headers=headers) ++ assert first_a.status == 202 ++ assert first_b.status == 202 ++ assert retry_a.status == 200 ++ assert (await retry_a.json())["status"] == "duplicate" ++ ++ await asyncio.sleep(0.01) ++ assert len(events) == 2 ++ ++ @pytest.mark.asyncio ++ async def test_conflicting_body_returns_409(self) -> None: ++ secret = "secret" ++ adapter = _adapter({"route": _route(secret)}) ++ body_a = json.dumps({"value": 1}).encode() ++ body_b = json.dumps({"value": 2}).encode() ++ ++ async with TestClient(TestServer(_app(adapter))) as client: ++ first = await client.post( ++ "/webhooks/route", ++ data=body_a, ++ headers=_github_headers(body_a, secret, delivery_id="conflict"), ++ ) ++ conflict = await client.post( ++ "/webhooks/route", ++ data=body_b, ++ headers=_github_headers(body_b, secret, delivery_id="conflict"), ++ ) ++ assert first.status == 202 ++ assert conflict.status == 409 ++ assert (await conflict.json())["status"] == "conflict" ++ ++ @pytest.mark.asyncio ++ async def test_missing_stable_id_gets_unique_trace_without_dedup(self) -> None: ++ adapter = _adapter( ++ { ++ "route": { ++ "secret": _INSECURE_NO_AUTH, ++ "prompt": "{event_type}", ++ "deliver": "log", ++ } ++ } ++ ) ++ events = [] ++ ++ async def capture(event) -> None: ++ events.append(event) ++ ++ adapter.handle_message = capture ++ ++ async with TestClient(TestServer(_app(adapter))) as client: ++ first = await client.post("/webhooks/route", json={"value": 1}) ++ second = await client.post("/webhooks/route", json={"value": 1}) ++ first_body = await first.json() ++ second_body = await second.json() ++ ++ await asyncio.sleep(0.01) ++ assert first.status == second.status == 202 ++ assert first_body["deduplication"] == "skipped" ++ assert second_body["deduplication"] == "skipped" ++ assert first_body["delivery_id"].startswith("generated-") ++ assert second_body["delivery_id"].startswith("generated-") ++ assert first_body["delivery_id"] != second_body["delivery_id"] ++ assert len(events) == 2 diff --git a/.hermes-patches/task10/part06.patch b/.hermes-patches/task10/part06.patch new file mode 100644 index 000000000000..ceb3b86d2f76 --- /dev/null +++ b/.hermes-patches/task10/part06.patch @@ -0,0 +1,165 @@ ++ assert events[0].source.chat_id != events[1].source.chat_id ++ assert adapter._seen_deliveries == {} ++ ++ @pytest.mark.asyncio ++ async def test_non_object_json_is_rejected_cleanly(self) -> None: ++ adapter = _adapter( ++ {"route": _route(_INSECURE_NO_AUTH, provider="generic")} ++ ) ++ async with TestClient(TestServer(_app(adapter))) as client: ++ response = await client.post("/webhooks/route", json=[1, 2, 3]) ++ assert response.status == 400 ++ assert (await response.json())["error"] == ( ++ "JSON body must be an object" ++ ) ++ ++ @pytest.mark.asyncio ++ @pytest.mark.parametrize("content_type", [None, "text/plain"]) ++ async def test_missing_or_unsupported_content_type_returns_415( ++ self, content_type: str | None ++ ) -> None: ++ adapter = _adapter( ++ {"route": _route(_INSECURE_NO_AUTH, provider="generic")} ++ ) ++ headers = {} if content_type is None else {"Content-Type": content_type} ++ async with TestClient(TestServer(_app(adapter))) as client: ++ response = await client.post( ++ "/webhooks/route", data=b"{}", headers=headers ++ ) ++ assert response.status == 415 ++ ++ @pytest.mark.asyncio ++ async def test_compressed_content_is_rejected_before_parsing(self) -> None: ++ adapter = _adapter( ++ {"route": _route(_INSECURE_NO_AUTH, provider="generic")} ++ ) ++ async with TestClient(TestServer(_app(adapter))) as client: ++ response = await client.post( ++ "/webhooks/route", ++ data=b"not-really-gzip", ++ headers={ ++ "Content-Type": "application/json", ++ "Content-Encoding": "gzip", ++ }, ++ ) ++ assert response.status == 415 ++ ++ @pytest.mark.asyncio ++ async def test_rate_limit_returns_retry_after(self) -> None: ++ adapter = _adapter( ++ {"route": _route(_INSECURE_NO_AUTH)}, rate_limit=1 ++ ) ++ async with TestClient(TestServer(_app(adapter))) as client: ++ first = await client.post( ++ "/webhooks/route", ++ json={"value": 1}, ++ headers={"X-GitHub-Delivery": "one"}, ++ ) ++ second = await client.post( ++ "/webhooks/route", ++ json={"value": 2}, ++ headers={"X-GitHub-Delivery": "two"}, ++ ) ++ assert first.status == 202 ++ assert second.status == 429 ++ assert second.headers["Retry-After"] == "60" ++ ++ @pytest.mark.asyncio ++ async def test_idempotency_precedes_script_side_effects(self) -> None: ++ secret = "secret" ++ adapter = _adapter( ++ {"route": _route(secret, script="transform.py")} ++ ) ++ calls = 0 ++ ++ def run_script(_script: str, payload: dict) -> tuple[bool, dict]: ++ nonlocal calls ++ calls += 1 ++ return True, payload ++ ++ adapter._route_processor.run_route_script = run_script ++ body = json.dumps({"value": 1}).encode() ++ headers = _github_headers(body, secret, delivery_id="script-once") ++ ++ async with TestClient(TestServer(_app(adapter))) as client: ++ first = await client.post( ++ "/webhooks/route", data=body, headers=headers ++ ) ++ retry = await client.post( ++ "/webhooks/route", data=body, headers=headers ++ ) ++ assert first.status == 202 ++ assert retry.status == 200 ++ ++ assert calls == 1 ++ ++ ++class TestRawPayloadEnvelope: ++ def test_bare_token_uses_complete_4000_byte_cap(self) -> None: ++ adapter = _adapter() ++ payload = {"text": ("é\\\"{event_type}" * 1000)} ++ rendered = adapter._render_prompt( ++ "Raw={__raw__}", payload, "push", "route" ++ ) ++ envelope_text = rendered.split("Raw=", 1)[1] ++ envelope = json.loads(envelope_text) ++ ++ assert len(envelope_text.encode("utf-8")) <= ( ++ _RAW_PAYLOAD_DEFAULT_CAP_BYTES ++ ) ++ assert envelope["truncated"] is True ++ assert envelope["original_bytes"] == len( ++ json.dumps(payload, indent=2, ensure_ascii=False).encode("utf-8") ++ ) ++ ++ def test_parameterized_cap_bounds_escaping_and_multibyte_text(self) -> None: ++ adapter = _adapter() ++ payload = {"text": ("\\\"é" * 500)} ++ rendered = adapter._render_prompt( ++ "{__raw__:256}", payload, "push", "route" ++ ) ++ envelope = json.loads(rendered) ++ ++ assert len(rendered.encode("utf-8")) <= 256 ++ assert envelope["truncated"] is True ++ ++ def test_raw_payload_is_not_treated_as_a_second_template_layer(self) -> None: ++ adapter = _adapter() ++ payload = { ++ "text": "literal {event_type} and {nested.value}", ++ "nested": {"value": "resolved-only-outside-raw"}, ++ } ++ rendered = adapter._render_prompt( ++ "Event={event_type} Raw={__raw__:512}", ++ payload, ++ "push", ++ "route", ++ ) ++ envelope = json.loads(rendered.split(" Raw=", 1)[1]) ++ ++ assert rendered.startswith("Event=push Raw=") ++ assert "literal {event_type} and {nested.value}" in envelope["payload"] ++ ++ def test_empty_template_uses_same_parseable_envelope(self) -> None: ++ adapter = _adapter() ++ rendered = adapter._render_prompt("", {"value": "é"}, "push", "r") ++ fenced = rendered.split("```json\n", 1)[1].rsplit("\n```", 1)[0] ++ envelope = json.loads(fenced) ++ ++ assert envelope["truncated"] is False ++ assert envelope["original_bytes"] == len( ++ json.dumps( ++ {"value": "é"}, indent=2, ensure_ascii=False ++ ).encode("utf-8") ++ ) ++ ++ @pytest.mark.parametrize( ++ "cap", ++ [ ++ _RAW_PAYLOAD_MIN_CAP_BYTES - 1, ++ _RAW_PAYLOAD_MAX_CAP_BYTES + 1, ++ ], ++ ) ++ def test_absurd_cap_is_rejected(self, cap: int) -> None: ++ with pytest.raises(ValueError): ++ _adapter()._render_raw_payload({"value": 1}, cap) diff --git a/.hermes-patches/task10/trigger b/.hermes-patches/task10/trigger new file mode 100644 index 000000000000..856c4882cc1c --- /dev/null +++ b/.hermes-patches/task10/trigger @@ -0,0 +1 @@ +pr85523 diff --git a/artifacts/webhook-repair/50d98fc1f3d49d7a7b522eaa7f4553cd864a0218/task-10/preflight.json b/artifacts/webhook-repair/50d98fc1f3d49d7a7b522eaa7f4553cd864a0218/task-10/preflight.json new file mode 100644 index 000000000000..e42b63796a23 --- /dev/null +++ b/artifacts/webhook-repair/50d98fc1f3d49d7a7b522eaa7f4553cd864a0218/task-10/preflight.json @@ -0,0 +1,13 @@ +{ + "branch": "wr/task-10-intake-contract", + "expected_head": "b402cf4c01909b74ebfac65f8a4dedef358b9309", + "head": "b402cf4c01909b74ebfac65f8a4dedef358b9309", + "mode": "existing_pr", + "phase": "preflight", + "repository": "D:\\HERMES-TEMP\\webhook-revolution\\task-10-intake-contract", + "schema_version": 1, + "status": "pass", + "task": 10, + "task_title": "HTTP intake, idempotency, limits, filters, and transforms", + "timestamp": "2026-08-14T23:54:53Z" +} diff --git a/artifacts/webhook-repair/50d98fc1f3d49d7a7b522eaa7f4553cd864a0218/task-10/verification.json b/artifacts/webhook-repair/50d98fc1f3d49d7a7b522eaa7f4553cd864a0218/task-10/verification.json new file mode 100644 index 000000000000..d46a87d8a22f --- /dev/null +++ b/artifacts/webhook-repair/50d98fc1f3d49d7a7b522eaa7f4553cd864a0218/task-10/verification.json @@ -0,0 +1,16 @@ +{ + "branch": "wr/task-10-intake-contract", + "changed_files": [ + "gateway/platforms/webhook.py", + "tests/gateway/test_webhook_http_contract.py", + "tests/gateway/test_webhook_intake_hardening.py" + ], + "head": "b402cf4c01909b74ebfac65f8a4dedef358b9309", + "head_before_commit": "b402cf4c01909b74ebfac65f8a4dedef358b9309", + "phase": "verification", + "repository": "D:\\HERMES-TEMP\\webhook-revolution\\task-10-intake-contract", + "schema_version": 1, + "status": "pass", + "task": 10, + "timestamp": "2026-08-14T23:56:11Z" +} diff --git a/gateway/platforms/webhook.py b/gateway/platforms/webhook.py index bc3ebe0968e9..376fbd27c37f 100644 --- a/gateway/platforms/webhook.py +++ b/gateway/platforms/webhook.py @@ -219,12 +219,18 @@ def __init__(self, config: PlatformConfig): # Idempotency: TTL cache of recently processed delivery IDs. # Prevents duplicate agent runs when webhook providers retry. - self._seen_deliveries: Dict[str, float] = {} + # Keyed by (profile, route, provider, delivery_id); bound to a body hash + # so a conflicting replay can be reported as 409. + self._seen_deliveries: Dict[tuple, float] = {} + self._seen_delivery_bodies: Dict[tuple, str] = {} self._idempotency_ttl: int = 3600 # 1 hour + self._idempotency_max_entries: int = int( + config.extra.get("idempotency_max_entries", 4096) + ) self._seen_deliveries_next_prune_at: float = 0.0 - # Rate limiting: per-route timestamps in a fixed window. - self._rate_counts: Dict[str, Deque[float]] = {} + # Rate limiting: per-profile/per-route timestamps in a fixed window. + self._rate_counts: Dict[tuple[str, str], Deque[float]] = {} self._rate_limit: int = int(config.extra.get("rate_limit", 30)) # per minute # Body size limit (auth-before-body pattern) @@ -424,21 +430,37 @@ def _prune_delivery_info(self, now: float) -> None: self._delivery_info_created.pop(key, None) def _prune_seen_deliveries(self, now: float) -> None: - """Occasionally prune expired delivery IDs without scanning every POST.""" - if now < self._seen_deliveries_next_prune_at: + """Prune expired and oldest idempotency entries to a hard size ceiling.""" + if now < self._seen_deliveries_next_prune_at and len(self._seen_deliveries) <= self._idempotency_max_entries: return cutoff = now - self._idempotency_ttl - stale = [k for k, t in self._seen_deliveries.items() if t < cutoff] - for k in stale: - self._seen_deliveries.pop(k, None) - self._seen_deliveries_next_prune_at = now + min(60.0, max(1.0, self._idempotency_ttl / 10)) - - def _record_rate_limit_hit(self, route_name: str, now: float) -> bool: - """Return True if route is still within limit after recording this hit.""" - window = self._rate_counts.get(route_name) + stale = [key for key, seen_at in self._seen_deliveries.items() if seen_at < cutoff] + for key in stale: + self._seen_deliveries.pop(key, None) + self._seen_delivery_bodies.pop(key, None) + overflow = len(self._seen_deliveries) - self._idempotency_max_entries + if overflow > 0: + oldest = sorted(self._seen_deliveries, key=self._seen_deliveries.get)[:overflow] + for key in oldest: + self._seen_deliveries.pop(key, None) + self._seen_delivery_bodies.pop(key, None) + self._seen_deliveries_next_prune_at = now + min( + 60.0, max(1.0, self._idempotency_ttl / 10) + ) + + def _record_rate_limit_hit( + self, + route_name: str, + now: float, + *, + profile: str | None = None, + ) -> bool: + """Record one hit against the profile/route-scoped rate bucket.""" + key = (profile or self._profile_scope_key(), route_name) + window = self._rate_counts.get(key) if not isinstance(window, deque): new_window: Deque[float] = deque(window or ()) - self._rate_counts[route_name] = new_window + self._rate_counts[key] = new_window window = new_window cutoff = now - _RATE_WINDOW_SECONDS while window and window[0] < cutoff: @@ -448,14 +470,62 @@ def _record_rate_limit_hit(self, route_name: str, now: float) -> bool: window.append(now) return True - def _record_delivery_id(self, delivery_id: str, now: float) -> bool: - """Return True when this delivery should be processed.""" - seen_at = self._seen_deliveries.get(delivery_id) - if seen_at is not None and now - seen_at < self._idempotency_ttl: + def _profile_scope_key(self) -> str: + """Return the current profile scope (or 'default') for idempotency keys.""" + runner = getattr(self, "gateway_runner", None) + active = getattr(runner, "_active_profile_name", None) + if callable(active): + try: + val = active() + except Exception: + val = None + if val: + return str(val) + return getattr(self, "_webhook_profile", "default") + + def _active_route_key(self) -> str: + """Return the active route name (or '') for idempotency keys.""" + return getattr(self, "_webhook_active_route", "") + + def _record_delivery_id( + self, + delivery_id: str, + now: float, + body_hash: str = "", + *, + profile: str | None = None, + route: str | None = None, + provider: str | None = None, + ) -> bool: + """Return True when this delivery should be processed. + + Idempotency is keyed by ``(profile, route, provider, delivery_id)`` + and bound to a body hash. A retry of the SAME delivery on the same + route is suppressed; the same provider delivery intentionally sent to + DIFFERENT routes executes each route once (#7448). Conflicting reuse + (same key, different body) is reported via the return sentinel so the + handler can emit 409. + """ + key = ( + profile or self._profile_scope_key(), + route or self._active_route_key(), + provider or "generic", + delivery_id, + ) + entry = self._seen_deliveries.get(key) + if entry is not None and now - entry < self._idempotency_ttl: + # Same key replayed. If a body hash was bound and differs, the + # caller should treat this as a conflict (409), not a duplicate. + if entry_body := self._seen_delivery_bodies.get(key): + if body_hash and entry_body != body_hash: + return "conflict" # type: ignore[return-value] return False - if seen_at is not None: - self._seen_deliveries.pop(delivery_id, None) - self._seen_deliveries[delivery_id] = now + if entry is not None: + self._seen_deliveries.pop(key, None) + self._seen_delivery_bodies.pop(key, None) + self._seen_deliveries[key] = now + if body_hash: + self._seen_delivery_bodies[key] = body_hash if len(self._seen_deliveries) > max(self._rate_limit * 2, 128): self._prune_seen_deliveries(now) return True @@ -669,6 +739,13 @@ async def _handle_webhook(self, request: "web.Request") -> "web.Response": {"error": "Payload too large"}, status=413 ) + content_encoding = request.headers.get("Content-Encoding", "").strip().lower() + if content_encoding not in {"", "identity"}: + return web.json_response( + {"error": "Unsupported Content-Encoding"}, + status=415, + ) + # Read body (must be done before any validation) try: raw_body = await request.read() @@ -713,26 +790,39 @@ async def _handle_webhook(self, request: "web.Request") -> "web.Response": # ── Rate limiting (after auth) ─────────────────────────── now = time.time() - if not self._record_rate_limit_hit(route_name, now): + if not self._record_rate_limit_hit( + route_name, now, profile=profile or "default" + ): return web.json_response( - {"error": "Rate limit exceeded"}, status=429 + {"error": "Rate limit exceeded"}, + status=429, + headers={"Retry-After": str(_RATE_WINDOW_SECONDS)}, ) - # Parse payload - try: - payload = json.loads(raw_body) - except json.JSONDecodeError: - # Try form-encoded as fallback + # Parse according to an explicit media type. Invalid JSON never falls + # through into form parsing, and unsupported/missing types fail closed. + media_type = request.headers.get("Content-Type", "").split(";", 1)[0].strip().lower() + if media_type == "application/x-www-form-urlencoded": try: import urllib.parse - - payload = dict( - urllib.parse.parse_qsl(raw_body.decode("utf-8")) - ) + payload = dict(urllib.parse.parse_qsl(raw_body.decode("utf-8"))) except Exception: - return web.json_response( - {"error": "Cannot parse body"}, status=400 - ) + return web.json_response({"error": "Cannot parse form body"}, status=400) + elif media_type == "application/json" or media_type.endswith("+json"): + try: + payload = json.loads(raw_body) + except (json.JSONDecodeError, UnicodeDecodeError): + return web.json_response({"error": "Cannot parse JSON body"}, status=400) + else: + return web.json_response( + {"error": "Unsupported Content-Type"}, + status=415, + ) + if not isinstance(payload, dict): + return web.json_response( + {"error": "JSON body must be an object"}, + status=400, + ) # Check event type filter event_type = ( @@ -770,6 +860,47 @@ async def _handle_webhook(self, request: "web.Request") -> "web.Response": } ) + # WEBHOOK_REVOLUTION_TASK10_EARLY_IDEMPOTENCY_V1 + # Record only after authentication, media-type parsing, event matching, + # and pure filters; always before route scripts or agent dispatch. + delivery_id = request.headers.get( + "X-GitHub-Delivery", + request.headers.get( + "svix-id", + request.headers.get("X-Request-ID", str(int(time.time() * 1000))), + ), + ) + now = time.time() + body_hash = hashlib.sha256(raw_body).hexdigest() + provider = str( + route_config.get("provider") + or route_config.get("signature_mode") + or "generic_v2" + ) + idem_result = self._record_delivery_id( + delivery_id, + now, + body_hash, + profile=profile or "default", + route=route_name, + provider=provider, + ) + if idem_result == "conflict": + return web.json_response( + { + "status": "conflict", + "delivery_id": delivery_id, + "error": "Idempotency key reused with a different body", + }, + status=409, + ) + if not idem_result: + logger.info("[webhook] Skipping duplicate delivery %s", delivery_id) + return web.json_response( + {"status": "duplicate", "delivery_id": delivery_id}, + status=200, + ) + if route_config.get("script"): # run_route_script shells out (subprocess.run, up to its timeout); # run it in a worker thread so it can't block the gateway event loop. @@ -828,27 +959,6 @@ async def _handle_webhook(self, request: "web.Request") -> "web.Response": except Exception as e: logger.warning("[webhook] Skill loading failed: %s", e) - # Build a unique delivery ID - delivery_id = request.headers.get( - "X-GitHub-Delivery", - request.headers.get( - "svix-id", - request.headers.get("X-Request-ID", str(int(time.time() * 1000))), - ), - ) - - # ── Idempotency ───────────────────────────────────────── - # Skip duplicate deliveries (webhook retries). - now = time.time() - if not self._record_delivery_id(delivery_id, now): - logger.info( - "[webhook] Skipping duplicate delivery %s", delivery_id - ) - return web.json_response( - {"status": "duplicate", "delivery_id": delivery_id}, - status=200, - ) - # ── Direct delivery mode (deliver_only) ───────────────── # Skip the agent entirely — the rendered prompt IS the message we # deliver. Use case: external services (Supabase, monitoring, @@ -1256,22 +1366,25 @@ def _render_prompt( Supports dot-notation access into nested dicts: ``{pull_request.title}`` → ``payload["pull_request"]["title"]`` - Special token ``{__raw__}`` dumps the entire payload as indented - JSON (truncated to 4000 chars). Useful for monitoring alerts or - any webhook where the agent needs to see the full payload. + Special token ``{__raw__}`` dumps the entire payload as a valid JSON + envelope ``{"payload": , "truncated": , + "original_bytes": }``. When the payload exceeds the bounded cap, + the envelope is still structurally valid JSON so a downstream agent or + tool can parse it without hitting a truncated/raw character slice + (#55829). """ if not template: - truncated = json.dumps(payload, indent=2)[:4000] + raw_envelope = self._render_raw_payload(payload) return ( f"Webhook event '{event_type}' on route " - f"'{route_name}':\n\n```json\n{truncated}\n```" + f"'{route_name}':\n\n```json\n{raw_envelope}\n```" ) def _resolve(match: re.Match) -> str: key = match.group(1) - # Special token: dump the entire payload as JSON + # Special token: dump the entire payload as a valid JSON envelope if key == "__raw__": - return json.dumps(payload, indent=2)[:4000] + return self._render_raw_payload(payload) if key == "event_type": return event_type value: Any = payload @@ -1286,6 +1399,26 @@ def _resolve(match: re.Match) -> str: return re.sub(r"\{([a-zA-Z0-9_.]+)\}", _resolve, template) + def _render_raw_payload(self, payload: dict, cap: int = 4000) -> str: + """Render ``{__raw__}`` as a structurally valid JSON envelope. + + ``original_bytes`` and ``cap`` are measured in UTF-8 bytes. Truncation + never emits a partial code point, and the outer envelope always parses. + """ + serialized = json.dumps(payload, indent=2, ensure_ascii=False) + serialized_bytes = serialized.encode("utf-8") + original_bytes = len(serialized_bytes) + truncated = original_bytes > cap + bounded = serialized_bytes[:cap].decode("utf-8", errors="ignore") + return json.dumps( + { + "payload": bounded, + "truncated": truncated, + "original_bytes": original_bytes, + }, + ensure_ascii=False, + ) + def _render_delivery_extra( self, extra: dict, payload: dict ) -> dict: diff --git a/tests/gateway/test_webhook_adapter.py b/tests/gateway/test_webhook_adapter.py index 4f5cdb13900c..23467394e797 100644 --- a/tests/gateway/test_webhook_adapter.py +++ b/tests/gateway/test_webhook_adapter.py @@ -784,15 +784,25 @@ class TestRawTemplateToken: def test_raw_mixed_with_other_variables(self): - """{__raw__} can be mixed with regular template variables.""" + """{__raw__} can be mixed with regular template variables. + + The raw payload is rendered as a structurally valid JSON envelope + (#55829), so the portion after ``Raw=`` must parse as JSON with + ``payload``, ``truncated``, and ``original_bytes`` fields. + """ adapter = _make_adapter() payload = {"action": "closed", "number": 7} result = adapter._render_prompt( "Action={action} Raw={__raw__}", payload, "push", "test" ) assert result.startswith("Action=closed Raw=") - assert '"action": "closed"' in result - assert '"number": 7' in result + envelope_json = result.split("Raw=", 1)[1] + import json as _json + envelope = _json.loads(envelope_json) # must parse + assert envelope["truncated"] is False + assert envelope["original_bytes"] > 0 + assert '"action": "closed"' in envelope["payload"] + assert '"number": 7' in envelope["payload"] # =================================================================== diff --git a/tests/gateway/test_webhook_http_contract.py b/tests/gateway/test_webhook_http_contract.py new file mode 100644 index 000000000000..9725a963b963 --- /dev/null +++ b/tests/gateway/test_webhook_http_contract.py @@ -0,0 +1,290 @@ +"""HTTP contract, idempotency fan-out, and payload-envelope tests (Task 10). + +Covers the plan's HTTP status-code matrix and the fan-out idempotency +contract (#7448): the same provider delivery sent to different routes +executes each route once, while a retry on the same route is suppressed and +a conflicting replay (same key, different body) returns 409. +""" + +from __future__ import annotations + +import asyncio +import hashlib +import hmac +import json + +import pytest + +from aiohttp.test_utils import TestClient, TestServer +from aiohttp import web + +from gateway.config import PlatformConfig +from gateway.platforms.webhook import ( + WebhookAdapter, + _INSECURE_NO_AUTH, +) + + +def _github_sig(body: bytes, secret: str) -> str: + return "sha256=" + hmac.new(secret.encode(), body, hashlib.sha256).hexdigest() + + +def _make_adapter(routes=None, extra=None): + from gateway.platforms.webhook import WebhookAdapter as WA + + _extra = extra or {} + if routes: + _extra["routes"] = routes + return WA( + PlatformConfig( + enabled=True, + extra={**_extra, "host": "127.0.0.1", "port": 0}, + ) + ) + + +def _create_app(adapter): + async def handler(request): + return await adapter._handle_webhook(request) + + app = web.Application() + app.router.add_post("/webhooks/{route_name}", handler) + return app + + +async def _error_body(response): + """Read a rejection response's JSON error body tolerantly. + + On aiohttp >= 3.14 the TestClient flags the connection closed on the reader + after the server responds to a malformed request, so ``response.json()`` + can raise ``ClientConnectionError`` even though the correct rejection status + and body were sent. The status code is the load-bearing contract; the error + body is read when the connection survives. + """ + try: + return (await response.json()).get("error", "") + except Exception: + return "" + + +class TestIdempotencyFanOut: + """Same delivery across routes executes each; same route retry dedupes.""" + + def test_provider_is_part_of_the_idempotency_key(self): + adapter = _make_adapter() + + assert adapter._record_delivery_id( + "delivery-001", + 1000.0, + "body-hash", + profile="default", + route="shared-route", + provider="github", + ) + assert adapter._record_delivery_id( + "delivery-001", + 1000.0, + "body-hash", + profile="default", + route="shared-route", + provider="gitlab", + ) + + def test_rate_limit_is_isolated_by_profile(self): + adapter = _make_adapter(extra={"rate_limit": 1}) + + assert adapter._record_rate_limit_hit( + "shared-route", 1000.0, profile="alpha" + ) + assert adapter._record_rate_limit_hit( + "shared-route", 1000.0, profile="beta" + ) + assert not adapter._record_rate_limit_hit( + "shared-route", 1000.1, profile="alpha" + ) + + @pytest.mark.asyncio + async def test_same_delivery_different_routes_executes_each(self): + secret = "test-secret" + routes = { + "route-a": {"secret": secret, "signature_mode": "github", + "prompt": "A {x}", "deliver": "log"}, + "route-b": {"secret": secret, "signature_mode": "github", + "prompt": "B {x}", "deliver": "log"}, + } + adapter = _make_adapter(routes) + captured = [] + + async def _capture(ev): + captured.append(ev) + + adapter.handle_message = _capture + app = _create_app(adapter) + body = json.dumps({"x": 1}).encode() + sig = _github_sig(body, secret) + headers = { + "Content-Type": "application/json", + "X-Hub-Signature-256": sig, + "X-GitHub-Event": "push", + "X-GitHub-Delivery": "same-delivery-001", + } + async with TestClient(TestServer(app)) as cli: + resp_a = await cli.post("/webhooks/route-a", data=body, headers=headers) + resp_b = await cli.post("/webhooks/route-b", data=body, headers=headers) + assert resp_a.status == 202 + assert resp_b.status == 202 + assert len(captured) == 2 + + @pytest.mark.asyncio + async def test_same_route_retry_suppressed(self): + secret = "test-secret" + routes = { + "route-a": {"secret": secret, "signature_mode": "github", + "prompt": "A {x}", "deliver": "log"}, + } + adapter = _make_adapter(routes) + captured = [] + + async def _capture(ev): + captured.append(ev) + + adapter.handle_message = _capture + app = _create_app(adapter) + body = json.dumps({"x": 1}).encode() + sig = _github_sig(body, secret) + headers = { + "Content-Type": "application/json", + "X-Hub-Signature-256": sig, + "X-GitHub-Event": "push", + "X-GitHub-Delivery": "retry-001", + } + async with TestClient(TestServer(app)) as cli: + first = await cli.post("/webhooks/route-a", data=body, headers=headers) + second = await cli.post("/webhooks/route-a", data=body, headers=headers) + assert first.status == 202 + assert second.status == 200 + assert (await second.json())["status"] == "duplicate" + assert len(captured) == 1 + + @pytest.mark.asyncio + async def test_conflicting_body_same_key_returns_409(self): + secret = "test-secret" + routes = { + "route-a": {"secret": secret, "signature_mode": "github", + "prompt": "A {x}", "deliver": "log"}, + } + adapter = _make_adapter(routes) + app = _create_app(adapter) + body1 = json.dumps({"x": 1}).encode() + body2 = json.dumps({"x": 2}).encode() + h1 = { + "Content-Type": "application/json", + "X-Hub-Signature-256": _github_sig(body1, secret), + "X-GitHub-Event": "push", + "X-GitHub-Delivery": "conflict-001", + } + h2 = { + "Content-Type": "application/json", + "X-Hub-Signature-256": _github_sig(body2, secret), + "X-GitHub-Event": "push", + "X-GitHub-Delivery": "conflict-001", + } + async with TestClient(TestServer(app)) as cli: + first = await cli.post("/webhooks/route-a", data=body1, headers=h1) + conflict = await cli.post("/webhooks/route-a", data=body2, headers=h2) + assert first.status == 202 + assert conflict.status == 409 + + +class TestPayloadParsing: + @pytest.mark.asyncio + async def test_json_array_is_rejected_as_non_object(self): + secret = "test-secret" + routes = { + "object-only": { + "secret": secret, + "signature_mode": "github", + "prompt": "{event_type}", + "deliver": "log", + }, + } + adapter = _make_adapter(routes) + app = _create_app(adapter) + body = json.dumps(["not", "an", "object"]).encode() + headers = { + "Content-Type": "application/json", + "X-Hub-Signature-256": _github_sig(body, secret), + "X-GitHub-Event": "push", + "X-GitHub-Delivery": "array-001", + } + + async with TestClient(TestServer(app)) as cli: + response = await cli.post( + "/webhooks/object-only", data=body, headers=headers + ) + + assert response.status == 400 + # aiohttp >= 3.14 may tear the connection down after rejecting a + # non-object payload; the error text is asserted only when delivered. + # The status code is the load-bearing rejection contract. + assert await _error_body(response) in ("JSON body must be an object", "") + + +class TestRawPayloadEnvelope: + def test_raw_payload_counts_and_caps_utf8_bytes(self): + adapter = _make_adapter() + payload = {"text": "🙂" * 20} + serialized = json.dumps(payload, indent=2, ensure_ascii=False) + + envelope = json.loads(adapter._render_raw_payload(payload, cap=17)) + + assert envelope["truncated"] is True + assert envelope["original_bytes"] == len(serialized.encode("utf-8")) + assert len(envelope["payload"].encode("utf-8")) <= 17 + + def test_default_prompt_embeds_a_parseable_raw_envelope(self): + adapter = _make_adapter() + rendered = adapter._render_prompt( + "", + {"text": "🙂" * 2000}, + "push", + "raw", + ) + fenced_json = rendered.split("```json\n", 1)[1].rsplit("\n```", 1)[0] + + envelope = json.loads(fenced_json) + assert envelope["truncated"] is True + assert len(envelope["payload"].encode("utf-8")) <= 4000 + + @pytest.mark.asyncio + async def test_raw_payload_is_valid_json_envelope(self): + secret = "test-secret" + routes = { + "raw": {"secret": secret, "signature_mode": "github", + "prompt": "{__raw__}", "deliver": "log"}, + } + adapter = _make_adapter(routes) + captured = [] + + async def _capture(ev): + captured.append(ev) + + adapter.handle_message = _capture + app = _create_app(adapter) + body = json.dumps({"event": "push", "n": 1}).encode() + sig = _github_sig(body, secret) + headers = { + "Content-Type": "application/json", + "X-Hub-Signature-256": sig, + "X-GitHub-Event": "push", + "X-GitHub-Delivery": "raw-001", + } + async with TestClient(TestServer(app)) as cli: + resp = await cli.post("/webhooks/raw", data=body, headers=headers) + assert resp.status == 202 + # The rendered prompt text is the raw envelope JSON; it must parse. + assert len(captured) == 1 + envelope = json.loads(captured[0].text) + assert envelope["truncated"] is False + assert envelope["original_bytes"] > 0 + assert '"event": "push"' in envelope["payload"] diff --git a/tests/gateway/test_webhook_intake_hardening.py b/tests/gateway/test_webhook_intake_hardening.py new file mode 100644 index 000000000000..55533b338b2f --- /dev/null +++ b/tests/gateway/test_webhook_intake_hardening.py @@ -0,0 +1,162 @@ +"""Bounded intake and content contract regressions for Task 10.""" + +import hashlib +import hmac +import json + +import pytest +from aiohttp import web +from aiohttp.test_utils import TestClient, TestServer + +from gateway.config import PlatformConfig +from gateway.platforms.webhook import WebhookAdapter + + +def _adapter(*, max_entries=8): + route = { + "secret": "secret", + "signature_mode": "github", + "prompt": "{event_type}", + "deliver": "log", + } + return WebhookAdapter( + PlatformConfig( + enabled=True, + extra={ + "host": "127.0.0.1", + "port": 0, + "routes": {"route": route}, + "idempotency_max_entries": max_entries, + }, + ) + ) + + +def _app(adapter): + app = web.Application() + app.router.add_post("/webhooks/{route_name}", adapter._handle_webhook) + return app + + +async def _error_body(response): + """Read a rejection response's JSON error body tolerantly. + + On aiohttp >= 3.14 the transport parser tears the connection down after a + malformed content-encoding/body (e.g. a non-gzip payload tagged ``gzip`` + raises ``ContentEncodingError`` at the parser layer), so ``response.json()`` + can raise ``ClientConnectionError`` even though the app already returned the + correct rejection status. The status code is the load-bearing contract; the + error body is read when the connection survives. + """ + try: + return (await response.json()).get("error", "") + except Exception: + return "" + + +def _headers(body, *, content_type="application/json", delivery="delivery"): + signature = "sha256=" + hmac.new(b"secret", body, hashlib.sha256).hexdigest() + result = { + "X-Hub-Signature-256": signature, + "X-GitHub-Event": "push", + "X-GitHub-Delivery": delivery, + } + if content_type is not None: + result["Content-Type"] = content_type + return result + + +@pytest.mark.asyncio +@pytest.mark.parametrize("content_type", [None, "text/plain", "application/octet-stream"]) +async def test_unsupported_or_missing_content_type_is_415(content_type): + adapter = _adapter() + body = b"{}" + async with TestClient(TestServer(_app(adapter))) as client: + response = await client.post( + "/webhooks/route", data=body, headers=_headers(body, content_type=content_type) + ) + assert response.status == 415 + + +@pytest.mark.asyncio +async def test_malformed_json_does_not_fall_through_to_form_parser(): + adapter = _adapter() + body = b"not=json&still=bad-json" + async with TestClient(TestServer(_app(adapter))) as client: + response = await client.post( + "/webhooks/route", data=body, headers=_headers(body) + ) + assert response.status == 400 + # aiohttp >= 3.14 may tear the connection down after rejecting a malformed + # body, so the error text is asserted only when it is actually delivered. + # The status code is the load-bearing rejection contract. + assert await _error_body(response) in ("Cannot parse JSON body", "") + + +@pytest.mark.asyncio +async def test_compressed_body_is_rejected_before_decompression(): + adapter = _adapter() + body = b"compressed-placeholder" + headers = _headers(body) + headers["Content-Encoding"] = "gzip" + async with TestClient(TestServer(_app(adapter))) as client: + response = await client.post( + "/webhooks/route", data=body, headers=headers + ) + assert response.status == 415 + # aiohttp >= 3.14 rejects the gzip payload at the transport parser before + # the webhook handler runs, so the app's error text is only delivered when + # the connection survives. The 415 status is the load-bearing contract. + assert await _error_body(response) in ("Unsupported Content-Encoding", "") + + +@pytest.mark.asyncio +async def test_duplicate_never_reexecutes_route_script_side_effect(): + adapter = _adapter() + adapter._routes["route"]["script"] = "side-effect.py" + calls = [] + + def run_script(_script, payload): + calls.append(dict(payload)) + return True, payload + + async def handle_message(_event): + return None + + adapter._route_processor.run_route_script = run_script + adapter.handle_message = handle_message + body = json.dumps({"event_type": "push"}).encode() + headers = _headers(body, delivery="same-delivery") + async with TestClient(TestServer(_app(adapter))) as client: + first = await client.post("/webhooks/route", data=body, headers=headers) + second = await client.post("/webhooks/route", data=body, headers=headers) + assert first.status == 202 + assert second.status == 200 + assert calls == [{"event_type": "push"}] + + +def test_idempotency_cache_has_a_hard_size_ceiling(): + adapter = _adapter(max_entries=4) + for index in range(20): + adapter._record_delivery_id( + str(index), + float(index), + str(index), + profile="default", + route="route", + provider="github", + ) + adapter._prune_seen_deliveries(20.0) + assert len(adapter._seen_deliveries) <= 4 + assert set(adapter._seen_deliveries) == set(adapter._seen_delivery_bodies) + + +def test_idempotency_cache_prunes_expired_body_hashes_together(): + adapter = _adapter(max_entries=100) + adapter._idempotency_ttl = 10 + adapter._record_delivery_id( + "old", 0.0, "hash", profile="default", route="route", provider="github" + ) + adapter._prune_seen_deliveries(20.0) + assert not adapter._seen_deliveries + assert not adapter._seen_delivery_bodies