diff --git a/CHANGELOG.md b/CHANGELOG.md index 0c1ecb0..9436ab3 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -5,6 +5,17 @@ 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 + +- `TinybirdApi.sample_datasource()` starts a sample import job for S3/GCS/DynamoDB connected data sources, delegating to the installed Tinybird CLI client's `datasource_sample()`. For blob storage connectors `max_files` bounds the number of imported files; for DynamoDB the sample is bounded by either `rows` or `max_bytes` (mutually exclusive), or `full_export` triggers a full PITR export of the whole table instead of a bounded scan. This lets cloud branches and local import bounded DynamoDB samples, avoiding slow exports and unnecessary egress costs on branches. + +### Changed + +- Bumped the pinned `tinybird` dependency to `>=4.6.22,<4.7.0` (from `>=4.6.0,<4.7.0`) — the version that adds DynamoDB support to `datasource_sample()`. +- `TinybirdApi.truncate_datasource()`, `delete_datasource()`, and `sample_datasource()` now delegate to the installed Tinybird CLI client (`tinybird.tb.client.TinyB`) instead of issuing raw HTTP requests, avoiding duplicating request-building logic that's already implemented and tested there. As a result, the `timeout` option is no longer honored for these three methods — the underlying client manages its own request timeout/retries. + ## [0.4.0] - 2026-06-29 ### Added diff --git a/pyproject.toml b/pyproject.toml index 780ebd0..811070f 100644 --- a/pyproject.toml +++ b/pyproject.toml @@ -9,7 +9,7 @@ authors = [ ] requires-python = ">=3.11" dependencies = [ - "tinybird>=4.6.0,<4.7.0", + "tinybird>=4.6.22,<4.7.0", ] classifiers = [ "Development Status :: 4 - Beta", diff --git a/src/tinybird_sdk/api/api.py b/src/tinybird_sdk/api/api.py index 8ee18b5..a21bcc0 100644 --- a/src/tinybird_sdk/api/api.py +++ b/src/tinybird_sdk/api/api.py @@ -6,9 +6,18 @@ from email.utils import parsedate_to_datetime import math import time -from typing import Any +from typing import Any, Callable from urllib.parse import urlencode, urljoin +from tinybird.tb.client import ( + AuthException, + CanNotBeDeletedException, + DoesNotExistException, + OperationCanNotBePerformed, + TimeoutException, + TinyB, +) + from .._http import ( HTTPClientError, create_multipart_body, @@ -298,22 +307,15 @@ def delete_datasource( if not delete_condition: raise ValueError("'delete_condition' must be provided in options") - body = {"delete_condition": delete_condition} - dry_run = options.get("dry_run", api_options.get("dry_run")) - if dry_run is not None: - body["dry_run"] = str(dry_run).lower() + dry_run = bool(options.get("dry_run", api_options.get("dry_run")) or False) - response = self.request( - f"/v0/datasources/{datasource_name}/delete", - method="POST", - token=api_options.get("token"), - headers={"Content-Type": "application/x-www-form-urlencoded"}, - body=urlencode(body), - timeout=options.get("timeout", api_options.get("timeout")), + client = self._tinyb_client(api_options.get("token")) + result = self._call_tinyb( + lambda: client.datasource_delete_rows( + datasource_name, delete_condition, dry_run=dry_run + ) ) - if not response.ok: - self._raise_for_error(response.status_code, response.text) - return response.json() + return result or {} def truncate_datasource( self, @@ -321,23 +323,52 @@ def truncate_datasource( options: dict[str, Any] | None = None, api_options: dict[str, Any] | None = None, ) -> dict[str, Any]: + api_options = api_options or {} + client = self._tinyb_client(api_options.get("token")) + result = self._call_tinyb(lambda: client.datasource_truncate(datasource_name)) + return result or {} + + def sample_datasource( + self, + datasource_name: str, + options: dict[str, Any] | None = None, + api_options: dict[str, Any] | None = None, + ) -> dict[str, Any]: + """Start a sample import job for an S3/GCS/DynamoDB connected data source. + + For blob storage (S3/GCS) connectors, ``max_files`` bounds how many files + are imported (default 1, max 10). For DynamoDB, the sample is bounded by + either ``rows`` or ``max_bytes`` (mutually exclusive), or ``full_export`` + triggers a full PITR export of the whole table instead of a bounded scan. + + Options: + max_files: Maximum number of files to import for blob storage connectors. + rows: For DynamoDB, the maximum number of rows to scan and import + (mutually exclusive with ``max_bytes``). + max_bytes: For DynamoDB, the maximum approximate JSONEachRow bytes to + import, e.g. ``"500MB"`` (mutually exclusive with ``rows``). + full_export: For DynamoDB, trigger a full PITR export instead of a + bounded scan. + """ options = options or {} api_options = api_options or {} - response = self.request( - f"/v0/datasources/{datasource_name}/truncate", - method="POST", - token=api_options.get("token"), - timeout=options.get("timeout", api_options.get("timeout")), - ) - if not response.ok: - self._raise_for_error(response.status_code, response.text) - if not response.text.strip(): - return {} - try: - return response.json() - except json.JSONDecodeError: - return {} + rows = options.get("rows") + max_bytes = options.get("max_bytes") + if rows is not None and max_bytes is not None: + raise ValueError("'rows' and 'max_bytes' are mutually exclusive; pass only one") + + client = self._tinyb_client(api_options.get("token")) + result = self._call_tinyb( + lambda: client.datasource_sample( + datasource_name, + max_files=options.get("max_files", 1), + rows=rows, + max_bytes=max_bytes, + full_export=bool(options.get("full_export", False)), + ) + ) + return result or {} def create_token( self, @@ -363,6 +394,27 @@ def create_token( self._raise_for_error(response.status_code, response.text) return response.json() + def _tinyb_client(self, token: str | None = None) -> TinyB: + return TinyB(token=token or self._default_token, host=self._base_url) + + def _call_tinyb(self, fn: Callable[[], Any]) -> Any: + # The installed Tinybird CLI client manages its own request timeout/retries, + # so `options`/`api_options` timeout overrides don't apply to delegated calls. + try: + return fn() + except AuthException as error: + raise TinybirdApiError(str(error), 403) from error + except DoesNotExistException as error: + raise TinybirdApiError(str(error), 404) from error + except OperationCanNotBePerformed as error: + raise TinybirdApiError(str(error), 400) from error + except CanNotBeDeletedException as error: + raise TinybirdApiError(str(error), 409) from error + except TimeoutException as error: + raise TinybirdApiError(str(error), 599) from error + except Exception as error: + raise TinybirdApiError(str(error), 0) from error + 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/tests/test_api_datasource_sample.py b/tests/test_api_datasource_sample.py new file mode 100644 index 0000000..107131c --- /dev/null +++ b/tests/test_api_datasource_sample.py @@ -0,0 +1,103 @@ +from __future__ import annotations + +from typing import Any + +import pytest + +import tinybird_sdk.api.api as api_module +from tinybird_sdk.api.api import TinybirdApi, TinybirdApiError + + +def _make_api() -> TinybirdApi: + return TinybirdApi({"base_url": "https://api.tinybird.co", "token": "p.test"}) + + +def _patch_datasource_sample(monkeypatch: pytest.MonkeyPatch, captured: dict[str, Any]) -> None: + def fake_datasource_sample( + self: Any, + datasource_name: str, + max_files: int = 1, + rows: int | None = None, + max_bytes: str | None = None, + full_export: bool = False, + ) -> dict[str, Any]: + captured["datasource_name"] = datasource_name + captured["max_files"] = max_files + captured["rows"] = rows + captured["max_bytes"] = max_bytes + captured["full_export"] = full_export + return {"job_id": "job-1", "status": "waiting"} + + monkeypatch.setattr(api_module.TinyB, "datasource_sample", fake_datasource_sample) + + +def test_sample_datasource_defaults(monkeypatch: pytest.MonkeyPatch) -> None: + captured: dict[str, Any] = {} + _patch_datasource_sample(monkeypatch, captured) + + result = _make_api().sample_datasource("events") + + assert result == {"job_id": "job-1", "status": "waiting"} + assert captured == { + "datasource_name": "events", + "max_files": 1, + "rows": None, + "max_bytes": None, + "full_export": False, + } + + +def test_sample_datasource_forwards_max_files(monkeypatch: pytest.MonkeyPatch) -> None: + captured: dict[str, Any] = {} + _patch_datasource_sample(monkeypatch, captured) + + _make_api().sample_datasource("events", {"max_files": 3}) + + assert captured["max_files"] == 3 + + +def test_sample_datasource_forwards_dynamodb_rows(monkeypatch: pytest.MonkeyPatch) -> None: + captured: dict[str, Any] = {} + _patch_datasource_sample(monkeypatch, captured) + + _make_api().sample_datasource("ddb_ds", {"rows": 100000}) + + assert captured["rows"] == 100000 + assert captured["max_bytes"] is None + + +def test_sample_datasource_forwards_dynamodb_max_bytes(monkeypatch: pytest.MonkeyPatch) -> None: + captured: dict[str, Any] = {} + _patch_datasource_sample(monkeypatch, captured) + + _make_api().sample_datasource("ddb_ds", {"max_bytes": "1GB"}) + + assert captured["max_bytes"] == "1GB" + assert captured["rows"] is None + + +def test_sample_datasource_forwards_full_export(monkeypatch: pytest.MonkeyPatch) -> None: + captured: dict[str, Any] = {} + _patch_datasource_sample(monkeypatch, captured) + + _make_api().sample_datasource("ddb_ds", {"full_export": True}) + + assert captured["full_export"] is True + + +def test_sample_datasource_rows_and_max_bytes_mutually_exclusive() -> None: + with pytest.raises(ValueError, match="mutually exclusive"): + _make_api().sample_datasource("ddb_ds", {"rows": 10, "max_bytes": "1GB"}) + + +def test_sample_datasource_translates_tinyb_errors(monkeypatch: pytest.MonkeyPatch) -> None: + from tinybird.tb.client import OperationCanNotBePerformed + + def fake_datasource_sample(self: Any, datasource_name: str, **_kwargs: Any) -> None: + raise OperationCanNotBePerformed("bad request") + + monkeypatch.setattr(api_module.TinyB, "datasource_sample", fake_datasource_sample) + + with pytest.raises(TinybirdApiError) as exc: + _make_api().sample_datasource("ddb_ds") + assert exc.value.status_code == 400 diff --git a/tests/test_api_parity.py b/tests/test_api_parity.py index 6f13da0..0049339 100644 --- a/tests/test_api_parity.py +++ b/tests/test_api_parity.py @@ -97,13 +97,11 @@ def bad_fetch(_url: str, **_kwargs: Any) -> _FakeResponse: api.sql("SELECT 2") -def test_append_delete_and_truncate_paths(monkeypatch: pytest.MonkeyPatch, tmp_path: Path) -> None: +def test_append_datasource_paths(monkeypatch: pytest.MonkeyPatch, tmp_path: Path) -> None: calls: list[tuple[str, dict[str, Any]]] = [] def fake_fetch(url: str, **kwargs: Any) -> _FakeResponse: calls.append((url, kwargs)) - if url.endswith("/truncate"): - return _FakeResponse(200, text="") return _FakeResponse(200, {"ok": True}) monkeypatch.setattr(api_module, "tinybird_fetch", fake_fetch) @@ -118,10 +116,48 @@ def fake_fetch(url: str, **kwargs: Any) -> _FakeResponse: api.append_datasource("events", {"file": str(local_file)}) assert "multipart/form-data;" in calls[1][1]["headers"]["Content-Type"] - api.delete_datasource("events", {"delete_condition": "id > 0"}) - assert calls[2][0].endswith("/v0/datasources/events/delete") + +def test_delete_and_truncate_datasource_delegate_to_tinyb(monkeypatch: pytest.MonkeyPatch) -> None: + delete_calls: list[tuple[str, str, bool]] = [] + + def fake_delete_rows( + self: Any, datasource_name: str, delete_condition: str, dry_run: bool = False + ) -> dict[str, Any]: + delete_calls.append((datasource_name, delete_condition, dry_run)) + return {"job_id": "job-1"} + + monkeypatch.setattr(api_module.TinyB, "datasource_delete_rows", fake_delete_rows) + + truncate_calls: list[str] = [] + + def fake_truncate(self: Any, datasource_name: str) -> None: + truncate_calls.append(datasource_name) + return None + + monkeypatch.setattr(api_module.TinyB, "datasource_truncate", fake_truncate) + + api = TinybirdApi({"base_url": "https://api.tinybird.co", "token": "p.token"}) + + result = api.delete_datasource("events", {"delete_condition": "id > 0", "dry_run": True}) + assert result == {"job_id": "job-1"} + assert delete_calls == [("events", "id > 0", True)] assert api.truncate_datasource("events") == {} + assert truncate_calls == ["events"] + + +def test_tinyb_errors_translate_to_tinybird_api_error(monkeypatch: pytest.MonkeyPatch) -> None: + from tinybird.tb.client import DoesNotExistException + + def fake_truncate(self: Any, datasource_name: str) -> None: + raise DoesNotExistException("not found") + + monkeypatch.setattr(api_module.TinyB, "datasource_truncate", fake_truncate) + + api = TinybirdApi({"base_url": "https://api.tinybird.co", "token": "p.token"}) + with pytest.raises(TinybirdApiError) as exc: + api.truncate_datasource("missing") + assert exc.value.status_code == 404 def test_append_requires_either_url_or_file() -> None: diff --git a/uv.lock b/uv.lock index c375471..2d3b3d8 100644 --- a/uv.lock +++ b/uv.lock @@ -1152,7 +1152,7 @@ wheels = [ [[package]] name = "tinybird" -version = "4.6.0" +version = "4.6.22" source = { registry = "https://pypi.org/simple" } dependencies = [ { name = "aiofiles" }, @@ -1183,9 +1183,9 @@ dependencies = [ { name = "watchdog" }, { name = "wheel" }, ] -sdist = { url = "https://files.pythonhosted.org/packages/cc/95/3fc6bff68f18b685e8b835f349901064f6652cf7187a1b9ddfc879683a6a/tinybird-4.6.0.tar.gz", hash = "sha256:a5fa831155978f7ebbdc4d33c6daedd08b8c7edb47439694e4ff4aed9e993173", size = 466849, upload-time = "2026-06-12T10:33:15.269Z" } +sdist = { url = "https://files.pythonhosted.org/packages/a5/1d/afa5179e655428c6c6574a3eb9bf73df863226419adf3bb8bdecedd2ff63/tinybird-4.6.22.tar.gz", hash = "sha256:4bf99abb27fce22409b6285f9d4dbd3634281c24dccc4da33e8f252f4f707fa2", size = 477649, upload-time = "2026-10-02T08:40:55.569Z" } wheels = [ - { url = "https://files.pythonhosted.org/packages/b0/4a/8328ff311282a36996e23477161bdcb766374a62113642e7233aaf0b78a1/tinybird-4.6.0-py3-none-any.whl", hash = "sha256:b883cf9501503bf12146b64d219807819b48cd054641d10b0b8f733c9257b2a3", size = 452823, upload-time = "2026-06-12T10:33:13.276Z" }, + { url = "https://files.pythonhosted.org/packages/fe/5c/8f8081be6d897b122510082cf9ebf93e0d2fcae9af27d379861a76d28a72/tinybird-4.6.22-py3-none-any.whl", hash = "sha256:9859bc581afd351abf220b6fd6555e66a4444790135ef6ad2efda3f69cac2f76", size = 458710, upload-time = "2026-10-02T08:40:53.308Z" }, ] [[package]] @@ -1205,7 +1205,7 @@ dev = [ ] [package.metadata] -requires-dist = [{ name = "tinybird", specifier = ">=4.6.0,<4.7.0" }] +requires-dist = [{ name = "tinybird", specifier = ">=4.6.22,<4.7.0" }] [package.metadata.requires-dev] dev = [ @@ -1226,9 +1226,20 @@ wheels = [ [[package]] name = "tornado" -version = "6.0.4" +version = "6.5.10" source = { registry = "https://pypi.org/simple" } -sdist = { url = "https://files.pythonhosted.org/packages/95/84/119a46d494f008969bf0c775cb2c6b3579d3c4cc1bb1b41a022aa93ee242/tornado-6.0.4.tar.gz", hash = "sha256:0fe2d45ba43b00a41cd73f8be321a44936dc1aba233dee979f17a042b83eb6dc", size = 496204, upload-time = "2020-03-04T02:32:39.825Z" } +sdist = { url = "https://files.pythonhosted.org/packages/06/61/53d562a57b28c08eda40b258c0f975e360541943ad7c7bef897a40caafda/tornado-6.5.10.tar.gz", hash = "sha256:a6b1ccd08c04b4a06fb5aeb381be99de5ad1e5375c1785e31d78c880feb57687", size = 537910, upload-time = "2026-09-15T13:47:48.73Z" } +wheels = [ + { url = "https://files.pythonhosted.org/packages/cd/5b/ff5fc58fa2427c30dea74c90053f4fc5eda1e7f3833ed3ecc7147fe2b311/tornado-6.5.10-cp39-abi3-macosx_10_9_universal2.whl", hash = "sha256:9261783640e23258694a9ff0795df430a5a7b0a651d3dd53dd0969ad6be16da7", size = 465883, upload-time = "2026-09-15T13:47:35.463Z" }, + { url = "https://files.pythonhosted.org/packages/ad/f5/cd7be26c34a3315532f3aef5f092465da8f59c334dd439d3c14aaef16461/tornado-6.5.10-cp39-abi3-macosx_10_9_x86_64.whl", hash = "sha256:83e6cf438b106c6b3852d70960967bb1b70c87438050dca0981e4b9aa751a4c1", size = 464046, upload-time = "2026-09-15T13:47:37.178Z" }, + { url = "https://files.pythonhosted.org/packages/60/33/df6d7d04854a58619f8349a51e3edb138324130a7562b0bb21f115bb940f/tornado-6.5.10-cp39-abi3-manylinux1_x86_64.manylinux_2_28_x86_64.manylinux_2_5_x86_64.whl", hash = "sha256:bdf942448169e5336451d0494d7e3d81cfa726d5aa312affdc4682dd62a62f6d", size = 467096, upload-time = "2026-09-15T13:47:38.559Z" }, + { url = "https://files.pythonhosted.org/packages/29/17/cc35dff68272d685cffd8600ffafbd8067e7d05e7348d9f80caddffbbd5f/tornado-6.5.10-cp39-abi3-manylinux2014_aarch64.manylinux_2_17_aarch64.manylinux_2_28_aarch64.whl", hash = "sha256:69acca6501eed74582b76dbbceee2a91613f54728e3e418346000d7103101676", size = 468067, upload-time = "2026-09-15T13:47:40.085Z" }, + { url = "https://files.pythonhosted.org/packages/c3/01/6e5349b4e1a53a4b4972a6716785e1fe7407f312063c3972690af8ff301b/tornado-6.5.10-cp39-abi3-musllinux_1_2_aarch64.whl", hash = "sha256:66aaa3f57d30c6e6becee83ff28055d5930ac724214bde99393eefda83d5e015", size = 467901, upload-time = "2026-09-15T13:47:41.576Z" }, + { url = "https://files.pythonhosted.org/packages/28/5e/b4facf94370dba006819c8d304376f8b9fbec6b935b5e51bf45823a9790b/tornado-6.5.10-cp39-abi3-musllinux_1_2_x86_64.whl", hash = "sha256:4bd192b959f9128fb99b8898148070ba4574c9589b78bce42d1851131fe85828", size = 467308, upload-time = "2026-09-15T13:47:43.145Z" }, + { url = "https://files.pythonhosted.org/packages/56/ae/047938e828cafc8eca4c908fafb6588fee944e3af39a0af9d7b602499ae5/tornado-6.5.10-cp39-abi3-win32.whl", hash = "sha256:302eb1e0e3e159314eb591920529fdea80acca92df5510a2cec5bbd4f099ec72", size = 468387, upload-time = "2026-09-15T13:47:44.556Z" }, + { url = "https://files.pythonhosted.org/packages/d8/d4/5901517f05affd752490f6a654ba31b7474664e8dd80bd045a00c220bd88/tornado-6.5.10-cp39-abi3-win_amd64.whl", hash = "sha256:37ae8f150cecfdbf747fc4e12f5e9a97ecd8cf1d4cdb3f119e2de84b11196918", size = 468828, upload-time = "2026-09-15T13:47:45.961Z" }, + { url = "https://files.pythonhosted.org/packages/f3/1a/fd497f3a7f7b74bb04f4b94536b5c9f80742b5d50501fd27977652ddec16/tornado-6.5.10-cp39-abi3-win_arm64.whl", hash = "sha256:ce045d3c298fddd30e89a2777f97039d1b641eb9518ac7b26a4721903539c694", size = 467847, upload-time = "2026-09-15T13:47:47.283Z" }, +] [[package]] name = "tqdm"