Source code for sqlspec.storage.pipeline

"""Storage pipeline scaffolding for driver-aware storage bridge."""

from collections import deque
from functools import partial
from pathlib import Path
from time import perf_counter, time
from typing import TYPE_CHECKING, Any, NamedTuple, TypeAlias, cast

from mypy_extensions import mypyc_attr
from typing_extensions import NotRequired, TypedDict

from sqlspec.exceptions import ImproperConfigurationError, StorageCapabilityError
from sqlspec.storage._arrow_payload import StorageFormat, decode_arrow_payload, encode_arrow_payload
from sqlspec.storage.errors import execute_async_storage_operation, execute_sync_storage_operation
from sqlspec.storage.registry import StorageRegistry, storage_registry
from sqlspec.utils.serializers import get_serializer_metrics, serialize_collection, to_json
from sqlspec.utils.sync_tools import async_
from sqlspec.utils.type_guards import supports_async_delete, supports_async_read_bytes, supports_async_write_bytes
from sqlspec.utils.uuids import uuid4

if TYPE_CHECKING:
    from collections.abc import AsyncIterator, Iterator

    from sqlspec.protocols import ObjectStoreProtocol
    from sqlspec.typing import ArrowTable


__all__ = (
    "AsyncStoragePipeline",
    "PartitionStrategyConfig",
    "StagedArtifact",
    "StorageBridgeJob",
    "StorageCapabilities",
    "StorageDestination",
    "StorageDiagnostics",
    "StorageFormat",
    "StorageLoadRequest",
    "StorageTelemetry",
    "SyncStoragePipeline",
    "create_storage_bridge_job",
    "get_recent_storage_events",
    "get_storage_bridge_diagnostics",
    "get_storage_bridge_metrics",
    "record_storage_diagnostic_event",
    "reset_storage_bridge_events",
    "reset_storage_bridge_metrics",
)

StorageDestination: TypeAlias = str | Path
StorageDiagnostics: TypeAlias = dict[str, float]


class StorageCapabilities(TypedDict):
    """Runtime-evaluated driver storage capabilities."""

    arrow_export_enabled: bool
    arrow_import_enabled: bool
    parquet_export_enabled: bool
    parquet_import_enabled: bool
    requires_staging_for_load: bool
    staging_protocols: "list[str]"
    partition_strategies: "list[str]"
    default_storage_profile: NotRequired[str | None]
class PartitionStrategyConfig(TypedDict, total=False): """Configuration for partition fan-out strategies.""" kind: str partitions: int rows_per_chunk: int manifest_path: str