"""Versioned, resumable Strava summary-history backfill contracts and orchestration."""

from __future__ import annotations

import hashlib
import time
from dataclasses import dataclass
from typing import Callable

from mountain_twin.activity import User
from mountain_twin.activity.persistence import (
    AcquiredProviderPage,
    DurableSyncCoordinator,
    OpaqueProviderState,
    ProviderAuthError,
    ProviderRateLimitError,
    ProviderRetryableError,
    SqliteActivityIdentityStore,
    SyncFrontierPolicy,
    SyncMode,
    SyncRunStatus,
)

from .client import (
    StravaApiError,
    StravaAuthenticationError,
    StravaClient,
    StravaInsufficientScopeError,
    StravaRateLimitError,
)
from .ingestion import PersistentStravaIngestionService
from .oauth import READ_ALL_SCOPE

INITIAL_BACKFILL_PLAN_VERSION = "strava_initial_backfill_plan_v1"
INITIAL_BACKFILL_RESUME_VERSION = "strava_initial_backfill_resume_v1"
INITIAL_BACKFILL_COVERAGE_VERSION = "strava_initial_backfill_coverage_v1"
INCREMENTAL_CONTINUATION_VERSION = "strava_incremental_continuation_v1"
INITIAL_BACKFILL_PER_PAGE = 100
HISTORICAL_STRATEGY = "frozen_before_page_until_empty_v1"
OVERLAP_STRATEGY = "strava_anchor_tick_replay_v1"
_HISTORICAL = "HISTORICAL"
_RECONCILIATION = "RECONCILIATION"
_RESUMABLE_INITIAL_BACKFILL_STATUSES = {
    SyncRunStatus.PLANNED,
    SyncRunStatus.RUNNING,
    SyncRunStatus.PAUSED_RATE_LIMIT,
    SyncRunStatus.FAILED_RETRYABLE,
    SyncRunStatus.FAILED_ACTIVITY,
    SyncRunStatus.BLOCKED_AUTH,
}


def new_initial_backfill_plan(
    *,
    provider_subject_id: str,
    historical_before: int,
    frontier_policy: SyncFrontierPolicy = SyncFrontierPolicy.PUBLISH_INCREMENTAL_CONTINUATION,
) -> OpaqueProviderState:
    """Construct an uncapped private-summary initial-history plan with explicit authority intent."""
    if not provider_subject_id or historical_before <= 0:
        raise ValueError(
            "initial backfill requires provider subject and positive historical boundary"
        )
    return OpaqueProviderState(
        INITIAL_BACKFILL_PLAN_VERSION,
        {
            "activity_cap": None,
            "frontier_policy": frontier_policy.value,
            "historical_before": historical_before,
            "historical_strategy": HISTORICAL_STRATEGY,
            "per_page": INITIAL_BACKFILL_PER_PAGE,
            "provider": "strava",
            "provider_subject_id": provider_subject_id,
            "query_scope": "private_inclusive_summary_history",
            "required_scope": READ_ALL_SCOPE,
        },
    )


def new_initial_backfill_resume(plan: OpaqueProviderState) -> OpaqueProviderState:
    values = validate_initial_backfill_plan(plan)
    return OpaqueProviderState(
        INITIAL_BACKFILL_RESUME_VERSION,
        {"historical_before": values["historical_before"], "page": 1, "phase": _HISTORICAL},
    )


def validate_initial_backfill_plan(
    plan: OpaqueProviderState,
    *,
    expected_subject_id: str | None = None,
    expected_frontier_policy: SyncFrontierPolicy | None = None,
) -> dict:
    if plan.version != INITIAL_BACKFILL_PLAN_VERSION:
        raise ValueError("unknown initial backfill plan version")
    values = plan.value
    required = {
        "activity_cap": None,
        "historical_strategy": HISTORICAL_STRATEGY,
        "per_page": INITIAL_BACKFILL_PER_PAGE,
        "provider": "strava",
        "query_scope": "private_inclusive_summary_history",
        "required_scope": READ_ALL_SCOPE,
    }
    if any(values.get(key) != value for key, value in required.items()):
        raise ValueError("initial backfill plan is not authoritative and uncapped")
    if not isinstance(values.get("historical_before"), int) or values["historical_before"] <= 0:
        raise ValueError("initial backfill plan requires positive frozen historical boundary")
    if not isinstance(values.get("provider_subject_id"), str) or not values["provider_subject_id"]:
        raise ValueError("initial backfill plan requires provider subject")
    try:
        plan_policy = SyncFrontierPolicy(values.get("frontier_policy"))
    except ValueError as error:
        raise ValueError("initial backfill plan requires explicit frontier policy") from error
    if expected_frontier_policy is not None and plan_policy is not expected_frontier_policy:
        raise ValueError("initial backfill plan frontier policy differs from durable run")
    if expected_subject_id is not None and values["provider_subject_id"] != expected_subject_id:
        raise ValueError("initial backfill provider subject differs from sync source")
    return values


def validate_initial_backfill_resume(
    plan: OpaqueProviderState, resume: OpaqueProviderState
) -> dict:
    plan_values = validate_initial_backfill_plan(plan)
    if resume.version != INITIAL_BACKFILL_RESUME_VERSION:
        raise ValueError("unknown initial backfill resume version")
    values = resume.value
    if values.get("phase") not in {_HISTORICAL, _RECONCILIATION} or not isinstance(
        values.get("page"), int
    ):
        raise ValueError("initial backfill resume has invalid phase or page")
    if values["page"] <= 0:
        raise ValueError("initial backfill resume page must be positive")
    if values.get("historical_before") != plan_values["historical_before"]:
        raise ValueError("initial backfill resume differs from immutable historical boundary")
    if values["phase"] == _RECONCILIATION:
        for key in ("historical_empty_page", "reconciliation_after", "reconciliation_before"):
            if not isinstance(values.get(key), int):
                raise ValueError("reconciliation resume is incomplete")
        if values["reconciliation_after"] >= values["reconciliation_before"]:
            raise ValueError("reconciliation bounds are invalid")
    return values


class StravaInitialBackfillAdapter:
    """Adapter-owned two-phase enumeration; no snapshot guarantee is implied."""

    def __init__(
        self,
        client: StravaClient,
        *,
        provider_subject_id: str,
        clock: Callable[[], int] | None = None,
    ):
        self.client = client
        self.provider_subject_id = provider_subject_id
        self.clock = clock or (lambda: int(time.time()))

    def acquire_page(
        self, plan: OpaqueProviderState, resume: OpaqueProviderState
    ) -> AcquiredProviderPage:
        try:
            plan_values = validate_initial_backfill_plan(
                plan, expected_subject_id=self.provider_subject_id
            )
            resume_values = validate_initial_backfill_resume(plan, resume)
            self.client.ensure_access(READ_ALL_SCOPE)
            if str(self.client.token.athlete_id) != self.provider_subject_id:
                raise ProviderAuthError("STRAVA_PROVIDER_SUBJECT_MISMATCH")
            phase, page = resume_values["phase"], resume_values["page"]
            if phase == _HISTORICAL:
                values = self.client.activity_page(
                    page=page,
                    per_page=plan_values["per_page"],
                    before=plan_values["historical_before"],
                )
                return self._historical_page(plan, resume, resume_values, values)
            values = self.client.activity_page(
                page=page,
                per_page=plan_values["per_page"],
                after=resume_values["reconciliation_after"],
                before=resume_values["reconciliation_before"],
            )
            return self._reconciliation_page(plan, resume, resume_values, values)
        except ProviderAuthError:
            raise
        except StravaRateLimitError as error:
            raise ProviderRateLimitError("STRAVA_RATE_LIMITED") from error
        except (StravaAuthenticationError, StravaInsufficientScopeError) as error:
            raise ProviderAuthError("STRAVA_AUTH_BLOCKED") from error
        except (StravaApiError, ValueError) as error:
            raise ProviderRetryableError("STRAVA_INITIAL_BACKFILL_PAGE_UNAVAILABLE") from error

    def _historical_page(
        self,
        plan: OpaqueProviderState,
        resume: OpaqueProviderState,
        values: dict,
        payloads: tuple[dict, ...],
    ) -> AcquiredProviderPage:
        if payloads:
            return AcquiredProviderPage(
                payloads,
                OpaqueProviderState(resume.version, {**values, "page": values["page"] + 1}),
                enumeration_exhausted=False,
            )
        reconciliation_before = max(int(self.clock()), values["historical_before"] + 1)
        reconciliation_after = values["historical_before"] - 1
        reconciliation = OpaqueProviderState(
            resume.version,
            {
                "historical_before": values["historical_before"],
                "historical_empty_page": values["page"],
                "page": 1,
                "phase": _RECONCILIATION,
                "reconciliation_after": reconciliation_after,
                "reconciliation_before": reconciliation_before,
            },
        )
        return AcquiredProviderPage(
            (), reconciliation, enumeration_exhausted=True, plan_complete=False
        )

    def _reconciliation_page(
        self,
        plan: OpaqueProviderState,
        resume: OpaqueProviderState,
        values: dict,
        payloads: tuple[dict, ...],
    ) -> AcquiredProviderPage:
        if payloads:
            return AcquiredProviderPage(
                payloads,
                OpaqueProviderState(resume.version, {**values, "page": values["page"] + 1}),
                enumeration_exhausted=False,
            )
        coverage = _coverage_state(plan, values)
        plan_values = validate_initial_backfill_plan(plan)
        candidate = None
        if (
            plan_values["frontier_policy"]
            == SyncFrontierPolicy.PUBLISH_INCREMENTAL_CONTINUATION.value
        ):
            candidate = OpaqueProviderState(
                INCREMENTAL_CONTINUATION_VERSION,
                {
                    "coverage_evidence_digest": _state_digest(coverage),
                    "established_before": values["reconciliation_before"],
                    "incremental_after": values["reconciliation_after"],
                    "overlap_strategy": OVERLAP_STRATEGY,
                },
            )
        return AcquiredProviderPage(
            (),
            resume,
            enumeration_exhausted=True,
            terminal_provider_state=coverage,
            plan_complete=True,
            frontier_candidate=candidate,
        )


def _coverage_state(plan: OpaqueProviderState, resume: dict) -> OpaqueProviderState:
    plan_values = validate_initial_backfill_plan(plan)
    if resume.get("phase") != _RECONCILIATION:
        raise ValueError("coverage requires completed reconciliation")
    return OpaqueProviderState(
        INITIAL_BACKFILL_COVERAGE_VERSION,
        {
            "historical_before": plan_values["historical_before"],
            "historical_empty_page": resume["historical_empty_page"],
            "per_page": plan_values["per_page"],
            "plan_digest": _state_digest(plan),
            "reconciliation_after": resume["reconciliation_after"],
            "reconciliation_before": resume["reconciliation_before"],
            "reconciliation_empty_page": resume["page"],
        },
    )


def _state_digest(state: OpaqueProviderState) -> str:
    return hashlib.sha256(f"{state.version}:{state.serialized}".encode()).hexdigest()


@dataclass(frozen=True)
class InitialBackfillResult:
    sync_source_id: str
    sync_run_id: str
    status: SyncRunStatus
    pages_completed: int
    activities_persisted: int
    coverage_evidence_present: bool
    authoritative_frontier_present: bool
    invocation_page_budget_reached: bool = False


def create_initial_backfill_run(
    store: SqliteActivityIdentityStore,
    *,
    user: User,
    provider_subject_id: str,
    historical_before: int,
    started_at: str,
    frontier_policy: SyncFrontierPolicy = SyncFrontierPolicy.PUBLISH_INCREMENTAL_CONTINUATION,
):
    source = store.resolve_provider_sync_source(user, "strava", provider_subject_id, started_at)
    plan = new_initial_backfill_plan(
        provider_subject_id=provider_subject_id,
        historical_before=historical_before,
        frontier_policy=frontier_policy,
    )
    return store.create_sync_run(
        source.sync_source_id,
        SyncMode.INITIAL_BACKFILL,
        plan,
        new_initial_backfill_resume(plan),
        started_at,
        frontier_policy=frontier_policy,
    )


def run_initial_backfill(
    store: SqliteActivityIdentityStore,
    *,
    sync_run_id: str,
    user: User,
    adapter: StravaInitialBackfillAdapter,
    updated_at: Callable[[], str],
    page_budget: int | None = None,
) -> InitialBackfillResult:
    """Resume one durable initial run until a safe terminal state; acquisition remains adapter-owned."""
    run = store.get_sync_run(sync_run_id)
    if page_budget is not None and page_budget <= 0:
        raise ValueError("invocation page budget must be positive")
    if page_budget is not None and run.frontier_policy is not SyncFrontierPolicy.NONE:
        raise ValueError("page budget is allowed only for non-authoritative initial backfill")
    if not initial_backfill_resume_available(store, run):
        raise ValueError("initial backfill run is not resumable")
    if run.status is not SyncRunStatus.RUNNING:
        run = store.transition_sync_run(run.sync_run_id, SyncRunStatus.RUNNING, updated_at())
    ingestion = PersistentStravaIngestionService(store, user)
    coordinator = DurableSyncCoordinator(
        store,
        adapter,
        lambda payload: ingestion.ingest_activity(payload, fetched_at=updated_at()),
    )
    durable_pages_this_invocation = 0
    while run.status is SyncRunStatus.RUNNING:
        if page_budget is not None and durable_pages_this_invocation >= page_budget:
            break
        pages_before = run.pages_completed
        status = coordinator.execute_next_page(run.sync_run_id, updated_at())
        run = store.get_sync_run(run.sync_run_id)
        durable_pages_this_invocation += run.pages_completed - pages_before
        if status is not SyncRunStatus.RUNNING:
            break
    try:
        store.get_authoritative_frontier(run.sync_source_id)
        frontier_present = True
    except KeyError:
        frontier_present = False
    return InitialBackfillResult(
        run.sync_source_id,
        run.sync_run_id,
        run.status,
        run.pages_completed,
        run.activities_persisted,
        run.terminal_provider_state is not None,
        frontier_present,
        page_budget is not None
        and durable_pages_this_invocation >= page_budget
        and run.status is SyncRunStatus.RUNNING,
    )


def initial_backfill_resume_available(store: SqliteActivityIdentityStore, run) -> bool:
    """True only when the same run can enter the existing initial-backfill resume path."""
    if run.status not in _RESUMABLE_INITIAL_BACKFILL_STATUSES:
        return False
    if run.mode is not SyncMode.INITIAL_BACKFILL:
        return False
    try:
        source = store.get_provider_sync_source(run.sync_source_id)
        validate_initial_backfill_plan(
            run.plan,
            expected_subject_id=source.external_subject_id,
            expected_frontier_policy=run.frontier_policy,
        )
        validate_initial_backfill_resume(run.plan, run.resume)
    except (KeyError, ValueError):
        return False
    return True


def diagnostics_for_initial_backfill(
    store: SqliteActivityIdentityStore, result: InitialBackfillResult
) -> dict[str, object]:
    """Return only safe lifecycle facts; opaque state and provider payloads remain private."""
    run = store.get_sync_run(result.sync_run_id)
    return {
        "sync_source_id": result.sync_source_id,
        "sync_run_id": result.sync_run_id,
        "status": result.status.value,
        "phase": run.resume.value.get("phase"),
        "pages_completed": result.pages_completed,
        "activities_persisted": result.activities_persisted,
        "coverage_evidence_present": result.coverage_evidence_present,
        "authoritative_frontier_present": result.authoritative_frontier_present,
        "resume_available": initial_backfill_resume_available(store, run),
        "invocation_page_budget_reached": getattr(result, "invocation_page_budget_reached", False),
    }
