"""Analysis pipeline for the analysis framework.

This module defines :class:`AnalysisPipeline`, which executes
registered analyzers against an :class:`~analyzer.context.AnalysisContext`
and collects results.

The pipeline isolates analyzer failures: if one analyzer raises an
exception, the pipeline logs the error and continues with the next
analyzer.  The pipeline never crashes because of an analyzer failure.

When a database session is provided, the pipeline automatically
persists results: AnalyzerRun records and TrackFeature rows are
created for each successful analyzer execution.  Idempotency is
enforced by skipping analyzers whose version has already run
successfully on the given track.
"""

from __future__ import annotations

import time

from sqlalchemy.orm import Session

from analyzer.context import AnalysisContext
from analyzer.feature import FeatureSet
from analyzer.registry import AnalyzerRegistry
from analyzer.result import AnalysisResult
from log.logger import get_logger
from storage.analysis_repository import AnalysisRepository
from storage.feature_repository import FeatureRepository

_logger = get_logger(__name__)


class AnalysisPipeline:
    """Executes registered analyzers against an analysis context.

    Attributes:
        registry: The :class:`AnalyzerRegistry` providing analyzers.
    """

    def __init__(self, registry: AnalyzerRegistry) -> None:
        """Initialize the pipeline with a registry.

        Args:
            registry: The registry containing analyzers to execute.
        """
        self.registry = registry

    def run(
        self,
        context: AnalysisContext,
        session: Session | None = None,
    ) -> list[AnalysisResult]:
        """Execute all registered analyzers against the context.

        Analyzers run in registration order.  If an analyzer raises
        an exception, a failed :class:`AnalysisResult` is created
        and execution continues with the next analyzer.

        When ``session`` is provided, results are persisted:
        - Successful runs create AnalyzerRun + TrackFeature rows.
        - Already-executed analyzer versions are skipped (idempotency).
        - Failed runs create a failed AnalyzerRun record.

        Args:
            context: The analysis context to analyze.
            session: Optional SQLAlchemy session for persistence.
                When ``None``, no persistence occurs.

        Returns:
            A list of :class:`AnalysisResult` objects, one per
            registered analyzer (including skipped ones).
        """
        analyzers = self.registry.list_analyzers()
        results: list[AnalysisResult] = []
        track_id = context.track.id

        _logger.info(
            "pipeline.started",
            analyzer_count=len(analyzers),
            track=context.track.relative_path,
        )

        for analyzer in analyzers:
            if session is not None and FeatureRepository.has_analyzer_version(
                session,
                track_id,
                analyzer.name,
                analyzer.version,
            ):
                _logger.info(
                    "analysis.skipped",
                    analyzer=analyzer.name,
                    version=analyzer.version,
                    track=context.track.relative_path,
                )
                results.append(
                    AnalysisResult(
                        analyzer_name=analyzer.name,
                        analyzer_version=analyzer.version,
                        execution_time_ms=0.0,
                        success=True,
                        warnings=(),
                        feature_set=FeatureSet(),
                    )
                )
                continue

            _logger.info(
                "analyzer.started",
                analyzer=analyzer.name,
                version=analyzer.version,
                track=context.track.relative_path,
            )

            run = None
            if session is not None:
                run = AnalysisRepository.create_run(
                    session,
                    track_id,
                    analyzer.name,
                    analyzer.version,
                )

            start = time.perf_counter()
            try:
                result = analyzer.analyze(context)
                elapsed_ms = (time.perf_counter() - start) * 1000.0

                _logger.info(
                    "analyzer.finished",
                    analyzer=analyzer.name,
                    track=context.track.relative_path,
                    execution_time_ms=round(elapsed_ms, 3),
                    success=result.success,
                    feature_count=len(result.feature_set),
                )

                if session is not None and run is not None:
                    AnalysisRepository.finish_run(
                        session,
                        run,
                        execution_time_ms=elapsed_ms,
                        success=result.success,
                        warnings=result.warnings,
                    )

                    if result.success and len(result.feature_set) > 0:
                        saved = FeatureRepository.save_feature_set(
                            session,
                            track_id,
                            run.id,
                            result.feature_set,
                        )
                        _logger.info(
                            "feature.saved",
                            analyzer=analyzer.name,
                            track=context.track.relative_path,
                            run_id=run.id,
                        )
                        _logger.info(
                            "feature.count",
                            analyzer=analyzer.name,
                            track=context.track.relative_path,
                            count=saved,
                        )

                    _logger.info(
                        "analysis.saved",
                        analyzer=analyzer.name,
                        version=analyzer.version,
                        track=context.track.relative_path,
                        run_id=run.id,
                    )

                results.append(result)

            except Exception as exc:
                elapsed_ms = (time.perf_counter() - start) * 1000.0

                _logger.warning(
                    "analyzer.failed",
                    analyzer=analyzer.name,
                    track=context.track.relative_path,
                    error=str(exc),
                    execution_time_ms=round(elapsed_ms, 3),
                )

                if session is not None and run is not None:
                    AnalysisRepository.mark_failed(
                        session,
                        run,
                        execution_time_ms=elapsed_ms,
                        error=str(exc),
                    )
                    _logger.info(
                        "analysis.failed",
                        analyzer=analyzer.name,
                        version=analyzer.version,
                        track=context.track.relative_path,
                        run_id=run.id,
                    )

                failed_result = AnalysisResult(
                    analyzer_name=analyzer.name,
                    analyzer_version=analyzer.version,
                    execution_time_ms=elapsed_ms,
                    success=False,
                    warnings=(str(exc),),
                    feature_set=FeatureSet(),
                )
                results.append(failed_result)

        _logger.info(
            "pipeline.finished",
            track=context.track.relative_path,
            total_results=len(results),
            successes=sum(1 for r in results if r.success),
            failures=sum(1 for r in results if not r.success),
        )

        return results
