Source code for sqlspec.adapters.cockroach_psycopg.config

"""CockroachDB configuration using psycopg."""

import re
from contextlib import suppress
from typing import TYPE_CHECKING, Any, ClassVar, Final, Literal, TypedDict, cast

from typing_extensions import NotRequired

from sqlspec.adapters.cockroach_psycopg._typing import (
    CockroachAsyncConnection,
    CockroachPsycopgAsyncSessionContext,
    CockroachPsycopgSyncSessionContext,
    CockroachSyncConnection,
)
from sqlspec.adapters.cockroach_psycopg._typing import CockroachPsycopgAsyncConnectionPool as AsyncConnectionPool
from sqlspec.adapters.cockroach_psycopg._typing import CockroachPsycopgConnectionPool as ConnectionPool
from sqlspec.adapters.cockroach_psycopg._typing import cockroach_psycopg_crdb as psycopg_crdb
from sqlspec.adapters.cockroach_psycopg.core import (
    apply_driver_features,
    build_connection_config,
    build_statement_config,
    validate_follower_read_staleness,
)
from sqlspec.adapters.cockroach_psycopg.driver import (
    CockroachPsycopgAsyncDriver,
    CockroachPsycopgAsyncExceptionHandler,
    CockroachPsycopgSyncDriver,
    CockroachPsycopgSyncExceptionHandler,
)
from sqlspec.config import AsyncDatabaseConfig, ExtensionConfigs, SyncDatabaseConfig
from sqlspec.core.capabilities import TypeCoercionCapabilities
from sqlspec.core.config_runtime import resolve_runtime_statement_config
from sqlspec.driver import (
    AsyncPoolConnectionContext,
    AsyncPoolSessionFactory,
    SyncPoolConnectionContext,
    SyncPoolSessionFactory,
)
from sqlspec.exceptions import ImproperConfigurationError
from sqlspec.extensions.events import EventRuntimeHints

if TYPE_CHECKING:
    from collections.abc import Awaitable, Callable
    from types import TracebackType

    from sqlspec.core import StatementConfig
    from sqlspec.observability import ObservabilityConfig

__all__ = (
    "CockroachPsycopgAsyncConfig",
    "CockroachPsycopgConnectionConfig",
    "CockroachPsycopgDriverFeatures",
    "CockroachPsycopgPoolConfig",
    "CockroachPsycopgSyncConfig",
    "build_connection_config",
)

default_statement_config = build_statement_config()


[docs] class CockroachPsycopgConnectionConfig(TypedDict): """CockroachDB connection parameters.""" conninfo: NotRequired[str] host: NotRequired[str] port: NotRequired[int] user: NotRequired[str] password: NotRequired[str] dbname: NotRequired[str] connect_timeout: NotRequired[int] options: NotRequired[str] default_transaction_use_follower_reads: NotRequired[bool] results_buffer_size: NotRequired[int] statement_timeout: NotRequired[int] idle_in_transaction_session_timeout: NotRequired[int] application_name: NotRequired[str] sslmode: NotRequired[str] sslcert: NotRequired[str] sslkey: NotRequired[str] sslrootcert: NotRequired[str] autocommit: NotRequired[bool] prepare_threshold: NotRequired[int | None] row_factory: NotRequired["Callable[..., Any]"] cursor_factory: NotRequired["type[Any]"] context: NotRequired[Any] cluster: NotRequired[str] extra: NotRequired["dict[str, Any]"]
[docs] class CockroachPsycopgPoolConfig(CockroachPsycopgConnectionConfig): """CockroachDB pool parameters.""" min_size: NotRequired[int] max_size: NotRequired[int] name: NotRequired[str] timeout: NotRequired[float] max_waiting: NotRequired[int] max_lifetime: NotRequired[float] max_idle: NotRequired[float] reconnect_timeout: NotRequired[float] reconnect_failed: NotRequired["Callable[..., Any]"] num_workers: NotRequired[int] configure: NotRequired["Callable[..., Any]"] check: NotRequired["Callable[..., Any]"] reset: NotRequired["Callable[..., Any]"] open: NotRequired[bool | None] close_returns: NotRequired[bool] kwargs: NotRequired["dict[str, Any]"]
class _NativeStorageCSVOptions(TypedDict): """Explicit CSV conventions; no header or NULL marker is inferred.""" nullas: NotRequired[str] nullif: NotRequired[str] skip: NotRequired[int]
[docs] class CockroachPsycopgDriverFeatures(TypedDict): """CockroachDB driver feature configuration. enable_native_storage: Opt in to server-side storage operations. Exports create generated files under a prefix; imports temporarily take the table offline. native_storage_csv_options: Explicit nullas/nullif markers and skip count. Export is headerless; native CSV import requires an explicit skip count. Markers must not occur as literal data. No NULL marker is inferred. on_connection_create: Callback executed when a connection is acquired from pool. For sync: Callable[[CockroachSyncConnection], None] For async: Callable[[CockroachAsyncConnection], Awaitable[None]] Called after internal setup. """ enable_native_storage: NotRequired[bool] native_storage_csv_options: NotRequired[_NativeStorageCSVOptions] enable_auto_retry: NotRequired[bool] max_retries: NotRequired[int] retry_delay_base_ms: NotRequired[float] retry_delay_max_ms: NotRequired[float] enable_retry_logging: NotRequired[bool] enable_follower_reads: NotRequired[bool] default_staleness: NotRequired[str] json_serializer: NotRequired["Callable[[Any], str]"] json_deserializer: NotRequired["Callable[[str], Any]"] on_connection_create: "NotRequired[Callable[..., Any]]" enable_events: NotRequired[bool] events_backend: NotRequired[Literal["poll_queue"]]
class CockroachPsycopgSyncConnectionContext(SyncPoolConnectionContext): """Context manager for CockroachDB psycopg connections.""" __slots__ = () def __init__(self, config: "CockroachPsycopgSyncConfig") -> None: super().__init__(config) def __enter__(self) -> "CockroachSyncConnection": self._ctx = self._config.provide_pool().connection() return cast("CockroachSyncConnection", self._ctx.__enter__()) def __exit__( self, exc_type: "type[BaseException] | None", exc_val: "BaseException | None", exc_tb: "TracebackType | None" ) -> bool | None: if self._ctx is None: return None return cast("bool | None", self._ctx.__exit__(exc_type, exc_val, exc_tb)) class _CockroachPsycopgSyncSessionConnectionHandler(SyncPoolSessionFactory): __slots__ = () def acquire_connection(self) -> "CockroachSyncConnection": self._ctx = self._config.provide_pool().connection() return cast("CockroachSyncConnection", self._ctx.__enter__()) def release_connection(self, _conn: "CockroachSyncConnection", **kwargs: Any) -> None: if self._ctx is not None: self._ctx.__exit__(kwargs.get("exc_type"), kwargs.get("exc_val"), kwargs.get("exc_tb")) self._ctx = None
[docs] class CockroachPsycopgSyncConfig( SyncDatabaseConfig[CockroachSyncConnection, ConnectionPool, CockroachPsycopgSyncDriver] ): """Configuration for CockroachDB synchronous connections using psycopg.""" driver_type: "ClassVar[type[CockroachPsycopgSyncDriver]]" = CockroachPsycopgSyncDriver connection_type: "ClassVar[type[CockroachSyncConnection]]" = CockroachSyncConnection supports_transactional_ddl: "ClassVar[bool]" = False supports_migration_schemas: "ClassVar[bool]" = True supports_native_arrow_export: "ClassVar[bool]" = True supports_native_arrow_import: "ClassVar[bool]" = True supports_native_parquet_export: "ClassVar[bool]" = True supports_native_parquet_import: "ClassVar[bool]" = True supports_native_row_streaming: "ClassVar[bool]" = True type_coercion_capabilities: "ClassVar[TypeCoercionCapabilities]" = TypeCoercionCapabilities( datetime_binding="native", timestamp_precision="microsecond", json_columns_decoded=True, uuid_binding="native" ) _connection_context_class: "ClassVar[type[CockroachPsycopgSyncConnectionContext]]" = ( CockroachPsycopgSyncConnectionContext ) _session_factory_class: "ClassVar[type[_CockroachPsycopgSyncSessionConnectionHandler]]" = ( _CockroachPsycopgSyncSessionConnectionHandler ) _session_context_class: "ClassVar[type[CockroachPsycopgSyncSessionContext]]" = CockroachPsycopgSyncSessionContext _default_statement_config = default_statement_config
[docs] def __init__( self, *, connection_config: "CockroachPsycopgPoolConfig | dict[str, Any] | None" = None, connection_instance: "ConnectionPool | None" = None, migration_config: "dict[str, Any] | None" = None, statement_config: "StatementConfig | None" = None, driver_features: "CockroachPsycopgDriverFeatures | dict[str, Any] | None" = None, bind_key: "str | None" = None, extension_config: "ExtensionConfigs | None" = None, observability_config: "ObservabilityConfig | None" = None, **kwargs: Any, ) -> None: connection_config = build_connection_config(connection_config) statement_config = statement_config or default_statement_config _validate_driver_features(driver_features) statement_config, driver_features = apply_driver_features(statement_config, driver_features) driver_features.setdefault("enable_native_storage", False) _validate_native_storage_options(driver_features) _validate_follower_read_features(driver_features) _validate_native_storage_autocommit(connection_config, driver_features) driver_features.setdefault("enable_auto_retry", True) features_dict = dict(driver_features) if driver_features else {} self._user_connection_hook: Callable[[CockroachSyncConnection], None] | None = features_dict.pop( "on_connection_create", None ) super().__init__( connection_config=connection_config, connection_instance=connection_instance, migration_config=migration_config, statement_config=statement_config, driver_features=features_dict, bind_key=bind_key, extension_config=extension_config, observability_config=observability_config, **kwargs, )
def _create_pool(self) -> "ConnectionPool": all_config = dict(self.connection_config) pool_parameters = { "min_size": all_config.pop("min_size", 4), "max_size": all_config.pop("max_size", None), "name": all_config.pop("name", None), "timeout": all_config.pop("timeout", 30.0), "max_waiting": all_config.pop("max_waiting", 0), "max_lifetime": all_config.pop("max_lifetime", 3600.0), "max_idle": all_config.pop("max_idle", 600.0), "reconnect_timeout": all_config.pop("reconnect_timeout", 300.0), "reconnect_failed": all_config.pop("reconnect_failed", None), "num_workers": all_config.pop("num_workers", 3), "check": all_config.pop("check", None), "reset": all_config.pop("reset", None), "open": all_config.pop("open", True), "close_returns": all_config.pop("close_returns", False), } pool_parameters["configure"] = all_config.pop("configure", self._configure_connection) pool_parameters = {k: v for k, v in pool_parameters.items() if v is not None} conninfo = all_config.pop("conninfo", None) kwargs = all_config.pop("kwargs", {}) all_config.update(kwargs) _apply_cluster_option(all_config) if conninfo: return ConnectionPool( conninfo, kwargs=all_config, connection_class=psycopg_crdb.CrdbConnection, **pool_parameters ) return ConnectionPool("", kwargs=all_config, connection_class=psycopg_crdb.CrdbConnection, **pool_parameters) def _configure_connection(self, conn: "CockroachSyncConnection") -> None: autocommit_setting = self.connection_config.get("autocommit") if autocommit_setting is not None: conn.autocommit = autocommit_setting # Call user-provided callback after internal setup if self._user_connection_hook is not None: self._user_connection_hook(conn) def _close_pool(self) -> None: if not self.connection_instance: return try: self.connection_instance.close() finally: self.connection_instance = None
[docs] def create_connection(self) -> "CockroachSyncConnection": """Open a standalone connection owned by the caller. Returns: A connection that is not bound to the pool and can be closed directly. """ conninfo, connection_kwargs = _standalone_connection_kwargs(dict(self.connection_config)) connection = psycopg_crdb.CrdbConnection.connect(conninfo, **connection_kwargs) configure = self.connection_config.get("configure", self._configure_connection) try: configure(cast("CockroachSyncConnection", connection)) except BaseException: with suppress(Exception): connection.close() raise return cast("CockroachSyncConnection", connection)
[docs] def provide_session( self, *_args: Any, statement_config: "StatementConfig | None" = None, follower_reads: bool | None = None, staleness: str | None = None, **_kwargs: Any, ) -> "CockroachPsycopgSyncSessionContext": handler = _CockroachPsycopgSyncSessionConnectionHandler(self) driver_features = dict(self.driver_features) if follower_reads is not None: driver_features["enable_follower_reads"] = follower_reads if staleness is not None: driver_features["default_staleness"] = validate_follower_read_staleness(staleness) return CockroachPsycopgSyncSessionContext( acquire_connection=handler.acquire_connection, release_connection=handler.release_connection, statement_config=statement_config or (lambda: resolve_runtime_statement_config(None, self.statement_config, default_statement_config)), driver_features=driver_features, prepare_driver=self._prepare_driver, )
[docs] def provide_pool(self, *args: Any, **kwargs: Any) -> "ConnectionPool": if not self.connection_instance: self.connection_instance = self.create_pool() return self.connection_instance
[docs] def get_signature_namespace(self) -> "dict[str, Any]": namespace = super().get_signature_namespace() namespace.update({ "CockroachPsycopgConnectionConfig": CockroachPsycopgConnectionConfig, "CockroachPsycopgPoolConfig": CockroachPsycopgPoolConfig, "CockroachSyncConnection": CockroachSyncConnection, "CockroachPsycopgSyncDriver": CockroachPsycopgSyncDriver, "CockroachPsycopgSyncExceptionHandler": CockroachPsycopgSyncExceptionHandler, "CockroachPsycopgSyncSessionContext": CockroachPsycopgSyncSessionContext, }) return namespace
[docs] def get_event_runtime_hints(self) -> "EventRuntimeHints": return EventRuntimeHints(poll_interval=0.5, select_for_update=True, skip_locked=True)
class CockroachPsycopgAsyncConnectionContext(AsyncPoolConnectionContext): """Async context manager for CockroachDB psycopg connections.""" __slots__ = ("_ctx",) def __init__(self, config: "CockroachPsycopgAsyncConfig") -> None: super().__init__(config) self._ctx: Any = None async def __aenter__(self) -> "CockroachAsyncConnection": if self._config.connection_instance is None: self._config.connection_instance = await self._config.create_pool() if self._config.connection_instance: self._ctx = self._config.connection_instance.connection() return cast("CockroachAsyncConnection", await self._ctx.__aenter__()) msg = "Connection pool is not initialized" raise ImproperConfigurationError(msg) async def __aexit__( self, exc_type: "type[BaseException] | None", exc_val: "BaseException | None", exc_tb: "TracebackType | None" ) -> bool | None: if self._ctx: return cast("bool | None", await self._ctx.__aexit__(exc_type, exc_val, exc_tb)) return None class _CockroachPsycopgAsyncSessionConnectionHandler(AsyncPoolSessionFactory): __slots__ = ("_ctx",) def __init__(self, config: "CockroachPsycopgAsyncConfig") -> None: super().__init__(config) self._ctx: Any = None async def acquire_connection(self) -> "CockroachAsyncConnection": pool = self._config.connection_instance if pool is None: pool = await self._config.create_pool() self._config.connection_instance = pool ctx = pool.connection() self._ctx = ctx return cast("CockroachAsyncConnection", await ctx.__aenter__()) async def release_connection(self, _conn: "CockroachAsyncConnection", **kwargs: Any) -> None: if self._ctx is not None: await self._ctx.__aexit__(kwargs.get("exc_type"), kwargs.get("exc_val"), kwargs.get("exc_tb")) self._ctx = None
[docs] class CockroachPsycopgAsyncConfig( AsyncDatabaseConfig[CockroachAsyncConnection, AsyncConnectionPool, CockroachPsycopgAsyncDriver] ): """Configuration for CockroachDB async connections using psycopg.""" driver_type: "ClassVar[type[CockroachPsycopgAsyncDriver]]" = CockroachPsycopgAsyncDriver connection_type: "ClassVar[type[CockroachAsyncConnection]]" = CockroachAsyncConnection supports_transactional_ddl: "ClassVar[bool]" = False supports_migration_schemas: "ClassVar[bool]" = True supports_native_arrow_export: "ClassVar[bool]" = True supports_native_arrow_import: "ClassVar[bool]" = True supports_native_parquet_export: "ClassVar[bool]" = True supports_native_parquet_import: "ClassVar[bool]" = True supports_native_row_streaming: "ClassVar[bool]" = True type_coercion_capabilities: "ClassVar[TypeCoercionCapabilities]" = TypeCoercionCapabilities( datetime_binding="native", timestamp_precision="microsecond", json_columns_decoded=True, uuid_binding="native" ) _connection_context_class: "ClassVar[type[CockroachPsycopgAsyncConnectionContext]]" = ( CockroachPsycopgAsyncConnectionContext ) _session_factory_class: "ClassVar[type[_CockroachPsycopgAsyncSessionConnectionHandler]]" = ( _CockroachPsycopgAsyncSessionConnectionHandler ) _session_context_class: "ClassVar[type[CockroachPsycopgAsyncSessionContext]]" = CockroachPsycopgAsyncSessionContext _default_statement_config = default_statement_config
[docs] def __init__( self, *, connection_config: "CockroachPsycopgPoolConfig | dict[str, Any] | None" = None, connection_instance: "AsyncConnectionPool | None" = None, migration_config: "dict[str, Any] | None" = None, statement_config: "StatementConfig | None" = None, driver_features: "CockroachPsycopgDriverFeatures | dict[str, Any] | None" = None, bind_key: "str | None" = None, extension_config: "ExtensionConfigs | None" = None, observability_config: "ObservabilityConfig | None" = None, **kwargs: Any, ) -> None: connection_config = build_connection_config(connection_config) statement_config = statement_config or default_statement_config _validate_driver_features(driver_features) statement_config, driver_features = apply_driver_features(statement_config, driver_features) driver_features.setdefault("enable_native_storage", False) _validate_native_storage_options(driver_features) _validate_follower_read_features(driver_features) _validate_native_storage_autocommit(connection_config, driver_features) driver_features.setdefault("enable_auto_retry", True) features_dict = dict(driver_features) if driver_features else {} self._user_connection_hook: Callable[[CockroachAsyncConnection], Awaitable[None]] | None = features_dict.pop( "on_connection_create", None ) super().__init__( connection_config=connection_config, connection_instance=connection_instance, migration_config=migration_config, statement_config=statement_config, driver_features=features_dict, bind_key=bind_key, extension_config=extension_config, observability_config=observability_config, **kwargs, )
async def _create_pool(self) -> "AsyncConnectionPool": all_config = dict(self.connection_config) pool_parameters = { "min_size": all_config.pop("min_size", 4), "max_size": all_config.pop("max_size", None), "name": all_config.pop("name", None), "timeout": all_config.pop("timeout", 30.0), "max_waiting": all_config.pop("max_waiting", 0), "max_lifetime": all_config.pop("max_lifetime", 3600.0), "max_idle": all_config.pop("max_idle", 600.0), "reconnect_timeout": all_config.pop("reconnect_timeout", 300.0), "reconnect_failed": all_config.pop("reconnect_failed", None), "num_workers": all_config.pop("num_workers", 3), "check": all_config.pop("check", None), "reset": all_config.pop("reset", None), "close_returns": all_config.pop("close_returns", False), } open_pool = all_config.pop("open", True) pool_parameters["open"] = False if open_pool is True else open_pool pool_parameters["configure"] = all_config.pop("configure", self._configure_async_connection) pool_parameters = {k: v for k, v in pool_parameters.items() if v is not None} conninfo = all_config.pop("conninfo", None) kwargs = all_config.pop("kwargs", {}) all_config.update(kwargs) _apply_cluster_option(all_config) if conninfo: pool = AsyncConnectionPool( conninfo, kwargs=all_config, connection_class=psycopg_crdb.AsyncCrdbConnection, **pool_parameters ) else: pool = AsyncConnectionPool( "", kwargs=all_config, connection_class=psycopg_crdb.AsyncCrdbConnection, **pool_parameters ) if open_pool is True: await pool.open() return cast("AsyncConnectionPool", pool) async def _configure_async_connection(self, conn: "CockroachAsyncConnection") -> None: autocommit_setting = self.connection_config.get("autocommit") if autocommit_setting is not None: await conn.set_autocommit(autocommit_setting) # Call user-provided callback after internal setup if self._user_connection_hook is not None: await self._user_connection_hook(conn) async def _close_pool(self) -> None: if not self.connection_instance: return try: await self.connection_instance.close() finally: self.connection_instance = None
[docs] async def create_connection(self) -> "CockroachAsyncConnection": """Open a standalone connection owned by the caller. Returns: A connection that is not bound to the pool and can be closed directly. """ conninfo, connection_kwargs = _standalone_connection_kwargs(dict(self.connection_config)) connection = await psycopg_crdb.AsyncCrdbConnection.connect(conninfo, **connection_kwargs) configure = self.connection_config.get("configure", self._configure_async_connection) try: await configure(cast("CockroachAsyncConnection", connection)) except BaseException: with suppress(Exception): await connection.close() raise return cast("CockroachAsyncConnection", connection)
[docs] def provide_session( self, *_args: Any, statement_config: "StatementConfig | None" = None, follower_reads: bool | None = None, staleness: str | None = None, **_kwargs: Any, ) -> "CockroachPsycopgAsyncSessionContext": handler = _CockroachPsycopgAsyncSessionConnectionHandler(self) driver_features = dict(self.driver_features) if follower_reads is not None: driver_features["enable_follower_reads"] = follower_reads if staleness is not None: driver_features["default_staleness"] = validate_follower_read_staleness(staleness) return CockroachPsycopgAsyncSessionContext( acquire_connection=handler.acquire_connection, release_connection=handler.release_connection, statement_config=statement_config or (lambda: resolve_runtime_statement_config(None, self.statement_config, default_statement_config)), driver_features=driver_features, prepare_driver=self._prepare_driver, )
[docs] async def provide_pool(self, *args: Any, **kwargs: Any) -> "AsyncConnectionPool": if not self.connection_instance: self.connection_instance = await self.create_pool() return self.connection_instance
[docs] def get_signature_namespace(self) -> "dict[str, Any]": namespace = super().get_signature_namespace() namespace.update({ "CockroachAsyncConnection": CockroachAsyncConnection, "CockroachPsycopgAsyncDriver": CockroachPsycopgAsyncDriver, "CockroachPsycopgAsyncExceptionHandler": CockroachPsycopgAsyncExceptionHandler, "CockroachPsycopgAsyncSessionContext": CockroachPsycopgAsyncSessionContext, "CockroachPsycopgConnectionConfig": CockroachPsycopgConnectionConfig, "CockroachPsycopgPoolConfig": CockroachPsycopgPoolConfig, }) return namespace
[docs] def get_event_runtime_hints(self) -> "EventRuntimeHints": return EventRuntimeHints(poll_interval=0.5, select_for_update=True, skip_locked=True)
_POOL_ONLY_CONFIG_KEYS: Final[frozenset[str]] = frozenset({ "check", "close_returns", "configure", "max_idle", "max_lifetime", "max_size", "max_waiting", "min_size", "name", "num_workers", "open", "reconnect_failed", "reconnect_timeout", "reset", "timeout", }) def _standalone_connection_kwargs(connection_config: "dict[str, Any]") -> "tuple[str, dict[str, Any]]": """Split a pooled connection config into conninfo plus connect keyword arguments. Args: connection_config: The adapter's normalized connection configuration. Returns: The conninfo string and the keyword arguments accepted by ``connect``. """ all_config = {k: v for k, v in connection_config.items() if k not in _POOL_ONLY_CONFIG_KEYS} conninfo = all_config.pop("conninfo", "") all_config.update(all_config.pop("kwargs", {})) all_config.update(all_config.pop("extra", {})) _apply_cluster_option(all_config) return str(conninfo or ""), all_config _CLUSTER_NAME_PATTERN: "re.Pattern[str]" = re.compile(r"[A-Za-z0-9._-]+") def _apply_cluster_option(connection_kwargs: "dict[str, Any]") -> None: """Fold a CockroachDB Cloud cluster name into the libpq options string. ``cluster`` is not a libpq connection parameter; multi-tenant routing is expressed through ``options=--cluster=<name>``. Raises: ImproperConfigurationError: If the name carries anything but the characters a cluster name may contain. The value is interpolated into an options string, where whitespace would start a new option. """ cluster = connection_kwargs.pop("cluster", None) if not cluster: return name = str(cluster) if _CLUSTER_NAME_PATTERN.fullmatch(name) is None: msg = f"CockroachDB cluster name {name!r} may contain only letters, digits, hyphens, underscores, and periods." raise ImproperConfigurationError(msg) routing = f"--cluster={name}" existing = connection_kwargs.get("options") connection_kwargs["options"] = f"{existing} {routing}" if existing else routing def _validate_driver_features(driver_features: "CockroachPsycopgDriverFeatures | dict[str, Any] | None") -> None: if driver_features and "prefer_uuid_keys" in driver_features: msg = "CockroachDB psycopg driver_features no longer supports unused 'prefer_uuid_keys'." raise ImproperConfigurationError(msg) def _validate_native_storage_autocommit(connection_config: "dict[str, Any]", driver_features: "dict[str, Any]") -> None: if not driver_features.get("enable_native_storage"): return autocommit = connection_config.get("autocommit") if autocommit is None: autocommit = connection_config.get("kwargs", {}).get("autocommit") if autocommit is not True: msg = "enable_native_storage requires connection_config autocommit=True." raise ImproperConfigurationError(msg) def _validate_follower_read_features(driver_features: "dict[str, Any]") -> None: staleness = driver_features.get("default_staleness") if staleness is None: return if not isinstance(staleness, str): msg = "default_staleness must be a string." raise ImproperConfigurationError(msg) driver_features["default_staleness"] = validate_follower_read_staleness(staleness) def _validate_native_storage_options(driver_features: "dict[str, Any]") -> None: if "native_storage_csv_options" not in driver_features: return options = driver_features["native_storage_csv_options"] if not isinstance(options, dict) or options.keys() - {"nullas", "nullif", "skip"}: msg = "native_storage_csv_options must contain only nullas, nullif, and skip." raise ImproperConfigurationError(msg) for key in ("nullas", "nullif"): if key in options and not isinstance(options[key], str): msg = "native_storage_csv_options nullas and nullif must be strings." raise ImproperConfigurationError(msg) if "skip" in options and (type(options["skip"]) is not int or options["skip"] < 0): msg = "native_storage_csv_options skip must be a nonnegative integer, excluding bool." raise ImproperConfigurationError(msg) driver_features["native_storage_csv_options"] = dict(options)