//go:build integration

package redis_test

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

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

func TestPublish_ProducerStreamsIntegration(t *testing.T) {
	client, cfg := testRedisClient(t)
	logger := slog.New(slog.NewTextHandler(io.Discard, nil))
	publisher := redisclient.NewPublisher(client, cfg, logger)

	ctx := context.Background()

	chunksEvent := contractevents.NewChunksReadyEvent(
		"evt_integration_chunks",
		"ad/example",
		"snap_int_01",
		"abc123",
		"main.go",
		"go",
		[]contractevents.CodeChunk{
			{ChunkID: "c1", SymbolName: "main", SymbolType: contractevents.SymbolTypeFunction, StartLine: 1, EndLine: 10},
		},
	)
	if err := publisher.PublishChunksReady(ctx, chunksEvent); err != nil {
		t.Fatalf("publish chunks.ready: %v", err)
	}

	graphEvent := contractevents.NewGraphArtifactReadyEvent(
		"evt_integration_graph",
		"ad/example",
		"snap_int_01",
		"abc123",
		"art_int_01",
		"file:///data/artifacts/art_int_01.json",
		"v0.1.0",
		graph.SchemaVersionV1,
	)
	if err := publisher.PublishGraphArtifactReady(ctx, graphEvent); err != nil {
		t.Fatalf("publish graph.artifact.ready: %v", err)
	}

	diagEvent := contractevents.NewParserDiagnosticsReadyEvent(
		"evt_integration_diag",
		"ad/example",
		"snap_int_01",
		"abc123",
		"job_int_01",
		[]contractevents.ParserDiagnostic{
			{
				Code:     "PARSE_WARN",
				Message:  "unused import",
				Severity: contractevents.DiagnosticSeverityWarning,
				FilePath: "main.go",
				Line:     3,
			},
		},
	)
	if err := publisher.PublishParserDiagnostics(ctx, diagEvent); err != nil {
		t.Fatalf("publish parser.diagnostics.ready: %v", err)
	}

	assertStreamEventID(t, client, cfg.StreamName(contractevents.StreamChunksReady), chunksEvent.EventID)
	assertStreamEventID(t, client, cfg.StreamName(contractevents.StreamGraphArtifactReady), graphEvent.EventID)
	assertStreamEventID(t, client, cfg.StreamName(contractevents.StreamParserDiagnostics), diagEvent.EventID)
}

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

	msgs, err := client.Underlying().XRange(context.Background(), stream, "-", "+").Result()
	if err != nil {
		t.Fatalf("xrange %s: %v", stream, err)
	}
	if len(msgs) == 0 {
		t.Fatalf("expected at least 1 message on %s", stream)
	}

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

	var envelope struct {
		EventID string `json:"event_id"`
	}
	if err := json.Unmarshal([]byte(payload), &envelope); err != nil {
		t.Fatalf("unmarshal payload on %s: %v", stream, err)
	}
	if envelope.EventID != wantEventID {
		t.Fatalf("stream %s event_id: got %q, want %q", stream, envelope.EventID, wantEventID)
	}
}
