Unbound Release / Check Preconditions (push) Successful in 34s
Unbound Release / Create Release (push) Skipped
schemas / vulnerabilities (push) Successful in 1m9s
schemas / check-release (push) Successful in 57s
Unbound Release / Generate Changelog and Handle PR (push) Successful in 39s
Unbound Release / Create Tag (push) Successful in 33s
Release / release (push) Successful in 1m25s
schemas / check (push) Successful in 1m34s
pre-commit / pre-commit (push) Successful in 3m17s
schemas / build (push) Successful in 6m34s
schemas / deploy-prod (push) Successful in 1m11s
86 lines
2.5 KiB
Go
86 lines
2.5 KiB
Go
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")
|
|
}
|