diff --git a/.release-please-manifest.json b/.release-please-manifest.json index bd34dce0..3e05ae56 100644 --- a/.release-please-manifest.json +++ b/.release-please-manifest.json @@ -1,3 +1,3 @@ { - ".": "0.107.0" + ".": "0.108.0" } \ No newline at end of file diff --git a/CHANGELOG.md b/CHANGELOG.md index 235126cc..1417baaa 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -1,5 +1,12 @@ # Changelog +## [0.108.0](https://github.com/kernel/kernel-python-sdk/compare/v0.107.0...v0.108.0) (2026-09-17) + + +### Features + +* feat: add config registry analysis waiter ([e38c887](https://github.com/kernel/kernel-python-sdk/commit/e38c8878f9f80f5d28fce9f92df6f0a18a48e2bb)) + ## [0.107.0](https://github.com/kernel/kernel-python-sdk/compare/v0.106.0...v0.107.0) (2026-09-16) diff --git a/pyproject.toml b/pyproject.toml index 885ce9e3..3a9c0dee 100644 --- a/pyproject.toml +++ b/pyproject.toml @@ -1,6 +1,6 @@ [project] name = "kernel" -version = "0.107.0" +version = "0.108.0" description = "The official Python library for the kernel API" dynamic = ["readme"] license = "Apache-2.0" diff --git a/src/kernel/_version.py b/src/kernel/_version.py index 9445a514..5c651d4a 100644 --- a/src/kernel/_version.py +++ b/src/kernel/_version.py @@ -1,4 +1,4 @@ # File generated from our OpenAPI spec by Stainless. See CONTRIBUTING.md for details. __title__ = "kernel" -__version__ = "0.107.0" # x-release-please-version +__version__ = "0.108.0" # x-release-please-version diff --git a/src/kernel/lib/config_registry_wait.py b/src/kernel/lib/config_registry_wait.py new file mode 100644 index 00000000..fbf399a8 --- /dev/null +++ b/src/kernel/lib/config_registry_wait.py @@ -0,0 +1,63 @@ +from __future__ import annotations + +import math +import time +import random +from datetime import datetime + +from .._types import Omit, Headers +from .._exceptions import KernelError +from ..types.config_registry_response import ConfigRegistryResponse + +DEFAULT_CONFIG_REGISTRY_POLL_INTERVAL = 5.0 +_TERMINAL_STATUSES = frozenset({"completed", "failed", "canceled", "expired"}) + + +def validate_wait_options(poll_interval: float, max_wait_seconds: float | None) -> None: + if not math.isfinite(poll_interval) or poll_interval <= 0: + raise ValueError("Expected a finite, positive value for `poll_interval`") + if max_wait_seconds is not None and (not math.isfinite(max_wait_seconds) or max_wait_seconds < 0): + raise ValueError("Expected a finite, non-negative value for `max_wait_seconds`") + + +def poll_headers(extra_headers: Headers | None) -> Headers: + headers: dict[str, str | Omit] = dict(extra_headers or {}) + headers["X-Stainless-Poll-Helper"] = "true" + return headers + + +def poll_delay(poll_interval: float) -> float: + return poll_interval * random.uniform(0.9, 1.1) + + +def wait_timeout_error(id: str, polls: int, last_status: str | None, started_at: float) -> TimeoutError: + elapsed = time.monotonic() - started_at + return TimeoutError( + f"Timed out waiting for config registry analysis {id!r} after {elapsed:.1f}s " + f"and {polls} polls; last status was {last_status!r}" + ) + + +def analysis_finished(response: ConfigRegistryResponse, requested_id: str) -> tuple[bool, str]: + analysis = response.analysis + if analysis is None: + raise KernelError(f"Config registry response for {requested_id!r} is missing an analysis") + + analysis_id = getattr(analysis, "id", None) + if not isinstance(analysis_id, str) or not analysis_id: + raise KernelError(f"Config registry response for {requested_id!r} has no valid analysis ID") + if analysis_id != requested_id: + raise KernelError(f"Config registry response for {requested_id!r} returned analysis {analysis_id!r}") + + status = getattr(analysis, "status", None) + if not isinstance(status, str) or not status: + raise KernelError(f"Config registry analysis {requested_id!r} has no valid status") + + if "finished_at" not in analysis.model_fields_set: + raise KernelError(f"Config registry analysis {requested_id!r} is missing `finished_at`") + + finished_at = getattr(analysis, "finished_at", None) + if finished_at is not None and not isinstance(finished_at, datetime): + raise KernelError(f"Config registry analysis {requested_id!r} has an invalid `finished_at`") + + return finished_at is not None or status in _TERMINAL_STATUSES, status diff --git a/src/kernel/resources/config_registry/analyses.py b/src/kernel/resources/config_registry/analyses.py index 3d975f32..bbaeff64 100644 --- a/src/kernel/resources/config_registry/analyses.py +++ b/src/kernel/resources/config_registry/analyses.py @@ -2,6 +2,8 @@ from __future__ import annotations +import time + import httpx from ..._types import Body, Omit, Query, Headers, NotGiven, omit, not_given @@ -18,6 +20,14 @@ from ..._base_client import AsyncPaginator, make_request_options from ...types.config_registry import analysis_list_params from ...types.analysis_summary import AnalysisSummary +from ...lib.config_registry_wait import ( + DEFAULT_CONFIG_REGISTRY_POLL_INTERVAL, + poll_delay, + poll_headers, + analysis_finished, + wait_timeout_error, + validate_wait_options, +) from ...types.config_registry_response import ConfigRegistryResponse __all__ = ["AnalysesResource", "AsyncAnalysesResource"] @@ -79,6 +89,54 @@ def retrieve( cast_to=ConfigRegistryResponse, ) + def wait_for_result( + self, + id: str, + *, + poll_interval: float = DEFAULT_CONFIG_REGISTRY_POLL_INTERVAL, + max_wait_seconds: float | None = None, + extra_headers: Headers | None = None, + extra_query: Query | None = None, + extra_body: Body | None = None, + timeout: float | httpx.Timeout | None | NotGiven = not_given, + ) -> ConfigRegistryResponse: + """Wait for an analysis to finish and return its complete result. + + The first retrieval happens immediately. ``max_wait_seconds`` is a soft + polling deadline: an in-flight request and its normal retries may finish + after it. Timing out does not cancel the remote analysis. + """ + validate_wait_options(poll_interval, max_wait_seconds) + started_at = time.monotonic() + deadline = started_at + max_wait_seconds if max_wait_seconds is not None else None + headers = poll_headers(extra_headers) + polls = 0 + last_status: str | None = None + + while True: + if polls > 0 and deadline is not None and time.monotonic() >= deadline: + raise wait_timeout_error(id, polls, last_status, started_at) + + response = self.retrieve( + id, + extra_headers=headers, + extra_query=extra_query, + extra_body=extra_body, + timeout=timeout, + ) + polls += 1 + finished, last_status = analysis_finished(response, id) + if finished: + return response + + delay = poll_delay(poll_interval) + if deadline is not None: + remaining = deadline - time.monotonic() + if remaining <= 0: + raise wait_timeout_error(id, polls, last_status, started_at) + delay = min(delay, remaining) + self._sleep(delay) + def list( self, *, @@ -220,6 +278,55 @@ async def retrieve( cast_to=ConfigRegistryResponse, ) + async def wait_for_result( + self, + id: str, + *, + poll_interval: float = DEFAULT_CONFIG_REGISTRY_POLL_INTERVAL, + max_wait_seconds: float | None = None, + extra_headers: Headers | None = None, + extra_query: Query | None = None, + extra_body: Body | None = None, + timeout: float | httpx.Timeout | None | NotGiven = not_given, + ) -> ConfigRegistryResponse: + """Wait for an analysis to finish and return its complete result. + + The first retrieval happens immediately. ``max_wait_seconds`` is a soft + polling deadline: an in-flight request and its normal retries may finish + after it. Timing out does not cancel the remote analysis. Cancelling the + calling task stops the wait without cancelling the remote analysis. + """ + validate_wait_options(poll_interval, max_wait_seconds) + started_at = time.monotonic() + deadline = started_at + max_wait_seconds if max_wait_seconds is not None else None + headers = poll_headers(extra_headers) + polls = 0 + last_status: str | None = None + + while True: + if polls > 0 and deadline is not None and time.monotonic() >= deadline: + raise wait_timeout_error(id, polls, last_status, started_at) + + response = await self.retrieve( + id, + extra_headers=headers, + extra_query=extra_query, + extra_body=extra_body, + timeout=timeout, + ) + polls += 1 + finished, last_status = analysis_finished(response, id) + if finished: + return response + + delay = poll_delay(poll_interval) + if deadline is not None: + remaining = deadline - time.monotonic() + if remaining <= 0: + raise wait_timeout_error(id, polls, last_status, started_at) + delay = min(delay, remaining) + await self._sleep(delay) + def list( self, *, diff --git a/tests/test_config_registry_wait.py b/tests/test_config_registry_wait.py new file mode 100644 index 00000000..e668aa8c --- /dev/null +++ b/tests/test_config_registry_wait.py @@ -0,0 +1,181 @@ +from __future__ import annotations + +import asyncio +from typing import Any +from collections.abc import Iterator + +import httpx +import pytest + +from kernel import Kernel, AsyncKernel +from kernel._exceptions import KernelError + + +def analysis_response( + status: str, + *, + analysis_id: str = "analysis-1", + finished_at: str | None = None, +) -> dict[str, Any]: + return { + "analysis": { + "id": analysis_id, + "created_at": "2026-09-16T00:00:00Z", + "expires_at": "2026-09-16T00:45:00Z", + "failure": None, + "finished_at": finished_at, + "status": status, + "intent": None, + }, + "recommendation": None, + "target": { + "domain": "example.com", + "host": "example.com", + "normalized": "https://example.com/", + }, + "working_configurations": [], + "guidance": None, + "workload_outcome": None, + } + + +def response_sequence(payloads: list[dict[str, Any]]) -> tuple[httpx.MockTransport, list[httpx.Request]]: + remaining: Iterator[dict[str, Any]] = iter(payloads) + requests: list[httpx.Request] = [] + + def handler(request: httpx.Request) -> httpx.Response: + requests.append(request) + return httpx.Response(200, json=next(remaining)) + + return httpx.MockTransport(handler), requests + + +def test_wait_for_result_polls_unknown_unfinished_status_until_finished() -> None: + transport, requests = response_sequence( + [ + analysis_response("running"), + analysis_response("queued"), + analysis_response("archived", finished_at="2026-09-16T00:01:00Z"), + ] + ) + + with httpx.Client(transport=transport) as http_client: + client = Kernel(api_key="test", base_url="https://api.example", http_client=http_client) + result = client.config_registry.analyses.wait_for_result( + "analysis-1", + poll_interval=0.001, + extra_headers={"X-Test": "preserved", "X-Stainless-Poll-Helper": "caller"}, + ) + + assert result.analysis is not None + assert result.analysis.status == "archived" + assert len(requests) == 3 + assert all(request.headers["X-Test"] == "preserved" for request in requests) + assert all(request.headers["X-Stainless-Poll-Helper"] == "true" for request in requests) + + +@pytest.mark.parametrize("status", ["completed", "failed", "canceled", "expired"]) +def test_wait_for_result_returns_known_terminal_status_without_finished_at(status: str) -> None: + transport, requests = response_sequence([analysis_response(status)]) + + with httpx.Client(transport=transport) as http_client: + client = Kernel(api_key="test", base_url="https://api.example", http_client=http_client) + result = client.config_registry.analyses.wait_for_result("analysis-1") + + assert result.analysis is not None + assert result.analysis.status == status + assert len(requests) == 1 + + +def test_wait_for_result_zero_max_wait_reads_once_then_times_out() -> None: + transport, requests = response_sequence([analysis_response("running")]) + + with httpx.Client(transport=transport) as http_client: + client = Kernel(api_key="test", base_url="https://api.example", http_client=http_client) + with pytest.raises(TimeoutError, match="analysis-1"): + client.config_registry.analyses.wait_for_result("analysis-1", max_wait_seconds=0) + + assert len(requests) == 1 + + +@pytest.mark.parametrize( + "payload, message", + [ + ({"analysis": None}, "missing an analysis"), + (analysis_response("running", analysis_id="analysis-2"), "analysis-2"), + ], +) +def test_wait_for_result_rejects_invalid_analysis_response(payload: dict[str, Any], message: str) -> None: + transport, _ = response_sequence([payload]) + + with httpx.Client(transport=transport) as http_client: + client = Kernel(api_key="test", base_url="https://api.example", http_client=http_client) + with pytest.raises(KernelError, match=message): + client.config_registry.analyses.wait_for_result("analysis-1") + + +def test_wait_for_result_rejects_missing_finished_at() -> None: + payload = analysis_response("running") + analysis = payload["analysis"] + assert isinstance(analysis, dict) + del analysis["finished_at"] + transport, _ = response_sequence([payload]) + + with httpx.Client(transport=transport) as http_client: + client = Kernel(api_key="test", base_url="https://api.example", http_client=http_client) + with pytest.raises(KernelError, match="finished_at"): + client.config_registry.analyses.wait_for_result("analysis-1") + + +@pytest.mark.parametrize( + "kwargs", + [ + {"poll_interval": 0}, + {"poll_interval": float("nan")}, + {"max_wait_seconds": -1}, + {"max_wait_seconds": float("inf")}, + ], +) +def test_wait_for_result_rejects_invalid_timing_options(kwargs: dict[str, float]) -> None: + client = Kernel(api_key="test", base_url="https://api.example") + + with pytest.raises(ValueError): + client.config_registry.analyses.wait_for_result("analysis-1", **kwargs) # type: ignore[arg-type] + + +async def test_async_wait_for_result_matches_sync_behavior() -> None: + transport, requests = response_sequence( + [ + analysis_response("running"), + analysis_response("completed", finished_at="2026-09-16T00:01:00Z"), + ] + ) + + async with httpx.AsyncClient(transport=transport) as http_client: + client = AsyncKernel(api_key="test", base_url="https://api.example", http_client=http_client) + result = await client.config_registry.analyses.wait_for_result("analysis-1", poll_interval=0.001) + + assert result.analysis is not None + assert result.analysis.status == "completed" + assert len(requests) == 2 + + +async def test_async_wait_for_result_is_cancellable_during_poll_sleep() -> None: + first_request = asyncio.Event() + requests: list[httpx.Request] = [] + + def handler(request: httpx.Request) -> httpx.Response: + requests.append(request) + first_request.set() + return httpx.Response(200, json=analysis_response("running")) + + async with httpx.AsyncClient(transport=httpx.MockTransport(handler)) as http_client: + client = AsyncKernel(api_key="test", base_url="https://api.example", http_client=http_client) + task = asyncio.create_task(client.config_registry.analyses.wait_for_result("analysis-1", poll_interval=60)) + await first_request.wait() + await asyncio.sleep(0) + task.cancel() + with pytest.raises(asyncio.CancelledError): + await task + + assert len(requests) == 1