from __future__ import annotations

import json
import re
from typing import Any

from ..config import (
    get_repository_develop_branch,
    get_repository_generate_default_agent_id,
    get_repository_generate_default_route_prefix,
    get_repository_generate_default_token_name,
    get_repository_internal_base_url,
    get_repository_internal_share_token,
    get_repository_internal_share_token_key,
    get_repository_preview_address_template,
    get_repository_prod_branch,
)
from ..errors import QingflowApiError, raise_tool_error
from ..json_types import JSONObject, JSONValue
from ..repository_store import RepositoryMetadataStore
from ..session_store import SessionProfile
from .base import ToolBase


_STREAM_EVENT_RE = re.compile(r"^type=(?P<type>[^,]+),timestamp=(?P<timestamp>[^,]+),data=(?P<data>.*)$")
_STREAM_NEWLINE_TOKEN = "<wingsBr>"


class RepositoryDevTools(ToolBase):
    """仓库开发工具（中文名：仓库初始化与绑定）。

    类型：开发辅助工具。
    主要职责：
    1. 初始化仓库开发元数据；
    2. 维护仓库模板与分组映射；
    3. 为 AI Builder 的仓库流程提供状态读写。
    """

    def __init__(self, sessions, backend, *, metadata_store: RepositoryMetadataStore | None = None) -> None:
        """执行内部辅助逻辑。"""
        super().__init__(sessions, backend)
        self._metadata = metadata_store or RepositoryMetadataStore()

    def repository_init(self, *, profile: str, group_name: str, repo_template: str) -> JSONObject:
        """执行工具方法逻辑。"""
        normalized_group = str(group_name or "").strip()
        normalized_template = str(repo_template or "").strip()
        if not normalized_group:
            raise_tool_error(QingflowApiError.config_error("group_name is required"))
        if not normalized_template:
            raise_tool_error(QingflowApiError.config_error("repo_template is required"))

        def runner(session_profile: SessionProfile, context):
            payload = self._request_custom_page_json(
                session_profile,
                context,
                "POST",
                "/ultron/custom_page/v1/_init",
                params={"groupName": normalized_group, "repoTemplate": normalized_template},
            )
            repo_name = _extract_repo_name(payload)
            preview_address = _format_preview_address(repo_name)
            stored = self._metadata.put(
                repo_name,
                {
                    "group_name": normalized_group,
                    "repo_template": normalized_template,
                    "preview_address": preview_address,
                },
            )
            return {
                "status": "success",
                "repo_name": repo_name,
                "group_name": stored.get("group_name"),
                "repo_template": stored.get("repo_template"),
                "preview_address": preview_address,
                "verification": {"repo_initialized": True},
                "warnings": self._custom_page_route_warning(),
            }

        return self._run(profile, runner, require_workspace=True, tool_name='仓库初始化')

    def repository_generate(
        self,
        *,
        profile: str,
        repo_name: str,
        query: str,
        tag_id: int | None = None,
        app_keys: list[str] | None = None,
        extra_info: JSONObject | None = None,
        file_messages: list[JSONObject] | None = None,
        being_trace_log_enabled: bool = True,
        agent_id: int | None = None,
        allow_create_table: bool = False,
        route_prefix: str | None = None,
        token_name: str | None = None,
        session_id: str | None = None,
        round_version: int | None = None,
    ) -> JSONObject:
        """执行工具方法逻辑。"""
        normalized_repo = str(repo_name or "").strip()
        normalized_query = str(query or "").strip()
        if not normalized_repo:
            raise_tool_error(QingflowApiError.config_error("repo_name is required"))
        if not normalized_query:
            raise_tool_error(QingflowApiError.config_error("query is required"))

        metadata = self._metadata.get(normalized_repo) or {}
        normalized_tag_id = _normalize_optional_int(tag_id, "tag_id") or _normalize_optional_int(metadata.get("tag_id"), "tag_id")
        normalized_app_keys = _normalize_string_list(app_keys or metadata.get("app_keys") or [], field_name="app_keys")
        normalized_extra_info = _normalize_optional_object(extra_info, field_name="extra_info")
        normalized_file_messages = _normalize_optional_list_of_objects(file_messages, field_name="file_messages")
        normalized_agent_id = _normalize_optional_int(agent_id, "agent_id") or get_repository_generate_default_agent_id()
        normalized_route_prefix = _normalize_optional_string(route_prefix) or _normalize_optional_string(metadata.get("route_prefix")) or get_repository_generate_default_route_prefix()
        normalized_token_name = _normalize_optional_string(token_name) or _normalize_optional_string(metadata.get("token_name")) or get_repository_generate_default_token_name()
        normalized_session_id = _normalize_optional_string(session_id)
        normalized_round_version = _normalize_optional_int(round_version, "round_version")

        def runner(session_profile: SessionProfile, context):
            payload: JSONObject = {
                "query": normalized_query,
                "repoName": normalized_repo,
                "uid": session_profile.uid,
                "beingTraceLogEnabled": bool(being_trace_log_enabled),
                "allowCreateTable": bool(allow_create_table),
            }
            if normalized_tag_id is not None:
                payload["tagId"] = normalized_tag_id
            if normalized_app_keys:
                payload["appKeys"] = normalized_app_keys
            if normalized_extra_info is not None:
                payload["extraInfo"] = normalized_extra_info
            if normalized_file_messages:
                payload["fileMessages"] = normalized_file_messages
            if normalized_agent_id is not None:
                payload["agentId"] = normalized_agent_id
            if normalized_route_prefix is not None:
                payload["routePrefix"] = normalized_route_prefix
            if normalized_token_name is not None:
                payload["tokenName"] = normalized_token_name
            if normalized_session_id is not None:
                payload["sessionId"] = normalized_session_id
            if normalized_round_version is not None:
                payload["roundVersion"] = normalized_round_version

            stream_lines = self._request_custom_page_stream(
                session_profile,
                context,
                "POST",
                "/ultron/custom_page/v1/_generate",
                json_body=payload,
            )
            summarized = _summarize_generate_stream(
                repo_name=normalized_repo,
                query=normalized_query,
                stream_lines=stream_lines,
                route_warning=self._custom_page_route_warning(),
            )
            self._metadata.put(
                normalized_repo,
                {
                    "tag_id": normalized_tag_id,
                    "app_keys": normalized_app_keys,
                    "route_prefix": normalized_route_prefix,
                    "token_name": normalized_token_name,
                    "agent_id": normalized_agent_id,
                    "last_generate_query": normalized_query,
                    "last_generate_session_id": summarized.get("session_id"),
                    "last_generate_round_version": summarized.get("round_version"),
                },
            )
            return summarized

        return self._run(profile, runner, require_workspace=True, tool_name='仓库生成')

    def repository_publish_prod(
        self,
        *,
        profile: str,
        repo_name: str,
        confirm: bool = False,
    ) -> JSONObject:
        """执行工具方法逻辑。"""
        normalized_repo = str(repo_name or "").strip()
        if not normalized_repo:
            raise_tool_error(QingflowApiError.config_error("repo_name is required"))
        if confirm is not True:
            raise_tool_error(QingflowApiError.config_error("confirm=true is required for repository_publish_prod"))

        def runner(session_profile: SessionProfile, context):
            self._request_custom_page_json(
                session_profile,
                context,
                "POST",
                "/ultron/custom_page/v1/_publish",
                params={"repoName": normalized_repo},
            )
            preview_address = _lookup_preview_address(normalized_repo, self._metadata)
            return {
                "status": "success",
                "repo_name": normalized_repo,
                "source_branch": get_repository_develop_branch(),
                "target_branch": get_repository_prod_branch(),
                "merge_status": "success",
                "pipeline_status": "success",
                "published": True,
                "preview_address": preview_address,
                "verification": {"publish_confirmed": True, "pipeline_verified": True},
                "warnings": self._custom_page_route_warning(),
            }

        return self._run(profile, runner, require_workspace=True, tool_name='仓库发布生产')

    def _custom_page_route_warning(self) -> list[JSONObject]:
        """执行内部辅助逻辑。"""
        if get_repository_internal_base_url() and get_repository_internal_share_token():
            return []
        return [
            {
                "code": "CUSTOM_PAGE_ROUTE_PUBLIC_FALLBACK",
                "message": (
                    "repository tools are using the current Qingflow session route. "
                    "If the official custom-page internal route is not exposed in this environment, "
                    "configure repository.internal_base_url and repository.internal_share_token."
                ),
            }
        ]

    def _request_custom_page_json(
        self,
        session_profile: SessionProfile,
        context,
        method: str,
        path: str,
        *,
        params: JSONObject | None = None,
        json_body: JSONValue = None,
    ) -> JSONValue:
        """执行内部辅助逻辑。"""
        internal_route = _resolve_internal_custom_page_route(session_profile)
        if internal_route is not None:
            return self.backend.public_request_with_headers(
                method,
                internal_route["base_url"],
                path,
                params=params,
                json_body=json_body,
                headers=internal_route["headers"],
                qf_version=context.qf_version,
            ).data
        return self.backend.request(method, context, path, params=params, json_body=json_body)

    def _request_custom_page_stream(
        self,
        session_profile: SessionProfile,
        context,
        method: str,
        path: str,
        *,
        params: JSONObject | None = None,
        json_body: JSONValue = None,
    ) -> list[str]:
        """执行内部辅助逻辑。"""
        internal_route = _resolve_internal_custom_page_route(session_profile)
        if internal_route is not None:
            return self.backend.public_stream_request(
                method,
                internal_route["base_url"],
                path,
                params=params,
                json_body=json_body,
                headers=internal_route["headers"],
                qf_version=context.qf_version,
            )
        return self.backend.stream_request(method, context, path, params=params, json_body=json_body)


def _extract_repo_name(payload: Any) -> str:
    if isinstance(payload, str) and payload.strip():
        return payload.strip()
    if isinstance(payload, dict):
        for key in ("repoName", "repo_name"):
            value = payload.get(key)
            if isinstance(value, str) and value.strip():
                return value.strip()
    raise_tool_error(QingflowApiError(category="runtime", message="repository init did not return repo_name"))
    raise AssertionError("unreachable")


def _format_preview_address(repo_name: str) -> str:
    template = get_repository_preview_address_template()
    short_name = repo_name.split("/")[-1]
    if "%s" in template:
        return template % short_name
    return template.format(repo=short_name)


def _lookup_preview_address(repo_name: str, store: RepositoryMetadataStore) -> str | None:
    normalized = repo_name.strip()
    short_name = normalized.split("/")[-1]
    stored = store.get(short_name) or {}
    preview_address = stored.get("preview_address")
    return preview_address if isinstance(preview_address, str) and preview_address.strip() else None


def _resolve_internal_custom_page_route(session_profile: SessionProfile) -> dict[str, Any] | None:
    base_url = get_repository_internal_base_url()
    share_token = get_repository_internal_share_token()
    token_key = get_repository_internal_share_token_key()
    if not base_url and not share_token:
        return None
    if not base_url or not share_token:
        raise_tool_error(
            QingflowApiError.config_error(
                "repository.internal_base_url and repository.internal_share_token must be configured together"
            )
        )
    if session_profile.selected_ws_id is None:
        raise_tool_error(
            QingflowApiError.config_error("auth_use_credential must return a valid wsId before using the internal custom-page route")
        )
    return {
        "base_url": base_url,
        "headers": {
            token_key: share_token,
            "wsId": str(session_profile.selected_ws_id),
        },
    }


def _normalize_optional_string(value: Any) -> str | None:
    normalized = str(value or "").strip()
    return normalized or None


def _normalize_optional_int(value: Any, field_name: str) -> int | None:
    if value is None or value == "":
        return None
    try:
        return int(value)
    except (TypeError, ValueError):
        raise_tool_error(QingflowApiError.config_error(f"{field_name} must be an integer"))
    raise AssertionError("unreachable")


def _normalize_string_list(values: list[Any], *, field_name: str) -> list[str]:
    normalized: list[str] = []
    for item in values:
        text = str(item or "").strip()
        if not text:
            continue
        normalized.append(text)
    return normalized


def _normalize_optional_object(value: Any, *, field_name: str) -> JSONObject | None:
    if value is None:
        return None
    if not isinstance(value, dict):
        raise_tool_error(QingflowApiError.config_error(f"{field_name} must be an object"))
    return value


def _normalize_optional_list_of_objects(value: Any, *, field_name: str) -> list[JSONObject]:
    if value is None:
        return []
    if not isinstance(value, list):
        raise_tool_error(QingflowApiError.config_error(f"{field_name} must be a list"))
    normalized: list[JSONObject] = []
    for item in value:
        if not isinstance(item, dict):
            raise_tool_error(QingflowApiError.config_error(f"each {field_name} item must be an object"))
        normalized.append(item)
    return normalized


def _parse_stream_line(raw_line: str) -> JSONObject | None:
    line = str(raw_line or "").strip()
    if not line:
        return None
    if not line.startswith("data: "):
        return {"type": "raw", "data": line}
    body = line[6:]
    if body == "[done]":
        return {"type": "done", "data": "[done]"}
    match = _STREAM_EVENT_RE.match(body)
    if not match:
        return {"type": "raw", "data": body}
    event_type = match.group("type")
    timestamp_raw = str(match.group("timestamp") or "").strip()
    try:
        timestamp = int(timestamp_raw) if timestamp_raw else None
    except ValueError:
        timestamp = None
    data_text = match.group("data").replace(_STREAM_NEWLINE_TOKEN, "\n")
    data = _maybe_parse_json(data_text)
    return {
        "type": event_type,
        "timestamp": timestamp,
        "data": data,
    }


def _maybe_parse_json(value: str) -> JSONValue:
    text = str(value or "").strip()
    if not text:
        return ""
    try:
        return json.loads(text)
    except json.JSONDecodeError:
        return text


def _summarize_generate_stream(
    *,
    repo_name: str,
    query: str,
    stream_lines: list[str],
    route_warning: list[JSONObject],
) -> JSONObject:
    parsed_events = [event for event in (_parse_stream_line(line) for line in stream_lines) if event is not None]
    done_seen = any(str(event.get("type")) == "done" for event in parsed_events)
    result_event = next((event for event in reversed(parsed_events) if event.get("type") == "result"), None)
    error_event = next((event for event in reversed(parsed_events) if event.get("type") == "error"), None)
    repository_events = [
        event.get("data")
        for event in parsed_events
        if event.get("type") == "event"
        and isinstance(event.get("data"), dict)
        and event["data"].get("action") == "REPOSITORY_COMMIT"
    ]
    other_events = [
        event.get("data")
        for event in parsed_events
        if event.get("type") == "event"
        and not (
            isinstance(event.get("data"), dict)
            and event["data"].get("action") == "REPOSITORY_COMMIT"
        )
    ]
    token_count = sum(1 for event in parsed_events if event.get("type") == "token")
    running_count = sum(1 for event in parsed_events if event.get("type") == "runing")
    waiting_count = sum(1 for event in parsed_events if event.get("type") == "waiting")
    raw_count = sum(1 for event in parsed_events if event.get("type") == "raw")
    warnings = list(route_warning)
    if raw_count:
        warnings.append(
            {
                "code": "GENERATE_STREAM_UNPARSED_CHUNKS",
                "message": f"generate stream contained {raw_count} unparsed chunk(s); raw chunks are returned for debugging",
            }
        )
    if not done_seen:
        warnings.append(
            {
                "code": "GENERATE_STREAM_DONE_MISSING",
                "message": "generate stream ended without an explicit [done] marker",
            }
        )

    if error_event is not None:
        error_data = error_event.get("data")
        error_code: str | None = None
        message = "repository generation failed"
        if isinstance(error_data, dict):
            error_code = str(error_data.get("errorCode") or error_data.get("code") or "").strip() or None
            message = str(error_data.get("errorMessage") or error_data.get("message") or message)
        elif isinstance(error_data, str) and error_data.strip():
            message = error_data.strip()
        return {
            "status": "failed",
            "error_code": error_code or "REPOSITORY_GENERATE_FAILED",
            "message": message,
            "repo_name": repo_name,
            "query": query,
            "repository_events": repository_events,
            "events": other_events,
            "stream_summary": {
                "token_events": token_count,
                "running_events": running_count,
                "waiting_events": waiting_count,
                "done_seen": done_seen,
                "raw_chunks": raw_count,
            },
            "verification": {
                "stream_done": done_seen,
                "result_received": False,
                "repository_commit_detected": bool(repository_events),
            },
            "warnings": warnings,
        }

    result_payload = result_event.get("data") if isinstance(result_event, dict) else None
    if not isinstance(result_payload, dict):
        return {
            "status": "failed",
            "error_code": "REPOSITORY_GENERATE_EMPTY_RESULT",
            "message": "repository generation stream completed without a structured result payload",
            "repo_name": repo_name,
            "query": query,
            "repository_events": repository_events,
            "events": other_events,
            "stream_summary": {
                "token_events": token_count,
                "running_events": running_count,
                "waiting_events": waiting_count,
                "done_seen": done_seen,
                "raw_chunks": raw_count,
            },
            "verification": {
                "stream_done": done_seen,
                "result_received": False,
                "repository_commit_detected": bool(repository_events),
            },
            "warnings": warnings,
        }

    response_status = str(result_payload.get("responseStatus") or "SUCCESS")
    return {
        "status": "success" if response_status.upper() == "SUCCESS" else "failed",
        "repo_name": repo_name,
        "query": query,
        "response_status": response_status,
        "session_id": result_payload.get("sessionId"),
        "round_version": result_payload.get("roundVersion"),
        "thread_id": result_payload.get("threadId"),
        "answer": result_payload.get("answer"),
        "answer_json": result_payload.get("answerJson"),
        "trace_logs": result_payload.get("traceLogs"),
        "total_tokens": result_payload.get("totalTokens"),
        "credit_consume": result_payload.get("creditConsume"),
        "history_messages": result_payload.get("historyMessages"),
        "repository_events": repository_events,
        "events": other_events,
        "stream_summary": {
            "token_events": token_count,
            "running_events": running_count,
            "waiting_events": waiting_count,
            "done_seen": done_seen,
            "raw_chunks": raw_count,
        },
        "verification": {
            "stream_done": done_seen,
            "result_received": True,
            "repository_commit_detected": bool(repository_events),
        },
        "warnings": warnings,
    }
