"""CheckThisFile integrity v1 / Python 3.11+ / httpx.

Server only. No automatic retry, upload, purchase, review or approval.
Reuse one AsyncClient in a FastAPI lifespan. Keep a durable outbox in YOUR DB.
Official HTTPX guidance: https://www.python-httpx.org/async/
"""
from __future__ import annotations
import hashlib
import json
from dataclasses import dataclass
from typing import Any
from urllib.parse import quote, urlparse
import httpx


@dataclass
class CTFRejected(Exception):
    status: int
    code: str
    retry_after: int | None = None


@dataclass
class CTFUnavailable(Exception):
    phase: str
    ambiguous: bool  # A write may have committed: use lookup, do not invent a new version.
    retry_after: int | None = None


def final_bytes_declaration(reference: dict[str, Any], pdf: bytes, generated_at: str,
                            previous_version: int | None = None, reservation_id: str | None = None) -> dict[str, Any]:
    result = {**reference, "sha256": hashlib.sha256(pdf).hexdigest(), "sizeBytes": len(pdf),
              "mimeType": "application/pdf", "generatedAt": generated_at, "previousVersion": previous_version}
    if reservation_id is not None:
        result["reservationId"] = reservation_id
    return result


def retry_seconds(response: httpx.Response) -> int | None:
    value = response.headers.get("Retry-After", "")
    return int(value) if value.isdigit() else None


class IntegrityClient:
    def __init__(self, origin: str, api_key: str, *, transport: httpx.AsyncBaseTransport | None = None):
        parsed = urlparse(origin)
        if (parsed.scheme != "https" and not (parsed.scheme == "http" and parsed.hostname in ("localhost", "127.0.0.1"))) or parsed.username or parsed.password or parsed.query or parsed.fragment or parsed.path not in ("", "/"):
            raise ValueError("Use an HTTPS origin, without /api/v1, credentials or query")
        self._client = httpx.AsyncClient(base_url=origin.rstrip("/"),
            headers={"Authorization": f"Bearer {api_key}"}, follow_redirects=False,
            timeout=httpx.Timeout(10.0, connect=5.0), transport=transport)

    async def close(self) -> None:
        await self._client.aclose()

    async def _request(self, method: str, route: str, *, body: dict | None = None,
                       params: dict | None = None, raw: bool = False) -> dict:
        try:
            # Explicit JSON bytes keep Content-Length deterministic; never send PDF bytes.
            kwargs: dict = {"params": params}
            if body is not None:
                kwargs.update(content=json.dumps(body, allow_nan=False).encode(), headers={"Content-Type": "application/json"})
            response = await self._client.request(method, "/api/v1/integrity/" + route, **kwargs)
        except (httpx.ConnectError, httpx.ConnectTimeout, httpx.PoolTimeout) as exc:
            raise CTFUnavailable("connect_or_pool", False) from exc
        except httpx.RequestError as exc:
            raise CTFUnavailable("transport", method == "POST") from exc
        if response.status_code >= 500:
            raise CTFUnavailable("service", method == "POST", retry_seconds(response))
        try:
            result = response.json()
        except ValueError as exc:
            raise CTFUnavailable("invalid_response", method == "POST") from exc
        if not isinstance(result, dict):
            raise CTFUnavailable("invalid_response", method == "POST")
        if not 200 <= response.status_code < 300:
            error = result.get("error")
            code = error.get("code", "HTTP_ERROR") if isinstance(error, dict) else "HTTP_ERROR"
            raise CTFRejected(response.status_code, str(code), retry_seconds(response))
        if raw:
            return result
        if not isinstance(result.get("data"), dict) or result.get("error") is not None:
            raise CTFUnavailable("invalid_response", method == "POST")
        return result["data"]

    async def reserve(self, reference: dict) -> dict:
        return await self._request("POST", "reservations", body=reference)

    async def create_company(self, reference: str) -> dict:
        """Admin scope and owner/admin role. Does not grant the worker access."""
        return await self._request("POST", "companies", body={"reference": reference})

    async def grant_key(self, company_id: str, key_id: str, enabled: bool = True) -> dict:
        return await self._request("POST", "companies/" + quote(company_id, safe="") + "/access",
                                   body={"keyId": key_id, "enabled": enabled})

    async def grant_webhook(self, company_id: str, webhook_id: str, enabled: bool = True) -> dict:
        return await self._request("POST", "companies/" + quote(company_id, safe="") + "/webhooks",
                                   body={"webhookId": webhook_id, "enabled": enabled})

    async def register(self, declaration: dict) -> dict:
        return await self._request("POST", "records", body=declaration)

    async def lookup(self, reference: dict) -> dict:
        return await self._request("GET", "lookup", params=reference)

    async def evidence(self, public_id: str) -> dict:
        return await self._request("GET", "records/" + quote(public_id, safe="") + "/evidence", raw=True)

    async def changes(self, company_reference: str, after: int = 0, limit: int = 100) -> dict:
        return await self._request("GET", "changes", params={"companyReference": company_reference, "after": after, "limit": limit})

    async def usage(self, company_reference: str) -> dict:
        return await self._request("GET", "usage", params={"companyReference": company_reference})

    async def revoke(self, public_id: str, reason: str) -> dict:
        return await self._request("POST", "records/" + quote(public_id, safe="") + "/revoke", body={"reason": reason})


async def reconcile_pending(client: IntegrityClient, declaration: dict) -> dict:
    """Call from a durable worker, not inside DECA's Emitir HTTP request.

    Lookup first after any ambiguous write. A 404 then permits replay of the EXACT
    saved declaration. A 409 is an incident, not a reason to overwrite/change keys.
    The caller schedules backoff/jitter honoring retry_after and commits results.
    """
    reference = {key: declaration[key] for key in ("companyReference", "source", "documentReference", "version")}
    try:
        existing = await client.lookup(reference)
    except CTFRejected as exc:
        if exc.status != 404:
            raise
        return await client.register(declaration)
    # Always replay the same declarations to validate all immutable fields, not just hash.
    if existing["sha256"] != declaration["sha256"].lower():
        raise CTFRejected(409, "VERSION_CONFLICT")
    return await client.register(declaration)
