Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
6 changes: 6 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
29 changes: 29 additions & 0 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Expand Down
59 changes: 59 additions & 0 deletions src/tinybird_sdk/api/api.py
Original file line number Diff line number Diff line change
Expand Up @@ -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)
Expand Down
50 changes: 50 additions & 0 deletions src/tinybird_sdk/client/base.py
Original file line number Diff line number Diff line change
Expand Up @@ -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"):
Expand All @@ -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"],
Expand Down Expand Up @@ -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

Expand Down
181 changes: 181 additions & 0 deletions tests/test_jobs.py
Original file line number Diff line number Diff line change
@@ -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")
Loading