feat!: consume privilege events with go-messaging-amqp
authz_client / test (push) Skipped
authz_client / vulnerabilities (push) Skipped
pre-commit / pre-commit (push) Skipped

eventsourced and sloth moved to Codeberg, where their AMQP modules use codeberg.org/messaging/go-messaging-amqp instead of goamqp; services migrating to them need authz_client on the same client. Setup() now returns go-messaging-amqp setups: the same four per-replica (transient) consumers, with typed handlers that feed Process.

BREAKING CHANGE: Setup returns go-messaging-amqp setups, and Process(msg any) error replaces Process(msg, goamqp.Headers) (any, error). Services stay on v0.5.x until they migrate (see docs/design/codeberg-migration.md).

Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01SV1epy2pKvBx3yDfhxATk2
This commit is contained in:
argoyleandClaude Opus 5 committed 2026-09-11 22:47:30 +02:00
1 parent f5e9eb52a7
commit 3750416d8d
5 files changed
+167 -68

No files matched your search

+20 -11
View File
@@ -1,6 +1,7 @@
package client
import (
"context"
"encoding/json"
"fmt"
"io"
@@ -8,7 +9,8 @@ import (
"reflect"
"sync"
"github.com/sparetimecoders/goamqp"
goamqp "codeberg.org/messaging/go-messaging-amqp"
spec "codeberg.org/messaging/messaging"
)
// CompanyPrivileges contains the privileges for a combination of email address and company id
@@ -95,15 +97,22 @@ func (h *PrivilegeHandler) Fetch() error {
func (h *PrivilegeHandler) Setup() []goamqp.Setup {
return []goamqp.Setup{
goamqp.TransientEventStreamConsumer("User.Added", h.Process, UserAdded{}),
goamqp.TransientEventStreamConsumer("User.Removed", h.Process, UserRemoved{}),
goamqp.TransientEventStreamConsumer("Privilege.Added", h.Process, PrivilegeAdded{}),
goamqp.TransientEventStreamConsumer("Privilege.Removed", h.Process, PrivilegeRemoved{}),
goamqp.TransientEventStreamConsumer("User.Added", process[UserAdded](h)),
goamqp.TransientEventStreamConsumer("User.Removed", process[UserRemoved](h)),
goamqp.TransientEventStreamConsumer("Privilege.Added", process[PrivilegeAdded](h)),
goamqp.TransientEventStreamConsumer("Privilege.Removed", process[PrivilegeRemoved](h)),
}
}
// process adapts Process to a typed go-messaging-amqp handler.
func process[T any](h *PrivilegeHandler) spec.EventHandler[T] {
return func(_ context.Context, event spec.ConsumableEvent[T]) error {
return h.Process(&event.Payload)
}
}
// Process privilege-related events and update the internal state
func (h *PrivilegeHandler) Process(msg interface{}, _ goamqp.Headers) (interface{}, error) {
func (h *PrivilegeHandler) Process(msg any) error {
h.Lock()
defer h.Unlock()
@@ -116,21 +125,21 @@ func (h *PrivilegeHandler) Process(msg interface{}, _ goamqp.Headers) (interface
ev.CompanyID: {},
}
}
return nil, nil
return nil
case *UserRemoved:
if priv, exists := h.privileges[ev.Email]; exists {
delete(priv, ev.CompanyID)
}
return nil, nil
return nil
case *PrivilegeAdded:
h.setPrivileges(ev.Email, ev.CompanyID, ev.Privilege, true)
return nil, nil
return nil
case *PrivilegeRemoved:
h.setPrivileges(ev.Email, ev.CompanyID, ev.Privilege, false)
return nil, nil
return nil
default:
fmt.Printf("Got unexpected message type (%s): '%+v'\n", reflect.TypeOf(msg).String(), msg)
return nil, fmt.Errorf("unexpected event type: '%s'", reflect.TypeOf(msg))
return fmt.Errorf("unexpected event type: '%s'", reflect.TypeOf(msg))
}
}