193 lines
4.7 KiB
Go
193 lines
4.7 KiB
Go
package pubsub
|
|
|
|
import (
|
|
"context"
|
|
"sync"
|
|
"testing"
|
|
|
|
"github.com/stretchr/testify/assert"
|
|
)
|
|
|
|
// --- InMemoryPubSub Publish/Subscribe Tests ---
|
|
|
|
func TestInMemoryPubSub_New(t *testing.T) {
|
|
ps := NewInMemoryPubSub()
|
|
assert.NotNil(t, ps)
|
|
assert.NotNil(t, ps.subscribers)
|
|
}
|
|
|
|
func TestInMemoryPubSub_PublishToNoSubscribers(t *testing.T) {
|
|
ps := NewInMemoryPubSub()
|
|
ctx := context.Background()
|
|
|
|
event := Event{
|
|
Type: "message_created",
|
|
AccountID: 1,
|
|
Payload: map[string]interface{}{"content": "hello"},
|
|
}
|
|
|
|
err := ps.Publish(ctx, "topic1", event)
|
|
assert.NoError(t, err)
|
|
}
|
|
|
|
func TestInMemoryPubSub_SubscribeAndPublish(t *testing.T) {
|
|
ps := NewInMemoryPubSub()
|
|
ctx := context.Background()
|
|
|
|
var receivedEvent Event
|
|
err := ps.Subscribe(ctx, "conversation_updated", func(event Event) {
|
|
receivedEvent = event
|
|
})
|
|
assert.NoError(t, err)
|
|
|
|
event := Event{
|
|
Type: "conversation_updated",
|
|
AccountID: 5,
|
|
Payload: map[string]interface{}{"status": "resolved"},
|
|
}
|
|
|
|
err = ps.Publish(ctx, "conversation_updated", event)
|
|
assert.NoError(t, err)
|
|
|
|
assert.Equal(t, "conversation_updated", receivedEvent.Type)
|
|
assert.Equal(t, uint(5), receivedEvent.AccountID)
|
|
assert.Equal(t, "resolved", receivedEvent.Payload["status"])
|
|
}
|
|
|
|
func TestInMemoryPubSub_MultipleSubscribers(t *testing.T) {
|
|
ps := NewInMemoryPubSub()
|
|
ctx := context.Background()
|
|
|
|
var mu sync.Mutex
|
|
receivedCount := 0
|
|
|
|
handler := func(event Event) {
|
|
mu.Lock()
|
|
receivedCount++
|
|
mu.Unlock()
|
|
}
|
|
|
|
err := ps.Subscribe(ctx, "message_created", handler)
|
|
assert.NoError(t, err)
|
|
err = ps.Subscribe(ctx, "message_created", handler)
|
|
assert.NoError(t, err)
|
|
|
|
event := Event{
|
|
Type: "message_created",
|
|
AccountID: 1,
|
|
Payload: map[string]interface{}{"content": "test"},
|
|
}
|
|
|
|
err = ps.Publish(ctx, "message_created", event)
|
|
assert.NoError(t, err)
|
|
|
|
mu.Lock()
|
|
assert.Equal(t, 2, receivedCount)
|
|
mu.Unlock()
|
|
}
|
|
|
|
func TestInMemoryPubSub_DifferentTopics(t *testing.T) {
|
|
ps := NewInMemoryPubSub()
|
|
ctx := context.Background()
|
|
|
|
var topic1Events []Event
|
|
var topic2Events []Event
|
|
|
|
err := ps.Subscribe(ctx, "topic1", func(event Event) {
|
|
topic1Events = append(topic1Events, event)
|
|
})
|
|
assert.NoError(t, err)
|
|
|
|
err = ps.Subscribe(ctx, "topic2", func(event Event) {
|
|
topic2Events = append(topic2Events, event)
|
|
})
|
|
assert.NoError(t, err)
|
|
|
|
event1 := Event{Type: "type1", AccountID: 1, Payload: map[string]interface{}{"key": "val1"}}
|
|
event2 := Event{Type: "type2", AccountID: 2, Payload: map[string]interface{}{"key": "val2"}}
|
|
|
|
err = ps.Publish(ctx, "topic1", event1)
|
|
assert.NoError(t, err)
|
|
err = ps.Publish(ctx, "topic2", event2)
|
|
assert.NoError(t, err)
|
|
|
|
assert.Len(t, topic1Events, 1)
|
|
assert.Len(t, topic2Events, 1)
|
|
assert.Equal(t, "type1", topic1Events[0].Type)
|
|
assert.Equal(t, "type2", topic2Events[0].Type)
|
|
}
|
|
|
|
func TestInMemoryPubSub_Unsubscribe(t *testing.T) {
|
|
ps := NewInMemoryPubSub()
|
|
ctx := context.Background()
|
|
|
|
received := false
|
|
err := ps.Subscribe(ctx, "test_topic", func(event Event) {
|
|
received = true
|
|
})
|
|
assert.NoError(t, err)
|
|
|
|
err = ps.Unsubscribe(ctx, "test_topic")
|
|
assert.NoError(t, err)
|
|
|
|
event := Event{Type: "test_event", AccountID: 1}
|
|
err = ps.Publish(ctx, "test_topic", event)
|
|
assert.NoError(t, err)
|
|
|
|
assert.False(t, received, "handler should not receive event after unsubscribe")
|
|
}
|
|
|
|
// --- Event Structure Tests ---
|
|
|
|
func TestEvent_Fields(t *testing.T) {
|
|
event := Event{
|
|
Type: "message_created",
|
|
AccountID: 42,
|
|
Payload: map[string]interface{}{"content": "hello", "sender_type": "contact"},
|
|
}
|
|
assert.Equal(t, "message_created", event.Type)
|
|
assert.Equal(t, uint(42), event.AccountID)
|
|
assert.Equal(t, "hello", event.Payload["content"])
|
|
assert.Equal(t, "contact", event.Payload["sender_type"])
|
|
}
|
|
|
|
func TestEvent_EmptyPayload(t *testing.T) {
|
|
event := Event{
|
|
Type: "conversation_resolved",
|
|
AccountID: 1,
|
|
Payload: nil,
|
|
}
|
|
assert.Equal(t, "conversation_resolved", event.Type)
|
|
assert.Nil(t, event.Payload)
|
|
}
|
|
|
|
func TestInMemoryPubSub_SubscribeToSameTopicMultipleHandlers(t *testing.T) {
|
|
ps := NewInMemoryPubSub()
|
|
ctx := context.Background()
|
|
|
|
results := make([]string, 0, 3)
|
|
|
|
err := ps.Subscribe(ctx, "alerts", func(event Event) {
|
|
results = append(results, "handler1")
|
|
})
|
|
assert.NoError(t, err)
|
|
|
|
err = ps.Subscribe(ctx, "alerts", func(event Event) {
|
|
results = append(results, "handler2")
|
|
})
|
|
assert.NoError(t, err)
|
|
|
|
err = ps.Subscribe(ctx, "alerts", func(event Event) {
|
|
results = append(results, "handler3")
|
|
})
|
|
assert.NoError(t, err)
|
|
|
|
event := Event{Type: "alert", AccountID: 1}
|
|
err = ps.Publish(ctx, "alerts", event)
|
|
assert.NoError(t, err)
|
|
|
|
assert.Len(t, results, 3)
|
|
assert.Contains(t, results, "handler1")
|
|
assert.Contains(t, results, "handler2")
|
|
assert.Contains(t, results, "handler3")
|
|
} |