package scrape

import (
	"context"
	"errors"
	"math/rand/v2"
	"sync"
	"time"

	"github.com/rs/zerolog"

	"github.com/operator/command-center/internal/db"
	"github.com/operator/command-center/internal/observability"
)

// Scheduler runs each configured tracker's scrape on its own cadence.
// Construct one with NewScheduler; call Start with a context whose cancel
// signals shutdown. Stop blocks until all worker goroutines have exited.
type Scheduler struct {
	writer   *db.SnapshotWriter
	events   *observability.EventRecorder
	logger   zerolog.Logger

	mu       sync.Mutex
	cancels  map[string]context.CancelFunc
	refreshCh map[string]chan struct{}
	wg       sync.WaitGroup
	rng      *rand.Rand
}

// NewScheduler constructs a stopped scheduler.
func NewScheduler(writer *db.SnapshotWriter, events *observability.EventRecorder, logger zerolog.Logger) *Scheduler {
	return &Scheduler{
		writer:    writer,
		events:    events,
		logger:    logger.With().Str("component", "scheduler").Logger(),
		cancels:   make(map[string]context.CancelFunc),
		refreshCh: make(map[string]chan struct{}),
		rng:       rand.New(rand.NewPCG(uint64(time.Now().UnixNano()), 0)),
	}
}

// Job is one tracker the scheduler should poll.
type Job struct {
	Adapter     Adapter
	TrackerID   string
	IntervalSec int
	JitterSec   int
}

// Replace stops any workers whose tracker is not in jobs, leaves workers for
// trackers whose interval/adapter is unchanged, and starts workers for new
// trackers. The caller invokes Replace on initial boot and on every
// trackers.yaml hot-reload.
func (s *Scheduler) Replace(parent context.Context, jobs []Job) {
	s.mu.Lock()
	defer s.mu.Unlock()

	wanted := make(map[string]Job, len(jobs))
	for _, j := range jobs {
		wanted[j.TrackerID] = j
	}

	// Stop removed.
	for id, cancel := range s.cancels {
		if _, keep := wanted[id]; !keep {
			cancel()
			delete(s.cancels, id)
			delete(s.refreshCh, id)
			s.logger.Info().Str("tracker", id).Msg("stopped scrape worker")
		}
	}

	// Start added (or replaced). For simplicity, every Replace re-creates
	// the worker; cheap goroutines, makes interval changes obvious.
	for id, job := range wanted {
		if _, already := s.cancels[id]; already {
			// Worker already running for this id; conservative for Phase 1:
			// keep it (don't churn on every reload).
			continue
		}
		ctx, cancel := context.WithCancel(parent)
		ch := make(chan struct{}, 1)
		s.cancels[id] = cancel
		s.refreshCh[id] = ch
		s.wg.Add(1)
		go s.runOne(ctx, job, ch)
	}
}

// Refresh signals an immediate scrape for trackerID. Non-blocking; if a
// refresh is already pending, this is a no-op.
func (s *Scheduler) Refresh(trackerID string) error {
	s.mu.Lock()
	ch, ok := s.refreshCh[trackerID]
	s.mu.Unlock()
	if !ok {
		return errors.New("scheduler: tracker not scheduled: " + trackerID)
	}
	select {
	case ch <- struct{}{}:
	default:
		// Refresh already queued; one is enough.
	}
	return nil
}

// Stop terminates all workers and blocks until they exit.
func (s *Scheduler) Stop() {
	s.mu.Lock()
	for id, cancel := range s.cancels {
		cancel()
		delete(s.cancels, id)
	}
	s.mu.Unlock()
	s.wg.Wait()
}

func (s *Scheduler) runOne(ctx context.Context, job Job, refresh <-chan struct{}) {
	defer s.wg.Done()
	logger := s.logger.With().Str("tracker", job.TrackerID).Logger()
	logger.Info().Int("interval_sec", job.IntervalSec).Msg("scrape worker started")

	// Stagger the very first scrape by a small random delay so all trackers
	// don't scrape at boot at the exact same instant. ~0-1s is enough; longer
	// just delays the first datapoint without meaningfully spreading load at
	// single-operator scale.
	initial := time.Duration(s.rng.IntN(1_000)) * time.Millisecond
	select {
	case <-time.After(initial):
	case <-ctx.Done():
		return
	}

	for {
		s.scrapeOnce(ctx, job, logger)

		wait := s.nextDelay(job)
		timer := time.NewTimer(wait)
		select {
		case <-timer.C:
		case <-refresh:
			if !timer.Stop() {
				<-timer.C
			}
			logger.Debug().Msg("scrape worker triggered by manual refresh")
		case <-ctx.Done():
			timer.Stop()
			logger.Info().Msg("scrape worker stopped")
			return
		}
	}
}

func (s *Scheduler) nextDelay(job Job) time.Duration {
	base := time.Duration(job.IntervalSec) * time.Second
	if job.JitterSec <= 0 {
		return base
	}
	jitterRange := time.Duration(job.JitterSec) * time.Second
	// Range is [-jitter, +jitter].
	offset := time.Duration(s.rng.Int64N(int64(2*jitterRange))) - jitterRange
	d := base + offset
	if d < time.Second {
		d = time.Second
	}
	return d
}

func (s *Scheduler) scrapeOnce(ctx context.Context, job Job, logger zerolog.Logger) {
	start := time.Now()
	snap, err := job.Adapter.FetchRatio(ctx)
	if err != nil {
		logger.Warn().Err(err).Msg("scrape failed")
		if s.events != nil {
			s.events.Record(ctx, observability.LevelError, "scrape", "fetch failed",
				map[string]any{"tracker": job.TrackerID, "error": err.Error()}, "")
		}
		return
	}

	row := db.RatioSnapshotRow{
		TrackerID:                job.TrackerID,
		Timestamp:                time.Now(),
		RealUploadedBytes:        snap.RealUploadedBytes,
		RealDownloadedBytes:      snap.RealDownloadedBytes,
		RealRatio:                snap.RealRatio,
		DisplayedUploadedBytes:   snap.DisplayedUploadedBytes,
		DisplayedDownloadedBytes: snap.DisplayedDownloadedBytes,
		DisplayedRatio:           snap.DisplayedRatio,
		BonusPoints:              snap.BonusPoints,
		UnsatCount:               snap.UnsatCount,
		UnsatLimit:               snap.UnsatLimit,
		ClassOrRank:              snap.ClassOrRank,
		RawJSON:                  snap.RawJSON,
	}
	if err := s.writer.WriteRatio(ctx, row); err != nil {
		logger.Error().Err(err).Msg("snapshot write failed")
		if s.events != nil {
			s.events.Record(ctx, observability.LevelError, "scrape", "write failed",
				map[string]any{"tracker": job.TrackerID, "error": err.Error()}, "")
		}
		return
	}
	logger.Info().
		Dur("duration", time.Since(start)).
		Interface("real_ratio", snap.RealRatio).
		Msg("scrape ok")
}
