feat!: consume privilege events with go-messaging-amqp (#327)
Unbound Release / Create Tag (push) Skipped
Unbound Release / Check Preconditions (push) Successful in 29s
authz_client / vulnerabilities (push) Skipped
authz_client / test (push) Skipped
pre-commit / pre-commit (push) Successful in 2m36s
Unbound Release / Create Release (push) Successful in 40s
Unbound Release / Generate Changelog and Handle PR (push) Successful in 37s
Release / release (push) Successful in 2m24s

This commit was merged in pull request #327.
This commit is contained in:
argoyle committed 2026-09-11 21:00:41 +00:00
1 parent f5e9eb52a7
commit 6bdf6e1cd3
5 files changed
+193 -68

No files matched your search

+25 -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,27 @@ 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)),
}
}
// 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 privilegeEvent](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 +130,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))
}
}