From 3750416d8d5a940e2d43d33aedcaedbaeec9261a Mon Sep 17 00:00:00 2001 From: Joakim Olsson Date: Fri, 11 Sep 2026 22:47:30 +0200 Subject: [PATCH 1/2] feat!: consume privilege events with go-messaging-amqp 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) Claude-Session: https://claude.ai/code/session_01SV1epy2pKvBx3yDfhxATk2 --- CLAUDE.md | 2 +- client.go | 31 +++++++++----- client_test.go | 110 ++++++++++++++++++++++++++++--------------------- go.mod | 22 ++++++++-- go.sum | 70 ++++++++++++++++++++++++++++--- 5 files changed, 167 insertions(+), 68 deletions(-) diff --git a/CLAUDE.md b/CLAUDE.md index c30a006..23d8116 100644 --- a/CLAUDE.md +++ b/CLAUDE.md @@ -43,4 +43,4 @@ The `CompanyPrivileges` struct contains permission flags: ### Event Handling -Implements `goamqp` message handlers to receive privilege update events from the authz-service, 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. diff --git a/client.go b/client.go index 7089683..b876dc7 100644 --- a/client.go +++ b/client.go @@ -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)) } } diff --git a/client_test.go b/client_test.go index f179563..ed926b3 100644 --- a/client_test.go +++ b/client_test.go @@ -1,6 +1,7 @@ package client import ( + "context" "fmt" "net/http" "net/http/httptest" @@ -8,28 +9,27 @@ import ( "sync" "testing" - "github.com/sparetimecoders/goamqp" + goamqp "codeberg.org/messaging/go-messaging-amqp" + spec "codeberg.org/messaging/messaging" "github.com/stretchr/testify/assert" ) func TestPrivilegeHandler_Process_InvalidType(t *testing.T) { handler := New(WithBaseURL("base")) - result, err := handler.Process("abc", goamqp.Headers{}) + err := handler.Process("abc") - assert.Nil(t, result) assert.EqualError(t, err, "unexpected event type: 'string'") } func TestPrivilegeHandler_Process_PrivilegeRemoved(t *testing.T) { handler := New(WithBaseURL("base")) - result, err := handler.Process(&PrivilegeAdded{ + err := handler.Process(&PrivilegeAdded{ Email: "jim@example.org", CompanyID: "abc-123", Privilege: PrivilegeAdmin, - }, goamqp.Headers{}) - assert.Nil(t, result) + }) assert.NoError(t, err) companies := handler.CompaniesByUser("jim@example.org", func(privileges CompanyPrivileges) bool { @@ -38,12 +38,11 @@ func TestPrivilegeHandler_Process_PrivilegeRemoved(t *testing.T) { assert.Equal(t, []string{"abc-123"}, companies) - result, err = handler.Process(&PrivilegeRemoved{ + err = handler.Process(&PrivilegeRemoved{ Email: "jim@example.org", CompanyID: "abc-123", Privilege: PrivilegeAdmin, - }, goamqp.Headers{}) - assert.Nil(t, result) + }) assert.NoError(t, err) companies = handler.CompaniesByUser("jim@example.org", func(privileges CompanyPrivileges) bool { @@ -56,18 +55,16 @@ func TestPrivilegeHandler_Process_PrivilegeRemoved(t *testing.T) { func TestPrivilegeHandler_Process_UserAdded_And_UserRemoved(t *testing.T) { handler := New(WithBaseURL("base")) - result, err := handler.Process(&UserAdded{ + err := handler.Process(&UserAdded{ Email: "jim@example.org", CompanyID: "abc-123", - }, goamqp.Headers{}) - assert.Nil(t, result) + }) assert.NoError(t, err) - result, err = handler.Process(&UserAdded{ + err = handler.Process(&UserAdded{ Email: "jim@example.org", CompanyID: "abc-456", - }, goamqp.Headers{}) - assert.Nil(t, result) + }) assert.NoError(t, err) companies := handler.CompaniesByUser("jim@example.org", func(privileges CompanyPrivileges) bool { @@ -76,18 +73,16 @@ func TestPrivilegeHandler_Process_UserAdded_And_UserRemoved(t *testing.T) { sort.Strings(companies) assert.Equal(t, []string{"abc-123", "abc-456"}, companies) - result, err = handler.Process(&UserRemoved{ + err = handler.Process(&UserRemoved{ Email: "jim@example.org", CompanyID: "abc-123", - }, goamqp.Headers{}) - assert.Nil(t, result) + }) assert.NoError(t, err) - result, err = handler.Process(&UserRemoved{ + err = handler.Process(&UserRemoved{ Email: "jim@example.org", CompanyID: "abc-456", - }, goamqp.Headers{}) - assert.Nil(t, result) + }) assert.NoError(t, err) companies = handler.CompaniesByUser("jim@example.org", func(privileges CompanyPrivileges) bool { @@ -109,11 +104,10 @@ func TestPrivilegeHandler_GetCompanies_Email_Not_Found(t *testing.T) { func TestPrivilegeHandler_GetCompanies_No_Companies_Found(t *testing.T) { handler := New(WithBaseURL("base")) - result, err := handler.Process(&UserAdded{ + err := handler.Process(&UserAdded{ Email: "jim@example.org", CompanyID: "abc-123", - }, goamqp.Headers{}) - assert.Nil(t, result) + }) assert.NoError(t, err) companies := handler.CompaniesByUser("jim@example.org", func(privileges CompanyPrivileges) bool { @@ -128,11 +122,10 @@ func TestPrivilegeHandler_GetCompanies_No_Companies_Found(t *testing.T) { assert.Equal(t, []string{"abc-123"}, companies) - result, err = handler.Process(&UserRemoved{ + err = handler.Process(&UserRemoved{ Email: "jim@example.org", CompanyID: "abc-123", - }, goamqp.Headers{}) - assert.Nil(t, result) + }) assert.NoError(t, err) companies = handler.CompaniesByUser("jim@example.org", func(privileges CompanyPrivileges) bool { @@ -144,12 +137,11 @@ func TestPrivilegeHandler_GetCompanies_No_Companies_Found(t *testing.T) { func TestPrivilegeHandler_GetCompanies_Company_With_Company_Access_Found(t *testing.T) { handler := New(WithBaseURL("base")) - result, err := handler.Process(&PrivilegeAdded{ + err := handler.Process(&PrivilegeAdded{ Email: "jim@example.org", CompanyID: "abc-123", Privilege: PrivilegeCompany, - }, goamqp.Headers{}) - assert.Nil(t, result) + }) assert.NoError(t, err) companies := handler.CompaniesByUser("jim@example.org", func(privileges CompanyPrivileges) bool { @@ -162,12 +154,11 @@ func TestPrivilegeHandler_GetCompanies_Company_With_Company_Access_Found(t *test func TestPrivilegeHandler_GetCompanies_Company_With_Admin_Access_Found(t *testing.T) { handler := New(WithBaseURL("base")) - result, err := handler.Process(&PrivilegeAdded{ + err := handler.Process(&PrivilegeAdded{ Email: "jim@example.org", CompanyID: "abc-123", Privilege: PrivilegeConsumer, - }, goamqp.Headers{}) - assert.Nil(t, result) + }) assert.NoError(t, err) companies := handler.CompaniesByUser("jim@example.org", func(privileges CompanyPrivileges) bool { @@ -190,11 +181,11 @@ func TestPrivilegeHandler_IsAllowed_Return_False_If_No_Privileges(t *testing.T) func TestPrivilegeHandler_IsAllowed_Return_True_If_Privilege_Exists(t *testing.T) { handler := New(WithBaseURL("base")) - _, _ = handler.Process(&PrivilegeAdded{ + _ = handler.Process(&PrivilegeAdded{ Email: "jim@example.org", CompanyID: "abc-123", Privilege: PrivilegeTime, - }, goamqp.Headers{}) + }) result := handler.IsAllowed("jim@example.org", "abc-123", func(privileges CompanyPrivileges) bool { return privileges.Time @@ -202,11 +193,11 @@ func TestPrivilegeHandler_IsAllowed_Return_True_If_Privilege_Exists(t *testing.T assert.True(t, result) - _, _ = handler.Process(&PrivilegeAdded{ + _ = handler.Process(&PrivilegeAdded{ Email: "jim@example.org", CompanyID: "abc-123", Privilege: PrivilegeInvoicing, - }, goamqp.Headers{}) + }) result = handler.IsAllowed("jim@example.org", "abc-123", func(privileges CompanyPrivileges) bool { return privileges.Invoicing @@ -214,11 +205,11 @@ func TestPrivilegeHandler_IsAllowed_Return_True_If_Privilege_Exists(t *testing.T assert.True(t, result) - _, _ = handler.Process(&PrivilegeAdded{ + _ = handler.Process(&PrivilegeAdded{ Email: "jim@example.org", CompanyID: "abc-123", Privilege: PrivilegeAccounting, - }, goamqp.Headers{}) + }) result = handler.IsAllowed("jim@example.org", "abc-123", func(privileges CompanyPrivileges) bool { return privileges.Accounting @@ -226,11 +217,11 @@ func TestPrivilegeHandler_IsAllowed_Return_True_If_Privilege_Exists(t *testing.T assert.True(t, result) - _, _ = handler.Process(&PrivilegeAdded{ + _ = handler.Process(&PrivilegeAdded{ Email: "jim@example.org", CompanyID: "abc-123", Privilege: PrivilegeSupplier, - }, goamqp.Headers{}) + }) result = handler.IsAllowed("jim@example.org", "abc-123", func(privileges CompanyPrivileges) bool { return privileges.Supplier @@ -238,11 +229,11 @@ func TestPrivilegeHandler_IsAllowed_Return_True_If_Privilege_Exists(t *testing.T assert.True(t, result) - _, _ = handler.Process(&PrivilegeAdded{ + _ = handler.Process(&PrivilegeAdded{ Email: "jim@example.org", CompanyID: "abc-123", Privilege: PrivilegeSalary, - }, goamqp.Headers{}) + }) result = handler.IsAllowed("jim@example.org", "abc-123", func(privileges CompanyPrivileges) bool { return privileges.Salary @@ -522,11 +513,11 @@ func TestPrivilegeHandler_Concurrent_Process_And_Read(t *testing.T) { companyID := fmt.Sprintf("company-%d", i%10) go func(id string) { defer wg.Done() - _, _ = handler.Process(&PrivilegeAdded{ + _ = handler.Process(&PrivilegeAdded{ Email: "jim@example.org", CompanyID: id, Privilege: PrivilegeAdmin, - }, goamqp.Headers{}) + }) }(companyID) } @@ -598,11 +589,11 @@ func TestPrivilegeHandler_Concurrent_Multiple_Operations(t *testing.T) { wg.Add(1) go func(idx int) { defer wg.Done() - _, _ = handler.Process(&PrivilegeAdded{ + _ = handler.Process(&PrivilegeAdded{ Email: "jane@example.org", CompanyID: fmt.Sprintf("company-%d", idx%5), Privilege: PrivilegeCompany, - }, goamqp.Headers{}) + }) }(i) } @@ -648,3 +639,28 @@ func TestPrivilegeHandler_Concurrent_Multiple_Operations(t *testing.T) { expectedJane := []string{"company-0", "company-1", "company-2", "company-3", "company-4"} assert.Equal(t, expectedJane, janeCompanies) } + +func TestPrivilegeHandler_Setup(t *testing.T) { + handler := New(WithBaseURL("base")) + + topology, err := goamqp.CollectTopology("some-service", handler.Setup()...) + assert.NoError(t, err) + var keys []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) + } + assert.Equal(t, []string{"User.Added", "User.Removed", "Privilege.Added", "Privilege.Removed"}, keys) +} + +func TestPrivilegeHandler_process(t *testing.T) { + handler := New(WithBaseURL("base")) + + err := process[UserAdded](handler)(context.Background(), spec.ConsumableEvent[UserAdded]{ + Payload: UserAdded{Email: "jim@example.org", CompanyID: "abc-123"}, + }) + assert.NoError(t, err) + + assert.Equal(t, []string{"abc-123"}, handler.CompaniesByUser("jim@example.org", func(CompanyPrivileges) bool { return true })) +} diff --git a/go.mod b/go.mod index cf460e2..ddfc797 100644 --- a/go.mod +++ b/go.mod @@ -3,13 +3,29 @@ module gitea.unbound.se/shiny/authz_client go 1.26.2 require ( - github.com/sparetimecoders/goamqp v0.3.3 + codeberg.org/messaging/go-messaging-amqp v0.0.4 + codeberg.org/messaging/messaging v0.0.5 github.com/stretchr/testify v1.12.1 ) require ( + github.com/beorn7/perks v1.0.1 // indirect + github.com/cespare/xxhash/v2 v2.3.0 // indirect + github.com/go-logr/logr v1.4.3 // indirect + github.com/go-logr/stdr v1.2.2 // indirect github.com/google/uuid v1.6.0 // indirect - github.com/pkg/errors v0.9.1 // indirect - github.com/rabbitmq/amqp091-go v1.10.0 // indirect + github.com/munnerz/goautoneg v0.0.0-20191010083416-a7dc8b61c822 // indirect + github.com/prometheus/client_golang v1.23.2 // indirect + github.com/prometheus/client_model v0.6.2 // indirect + github.com/prometheus/common v0.66.1 // indirect + github.com/prometheus/procfs v0.16.1 // indirect + github.com/rabbitmq/amqp091-go v1.12.0 // indirect + go.opentelemetry.io/auto/sdk v1.2.1 // indirect + go.opentelemetry.io/otel v1.44.0 // indirect + go.opentelemetry.io/otel/metric v1.44.0 // indirect + go.opentelemetry.io/otel/trace v1.44.0 // indirect + go.yaml.in/yaml/v2 v2.4.2 // indirect go.yaml.in/yaml/v3 v3.0.5 // indirect + golang.org/x/sys v0.45.0 // indirect + google.golang.org/protobuf v1.36.8 // indirect ) diff --git a/go.sum b/go.sum index 57397a7..6ae674b 100644 --- a/go.sum +++ b/go.sum @@ -1,14 +1,72 @@ +codeberg.org/messaging/go-messaging-amqp v0.0.4 h1:MwY/kU1lBdCL7ywKAg2k6aHuJf8f78MRZQSuyrNFvQ8= +codeberg.org/messaging/go-messaging-amqp v0.0.4/go.mod h1:6bSCIkKH0V/oI7vaIq8ttINKYwov9c8ckZrwjfvp7kc= +codeberg.org/messaging/messaging v0.0.5 h1:/ueH90F4RNPUeeJIwdG9oVhDW7+XjCzwOYq8xc0kfQk= +codeberg.org/messaging/messaging v0.0.5/go.mod h1:xyWLUcfaVzcN6GWk9uaQsxPVUC2Vm2S89mYZGdiWTWg= +codeberg.org/messaging/messaging/tck v0.0.3 h1:a4Nr7ytFqEJloTOyVeAtMNgO9Fig1ZUirjW4FVlK0lc= +codeberg.org/messaging/messaging/tck v0.0.3/go.mod h1:RiOsKXGAhNQK2STBVYGjhN0CFtmcrsEne1qUe15n8Q4= +github.com/beorn7/perks v1.0.1 h1:VlbKKnNfV8bJzeqoa4cOKqO6bYr3WgKZxO8Z16+hsOM= +github.com/beorn7/perks v1.0.1/go.mod h1:G2ZrVWU2WbWT9wwq4/hrbKbnv/1ERSJQ0ibhJ6rlkpw= +github.com/cespare/xxhash/v2 v2.3.0 h1:UL815xU9SqsFlibzuggzjXhog7bL6oX9BbNZnL2UFvs= +github.com/cespare/xxhash/v2 v2.3.0/go.mod h1:VGX0DQ3Q6kWi7AoAeZDth3/j3BFtOZR5XLFGgcrjCOs= +github.com/go-logr/logr v1.2.2/go.mod h1:jdQByPbusPIv2/zmleS9BjJVeZ6kBagPoEUsqbVz/1A= +github.com/go-logr/logr v1.4.3 h1:CjnDlHq8ikf6E492q6eKboGOC0T8CDaOvkHCIg8idEI= +github.com/go-logr/logr v1.4.3/go.mod h1:9T104GzyrTigFIr8wt5mBrctHMim0Nb2HLGrmQ40KvY= +github.com/go-logr/stdr v1.2.2 h1:hSWxHoqTgW2S2qGc0LTAI563KZ5YKYRhT3MFKZMbjag= +github.com/go-logr/stdr v1.2.2/go.mod h1:mMo/vtBO5dYbehREoey6XUKy/eSumjCCveDpRre4VKE= +github.com/google/go-cmp v0.7.0 h1:wk8382ETsv4JYUZwIsn6YpYiWiBsYLSJiTsyBybVuN8= +github.com/google/go-cmp v0.7.0/go.mod h1:pXiqmnSA92OHEEa9HXL2W4E7lf9JzCmGVUdgjX3N/iU= github.com/google/uuid v1.6.0 h1:NIvaJDMOsjHA8n1jAhLSgzrAzy1Hgr+hNrb57e+94F0= github.com/google/uuid v1.6.0/go.mod h1:TIyPZe4MgqvfeYDBFedMoGGpEw/LqOeaOT+nhxU+yHo= -github.com/pkg/errors v0.9.1 h1:FEBLx1zS214owpjy7qsBeixbURkuhQAwrK5UwLGTwt4= -github.com/pkg/errors v0.9.1/go.mod h1:bwawxfHBFNV+L2hUp1rHADufV3IMtnDRdf1r5NINEl0= -github.com/rabbitmq/amqp091-go v1.10.0 h1:STpn5XsHlHGcecLmMFCtg7mqq0RnD+zFr4uzukfVhBw= -github.com/rabbitmq/amqp091-go v1.10.0/go.mod h1:Hy4jKW5kQART1u+JkDTF9YYOQUHXqMuhrgxOEeS7G4o= -github.com/sparetimecoders/goamqp v0.3.3 h1:z/nfTPmrjeU/rIVuNOgsVLCimp3WFoNFvS3ZzXRJ6HE= -github.com/sparetimecoders/goamqp v0.3.3/go.mod h1:W9NRCpWLE+Vruv2dcRSbszNil2O826d2Nv6kAkETW5o= +github.com/klauspost/compress v1.18.5 h1:/h1gH5Ce+VWNLSWqPzOVn6XBO+vJbCNGvjoaGBFW2IE= +github.com/klauspost/compress v1.18.5/go.mod h1:cwPg85FWrGar70rWktvGQj8/hthj3wpl0PGDogxkrSQ= +github.com/kr/pretty v0.3.1 h1:flRD4NNwYAUpkphVc1HcthR4KEIFJ65n8Mw5qdRn3LE= +github.com/kr/pretty v0.3.1/go.mod h1:hoEshYVHaxMs3cyo3Yncou5ZscifuDolrwPKZanG3xk= +github.com/kr/text v0.2.0 h1:5Nx0Ya0ZqY2ygV366QzturHI13Jq95ApcVaJBhpS+AY= +github.com/kr/text v0.2.0/go.mod h1:eLer722TekiGuMkidMxC/pM04lWEeraHUUmBw8l2grE= +github.com/munnerz/goautoneg v0.0.0-20191010083416-a7dc8b61c822 h1:C3w9PqII01/Oq1c1nUAm88MOHcQC9l5mIlSMApZMrHA= +github.com/munnerz/goautoneg v0.0.0-20191010083416-a7dc8b61c822/go.mod h1:+n7T8mK8HuQTcFwEeznm/DIxMOiR9yIdICNftLE1DvQ= +github.com/nats-io/nats.go v1.52.0 h1:n3avV4VBsCgsdwh71TppsTwtv+QdPs7ntSKM8qJLGsc= +github.com/nats-io/nats.go v1.52.0/go.mod h1:26HypzazeOkyO3/mqd1zZd53STJN0EjCYF9Uy2ZOBno= +github.com/nats-io/nkeys v0.4.15 h1:JACV5jRVO9V856KOapQ7x+EY8Jo3qw1vJt/9Jpwzkk4= +github.com/nats-io/nkeys v0.4.15/go.mod h1:CpMchTXC9fxA5zrMo4KpySxNjiDVvr8ANOSZdiNfUrs= +github.com/nats-io/nuid v1.0.1 h1:5iA8DT8V7q8WK2EScv2padNa/rTESc1KdnPw4TC2paw= +github.com/nats-io/nuid v1.0.1/go.mod h1:19wcPz3Ph3q0Jbyiqsd0kePYG7A95tJPxeL+1OSON2c= +github.com/prometheus/client_golang v1.23.2 h1:Je96obch5RDVy3FDMndoUsjAhG5Edi49h0RJWRi/o0o= +github.com/prometheus/client_golang v1.23.2/go.mod h1:Tb1a6LWHB3/SPIzCoaDXI4I8UHKeFTEQ1YCr+0Gyqmg= +github.com/prometheus/client_model v0.6.2 h1:oBsgwpGs7iVziMvrGhE53c/GrLUsZdHnqNwqPLxwZyk= +github.com/prometheus/client_model v0.6.2/go.mod h1:y3m2F6Gdpfy6Ut/GBsUqTWZqCUvMVzSfMLjcu6wAwpE= +github.com/prometheus/common v0.66.1 h1:h5E0h5/Y8niHc5DlaLlWLArTQI7tMrsfQjHV+d9ZoGs= +github.com/prometheus/common v0.66.1/go.mod h1:gcaUsgf3KfRSwHY4dIMXLPV0K/Wg1oZ8+SbZk/HH/dA= +github.com/prometheus/procfs v0.16.1 h1:hZ15bTNuirocR6u0JZ6BAHHmwS1p8B4P6MRqxtzMyRg= +github.com/prometheus/procfs v0.16.1/go.mod h1:teAbpZRB1iIAJYREa1LsoWUXykVXA1KlTmWl8x/U+Is= +github.com/rabbitmq/amqp091-go v1.12.0 h1:V0v14Iqfs+MwHWihJt/nGS5Ulu0vw572b2Co3mwunkI= +github.com/rabbitmq/amqp091-go v1.12.0/go.mod h1:Hy4jKW5kQART1u+JkDTF9YYOQUHXqMuhrgxOEeS7G4o= +github.com/rogpeppe/go-internal v1.14.1 h1:UQB4HGPB6osV0SQTLymcB4TgvyWu6ZyliaW0tI/otEQ= +github.com/rogpeppe/go-internal v1.14.1/go.mod h1:MaRKkUm5W0goXpeCfT7UZI6fk/L7L7so1lCWt35ZSgc= github.com/stretchr/testify v1.12.1 h1:EuwCh5fleGS7H32xRwO3wRGT7DxrDhLAT6FF8MpWDWE= github.com/stretchr/testify v1.12.1/go.mod h1:MDEgiDPPsNp5cuIrHPPCyornHKgEVbtFUmoNlxoYthg= +go.opentelemetry.io/auto/sdk v1.2.1 h1:jXsnJ4Lmnqd11kwkBV2LgLoFMZKizbCi5fNZ/ipaZ64= +go.opentelemetry.io/auto/sdk v1.2.1/go.mod h1:KRTj+aOaElaLi+wW1kO/DZRXwkF4C5xPbEe3ZiIhN7Y= +go.opentelemetry.io/otel v1.44.0 h1:JjwHmHpA4iZ3wBxluu2fbbE7j4kqlE8jXyAyPXH7HqU= +go.opentelemetry.io/otel v1.44.0/go.mod h1:BMgjTHL9WPRlRjL2oZCBTL4whCGtXch2H4BhOPIAyYc= +go.opentelemetry.io/otel/metric v1.44.0 h1:1w0gILTcHdr3YI+ixLyjemwrVnsMURbTZFrSYCdDdmc= +go.opentelemetry.io/otel/metric v1.44.0/go.mod h1:8O7hanEPBNgEMmybD3s2VBKcgWOCsA6tzHBPODAiquo= +go.opentelemetry.io/otel/sdk v1.44.0 h1:nHYwb9lK+fJPU/dnT6s7W7Z8itMWyqrnVfbheVYrZ58= +go.opentelemetry.io/otel/sdk v1.44.0/go.mod h1:Osuydd3Se74nqjAKxid74N5eC+jfEqfTegHRnq58oK0= +go.opentelemetry.io/otel/trace v1.44.0 h1:jxF5CsGYCe74MCRx2X4g7WsY/VBKRqqpNvXlX/6gtIk= +go.opentelemetry.io/otel/trace v1.44.0/go.mod h1:oLl1jrMQAVo6v3GAggN+1VH9VIz9iUSvW53sW1Q8PIE= go.uber.org/goleak v1.3.0 h1:2K3zAYmnTNqV73imy9J1T3WC+gmCePx2hEGkimedGto= go.uber.org/goleak v1.3.0/go.mod h1:CoHD4mav9JJNrW/WLlf7HGZPjdw8EucARQHekz1X6bE= +go.yaml.in/yaml/v2 v2.4.2 h1:DzmwEr2rDGHl7lsFgAHxmNz/1NlQ7xLIrlN2h5d1eGI= +go.yaml.in/yaml/v2 v2.4.2/go.mod h1:081UH+NErpNdqlCXm3TtEran0rJZGxAYx9hb/ELlsPU= go.yaml.in/yaml/v3 v3.0.5 h1:N6y/pJk8buWs9NY5ERU2HSMfm+IuD/OtfdAnq6kESPw= go.yaml.in/yaml/v3 v3.0.5/go.mod h1:HVTZu1O7/Vkt2N+BFy8Zza+lnLsABggaTM2ZpNIGuKg= +golang.org/x/crypto v0.49.0 h1:+Ng2ULVvLHnJ/ZFEq4KdcDd/cfjrrjjNSXNzxg0Y4U4= +golang.org/x/crypto v0.49.0/go.mod h1:ErX4dUh2UM+CFYiXZRTcMpEcN8b/1gxEuv3nODoYtCA= +golang.org/x/sys v0.45.0 h1:dO4czNzziLiiXplLQgBCEpCvXQ3dnkn0SdaZSYdQ+FY= +golang.org/x/sys v0.45.0/go.mod h1:4GL1E5IUh+htKOUEOaiffhrAeqysfVGipDYzABqnCmw= +google.golang.org/protobuf v1.36.8 h1:xHScyCOEuuwZEc6UtSOvPbAT4zRh0xcNRYekJwfqyMc= +google.golang.org/protobuf v1.36.8/go.mod h1:fuxRtAxBytpl4zzqUh6/eyUujkJdNiuEkXntxiD/uRU= +gopkg.in/check.v1 v0.0.0-20161208181325-20d25e280405/go.mod h1:Co6ibVJAznAaIkqp8huTwlJQCZ016jof/cbN4VW5Yz0= +gopkg.in/check.v1 v1.0.0-20201130134442-10cb98267c6c h1:Hei/4ADfdWqJk1ZMxUNpqntNwaWcugrBjAiHlqqRiVk= +gopkg.in/check.v1 v1.0.0-20201130134442-10cb98267c6c/go.mod h1:JHkPIbrfpd72SG/EVd6muEfDQjcINNoR0C8j2r3qZ4Q= -- 2.54.0 From fc0a4a0e315299ceace1926badf3a63123dededd Mon Sep 17 00:00:00 2001 From: Joakim Olsson Date: Fri, 11 Sep 2026 22:57:25 +0200 Subject: [PATCH 2/2] test: assert which event type each privilege routing key decodes to Swapping the handler for Privilege.Removed to PrivilegeAdded passed the suite, and would turn every revocation into a grant. The Setup test now asserts each key's message type, the adapter is exercised for all four events (grant then revoke), and process only accepts the four privilege event types. CLAUDE.md warns against combining Setup() with WithReconnect (per-replica queues are re-created on reconnect, losing revocations sent meanwhile). Co-Authored-By: Claude Opus 5 (1M context) Claude-Session: https://claude.ai/code/session_01SV1epy2pKvBx3yDfhxATk2 --- CLAUDE.md | 2 +- client.go | 7 ++++++- client_test.go | 35 ++++++++++++++++++++++++++++------- 3 files changed, 35 insertions(+), 9 deletions(-) diff --git a/CLAUDE.md b/CLAUDE.md index 23d8116..79cd1d8 100644 --- a/CLAUDE.md +++ b/CLAUDE.md @@ -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. diff --git a/client.go b/client.go index b876dc7..03dddf1 100644 --- a/client.go +++ b/client.go @@ -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) } diff --git a/client_test.go b/client_test.go index ed926b3..ac69459 100644 --- a/client_test.go +++ b/client_test.go @@ -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 })) } -- 2.54.0