Source code for sqlspec.extensions.events._channel

"""Event channel API with separate sync and async implementations."""

import asyncio
import importlib
import inspect
import logging
import threading
from collections.abc import AsyncIterator, Iterator, Sequence
from contextlib import contextmanager, suppress
from dataclasses import dataclass
from time import perf_counter
from typing import TYPE_CHECKING, Any, cast

from sqlspec.exceptions import ImproperConfigurationError, MissingDependencyError
from sqlspec.extensions.events._buffer import validate_queue_capacity
from sqlspec.extensions.events._hints import get_runtime_hints, resolve_adapter_name
from sqlspec.extensions.events._models import EventMessage
from sqlspec.extensions.events._names import normalize_event_channel_name
from sqlspec.extensions.events._protocols import AsyncEventBackendProtocol, SyncEventBackendProtocol
from sqlspec.extensions.events._queue import build_queue_backend
from sqlspec.utils.logging import get_logger, log_with_context
from sqlspec.utils.type_guards import has_span_attribute
from sqlspec.utils.uuids import uuid4

if TYPE_CHECKING:
    from sqlspec.config import AsyncDatabaseConfig, SyncDatabaseConfig
    from sqlspec.extensions.events._hints import EventRuntimeHints
    from sqlspec.extensions.events._protocols import AsyncEventHandler, SyncEventHandler
    from sqlspec.observability import ObservabilityRuntime

__all__ = (
    "AsyncEventChannel",
    "AsyncEventListener",
    "EventMessage",
    "SyncEventChannel",
    "SyncEventListener",
    "load_native_backend",
    "resolve_event_poll_interval",
    "resolve_poll_interval",
)

logger = get_logger("sqlspec.events.channel")
_LISTENER_SHUTDOWN_TIMEOUT = 0.5


@dataclass(slots=True)
class AsyncEventListener:
    """Represents a running async listener task."""

    id: str
    channel: str
    task: "asyncio.Task[Any]"
    stop_event: "asyncio.Event"
    poll_interval: float

    async def stop(self) -> None:
        """Signal the listener to stop and await task completion."""
        self.stop_event.set()
        if not self.task.done():
            self.task.cancel()
        with suppress(asyncio.CancelledError):
            await self.task