package redis

import (
	"context"
	"encoding/json"
	"fmt"
	"log/slog"
	"time"

	goredis "github.com/redis/go-redis/v9"

	"bit.admedia.com/scm/ad/adpilot-indexing-code-parser.com/internal/config"
	"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/validate"
	"bit.admedia.com/scm/ad/adpilot-indexing-code-parser.com/internal/service/parser"
)

// Publisher publishes validated parser output events to Redis Streams.
type Publisher struct {
	redis    *Client
	redisCfg config.RedisConfig
	logger   *slog.Logger
}

// NewPublisher creates a Redis Streams event publisher.
func NewPublisher(client *Client, redisCfg config.RedisConfig, logger *slog.Logger) *Publisher {
	return &Publisher{
		redis:    client,
		redisCfg: redisCfg,
		logger:   logger,
	}
}

var _ parser.EventPublisher = (*Publisher)(nil)

// PublishChunksReady publishes a chunks.ready event.
func (p *Publisher) PublishChunksReady(ctx context.Context, event events.ChunksReadyEvent) error {
	if err := validate.ChunksReady(event); err != nil {
		return fmt.Errorf("validate chunks.ready: %w", err)
	}
	stream := p.redisCfg.StreamName(events.StreamChunksReady)
	return p.publishJSON(ctx, stream, event.EventID, event)
}

// PublishGraphArtifactReady publishes a graph.artifact.ready event.
func (p *Publisher) PublishGraphArtifactReady(ctx context.Context, event events.GraphArtifactReadyEvent) error {
	if err := validate.GraphArtifactReady(event); err != nil {
		return fmt.Errorf("validate graph.artifact.ready: %w", err)
	}
	stream := p.redisCfg.StreamName(events.StreamGraphArtifactReady)
	return p.publishJSON(ctx, stream, event.EventID, event)
}

// PublishGraphDeltaReady publishes a graph.delta.ready event.
func (p *Publisher) PublishGraphDeltaReady(ctx context.Context, event events.GraphDeltaReadyEvent) error {
	if err := validate.GraphDeltaReady(event); err != nil {
		return fmt.Errorf("validate graph.delta.ready: %w", err)
	}
	stream := p.redisCfg.StreamName(events.StreamGraphDeltaReady)
	return p.publishJSON(ctx, stream, event.EventID, event)
}

// PublishParserDiagnostics publishes a parser.diagnostics.ready event.
func (p *Publisher) PublishParserDiagnostics(ctx context.Context, event events.ParserDiagnosticsReadyEvent) error {
	if err := validate.ParserDiagnosticsReady(event); err != nil {
		return fmt.Errorf("validate parser.diagnostics.ready: %w", err)
	}
	stream := p.redisCfg.StreamName(events.StreamParserDiagnostics)
	return p.publishJSON(ctx, stream, event.EventID, event)
}

func (p *Publisher) publishJSON(ctx context.Context, stream, eventID string, event any) error {
	payload, err := json.Marshal(event)
	if err != nil {
		return fmt.Errorf("marshal event: %w", err)
	}

	started := time.Now()
	entryID, err := p.redis.Underlying().XAdd(ctx, &goredis.XAddArgs{
		Stream: stream,
		Values: map[string]any{"payload": string(payload)},
	}).Result()
	if err != nil {
		p.logger.Error("failed to publish event",
			"stream", stream,
			"event_id", eventID,
			"error", err,
		)
		return fmt.Errorf("xadd stream %s: %w", stream, err)
	}

	p.logger.Info("published event",
		"stream", stream,
		"event_id", eventID,
		"entry_id", entryID,
		"publish_duration_ms", time.Since(started).Milliseconds(),
	)
	return nil
}
