//go:build integration

package redis_test

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

	"bit.admedia.com/scm/ad/adpilot-indexing-repo-sync.com/internal/config"
	"bit.admedia.com/scm/ad/adpilot-indexing-repo-sync.com/internal/contracts/events"
	redisclient "bit.admedia.com/scm/ad/adpilot-indexing-repo-sync.com/internal/redis"
)

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
}

func streamLatestPayload(t *testing.T, client *redisclient.Client, stream string) string {
	t.Helper()

	msgs, err := client.Underlying().XRevRangeN(context.Background(), stream, "+", "-", 1).Result()
	if err != nil {
		t.Fatalf("xrevrange %s: %v", stream, err)
	}
	if len(msgs) != 1 {
		t.Fatalf("expected 1 latest message on %s, got %d", stream, len(msgs))
	}

	payload, ok := msgs[0].Values["payload"].(string)
	if !ok {
		t.Fatalf("expected string payload on %s", stream)
	}
	return payload
}

func streamMessageCount(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 TestPublishRepoSnapshotReadyIntegration(t *testing.T) {
	client, cfg := testRedisClient(t)
	logger := slog.New(slog.NewTextHandler(io.Discard, nil))
	publisher := redisclient.NewPublisher(client, cfg, logger)

	stream := cfg.StreamName(events.StreamRepoSnapshotReady)
	before := streamMessageCount(t, client, stream)

	event := events.NewRepoSnapshotReadyEvent(
		"evt_tw23_snapshot_ready",
		"ad/integration-repo",
		"snap_tw23integration",
		"deadbeef",
		"master",
	)

	if err := publisher.PublishRepoSnapshotReady(context.Background(), event); err != nil {
		t.Fatalf("publish: %v", err)
	}

	after := streamMessageCount(t, client, stream)
	if after != before+1 {
		t.Fatalf("stream length: got %d, want %d", after, before+1)
	}

	payload := streamLatestPayload(t, client, stream)
	var decoded events.RepoSnapshotReadyEvent
	if err := json.Unmarshal([]byte(payload), &decoded); err != nil {
		t.Fatalf("unmarshal: %v", err)
	}
	if decoded.EventID != event.EventID {
		t.Fatalf("event_id: got %q, want %q", decoded.EventID, event.EventID)
	}
	if decoded.RepoID != "ad/integration-repo" || decoded.SnapshotID != "snap_tw23integration" {
		t.Fatalf("unexpected payload: %+v", decoded)
	}
}

func TestPublishFileChangedWithStreamPrefixIntegration(t *testing.T) {
	client, baseCfg := testRedisClient(t)
	cfg := config.RedisConfig{
		URL:          baseCfg.URL,
		StreamPrefix: "tw23test.",
	}
	logger := slog.New(slog.NewTextHandler(io.Discard, nil))
	publisher := redisclient.NewPublisher(client, cfg, logger)

	stream := cfg.StreamName(events.StreamFilesChanged)
	if stream != "tw23test.files.changed" {
		t.Fatalf("unexpected stream name: %s", stream)
	}

	event := events.NewFileChangedEvent(
		"evt_tw23_file",
		"ad/integration-repo",
		"snap_tw23integration",
		"deadbeef",
		"main.go",
		"go",
		events.FileKindCode,
		events.ChangeTypeAdded,
	)

	if err := publisher.PublishFileChanged(context.Background(), event); err != nil {
		t.Fatalf("publish: %v", err)
	}

	payload := streamLatestPayload(t, client, stream)
	var decoded events.FileChangedEvent
	if err := json.Unmarshal([]byte(payload), &decoded); err != nil {
		t.Fatalf("unmarshal: %v", err)
	}
	if decoded.EventID != event.EventID {
		t.Fatalf("event_id: got %q, want %q", decoded.EventID, event.EventID)
	}
}
