"""Provider-neutral durable page orchestration; acquisition stays outside SQLite transactions."""

from __future__ import annotations

from dataclasses import dataclass
from typing import Callable, Protocol

from .ingestion import PersistentIngestionResult
from .repository import SqliteActivityIdentityStore
from .sync import OpaqueProviderState, SyncFrontierPolicy, SyncRunStatus


class ProviderPageError(RuntimeError):
    code = "PROVIDER_FAILURE"


class ProviderRateLimitError(ProviderPageError):
    code = "PROVIDER_RATE_LIMITED"


class ProviderAuthError(ProviderPageError):
    code = "PROVIDER_AUTH_BLOCKED"


class ProviderRetryableError(ProviderPageError):
    code = "PROVIDER_RETRYABLE_FAILURE"


@dataclass(frozen=True)
class AcquiredProviderPage:
    payloads: tuple[dict, ...]
    next_resume: OpaqueProviderState
    enumeration_exhausted: bool
    terminal_provider_state: OpaqueProviderState | None = None
    plan_complete: bool | None = None
    frontier_candidate: OpaqueProviderState | None = None

    def __post_init__(self) -> None:
        # Retain the WP5D.4b positional form while making plan completion explicit.
        if self.plan_complete is None:
            object.__setattr__(
                self,
                "plan_complete",
                self.enumeration_exhausted and self.terminal_provider_state is not None,
            )
        if (
            self.terminal_provider_state is not None or self.frontier_candidate is not None
        ) and not self.plan_complete:
            raise ValueError(
                "terminal provider state and frontier candidate require completed provider plan"
            )
        if self.frontier_candidate is not None and self.terminal_provider_state is None:
            raise ValueError("frontier candidate requires terminal provider evidence")

    @property
    def exhausted(self) -> bool:
        """Compatibility alias; this is query exhaustion, not necessarily plan completion."""
        return self.enumeration_exhausted


class ProviderPageAdapter(Protocol):
    def acquire_page(
        self, plan: OpaqueProviderState, resume: OpaqueProviderState
    ) -> AcquiredProviderPage: ...


class DurableSyncCoordinator:
    """Execute one page at a time without interpreting adapter-owned opaque state."""

    def __init__(
        self,
        store: SqliteActivityIdentityStore,
        adapter: ProviderPageAdapter,
        ingest_payload: Callable[[dict], PersistentIngestionResult],
    ):
        self.store, self.adapter, self.ingest_payload = store, adapter, ingest_payload

    def execute_next_page(
        self,
        sync_run_id: str,
        updated_at: str,
        *,
        cancelled: Callable[[], bool] | None = None,
        before_progress: Callable[[], None] | None = None,
        complete_requested_plan: bool = False,
    ) -> SyncRunStatus:
        """Process one page; callers may explicitly complete a bounded plan without provider exhaustion."""
        run = self.store.get_sync_run(sync_run_id)
        if run.status is not SyncRunStatus.RUNNING:
            raise ValueError("durable coordinator requires a running sync run")
        if cancelled and cancelled():
            return self.store.transition_sync_run(
                sync_run_id, SyncRunStatus.CANCELLED, updated_at
            ).status
        try:
            page = self.adapter.acquire_page(run.plan, run.resume)
        except ProviderRateLimitError as error:
            return self.store.transition_sync_run(
                sync_run_id, SyncRunStatus.PAUSED_RATE_LIMIT, updated_at, error.code
            ).status
        except ProviderAuthError as error:
            return self.store.transition_sync_run(
                sync_run_id, SyncRunStatus.BLOCKED_AUTH, updated_at, error.code
            ).status
        except ProviderPageError as error:
            return self.store.transition_sync_run(
                sync_run_id, SyncRunStatus.FAILED_RETRYABLE, updated_at, error.code
            ).status
        try:
            for payload in page.payloads:
                result = self.ingest_payload(payload)
                self.store.record_sync_run_observation(
                    sync_run_id, result.source_link_id, result.raw_record_id, updated_at
                )
        except Exception:
            return self.store.transition_sync_run(
                sync_run_id,
                SyncRunStatus.FAILED_ACTIVITY,
                updated_at,
                "ACTIVITY_PERSISTENCE_FAILED",
            ).status
        if cancelled and cancelled():
            return self.store.transition_sync_run(
                sync_run_id, SyncRunStatus.CANCELLED, updated_at
            ).status
        if before_progress:
            before_progress()
        if complete_requested_plan and run.frontier_policy is not SyncFrontierPolicy.NONE:
            raise ValueError("explicit bounded completion cannot bypass an eligible provider plan")
        if page.plan_complete or complete_requested_plan:
            if (
                run.frontier_policy is SyncFrontierPolicy.PUBLISH_INCREMENTAL_CONTINUATION
                and page.frontier_candidate is None
            ):
                return self.store.transition_sync_run(
                    sync_run_id,
                    SyncRunStatus.FAILED_ACTIVITY,
                    updated_at,
                    "AUTHORITATIVE_COMPLETION_EVIDENCE_MISSING",
                ).status
            if (
                page.frontier_candidate is not None
                and run.frontier_policy is SyncFrontierPolicy.NONE
            ):
                raise ValueError("ineligible sync run cannot complete with a frontier candidate")
            self.store.complete_sync_run(
                sync_run_id,
                updated_at,
                pages_completed=run.pages_completed + int(bool(page.payloads)),
                activities_persisted=run.activities_persisted + len(page.payloads),
                terminal_provider_state=page.terminal_provider_state,
                frontier_candidate=page.frontier_candidate
                if run.frontier_policy is SyncFrontierPolicy.PUBLISH_INCREMENTAL_CONTINUATION
                else None,
            )
            return SyncRunStatus.COMPLETED
        return self.store.advance_sync_run_resume(
            sync_run_id,
            page.next_resume,
            updated_at,
            pages_completed=run.pages_completed + int(bool(page.payloads)),
            activities_persisted=run.activities_persisted + len(page.payloads),
        ).status
