diff --git a/CHANGELOG.md b/CHANGELOG.md index 0c1ecb0..f97d07a 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -5,6 +5,12 @@ All notable changes to this project will be documented in this file. The format is based on [Keep a Changelog](https://keepachangelog.com/en/1.1.0/), and this project follows [Semantic Versioning](https://semver.org/spec/v2.0.0.html). +## [Unreleased] + +### Added + +- `TinybirdClient.jobs` namespace (`list`/`get`/`cancel`/`retry`) and matching `TinybirdApi.list_jobs()`/`get_job()`/`cancel_job()`/`retry_job()` methods, wrapping `/v0/jobs`. Lets code that triggers async jobs itself (`append`/`replace`, connector syncs) check status, cancel, or retry them afterward, matching `tb job ls/details/cancel/retry`. Retry eligibility (job kind/state) is enforced by the API, not duplicated client-side. + ## [0.4.0] - 2026-06-29 ### Added diff --git a/README.md b/README.md index e42f34b..45d2c3a 100644 --- a/README.md +++ b/README.md @@ -350,6 +350,35 @@ jwt_token = result["token"] - **`fixed_params`**: For pipes, embed parameters that cannot be overridden by the caller. - **`filter`**: For datasources, append a SQL WHERE clause (for example, `"org_id = 'acme'"`). +## Job Management + +Check status, cancel, or retry async jobs started by `append`/`replace` or connector syncs. + +```python +from tinybird_sdk import create_client + +client = create_client( + { + "base_url": "https://api.tinybird.co", + "token": "p.your_admin_token", + } +) + +# List jobs, optionally filtered +jobs = client.jobs.list({"status": "error", "kind": "import"}) + +# Get a single job's status +job = client.jobs.get("job_id") + +# Cancel a running job +client.jobs.cancel("job_id") + +# Retry an eligible import/S3/GCS sync job in error or cancelled state +client.jobs.retry("job_id") +``` + +`TinybirdApi` exposes the same operations directly as `list_jobs()`, `get_job()`, `cancel_job()`, and `retry_job()`. Retry eligibility (job kind and state) is enforced by the Tinybird API, not duplicated client-side — an ineligible retry raises a `TinybirdError`/`TinybirdApiError` with the API's own message. + ## CLI Commands This package installs `tinybird` as a runtime dependency. diff --git a/src/tinybird_sdk/api/api.py b/src/tinybird_sdk/api/api.py index 8ee18b5..638f688 100644 --- a/src/tinybird_sdk/api/api.py +++ b/src/tinybird_sdk/api/api.py @@ -363,6 +363,65 @@ def create_token( self._raise_for_error(response.status_code, response.text) return response.json() + def list_jobs(self, options: dict[str, Any] | None = None) -> dict[str, Any]: + options = options or {} + + query: dict[str, str] = {} + if options.get("status"): + query["status"] = options["status"] + if options.get("kind"): + query["kind"] = options["kind"] + + path = "/v0/jobs" + if query: + path = f"{path}?{urlencode(query)}" + + response = self.request( + path, + method="GET", + token=options.get("token"), + timeout=options.get("timeout"), + ) + if not response.ok: + self._raise_for_error(response.status_code, response.text) + return response.json() + + def get_job(self, job_id: str, options: dict[str, Any] | None = None) -> dict[str, Any]: + options = options or {} + response = self.request( + f"/v0/jobs/{job_id}", + method="GET", + token=options.get("token"), + timeout=options.get("timeout"), + ) + if not response.ok: + self._raise_for_error(response.status_code, response.text) + return response.json() + + def cancel_job(self, job_id: str, options: dict[str, Any] | None = None) -> dict[str, Any]: + options = options or {} + response = self.request( + f"/v0/jobs/{job_id}/cancel", + method="POST", + token=options.get("token"), + timeout=options.get("timeout"), + ) + if not response.ok: + self._raise_for_error(response.status_code, response.text) + return response.json() + + def retry_job(self, job_id: str, options: dict[str, Any] | None = None) -> dict[str, Any]: + options = options or {} + response = self.request( + f"/v0/jobs/{job_id}/retry", + method="POST", + token=options.get("token"), + timeout=options.get("timeout"), + ) + if not response.ok: + self._raise_for_error(response.status_code, response.text) + return response.json() + def _timeout_seconds(self, timeout_ms: int | None) -> float: timeout = timeout_ms if timeout_ms is not None else self._default_timeout return max(timeout / 1000.0, 0.001) diff --git a/src/tinybird_sdk/client/base.py b/src/tinybird_sdk/client/base.py index b48ba09..588753e 100644 --- a/src/tinybird_sdk/client/base.py +++ b/src/tinybird_sdk/client/base.py @@ -35,6 +35,23 @@ def truncate( return self._client._truncate_datasource(datasource_name, options or {}) +class _JobsNamespace: + def __init__(self, client: "TinybirdClient"): + self._client = client + + def list(self, options: dict[str, Any] | None = None) -> dict[str, Any]: + return self._client._list_jobs(options or {}) + + def get(self, job_id: str, options: dict[str, Any] | None = None) -> dict[str, Any]: + return self._client._get_job(job_id, options or {}) + + def cancel(self, job_id: str, options: dict[str, Any] | None = None) -> dict[str, Any]: + return self._client._cancel_job(job_id, options or {}) + + def retry(self, job_id: str, options: dict[str, Any] | None = None) -> dict[str, Any]: + return self._client._retry_job(job_id, options or {}) + + class TinybirdClient: def __init__(self, config: dict[str, Any]): if not config.get("base_url"): @@ -47,6 +64,7 @@ def __init__(self, config: dict[str, Any]): self._resolved_context: ClientContext | None = None self.datasources = _DatasourcesNamespace(self) + self.jobs = _JobsNamespace(self) self.tokens = TokensNamespace( self._get_token, self._config["base_url"], @@ -97,6 +115,38 @@ def _ingest_datasource( self._rethrow_api_error(error) raise AssertionError("unreachable") + def _list_jobs(self, options: dict[str, Any]) -> dict[str, Any]: + token = self._get_token() + try: + return self._get_api(token).list_jobs(options) + except Exception as error: + self._rethrow_api_error(error) + raise AssertionError("unreachable") + + def _get_job(self, job_id: str, options: dict[str, Any]) -> dict[str, Any]: + token = self._get_token() + try: + return self._get_api(token).get_job(job_id, options) + except Exception as error: + self._rethrow_api_error(error) + raise AssertionError("unreachable") + + def _cancel_job(self, job_id: str, options: dict[str, Any]) -> dict[str, Any]: + token = self._get_token() + try: + return self._get_api(token).cancel_job(job_id, options) + except Exception as error: + self._rethrow_api_error(error) + raise AssertionError("unreachable") + + def _retry_job(self, job_id: str, options: dict[str, Any]) -> dict[str, Any]: + token = self._get_token() + try: + return self._get_api(token).retry_job(job_id, options) + except Exception as error: + self._rethrow_api_error(error) + raise AssertionError("unreachable") + def _get_token(self) -> str: return self._resolve_context().token diff --git a/tests/test_jobs.py b/tests/test_jobs.py new file mode 100644 index 0000000..d6aa335 --- /dev/null +++ b/tests/test_jobs.py @@ -0,0 +1,181 @@ +from __future__ import annotations + +from typing import Any +from urllib.parse import parse_qs, urlparse + +import pytest + +import tinybird_sdk.api.api as api_module +from tinybird_sdk.api.api import TinybirdApi, TinybirdApiError +from tinybird_sdk.client.base import TinybirdClient +from tinybird_sdk.client.types import TinybirdError + + +class _FakeResponse: + def __init__(self, status_code: int, payload: dict[str, Any]): + self.status_code = status_code + self._payload = payload + self.text = "" + + @property + def ok(self) -> bool: + return 200 <= self.status_code < 300 + + def json(self) -> dict[str, Any]: + return self._payload + + +def _capture_fetch(captured: dict[str, Any], payload: dict[str, Any]) -> Any: + def fake_fetch(url: str, **kwargs: Any) -> _FakeResponse: + captured["url"] = url + captured["method"] = kwargs.get("method") + return _FakeResponse(200, payload) + + return fake_fetch + + +def _make_api() -> TinybirdApi: + return TinybirdApi({"base_url": "https://api.tinybird.co", "token": "p.test"}) + + +def test_list_jobs_defaults(monkeypatch: pytest.MonkeyPatch) -> None: + captured: dict[str, Any] = {} + monkeypatch.setattr( + api_module, "tinybird_fetch", _capture_fetch(captured, {"jobs": [{"id": "job-1"}]}) + ) + + result = _make_api().list_jobs() + + assert result == {"jobs": [{"id": "job-1"}]} + assert captured["method"] == "GET" + parsed = urlparse(captured["url"]) + assert parsed.path == "/v0/jobs" + assert parsed.query == "" + + +def test_list_jobs_with_filters(monkeypatch: pytest.MonkeyPatch) -> None: + captured: dict[str, Any] = {} + monkeypatch.setattr(api_module, "tinybird_fetch", _capture_fetch(captured, {"jobs": []})) + + _make_api().list_jobs({"status": "error", "kind": "import"}) + + parsed = urlparse(captured["url"]) + assert parse_qs(parsed.query) == {"status": ["error"], "kind": ["import"]} + + +def test_get_job(monkeypatch: pytest.MonkeyPatch) -> None: + captured: dict[str, Any] = {} + monkeypatch.setattr( + api_module, "tinybird_fetch", _capture_fetch(captured, {"id": "job-1", "status": "done"}) + ) + + result = _make_api().get_job("job-1") + + assert result == {"id": "job-1", "status": "done"} + assert captured["method"] == "GET" + assert captured["url"].endswith("/v0/jobs/job-1") + + +def test_cancel_job(monkeypatch: pytest.MonkeyPatch) -> None: + captured: dict[str, Any] = {} + monkeypatch.setattr( + api_module, + "tinybird_fetch", + _capture_fetch(captured, {"id": "job-1", "status": "cancelling"}), + ) + + result = _make_api().cancel_job("job-1") + + assert result == {"id": "job-1", "status": "cancelling"} + assert captured["method"] == "POST" + assert captured["url"].endswith("/v0/jobs/job-1/cancel") + + +def test_retry_job(monkeypatch: pytest.MonkeyPatch) -> None: + captured: dict[str, Any] = {} + monkeypatch.setattr( + api_module, "tinybird_fetch", _capture_fetch(captured, {"id": "job-2", "status": "waiting"}) + ) + + result = _make_api().retry_job("job-1") + + assert result == {"id": "job-2", "status": "waiting"} + assert captured["method"] == "POST" + assert captured["url"].endswith("/v0/jobs/job-1/retry") + + +def test_job_error_response_raises_api_error(monkeypatch: pytest.MonkeyPatch) -> None: + def fake_fetch(_url: str, **_kwargs: Any) -> _FakeResponse: + response = _FakeResponse(404, {"error": "Job not found"}) + response.text = '{"error": "Job not found"}' + return response + + monkeypatch.setattr(api_module, "tinybird_fetch", fake_fetch) + + with pytest.raises(TinybirdApiError, match="Job not found"): + _make_api().get_job("missing") + + +def test_client_jobs_namespace_list_and_get(monkeypatch: pytest.MonkeyPatch) -> None: + class FakeApi: + def __init__(self, _config: dict[str, Any]): + pass + + def list_jobs(self, options: dict[str, Any]) -> dict[str, Any]: + return {"jobs": [], "options": options} + + def get_job(self, job_id: str, options: dict[str, Any]) -> dict[str, Any]: + return {"id": job_id, "options": options} + + import tinybird_sdk.client.base as client_base + + monkeypatch.setattr(client_base, "TinybirdApi", FakeApi) + + client = TinybirdClient({"base_url": "https://api.tinybird.co", "token": "p.test"}) + + assert client.jobs.list({"status": "error"}) == {"jobs": [], "options": {"status": "error"}} + assert client.jobs.get("job-1") == {"id": "job-1", "options": {}} + + +def test_client_jobs_namespace_cancel_retry(monkeypatch: pytest.MonkeyPatch) -> None: + class FakeApi: + def __init__(self, _config: dict[str, Any]): + pass + + def cancel_job(self, job_id: str, options: dict[str, Any]) -> dict[str, Any]: + return {"id": job_id, "status": "cancelling"} + + def retry_job(self, job_id: str, options: dict[str, Any]) -> dict[str, Any]: + return {"id": job_id, "status": "waiting"} + + import tinybird_sdk.client.base as client_base + + monkeypatch.setattr(client_base, "TinybirdApi", FakeApi) + + client = TinybirdClient({"base_url": "https://api.tinybird.co", "token": "p.test"}) + + assert client.jobs.cancel("job-1") == {"id": "job-1", "status": "cancelling"} + assert client.jobs.retry("job-1") == {"id": "job-1", "status": "waiting"} + + +def test_client_jobs_namespace_wraps_api_errors(monkeypatch: pytest.MonkeyPatch) -> None: + class FakeApi: + def __init__(self, _config: dict[str, Any]): + pass + + def retry_job(self, job_id: str, options: dict[str, Any]) -> dict[str, Any]: + raise TinybirdApiError( + "Job is not eligible for retry", + 400, + '{"error":"Job is not eligible for retry"}', + {"error": "Job is not eligible for retry"}, + ) + + import tinybird_sdk.client.base as client_base + + monkeypatch.setattr(client_base, "TinybirdApi", FakeApi) + + client = TinybirdClient({"base_url": "https://api.tinybird.co", "token": "p.test"}) + + with pytest.raises(TinybirdError, match="not eligible for retry"): + client.jobs.retry("job-1")