from __future__ import annotations

import json
from typing import Any, TypeVar

from pydantic import BaseModel

from app.models.embedding_record import EmbeddingRecord, validate_embedding_record
from app.models.stream_events import CommitAnalysisReadyEvent, GraphArtifactReadyEvent

T = TypeVar("T", bound=BaseModel)


def parse_stream_payload(raw: str | bytes | dict[str, Any]) -> dict[str, Any]:
    """
    Parse Redis stream payload field into a JSON object.

    Pipeline services store the event JSON in a hash field named ``payload``.
    """
    if isinstance(raw, dict):
        if "payload" in raw and len(raw) == 1:
            nested = raw["payload"]
            if isinstance(nested, str):
                return parse_stream_payload(nested)
            if isinstance(nested, dict):
                return nested
        return raw

    if isinstance(raw, bytes):
        raw = raw.decode("utf-8")

    if not isinstance(raw, str):
        raise ValueError("stream payload must be a string, bytes, or dict")

    text = raw.strip()
    if not text:
        raise ValueError("stream payload is empty")

    try:
        parsed = json.loads(text)
    except json.JSONDecodeError as exc:
        raise ValueError("stream payload is not valid JSON") from exc

    if not isinstance(parsed, dict):
        raise ValueError("stream payload must be a JSON object")

    return parsed


def parse_stream_event(raw: str | bytes | dict[str, Any], model: type[T]) -> T:
    """Parse and validate a Redis stream payload against a typed event model."""
    return model.model_validate(parse_stream_payload(raw))


def parse_graph_artifact_ready(raw: str | bytes | dict[str, Any]) -> GraphArtifactReadyEvent:
    return parse_stream_event(raw, GraphArtifactReadyEvent)


def parse_commit_analysis_ready(raw: str | bytes | dict[str, Any]) -> CommitAnalysisReadyEvent:
    return parse_stream_event(raw, CommitAnalysisReadyEvent)


__all__ = [
    "parse_stream_payload",
    "parse_stream_event",
    "parse_graph_artifact_ready",
    "parse_commit_analysis_ready",
    "validate_embedding_record",
    "EmbeddingRecord",
]
