package parser

import (
	"context"
	"errors"
	"fmt"
	"strings"
	"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/contracts/metadata"
	"bit.admedia.com/scm/ad/adpilot-indexing-code-parser.com/internal/contracts/validate"
	"bit.admedia.com/scm/ad/adpilot-indexing-code-parser.com/internal/store"
)

// ErrCommitOrderNotReady indicates the parent commit delta is not yet available.
var ErrCommitOrderNotReady = errors.New("parent commit delta not ready")

// ProcessCommitsChanged resolves a technical commit diff and publishes graph.delta.ready.
func (s *Service) ProcessCommitsChanged(ctx context.Context, event contractevents.CommitsChangedEvent) error {
	started := time.Now()
	jobID := CommitJobIDFromEvent(event)
	s.setJob(commitJobFromEvent(event, jobID, JobStatusRunning))

	if s.commitOrder != nil {
		ready, err := s.commitOrder.ParentCommitReady(ctx, event.RepoID, event.SnapshotID, event.ParentSHA)
		if err != nil {
			s.updateJobStatus(jobID, JobStatusFailed)
			return fmt.Errorf("check parent commit readiness: %w", err)
		}
		if !ready {
			s.updateJobStatus(jobID, JobStatusPending)
			return ErrCommitOrderNotReady
		}
	}

	codeReady, err := s.codeParseReady(ctx, event.RepoID, event.SnapshotID)
	if err != nil {
		s.updateJobStatus(jobID, JobStatusFailed)
		return err
	}
	if !codeReady {
		s.updateJobStatus(jobID, JobStatusPending)
		return ErrCommitArtifactsNotReady
	}

	diff, err := ResolveCommitDiff(s.cfg.WorkspacePath, event)
	if err != nil {
		if errors.Is(err, ErrNoCodeFilesChanged) {
			nonCodeDiff, nonCodeErr := ResolveNonCodeCommitDiff(s.cfg.WorkspacePath, event)
			if nonCodeErr != nil {
				if errors.Is(nonCodeErr, ErrNoCodeFilesChanged) {
					if markErr := s.markCommitFanInProgress(ctx, event.RepoID, event.SnapshotID, event.CommitSHA, event.ParentSHA); markErr != nil {
						s.updateJobStatus(jobID, JobStatusFailed)
						return markErr
					}
					s.updateJobStatus(jobID, JobStatusSkipped)
					s.logger().Info("commit skipped: no changed files",
						"job_id", jobID,
						"event_id", event.EventID,
						"repo_id", event.RepoID,
						"snapshot_id", event.SnapshotID,
						"commit_sha", event.CommitSHA,
						"parent_sha", event.ParentSHA,
					)
					return nil
				}
				s.updateJobStatus(jobID, JobStatusFailed)
				return nonCodeErr
			}
			if err := s.publishNonCodeCommitDelta(ctx, event, nonCodeDiff, jobID); err != nil {
				s.updateJobStatus(jobID, JobStatusFailed)
				return err
			}
			s.updateJobStatus(jobID, JobStatusCompleted)
			return nil
		}
		s.logger().Error("commit diff resolution failed",
			"job_id", jobID,
			"event_id", event.EventID,
			"repo_id", event.RepoID,
			"snapshot_id", event.SnapshotID,
			"commit_sha", event.CommitSHA,
			"parent_sha", event.ParentSHA,
			"error", err,
		)
		_ = s.publishCommitDiagnostics(ctx, event, jobID, []contractevents.ParserDiagnostic{
			{
				Code:     "commit_diff_failed",
				Message:  err.Error(),
				Severity: contractevents.DiagnosticSeverityError,
			},
		})
		s.updateJobStatus(jobID, JobStatusFailed)
		return err
	}

	if s.artifacts == nil {
		err := fmt.Errorf("artifact store is not configured")
		_ = s.publishCommitDiagnostics(ctx, event, jobID, []contractevents.ParserDiagnostic{
			{
				Code:     "artifact_store_unavailable",
				Message:  err.Error(),
				Severity: contractevents.DiagnosticSeverityError,
			},
		})
		s.updateJobStatus(jobID, JobStatusFailed)
		return err
	}

	fileDeltas := make([]graph.FileDelta, 0, len(diff.ChangedFiles))
	var missingArtifacts []string
	rootCommit := isRootCommitBase(diff.BaseSHA)

	for _, filePath := range diff.ChangedFiles {
		if err := ctx.Err(); err != nil {
			return err
		}

		baseArtifact, baseID, err := s.findBaseArtifact(ctx, diff.RepoID, diff.SnapshotID, diff.BaseSHA, filePath, rootCommit)
		if err != nil {
			s.updateJobStatus(jobID, JobStatusFailed)
			return fmt.Errorf("resolve base artifact for %q: %w", filePath, err)
		}

		targetArtifact, targetID, err := s.findTargetArtifact(ctx, diff.RepoID, diff.SnapshotID, diff.TargetSHA, filePath)
		if err != nil {
			if errors.Is(err, store.ErrArtifactNotFound) {
				missingArtifacts = append(missingArtifacts, fmt.Sprintf("target:%s", filePath))
				continue
			}
			s.updateJobStatus(jobID, JobStatusFailed)
			return fmt.Errorf("resolve target artifact for %q: %w", filePath, err)
		}

		fileDeltas = append(fileDeltas, graph.ComputeFileDelta(baseArtifact, targetArtifact, filePath, baseID, targetID))
	}

	if len(fileDeltas) == 0 {
		if len(missingArtifacts) > 0 && missingArtifactsAllTargets(missingArtifacts) {
			s.updateJobStatus(jobID, JobStatusPending)
			return ErrCommitArtifactsNotReady
		}

		message := fmt.Sprintf("no comparable graph artifacts for commit pair %s..%s", diff.BaseSHA, diff.TargetSHA)
		if len(missingArtifacts) > 0 {
			message += "; missing: " + strings.Join(missingArtifacts, ", ")
		}
		err := fmt.Errorf("%s", message)
		_ = s.publishCommitDiagnostics(ctx, event, jobID, []contractevents.ParserDiagnostic{
			{
				Code:     "graph_delta_unavailable",
				Message:  message,
				Severity: contractevents.DiagnosticSeverityWarning,
			},
		})
		s.updateJobStatus(jobID, JobStatusFailed)
		return err
	}

	delta := graph.AggregateDelta(diff.RepoID, diff.SnapshotID, diff.BaseSHA, diff.TargetSHA, diff.ChangedFiles, fileDeltas)
	delta.GitCommit = gitCommitMetadataFromEvent(event)
	deltaID, deltaURI, err := s.artifacts.SaveDelta(ctx, delta)
	if err != nil {
		_ = s.publishCommitDiagnostics(ctx, event, jobID, []contractevents.ParserDiagnostic{
			{
				Code:     "delta_save_failed",
				Message:  err.Error(),
				Severity: contractevents.DiagnosticSeverityError,
			},
		})
		s.updateJobStatus(jobID, JobStatusFailed)
		return fmt.Errorf("save delta: %w", err)
	}

	if err := s.markCommitFanInProgress(ctx, diff.RepoID, diff.SnapshotID, diff.TargetSHA, diff.BaseSHA); err != nil {
		s.updateJobStatus(jobID, JobStatusFailed)
		return err
	}

	if err := s.publishGraphDeltaReady(ctx, event, diff, deltaID, deltaURI, delta); err != nil {
		_ = s.publishCommitDiagnostics(ctx, event, jobID, []contractevents.ParserDiagnostic{
			{
				Code:     "publish_failed",
				Message:  err.Error(),
				Severity: contractevents.DiagnosticSeverityError,
			},
		})
		s.updateJobStatus(jobID, JobStatusFailed)
		return err
	}

	s.updateJobStatus(jobID, JobStatusCompleted)

	addedNodes, removedNodes, changedEdges := delta.Counts()
	s.logger().Info("commit graph delta completed",
		"job_id", jobID,
		"event_id", event.EventID,
		"repo_id", diff.RepoID,
		"snapshot_id", diff.SnapshotID,
		"base_commit_sha", diff.BaseSHA,
		"target_commit_sha", diff.TargetSHA,
		"delta_artifact_id", deltaID,
		"changed_files", len(diff.ChangedFiles),
		"compared_files", len(fileDeltas),
		"added_nodes", addedNodes,
		"removed_nodes", removedNodes,
		"changed_edges", changedEdges,
		"duration_ms", time.Since(started).Milliseconds(),
	)

	if len(missingArtifacts) > 0 {
		s.logger().Warn("commit graph delta partial",
			"job_id", jobID,
			"event_id", event.EventID,
			"missing_artifacts", strings.Join(missingArtifacts, ", "),
		)
	}

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

	return nil
}

func (s *Service) publishNonCodeCommitDelta(
	ctx context.Context,
	event contractevents.CommitsChangedEvent,
	diff CommitDiff,
	jobID string,
) error {
	if s.artifacts == nil {
		return fmt.Errorf("artifact store is not configured")
	}

	delta := graph.AggregateDelta(
		diff.RepoID,
		diff.SnapshotID,
		diff.BaseSHA,
		diff.TargetSHA,
		diff.ChangedFiles,
		nil,
	)
	delta.GitCommit = gitCommitMetadataFromEvent(event)

	deltaID, deltaURI, err := s.artifacts.SaveDelta(ctx, delta)
	if err != nil {
		_ = s.publishCommitDiagnostics(ctx, event, jobID, []contractevents.ParserDiagnostic{
			{
				Code:     "delta_save_failed",
				Message:  err.Error(),
				Severity: contractevents.DiagnosticSeverityError,
			},
		})
		return fmt.Errorf("save non-code delta: %w", err)
	}

	if err := s.markCommitFanInProgress(ctx, diff.RepoID, diff.SnapshotID, diff.TargetSHA, diff.BaseSHA); err != nil {
		return err
	}

	if err := s.publishGraphDeltaReady(ctx, event, diff, deltaID, deltaURI, delta); err != nil {
		_ = s.publishCommitDiagnostics(ctx, event, jobID, []contractevents.ParserDiagnostic{
			{
				Code:     "publish_failed",
				Message:  err.Error(),
				Severity: contractevents.DiagnosticSeverityError,
			},
		})
		return err
	}

	s.logger().Info("non-code commit delta completed",
		"job_id", jobID,
		"event_id", event.EventID,
		"repo_id", diff.RepoID,
		"snapshot_id", diff.SnapshotID,
		"base_commit_sha", diff.BaseSHA,
		"target_commit_sha", diff.TargetSHA,
		"changed_files", len(diff.ChangedFiles),
	)
	return nil
}

func (s *Service) publishGraphDeltaReady(
	ctx context.Context,
	event contractevents.CommitsChangedEvent,
	diff CommitDiff,
	deltaID, deltaURI string,
	delta graph.Delta,
) error {
	if s.events == nil {
		return nil
	}

	addedNodes, removedNodes, changedEdges := delta.Counts()
	graphEvent := contractevents.NewGraphDeltaReadyEvent(
		graphDeltaEventID(diff.SnapshotID, diff.BaseSHA, diff.TargetSHA),
		diff.RepoID,
		diff.BaseSHA,
		diff.TargetSHA,
		deltaID,
		deltaURI,
		GraphSchemaVersion,
	)
	graphEvent.AddedNodes = addedNodes
	graphEvent.RemovedNodes = removedNodes
	graphEvent.ChangedEdges = changedEdges
	graphEvent.GitCommitMetadata = gitCommitMetadataFromEvent(event)

	if err := validate.GraphDeltaReady(graphEvent); err != nil {
		return fmt.Errorf("validate graph.delta.ready: %w", err)
	}
	if err := s.events.PublishGraphDeltaReady(ctx, graphEvent); err != nil {
		return fmt.Errorf("publish graph.delta.ready: %w", err)
	}
	return nil
}

func (s *Service) publishCommitDiagnostics(
	ctx context.Context,
	event contractevents.CommitsChangedEvent,
	jobID string,
	diagnostics []contractevents.ParserDiagnostic,
) error {
	if len(diagnostics) > 0 {
		s.setJobDiagnostics(jobID, diagnostics)
	}
	if s.finalize != nil && len(diagnostics) > 0 {
		if err := s.finalize.SaveIndexingDiagnostics(ctx, event.RepoID, event.SnapshotID, "commit_delta", "commit", event.CommitSHA, jobID, event.EventID, diagnostics); err != nil {
			return fmt.Errorf("persist indexing diagnostics: %w", err)
		}
	}
	if s.events == nil || len(diagnostics) == 0 {
		return nil
	}

	diagEvent := contractevents.NewParserDiagnosticsReadyEvent(
		diagnosticsEventID(jobID),
		event.RepoID,
		event.SnapshotID,
		event.CommitSHA,
		jobID,
		diagnostics,
	)
	if err := validate.ParserDiagnosticsReady(diagEvent); err != nil {
		return fmt.Errorf("validate parser.diagnostics.ready: %w", err)
	}
	return s.events.PublishParserDiagnostics(ctx, diagEvent)
}

// CommitJobIDFromEvent returns a deterministic job ID for a commits.changed event.
func CommitJobIDFromEvent(event contractevents.CommitsChangedEvent) string {
	return fmt.Sprintf("commit_job_%s_%s", event.SnapshotID, hashKey(event.CommitSHA, event.ParentSHA))
}

func graphDeltaEventID(snapshotID, baseSHA, targetSHA string) string {
	return fmt.Sprintf("evt_%s_delta_%s", snapshotID, hashKey(baseSHA, targetSHA))
}

func commitJobFromEvent(event contractevents.CommitsChangedEvent, jobID, status string) ParseJob {
	return ParseJob{
		JobID:      jobID,
		RepoID:     event.RepoID,
		SnapshotID: event.SnapshotID,
		CommitSHA:  event.CommitSHA,
		Status:     status,
	}
}

func gitCommitMetadataFromEvent(event contractevents.CommitsChangedEvent) metadata.GitCommitMetadata {
	return event.GitCommitMetadata
}

func (s *Service) markCommitFanInProgress(ctx context.Context, repoID, snapshotID, commitSHA, parentSHA string) error {
	if s.commitOrder == nil {
		return nil
	}
	if err := s.commitOrder.MarkCommitCompleted(ctx, repoID, snapshotID, commitSHA, parentSHA); err != nil {
		return fmt.Errorf("mark commit completed: %w", err)
	}
	if err := s.commitOrder.SyncCommitFanIn(ctx, repoID, snapshotID); err != nil {
		return fmt.Errorf("sync commit fan-in: %w", err)
	}
	return nil
}
