"""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