//go:build integration

package integration_test

import (
	"context"
	"crypto/rand"
	"crypto/sha256"
	"encoding/hex"
	"encoding/json"
	"fmt"
	"io"
	"log/slog"
	"os"
	"testing"
	"time"

	"go.mongodb.org/mongo-driver/bson"
	"go.mongodb.org/mongo-driver/mongo"
	"go.mongodb.org/mongo-driver/mongo/options"

	"bit.admedia.com/scm/ad/adpilot-indexing-repo-sync.com/internal/config"
	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/contracts/events"
	"bit.admedia.com/scm/ad/adpilot-indexing-repo-sync.com/internal/domain"
	redisclient "bit.admedia.com/scm/ad/adpilot-indexing-repo-sync.com/internal/redis"
	"bit.admedia.com/scm/ad/adpilot-indexing-repo-sync.com/internal/service"
	mongostore "bit.admedia.com/scm/ad/adpilot-indexing-repo-sync.com/internal/store/mongo"
)

type mockGitProvider struct {
	result gitcontract.SyncResult
}

func (m *mockGitProvider) SyncRepository(_ context.Context, repoID, ref string) (gitcontract.SyncResult, error) {
	result := m.result
	result.Snapshot.RepoID = repoID
	result.Snapshot.Ref = ref
	return result, nil
}

func (m *mockGitProvider) SyncRepositoryWithAuth(_ context.Context, opts gitcontract.SyncAuthOptions) (gitcontract.SyncResult, error) {
	return m.SyncRepository(context.Background(), opts.RepoID, opts.Ref)
}

func testMongoStore(t *testing.T) *mongostore.Store {
	t.Helper()

	uri := os.Getenv("MONGO_URI")
	if uri == "" {
		uri = "mongodb://localhost:27018/adpilot_repo_sync"
	}
	database := os.Getenv("MONGODB_DATABASE")
	if database == "" {
		database = "adpilot_repo_sync"
	}

	store, err := mongostore.New(config.MongoConfig{URI: uri, Database: database})
	if err != nil {
		t.Skipf("mongodb unavailable: %v", err)
	}

	ctx, cancel := context.WithTimeout(context.Background(), 3*time.Second)
	defer cancel()
	if !store.Ready(ctx) {
		_ = store.Close(context.Background())
		t.Skip("mongodb not reachable")
	}
	if err := store.EnsureIndexes(ctx); err != nil {
		t.Fatalf("ensure indexes: %v", err)
	}

	t.Cleanup(func() {
		_ = store.Close(context.Background())
	})
	return store
}

func testRedisClient(t *testing.T, streamPrefix string) (*redisclient.Client, config.RedisConfig) {
	t.Helper()

	url := os.Getenv("REDIS_URL")
	if url == "" {
		url = "redis://localhost:6379"
	}
	cfg := config.RedisConfig{URL: url, StreamPrefix: streamPrefix}

	client, err := redisclient.New(cfg)
	if err != nil {
		t.Skipf("redis unavailable: %v", err)
	}

	ctx, cancel := context.WithTimeout(context.Background(), 3*time.Second)
	defer cancel()
	if !client.Ready(ctx) {
		_ = client.Close()
		t.Skip("redis not reachable")
	}

	t.Cleanup(func() {
		_ = client.Close()
	})
	return client, cfg
}

func randomSuffix(t *testing.T) string {
	t.Helper()
	var b [4]byte
	if _, err := rand.Read(b[:]); err != nil {
		t.Fatalf("random suffix: %v", err)
	}
	return hex.EncodeToString(b[:])
}

func cleanupTestRepo(t *testing.T, repoID string) {
	t.Helper()
	ctx := context.Background()

	uri := os.Getenv("MONGO_URI")
	if uri == "" {
		uri = "mongodb://localhost:27018/adpilot_repo_sync"
	}
	database := os.Getenv("MONGODB_DATABASE")
	if database == "" {
		database = "adpilot_repo_sync"
	}

	client, err := mongo.Connect(ctx, options.Client().ApplyURI(uri))
	if err != nil {
		t.Fatalf("cleanup connect: %v", err)
	}
	defer func() {
		_ = client.Disconnect(ctx)
	}()

	filter := bson.M{"repo_id": repoID}
	db := client.Database(database)
	for _, collection := range []string{"repositories", "snapshots", "sync_jobs", "indexing_runs", "snapshot_files"} {
		if _, err := db.Collection(collection).DeleteMany(ctx, filter); err != nil {
			t.Fatalf("delete %s for %q: %v", collection, repoID, err)
		}
	}
}

func expectedCommitEventID(snapshotID, commitSHA string) string {
	sum := sha256.Sum256([]byte(commitSHA))
	return fmt.Sprintf("evt_%s_commit_%s", snapshotID, hex.EncodeToString(sum[:4]))
}

func expectedFileEventID(snapshotID string, change gitcontract.FileChange) string {
	key := change.Path
	if change.OldPath != "" {
		key = change.OldPath + "->" + change.Path
	}
	sum := sha256.Sum256([]byte(key))
	return fmt.Sprintf("evt_%s_file_%s", snapshotID, hex.EncodeToString(sum[:4]))
}

func streamLen(t *testing.T, client *redisclient.Client, stream string) int64 {
	t.Helper()
	count, err := client.Underlying().XLen(context.Background(), stream).Result()
	if err != nil {
		t.Fatalf("xlen %s: %v", stream, err)
	}
	return count
}

func latestStreamPayloads(t *testing.T, client *redisclient.Client, stream string, count int64) []string {
	t.Helper()
	msgs, err := client.Underlying().XRevRangeN(context.Background(), stream, "+", "-", count).Result()
	if err != nil {
		t.Fatalf("xrevrange %s: %v", stream, err)
	}
	if int64(len(msgs)) != count {
		t.Fatalf("expected %d messages on %s, got %d", count, stream, len(msgs))
	}

	payloads := make([]string, len(msgs))
	for i, msg := range msgs {
		payload, ok := msg.Values["payload"].(string)
		if !ok {
			t.Fatalf("expected string payload on %s", stream)
		}
		payloads[len(msgs)-1-i] = payload
	}
	return payloads
}

func e2eSyncResult(repoID string) gitcontract.SyncResult {
	return gitcontract.SyncResult{
		CloneURL: "https://bit.admedia.com/scm/ad/integration-repo.git",
		Snapshot: gitcontract.RepositorySnapshot{
			RepoID:     repoID,
			SnapshotID: "snap_tw23e2eintegration",
			CommitSHA:  "deadbeefcafebabe",
			Ref:        "master",
			Status:     gitcontract.SnapshotStatusReady,
		},
		FileChanges: []gitcontract.FileChange{
			{
				Path:       "main.go",
				ChangeType: events.ChangeTypeAdded,
				Language:   "go",
				FileKind:   events.FileKindCode,
			},
			{
				Path:       "docs/readme.md",
				ChangeType: events.ChangeTypeAdded,
				Language:   "markdown",
				FileKind:   events.FileKindDocs,
			},
		},
		CommitChanges: []gitcontract.CommitChange{
			{CommitSHA: "aaa111"},
			{CommitSHA: "bbb222", ParentSHA: "aaa111"},
		},
		IsFirstSync: true,
	}
}

func TestSyncFlowE2E(t *testing.T) {
	suffix := randomSuffix(t)
	repoID := fmt.Sprintf("ad/tw23-e2e-%s", suffix)
	streamPrefix := fmt.Sprintf("tw23e2e.%s.", suffix)

	mongoStore := testMongoStore(t)
	t.Cleanup(func() { cleanupTestRepo(t, repoID) })

	redisClient, redisCfg := testRedisClient(t, streamPrefix)
	logger := slog.New(slog.NewTextHandler(io.Discard, nil))
	publisher := redisclient.NewPublisher(redisClient, redisCfg, logger)

	git := &mockGitProvider{result: e2eSyncResult(repoID)}
	svc := service.NewSyncService(git, publisher, mongoStore, logger)

	readyStream := redisCfg.StreamName(events.StreamRepoSnapshotReady)
	filesStream := redisCfg.StreamName(events.StreamFilesChanged)
	commitsStream := redisCfg.StreamName(events.StreamCommitsChanged)
	readyBefore := streamLen(t, redisClient, readyStream)
	filesBefore := streamLen(t, redisClient, filesStream)
	commitsBefore := streamLen(t, redisClient, commitsStream)

	ctx := context.Background()
	if err := svc.SyncRepository(ctx, repoID, "master"); err != nil {
		t.Fatalf("SyncRepository: %v", err)
	}

	repo, ok, err := mongoStore.GetRepository(ctx, repoID)
	if err != nil {
		t.Fatalf("GetRepository: %v", err)
	}
	if !ok {
		t.Fatal("expected repository in mongo")
	}
	if repo.LastCommitSHA != "deadbeefcafebabe" {
		t.Fatalf("last commit: got %q", repo.LastCommitSHA)
	}
	if repo.DefaultRef != "master" {
		t.Fatalf("default ref: got %q", repo.DefaultRef)
	}

	snap, ok, err := mongoStore.GetSnapshot(ctx, repoID, "snap_tw23e2eintegration")
	if err != nil {
		t.Fatalf("GetSnapshot: %v", err)
	}
	if !ok {
		t.Fatal("expected snapshot in mongo")
	}
	if snap.FileCount != 2 {
		t.Fatalf("file count: got %d, want 2", snap.FileCount)
	}

	if got := streamLen(t, redisClient, readyStream); got != readyBefore+1 {
		t.Fatalf("ready stream length: got %d, want %d", got, readyBefore+1)
	}
	if got := streamLen(t, redisClient, filesStream); got != filesBefore+2 {
		t.Fatalf("files stream length: got %d, want %d", got, filesBefore+2)
	}
	if got := streamLen(t, redisClient, commitsStream); got != commitsBefore+2 {
		t.Fatalf("commits stream length: got %d, want %d", got, commitsBefore+2)
	}

	run, ok, err := mongoStore.GetIndexingRunBySnapshot(ctx, repoID, "snap_tw23e2eintegration")
	if err != nil {
		t.Fatalf("GetIndexingRunBySnapshot: %v", err)
	}
	if !ok {
		t.Fatal("expected indexing run in mongo")
	}
	if run.ExpectedCodeFiles != 1 || run.ExpectedDocsFiles != 1 || run.ExpectedCommits != 2 {
		t.Fatalf("unexpected expected counts: %+v", run)
	}
	if run.Status != domain.IndexingRunStatusCompleted {
		t.Fatalf("indexing run status: got %q", run.Status)
	}

	files, err := mongoStore.ListSnapshotFiles(ctx, repoID, snap.SnapshotID)
	if err != nil {
		t.Fatalf("ListSnapshotFiles: %v", err)
	}
	if len(files) != 2 {
		t.Fatalf("snapshot files: got %d, want 2", len(files))
	}

	readyPayload := latestStreamPayloads(t, redisClient, readyStream, 1)[0]
	var readyEvent events.RepoSnapshotReadyEvent
	if err := json.Unmarshal([]byte(readyPayload), &readyEvent); err != nil {
		t.Fatalf("unmarshal ready event: %v", err)
	}
	wantReadyID := fmt.Sprintf("evt_%s_ready", snap.SnapshotID)
	if readyEvent.EventID != wantReadyID {
		t.Fatalf("ready event_id: got %q, want %q", readyEvent.EventID, wantReadyID)
	}
	if readyEvent.RepoID != repoID || readyEvent.SnapshotID != snap.SnapshotID {
		t.Fatalf("unexpected ready payload: %+v", readyEvent)
	}

	filePayloads := latestStreamPayloads(t, redisClient, filesStream, 2)
	changes := e2eSyncResult(repoID).FileChanges
	wantFileIDs := []string{
		expectedFileEventID(snap.SnapshotID, changes[0]),
		expectedFileEventID(snap.SnapshotID, changes[1]),
	}

	for i, payload := range filePayloads {
		var fileEvent events.FileChangedEvent
		if err := json.Unmarshal([]byte(payload), &fileEvent); err != nil {
			t.Fatalf("unmarshal file event %d: %v", i, err)
		}
		if fileEvent.EventID != wantFileIDs[i] {
			t.Fatalf("file event_id[%d]: got %q, want %q", i, fileEvent.EventID, wantFileIDs[i])
		}
		if fileEvent.RepoID != repoID || fileEvent.SnapshotID != snap.SnapshotID {
			t.Fatalf("unexpected file payload[%d]: %+v", i, fileEvent)
		}
	}

	commitPayloads := latestStreamPayloads(t, redisClient, commitsStream, 2)
	commits := e2eSyncResult(repoID).CommitChanges
	wantCommitIDs := []string{
		expectedCommitEventID(snap.SnapshotID, commits[0].CommitSHA),
		expectedCommitEventID(snap.SnapshotID, commits[1].CommitSHA),
	}
	for i, payload := range commitPayloads {
		var commitEvent events.CommitsChangedEvent
		if err := json.Unmarshal([]byte(payload), &commitEvent); err != nil {
			t.Fatalf("unmarshal commit event %d: %v", i, err)
		}
		if commitEvent.EventID != wantCommitIDs[i] {
			t.Fatalf("commit event_id[%d]: got %q, want %q", i, commitEvent.EventID, wantCommitIDs[i])
		}
		if commitEvent.RepoID != repoID || commitEvent.SnapshotID != snap.SnapshotID {
			t.Fatalf("unexpected commit payload[%d]: %+v", i, commitEvent)
		}
	}

	if svc.Status() != service.StatusIdle {
		t.Fatalf("status: got %q, want idle", svc.Status())
	}

	gotRepo, ok := svc.GetRepository(repoID)
	if !ok {
		t.Fatal("expected GetRepository via service")
	}
	if gotRepo.RepoID != repoID {
		t.Fatalf("service repo: got %q", gotRepo.RepoID)
	}

	gotSnap, ok := svc.GetSnapshotMetadata(repoID, snap.SnapshotID)
	if !ok {
		t.Fatal("expected GetSnapshotMetadata via service")
	}
	if gotSnap.Status != domain.SnapshotStatusReady {
		t.Fatalf("snapshot status: got %q", gotSnap.Status)
	}
}
