package main import ( "context" "log/slog" "testing" "codeberg.org/eventsourced/eventsourced" goamqp "codeberg.org/messaging/go-messaging-amqp" spec "codeberg.org/messaging/messaging" "github.com/stretchr/testify/assert" "github.com/stretchr/testify/require" "gitea.unbound.se/unboundsoftware/schemas/cache" "gitea.unbound.se/unboundsoftware/schemas/domain" ) // ownKeys is a literal list on purpose: deriving it from the code would keep a renamed key green. var ownKeys = []string{ "SubGraph.Updated", "Organization.Added", "Organization.UserAdded", "Organization.APIKeyAdded", "Organization.APIKeyRemoved", "Organization.Removed", } func Test_amqpSetups(t *testing.T) { setups := amqpSetups(slog.Default(), make(chan error, 1), goamqp.NewPublisher(), cache.New(slog.Default())) topology, err := goamqp.CollectTopology(serviceName, setups...) require.NoError(t, err) var consumed, publishesTo []string for _, e := range topology.Endpoints { switch e.Direction { case spec.DirectionConsume: assert.True(t, e.Ephemeral, "cache consumers must be transient, %s is durable", e.QueueName) consumed = append(consumed, e.RoutingKey) case spec.DirectionPublish: publishesTo = append(publishesTo, e.ExchangeName) } } assert.ElementsMatch(t, ownKeys, consumed) assert.Equal(t, []string{"events.topic.exchange"}, publishesTo) } type recordingPublisher struct { keys []string ctxs []context.Context } func (r *recordingPublisher) Publish(ctx context.Context, routingKey string, _ any, _ ...goamqp.Header) error { r.keys = append(r.keys, routingKey) r.ctxs = append(r.ctxs, ctx) return nil } func Test_newEventPublisher(t *testing.T) { rec := &recordingPublisher{} p, err := newEventPublisher(rec) require.NoError(t, err) for _, e := range []eventsourced.Event{ &domain.SubGraphUpdated{}, &domain.OrganizationAdded{}, &domain.UserAddedToOrganization{}, &domain.APIKeyAdded{}, &domain.APIKeyRemoved{}, &domain.OrganizationRemoved{}, } { require.NoError(t, p.Publish(context.Background(), e)) } assert.Equal(t, ownKeys, rec.keys, "published keys must match what the cache consumers bind") } func Test_uncancelledPublisher(t *testing.T) { ctx, cancel := context.WithCancel(context.Background()) cancel() rec := &recordingPublisher{} p, err := newEventPublisher(uncancelledPublisher{rec}) require.NoError(t, err) require.NoError(t, p.Publish(ctx, &domain.OrganizationAdded{})) assert.NoError(t, rec.ctxs[0].Err(), "publish must not see the request's cancellation") }