Source code for sqlspec.adapters.oracledb.events.store

"""Oracle event queue stores with auto-detected JSON storage.

JSON storage is automatically detected based on Oracle version:
    - json: Native JSON type (Oracle 21c+ with COMPATIBLE >= 20)
    - blob_json: BLOB with IS JSON constraint (Oracle 12c+, recommended)
    - blob: Plain BLOB without constraint (Oracle 11g and earlier)

Note: CLOB should be avoided for JSON storage in Oracle. Oracle recommends
BLOB over CLOB for JSON data as BLOB performs significantly better.

Configuration (optional override):
    extension_config={
        "events": {
            "json_storage": "blob_json",
            "in_memory": False,
        }
    }
"""

import logging
from typing import TYPE_CHECKING, Any, Literal

from typing_extensions import NotRequired, TypedDict

from sqlspec.adapters.oracledb._storage import oracle_table_feature_report
from sqlspec.adapters.oracledb.data_dictionary import JSONStorageType, storage_type_from_version
from sqlspec.config import EventsConfig
from sqlspec.extensions.events import BaseEventQueueStore
from sqlspec.utils.logging import get_logger, log_with_context

if TYPE_CHECKING:
    from sqlspec.adapters.oracledb.config import OracleAsyncConfig, OracleSyncConfig
    from sqlspec.adapters.oracledb.data_dictionary import OracleVersionInfo

__all__ = (
    "OracleAsyncEventQueueStore",
    "OracleEventsCompressionConfig",
    "OracleEventsConfig",
    "OracleEventsPartitionConfig",
    "OracleSyncEventQueueStore",
)

logger = get_logger("sqlspec.adapters.oracledb.events.store")


[docs] class OracleEventsCompressionConfig(TypedDict): """Oracle queue-table compression settings.""" enabled: NotRequired[bool] """Enable compression. Default: False.""" algorithm: NotRequired[str] """Oracle compression algorithm key. Default: advanced."""
[docs] class OracleEventsPartitionConfig(TypedDict): """Oracle queue-table hash or range partition settings.""" strategy: NotRequired[str] """Partitioning strategy: hash or range; omitted to disable partitioning.""" partition_count: NotRequired[int] """Hash partition count, at least two. Default: 16.""" partitions: NotRequired[int] """Alternative name for partition_count.""" interval: NotRequired[str] """Range interval: day, week, month or year. Default: month.""" initial_less_than: NotRequired[str] """Initial range upper bound. Default: TIMESTAMP '2000-01-01 00:00:00'.""" partition_key: NotRequired[str] """Partition key; defaults to event_id for hash and available_at for range.""" queue_partition_key: NotRequired[str] """Queue-specific override for partition_key."""
[docs] class OracleEventsConfig(EventsConfig): """Oracle events settings for queue storage and supported native transports.""" in_memory: NotRequired[bool] """Emit Oracle INMEMORY PRIORITY HIGH. Default: False.""" partitioning: NotRequired[OracleEventsPartitionConfig] """Oracle hash/range queue partition settings; omitted for an unpartitioned table.""" compression: NotRequired[OracleEventsCompressionConfig] """Oracle queue-table compression; disabled unless explicitly enabled.""" table_options: NotRequired[str] """Additional Oracle queue-table DDL options; omitted by default.""" json_storage: NotRequired[Literal["json", "blob_json", "blob_plain"]] """Override Oracle JSON storage detection with json, blob_json or blob_plain.""" aq_queue: NotRequired[str] """Oracle AQ queue name or channel template. Default: SQLSPEC_EVENTS_QUEUE.""" aq_visibility: NotRequired[int | Literal["DEQ_IMMEDIATE", "DEQ_ON_COMMIT"]] """Oracle AQ visibility integer or DEQ_IMMEDIATE/DEQ_ON_COMMIT name.""" aq_wait_seconds: NotRequired[int] """Oracle AQ dequeue wait in seconds. Default: 5."""
class OracleSyncEventQueueStore(BaseEventQueueStore["OracleSyncConfig"]): """Oracle sync event queue store with auto-detected JSON storage. Automatically detects the Oracle version and uses the optimal JSON storage: - Oracle 21c+: Native JSON type - Oracle 12c-19c: BLOB with IS JSON constraint - Oracle 11g: Plain BLOB Args: config: OracleSyncConfig with extension_config["events"] settings. """ __slots__ = ("_in_memory", "_json_storage_override") extension_config_options = BaseEventQueueStore.extension_config_options | frozenset({ "aq_queue", "aq_visibility", "aq_wait_seconds", "compression", "in_memory", "json_storage", "partitioning", "table_options", }) def __init__(self, config: "OracleSyncConfig") -> None: """Initialize Oracle sync event queue store.""" super().__init__(config) self._in_memory, self._json_storage_override = _init_oracle_settings(self._extension_settings) def create_statements(self) -> "list[str]": """Return single PL/SQL script for table and index creation. Uses the configured storage override when set, otherwise defaults to BLOB_JSON. For auto-detection, use create_table() instead. """ storage_type = self._json_storage_override or JSONStorageType.BLOB_JSON return [_oracle_table_ddl(self.table_name, storage_type, self._table_feature_clause(), self._index_name())] def drop_statements(self) -> "list[str]": """Return drop statements in reverse dependency order.""" return _oracle_drop_sql(self.table_name, self._index_name()) def create_table(self) -> None: """Create the event queue table with auto-detected storage type.""" storage_type = self._detect_json_storage_type() log_with_context( logger, logging.DEBUG, "events.queue.create", storage_type=storage_type.value, table_name=self.table_name ) with self._config.provide_session() as driver: sql = _oracle_table_ddl(self.table_name, storage_type, self._table_feature_clause(), self._index_name()) driver.execute_script(sql) def drop_table(self) -> None: """Drop the event queue table and index.""" with self._config.provide_session() as driver: for stmt in _oracle_drop_sql(self.table_name, self._index_name()): driver.execute_script(stmt) def _column_types(self) -> "tuple[str, str, str]": """Return Oracle column types based on storage mode.""" storage = self._json_storage_override or JSONStorageType.BLOB_JSON if storage == JSONStorageType.JSON_NATIVE: return "JSON", "JSON", "TIMESTAMP" return "BLOB", "BLOB", "TIMESTAMP" def _string_type(self, length: int) -> str: """Return Oracle VARCHAR2 type syntax.""" return f"VARCHAR2({length})" def _index_name(self) -> str: """Return index name truncated to Oracle's 30-character limit.""" base_name = f"idx_{self.table_name.replace('.', '_')}_channel_status" return base_name[:30] def _detect_json_storage_type(self) -> JSONStorageType: """Resolve the JSON storage type from the configured override or server version.""" if self._json_storage_override is not None: return self._json_storage_override return storage_type_from_version(self._get_version_info()) def _get_version_info(self) -> "OracleVersionInfo | None": """Return the pool-scoped Oracle version through the data dictionary.""" cache = self._config._oracle_version_cache if cache.resolved: return cache.version with self._config.provide_session() as driver: version_info = driver.data_dictionary.get_version(driver) if version_info is None: log_with_context( logger, logging.WARNING, "events.queue.storage.fallback", storage_type=JSONStorageType.BLOB_JSON.value, reason="version_detection_failed", ) return version_info def _table_feature_clause(self) -> str: report = oracle_table_feature_report( self._config, "events", self._extension_settings, "queue", in_memory=self._in_memory, hash_partition_key="event_id", range_partition_key="available_at", table_options_key="table_options", ) return report["clause"] class OracleAsyncEventQueueStore(BaseEventQueueStore["OracleAsyncConfig"]): """Oracle async event queue store with auto-detected JSON storage. Automatically detects the Oracle version and uses the optimal JSON storage: - Oracle 21c+: Native JSON type - Oracle 12c-19c: BLOB with IS JSON constraint - Oracle 11g: Plain BLOB Args: config: OracleAsyncConfig with extension_config["events"] settings. """ __slots__ = ("_in_memory", "_json_storage_override") extension_config_options = BaseEventQueueStore.extension_config_options | frozenset({ "aq_queue", "aq_visibility", "aq_wait_seconds", "compression", "in_memory", "json_storage", "partitioning", "table_options", }) def __init__(self, config: "OracleAsyncConfig") -> None: """Initialize Oracle async event queue store.""" super().__init__(config) self._in_memory, self._json_storage_override = _init_oracle_settings(self._extension_settings) def create_statements(self) -> "list[str]": """Return single PL/SQL script for table and index creation. Uses the configured storage override when set, otherwise defaults to BLOB_JSON. For auto-detection, use create_table() instead. """ storage_type = self._json_storage_override or JSONStorageType.BLOB_JSON return [_oracle_table_ddl(self.table_name, storage_type, self._table_feature_clause(), self._index_name())] def drop_statements(self) -> "list[str]": """Return drop statements in reverse dependency order.""" return _oracle_drop_sql(self.table_name, self._index_name()) async def create_table(self) -> None: """Create the event queue table with auto-detected storage type.""" storage_type = await self._detect_json_storage_type() log_with_context( logger, logging.DEBUG, "events.queue.create", storage_type=storage_type.value, table_name=self.table_name ) async with self._config.provide_session() as driver: sql = _oracle_table_ddl(self.table_name, storage_type, self._table_feature_clause(), self._index_name()) await driver.execute_script(sql) async def drop_table(self) -> None: """Drop the event queue table and index.""" async with self._config.provide_session() as driver: for stmt in _oracle_drop_sql(self.table_name, self._index_name()): await driver.execute_script(stmt) def _column_types(self) -> "tuple[str, str, str]": """Return Oracle column types based on storage mode.""" storage = self._json_storage_override or JSONStorageType.BLOB_JSON if storage == JSONStorageType.JSON_NATIVE: return "JSON", "JSON", "TIMESTAMP" return "BLOB", "BLOB", "TIMESTAMP" def _string_type(self, length: int) -> str: """Return Oracle VARCHAR2 type syntax.""" return f"VARCHAR2({length})" def _index_name(self) -> str: """Return index name truncated to Oracle's 30-character limit.""" base_name = f"idx_{self.table_name.replace('.', '_')}_channel_status" return base_name[:30] async def _detect_json_storage_type(self) -> JSONStorageType: """Resolve the JSON storage type from the configured override or server version.""" if self._json_storage_override is not None: return self._json_storage_override return storage_type_from_version(await self._get_version_info()) async def _get_version_info(self) -> "OracleVersionInfo | None": """Return the pool-scoped Oracle version through the data dictionary.""" cache = self._config._oracle_version_cache if cache.resolved: return cache.version async with self._config.provide_session() as driver: version_info = await driver.data_dictionary.get_version(driver) if version_info is None: log_with_context( logger, logging.WARNING, "events.queue.storage.fallback", storage_type=JSONStorageType.BLOB_JSON.value, reason="version_detection_failed", ) return version_info def _table_feature_clause(self) -> str: report = oracle_table_feature_report( self._config, "events", self._extension_settings, "queue", in_memory=self._in_memory, hash_partition_key="event_id", range_partition_key="available_at", table_options_key="table_options", ) return report["clause"] def _init_oracle_settings(extension_settings: "dict[str, Any]") -> "tuple[bool, JSONStorageType | None]": """Initialize Oracle-specific settings from extension config. Args: extension_settings: The events extension settings dict. Returns: Tuple of (in_memory, json_storage) settings. """ in_memory = bool(extension_settings.get("in_memory", False)) json_storage_override = extension_settings.get("json_storage") if json_storage_override == "json": json_storage: JSONStorageType | None = JSONStorageType.JSON_NATIVE elif json_storage_override == "blob_json": json_storage = JSONStorageType.BLOB_JSON elif json_storage_override in {"blob", "blob_plain"}: json_storage = JSONStorageType.BLOB_PLAIN else: json_storage = None return in_memory, json_storage def _oracle_table_ddl(table_name: str, storage_type: "JSONStorageType", table_clause: str, index_name: str) -> str: """Build Oracle CREATE TABLE and INDEX SQL as a single PL/SQL script. Args: table_name: The queue table name. storage_type: JSON storage type (native, blob_json, or blob). table_clause: Capability-gated table storage clauses. index_name: The index name to create. Returns: PL/SQL script for creating table and index. """ if storage_type == JSONStorageType.JSON_NATIVE: payload_col = "payload_json JSON NOT NULL" metadata_col = "metadata_json JSON" elif storage_type == JSONStorageType.BLOB_JSON: payload_col = "payload_json BLOB CHECK (payload_json IS JSON) NOT NULL" metadata_col = "metadata_json BLOB CHECK (metadata_json IS JSON)" else: payload_col = "payload_json BLOB NOT NULL" metadata_col = "metadata_json BLOB" return f""" BEGIN EXECUTE IMMEDIATE 'CREATE TABLE {table_name} ( event_id VARCHAR2(64) PRIMARY KEY, channel VARCHAR2(128) NOT NULL, {payload_col}, {metadata_col}, status VARCHAR2(32) DEFAULT ''pending'' NOT NULL, available_at TIMESTAMP DEFAULT SYSTIMESTAMP NOT NULL, lease_expires_at TIMESTAMP, attempts NUMBER(10) DEFAULT 0 NOT NULL, created_at TIMESTAMP DEFAULT SYSTIMESTAMP NOT NULL, acknowledged_at TIMESTAMP ){table_clause}'; EXCEPTION WHEN OTHERS THEN IF SQLCODE != -955 THEN RAISE; END IF; END; BEGIN EXECUTE IMMEDIATE 'CREATE INDEX {index_name} ON {table_name}(channel, status, available_at)'; EXCEPTION WHEN OTHERS THEN IF SQLCODE != -955 THEN RAISE; END IF; END; """ def _oracle_drop_sql(table_name: str, index_name: str) -> "list[str]": """Build Oracle DROP TABLE SQL with PL/SQL error handling. Args: table_name: The queue table name. index_name: The index name to drop. Returns: List of PL/SQL scripts for dropping index and table. """ return [ f""" BEGIN EXECUTE IMMEDIATE 'DROP INDEX {index_name}'; EXCEPTION WHEN OTHERS THEN IF SQLCODE != -1418 THEN RAISE; END IF; END; """, f""" BEGIN EXECUTE IMMEDIATE 'DROP TABLE {table_name}'; EXCEPTION WHEN OTHERS THEN IF SQLCODE != -942 THEN RAISE; END IF; END; """, ]