feat!: consume privilege events with go-messaging-amqp #327

Merged
argoyle merged 2 commits from feat/go-messaging-amqp into main 2026-09-11 21:00:43 +00:00
3 changed files with 35 additions and 9 deletions
Showing only changes of commit fc0a4a0e31 - Show all commits

No files matched your search

+1 -1
View File
@@ -43,4 +43,4 @@ The `CompanyPrivileges` struct contains permission flags:
### Event Handling
Registers per-replica (transient) go-messaging-amqp consumers for privilege update events from the authz-service (`Setup()`), keeping the local privilege cache up-to-date.
Registers per-replica (transient) go-messaging-amqp consumers for privilege update events from the authz-service (`Setup()`), keeping the local privilege cache up-to-date. Don't combine `Setup()` with go-messaging-amqp's `WithReconnect`: a reconnect declares new per-replica queues, so revocations published during the outage are lost unless `Fetch()` runs again. Services exit on connection loss (`CloseListener`) and re-fetch on start.
+6 -1
View File
@@ -104,8 +104,13 @@ func (h *PrivilegeHandler) Setup() []goamqp.Setup {
}
}
// privilegeEvent is the set of events Process handles.
type privilegeEvent interface {
UserAdded | UserRemoved | PrivilegeAdded | PrivilegeRemoved
}
// process adapts Process to a typed go-messaging-amqp handler.
func process[T any](h *PrivilegeHandler) spec.EventHandler[T] {
func process[T privilegeEvent](h *PrivilegeHandler) spec.EventHandler[T] {
return func(_ context.Context, event spec.ConsumableEvent[T]) error {
return h.Process(&event.Payload)
}
+28 -7
View File
@@ -645,22 +645,43 @@ func TestPrivilegeHandler_Setup(t *testing.T) {
topology, err := goamqp.CollectTopology("some-service", handler.Setup()...)
assert.NoError(t, err)
var keys []string
wiring := map[string]string{}
for _, e := range topology.Endpoints {
assert.Equal(t, spec.DirectionConsume, e.Direction)
assert.True(t, e.Ephemeral, "%s must be a per-replica consumer", e.RoutingKey)
keys = append(keys, e.RoutingKey)
wiring[e.RoutingKey] = e.MessageType
}
assert.Equal(t, []string{"User.Added", "User.Removed", "Privilege.Added", "Privilege.Removed"}, keys)
// A key wired to the wrong type could turn a revocation into a grant.
assert.Equal(t, map[string]string{
"User.Added": "client.UserAdded",
"User.Removed": "client.UserRemoved",
"Privilege.Added": "client.PrivilegeAdded",
"Privilege.Removed": "client.PrivilegeRemoved",
}, wiring)
}
func TestPrivilegeHandler_process(t *testing.T) {
ctx := context.Background()
handler := New(WithBaseURL("base"))
admin := func(p CompanyPrivileges) bool { return p.Admin }
err := process[UserAdded](handler)(context.Background(), spec.ConsumableEvent[UserAdded]{
assert.NoError(t, process[UserAdded](handler)(ctx, spec.ConsumableEvent[UserAdded]{
Payload: UserAdded{Email: "jim@example.org", CompanyID: "abc-123"},
})
assert.NoError(t, err)
}))
assert.False(t, handler.IsAllowed("jim@example.org", "abc-123", admin))
assert.Equal(t, []string{"abc-123"}, handler.CompaniesByUser("jim@example.org", func(CompanyPrivileges) bool { return true }))
assert.NoError(t, process[PrivilegeAdded](handler)(ctx, spec.ConsumableEvent[PrivilegeAdded]{
Payload: PrivilegeAdded{Email: "jim@example.org", CompanyID: "abc-123", Privilege: PrivilegeAdmin},
}))
assert.True(t, handler.IsAllowed("jim@example.org", "abc-123", admin))
assert.NoError(t, process[PrivilegeRemoved](handler)(ctx, spec.ConsumableEvent[PrivilegeRemoved]{
Payload: PrivilegeRemoved{Email: "jim@example.org", CompanyID: "abc-123", Privilege: PrivilegeAdmin},
}))
assert.False(t, handler.IsAllowed("jim@example.org", "abc-123", admin))
assert.NoError(t, process[UserRemoved](handler)(ctx, spec.ConsumableEvent[UserRemoved]{
Payload: UserRemoved{Email: "jim@example.org", CompanyID: "abc-123"},
}))
assert.Empty(t, handler.CompaniesByUser("jim@example.org", func(CompanyPrivileges) bool { return true }))
}