package redis_test

import (
	"context"
	"encoding/json"
	"errors"
	"io"
	"log/slog"
	"sync"
	"sync/atomic"
	"testing"
	"time"

	miniredis "github.com/alicebob/miniredis/v2"
	goredis "github.com/redis/go-redis/v9"

	"bit.admedia.com/scm/ad/adpilot-indexing-code-parser.com/internal/config"
	contractevents "bit.admedia.com/scm/ad/adpilot-indexing-code-parser.com/internal/contracts/events"
	redisclient "bit.admedia.com/scm/ad/adpilot-indexing-code-parser.com/internal/redis"
)

type mockFileHandler struct {
	mu    sync.Mutex
	calls []contractevents.FileChangedEvent
}

func (m *mockFileHandler) ProcessFileChanged(_ context.Context, event contractevents.FileChangedEvent) error {
	m.mu.Lock()
	defer m.mu.Unlock()
	m.calls = append(m.calls, event)
	return nil
}

func (m *mockFileHandler) callCount() int {
	m.mu.Lock()
	defer m.mu.Unlock()
	return len(m.calls)
}

type mockCommitsHandler struct {
	mu    sync.Mutex
	calls []contractevents.CommitsChangedEvent
}

func (m *mockCommitsHandler) ProcessCommitsChanged(_ context.Context, event contractevents.CommitsChangedEvent) error {
	m.mu.Lock()
	defer m.mu.Unlock()
	m.calls = append(m.calls, event)
	return nil
}

func (m *mockCommitsHandler) callCount() int {
	m.mu.Lock()
	defer m.mu.Unlock()
	return len(m.calls)
}

func setupMiniRedisConsumer(t *testing.T, handler redisclient.FileChangedProcessor) (*miniredis.Miniredis, *redisclient.Consumer, config.RedisConfig) {
	return setupMiniRedisConsumerWithCommits(t, handler, nil)
}

func setupMiniRedisConsumerWithCommits(t *testing.T, fileHandler redisclient.FileChangedProcessor, commitsHandler redisclient.CommitsChangedProcessor) (*miniredis.Miniredis, *redisclient.Consumer, config.RedisConfig) {
	t.Helper()

	server, err := miniredis.Run()
	if err != nil {
		t.Fatalf("start miniredis: %v", err)
	}

	redisCfg := config.RedisConfig{URL: "redis://" + server.Addr()}
	appCfg := config.Config{
		Redis: redisCfg,
		Stream: config.StreamConsumerConfig{
			Group:      "code-parser-service",
			MaxRetries: 5,
			BlockMs:    50,
		},
	}
	client, err := redisclient.New(redisCfg)
	if err != nil {
		server.Close()
		t.Fatalf("new client: %v", err)
	}

	logger := slog.New(slog.NewTextHandler(io.Discard, nil))
	consumer := redisclient.NewConsumerWithHandlers(client, appCfg, fileHandler, nil, nil, commitsHandler, logger)

	t.Cleanup(func() {
		_ = client.Close()
		server.Close()
	})

	return server, consumer, redisCfg
}

func xaddEvent(t *testing.T, server *miniredis.Miniredis, stream string, event any) {
	t.Helper()

	payload, err := json.Marshal(event)
	if err != nil {
		t.Fatalf("marshal event: %v", err)
	}

	if _, err := server.XAdd(stream, "*", []string{"payload", string(payload)}); err != nil {
		t.Fatalf("xadd: %v", err)
	}
}

func pendingCount(t *testing.T, server *miniredis.Miniredis, stream, group string) int {
	t.Helper()

	client := goredis.NewClient(&goredis.Options{Addr: server.Addr()})
	defer client.Close()

	pending, err := client.XPending(context.Background(), stream, group).Result()
	if err != nil {
		t.Fatalf("xpending: %v", err)
	}
	return int(pending.Count)
}

func startFilesConsumer(t *testing.T, consumer *redisclient.Consumer) (context.Context, context.CancelFunc) {
	t.Helper()

	ctx, cancel := context.WithCancel(context.Background())
	go func() { _ = consumer.ConsumeFilesChanged(ctx) }()
	time.Sleep(50 * time.Millisecond)
	return ctx, cancel
}

func startSnapshotConsumer(t *testing.T, consumer *redisclient.Consumer) (context.Context, context.CancelFunc) {
	t.Helper()

	ctx, cancel := context.WithCancel(context.Background())
	go func() { _ = consumer.ConsumeRepoSnapshotReady(ctx) }()
	time.Sleep(50 * time.Millisecond)
	return ctx, cancel
}

func startCommitsConsumer(t *testing.T, consumer *redisclient.Consumer) (context.Context, context.CancelFunc) {
	t.Helper()

	ctx, cancel := context.WithCancel(context.Background())
	go func() { _ = consumer.ConsumeCommitsChanged(ctx) }()
	time.Sleep(50 * time.Millisecond)
	return ctx, cancel
}

func waitForCommitCalls(t *testing.T, handler *mockCommitsHandler, want int) {
	t.Helper()

	deadline := time.After(3 * time.Second)
	for handler.callCount() < want {
		select {
		case <-deadline:
			t.Fatalf("timeout waiting for %d ProcessCommitsChanged calls, got %d", want, handler.callCount())
		default:
			time.Sleep(10 * time.Millisecond)
		}
	}
}

func waitForCalls(t *testing.T, handler *mockFileHandler, want int) {
	t.Helper()

	deadline := time.After(3 * time.Second)
	for handler.callCount() < want {
		select {
		case <-deadline:
			t.Fatalf("timeout waiting for %d ProcessFileChanged calls, got %d", want, handler.callCount())
		default:
			time.Sleep(10 * time.Millisecond)
		}
	}
}

func TestConsumer_FilesChangedDelegates(t *testing.T) {
	handler := &mockFileHandler{}
	server, consumer, cfg := setupMiniRedisConsumer(t, handler)

	ctx, cancel := startFilesConsumer(t, consumer)
	defer cancel()

	event := contractevents.NewFileChangedEvent(
		"evt_file_01",
		"ad/example",
		"snap_01",
		"abc123",
		"main.go",
		"go",
		contractevents.FileKindCode,
		contractevents.ChangeTypeAdded,
	)
	xaddEvent(t, server, cfg.StreamName(contractevents.StreamFilesChanged), event)

	waitForCalls(t, handler, 1)

	if handler.calls[0].FilePath != "main.go" {
		t.Fatalf("unexpected file path: %q", handler.calls[0].FilePath)
	}
	_ = ctx
}

func TestConsumer_FilesChangedSkipsNonCode(t *testing.T) {
	handler := &mockFileHandler{}
	server, consumer, cfg := setupMiniRedisConsumer(t, handler)

	_, cancel := startFilesConsumer(t, consumer)
	defer cancel()

	event := contractevents.NewFileChangedEvent(
		"evt_docs_01",
		"ad/example",
		"snap_01",
		"abc123",
		"README.md",
		"markdown",
		contractevents.FileKindDocs,
		contractevents.ChangeTypeAdded,
	)
	xaddEvent(t, server, cfg.StreamName(contractevents.StreamFilesChanged), event)

	time.Sleep(300 * time.Millisecond)

	if handler.callCount() != 0 {
		t.Fatalf("expected 0 calls, got %d", handler.callCount())
	}
	if pending := pendingCount(t, server, cfg.StreamName(contractevents.StreamFilesChanged), "code-parser-service"); pending != 0 {
		t.Fatalf("expected 0 pending messages, got %d", pending)
	}
}

func TestConsumer_FilesChangedIdempotent(t *testing.T) {
	handler := &mockFileHandler{}
	server, consumer, cfg := setupMiniRedisConsumer(t, handler)

	_, cancel := startFilesConsumer(t, consumer)
	defer cancel()

	stream := cfg.StreamName(contractevents.StreamFilesChanged)
	event := contractevents.NewFileChangedEvent(
		"evt_dup_01",
		"ad/example",
		"snap_01",
		"abc123",
		"main.go",
		"go",
		contractevents.FileKindCode,
		contractevents.ChangeTypeAdded,
	)
	xaddEvent(t, server, stream, event)
	xaddEvent(t, server, stream, event)

	waitForCalls(t, handler, 1)
	time.Sleep(300 * time.Millisecond)

	if handler.callCount() != 1 {
		t.Fatalf("expected 1 call for duplicate event_id, got %d", handler.callCount())
	}
}

func TestConsumer_MalformedPayloadAcked(t *testing.T) {
	handler := &mockFileHandler{}
	server, consumer, cfg := setupMiniRedisConsumer(t, handler)

	_, cancel := startFilesConsumer(t, consumer)
	defer cancel()

	stream := cfg.StreamName(contractevents.StreamFilesChanged)
	if _, err := server.XAdd(stream, "*", []string{"payload", "{not-json"}); err != nil {
		t.Fatalf("xadd: %v", err)
	}

	time.Sleep(300 * time.Millisecond)

	if handler.callCount() != 0 {
		t.Fatalf("expected 0 calls, got %d", handler.callCount())
	}
	if pending := pendingCount(t, server, stream, "code-parser-service"); pending != 0 {
		t.Fatalf("expected malformed message acked, pending=%d", pending)
	}
}

func TestConsumer_SnapshotAckOnly(t *testing.T) {
	handler := &mockFileHandler{}
	server, consumer, cfg := setupMiniRedisConsumer(t, handler)

	_, cancel := startSnapshotConsumer(t, consumer)
	defer cancel()

	event := contractevents.NewRepoSnapshotReadyEvent(
		"evt_snap_01",
		"ad/example",
		"snap_01",
		"abc123",
		"master",
	)
	event.FileCount = 42
	xaddEvent(t, server, cfg.StreamName(contractevents.StreamRepoSnapshotReady), event)

	time.Sleep(300 * time.Millisecond)

	if handler.callCount() != 0 {
		t.Fatalf("expected snapshot to not call ProcessFileChanged, got %d", handler.callCount())
	}
	if pending := pendingCount(t, server, cfg.StreamName(contractevents.StreamRepoSnapshotReady), "code-parser-service"); pending != 0 {
		t.Fatalf("expected snapshot message acked, pending=%d", pending)
	}
}

func TestConsumer_ContextCancel(t *testing.T) {
	handler := &mockFileHandler{}
	_, consumer, _ := setupMiniRedisConsumer(t, handler)

	ctx, cancel := context.WithCancel(context.Background())
	done := make(chan struct{})
	go func() {
		defer close(done)
		_ = consumer.ConsumeFilesChanged(ctx)
	}()

	cancel()

	select {
	case <-done:
	case <-time.After(3 * time.Second):
		t.Fatal("consumer did not stop after context cancel")
	}
}

func TestConsumer_RunStopsOnCancel(t *testing.T) {
	handler := &mockFileHandler{}
	_, consumer, _ := setupMiniRedisConsumer(t, handler)

	ctx, cancel := context.WithCancel(context.Background())
	var running atomic.Bool
	running.Store(true)

	done := make(chan struct{})
	go func() {
		defer close(done)
		_ = consumer.Run(ctx)
		running.Store(false)
	}()

	cancel()

	select {
	case <-done:
	case <-time.After(3 * time.Second):
		t.Fatal("Run did not stop after context cancel")
	}
	if running.Load() {
		t.Fatal("expected Run to finish")
	}
}

type failingFileHandler struct {
	mockFileHandler
}

func (f *failingFileHandler) ProcessFileChanged(context.Context, contractevents.FileChangedEvent) error {
	return errors.New("transient parse failure")
}

func TestConsumer_TransientFailureLeavesPending(t *testing.T) {
	handler := &failingFileHandler{}
	server, consumer, cfg := setupMiniRedisConsumer(t, handler)

	_, cancel := startFilesConsumer(t, consumer)
	defer cancel()

	stream := cfg.StreamName(contractevents.StreamFilesChanged)
	event := contractevents.NewFileChangedEvent(
		"evt_fail_01",
		"ad/example",
		"snap_01",
		"abc123",
		"main.go",
		"go",
		contractevents.FileKindCode,
		contractevents.ChangeTypeAdded,
	)
	xaddEvent(t, server, stream, event)

	time.Sleep(500 * time.Millisecond)

	if pending := pendingCount(t, server, stream, "code-parser-service"); pending == 0 {
		t.Fatal("expected transient failure to leave message pending")
	}
}

func TestConsumer_CommitsChanged(t *testing.T) {
	fileHandler := &mockFileHandler{}
	commitsHandler := &mockCommitsHandler{}
	server, consumer, cfg := setupMiniRedisConsumerWithCommits(t, fileHandler, commitsHandler)

	ctx, cancel := startCommitsConsumer(t, consumer)
	defer cancel()

	event := contractevents.NewCommitsChangedEvent("evt_commit_1", "ad/example", "snap_1", "abc123")
	event.ParentSHA = "def456"
	event.ChangedFiles = []string{"main.go"}

	xaddEvent(t, server, cfg.StreamName(contractevents.StreamCommitsChanged), event)
	waitForCommitCalls(t, commitsHandler, 1)

	if commitsHandler.callCount() != 1 {
		t.Fatalf("expected 1 commit handler call, got %d", commitsHandler.callCount())
	}
	if fileHandler.callCount() != 0 {
		t.Fatalf("expected files handler untouched, got %d", fileHandler.callCount())
	}
	if pending := pendingCount(t, server, cfg.StreamName(contractevents.StreamCommitsChanged), "code-parser-service"); pending != 0 {
		t.Fatalf("expected commits.changed message acked, pending=%d", pending)
	}

	_ = ctx
}
