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

- `AsyncTinybirdClient`/`AsyncTinybirdApi` (and an `AsyncTinybird` facade counterpart to `Tinybird`), covering the same `query`/`ingest`/`ingest_batch`/`sql`/`datasources.*`/`tokens.create_jwt` surface as the sync client, backed by `httpx.AsyncClient` instead of blocking `urllib` calls, so the SDK no longer blocks the event loop when used from async frameworks (FastAPI, aiohttp). Added `httpx` as a runtime dependency. `TinybirdClient`/`TinybirdApi`'s own request/retry/error-handling logic was refactored to share its pure helpers (`src/tinybird_sdk/api/_shared.py`) with the new async client, with no behavior change to the sync path.

## [0.4.0] - 2026-06-29

### Added
Expand Down
64 changes: 64 additions & 0 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -303,6 +303,70 @@ workspace_response = api.request_json(

This Tinybird API is standalone and can be used without `create_client()` or `Tinybird(...)`.

## Async Client

Every client layer (`create_client`/`TinybirdClient`, `Tinybird`, and the low-level
`create_tinybird_api`/`TinybirdApi`) has an async counterpart backed by `httpx.AsyncClient`
instead of blocking `urllib` calls, so it won't block the event loop when called from
async frameworks like FastAPI or aiohttp. Method names and options are identical to the
sync client — just `await` them:

```python
# app.py
from fastapi import FastAPI
from tinybird_sdk import create_async_client

app = FastAPI()
client = create_async_client(
{
"base_url": "https://api.tinybird.co",
"token": "p.your_token",
}
)


@app.post("/events")
async def ingest_event(event: dict):
return await client.ingest("events", event)


@app.get("/top_pages")
async def top_pages(start_date: str, end_date: str):
return await client.query("top_pages", {"start_date": start_date, "end_date": end_date})


@app.on_event("shutdown")
async def shutdown() -> None:
await client.aclose()
```

`AsyncTinybird` is the async counterpart to the generated `Tinybird(...)` facade — construct
it directly with the same `datasources`/`pipes` dicts your project's generated code already
exposes (the code generator still emits the sync `Tinybird(...)` by default):

```python
from fastapi import FastAPI
from tinybird_sdk import AsyncTinybird
from lib.datasources import events
from lib.pipes import top_pages

app = FastAPI()
tinybird = AsyncTinybird({"datasources": {"events": events}, "pipes": {"top_pages": top_pages}})


@app.post("/events")
async def ingest_event(event: dict):
return await tinybird.events.ingest(event)
```

The low-level `AsyncTinybirdApi` (via `create_async_tinybird_api`) is also available
standalone, mirroring `create_tinybird_api()`'s methods as `async def`.

Note: creating or resolving a Tinybird branch (`dev_mode=True`) still goes through the
synchronous branch-management API under the hood (branch creation can take up to a couple
of minutes to provision), offloaded via a background thread so it doesn't block the event
loop — it's a one-time setup cost, not a per-request one.

## JWT Token Creation

Create short-lived JWT tokens for secure scoped access to Tinybird resources.
Expand Down
1 change: 1 addition & 0 deletions pyproject.toml
Original file line number Diff line number Diff line change
Expand Up @@ -10,6 +10,7 @@ authors = [
requires-python = ">=3.11"
dependencies = [
"tinybird>=4.6.0,<4.7.0",
"httpx>=0.27,<1",
]
classifiers = [
"Development Status :: 4 - Beta",
Expand Down
6 changes: 6 additions & 0 deletions src/tinybird_sdk/__init__.py
Original file line number Diff line number Diff line change
Expand Up @@ -18,6 +18,7 @@
"sql": ("tinybird_sdk.schema", "sql"),
"define_project": ("tinybird_sdk.schema", "define_project"),
"Tinybird": ("tinybird_sdk.schema", "Tinybird"),
"AsyncTinybird": ("tinybird_sdk.schema", "AsyncTinybird"),
# Schema helpers
"column": ("tinybird_sdk.schema", "column"),
"get_column_type": ("tinybird_sdk.schema", "get_column_type"),
Expand Down Expand Up @@ -70,6 +71,8 @@
# Client
"TinybirdClient": ("tinybird_sdk.client", "TinybirdClient"),
"create_client": ("tinybird_sdk.client", "create_client"),
"AsyncTinybirdClient": ("tinybird_sdk.client", "AsyncTinybirdClient"),
"create_async_client": ("tinybird_sdk.client", "create_async_client"),
"TinybirdError": ("tinybird_sdk.client", "TinybirdError"),
"is_preview_environment": ("tinybird_sdk.client", "is_preview_environment"),
"get_preview_branch_name": ("tinybird_sdk.client", "get_preview_branch_name"),
Expand All @@ -80,7 +83,10 @@
"create_tinybird_api": ("tinybird_sdk.api.api", "create_tinybird_api"),
"create_tinybird_api_wrapper": ("tinybird_sdk.api.api", "create_tinybird_api_wrapper"),
"TinybirdApiError": ("tinybird_sdk.api.api", "TinybirdApiError"),
"AsyncTinybirdApi": ("tinybird_sdk.api.async_api", "AsyncTinybirdApi"),
"create_async_tinybird_api": ("tinybird_sdk.api.async_api", "create_async_tinybird_api"),
"create_jwt": ("tinybird_sdk.api.tokens", "create_jwt"),
"create_jwt_async": ("tinybird_sdk.api.tokens", "create_jwt_async"),
"TokenApiError": ("tinybird_sdk.api.tokens", "TokenApiError"),
"parse_api_url": ("tinybird_sdk.api.dashboard", "parse_api_url"),
"get_dashboard_url": ("tinybird_sdk.api.dashboard", "get_dashboard_url"),
Expand Down
39 changes: 38 additions & 1 deletion src/tinybird_sdk/_http.py
Original file line number Diff line number Diff line change
Expand Up @@ -7,11 +7,14 @@
from dataclasses import dataclass
from datetime import date, datetime
from decimal import Decimal
from typing import Any, Mapping
from typing import TYPE_CHECKING, Any, Mapping
from urllib.error import HTTPError
from urllib.parse import parse_qsl, urlencode, urlparse, urlunparse
from urllib.request import Request, urlopen

if TYPE_CHECKING:
import httpx

TINYBIRD_FROM_PARAM = "python-sdk"


Expand Down Expand Up @@ -92,6 +95,40 @@ def tinybird_fetch(
raise HTTPClientError(str(error)) from error


async def tinybird_fetch_async(
client: "httpx.AsyncClient",
url: str,
*,
method: str = "GET",
headers: Mapping[str, str] | None = None,
body: bytes | str | None = None,
timeout: float | None = None,
) -> HTTPResponse:
"""Async counterpart to `tinybird_fetch`, sharing its request-building (via
`with_tinybird_from_param`) and response shape (`HTTPResponse`). Takes an
explicit `httpx.AsyncClient` so callers (AsyncTinybirdApi) own the client's
lifecycle and connection pool rather than opening one per call.
"""
import httpx

request_body = body.encode("utf-8") if isinstance(body, str) else body
try:
response = await client.request(
method,
with_tinybird_from_param(url),
headers=dict(headers or {}),
content=request_body,
timeout=timeout,
)
return HTTPResponse(
status_code=response.status_code,
headers=dict(response.headers),
body=response.content,
)
except httpx.HTTPError as error:
raise HTTPClientError(str(error)) from error


def normalize_base_url(base_url: str) -> str:
return base_url.rstrip("/")

Expand Down
3 changes: 3 additions & 0 deletions src/tinybird_sdk/api/__init__.py
Original file line number Diff line number Diff line change
Expand Up @@ -5,6 +5,7 @@
create_tinybird_api,
create_tinybird_api_wrapper,
)
from .async_api import AsyncTinybirdApi, create_async_tinybird_api
from .fetcher import (
TINYBIRD_FROM_PARAM,
create_tinybird_fetcher,
Expand Down Expand Up @@ -74,6 +75,8 @@
"TinybirdApiError",
"create_tinybird_api",
"create_tinybird_api_wrapper",
"AsyncTinybirdApi",
"create_async_tinybird_api",
"create_jwt",
"TokenApiError",
"build_to_tinybird",
Expand Down
153 changes: 153 additions & 0 deletions src/tinybird_sdk/api/_shared.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,153 @@
from __future__ import annotations

import json
import math
import time
from dataclasses import dataclass
from datetime import timezone
from email.utils import parsedate_to_datetime
from typing import Any

"""Transport-agnostic helpers shared by TinybirdApi and AsyncTinybirdApi.

Kept separate from api.py so the sync and async clients can share identical
retry/error semantics without one importing internals from the other.
"""


@dataclass(frozen=True, slots=True)
class ApiErrorInfo:
message: str
status_code: int
response_body: str | None
response: dict[str, Any] | None


def resolve_ingest_max_retries(options: dict[str, Any]) -> int | None:
value = options.get("maxRetries")
if value is None:
value = options.get("max_retries")
if value is None:
return None
if isinstance(value, bool) or not isinstance(value, (int, float)) or not math.isfinite(value):
raise ValueError("'maxRetries' must be a finite number")
return max(0, math.floor(value))


def get_header(headers: dict[str, str] | Any, header_name: str) -> str | None:
if hasattr(headers, "get"):
value = headers.get(header_name)
if value is None:
value = headers.get(header_name.lower())
if value is None:
value = headers.get(header_name.title())
if isinstance(value, str):
return value

for key, value in dict(headers).items():
if isinstance(key, str) and key.lower() == header_name.lower() and isinstance(value, str):
return value
return None


def parse_retry_after_delay_ms(value: str | None) -> int | None:
if not value:
return None

trimmed = value.strip()
try:
seconds = float(trimmed)
if math.isfinite(seconds):
return max(0, math.floor(seconds * 1000))
except ValueError:
pass

try:
parsed_date = parsedate_to_datetime(trimmed)
except (TypeError, ValueError):
return None

if parsed_date.tzinfo is None:
parsed_date = parsed_date.replace(tzinfo=timezone.utc)

return max(0, math.floor((parsed_date.timestamp() - time.time()) * 1000))


def parse_rate_limit_reset_delay_ms(value: str | None) -> int | None:
if not value:
return None
try:
numeric_value = float(value.strip())
except ValueError:
return None
if not math.isfinite(numeric_value):
return None
return max(0, math.floor(numeric_value * 1000))


def resolve_retry_delay_from_headers(headers: dict[str, str] | Any) -> int | None:
retry_after = get_header(headers, "retry-after")
retry_after_delay_ms = parse_retry_after_delay_ms(retry_after)
if retry_after_delay_ms is not None:
return retry_after_delay_ms

rate_limit_reset = get_header(headers, "x-ratelimit-reset")
return parse_rate_limit_reset_delay_ms(rate_limit_reset)


def resolve_retry_429_delay_ms(
status_code: int,
headers: dict[str, str] | Any,
max_retries: int | None,
retry_count: int,
) -> int | None:
if max_retries is None or status_code != 429 or retry_count >= max_retries:
return None
return resolve_retry_delay_from_headers(headers)


def calculate_retry_503_delay_ms(
retry_count: int,
*,
base_delay_ms: int,
max_delay_ms: int,
) -> int:
return min(max_delay_ms, base_delay_ms * (2**retry_count))


def resolve_retry_503_delay_ms(
status_code: int,
max_retries: int | None,
retry_count: int,
*,
base_delay_ms: int,
max_delay_ms: int,
) -> int | None:
if max_retries is None or status_code != 503 or retry_count >= max_retries:
return None
return calculate_retry_503_delay_ms(
retry_count, base_delay_ms=base_delay_ms, max_delay_ms=max_delay_ms
)


def build_api_error_info(status_code: int, body: str) -> ApiErrorInfo:
parsed: dict[str, Any] | None = None
try:
parsed = json.loads(body) if body else None
except json.JSONDecodeError:
parsed = None

message = ""
if parsed and parsed.get("error"):
message = str(parsed["error"])
elif body:
message = f"Request failed with status {status_code}: {body}"
else:
message = f"Request failed with status {status_code}"

return ApiErrorInfo(
message=message,
status_code=status_code,
response_body=body or None,
response=parsed,
)
Loading
Loading