package stream

import (
	"context"
	"log"
	"sync"
)

const (
	GenerateChannel = "generate-channel"
	RadioChannel    = "radio-channel"

	ReplyChannelPrefix = "reply-channel-"
)

type StreamWorker interface {
	Start(ctx context.Context) error
}

type StreamServer struct {
	workers []StreamWorker
}

func NewStreamServer(workers []StreamWorker) *StreamServer {
	return &StreamServer{
		workers: workers,
	}
}

func (s *StreamServer) Start(ctx context.Context) {
	var wg sync.WaitGroup

	for _, worker := range s.workers {
		wg.Add(1)
		go func(w StreamWorker) {
			defer wg.Done()
			if err := w.Start(ctx); err != nil {
				log.Printf("Worker error: %v", err)
			}
		}(worker)
	}

	// Wait for context cancellation
	<-ctx.Done()
	log.Println("Shutting down workers...")
	wg.Wait()
	log.Println("All workers stopped")
}
