package eventbus

import (
	"sync"
	"testing"
	"time"

	"github.com/rs/zerolog"
)

func TestPublishSubscribe(t *testing.T) {
	b := New(zerolog.Nop())
	s := b.Subscribe("t")
	defer s.Close()

	if got := b.Publish(Event{Topic: "t", Payload: "hi"}); got != 1 {
		t.Fatalf("Publish delivered %d, want 1", got)
	}
	select {
	case ev := <-s.Chan():
		if ev.Payload != "hi" {
			t.Errorf("payload: %v", ev.Payload)
		}
	case <-time.After(time.Second):
		t.Fatal("did not receive")
	}
}

func TestMultipleSubscribersOneTopic(t *testing.T) {
	b := New(zerolog.Nop())
	s1 := b.Subscribe("t")
	s2 := b.Subscribe("t")
	defer s1.Close()
	defer s2.Close()

	if got := b.Publish(Event{Topic: "t"}); got != 2 {
		t.Fatalf("delivered %d, want 2", got)
	}
}

func TestDifferentTopicsIsolated(t *testing.T) {
	b := New(zerolog.Nop())
	a := b.Subscribe("a")
	x := b.Subscribe("x")
	defer a.Close()
	defer x.Close()

	if got := b.Publish(Event{Topic: "a", Payload: 1}); got != 1 {
		t.Errorf("delivered %d to 'a', want 1", got)
	}
	select {
	case <-x.Chan():
		t.Error("'x' subscriber received an 'a' event")
	case <-time.After(50 * time.Millisecond):
	}
}

func TestSlowSubscriberDropped(t *testing.T) {
	b := New(zerolog.Nop())
	b.BufferSize = 2
	s := b.Subscribe("t")
	defer s.Close()
	// Don't drain the channel; publish 5 events. The 3rd, 4th, 5th drop.
	delivered := 0
	for i := 0; i < 5; i++ {
		delivered += b.Publish(Event{Topic: "t", Payload: i})
	}
	if delivered != 2 {
		t.Errorf("delivered %d total, want 2", delivered)
	}
}

func TestUnsubscribeStopsDelivery(t *testing.T) {
	b := New(zerolog.Nop())
	s := b.Subscribe("t")
	s.Close()
	if got := b.Publish(Event{Topic: "t"}); got != 0 {
		t.Errorf("delivered %d after close, want 0", got)
	}
	if got := b.SubscriberCount("t"); got != 0 {
		t.Errorf("SubscriberCount: %d, want 0", got)
	}
	// Double-close is safe.
	s.Close()
}

func TestConcurrentSubscribeAndPublish(t *testing.T) {
	b := New(zerolog.Nop())
	var wg sync.WaitGroup
	for i := 0; i < 8; i++ {
		wg.Add(1)
		go func() {
			defer wg.Done()
			s := b.Subscribe("t")
			defer s.Close()
			for j := 0; j < 100; j++ {
				select {
				case <-s.Chan():
				case <-time.After(10 * time.Millisecond):
				}
			}
		}()
	}
	for i := 0; i < 500; i++ {
		b.Publish(Event{Topic: "t", Payload: i})
	}
	wg.Wait()
	// No assertion; the test passes if it doesn't deadlock or race-detect.
}
