"""Provider-neutral durable synchronization state; never provider acquisition or credentials."""

from __future__ import annotations

import hashlib
import json
from dataclasses import dataclass
from datetime import datetime
from enum import Enum
from uuid import uuid4


class SyncMode(str, Enum):
    INITIAL_BACKFILL = "INITIAL_BACKFILL"
    INCREMENTAL = "INCREMENTAL"
    BOUNDED = "BOUNDED"
    RECONCILIATION = "RECONCILIATION"


class SyncFrontierPolicy(str, Enum):
    """Eligibility fixed before acquisition; never inferred from provider state."""

    NONE = "NONE"
    PUBLISH_INCREMENTAL_CONTINUATION = "PUBLISH_INCREMENTAL_CONTINUATION"


class SyncFrontierScope(str, Enum):
    """Meaning of a persisted authoritative source frontier."""

    LEGACY_AMBIGUOUS = "LEGACY_AMBIGUOUS"
    INCREMENTAL_CONTINUATION = "INCREMENTAL_CONTINUATION"


class SyncRunStatus(str, Enum):
    PLANNED = "PLANNED"
    RUNNING = "RUNNING"
    COMPLETED = "COMPLETED"
    PAUSED_RATE_LIMIT = "PAUSED_RATE_LIMIT"
    FAILED_RETRYABLE = "FAILED_RETRYABLE"
    FAILED_ACTIVITY = "FAILED_ACTIVITY"
    BLOCKED_AUTH = "BLOCKED_AUTH"
    CANCELLED = "CANCELLED"


_SENSITIVE_STATE_KEY_PARTS = ("token", "secret", "authorization", "payload", "coordinate", "location", "name")


def _aware(value: str, label: str) -> None:
    instant = datetime.fromisoformat(value)
    if instant.tzinfo is None or instant.utcoffset() is None:
        raise ValueError(f"{label} must be timezone-aware")


@dataclass(frozen=True)
class OpaqueProviderState:
    """Bounded deterministic adapter state; generic persistence does not interpret its keys."""

    version: str
    value: dict

    def __post_init__(self) -> None:
        if not self.version or not isinstance(self.value, dict):
            raise ValueError("opaque provider state requires a versioned object")
        _validate_safe_state(self.value)
        if len(self.serialized) > 16_384:
            raise ValueError("opaque provider state exceeds bounded storage")

    @property
    def serialized(self) -> str:
        return json.dumps(self.value, sort_keys=True, separators=(",", ":"), allow_nan=False)

    @classmethod
    def from_serialized(cls, version: str, value: str) -> "OpaqueProviderState":
        decoded = json.loads(value)
        return cls(version, decoded)


@dataclass(frozen=True)
class ProviderSyncSource:
    sync_source_id: str
    user_id: str
    provider: str
    external_subject_id: str
    created_at: str

    def __post_init__(self) -> None:
        if not all((self.sync_source_id, self.user_id, self.provider, self.external_subject_id)):
            raise ValueError("provider sync source requires stable owner and provider account identity")
        _aware(self.created_at, "provider sync source creation")
        if self.sync_source_id != provider_sync_source_id_for(
            self.user_id, self.provider, self.external_subject_id
        ):
            raise ValueError("provider sync source identity must match its logical source")


@dataclass(frozen=True)
class ProviderSyncRun:
    sync_run_id: str
    sync_source_id: str
    user_id: str
    mode: SyncMode
    status: SyncRunStatus
    plan: OpaqueProviderState
    resume: OpaqueProviderState
    started_at: str
    updated_at: str
    completed_at: str | None = None
    pages_completed: int = 0
    activities_persisted: int = 0
    failure_code: str | None = None
    frontier_policy: SyncFrontierPolicy = SyncFrontierPolicy.NONE
    terminal_provider_state: OpaqueProviderState | None = None

    def __post_init__(self) -> None:
        if not all((self.sync_run_id, self.sync_source_id, self.user_id)):
            raise ValueError("provider sync run requires identity and owner")
        _aware(self.started_at, "sync run start")
        _aware(self.updated_at, "sync run update")
        if self.completed_at is not None:
            _aware(self.completed_at, "sync run completion")
        if self.pages_completed < 0 or self.activities_persisted < 0:
            raise ValueError("sync progress cannot be negative")
        if self.status is SyncRunStatus.COMPLETED and self.completed_at is None:
            raise ValueError("completed sync run requires completion time")
        if self.status is not SyncRunStatus.COMPLETED and self.completed_at is not None:
            raise ValueError("only completed sync run may have completion time")
        if self.terminal_provider_state is not None and self.status is not SyncRunStatus.COMPLETED:
            raise ValueError("only completed sync run may retain terminal provider state")
        if self.mode in {SyncMode.BOUNDED, SyncMode.RECONCILIATION} and self.frontier_policy is not SyncFrontierPolicy.NONE:
            raise ValueError("bounded and reconciliation sync runs cannot publish an authoritative frontier")


@dataclass(frozen=True)
class ProviderSyncCheckpoint:
    sync_source_id: str
    state: OpaqueProviderState
    completed_sync_run_id: str
    published_at: str
    frontier_scope: SyncFrontierScope

    def __post_init__(self) -> None:
        if not all((self.sync_source_id, self.completed_sync_run_id)):
            raise ValueError("sync checkpoint requires source and completed run")
        _aware(self.published_at, "sync checkpoint publication")


@dataclass(frozen=True)
class SyncRunSourceObservation:
    sync_run_id: str
    sync_source_id: str
    user_id: str
    source_link_id: str
    raw_record_id: str
    observed_at: str

    def __post_init__(self) -> None:
        if not all((self.sync_run_id, self.sync_source_id, self.user_id, self.source_link_id, self.raw_record_id)):
            raise ValueError("sync observation requires run, source and raw identities")
        _aware(self.observed_at, "sync observation time")


def provider_sync_source_id_for(user_id: str, provider: str, external_subject_id: str) -> str:
    identity = json.dumps(
        {"format": "provider_sync_source_identity_v0_1", "user_id": user_id,
         "provider": provider, "external_subject_id": external_subject_id},
        sort_keys=True, separators=(",", ":"),
    )
    return f"provider-sync-source-v0_1:{hashlib.sha256(identity.encode()).hexdigest()}"


def new_provider_sync_run_id() -> str:
    return f"provider-sync-run-v0_1:{uuid4()}"


def _validate_safe_state(value: object) -> None:
    if isinstance(value, dict):
        for key, nested in value.items():
            if not isinstance(key, str) or any(part in key.lower() for part in _SENSITIVE_STATE_KEY_PARTS):
                raise ValueError("opaque provider state contains prohibited diagnostic key")
            _validate_safe_state(nested)
    elif isinstance(value, list):
        for nested in value:
            _validate_safe_state(nested)
    elif value is not None and not isinstance(value, (str, int, float, bool)):
        raise ValueError("opaque provider state must be JSON-compatible")
