package parser

import (
	"context"
	"fmt"
	"log/slog"
	"os"
	"runtime"
	"strings"
	"sync"
	"time"

	contractevents "bit.admedia.com/scm/ad/adpilot-indexing-code-parser.com/internal/contracts/events"
	"bit.admedia.com/scm/ad/adpilot-indexing-code-parser.com/internal/contracts/graph"
	"bit.admedia.com/scm/ad/adpilot-indexing-code-parser.com/internal/service/parser/analyzererrors"
)

type fileJob struct {
	jobID string
	event contractevents.FileChangedEvent
}

type parseResult struct {
	jobID         string
	event         contractevents.FileChangedEvent
	status        string
	artifact      graph.Artifact
	chunks        []contractevents.CodeChunk
	artifactID    string
	parseDuration time.Duration
	usedFallback  bool
	err           error
}

type skipDecision struct {
	reason              string
	diagnostics         []contractevents.ParserDiagnostic
	countTowardExpected bool
}

// ProcessFileChanged runs parse orchestration for a single files.changed event.
func (s *Service) ProcessFileChanged(ctx context.Context, event contractevents.FileChangedEvent) error {
	return s.ProcessFiles(ctx, []contractevents.FileChangedEvent{event})
}

// ProcessFiles runs parse orchestration for a batch of files.changed events using a bounded worker pool.
func (s *Service) ProcessFiles(ctx context.Context, events []contractevents.FileChangedEvent) error {
	if len(events) == 0 {
		return nil
	}

	workers := s.cfg.Workers
	if workers <= 0 {
		workers = runtime.NumCPU()
	}
	if workers > len(events) {
		workers = len(events)
	}

	jobs := make(chan fileJob)
	results := make(chan parseResult, len(events))

	var workerWG sync.WaitGroup
	for i := 0; i < workers; i++ {
		workerWG.Add(1)
		go func() {
			defer workerWG.Done()
			for job := range jobs {
				results <- s.runFileJob(ctx, job)
			}
		}()
	}

	var coordWG sync.WaitGroup
	var resultErr error
	var resultErrMu sync.Mutex
	coordWG.Add(1)
	go func() {
		defer coordWG.Done()
		for result := range results {
			if err := s.handleParseResult(ctx, result); err != nil {
				resultErrMu.Lock()
				if resultErr == nil {
					resultErr = err
				}
				resultErrMu.Unlock()
			}
		}
	}()

enqueue:
	for _, event := range events {
		if ctx.Err() != nil {
			break enqueue
		}

		jobID := jobIDFromEvent(event)
		s.setJob(s.parseJobFromEvent(event, jobID, JobStatusPending))

		if skip := s.shouldSkip(event); skip != nil {
			s.updateJobStatus(jobID, JobStatusSkipped)
			_ = s.publishDiagnostics(ctx, event, jobID, skip.diagnostics)
			s.logger().Info("parse job skipped",
				"job_id", jobID,
				"event_id", event.EventID,
				"repo_id", event.RepoID,
				"snapshot_id", event.SnapshotID,
				"file_path", event.FilePath,
				"status", JobStatusSkipped,
				"reason", skip.reason,
			)
			if skip.countTowardExpected {
				if err := s.recordCodeFileOutcome(ctx, event); err != nil {
					return err
				}
			}
			continue
		}

		if _, ok := s.analyzers[event.Language]; !ok && s.fallback == nil {
			s.updateJobStatus(jobID, JobStatusSkipped)
			_ = s.publishDiagnostics(ctx, event, jobID, []contractevents.ParserDiagnostic{
				{
					Code:     "unsupported_language",
					Message:  fmt.Sprintf("no analyzer for language %q", event.Language),
					Severity: contractevents.DiagnosticSeverityWarning,
					FilePath: event.FilePath,
				},
			})
			if err := s.recordCodeFileOutcome(ctx, event); err != nil {
				return err
			}
			continue
		}

		jobs <- fileJob{jobID: jobID, event: event}
	}
	close(jobs)

	workerWG.Wait()
	close(results)
	coordWG.Wait()

	if err := ctx.Err(); err != nil {
		return err
	}
	return resultErr
}

func (s *Service) shouldSkip(event contractevents.FileChangedEvent) *skipDecision {
	// Non-code files must be skipped WITHOUT counting toward the code fan-in.
	// This must run BEFORE IsIgnoredPath below: an ignored file that is not
	// code-kind (e.g. a data/asset file that also matches an ignore rule) would
	// otherwise be counted toward processed_code_files, while repo-sync's
	// expected_code_files only counts file_kind==code. That inflates processed
	// past expected and trips premature snapshot-graph finalization, leaving the
	// repo under-indexed. Keep both counters over the same set (code-kind files).
	if event.FileKind != "" && event.FileKind != contractevents.FileKindCode {
		return &skipDecision{reason: "non_code_file_kind"}
	}

	if IsIgnoredPath(event.FilePath) {
		return &skipDecision{
			reason:              "ignored_path",
			countTowardExpected: true,
			diagnostics: []contractevents.ParserDiagnostic{
				{
					Code:     "ignored_path",
					Message:  fmt.Sprintf("skipping ignored file %q", event.FilePath),
					Severity: contractevents.DiagnosticSeverityInfo,
					FilePath: event.FilePath,
				},
			},
		}
	}

	if event.ChangeType == contractevents.ChangeTypeDeleted {
		return &skipDecision{
			reason:              "file_deleted",
			countTowardExpected: true,
			diagnostics: []contractevents.ParserDiagnostic{
				{
					Code:     "file_deleted",
					Message:  fmt.Sprintf("skipping parse for deleted file %q", event.FilePath),
					Severity: contractevents.DiagnosticSeverityInfo,
					FilePath: event.FilePath,
				},
			},
		}
	}

	return nil
}

func (s *Service) runFileJob(ctx context.Context, job fileJob) parseResult {
	s.updateJobStatus(job.jobID, JobStatusRunning)

	if err := ctx.Err(); err != nil {
		return parseResult{
			jobID:  job.jobID,
			event:  job.event,
			status: JobStatusFailed,
			err:    err,
		}
	}

	sourcePath, err := ResolveSourcePath(s.cfg.WorkspacePath, job.event.RepoID, job.event.FilePath)
	if err != nil {
		return parseResult{
			jobID:  job.jobID,
			event:  job.event,
			status: JobStatusFailed,
			err:    fmt.Errorf("resolve source path: %w", err),
		}
	}

	source, err := os.ReadFile(sourcePath)
	if err != nil {
		return parseResult{
			jobID:  job.jobID,
			event:  job.event,
			status: JobStatusFailed,
			err:    fmt.Errorf("read source file %q: %w", sourcePath, err),
		}
	}

	parseStarted := time.Now()
	artifact, chunks, err, usedFallback := s.parseSource(ctx, source, job.event)
	parseDuration := time.Since(parseStarted)
	if err != nil {
		return parseResult{
			jobID:         job.jobID,
			event:         job.event,
			status:        JobStatusFailed,
			parseDuration: parseDuration,
			err:           fmt.Errorf("parse file: %w", err),
		}
	}

	if artifact == nil {
		artifact = &graph.Artifact{SchemaVersion: graph.SchemaVersionV1}
	}

	enriched := *artifact
	enriched.RepoID = job.event.RepoID
	enriched.SnapshotID = job.event.SnapshotID
	enriched.CommitSHA = job.event.CommitSHA
	if enriched.SchemaVersion == "" {
		enriched.SchemaVersion = graph.SchemaVersionV1
	}

	hydratedChunks := hydrateChunkText(source, chunks)

	result := parseResult{
		jobID:         job.jobID,
		event:         job.event,
		status:        JobStatusCompleted,
		artifact:      enriched,
		chunks:        hydratedChunks,
		parseDuration: parseDuration,
		usedFallback:  usedFallback,
	}
	return result
}

func (s *Service) parseSource(ctx context.Context, source []byte, event contractevents.FileChangedEvent) (*graph.Artifact, []contractevents.CodeChunk, error, bool) {
	if analyzer, ok := s.analyzers[event.Language]; ok {
		artifact, chunks, err := analyzer.ParseFile(ctx, source, event.FilePath)
		return artifact, chunks, err, false
	}
	if s.fallback == nil {
		return nil, nil, fmt.Errorf("no analyzer for language %q", event.Language), false
	}
	artifact, chunks, err := s.fallback.ParseFile(ctx, source, event.FilePath, event.Language)
	return artifact, chunks, err, true
}

func hydrateChunkText(source []byte, chunks []contractevents.CodeChunk) []contractevents.CodeChunk {
	if len(chunks) == 0 {
		return chunks
	}

	lines := strings.Split(string(source), "\n")
	out := make([]contractevents.CodeChunk, len(chunks))
	for i, chunk := range chunks {
		out[i] = chunk
		if chunk.StartLine < 1 || chunk.EndLine < chunk.StartLine {
			continue
		}
		start := chunk.StartLine - 1
		end := chunk.EndLine
		if start >= len(lines) {
			continue
		}
		if end > len(lines) {
			end = len(lines)
		}
		out[i].Text = strings.Join(lines[start:end], "\n")
	}
	return out
}

func (s *Service) handleParseResult(ctx context.Context, result parseResult) error {
	switch result.status {
	case JobStatusCompleted:
		if s.artifacts == nil {
			s.updateJobStatus(result.jobID, JobStatusFailed)
			_ = s.publishDiagnostics(ctx, result.event, result.jobID, []contractevents.ParserDiagnostic{
				{
					Code:     "artifact_store_unavailable",
					Message:  "artifact store is not configured",
					Severity: contractevents.DiagnosticSeverityError,
					FilePath: result.event.FilePath,
				},
			})
			return fmt.Errorf("artifact store is not configured")
		}

		artifactID, artifactURI, err := s.artifacts.SaveArtifact(ctx, result.artifact)
		if err != nil {
			s.updateJobStatus(result.jobID, JobStatusFailed)
			_ = s.publishDiagnostics(ctx, result.event, result.jobID, []contractevents.ParserDiagnostic{
				{
					Code:     "artifact_save_failed",
					Message:  err.Error(),
					Severity: contractevents.DiagnosticSeverityError,
					FilePath: result.event.FilePath,
				},
			})
			return fmt.Errorf("save artifact: %w", err)
		}

		if s.finalize != nil && len(result.chunks) > 0 {
			if err := s.finalize.SaveCodeChunks(
				ctx,
				result.event.RepoID,
				result.event.SnapshotID,
				result.event.CommitSHA,
				result.event.FilePath,
				result.event.Language,
				artifactID,
				result.chunks,
			); err != nil {
				s.updateJobStatus(result.jobID, JobStatusFailed)
				_ = s.publishDiagnostics(ctx, result.event, result.jobID, []contractevents.ParserDiagnostic{
					{
						Code:     "code_chunks_save_failed",
						Message:  err.Error(),
						Severity: contractevents.DiagnosticSeverityError,
						FilePath: result.event.FilePath,
					},
				})
				return fmt.Errorf("save code chunks: %w", err)
			}
		}

		if s.finalize != nil {
			if err := s.finalize.IncrementProcessedCodeFiles(ctx, result.event.RepoID, result.event.SnapshotID, result.event.FilePath); err != nil {
				s.updateJobStatus(result.jobID, JobStatusFailed)
				_ = s.publishDiagnostics(ctx, result.event, result.jobID, []contractevents.ParserDiagnostic{
					{
						Code:     "indexing_run_increment_failed",
						Message:  err.Error(),
						Severity: contractevents.DiagnosticSeverityError,
						FilePath: result.event.FilePath,
					},
				})
				return fmt.Errorf("increment processed_code_files: %w", err)
			}
		}

		if err := s.publishSuccess(ctx, result.event, result.jobID, result.artifact, result.chunks, artifactID, artifactURI); err != nil {
			s.updateJobStatus(result.jobID, JobStatusFailed)
			_ = s.publishDiagnostics(ctx, result.event, result.jobID, []contractevents.ParserDiagnostic{
				{
					Code:     "publish_failed",
					Message:  err.Error(),
					Severity: contractevents.DiagnosticSeverityError,
					FilePath: result.event.FilePath,
				},
			})
			return fmt.Errorf("publish success: %w", err)
		}

		s.setJobArtifactID(result.jobID, artifactID)
		s.updateJobStatus(result.jobID, JobStatusCompleted)
		if result.usedFallback {
			_ = s.publishDiagnostics(ctx, result.event, result.jobID, []contractevents.ParserDiagnostic{
				{
					Code:     "fallback_parser_used",
					Message:  fmt.Sprintf("parsed %q using fallback parser for language %q", result.event.FilePath, result.event.Language),
					Severity: contractevents.DiagnosticSeverityInfo,
					FilePath: result.event.FilePath,
				},
			})
		}
		s.logParseResult(result, artifactID, JobStatusCompleted, "")

		if err := s.maybeFinalizeSnapshot(ctx, result.event.RepoID, result.event.SnapshotID); err != nil {
			return err
		}
		return nil

	case JobStatusFailed:
		s.updateJobStatus(result.jobID, JobStatusFailed)
		diagnostics := analyzererrors.DiagnosticsFromError(result.err)
		if len(diagnostics) == 0 {
			message := "parse failed"
			if result.err != nil {
				message = result.err.Error()
			}
			diagnostics = []contractevents.ParserDiagnostic{
				{
					Code:     "parse_failed",
					Message:  message,
					Severity: contractevents.DiagnosticSeverityError,
					FilePath: result.event.FilePath,
				},
			}
		}
		_ = s.publishDiagnostics(ctx, result.event, result.jobID, diagnostics)
		errMsg := ""
		if result.err != nil {
			errMsg = result.err.Error()
		}
		s.logParseResult(result, "", JobStatusFailed, errMsg)
		if result.err != nil {
			return result.err
		}
		return fmt.Errorf("parse failed")
	}
	return nil
}

func (s *Service) recordCodeFileOutcome(ctx context.Context, event contractevents.FileChangedEvent) error {
	if s.finalize == nil {
		return nil
	}
	if err := s.finalize.IncrementProcessedCodeFiles(ctx, event.RepoID, event.SnapshotID, event.FilePath); err != nil {
		return fmt.Errorf("increment processed_code_files: %w", err)
	}
	// This path (skipped / ignored / deleted / unsupported-language files) counts
	// toward the code fan-in but produces NO code_graph_artifact. Track it
	// separately so the finalize artifact-completeness gate can subtract these
	// from its target — otherwise the gate waits for artifacts that will never
	// exist and finalization hangs forever (no graph, no embeddings).
	if err := s.finalize.IncrementSkippedCodeFiles(ctx, event.RepoID, event.SnapshotID, event.FilePath); err != nil {
		return fmt.Errorf("increment skipped_code_files: %w", err)
	}
	return s.maybeFinalizeSnapshot(ctx, event.RepoID, event.SnapshotID)
}

func (s *Service) logger() *slog.Logger {
	if s.cfg.Logger != nil {
		return s.cfg.Logger
	}
	return slog.Default()
}

func (s *Service) logParseResult(result parseResult, artifactID, status, errMsg string) {
	attrs := []any{
		"job_id", result.jobID,
		"event_id", result.event.EventID,
		"repo_id", result.event.RepoID,
		"snapshot_id", result.event.SnapshotID,
		"file_path", result.event.FilePath,
		"status", status,
	}
	if result.parseDuration > 0 {
		attrs = append(attrs, "parse_duration_ms", result.parseDuration.Milliseconds())
	}
	if artifactID != "" {
		attrs = append(attrs, "artifact_id", artifactID)
	}
	if errMsg != "" {
		attrs = append(attrs, "error", errMsg)
	}

	switch status {
	case JobStatusCompleted:
		s.logger().Info("parse job completed", attrs...)
	case JobStatusFailed:
		s.logger().Error("parse job failed", attrs...)
	default:
		s.logger().Info("parse job finished", attrs...)
	}
}
