Skip to content
Merged
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
16 changes: 15 additions & 1 deletion docs/changelog.rst
Original file line number Diff line number Diff line change
Expand Up @@ -24,6 +24,13 @@ Unreleased
* ADBC FlightSQL adds options for TLS/mTLS, RPC timeouts, message size, cookies
and headers. Values set in native ``db_kwargs`` take precedence.

* PostgreSQL adapters expose native asyncpg custom codecs and per-query timeouts,
psycopg null pools and JSON codecs, supported CockroachDB startup settings,
and psqlpy dense-vector conversion. PgBouncer compatibility mode avoids
explicit prepared stack statements without weakening transaction cleanup.
Null pools preserve concurrency limits, and timeout forwarding retains
explicit zero values.

* Added an IBM Db2 adapter for Db2 LUW 11.5 and later with sync
(``Db2SyncConfig``) and async (``Db2AsyncConfig``) configurations built on
``ibm_db``. It includes connection pooling, catalog reflection, migrations,
Expand All @@ -38,14 +45,21 @@ Unreleased

**Fixed:**

* Psycopg reads COPY files in chunks, not all at once. ADK stores use
RETURNING to cut round trips. Psqlpy closes a connection if setup fails.

* Asyncpg stack telemetry reports sequential prepared execution rather than
native pipelining. Each statement still returns its own result.

* Builder upserts emit ``MERGE`` for the ``db2`` dialect.

* ADBC ADK stores reuse cached PostgreSQL placeholder conversion and preserve
question marks in quoted identifiers, literals, and comments.

* Arrow ODBC pagination reuses compiled placeholder positions instead of
parsing SQL again. ADBC keeps bound values in its ADK store queries.
DuckDB Arrow loads keep sparse dictionary fields and quote table names.

* Builder upserts emit ``MERGE`` for the ``db2`` dialect.
* The arrow-odbc adapter detects the SQL dialect from the ODBC driver name
only, so database, host, or user names no longer select the wrong dialect.

Expand Down
19 changes: 19 additions & 0 deletions docs/reference/adapters/asyncpg.rst
Original file line number Diff line number Diff line change
Expand Up @@ -89,3 +89,22 @@ All PostgreSQL adapters support ``data ? 'key'``, ``data ? $1``,
``data ? :key``, ``?|``, and ``?&`` without treating the operator as a
placeholder. Write ``data ?? other_col`` for identifier or function right-hand operands. Write parameterized intervals as ``? * interval '1 day'``
(or use ``$1``).

Native execution options
------------------------

Set ``execution_args={"timeout": seconds}`` on ``StatementConfig`` to forward
an asyncpg timeout to queries, batches, script statements, stack operations,
and stream fetches. ``command_timeout`` is accepted as an alias; ``timeout``
takes precedence. Explicit ``0`` and ``None`` values are preserved.

``driver_features={"pgbouncer": True}`` disables asyncpg's statement cache and
SQLSpec's explicit prepared statements in stacks, while preserving transaction
cleanup. The option is also accepted in ``connection_config``. Use this mode
when the proxy configuration does not support prepared statements; PgBouncer
can support them when configured to track protocol-level prepared statements.

``driver_features["type_codecs"]`` accepts a list of native codec specifications.
Each entry requires ``typename``, ``encoder``, and ``decoder``; ``schema`` defaults
to ``public`` and ``format`` to ``text``. Codecs register after SQLSpec's built-in
JSON and vector setup and before ``on_connection_create``.
8 changes: 8 additions & 0 deletions docs/reference/adapters/cockroach_asyncpg.rst
Original file line number Diff line number Diff line change
Expand Up @@ -130,3 +130,11 @@ namespace: ``"litestar"``, ``"events"``, or ``"adk"`` as supported by this adapt
.. autoclass:: sqlspec.adapters.cockroach_asyncpg.adk.CockroachAsyncpgADKConfig
:members:
:show-inheritance:

Native startup settings
-----------------------

``connection_config`` accepts ``application_name``,
``default_transaction_use_follower_reads`` (boolean), and ``results_buffer_size``
(non-negative integer bytes). These map to asyncpg ``server_settings``;
explicit entries in ``server_settings`` take precedence.
8 changes: 8 additions & 0 deletions docs/reference/adapters/cockroach_psycopg.rst
Original file line number Diff line number Diff line change
Expand Up @@ -254,3 +254,11 @@ namespace: ``"litestar"``, ``"events"``, or ``"adk"`` as supported by this adapt
.. autoclass:: sqlspec.adapters.cockroach_psycopg.adk.CockroachPsycopgADKConfig
:members:
:show-inheritance:

Native startup settings
-----------------------

``connection_config`` accepts ``default_transaction_use_follower_reads`` as a
boolean, plus non-negative integer ``results_buffer_size`` (bytes),
``statement_timeout`` and ``idle_in_transaction_session_timeout`` (milliseconds).
These settings append to the existing libpq ``options`` string.
9 changes: 9 additions & 0 deletions docs/reference/adapters/psqlpy.rst
Original file line number Diff line number Diff line change
Expand Up @@ -77,3 +77,12 @@ namespace: ``"litestar"``, ``"events"``, or ``"adk"`` as supported by this adapt
.. autoclass:: sqlspec.adapters.psqlpy.adk.PsqlpyADKConfig
:members:
:show-inheritance:

Dense vector parameters
-----------------------

Parameters explicitly cast to ``VECTOR`` accept lists, tuples, and objects
with ``tolist()``. SQLSpec wraps these values with psqlpy's native ``PgVector``;
already-wrapped values pass through. This conversion does not require the
Python ``pgvector`` package and does not apply the dense-vector encoder to
``HALFVEC`` or ``SPARSEVEC``.
13 changes: 13 additions & 0 deletions docs/reference/adapters/psycopg.rst
Original file line number Diff line number Diff line change
Expand Up @@ -110,3 +110,16 @@ namespace: ``"litestar"``, ``"events"``, or ``"adk"`` as supported by this adapt
.. autoclass:: sqlspec.adapters.psycopg.adk.PsycopgADKConfig
:members:
:show-inheritance:

Native null pools and JSON codecs
--------------------------------

Set ``null_pool=True`` in ``connection_config`` or ``driver_features`` to select
psycopg's native null pool. Returned connections close instead of remaining idle.
``max_size`` and ``max_waiting`` still control concurrency and waiting clients;
``min_size`` is omitted because null pools do not maintain idle connections.
The same option works for sync and async configurations.

``json_serializer`` and ``json_deserializer`` driver features also configure
psycopg's native JSON adaptation for each new connection. Codec setup errors
propagate to the caller.
18 changes: 13 additions & 5 deletions sqlspec/adapters/asyncpg/adk/store.py
Original file line number Diff line number Diff line change
Expand Up @@ -106,20 +106,28 @@ async def create_session(
INSERT INTO {self._session_table}
(id, app_name, user_id, {self._owner_id_column_name}, state, create_time, update_time)
VALUES ($1, $2, $3, $4, $5, CURRENT_TIMESTAMP, CURRENT_TIMESTAMP)
RETURNING id, app_name, user_id, state, create_time, update_time
"""
await conn.execute(sql, session_id, app_name, user_id, owner_id, state)
row = await conn.fetchrow(sql, session_id, app_name, user_id, owner_id, state)
else:
sql = f"""
INSERT INTO {self._session_table} (id, app_name, user_id, state, create_time, update_time)
VALUES ($1, $2, $3, $4, CURRENT_TIMESTAMP, CURRENT_TIMESTAMP)
RETURNING id, app_name, user_id, state, create_time, update_time
"""
await conn.execute(sql, session_id, app_name, user_id, state)
row = await conn.fetchrow(sql, session_id, app_name, user_id, state)

result = await self.get_session(app_name, user_id, session_id)
if result is None:
if row is None:
msg = "Failed to fetch created session"
raise RuntimeError(msg)
return result
return StoredSession(
id=row["id"],
app_name=row["app_name"],
user_id=row["user_id"],
state=row["state"],
create_time=row["create_time"],
update_time=row["update_time"],
)

async def get_session(
self, app_name: str, user_id: str, session_id: str, *, renew_for: "int | timedelta | None" = None
Expand Down
28 changes: 25 additions & 3 deletions sqlspec/adapters/asyncpg/config.py
Original file line number Diff line number Diff line change
Expand Up @@ -101,6 +101,7 @@ class AsyncpgConnectionConfig(TypedDict):
connect_timeout: NotRequired[float]
command_timeout: NotRequired[float]
statement_cache_size: NotRequired[int]
pgbouncer: NotRequired[bool]
max_cached_statement_lifetime: NotRequired[int]
max_cacheable_statement_size: NotRequired[int]
server_settings: NotRequired["dict[str, str]"]
Expand Down Expand Up @@ -191,6 +192,9 @@ class AsyncpgDriverFeatures(TypedDict):
- "notify_queue": Durable queue plus a PostgreSQL notification wakeup hint
- "poll_queue": Durable queue discovered by polling
Defaults to "notify".
pgbouncer: Enable PgBouncer transaction-pooling compatibility mode.
Disables server-side prepared statement caching (statement_cache_size=0).
type_codecs: Optional list of custom type codec specifications to register.
"""

json_serializer: NotRequired["Callable[[Any], str]"]
Expand All @@ -211,6 +215,8 @@ class AsyncpgDriverFeatures(TypedDict):
events_backend: NotRequired[Literal["notify", "notify_queue", "poll_queue"]]
connection_instance: NotRequired["AsyncpgPool"]
on_connection_create: NotRequired["Callable[[AsyncpgConnection], Awaitable[None]]"]
pgbouncer: NotRequired[bool]
type_codecs: NotRequired["list[dict[str, Any]]"]


class _AsyncpgCloudSqlConnector:
Expand Down Expand Up @@ -334,6 +340,9 @@ def __init__(
self._user_connection_hook: Callable[[AsyncpgConnection], Awaitable[None]] | None = features_dict.pop(
"on_connection_create", None
)
if connection_config and connection_config.get("pgbouncer"):
features_dict["pgbouncer"] = True
self._custom_type_codecs: list[dict[str, Any]] = list(features_dict.pop("type_codecs", None) or [])

super().__init__(
connection_config=build_connection_config(normalize_connection_config(connection_config)),
Expand Down Expand Up @@ -457,6 +466,9 @@ async def _create_pool(self) -> "Pool[Record]":
key: value for key, value in build_connection_config(self.connection_config).items() if value is not None
}

if self.connection_config.get("pgbouncer") or self.driver_features.get("pgbouncer"):
config["statement_cache_size"] = 0

if self.driver_features.get("enable_cloud_sql", False):
self._setup_cloud_sql_connector(config)
elif self.driver_features.get("enable_alloydb", False):
Expand All @@ -467,7 +479,7 @@ async def _create_pool(self) -> "Pool[Record]":
return await asyncpg_create_pool(**config)

async def _init_connection(self, connection: "AsyncpgConnection") -> None:
"""Initialize connection with JSON codecs, pgvector support, and user callback.
"""Initialize connection with JSON codecs, pgvector support, custom codecs, and user callback.

Args:
connection: AsyncPG connection to initialize.
Expand All @@ -479,7 +491,6 @@ async def _init_connection(self, connection: "AsyncpgConnection") -> None:
decoder=self.driver_features.get("json_deserializer", from_json),
)

# Detect extensions on first connection, update dialect
if self._pgvector_available is None:
detected_extensions: set[str] = set()
extensions = build_postgres_extension_probe_names(self.driver_features)
Expand All @@ -500,7 +511,15 @@ async def _init_connection(self, connection: "AsyncpgConnection") -> None:
if self._pgvector_available:
await register_pgvector_support(connection)

# Call user-provided callback after internal setup
for codec in self._custom_type_codecs:
codec_kwargs: dict[str, Any] = {
"schema": codec.get("schema", "public"),
"format": codec.get("format", "text"),
}
codec_kwargs["encoder"] = codec["encoder"]
codec_kwargs["decoder"] = codec["decoder"]
await connection.set_type_codec(codec["typename"], **codec_kwargs)

if self._user_connection_hook is not None:
await self._user_connection_hook(connection)

Expand Down Expand Up @@ -544,6 +563,9 @@ async def create_connection(self) -> "AsyncpgConnection":
for key in _POOL_ONLY_CONFIG_KEYS:
config.pop(key, None)

if self.driver_features.get("pgbouncer"):
config["statement_cache_size"] = 0

if self.driver_features.get("enable_cloud_sql", False):
self._setup_cloud_sql_connector(config)
elif self.driver_features.get("enable_alloydb", False):
Expand Down
38 changes: 26 additions & 12 deletions sqlspec/adapters/asyncpg/core.py
Original file line number Diff line number Diff line change
Expand Up @@ -131,6 +131,9 @@ def build_connection_config(connection_config: "Mapping[str, Any]") -> "dict[str
if found_user:
config["user"] = user_val

if config.pop("pgbouncer", False):
config["statement_cache_size"] = 0

return config


Expand Down Expand Up @@ -167,33 +170,34 @@ def configure_parameter_serializers(


async def invoke_prepared_statement(
prepared: Any, parameters: "tuple[Any, ...] | dict[str, Any] | list[Any] | None", *, fetch: bool
prepared: Any, parameters: "tuple[Any, ...] | dict[str, Any] | list[Any] | None", *, fetch: bool, **kwargs: Any
) -> Any:
"""Invoke an AsyncPG prepared statement with optional parameters.

Args:
prepared: AsyncPG prepared statement object.
parameters: Prepared parameters payload.
fetch: Whether to fetch rows.
**kwargs: Native execution options, including timeout.

Returns:
Query result or status message.
"""
if parameters is None:
if fetch:
return await prepared.fetch()
await prepared.fetch()
return await prepared.fetch(**kwargs)
await prepared.fetch(**kwargs)
return prepared.get_statusmsg()

if isinstance(parameters, dict):
if fetch:
return await prepared.fetch(**parameters)
await prepared.fetch(**parameters)
return await prepared.fetch(**parameters, **kwargs)
await prepared.fetch(**parameters, **kwargs)
return prepared.get_statusmsg()

if fetch:
return await prepared.fetch(*parameters)
await prepared.fetch(*parameters)
return await prepared.fetch(*parameters, **kwargs)
await prepared.fetch(*parameters, **kwargs)
return prepared.get_statusmsg()


Expand Down Expand Up @@ -426,9 +430,17 @@ def collect_rows(records: "list[Any] | None") -> "tuple[list[Any], list[str]]":
class AsyncpgStreamSource:
"""Compiled async chunk source streaming dict rows from an asyncpg cursor in a stream-owned transaction."""

__slots__ = ("_chunk_size", "_cursor", "_driver", "_parameters", "_sql", "_transaction")

def __init__(self, driver: Any, sql: str, parameters: "tuple[Any, ...]", chunk_size: int) -> None:
__slots__ = ("_chunk_size", "_cursor", "_driver", "_parameters", "_sql", "_timeout_args", "_transaction")

def __init__(
self,
driver: Any,
sql: str,
parameters: "tuple[Any, ...]",
chunk_size: int,
timeout_args: "dict[str, Any] | None" = None,
) -> None:
self._timeout_args = timeout_args or {}
self._driver = driver
self._sql = sql
self._parameters = parameters
Expand All @@ -446,15 +458,17 @@ async def _start(self) -> None:
await transaction.start()
self._transaction = transaction
try:
self._cursor = await self._driver.connection.cursor(self._sql, *self._parameters)
self._cursor = await self._driver.connection.cursor(self._sql, *self._parameters, **self._timeout_args)
except BaseException:
await transaction.rollback()
self._transaction = None
raise

async def fetch_chunk(self) -> "list[dict[str, Any]]":
handler = self._driver.handle_database_exceptions()
records = await self._driver._run_with_exception_handler(handler, self._cursor.fetch, self._chunk_size)
records = await self._driver._run_with_exception_handler(
handler, self._cursor.fetch, self._chunk_size, **self._timeout_args
)
self._driver._check_pending_exception(handler)
assert records is not None
return [dict(record) for record in records]
Expand Down
Loading
Loading