package service

import (
	"context"
	"errors"
	"fmt"
	"strings"
	"time"

	apicontract "bit.admedia.com/scm/ad/adpilot-indexing-repo-sync.com/internal/contracts/api"
	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/git"
	"bit.admedia.com/scm/ad/adpilot-indexing-repo-sync.com/internal/store"
)

// ErrInvalidImportPayload indicates the gateway import payload failed validation.
var ErrInvalidImportPayload = errors.New("invalid import payload")

// ErrDuplicateImport indicates the importJobId was already processed.
var ErrDuplicateImport = errors.New("duplicate import job")

// AcceptImport validates a gateway import request, persists the import job, and enqueues sync work.
func (s *SyncService) AcceptImport(ctx context.Context, req apicontract.ImportRequest) (apicontract.ImportAcceptedResponse, error) {
	if err := validateImportRequest(req); err != nil {
		return apicontract.ImportAcceptedResponse{}, fmt.Errorf("%w: %v", ErrInvalidImportPayload, err)
	}

	if existing, ok, err := s.metadata.GetImportJob(ctx, req.ImportJobID); err != nil {
		return apicontract.ImportAcceptedResponse{}, err
	} else if ok {
		return toImportResponse(existing), ErrDuplicateImport
	}

	now := time.Now().UTC()
	repoJobs := make([]domain.ImportRepoJob, 0, len(req.Repos))
	respJobs := make([]apicontract.ImportRepoJob, 0, len(req.Repos))

	for _, repo := range req.Repos {
		jobID, err := newJobID()
		if err != nil {
			return apicontract.ImportAcceptedResponse{}, err
		}
		repoID := git.BuildGatewayRepoID(req.Provider, repo.FullName)
		ref := strings.TrimSpace(repo.DefaultBranch)
		if ref == "" {
			ref = "main"
		}
		cloneURL := strings.TrimSpace(repo.CloneURL)
		if cloneURL == "" {
			return apicontract.ImportAcceptedResponse{}, fmt.Errorf("%w: cloneUrl required for %s", ErrInvalidImportPayload, repo.FullName)
		}

		repoJobs = append(repoJobs, domain.ImportRepoJob{
			JobID:          jobID,
			ProviderRepoID: repo.ProviderRepoID,
			FullName:       repo.FullName,
			RepoID:         repoID,
			Status:         domain.SyncStatusPending,
		})
		respJobs = append(respJobs, apicontract.ImportRepoJob{
			JobID:          jobID,
			ProviderRepoID: repo.ProviderRepoID,
			FullName:       repo.FullName,
			RepoID:         repoID,
			Status:         domain.SyncStatusPending,
		})
	}

	importJob := domain.ImportJob{
		ImportJobID:  req.ImportJobID,
		OrgID:        req.OrgID,
		ConnectionID: req.ConnectionID,
		Provider:     req.Provider,
		Status:       domain.SyncStatusPending,
		RepoJobs:     repoJobs,
		CreatedAt:    now,
	}

	if err := s.metadata.CreateImportJob(ctx, importJob); err != nil {
		if errors.Is(err, store.ErrDuplicateImportJob) {
			if existing, ok, getErr := s.metadata.GetImportJob(ctx, req.ImportJobID); getErr == nil && ok {
				return toImportResponse(existing), ErrDuplicateImport
			}
			return apicontract.ImportAcceptedResponse{}, ErrDuplicateImport
		}
		return apicontract.ImportAcceptedResponse{}, err
	}

	for _, rj := range repoJobs {
		if err := s.metadata.UpsertOrgRepositoryLink(ctx, domain.OrgRepositoryLink{
			OrgID:     req.OrgID,
			RepoID:    rj.RepoID,
			CreatedAt: now,
		}); err != nil {
			return apicontract.ImportAcceptedResponse{}, fmt.Errorf("upsert org repository link: %w", err)
		}
	}

	for i, repo := range req.Repos {
		ref := strings.TrimSpace(repo.DefaultBranch)
		if ref == "" {
			ref = "main"
		}
		cloneURL := strings.TrimSpace(repo.CloneURL)
		auth := cloneAuthFromRequest(req.CloneAuth)
		s.triggerImportSync(respJobs[i].RepoID, ref, cloneURL, auth, respJobs[i].JobID)
	}

	return apicontract.ImportAcceptedResponse{
		ImportJobID: req.ImportJobID,
		Status:      "accepted",
		Jobs:        respJobs,
	}, nil
}

func validateImportRequest(req apicontract.ImportRequest) error {
	if strings.TrimSpace(req.ImportJobID) == "" {
		return fmt.Errorf("importJobId is required")
	}
	if strings.TrimSpace(req.OrgID) == "" {
		return fmt.Errorf("orgId is required")
	}
	if strings.TrimSpace(req.ConnectionID) == "" {
		return fmt.Errorf("connectionId is required")
	}
	if strings.TrimSpace(req.Provider) == "" {
		return fmt.Errorf("provider is required")
	}
	if strings.TrimSpace(req.CloneAuth.Username) == "" || strings.TrimSpace(req.CloneAuth.Password) == "" {
		return fmt.Errorf("cloneAuth username and password are required")
	}
	if len(req.Repos) == 0 {
		return fmt.Errorf("repos must be a non-empty array")
	}
	for _, repo := range req.Repos {
		if strings.TrimSpace(repo.ProviderRepoID) == "" {
			return fmt.Errorf("providerRepoId is required")
		}
		if strings.TrimSpace(repo.FullName) == "" {
			return fmt.Errorf("fullName is required")
		}
		repoID := git.BuildGatewayRepoID(req.Provider, repo.FullName)
		if err := git.ValidateGatewayRepoID(repoID); err != nil {
			return err
		}
	}
	return nil
}

func cloneAuthFromRequest(auth apicontract.CloneAuth) gitcontract.SyncAuthOptions {
	return gitcontract.SyncAuthOptions{
		Username: auth.Username,
		Password: auth.Password,
	}
}

func toImportResponse(job domain.ImportJob) apicontract.ImportAcceptedResponse {
	jobs := make([]apicontract.ImportRepoJob, len(job.RepoJobs))
	for i, rj := range job.RepoJobs {
		jobs[i] = apicontract.ImportRepoJob{
			JobID:          rj.JobID,
			ProviderRepoID: rj.ProviderRepoID,
			FullName:       rj.FullName,
			RepoID:         rj.RepoID,
			Status:         rj.Status,
		}
	}
	return apicontract.ImportAcceptedResponse{
		ImportJobID: job.ImportJobID,
		Status:      job.Status,
		Jobs:        jobs,
	}
}

type importSyncParams struct {
	repoID   string
	ref      string
	cloneURL string
	auth     gitcontract.SyncAuthOptions
	jobID    string
}

func (s *SyncService) triggerImportSync(repoID, ref, cloneURL string, auth gitcontract.SyncAuthOptions, jobID string) {
	params := importSyncParams{
		repoID:   repoID,
		ref:      ref,
		cloneURL: cloneURL,
		auth: gitcontract.SyncAuthOptions{
			RepoID:   repoID,
			Ref:      ref,
			CloneURL: cloneURL,
			Username: auth.Username,
			Password: auth.Password,
		},
		jobID: jobID,
	}

	go func(p importSyncParams) {
		if err := s.runImportSync(context.Background(), p); err != nil {
			s.logger.Error("import sync failed",
				"repo_id", p.repoID,
				"ref", p.ref,
				"job_id", p.jobID,
				"error", err,
			)
		}
	}(params)
}

func (s *SyncService) runImportSync(ctx context.Context, p importSyncParams) (err error) {
	startedAt := time.Now().UTC()
	job := domain.SyncJob{
		JobID:     p.jobID,
		RepoID:    p.repoID,
		Ref:       p.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.runSyncWithAuth(ctx, p.repoID, p.ref, job, p.auth)
}

func (s *SyncService) runSyncWithAuth(ctx context.Context, repoID, ref string, job domain.SyncJob, auth gitcontract.SyncAuthOptions) (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("import sync failed",
			"repo_id", repoID,
			"ref", ref,
			"job_id", jobID,
			"duration_ms", time.Since(syncStarted).Milliseconds(),
			"error", cause,
		)
		return cause
	}

	auth.RepoID = repoID
	auth.Ref = ref
	result, err := s.git.SyncRepositoryWithAuth(ctx, auth)
	if err != nil {
		return failJob(err)
	}

	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("import sync completed",
		"repo_id", repoID,
		"ref", ref,
		"job_id", jobID,
		"snapshot_id", result.Snapshot.SnapshotID,
		"duration_ms", time.Since(syncStarted).Milliseconds(),
	)
	return nil
}
