package commitintel_test

import (
	"context"
	"encoding/json"
	"io"
	"log/slog"
	"os"
	"path/filepath"
	"testing"
	"time"

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

	"bit.admedia.com/scm/ad/adpilot-indexing-commit-intel.com/internal/config"
	contractevents "bit.admedia.com/scm/ad/adpilot-indexing-commit-intel.com/internal/contracts/events"
	redisclient "bit.admedia.com/scm/ad/adpilot-indexing-commit-intel.com/internal/redis"
	"bit.admedia.com/scm/ad/adpilot-indexing-commit-intel.com/internal/service/commitintel"
	"bit.admedia.com/scm/ad/adpilot-indexing-commit-intel.com/internal/store"
)

func TestPipeline_ProcessDeltaAndPublish(t *testing.T) {
	dir := t.TempDir()
	fixture, err := os.ReadFile(filepath.Join("..", "..", "store", "testdata", "delta_sample.json"))
	if err != nil {
		t.Fatalf("read fixture: %v", err)
	}
	deltaPath := filepath.Join(dir, "delta_target456_a1b2.json")
	if err := os.WriteFile(deltaPath, fixture, 0o644); err != nil {
		t.Fatalf("write fixture: %v", err)
	}

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

	redisCfg := config.RedisConfig{URL: "redis://" + server.Addr(), ConsumerGroup: "commit-intelligence-service"}
	client, err := redisclient.New(redisCfg)
	if err != nil {
		t.Fatalf("new redis client: %v", err)
	}
	defer client.Close()

	reader, err := store.NewFilesystemReader(dir)
	if err != nil {
		t.Fatalf("reader: %v", err)
	}

	jobs := store.NewAnalysisStore()
	logger := slog.New(slog.NewTextHandler(io.Discard, nil))
	publisher := redisclient.NewPublisher(client, redisCfg, logger)
	svc := commitintel.NewService(publisher, reader, commitintel.NewDeltaImpactAnalyzer(), jobs)

	event := validGraphDeltaEvent(t, "file://"+deltaPath)
	if err := svc.ProcessGraphDeltaReady(context.Background(), event); err != nil {
		t.Fatalf("ProcessGraphDeltaReady: %v", err)
	}

	analysisID := commitintel.AnalysisIDFromEvent(event)
	result, ok := jobs.GetResult(analysisID)
	if !ok {
		t.Fatal("expected stored result")
	}
	if result.Summary == "" {
		t.Fatal("expected summary")
	}

	stream := redisCfg.StreamName(contractevents.StreamCommitAnalysisReady)
	msgs, err := client.Underlying().XRange(context.Background(), stream, "-", "+").Result()
	if err != nil {
		t.Fatalf("xrange: %v", err)
	}
	if len(msgs) != 1 {
		t.Fatalf("expected 1 published event, got %d", len(msgs))
	}

	payload, ok := msgs[0].Values["payload"].(string)
	if !ok {
		t.Fatal("expected string payload")
	}
	var published contractevents.CommitAnalysisReadyEvent
	if err := json.Unmarshal([]byte(payload), &published); err != nil {
		t.Fatalf("unmarshal published: %v", err)
	}
	if published.AnalysisID != analysisID {
		t.Fatalf("published analysis_id: got %q want %q", published.AnalysisID, analysisID)
	}
}

func TestPipeline_RedisConsumerRoundTrip(t *testing.T) {
	dir := t.TempDir()
	fixture, err := os.ReadFile(filepath.Join("..", "..", "store", "testdata", "delta_sample.json"))
	if err != nil {
		t.Fatalf("read fixture: %v", err)
	}
	deltaPath := filepath.Join(dir, "delta_target456_a1b2.json")
	if err := os.WriteFile(deltaPath, fixture, 0o644); err != nil {
		t.Fatalf("write fixture: %v", err)
	}

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

	redisCfg := config.RedisConfig{URL: "redis://" + server.Addr(), ConsumerGroup: "commit-intelligence-service"}
	client, err := redisclient.New(redisCfg)
	if err != nil {
		t.Fatalf("new redis client: %v", err)
	}
	defer client.Close()

	reader, err := store.NewFilesystemReader(dir)
	if err != nil {
		t.Fatalf("reader: %v", err)
	}

	jobs := store.NewAnalysisStore()
	logger := slog.New(slog.NewTextHandler(io.Discard, nil))
	publisher := redisclient.NewPublisher(client, redisCfg, logger)
	svc := commitintel.NewService(publisher, reader, commitintel.NewDeltaImpactAnalyzer(), jobs)
	appCfg := config.Config{
		Redis:    redisCfg,
		Stream:   config.StreamConsumerConfig{Group: "commit-intelligence-service", MaxRetries: 5, BlockMs: 50},
		Analysis: config.AnalysisConfig{Workers: 1},
	}
	consumer := redisclient.NewConsumer(client, appCfg, svc, logger)

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

	event := validGraphDeltaEvent(t, "file://"+deltaPath)
	payload, err := json.Marshal(event)
	if err != nil {
		t.Fatalf("marshal: %v", err)
	}
	stream := redisCfg.StreamName(contractevents.StreamGraphDeltaReady)
	if _, err := server.XAdd(stream, "*", []string{"payload", string(payload)}); err != nil {
		t.Fatalf("xadd: %v", err)
	}

	deadline := time.After(3 * time.Second)
	for {
		if _, ok := jobs.GetResult(commitintel.AnalysisIDFromEvent(event)); ok {
			break
		}
		select {
		case <-deadline:
			t.Fatal("timeout waiting for consumer to process event")
		default:
			time.Sleep(10 * time.Millisecond)
		}
	}

	rdb := goredis.NewClient(&goredis.Options{Addr: server.Addr()})
	defer rdb.Close()
	msgs, err := rdb.XRange(context.Background(), redisCfg.StreamName(contractevents.StreamCommitAnalysisReady), "-", "+").Result()
	if err != nil {
		t.Fatalf("xrange output: %v", err)
	}
	if len(msgs) != 1 {
		t.Fatalf("expected 1 output message, got %d", len(msgs))
	}
}
