diff --git a/CHANGELOG.md b/CHANGELOG.md index 7a6c6648..b9f30074 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -2,6 +2,21 @@ ## Unreleased +### Cloud Security — code-scan pushes retry when the service is busy + +- `CloudSec.ingest_code_results` (and so `cloudsec code ingest` and + `cloudsec code scan --ingest`) now re-sends a push refused with HTTP 429 — + the organization already has as many pushes in progress as it may + (`error_code: "ingest_busy"`), or the request quota is spent. Nothing is + recorded for a refused push, so the re-send is safe. Each wait honours the + response's `Retry-After` (up to 120 seconds) as a floor, backs off + exponentially without one, and adds random jitter; at most 5 re-sends and + 10 minutes of waiting in total. `busy_retries=0` restores the old behaviour. +- `RateLimitError.retry_after` is now filled from the response's `Retry-After` + header when the client raises on the first 429 (without `--retry`, or with + `Client.request(..., retry_quota_errors=False)`, which overrides the client's + `--retry` setting for one call). + ### Cloud Security — SBOM route - `CloudSec.get_code_sbom` / `download_code_sbom` (and `cloudsec code sbom`) diff --git a/doc/cli/cloud-security.md b/doc/cli/cloud-security.md index 4e4f9dd9..d7b864dc 100644 --- a/doc/cli/cloud-security.md +++ b/doc/cli/cloud-security.md @@ -308,6 +308,8 @@ limacharlie cloudsec code autofix # open the upgrade PR limacharlie cloudsec code ingest --repo acme/api --source sarif --file report.sarif ``` +`code ingest` (and `code scan --ingest`) retries a push the service answers with HTTP 429, which means the organization already has as many pushes in progress as it may, or the request quota is spent. It waits at least the response's `Retry-After`, adds random jitter so a CI fan-out that was refused together does not come back together, and gives up after 5 retries or 10 minutes of waiting, exiting with the rate-limit error. A refused push recorded nothing, so the retry is safe. + `code repos` reports `scan_status` as `scanned`, `partial` or `unknown`. `partial` means the scan tripped a limit, so the finding set is INCOMPLETE — not a clean bill. `unknown` means this view has no scan state and says so rather than guessing; `code status` is the authoritative view of the run. `code capabilities` covers **GitHub connections only** — a GitLab or Bitbucket connection scans with its own read-only token and has no write plane to detect, so it never appears, not even as `unknown`. Use `provider manifest` for those. A capability of `available` means the control MAY be offered, not that anything fires on its own. diff --git a/limacharlie/client.py b/limacharlie/client.py index 17437f02..be11fe65 100644 --- a/limacharlie/client.py +++ b/limacharlie/client.py @@ -72,6 +72,38 @@ def _root_from_env(name: str, default: str) -> str: HTTP_TOO_MANY_REQUESTS = 429 HTTP_GATEWAY_TIMEOUT = 504 + +def parse_retry_after(value: str | None, now: float | None = None) -> int | None: + """Parse an HTTP ``Retry-After`` header into whole seconds from now. + + Both forms RFC 9110 allows are read: delay-seconds (``"30"``) and an + HTTP-date. A date in the past is 0. Anything else — absent, empty, + negative, garbage — is ``None``, meaning "the server did not say", which + is different from "retry now". + """ + if value is None: + return None + value = value.strip() + if not value: + return None + if value.isascii() and value.isdigit(): + # Clamped here so a hostile or broken header can neither overflow the int + # parser (Python refuses very long digit strings) nor mean "forever". + if len(value) > 9: + return 999_999_999 + return int(value) + try: + from email.utils import parsedate_to_datetime + + when = parsedate_to_datetime(value) + except (TypeError, ValueError, IndexError, OverflowError): + return None + if when is None or when.tzinfo is None: + return None + current = time.time() if now is None else now + return max(0, int(when.timestamp() - current + 0.999)) + + # Default limit for response body in debug output. Bodies longer than # this are truncated with a "[truncated]" marker. Use debug_full_response=True # on the Client (or --debug-full on the CLI) to disable truncation. @@ -614,13 +646,16 @@ def _call_jwt_endpoint(self, auth_data: dict[str, Any]) -> str: def _rest_call(self, url: str, verb: str, params: dict[str, Any] | None = None, alt_root: str | None = None, query_params: dict[str, Any] | list[tuple[str, str]] | None = None, raw_body: bytes | None = None, content_type: str | None = None, is_no_auth: bool = False, timeout: int | None = None, - extra_headers: dict[str, str] | None = None, raw_response: bool = False) -> tuple[int, Any]: + extra_headers: dict[str, str] | None = None, raw_response: bool = False, + response_headers: dict[str, str] | None = None) -> tuple[int, Any]: """Make a single HTTP request to the API. Args: raw_response: Return the response body as decoded text instead of parsed JSON (for non-JSON endpoints like CSV exports). Error (non-2xx) bodies are still parsed as JSON when possible. + response_headers: When given, filled with the headers of an error + (non-2xx) response, so the caller can read e.g. ``Retry-After``. Returns: tuple: (status_code, response_data) @@ -707,8 +742,11 @@ def _rest_call(self, url: str, verb: str, params: dict[str, Any] | None = None, resp = json.loads(error_str) except Exception: resp = error_str - resp_headers = e.headers.items() if hasattr(e, "headers") else [] - self._debug_response(e.code, list(resp_headers), error_str) + resp_headers = list(e.headers.items()) if getattr(e, "headers", None) is not None else [] + if response_headers is not None: + for h_name, h_val in resp_headers: + response_headers[h_name.lower()] = h_val + self._debug_response(e.code, resp_headers, error_str) return (e.code, resp) except ssl.SSLError as e: @@ -718,7 +756,7 @@ def _rest_call(self, url: str, verb: str, params: dict[str, Any] | None = None, def request(self, verb: str, url: str, params: dict[str, Any] | None = None, alt_root: str | None = None, query_params: dict[str, Any] | list[tuple[str, str]] | None = None, raw_body: bytes | None = None, content_type: str | None = None, is_no_auth: bool = False, max_retries: int = 3, timeout: int | None = None, extra_headers: dict[str, str] | None = None, - raw_response: bool = False) -> Any: + raw_response: bool = False, retry_quota_errors: bool | None = None) -> Any: """Make an API request with retry logic and JWT management. Args: @@ -735,13 +773,18 @@ def request(self, verb: str, url: str, params: dict[str, Any] | None = None, alt extra_headers: Additional HTTP headers to include. raw_response: Return the response body as decoded text instead of parsed JSON (for non-JSON endpoints like CSV exports). + retry_quota_errors: Override the client's ``is_retry_quota_errors`` + for this call. ``False`` raises :class:`RateLimitError` on the + first 429, for a caller that runs its own backoff. Returns: dict: Parsed JSON response (or str when raw_response is set). Raises: AuthenticationError: on 401 after JWT refresh. - RateLimitError: on 429 when retry is not enabled. + RateLimitError: on 429 when retry is not enabled. Its + ``retry_after`` carries the response's ``Retry-After`` in + seconds, or ``None`` when the response had none. ApiError: on other non-200 responses after retries. """ has_auth_refreshed = False @@ -753,16 +796,20 @@ def request(self, verb: str, url: str, params: dict[str, Any] | None = None, alt else: self.refresh_jwt() + if retry_quota_errors is None: + retry_quota_errors = self._is_retry_quota_errors + retries = 0 while retries < max_retries: retries += 1 + error_headers: dict[str, str] = {} code, data = self._rest_call( url, verb, params=params, alt_root=alt_root, query_params=query_params, raw_body=raw_body, content_type=content_type, is_no_auth=is_no_auth, timeout=timeout, extra_headers=extra_headers, - raw_response=raw_response, + raw_response=raw_response, response_headers=error_headers, ) if code == HTTP_OK: @@ -798,13 +845,14 @@ def request(self, verb: str, url: str, params: dict[str, Any] | None = None, alt break if code == HTTP_TOO_MANY_REQUESTS: - if self._is_retry_quota_errors: + if retry_quota_errors: wait = min(10 * (2 ** (retries - 1)), 60) self._debug(f"Rate limited, waiting {wait}s before retry...") time.sleep(wait) continue raise RateLimitError( f"Rate limit exceeded: {data}", + retry_after=parse_retry_after(error_headers.get("retry-after")), code=code, ) diff --git a/limacharlie/sdk/cloudsec.py b/limacharlie/sdk/cloudsec.py index 33c6f958..5c707802 100644 --- a/limacharlie/sdk/cloudsec.py +++ b/limacharlie/sdk/cloudsec.py @@ -66,12 +66,14 @@ import base64 import json +import random import re +import time from typing import Any, TYPE_CHECKING from urllib.parse import quote as _quote from urllib.request import urlopen as _urlopen -from ..errors import AuthenticationError +from ..errors import AuthenticationError, RateLimitError if TYPE_CHECKING: from .organization import Organization @@ -124,6 +126,41 @@ def _query_pairs(**params: Any) -> list[tuple[str, str]]: return pairs +# How a code-scan push backs off when it is told to come back later (429). +# +# A 429 on the ingest means the push was NOT processed: either the org already has as +# many pushes running and queued on the ingest service as it may ("ingest_busy", which +# carries a Retry-After), or the caller's per-identity request quota is spent. Both are +# refused before the document is processed, so nothing was recorded and re-sending the +# same push is safe. What matters is WHEN: a CI fan-out that is refused +# together must not come back together, so every wait is jittered, and the whole thing +# is bounded so a job that cannot get in fails instead of hanging. +INGEST_BUSY_RETRIES = 5 +# The first backoff when the server gives no Retry-After, doubled per attempt up to the cap. +_INGEST_BUSY_BASE_S = 5.0 +_INGEST_BUSY_CAP_S = 60.0 +# A Retry-After past this is clamped: one header must not park a CI job for an hour. +_INGEST_RETRY_AFTER_MAX_S = 120.0 +# Jitter adds up to this fraction of the wait on top of it, never below it, so a +# Retry-After is always honoured as a floor. +_INGEST_JITTER = 0.5 +# The most a push waits in total across its retries. Past it, the 429 is raised. +_INGEST_BUSY_BUDGET_S = 600.0 + + +def _ingest_busy_delay(attempt: int, retry_after: int | None, rand: Any = random) -> float: + """Seconds to wait before retry number ``attempt + 1`` of a refused push. + + The exponential backoff, raised to the server's ``Retry-After`` when it + asks for longer (clamped to :data:`_INGEST_RETRY_AFTER_MAX_S`), plus up to + :data:`_INGEST_JITTER` of that on top. + """ + delay = min(_INGEST_BUSY_CAP_S, _INGEST_BUSY_BASE_S * (2 ** attempt)) + if retry_after is not None: + delay = max(delay, min(float(retry_after), _INGEST_RETRY_AFTER_MAX_S)) + return delay + rand.uniform(0.0, delay * _INGEST_JITTER) + + def _inventory_account_selector( account_empty: bool | None, account_unscoped: bool | None, @@ -772,14 +809,20 @@ def _post( query_params: list[tuple[str, str]] | None = None, *, raw_response: bool = False, + raw_body: bytes | None = None, + retry_quota_errors: bool | None = None, ) -> Any: + kwargs: dict[str, Any] = {} + if retry_quota_errors is not None: + kwargs["retry_quota_errors"] = retry_quota_errors return self._org.client.request( "POST", f"cloudsec/{self.oid}/{path}", query_params=query_params or None, - raw_body=json.dumps(body).encode(), + raw_body=raw_body if raw_body is not None else json.dumps(body).encode(), content_type="application/json", raw_response=raw_response, + **kwargs, ) # ------------------------------------------------------------------ @@ -2873,6 +2916,7 @@ def ingest_code_results( ref: str | None = None, default_branch: str | None = None, provider: str | None = None, + busy_retries: int = INGEST_BUSY_RETRIES, ) -> dict[str, Any]: """Push results your own pipeline produced for one repository. @@ -2906,6 +2950,19 @@ def ingest_code_results( sending for a repository LimaCharlie does not collect — nothing else can state it there, and it is left unset rather than guessed when you do not know it. + busy_retries: how many times to re-send the push when it is + refused with 429 — the organization already has as many + pushes in progress as it may (``error_code: "ingest_busy"``), + or the request quota is spent. Nothing was recorded for a + refused push, so re-sending it is safe. Each wait honours the + response's ``Retry-After`` (up to 120 seconds) as a floor, backs + off exponentially otherwise, and adds random jitter so a + refused CI fan-out does not come back as one burst; the waits + total at most 10 minutes. ``0`` raises on the first 429. + + Raises: + RateLimitError: the push was still refused after ``busy_retries`` + re-sends, or the next wait would pass the 10-minute budget. Returns: ``{"result": {...}}`` — what landed: ``findings``, the SoR @@ -2964,7 +3021,39 @@ def ingest_code_results( body["default_branch"] = default_branch if provider is not None: body["provider"] = provider - return self._post("code/ingest", body) + # Serialized once: the document can be tens of megabytes, and a retry re-sends + # the same bytes. + raw = json.dumps(body).encode() + waited = 0.0 + attempt = 0 + while True: + try: + # The client's own 429 retry is off here: it would neither honour the + # Retry-After nor jitter, and stacked under this loop it would multiply + # the attempts. + return self._post("code/ingest", body, raw_body=raw, retry_quota_errors=False) + except RateLimitError as e: + delay = _ingest_busy_delay(attempt, e.retry_after) + if attempt >= busy_retries or waited + delay > _INGEST_BUSY_BUDGET_S: + if attempt == 0: + raise + # Not the default "use --retry" hint: --retry does not change this + # path, which has already backed off. + raise RateLimitError( + e.raw_message, + retry_after=e.retry_after, + suggestion=( + f"The push was refused {attempt + 1} times over {waited:.0f}s. " + "The organization has too many pushes in progress; push fewer " + "at once, or retry this one later."), + code=e.status_code, + ) from e + self._org.client._debug( + f"code ingest refused (429), retrying in {delay:.1f}s " + f"({attempt + 1}/{busy_retries})") + time.sleep(delay) + waited += delay + attempt += 1 def get_code_capabilities(self, *, repo: str | None = None) -> dict[str, Any]: """What each connected source-control organization may actually DO diff --git a/tests/unit/test_code_ingest_busy_retry.py b/tests/unit/test_code_ingest_busy_retry.py new file mode 100644 index 00000000..39e217ee --- /dev/null +++ b/tests/unit/test_code_ingest_busy_retry.py @@ -0,0 +1,224 @@ +"""A code-scan push that is told to come back later (429) comes back later. + +These drive the REAL client (only ``urlopen`` and the sleep are replaced), so +the ``Retry-After`` header travels the same way it does against the API: out +of the HTTP error, through ``Client.request`` into ``RateLimitError``, and into +the push's backoff. +""" + +import io +import json +import os +from http.client import HTTPMessage +from unittest.mock import MagicMock, patch +from urllib.error import HTTPError + +import pytest + +from limacharlie.client import Client, parse_retry_after +from limacharlie.errors import ApiError, RateLimitError +from limacharlie.sdk import cloudsec as cloudsec_mod +from limacharlie.sdk.cloudsec import CloudSec, _ingest_busy_delay + + +OID = "11111111-2222-3333-4444-555555555555" +BUSY_BODY = json.dumps({ + "error": "code_ingest: ingest_busy: this organization already has 2 pushes running " + "and more waiting on this replica; retry this push after its earlier pushes finish", + "error_code": "ingest_busy", + "retry": True, +}).encode() + + +@pytest.fixture(autouse=True) +def _isolate_caches(monkeypatch, tmp_path): + from limacharlie.config import _reset_config_cache + from limacharlie.jwt_cache import _reset_cache_disabled + from limacharlie.paths import _reset_path_cache + + config_dir = str(tmp_path / "lc_config") + os.makedirs(config_dir, exist_ok=True) + monkeypatch.setenv("LC_CONFIG_DIR", config_dir) + for var in ("LC_CREDS_FILE", "LC_LEGACY_CONFIG", "LC_EPHEMERAL_CREDS", "LC_NO_JWT_CACHE"): + monkeypatch.delenv(var, raising=False) + _reset_path_cache() + _reset_config_cache() + _reset_cache_disabled() + yield + _reset_path_cache() + _reset_config_cache() + _reset_cache_disabled() + + +def _http_429(retry_after=None, body=BUSY_BODY): + headers = HTTPMessage() + headers["Content-Type"] = "application/json" + if retry_after is not None: + headers["Retry-After"] = retry_after + return HTTPError("https://api.limacharlie.io/v1/x", 429, "Too Many Requests", headers, io.BytesIO(body)) + + +def _http_400(): + return HTTPError("https://api.limacharlie.io/v1/x", 400, "Bad Request", HTTPMessage(), + io.BytesIO(b'{"error": "code_ingest: repo_not_in_policy_scope: no"}')) + + +def _ok(payload): + resp = MagicMock() + resp.read.return_value = json.dumps(payload).encode() + resp.getheaders.return_value = [] + return resp + + +def _cloudsec(is_retry_quota_errors=False): + client = Client(oid=OID, jwt="jwt", is_retry_quota_errors=is_retry_quota_errors) + org = MagicMock() + org.oid = OID + org.client = client + return CloudSec(org) + + +def _bodies(mock_urlopen): + return [call.args[0].data for call in mock_urlopen.call_args_list] + + +class TestRetryAfterReachesTheError: + @patch("limacharlie.client.urlopen") + def test_429_carries_its_retry_after(self, mock_urlopen): + mock_urlopen.side_effect = [_http_429("30")] + client = Client(oid=OID, jwt="jwt") + with pytest.raises(RateLimitError) as e: + client.request("POST", "x") + assert e.value.retry_after == 30 + + @patch("limacharlie.client.urlopen") + def test_429_without_retry_after_is_none_not_zero(self, mock_urlopen): + mock_urlopen.side_effect = [_http_429(None)] + client = Client(oid=OID, jwt="jwt") + with pytest.raises(RateLimitError) as e: + client.request("POST", "x") + assert e.value.retry_after is None + + @patch("limacharlie.client.time") + @patch("limacharlie.client.urlopen") + def test_per_call_override_beats_the_client_retry_flag(self, mock_urlopen, mock_time): + """A caller running its own backoff gets the 429 at once even on a --retry client.""" + mock_urlopen.side_effect = [_http_429("30"), _ok({"ok": True})] + client = Client(oid=OID, jwt="jwt", is_retry_quota_errors=True) + with pytest.raises(RateLimitError): + client.request("POST", "x", retry_quota_errors=False) + assert mock_urlopen.call_count == 1 + mock_time.sleep.assert_not_called() + + def test_parse_retry_after_forms(self): + assert parse_retry_after("30") == 30 + assert parse_retry_after(" 7 ") == 7 + # An HTTP-date 90 seconds after "now". + assert parse_retry_after("Wed, 21 Oct 2015 07:29:30 GMT", now=1445412480.0) == 90 + assert parse_retry_after("Wed, 21 Oct 2015 07:28:00 GMT", now=1445412580.0) == 0 + for garbage in (None, "", "-5", "soon", "1.5e3", "\u00b2", "\u0663\u0660"): + assert parse_retry_after(garbage) is None + # Past Python's integer-string limit: clamped, not a ValueError out of request(). + assert parse_retry_after("9" * 5000) == 999_999_999 + + +class TestIngestRetriesBusy: + @patch.object(cloudsec_mod.time, "sleep") + @patch("limacharlie.client.urlopen") + def test_busy_then_accepted(self, mock_urlopen, mock_sleep): + mock_urlopen.side_effect = [_http_429("30"), _http_429("30"), _ok({"result": {"findings": 3}})] + out = _cloudsec().ingest_code_results("acme/api", "report", b"\x1f\x8bdoc", commit="c0ffee") + assert out == {"result": {"findings": 3}} + assert mock_urlopen.call_count == 3 + # The same push each time. + bodies = _bodies(mock_urlopen) + assert bodies[0] == bodies[1] == bodies[2] + assert json.loads(bodies[0])["repo"] == "acme/api" + # Each wait honours Retry-After as a floor, with at most 50% jitter on top. + waits = [c.args[0] for c in mock_sleep.call_args_list] + assert len(waits) == 2 + for w in waits: + assert 30.0 <= w <= 45.0 + + @patch.object(cloudsec_mod.time, "sleep") + @patch("limacharlie.client.urlopen") + def test_gives_up_after_busy_retries(self, mock_urlopen, mock_sleep): + mock_urlopen.side_effect = [_http_429("1") for _ in range(10)] + with pytest.raises(RateLimitError) as e: + _cloudsec().ingest_code_results("acme/api", "sarif", {"runs": []}, busy_retries=3) + assert mock_urlopen.call_count == 4 + assert mock_sleep.call_count == 3 + assert e.value.retry_after == 1 + # The hint names what happened, not a --retry flag that does not apply here. + assert "--retry" not in str(e.value) + assert "refused 4 times" in str(e.value) + + @patch.object(cloudsec_mod.time, "sleep") + @patch("limacharlie.client.urlopen") + def test_zero_retries_raises_at_once(self, mock_urlopen, mock_sleep): + mock_urlopen.side_effect = [_http_429("30"), _ok({"result": {}})] + with pytest.raises(RateLimitError): + _cloudsec().ingest_code_results("acme/api", "sarif", {"runs": []}, busy_retries=0) + assert mock_urlopen.call_count == 1 + mock_sleep.assert_not_called() + + @patch.object(cloudsec_mod.time, "sleep") + @patch("limacharlie.client.urlopen") + def test_total_wait_is_bounded(self, mock_urlopen, mock_sleep): + """A server that keeps asking for the maximum cannot hold a CI job past the budget.""" + mock_urlopen.side_effect = [_http_429("3600") for _ in range(50)] + with pytest.raises(RateLimitError): + _cloudsec().ingest_code_results("acme/api", "sarif", {"runs": []}, busy_retries=40) + waits = [c.args[0] for c in mock_sleep.call_args_list] + assert waits, "it never retried" + assert sum(waits) <= cloudsec_mod._INGEST_BUSY_BUDGET_S + # Each Retry-After was clamped, not obeyed for an hour. + assert max(waits) <= cloudsec_mod._INGEST_RETRY_AFTER_MAX_S * (1 + cloudsec_mod._INGEST_JITTER) + assert mock_urlopen.call_count < 41 + + @patch.object(cloudsec_mod.time, "sleep") + @patch("limacharlie.client.time") + @patch("limacharlie.client.urlopen") + def test_a_retry_client_does_not_stack_its_own_429_loop(self, mock_urlopen, mock_client_time, mock_sleep): + """With the CLI's --retry on, the push still makes exactly one request per attempt.""" + mock_urlopen.side_effect = [_http_429("30"), _ok({"result": {}})] + _cloudsec(is_retry_quota_errors=True).ingest_code_results("acme/api", "sarif", {"runs": []}) + assert mock_urlopen.call_count == 2 + mock_client_time.sleep.assert_not_called() + assert mock_sleep.call_count == 1 + + @patch.object(cloudsec_mod.time, "sleep") + @patch("limacharlie.client.urlopen") + def test_a_deterministic_refusal_is_not_retried(self, mock_urlopen, mock_sleep): + mock_urlopen.side_effect = [_http_400(), _ok({"result": {}})] + with pytest.raises(ApiError): + _cloudsec().ingest_code_results("acme/api", "sarif", {"runs": []}) + assert mock_urlopen.call_count == 1 + mock_sleep.assert_not_called() + + +class TestBackoffShape: + class _Rand: + def __init__(self, pick): + self.pick = pick + + def uniform(self, lo, hi): + return lo if self.pick == "lo" else hi + + def test_without_retry_after_it_backs_off_exponentially_to_the_cap(self): + lo = self._Rand("lo") + assert [_ingest_busy_delay(n, None, lo) for n in range(6)] == [5.0, 10.0, 20.0, 40.0, 60.0, 60.0] + + def test_retry_after_is_a_floor_and_is_clamped(self): + lo = self._Rand("lo") + assert _ingest_busy_delay(0, 30, lo) == 30.0 + # A longer backoff than the server asked for is kept. + assert _ingest_busy_delay(4, 30, lo) == 60.0 + assert _ingest_busy_delay(0, 3600, lo) == cloudsec_mod._INGEST_RETRY_AFTER_MAX_S + + def test_jitter_is_on_top_and_bounded(self): + assert _ingest_busy_delay(0, 30, self._Rand("hi")) == 45.0 + # Real randomness spreads a fan-out: 50 draws are not all the same wait. + draws = {_ingest_busy_delay(0, 30) for _ in range(50)} + assert len(draws) > 1 + assert all(30.0 <= d <= 45.0 for d in draws)