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
11 changes: 11 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
2 changes: 1 addition & 1 deletion pyproject.toml
Original file line number Diff line number Diff line change
Expand Up @@ -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",
Expand Down
110 changes: 81 additions & 29 deletions src/tinybird_sdk/api/api.py
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down Expand Up @@ -298,46 +307,68 @@ 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,
datasource_name: str,
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,
Expand All @@ -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)
Expand Down
103 changes: 103 additions & 0 deletions tests/test_api_datasource_sample.py
Original file line number Diff line number Diff line change
@@ -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
46 changes: 41 additions & 5 deletions tests/test_api_parity.py
Original file line number Diff line number Diff line change
Expand Up @@ -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)
Expand All @@ -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:
Expand Down
23 changes: 17 additions & 6 deletions uv.lock

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

Loading