//go:build integration

package redis_test

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

	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 integrationHandler struct {
	mu    sync.Mutex
	calls int
}

func (h *integrationHandler) ProcessFileChanged(_ context.Context, _ contractevents.FileChangedEvent) error {
	h.mu.Lock()
	defer h.mu.Unlock()
	h.calls++
	return nil
}

func (h *integrationHandler) count() int {
	h.mu.Lock()
	defer h.mu.Unlock()
	return h.calls
}

func TestConsumer_FilesChangedIntegration(t *testing.T) {
	client, cfg := testRedisClient(t)
	logger := slog.New(slog.NewTextHandler(io.Discard, nil))
	handler := &integrationHandler{}
	appCfg := config.Config{
		Redis: cfg,
		Stream: config.StreamConsumerConfig{Group: "code-parser-service"},
	}
	consumer := redisclient.NewConsumerWithHandler(client, appCfg, handler, logger)

	stream := cfg.StreamName(contractevents.StreamFilesChanged)

	ctx, cancel := context.WithTimeout(context.Background(), 5*time.Second)
	defer cancel()

	go func() { _ = consumer.ConsumeFilesChanged(ctx) }()
	time.Sleep(50 * time.Millisecond)

	event := contractevents.NewFileChangedEvent(
		"evt_integration_01",
		"ad/example",
		"snap_int_01",
		"abc123",
		"main.go",
		"go",
		contractevents.FileKindCode,
		contractevents.ChangeTypeAdded,
	)

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

	rdb := client.Underlying()
	if _, err := rdb.XAdd(context.Background(), &goredis.XAddArgs{
		Stream: stream,
		Values: map[string]any{"payload": string(payload)},
	}).Result(); err != nil {
		t.Fatalf("xadd: %v", err)
	}

	deadline := time.After(5 * time.Second)
	for handler.count() < 1 {
		select {
		case <-deadline:
			t.Fatal("timeout waiting for consumer")
		default:
			time.Sleep(50 * time.Millisecond)
		}
	}
}

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

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

	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
}
