package service

import (
	"context"
	"crypto/rand"
	"encoding/hex"
	"fmt"
	"log/slog"
	"strings"
	"sync"
	"time"

	"bit.admedia.com/scm/ad/adpilot-indexing-repo-sync.com/internal/client"
	"bit.admedia.com/scm/ad/adpilot-indexing-repo-sync.com/internal/config"
	"bit.admedia.com/scm/ad/adpilot-indexing-repo-sync.com/internal/contracts/events"
	gitcontract "bit.admedia.com/scm/ad/adpilot-indexing-repo-sync.com/internal/contracts/git"
	"bit.admedia.com/scm/ad/adpilot-indexing-repo-sync.com/internal/domain"
	"bit.admedia.com/scm/ad/adpilot-indexing-repo-sync.com/internal/purge"
	"bit.admedia.com/scm/ad/adpilot-indexing-repo-sync.com/internal/store"
)

const (
	StatusIdle    = "idle"
	StatusSyncing = "syncing"
	StatusError   = "error"

	// DefaultRef is the branch synced when no ref query param is supplied.
	DefaultRef = "master"
)

// EventPublisher defines Redis Stream publishing boundaries.
type EventPublisher interface {
	PublishRepoSnapshotReady(ctx context.Context, event events.RepoSnapshotReadyEvent) error
	PublishFileChanged(ctx context.Context, event events.FileChangedEvent) error
	PublishCommitsChanged(ctx context.Context, event events.CommitsChangedEvent) error
}

// SyncService coordinates repository sync workflows.
type SyncService struct {
	git      gitcontract.GitProvider
	events   EventPublisher
	metadata store.MetadataStore
	logger   *slog.Logger

	forceSyncCfg config.Config
	purger       *purge.PipelinePurger
	embedding    *client.EmbeddingEngineClient

	mu          sync.RWMutex
	status      string
	lastError   string
	activeSyncs int
	syncWG      sync.WaitGroup
}

// NewSyncService creates a sync service with injected dependencies.
func NewSyncService(
	git gitcontract.GitProvider,
	events EventPublisher,
	metadata store.MetadataStore,
	logger *slog.Logger,
) *SyncService {
	if logger == nil {
		logger = slog.Default()
	}
	return &SyncService{
		git:      git,
		events:   events,
		metadata: metadata,
		logger:   logger,
		status:   StatusIdle,
	}
}

// SyncRepository runs the full sync workflow synchronously: git sync, persist metadata, publish events.
func (s *SyncService) SyncRepository(ctx context.Context, repoID, ref string) (err error) {
	jobID, err := newJobID()
	if err != nil {
		s.recordError(err)
		return err
	}

	startedAt := time.Now().UTC()
	job := domain.SyncJob{
		JobID:     jobID,
		RepoID:    repoID,
		Ref:       ref,
		Status:    domain.SyncStatusRunning,
		StartedAt: startedAt,
	}
	if err := s.metadata.CreateSyncJob(ctx, job); err != nil {
		s.recordError(err)
		return fmt.Errorf("create sync job: %w", err)
	}

	return s.runSync(ctx, repoID, ref, job)
}

// TriggerSync accepts a repository sync for background processing and returns the job ID immediately.
func (s *SyncService) TriggerSync(repoID, ref string) (string, error) {
	repoID = strings.TrimSpace(repoID)
	ref = strings.TrimSpace(ref)
	if ref == "" {
		ref = DefaultRef
	}

	jobID, err := newJobID()
	if err != nil {
		return "", err
	}

	startedAt := time.Now().UTC()
	job := domain.SyncJob{
		JobID:     jobID,
		RepoID:    repoID,
		Ref:       ref,
		Status:    domain.SyncStatusRunning,
		StartedAt: startedAt,
	}
	ctx := context.Background()
	if err := s.metadata.CreateSyncJob(ctx, job); err != nil {
		return "", fmt.Errorf("create sync job: %w", err)
	}

	go func() {
		if err := s.runSync(context.Background(), repoID, ref, job); err != nil {
			s.logger.Error("background sync failed", "repo_id", repoID, "ref", ref, "job_id", jobID, "error", err)
		}
	}()

	return jobID, nil
}

// ActiveSyncCount returns the number of in-flight sync workflows.
func (s *SyncService) ActiveSyncCount() int {
	s.mu.RLock()
	defer s.mu.RUnlock()
	return s.activeSyncs
}

// WaitForSyncs blocks until all in-flight syncs finish or ctx is cancelled.
func (s *SyncService) WaitForSyncs(ctx context.Context) {
	if ctx == nil {
		ctx = context.Background()
	}
	done := make(chan struct{})
	go func() {
		s.syncWG.Wait()
		close(done)
	}()
	select {
	case <-done:
	case <-ctx.Done():
	}
}

func (s *SyncService) runSync(ctx context.Context, repoID, ref string, job domain.SyncJob) (err error) {
	s.syncWG.Add(1)
	defer s.syncWG.Done()

	syncStarted := time.Now()
	s.beginSync()
	defer func() { s.endSync(err == nil) }()

	jobID := job.JobID

	failJob := func(cause error) error {
		completedAt := time.Now().UTC()
		job.Status = domain.SyncStatusFailed
		job.Error = cause.Error()
		job.CompletedAt = &completedAt
		if updateErr := s.metadata.UpdateSyncJob(ctx, job); updateErr != nil {
			s.logger.Error("failed to update sync job", "job_id", jobID, "error", updateErr)
		}
		s.recordError(cause)
		s.logger.Error("sync failed",
			"repo_id", repoID,
			"ref", ref,
			"job_id", jobID,
			"duration_ms", time.Since(syncStarted).Milliseconds(),
			"error", cause,
		)
		return cause
	}

	result, err := s.git.SyncRepository(ctx, repoID, ref)
	if err != nil {
		s.logger.Error("git sync failed", "repo_id", repoID, "ref", ref, "error", err)
		return failJob(err)
	}

	if !result.IsFirstSync && len(result.FileChanges) == 0 && len(result.CommitChanges) == 0 {
		if err := s.persistRepositoryOnly(ctx, repoID, ref, result); err != nil {
			return failJob(fmt.Errorf("persist metadata: %w", err))
		}

		completedAt := time.Now().UTC()
		job.Status = domain.SyncStatusCompleted
		if prev, ok := s.GetLatestSnapshot(repoID); ok {
			job.SnapshotID = prev.SnapshotID
		}
		job.CompletedAt = &completedAt
		job.Error = ""
		if err := s.metadata.UpdateSyncJob(ctx, job); err != nil {
			s.recordError(err)
			return fmt.Errorf("update sync job: %w", err)
		}

		s.clearError()
		s.logger.Info("sync completed with no changes",
			"repo_id", repoID,
			"ref", ref,
			"job_id", jobID,
			"commit_sha", result.Snapshot.CommitSHA,
			"duration_ms", time.Since(syncStarted).Milliseconds(),
		)
		return nil
	}

	if err := s.persistMetadata(ctx, repoID, ref, result); err != nil {
		return failJob(fmt.Errorf("persist metadata: %w", err))
	}

	indexingRun, err := s.persistIndexingRun(ctx, job, result)
	if err != nil {
		return failJob(fmt.Errorf("persist indexing run: %w", err))
	}

	if err := s.publishSyncResult(ctx, result); err != nil {
		s.failIndexingRun(ctx, indexingRun, err)
		return failJob(fmt.Errorf("publish events: %w", err))
	}

	if err := s.completeIndexingRun(ctx, indexingRun); err != nil {
		return failJob(fmt.Errorf("complete indexing run: %w", err))
	}

	completedAt := time.Now().UTC()
	job.Status = domain.SyncStatusCompleted
	job.SnapshotID = result.Snapshot.SnapshotID
	job.CompletedAt = &completedAt
	job.Error = ""
	if err := s.metadata.UpdateSyncJob(ctx, job); err != nil {
		s.recordError(err)
		return fmt.Errorf("update sync job: %w", err)
	}

	s.clearError()
	s.logger.Info("sync completed",
		"repo_id", repoID,
		"ref", ref,
		"job_id", jobID,
		"indexing_run_id", indexingRun.IndexingRunID,
		"snapshot_id", result.Snapshot.SnapshotID,
		"commit_sha", result.Snapshot.CommitSHA,
		"file_changes", len(result.FileChanges),
		"commit_changes", len(result.CommitChanges),
		"expected_code_files", indexingRun.ExpectedCodeFiles,
		"expected_docs_files", indexingRun.ExpectedDocsFiles,
		"is_first_sync", result.IsFirstSync,
		"duration_ms", time.Since(syncStarted).Milliseconds(),
	)
	return nil
}

func (s *SyncService) persistRepositoryOnly(ctx context.Context, repoID, ref string, result gitcontract.SyncResult) error {
	now := time.Now().UTC()
	snap := result.Snapshot

	repo := domain.Repository{
		RepoID:        repoID,
		CloneURL:      result.CloneURL,
		DefaultRef:    ref,
		LastCommitSHA: snap.CommitSHA,
		LastSyncAt:    now,
		Status:        domain.RepositoryStatusActive,
		UpdatedAt:     now,
	}
	return s.metadata.UpsertRepository(ctx, repo)
}

func (s *SyncService) persistMetadata(ctx context.Context, repoID, ref string, result gitcontract.SyncResult) error {
	now := time.Now().UTC()
	snap := result.Snapshot

	repo := domain.Repository{
		RepoID:        repoID,
		CloneURL:      result.CloneURL,
		DefaultRef:    ref,
		LastCommitSHA: snap.CommitSHA,
		LastSyncAt:    now,
		Status:        domain.RepositoryStatusActive,
		UpdatedAt:     now,
	}
	if err := s.metadata.UpsertRepository(ctx, repo); err != nil {
		return err
	}

	domainSnap := domain.Snapshot{
		SnapshotID: snap.SnapshotID,
		RepoID:     snap.RepoID,
		CommitSHA:  snap.CommitSHA,
		Ref:        snap.Ref,
		FileCount:  len(result.FileChanges),
		Status:     snap.Status,
		CreatedAt:  now,
	}
	return s.metadata.SaveSnapshot(ctx, domainSnap)
}

// Status returns the current sync service status: idle, syncing, or error.
func (s *SyncService) Status() string {
	s.mu.RLock()
	defer s.mu.RUnlock()
	return s.status
}

// LastError returns the last sync error message when status is error.
func (s *SyncService) LastError() string {
	s.mu.RLock()
	defer s.mu.RUnlock()
	return s.lastError
}

// GetSnapshot returns a persisted snapshot by repo and snapshot ID.
func (s *SyncService) GetSnapshot(repoID, snapshotID string) (gitcontract.RepositorySnapshot, bool) {
	snap, ok, err := s.metadata.GetSnapshot(context.Background(), repoID, snapshotID)
	if err != nil || !ok {
		return gitcontract.RepositorySnapshot{}, false
	}
	return toGitSnapshot(snap), true
}

// GetLatestSnapshot returns the most recent snapshot for a repository.
func (s *SyncService) GetLatestSnapshot(repoID string) (gitcontract.RepositorySnapshot, bool) {
	snap, ok, err := s.metadata.GetLatestSnapshot(context.Background(), repoID)
	if err != nil || !ok {
		return gitcontract.RepositorySnapshot{}, false
	}
	return toGitSnapshot(snap), true
}

// GetRepository returns persisted repository metadata.
func (s *SyncService) GetRepository(repoID string) (domain.Repository, bool) {
	repo, ok, err := s.metadata.GetRepository(context.Background(), repoID)
	if err != nil || !ok {
		return domain.Repository{}, false
	}
	return repo, true
}

// GetSnapshotMetadata returns persisted snapshot metadata including timestamps and file count.
func (s *SyncService) GetSnapshotMetadata(repoID, snapshotID string) (domain.Snapshot, bool) {
	snap, ok, err := s.metadata.GetSnapshot(context.Background(), repoID, snapshotID)
	if err != nil || !ok {
		return domain.Snapshot{}, false
	}
	return snap, true
}

// GetLatestIndexingRun returns the most recent indexing run for a repository.
func (s *SyncService) GetLatestIndexingRun(repoID string) (domain.IndexingRun, bool) {
	run, ok, err := s.metadata.GetLatestIndexingRunByRepo(context.Background(), repoID)
	if err != nil || !ok {
		return domain.IndexingRun{}, false
	}
	return run, true
}

// GetIndexingRunBySnapshot returns the indexing run for a repo and snapshot.
func (s *SyncService) GetIndexingRunBySnapshot(repoID, snapshotID string) (domain.IndexingRun, bool) {
	run, ok, err := s.metadata.GetIndexingRunBySnapshot(context.Background(), repoID, snapshotID)
	if err != nil || !ok {
		return domain.IndexingRun{}, false
	}
	return run, true
}

// ListIndexingDiagnostics returns recent diagnostics for a repo snapshot.
func (s *SyncService) ListIndexingDiagnostics(repoID, snapshotID string, limit int) ([]domain.IndexingDiagnostic, bool) {
	diagnostics, err := s.metadata.ListIndexingDiagnostics(context.Background(), repoID, snapshotID, limit)
	if err != nil {
		return nil, false
	}
	return diagnostics, true
}

func toGitSnapshot(snap domain.Snapshot) gitcontract.RepositorySnapshot {
	return gitcontract.RepositorySnapshot{
		RepoID:     snap.RepoID,
		SnapshotID: snap.SnapshotID,
		CommitSHA:  snap.CommitSHA,
		Ref:        snap.Ref,
		Status:     snap.Status,
	}
}

func (s *SyncService) beginSync() {
	s.mu.Lock()
	defer s.mu.Unlock()
	s.activeSyncs++
	s.status = StatusSyncing
}

func (s *SyncService) endSync(success bool) {
	s.mu.Lock()
	defer s.mu.Unlock()
	s.activeSyncs--
	if s.activeSyncs > 0 {
		s.status = StatusSyncing
		return
	}
	if success && s.lastError == "" {
		s.status = StatusIdle
		return
	}
	if !success {
		s.status = StatusError
	}
}

func (s *SyncService) recordError(err error) {
	s.mu.Lock()
	defer s.mu.Unlock()
	s.lastError = err.Error()
	s.status = StatusError
}

func (s *SyncService) clearError() {
	s.mu.Lock()
	defer s.mu.Unlock()
	s.lastError = ""
}

func newJobID() (string, error) {
	var b [8]byte
	if _, err := rand.Read(b[:]); err != nil {
		return "", fmt.Errorf("generate job id: %w", err)
	}
	return "job_" + hex.EncodeToString(b[:]), nil
}
