refactor: migrate eventsourced to codeberg and go-messaging-amqp
schemas / check (push) Skipped
schemas / vulnerabilities (push) Skipped
schemas / check-release (push) Skipped
schemas / build (push) Skipped
pre-commit / pre-commit (push) Skipped
schemas / vulnerabilities (pull_request) Successful in 58s
schemas / check-release (pull_request) Successful in 59s
schemas / check (pull_request) Successful in 1m58s
pre-commit / pre-commit (pull_request) Successful in 2m54s
schemas / build (pull_request) Successful in 6m9s
schemas / deploy-prod (pull_request) Skipped
schemas / check (push) Skipped
schemas / vulnerabilities (push) Skipped
schemas / check-release (push) Skipped
schemas / build (push) Skipped
pre-commit / pre-commit (push) Skipped
schemas / vulnerabilities (pull_request) Successful in 58s
schemas / check-release (pull_request) Successful in 59s
schemas / check (pull_request) Successful in 1m58s
pre-commit / pre-commit (pull_request) Successful in 2m54s
schemas / build (pull_request) Successful in 6m9s
schemas / deploy-prod (pull_request) Skipped
Move eventsourced modules from gitlab.com/unboundsoftware to codeberg.org/eventsourced. The Codeberg amqp publisher is built on go-messaging-amqp, so the service switches from sparetimecoders/goamqp as well: - routing keys are mapped on the eventsourced publisher (newEventPublisher) - cache consumers are typed handlers; Cache.Update drops the goamqp shape - publishes detach from the request context (uncancelledPublisher) - closeEvents is buffered and the receiver also watches rootCtx - WithPrefetchLimit dropped (no-op, 20 is the default) - amqp091-go v1.15.0 (GO-2026-6372) - renovate groups codeberg.org/messaging with eventsourced Only transient consumer queues are used, so no queue cutover is needed. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01Dpb7Y47dPz246oY9bCzh7k
This commit is contained in:
18 files changed
+221
-73
No files matched your search
@@ -0,0 +1,85 @@
|
||||
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")
|
||||
}
|
||||
+61
-32
@@ -13,6 +13,11 @@ import (
|
||||
"syscall"
|
||||
"time"
|
||||
|
||||
"codeberg.org/eventsourced/amqp"
|
||||
"codeberg.org/eventsourced/eventsourced"
|
||||
"codeberg.org/eventsourced/pg/v2"
|
||||
goamqp "codeberg.org/messaging/go-messaging-amqp"
|
||||
spec "codeberg.org/messaging/messaging"
|
||||
"github.com/99designs/gqlgen/graphql/handler"
|
||||
"github.com/99designs/gqlgen/graphql/handler/extension"
|
||||
"github.com/99designs/gqlgen/graphql/handler/lru"
|
||||
@@ -20,11 +25,7 @@ import (
|
||||
"github.com/99designs/gqlgen/graphql/playground"
|
||||
"github.com/alecthomas/kong"
|
||||
"github.com/rs/cors"
|
||||
"github.com/sparetimecoders/goamqp"
|
||||
"github.com/vektah/gqlparser/v2/ast"
|
||||
"gitlab.com/unboundsoftware/eventsourced/amqp"
|
||||
"gitlab.com/unboundsoftware/eventsourced/eventsourced"
|
||||
"gitlab.com/unboundsoftware/eventsourced/pg/v2"
|
||||
|
||||
"gitea.unbound.se/unboundsoftware/schemas/cache"
|
||||
"gitea.unbound.se/unboundsoftware/schemas/domain"
|
||||
@@ -57,7 +58,8 @@ func main() {
|
||||
var cli CLI
|
||||
_ = kong.Parse(&cli)
|
||||
logger := logging.SetupLogger(cli.LogLevel, cli.LogFormat, serviceName, buildVersion)
|
||||
closeEvents := make(chan error)
|
||||
// buffered: go-messaging-amqp reports a lost connection without blocking
|
||||
closeEvents := make(chan error, 1)
|
||||
|
||||
if err := start(
|
||||
closeEvents,
|
||||
@@ -114,7 +116,7 @@ func start(closeEvents chan error, logger *slog.Logger, connectToAmqpFunc func(u
|
||||
}
|
||||
|
||||
publisher := goamqp.NewPublisher()
|
||||
eventPublisher, err := amqp.New(publisher)
|
||||
eventPublisher, err := newEventPublisher(uncancelledPublisher{publisher})
|
||||
if err != nil {
|
||||
return fmt.Errorf("failed to create event publisher: %v", err)
|
||||
}
|
||||
@@ -130,25 +132,7 @@ func start(closeEvents chan error, logger *slog.Logger, connectToAmqpFunc func(u
|
||||
if err := loadSubGraphs(rootCtx, eventStore, serviceCache); err != nil {
|
||||
return fmt.Errorf("caching subgraphs: %w", err)
|
||||
}
|
||||
setups := []goamqp.Setup{
|
||||
goamqp.UseLogger(func(s string) { logger.Error(s) }),
|
||||
goamqp.CloseListener(closeEvents),
|
||||
goamqp.WithPrefetchLimit(20),
|
||||
goamqp.EventStreamPublisher(publisher),
|
||||
goamqp.TransientEventStreamConsumer("SubGraph.Updated", serviceCache.Update, domain.SubGraphUpdated{}),
|
||||
goamqp.TransientEventStreamConsumer("Organization.Added", serviceCache.Update, domain.OrganizationAdded{}),
|
||||
goamqp.TransientEventStreamConsumer("Organization.UserAdded", serviceCache.Update, domain.UserAddedToOrganization{}),
|
||||
goamqp.TransientEventStreamConsumer("Organization.APIKeyAdded", serviceCache.Update, domain.APIKeyAdded{}),
|
||||
goamqp.TransientEventStreamConsumer("Organization.APIKeyRemoved", serviceCache.Update, domain.APIKeyRemoved{}),
|
||||
goamqp.TransientEventStreamConsumer("Organization.Removed", serviceCache.Update, domain.OrganizationRemoved{}),
|
||||
goamqp.WithTypeMapping("SubGraph.Updated", domain.SubGraphUpdated{}),
|
||||
goamqp.WithTypeMapping("Organization.Added", domain.OrganizationAdded{}),
|
||||
goamqp.WithTypeMapping("Organization.UserAdded", domain.UserAddedToOrganization{}),
|
||||
goamqp.WithTypeMapping("Organization.APIKeyAdded", domain.APIKeyAdded{}),
|
||||
goamqp.WithTypeMapping("Organization.APIKeyRemoved", domain.APIKeyRemoved{}),
|
||||
goamqp.WithTypeMapping("Organization.Removed", domain.OrganizationRemoved{}),
|
||||
}
|
||||
if err := conn.Start(rootCtx, setups...); err != nil {
|
||||
if err := conn.Start(rootCtx, amqpSetups(logger, closeEvents, publisher, serviceCache)...); err != nil {
|
||||
return fmt.Errorf("failed to setup AMQP: %v", err)
|
||||
}
|
||||
|
||||
@@ -181,10 +165,13 @@ func start(closeEvents chan error, logger *slog.Logger, connectToAmqpFunc func(u
|
||||
wg.Add(1)
|
||||
go func() {
|
||||
defer wg.Done()
|
||||
err := <-closeEvents
|
||||
if err != nil {
|
||||
logger.With("error", err).Error("received close from AMQP")
|
||||
rootCancel()
|
||||
select {
|
||||
case err := <-closeEvents:
|
||||
if err != nil {
|
||||
logger.With("error", err).Error("received close from AMQP")
|
||||
rootCancel()
|
||||
}
|
||||
case <-rootCtx.Done():
|
||||
}
|
||||
}()
|
||||
|
||||
@@ -200,7 +187,6 @@ func start(closeEvents chan error, logger *slog.Logger, connectToAmqpFunc func(u
|
||||
logger.With("error", err).Error("close http server")
|
||||
}
|
||||
close(sigint)
|
||||
close(closeEvents)
|
||||
}()
|
||||
wg.Add(1)
|
||||
go func() {
|
||||
@@ -300,7 +286,7 @@ func loadOrganizations(ctx context.Context, eventStore eventsourced.EventStore,
|
||||
if _, err := eventsourced.NewHandler(ctx, organization, eventStore); err != nil {
|
||||
return err
|
||||
}
|
||||
_, err := serviceCache.Update(organization, nil)
|
||||
err := serviceCache.Update(organization)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
@@ -318,7 +304,7 @@ func loadSubGraphs(ctx context.Context, eventStore eventsourced.EventStore, serv
|
||||
if _, err := eventsourced.NewHandler(ctx, subGraph, eventStore); err != nil {
|
||||
return err
|
||||
}
|
||||
_, err := serviceCache.Update(subGraph, nil)
|
||||
err := serviceCache.Update(subGraph)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
@@ -326,6 +312,49 @@ func loadSubGraphs(ctx context.Context, eventStore eventsourced.EventStore, serv
|
||||
return nil
|
||||
}
|
||||
|
||||
func newEventPublisher(p amqp.BackingPublisher) (*amqp.Amqp, error) {
|
||||
return amqp.New(
|
||||
p,
|
||||
amqp.WithTypeMapping("SubGraph.Updated", domain.SubGraphUpdated{}),
|
||||
amqp.WithTypeMapping("Organization.Added", domain.OrganizationAdded{}),
|
||||
amqp.WithTypeMapping("Organization.UserAdded", domain.UserAddedToOrganization{}),
|
||||
amqp.WithTypeMapping("Organization.APIKeyAdded", domain.APIKeyAdded{}),
|
||||
amqp.WithTypeMapping("Organization.APIKeyRemoved", domain.APIKeyRemoved{}),
|
||||
amqp.WithTypeMapping("Organization.Removed", domain.OrganizationRemoved{}),
|
||||
)
|
||||
}
|
||||
|
||||
func amqpSetups(logger *slog.Logger, closeEvents chan error, publisher *goamqp.Publisher, serviceCache *cache.Cache) []goamqp.Setup {
|
||||
return []goamqp.Setup{
|
||||
goamqp.WithLogger(logger),
|
||||
goamqp.CloseListener(closeEvents),
|
||||
goamqp.EventStreamPublisher(publisher),
|
||||
goamqp.TransientEventStreamConsumer("SubGraph.Updated", cacheUpdater[domain.SubGraphUpdated](serviceCache)),
|
||||
goamqp.TransientEventStreamConsumer("Organization.Added", cacheUpdater[domain.OrganizationAdded](serviceCache)),
|
||||
goamqp.TransientEventStreamConsumer("Organization.UserAdded", cacheUpdater[domain.UserAddedToOrganization](serviceCache)),
|
||||
goamqp.TransientEventStreamConsumer("Organization.APIKeyAdded", cacheUpdater[domain.APIKeyAdded](serviceCache)),
|
||||
goamqp.TransientEventStreamConsumer("Organization.APIKeyRemoved", cacheUpdater[domain.APIKeyRemoved](serviceCache)),
|
||||
goamqp.TransientEventStreamConsumer("Organization.Removed", cacheUpdater[domain.OrganizationRemoved](serviceCache)),
|
||||
}
|
||||
}
|
||||
|
||||
// uncancelledPublisher detaches publishes from the request context: eventsourced publishes
|
||||
// right after storing an event, and a stored event must be published even if the client left.
|
||||
type uncancelledPublisher struct {
|
||||
amqp.BackingPublisher
|
||||
}
|
||||
|
||||
func (p uncancelledPublisher) Publish(ctx context.Context, routingKey string, msg any, headers ...goamqp.Header) error {
|
||||
return p.BackingPublisher.Publish(context.WithoutCancel(ctx), routingKey, msg, headers...)
|
||||
}
|
||||
|
||||
// cacheUpdater passes the payload as a pointer, which is what Cache.Update switches on
|
||||
func cacheUpdater[T any](c *cache.Cache) spec.EventHandler[T] {
|
||||
return func(_ context.Context, e spec.ConsumableEvent[T]) error {
|
||||
return c.Update(&e.Payload)
|
||||
}
|
||||
}
|
||||
|
||||
func ConnectAMQP(url string) (Connection, error) {
|
||||
return goamqp.NewFromURL(serviceName, url)
|
||||
}
|
||||
|
||||
@@ -4,11 +4,11 @@ import (
|
||||
"context"
|
||||
"testing"
|
||||
|
||||
"codeberg.org/eventsourced/eventsourced"
|
||||
"github.com/99designs/gqlgen/graphql/handler/transport"
|
||||
"github.com/google/uuid"
|
||||
"github.com/stretchr/testify/assert"
|
||||
"github.com/stretchr/testify/require"
|
||||
"gitlab.com/unboundsoftware/eventsourced/eventsourced"
|
||||
|
||||
"gitea.unbound.se/unboundsoftware/schemas/domain"
|
||||
"gitea.unbound.se/unboundsoftware/schemas/hash"
|
||||
|
||||
Reference in new issue
Block a user