"""Python 3.11+ standard-library reference connector. Keep credentials out of tool parameters."""
import base64
import json
import os
from pathlib import Path
import secrets
import hashlib
import urllib.error
import urllib.parse
import urllib.request
import uuid


class NoRedirect(urllib.request.HTTPRedirectHandler):
    def redirect_request(self, *args, **kwargs):
        return None


class DoctrineClient:
    def __init__(self, provider=None, model=None, model_key=None, base="https://heng.lu", vault=None):
        url = urllib.parse.urlsplit(base)
        if (url.scheme != "https" and not (url.scheme == "http" and url.hostname in ("localhost", "127.0.0.1", "::1"))) or url.username or url.password or url.query or url.fragment or url.path not in ("", "/"):
            raise ValueError("invalid_service_origin")
        self.base = base.rstrip("/") + "/api/doctrine/v1"
        self.provider, self.model, self.model_key = provider, model, model_key
        self.vault = Path(vault or Path.home() / ".config/henglu-doctrine/python-credentials.json")
        self.vault.parent.mkdir(parents=True, exist_ok=True, mode=0o700)
        self.credentials = self.load()
        self.http = urllib.request.build_opener(NoRedirect)

    @staticmethod
    def private_json(path):
        if path.stat().st_mode & 0o077:
            raise ValueError("credential_vault_must_be_private")
        return json.loads(path.read_text())

    def load(self):
        entries = self.private_json(self.vault) if self.vault.exists() else {}
        directory = Path(str(self.vault) + ".d")
        if not directory.exists():
            return entries
        if directory.stat().st_mode & 0o077:
            raise ValueError("credential_vault_must_be_private")
        for path in directory.glob("*.entry.json"):
            value = self.private_json(path)
            binding = path.with_name(path.name.replace(".entry.json", ".binding.json"))
            entries[value["slot"]] = dict(value["entry"])
            if binding.exists():
                entries[value["slot"]]["analysis_id"] = self.private_json(binding)["analysis_id"]
        return entries

    def slot_path(self, slot, suffix):
        return Path(str(self.vault) + ".d") / (hashlib.sha256(slot.encode()).hexdigest() + "." + suffix + ".json")

    def publish(self, path, value):
        # Immutable files avoid a process-owned lock: a crash cannot lock out
        # another connector, and the atomic hard link selects one shared winner.
        path.parent.mkdir(parents=True, exist_ok=True, mode=0o700)
        if path.parent.stat().st_mode & 0o077:
            raise ValueError("credential_vault_must_be_private")
        pending = path.parent / (uuid.uuid4().hex + ".pending")
        try:
            fd = os.open(pending, os.O_CREAT | os.O_EXCL | os.O_WRONLY, 0o600)
            with os.fdopen(fd, "w") as stream:
                json.dump(value, stream)
                stream.flush()
                os.fsync(stream.fileno())
            try:
                os.link(pending, path)
            except FileExistsError:
                pass
            for directory in (path.parent, path.parent.parent):
                fd = os.open(directory, os.O_RDONLY)
                try:
                    os.fsync(fd)
                finally:
                    os.close(fd)
            return self.private_json(path)
        finally:
            pending.unlink(missing_ok=True)

    def entry(self, slot):
        entry = self.load().get(slot) or {"token": base64.urlsafe_b64encode(secrets.token_bytes(32)).decode().rstrip("="), "origin": self.base}
        return self.publish(self.slot_path(slot, "entry"), {"slot": slot, "entry": entry})["entry"]

    def bind(self, slot, analysis_id):
        value = self.publish(self.slot_path(slot, "binding"), {"analysis_id": analysis_id})
        if value["analysis_id"] != analysis_id:
            raise ValueError("credential_binding_conflict")

    def capability(self, analysis_id):
        self.credentials = self.load()
        for item in self.credentials.values():
            if item.get("analysis_id") == analysis_id and item["origin"] == self.base:
                return item["token"]
        raise ValueError("analysis_credential_not_found")

    def request(self, path, method="GET", token=None, body=None, key=None, on_event=None, sse=True):
        headers = {"Accept": "text/event-stream" if body and sse else "application/json"}
        if token:
            headers["Authorization"] = "Bearer " + token
        if body:
            headers.update({"Content-Type": "application/json", "X-Model-API-Key": self.model_key, "Idempotency-Key": key})
        request = urllib.request.Request(self.base + path, data=json.dumps(body).encode() if body else None, headers=headers, method=method)
        try:
            response = self.http.open(request, timeout=190)
        except urllib.error.HTTPError as exc:
            try:
                error = json.load(exc)
            except ValueError:
                error = {"error": "http_" + str(exc.code)}
            if on_event:
                on_event("error", error)
            raise RuntimeError(error.get("error", "request_failed")) from None
        with response:
            if response.status == 204:
                return None
            if "text/event-stream" not in response.headers.get("Content-Type", ""):
                return json.load(response)
            event, data, final = "", [], None
            for raw in response:
                line = raw.decode("utf-8").rstrip("\r\n")
                if line.startswith("event: "):
                    event = line[7:]
                elif line.startswith("data: "):
                    data.append(line[6:])
                elif not line and data:
                    value = json.loads("\n".join(data))
                    if on_event:
                        on_event(event, value)
                    if event == "error":
                        raise RuntimeError(value["error"])
                    if event == "result":
                        final = value
                    event, data = "", []
            if final is None:
                raise RuntimeError("stream_interrupted_retry_same_idempotency_key")
            return final

    def create(self, question, policy, language="zh", idempotency_key=None, on_event=None):
        key = idempotency_key or str(uuid.uuid4())
        slot = self.base + ":" + key
        entry = self.entry(slot)

        def bind(name, event):
            if event.get("analysis_id"):
                self.bind(slot, event["analysis_id"])
            if on_event:
                on_event(name, event)

        result = self.request("/analyses", "POST", entry["token"], {"question": question, "policy": policy, "language": language, "provider": self.provider, "model": self.model}, key, bind)
        bind("result", result)
        return result

    def turn(self, analysis_id, question, policy=None, language="zh", idempotency_key=None):
        body = {"question": question, "language": language, "provider": self.provider, "model": self.model}
        if policy is not None:
            body["policy"] = policy
        return self.request("/analyses/" + urllib.parse.quote(analysis_id, safe="") + "/turns", "POST", self.capability(analysis_id), body, idempotency_key or str(uuid.uuid4()))

    def read(self, analysis_id):
        return self.request("/analyses/" + urllib.parse.quote(analysis_id, safe=""), token=self.capability(analysis_id))

    def delete(self, analysis_id):
        return self.request("/analyses/" + urllib.parse.quote(analysis_id, safe=""), "DELETE", self.capability(analysis_id))
