Compare commits

..

2 Commits

Author SHA1 Message Date
dusan 1de563ba65 Fix AMQP wildcard
Signed-off-by: dusan <borovcanindusan1@gmail.com>
2026-07-04 19:24:36 +02:00
Arvindh ee762fd1db add hooks endpoints
Signed-off-by: Arvindh <arvindh91@gmail.com>
2026-07-02 15:30:55 +05:30
21 changed files with 188 additions and 474 deletions
+1 -1
View File
@@ -21,7 +21,7 @@ jobs:
cache-dependency-path: "go.sum"
- name: Run linters
uses: golangci/golangci-lint-action@v9.3.0
uses: golangci/golangci-lint-action@v9.2.1
with:
version: v2.10.1
args: --config ./tools/config/.golangci.yaml
-1
View File
@@ -57,7 +57,6 @@ type PageMetadata struct {
Limit uint64 `json:"limit" db:"limit"`
DomainID string `json:"domain_id" db:"domain_id"`
RuleID string `json:"rule_id" db:"rule_id"`
RuleIDs []string `json:"rule_ids" db:"rule_ids"`
ChannelID string `json:"channel_id" db:"channel_id"`
ClientID string `json:"client_id" db:"client_id"`
Subtopic string `json:"subtopic" db:"subtopic"`
+11 -136
View File
@@ -33,16 +33,6 @@ type authorizationMiddleware struct {
var _ alarms.Service = (*authorizationMiddleware)(nil)
const (
atomObjectKindResource = "resource"
atomObjectTypeResourceRule = "resource:" + atom.KindRule
atomAuthorizedRulePageLimit = 100
)
type atomAuthorizedObjectLister interface {
AuthorizedObjectIDs(ctx context.Context, q atom.AuthorizedObjectIDsQuery) (atom.AuthorizedObjectIDs, error)
}
func NewAuthorizationMiddleware(svc alarms.Service, authz smqauthz.Authorization, entitiesOps permissions.EntitiesOperations[permissions.Operation]) (alarms.Service, error) {
if err := entitiesOps.Validate(); err != nil {
return nil, err
@@ -72,19 +62,14 @@ func (am *authorizationMiddleware) CreateAlarm(ctx context.Context, alarm alarms
}
func (am *authorizationMiddleware) UpdateAlarm(ctx context.Context, session authn.Session, alarm alarms.Alarm) (alarms.Alarm, error) {
current, err := am.svc.ViewAlarm(ctx, session, alarm.ID)
if err != nil {
return alarms.Alarm{}, err
}
if len(alarm.Metadata) > 0 {
if err := am.authorizeAlarmOrRule(ctx, operations.OpUpdateAlarm, session, current); err != nil {
if err := am.authorize(ctx, operations.OpUpdateAlarm, session, policies.DomainType, session.DomainID); err != nil {
return alarms.Alarm{}, errors.Wrap(errDomainUpdateAlarms, err)
}
}
if alarm.AssigneeID != "" {
if err := am.authorizeAlarmOrRule(ctx, operations.OpAssignAlarm, session, current); err != nil {
if err := am.authorize(ctx, operations.OpAssignAlarm, session, policies.DomainType, session.DomainID); err != nil {
return alarms.Alarm{}, errors.Wrap(errDomainUpdateAlarms, err)
}
if am.atomAuthz == nil {
@@ -104,13 +89,13 @@ func (am *authorizationMiddleware) UpdateAlarm(ctx context.Context, session auth
}
if alarm.AcknowledgedBy != "" {
if err := am.authorizeAlarmOrRule(ctx, operations.OpAcknowledgeAlarm, session, current); err != nil {
if err := am.authorize(ctx, operations.OpAcknowledgeAlarm, session, policies.DomainType, session.DomainID); err != nil {
return alarms.Alarm{}, errors.Wrap(errDomainUpdateAlarms, err)
}
}
if alarm.ResolvedBy != "" {
if err := am.authorizeAlarmOrRule(ctx, operations.OpResolveAlarm, session, current); err != nil {
if err := am.authorize(ctx, operations.OpResolveAlarm, session, policies.DomainType, session.DomainID); err != nil {
return alarms.Alarm{}, errors.Wrap(errDomainUpdateAlarms, err)
}
}
@@ -119,11 +104,7 @@ func (am *authorizationMiddleware) UpdateAlarm(ctx context.Context, session auth
}
func (am *authorizationMiddleware) DeleteAlarm(ctx context.Context, session authn.Session, id string) error {
alarm, err := am.svc.ViewAlarm(ctx, session, id)
if err != nil {
return err
}
if err := am.authorizeAlarmOrRule(ctx, operations.OpDeleteAlarm, session, alarm); err != nil {
if err := am.authorize(ctx, operations.OpDeleteAlarm, session, policies.DomainType, session.DomainID); err != nil {
return errors.Wrap(errDomainDeleteAlarms, err)
}
@@ -139,25 +120,8 @@ func (am *authorizationMiddleware) ListAlarms(ctx context.Context, session authn
case err == nil:
session.SuperAdmin = true
case errors.Contains(err, svcerr.ErrSuperAdminAction):
if err := am.authorizeTenantAlarm(ctx, operations.OpViewAlarm, session); err != nil {
if pm.RuleID != "" {
if ruleErr := am.authorizeRuleAlarmRead(ctx, session, pm.RuleID); ruleErr != nil {
return alarms.AlarmsPage{}, errors.Wrap(errDomainViewAlarms, err)
}
break
}
ruleIDs, ruleErr := am.authorizedReadableRuleIDs(ctx, session)
if ruleErr != nil {
return alarms.AlarmsPage{}, errors.Wrap(errDomainViewAlarms, err)
}
if len(ruleIDs) == 0 {
return alarms.AlarmsPage{
Offset: pm.Offset,
Limit: pm.Limit,
Alarms: []alarms.Alarm{},
}, nil
}
pm.RuleIDs = ruleIDs
if err := am.authorize(ctx, operations.OpListAlarms, session, policies.DomainType, session.DomainID); err != nil {
return alarms.AlarmsPage{}, errors.Wrap(errDomainViewAlarms, err)
}
default:
return alarms.AlarmsPage{}, err
@@ -167,109 +131,20 @@ func (am *authorizationMiddleware) ListAlarms(ctx context.Context, session authn
}
func (am *authorizationMiddleware) ViewAlarm(ctx context.Context, session authn.Session, id string) (alarms.Alarm, error) {
alarm, err := am.svc.ViewAlarm(ctx, session, id)
if err != nil {
return alarms.Alarm{}, err
}
if err := am.authorizeViewAlarm(ctx, session, alarm); err != nil {
if err := am.authorize(ctx, operations.OpViewAlarm, session, policies.DomainType, session.DomainID); err != nil {
return alarms.Alarm{}, errors.Wrap(errDomainViewAlarms, err)
}
return alarm, nil
return am.svc.ViewAlarm(ctx, session, id)
}
func (am *authorizationMiddleware) authorizeAlarmOrRule(ctx context.Context, op permissions.Operation, session authn.Session, alarm alarms.Alarm) error {
tenantErr := am.authorizeTenantAlarm(ctx, op, session)
if tenantErr == nil {
return nil
}
if alarm.RuleID == "" {
return tenantErr
}
if err := am.authorize(ctx, op, session, policies.RulesType, alarm.RuleID, atom.KindRule); err != nil {
return tenantErr
}
return nil
}
func (am *authorizationMiddleware) authorizeViewAlarm(ctx context.Context, session authn.Session, alarm alarms.Alarm) error {
tenantErr := am.authorizeTenantAlarm(ctx, operations.OpViewAlarm, session)
if tenantErr == nil {
return nil
}
if alarm.RuleID == "" {
return tenantErr
}
if err := am.authorizeRuleAlarmRead(ctx, session, alarm.RuleID); err != nil {
return tenantErr
}
return nil
}
func (am *authorizationMiddleware) authorizeTenantAlarm(ctx context.Context, op permissions.Operation, session authn.Session) error {
return am.authorize(ctx, op, session, policies.DomainType, session.DomainID, atom.KindAlarm)
}
func (am *authorizationMiddleware) authorizeRuleAlarmRead(ctx context.Context, session authn.Session, ruleID string) error {
if am.atomAuthz != nil {
return am.authorize(ctx, operations.OpViewAlarm, session, policies.RulesType, ruleID, atom.KindRule)
}
perm, err := am.entitiesOps.GetPermission(operations.EntityType, operations.OpViewAlarm)
if err != nil {
return err
}
pr := smqauthz.PolicyReq{
Domain: session.DomainID,
SubjectType: policies.UserType,
SubjectKind: policies.UsersKind,
Subject: session.DomainUserID,
Object: ruleID,
ObjectType: policies.RulesType,
Permission: perm.String(),
}
return am.authz.Authorize(ctx, pr, nil)
}
func (am *authorizationMiddleware) authorizedReadableRuleIDs(ctx context.Context, session authn.Session) ([]string, error) {
lister, ok := am.atomAuthz.(atomAuthorizedObjectLister)
if !ok {
return nil, errors.ErrAuthorization
}
perm, err := am.entitiesOps.GetPermission(operations.EntityType, operations.OpViewAlarm)
if err != nil {
return nil, err
}
var ids []string
for offset := uint64(0); ; offset += atomAuthorizedRulePageLimit {
page, err := lister.AuthorizedObjectIDs(ctx, atom.AuthorizedObjectIDsQuery{
SubjectID: atom.SubjectID(session),
Action: atom.CapabilityName(perm.String()),
ObjectKind: atomObjectKindResource,
ObjectType: atomObjectTypeResourceRule,
TenantID: session.DomainID,
Limit: atomAuthorizedRulePageLimit,
Offset: offset,
})
if err != nil {
return nil, err
}
ids = append(ids, page.IDs...)
if uint64(len(page.IDs)) < atomAuthorizedRulePageLimit || offset+uint64(len(page.IDs)) >= page.Total {
break
}
}
return ids, nil
}
func (am *authorizationMiddleware) authorize(ctx context.Context, op permissions.Operation, session authn.Session, objType, obj, resourceKind string) error {
func (am *authorizationMiddleware) authorize(ctx context.Context, op permissions.Operation, session authn.Session, objType, obj string) error {
perm, err := am.entitiesOps.GetPermission(operations.EntityType, op)
if err != nil {
return err
}
if am.atomAuthz != nil {
return atom.Authorize(ctx, am.atomAuthz, session, perm.String(), objType, obj, resourceKind)
return atom.Authorize(ctx, am.atomAuthz, session, perm.String(), objType, obj, atom.KindAlarm)
}
pr := smqauthz.PolicyReq{
+11 -127
View File
@@ -12,6 +12,7 @@ import (
"github.com/absmach/magistrala/alarms/operations"
"github.com/absmach/magistrala/internal/atom"
"github.com/absmach/magistrala/pkg/authn"
pkgerrors "github.com/absmach/magistrala/pkg/errors"
"github.com/absmach/magistrala/pkg/permissions"
"github.com/stretchr/testify/assert"
"github.com/stretchr/testify/mock"
@@ -19,27 +20,16 @@ import (
)
type recordingAtomAuthorizer struct {
allowed bool
allow func(atom.AuthzRequest) bool
authorized atom.AuthorizedObjectIDs
reqs []atom.AuthzRequest
queries []atom.AuthorizedObjectIDsQuery
allowed bool
reqs []atom.AuthzRequest
}
func (a *recordingAtomAuthorizer) CheckAuthz(_ context.Context, req atom.AuthzRequest) (atom.AuthzResponse, error) {
a.reqs = append(a.reqs, req)
if a.allow != nil {
return atom.AuthzResponse{Allowed: a.allow(req)}, nil
}
return atom.AuthzResponse{Allowed: a.allowed}, nil
}
func (a *recordingAtomAuthorizer) AuthorizedObjectIDs(_ context.Context, q atom.AuthorizedObjectIDsQuery) (atom.AuthorizedObjectIDs, error) {
a.queries = append(a.queries, q)
return a.authorized, nil
}
func TestListAlarmsAuthorizesTenantAlarmReader(t *testing.T) {
func TestListAlarmsAuthorizesRegularUser(t *testing.T) {
svc := mocks.NewService(t)
pm := alarms.PageMetadata{Limit: 10}
expectedPM := pm
@@ -57,7 +47,7 @@ func TestListAlarmsAuthorizesTenantAlarmReader(t *testing.T) {
require.Len(t, authz.reqs, 1)
assert.Equal(t, atom.AuthzRequest{
SubjectID: "user-1",
Action: "alarm_read",
Action: "list",
ResourceID: "",
ObjectKind: "tenant",
ObjectID: "domain-1",
@@ -68,60 +58,16 @@ func TestListAlarmsAuthorizesTenantAlarmReader(t *testing.T) {
}, authz.reqs[0])
}
func TestListAlarmsFiltersToReadableRulesWhenTenantAlarmReadDenied(t *testing.T) {
func TestListAlarmsDeniedRegularUserDoesNotDelegate(t *testing.T) {
svc := mocks.NewService(t)
pm := alarms.PageMetadata{Limit: 10}
expectedPM := pm
expectedPM.DomainID = "domain-1"
expectedPM.RuleIDs = []string{"rule-1", "rule-2"}
authz := &recordingAtomAuthorizer{
allowed: false,
authorized: atom.AuthorizedObjectIDs{IDs: []string{"rule-1", "rule-2"}, Total: 2},
}
authz := &recordingAtomAuthorizer{allowed: false}
wrapped, err := NewAtomAuthorizationMiddleware(svc, authz, testEntitiesOps(t))
require.NoError(t, err)
svc.On("ListAlarms", mock.Anything, authn.Session{UserID: "user-1", DomainID: "domain-1"}, expectedPM).Return(alarms.AlarmsPage{Limit: 10}, nil).Once()
_, err = wrapped.ListAlarms(context.Background(), authn.Session{UserID: "user-1", DomainID: "domain-1"}, pm)
_, err = wrapped.ListAlarms(context.Background(), authn.Session{UserID: "user-1", DomainID: "domain-1"}, alarms.PageMetadata{})
require.NoError(t, err)
assert.True(t, pkgerrors.Contains(err, pkgerrors.ErrAuthorization))
require.Len(t, authz.reqs, 1)
assert.Equal(t, "alarm_read", authz.reqs[0].Action)
assert.Equal(t, "tenant", authz.reqs[0].ObjectKind)
require.Len(t, authz.queries, 1)
assert.Equal(t, atom.AuthorizedObjectIDsQuery{
SubjectID: "user-1",
Action: "alarm_read",
ObjectKind: "resource",
ObjectType: "resource:rule",
TenantID: "domain-1",
Limit: 100,
}, authz.queries[0])
}
func TestListAlarmsWithRuleFilterAuthorizesRuleRead(t *testing.T) {
svc := mocks.NewService(t)
pm := alarms.PageMetadata{Limit: 10, RuleID: "rule-1"}
expectedPM := pm
expectedPM.DomainID = "domain-1"
session := authn.Session{UserID: "user-1", DomainID: "domain-1"}
authz := &recordingAtomAuthorizer{
allow: func(req atom.AuthzRequest) bool {
return req.Action == "alarm_read" && req.ObjectKind == "resource" && req.ObjectID == "rule-1"
},
}
wrapped, err := NewAtomAuthorizationMiddleware(svc, authz, testEntitiesOps(t))
require.NoError(t, err)
svc.On("ListAlarms", mock.Anything, session, expectedPM).Return(alarms.AlarmsPage{Limit: 10}, nil).Once()
_, err = wrapped.ListAlarms(context.Background(), session, pm)
require.NoError(t, err)
require.Len(t, authz.reqs, 2)
assert.Equal(t, "alarm_read", authz.reqs[0].Action)
assert.Equal(t, "alarm_read", authz.reqs[1].Action)
assert.Equal(t, "resource", authz.reqs[1].ObjectKind)
assert.Equal(t, "rules", authz.reqs[1].Context["legacy_object_type"])
}
func TestListAlarmsSuperAdminSkipsListAuthorization(t *testing.T) {
@@ -144,56 +90,13 @@ func TestListAlarmsSuperAdminSkipsListAuthorization(t *testing.T) {
assert.Equal(t, "manage", authz.reqs[0].Action)
}
func TestAcknowledgeAlarmAuthorizesRuleAlarmActionWhenTenantDenied(t *testing.T) {
svc := mocks.NewService(t)
session := authn.Session{UserID: "user-1", DomainID: "domain-1"}
current := alarms.Alarm{ID: "alarm-1", RuleID: "rule-1", DomainID: "domain-1"}
update := alarms.Alarm{ID: "alarm-1", AcknowledgedBy: "user-1"}
authz := &recordingAtomAuthorizer{
allow: func(req atom.AuthzRequest) bool {
return req.Action == "alarm_acknowledge" && req.ObjectKind == "resource" && req.ObjectID == "rule-1"
},
}
wrapped, err := NewAtomAuthorizationMiddleware(svc, authz, testEntitiesOps(t))
require.NoError(t, err)
svc.On("ViewAlarm", mock.Anything, session, "alarm-1").Return(current, nil).Once()
svc.On("UpdateAlarm", mock.Anything, session, update).Return(update, nil).Once()
_, err = wrapped.UpdateAlarm(context.Background(), session, update)
require.NoError(t, err)
require.Len(t, authz.reqs, 2)
assert.Equal(t, atom.AuthzRequest{
SubjectID: "user-1",
Action: "alarm_acknowledge",
ResourceID: "",
ObjectKind: "tenant",
ObjectID: "domain-1",
Context: map[string]any{
"domain_id": "domain-1",
"legacy_object_type": "domain",
},
}, authz.reqs[0])
assert.Equal(t, atom.AuthzRequest{
SubjectID: "user-1",
Action: "alarm_acknowledge",
ResourceID: "rule-1",
ObjectKind: "resource",
ObjectID: "rule-1",
Context: map[string]any{
"domain_id": "domain-1",
"legacy_object_type": "rules",
},
}, authz.reqs[1])
}
func testEntitiesOps(t *testing.T) permissions.EntitiesOperations[permissions.Operation] {
t.Helper()
details := operations.OperationDetails()
perms := make(map[string]permissions.Permission, len(details))
for op, detail := range details {
for _, detail := range details {
if detail.PermissionRequired {
perms[detail.Name] = testPermission(op, detail.Name)
perms[detail.Name] = permissions.Permission(detail.Name)
}
}
entitiesOps, err := permissions.NewEntitiesOperations(
@@ -203,22 +106,3 @@ func testEntitiesOps(t *testing.T) permissions.EntitiesOperations[permissions.Op
require.NoError(t, err)
return entitiesOps
}
func testPermission(op permissions.Operation, fallback string) permissions.Permission {
switch op {
case operations.OpViewAlarm, operations.OpListAlarms:
return "alarm_read_permission"
case operations.OpUpdateAlarm:
return "alarm_update_permission"
case operations.OpDeleteAlarm:
return "alarm_delete_permission"
case operations.OpAssignAlarm:
return "alarm_assign_permission"
case operations.OpAcknowledgeAlarm:
return "alarm_acknowledge_permission"
case operations.OpResolveAlarm:
return "alarm_resolve_permission"
default:
return permissions.Permission(fallback)
}
}
-3
View File
@@ -464,9 +464,6 @@ func pageQueryConditions(pm alarms.PageMetadata) []string {
if pm.RuleID != "" {
query = append(query, "alarms.rule_id = :rule_id")
}
if len(pm.RuleIDs) > 0 {
query = append(query, "alarms.rule_id = ANY(:rule_ids)")
}
if pm.ChannelID != "" {
query = append(query, "alarms.channel_id = :channel_id")
}
+3 -14
View File
@@ -4,7 +4,7 @@
// Code generated by protoc-gen-go. DO NOT EDIT.
// versions:
// protoc-gen-go v1.36.11
// protoc v7.35.1
// protoc v6.33.0
// source: readers/v1/readers.proto
package v1
@@ -104,7 +104,6 @@ type PageMetadata struct {
Format string `protobuf:"bytes,17,opt,name=format,proto3" json:"format,omitempty"`
Order string `protobuf:"bytes,18,opt,name=order,proto3" json:"order,omitempty"`
Dir string `protobuf:"bytes,19,opt,name=dir,proto3" json:"dir,omitempty"`
Publishers []string `protobuf:"bytes,20,rep,name=publishers,proto3" json:"publishers,omitempty"`
unknownFields protoimpl.UnknownFields
sizeCache protoimpl.SizeCache
}
@@ -272,13 +271,6 @@ func (x *PageMetadata) GetDir() string {
return ""
}
func (x *PageMetadata) GetPublishers() []string {
if x != nil {
return x.Publishers
}
return nil
}
type ReadMessagesRes struct {
state protoimpl.MessageState `protogen:"open.v1"`
Total uint64 `protobuf:"varint,1,opt,name=total,proto3" json:"total,omitempty"`
@@ -730,7 +722,7 @@ var File_readers_v1_readers_proto protoreflect.FileDescriptor
const file_readers_v1_readers_proto_rawDesc = "" +
"\n" +
"\x18readers/v1/readers.proto\x12\n" +
"readers.v1\"\xac\x04\n" +
"readers.v1\"\x8c\x04\n" +
"\fPageMetadata\x12\x14\n" +
"\x05limit\x18\x01 \x01(\x04R\x05limit\x12\x16\n" +
"\x06offset\x18\x02 \x01(\x04R\x06offset\x12\x1a\n" +
@@ -755,10 +747,7 @@ const file_readers_v1_readers_proto_rawDesc = "" +
"comparator\x12\x16\n" +
"\x06format\x18\x11 \x01(\tR\x06format\x12\x14\n" +
"\x05order\x18\x12 \x01(\tR\x05order\x12\x10\n" +
"\x03dir\x18\x13 \x01(\tR\x03dir\x12\x1e\n" +
"\n" +
"publishers\x18\x14 \x03(\tR\n" +
"publishers\"\x97\x01\n" +
"\x03dir\x18\x13 \x01(\tR\x03dir\"\x97\x01\n" +
"\x0fReadMessagesRes\x12\x14\n" +
"\x05total\x18\x01 \x01(\x04R\x05total\x12=\n" +
"\rpage_metadata\x18\x02 \x01(\v2\x18.readers.v1.PageMetadataR\fpageMetadata\x12/\n" +
+1 -5
View File
@@ -1029,12 +1029,8 @@ services:
- ./templates/${MG_REPORTS_EMAIL_TEMPLATE}:/email.tmpl
pdf-generator:
image: gotenberg/gotenberg:8.34.0
image: gotenberg/gotenberg:8.25.1
container_name: magistrala-pdf
restart: on-failure
# Gotenberg leaks memory over time (upstream issues #987, #1169); the cap
# keeps the OOM kill inside this container instead of stalling the host.
mem_limit: 512m
ports:
- "4000:3000"
networks:
+17 -8
View File
@@ -20,8 +20,6 @@ const (
hookAuthOnPublish = "auth_on_publish"
hookAuthOnSubscribe = "auth_on_subscribe"
hookAuthOnUnsubscribe = "auth_on_unsubscribe"
hookProtocolAMQP091 = "amqp091"
)
type hookRequest struct {
@@ -85,9 +83,11 @@ func handleHook(ctx context.Context, parser messaging.TopicParser, req hookReque
}
func resolveHookTopic(ctx context.Context, parser messaging.TopicParser, req hookRequest) (string, error) {
if isInternalMessageConsumer(req) {
return "", nil
hook := strings.ToLower(strings.TrimSpace(req.Hook))
if isAMQP091MessageStreamConsume(req, hook) {
return strings.TrimPrefix(strings.TrimSpace(req.Topic), "/"), nil
}
if !isMessageTopic(req.Topic) {
return "", nil
}
@@ -99,7 +99,7 @@ func resolveHookTopic(ctx context.Context, parser messaging.TopicParser, req hoo
var topicType messaging.TopicType
var err error
switch strings.ToLower(strings.TrimSpace(req.Hook)) {
switch hook {
case hookAuthOnPublish:
domainID, channelID, subtopic, topicType, err = parser.ParsePublishTopic(ctx, req.Topic, true)
case hookAuthOnSubscribe, hookAuthOnUnsubscribe:
@@ -123,14 +123,23 @@ func isMessageTopic(topic string) bool {
return strings.HasPrefix(topic, string(messaging.MsgTopicPrefix)+"/")
}
func isInternalMessageConsumer(req hookRequest) bool {
hook := strings.ToLower(strings.TrimSpace(req.Hook))
// isAMQP091MessageStreamConsume reports whether the request is an AMQP 0-9-1
// stream-queue consume of the full message firehose (m/#), which is passed
// through without parsing because the topic parser cannot resolve a
// channel-level wildcard.
//
// SECURITY: in the default deployment the auth callout is disabled for
// amqp091 (docker/fluxmq/node*.yaml), so this allow is the only gate for
// stream consume. The amqp091 listener must remain network-restricted until
// identity-gated authorization for m/# lands in the gRPC Authorize path.
func isAMQP091MessageStreamConsume(req hookRequest, hook string) bool {
if hook != hookAuthOnSubscribe && hook != hookAuthOnUnsubscribe {
return false
}
if strings.ToLower(strings.TrimSpace(req.Protocol)) != hookProtocolAMQP091 {
if strings.ToLower(strings.TrimSpace(req.Protocol)) != "amqp091" {
return false
}
topic := strings.TrimPrefix(strings.TrimSpace(req.Topic), "/")
return topic == string(messaging.MsgTopicPrefix)+"/#"
}
+46 -21
View File
@@ -84,6 +84,52 @@ func TestHooksHandlerUsesSubscribeParserForSubscribeAndUnsubscribe(t *testing.T)
}
}
func TestHooksHandlerAllowsAMQP091MessageStreamWildcard(t *testing.T) {
cases := []struct {
desc string
protocol string
topic string
}{
{desc: "plain topic", protocol: "amqp091", topic: "m/#"},
{desc: "leading slash topic", protocol: "amqp091", topic: "/m/#"},
{desc: "uppercase protocol", protocol: "AMQP091", topic: "m/#"},
}
for _, tc := range cases {
for _, hook := range []string{hookAuthOnSubscribe, hookAuthOnUnsubscribe} {
parser := &fakeHookParser{err: errors.New("must not parse stream queue wildcard")}
req := httptest.NewRequest("POST", "/hooks", strings.NewReader(`{"hook":"`+hook+`","protocol":"`+tc.protocol+`","topic":"`+tc.topic+`"}`))
w := httptest.NewRecorder()
MakeHooksHandler(parser).ServeHTTP(w, req)
require.Equal(t, 200, w.Code, tc.desc)
require.False(t, parser.publishCalled, tc.desc)
require.False(t, parser.subscribeCalled, tc.desc)
var res hookResponse
require.NoError(t, json.NewDecoder(w.Body).Decode(&res), tc.desc)
require.Equal(t, hookResultOK, res.Result, tc.desc)
require.Equal(t, "m/#", res.Topic, tc.desc)
}
}
}
func TestHooksHandlerStillParsesMQTTMessageWildcard(t *testing.T) {
parser := &fakeHookParser{err: errors.New("malformed topic")}
req := httptest.NewRequest("POST", "/hooks", strings.NewReader(`{"hook":"auth_on_subscribe","protocol":"mqtt","topic":"m/#"}`))
w := httptest.NewRecorder()
MakeHooksHandler(parser).ServeHTTP(w, req)
require.Equal(t, 200, w.Code)
require.False(t, parser.publishCalled)
require.True(t, parser.subscribeCalled)
var res hookResponse
require.NoError(t, json.NewDecoder(w.Body).Decode(&res))
require.Equal(t, hookResultDeny, res.Result)
}
func TestHooksHandlerReturnsOKForNonMGTopic(t *testing.T) {
req := httptest.NewRequest("POST", "/hooks", strings.NewReader(`{"hook":"auth_on_publish","topic":"$SYS/broker/uptime"}`))
w := httptest.NewRecorder()
@@ -97,27 +143,6 @@ func TestHooksHandlerReturnsOKForNonMGTopic(t *testing.T) {
require.Empty(t, res.Topic)
}
func TestHooksHandlerAllowsInternalAMQPMessageConsumer(t *testing.T) {
parser := &fakeHookParser{err: errors.New("should not parse internal consumer")}
req := httptest.NewRequest("POST", "/hooks", strings.NewReader(`{
"hook":"auth_on_subscribe",
"protocol":"amqp091",
"topic":"m/#"
}`))
w := httptest.NewRecorder()
MakeHooksHandler(parser).ServeHTTP(w, req)
require.Equal(t, 200, w.Code)
require.False(t, parser.publishCalled)
require.False(t, parser.subscribeCalled)
var res hookResponse
require.NoError(t, json.NewDecoder(w.Body).Decode(&res))
require.Equal(t, hookResultOK, res.Result)
require.Empty(t, res.Topic)
}
func TestHooksHandlerDeniesUnresolvedMGTopic(t *testing.T) {
parser := &fakeHookParser{err: errors.New("failed to resolve channel route")}
req := httptest.NewRequest("POST", "/hooks", strings.NewReader(`{"hook":"auth_on_publish","topic":"m/d1/c/ch1/messages"}`))
+10 -10
View File
@@ -7,15 +7,15 @@ require (
connectrpc.com/otelconnect v0.9.0
github.com/0x6flab/namegenerator v1.4.0
github.com/absmach/callhome v0.18.2
github.com/absmach/fluxmq v0.40.0
github.com/absmach/fluxmq v0.30.0
github.com/absmach/senml v1.0.8
github.com/caarlos0/env/v10 v10.0.0
github.com/caarlos0/env/v11 v11.4.1
github.com/dgraph-io/ristretto/v2 v2.4.2
github.com/dgraph-io/ristretto/v2 v2.4.0
github.com/eclipse/paho.mqtt.golang v1.5.1
github.com/fatih/color v1.19.0
github.com/fiorix/go-smpp v0.0.0-20210403173735-2894b96e70ba
github.com/go-chi/chi/v5 v5.3.1
github.com/go-chi/chi/v5 v5.3.0
github.com/go-kit/kit v0.13.0
github.com/gofrs/uuid/v5 v5.4.0
github.com/google/uuid v1.6.0
@@ -27,7 +27,7 @@ require (
github.com/jackc/pgtype v1.14.4
github.com/jackc/pgx/v5 v5.10.0
github.com/jmoiron/sqlx v1.4.0
github.com/lestrrat-go/jwx/v2 v2.1.7
github.com/lestrrat-go/jwx/v2 v2.1.6
github.com/lib/pq v1.12.3
github.com/mitchellh/mapstructure v1.5.0
github.com/nats-io/nats.go v1.52.0
@@ -54,11 +54,11 @@ require (
go.opentelemetry.io/otel/exporters/otlp/otlptrace/otlptracehttp v1.44.0
go.opentelemetry.io/otel/sdk v1.44.0
go.opentelemetry.io/otel/trace v1.44.0
golang.org/x/crypto v0.54.0
golang.org/x/net v0.57.0
golang.org/x/sync v0.22.0
golang.org/x/crypto v0.53.0
golang.org/x/net v0.56.0
golang.org/x/sync v0.21.0
gonum.org/v1/gonum v0.17.0
google.golang.org/grpc v1.82.1
google.golang.org/grpc v1.81.1
google.golang.org/protobuf v1.36.11
gopkg.in/gomail.v2 v2.0.0-20160411212932-81ebce5c23df
gopkg.in/yaml.v3 v3.0.1
@@ -168,8 +168,8 @@ require (
go.uber.org/atomic v1.11.0 // indirect
go.yaml.in/yaml/v3 v3.0.4 // indirect
golang.org/x/exp v0.0.0-20260611194520-c48552f49976 // indirect
golang.org/x/sys v0.47.0 // indirect
golang.org/x/text v0.40.0 // indirect
golang.org/x/sys v0.46.0 // indirect
golang.org/x/text v0.38.0 // indirect
golang.org/x/time v0.15.0 // indirect
google.golang.org/genproto/googleapis/api v0.0.0-20260610212136-7ab31c22f7ad // indirect
gopkg.in/alexcesaro/quotedprintable.v3 v3.0.0-20150716171945-2caba252f4dc // indirect
+20 -20
View File
@@ -24,8 +24,8 @@ github.com/VividCortex/gohistogram v1.0.0 h1:6+hBz+qvs0JOrrNhhmR7lFxo5sINxBCGXrd
github.com/VividCortex/gohistogram v1.0.0/go.mod h1:Pf5mBqqDxYaXu3hDrrU+w6nw50o/4+TcAqDqk/vUH7g=
github.com/absmach/callhome v0.18.2 h1:dmopRHm2qTheHN1hdUKRRYpKwRrj7X9d8AWCFrb+K6s=
github.com/absmach/callhome v0.18.2/go.mod h1:LEXKhES9JJtj3tBgTZv7VPNjOi5ukJQB0mFic0QP60Q=
github.com/absmach/fluxmq v0.40.0 h1:J7s6PHXliWRfwpQpt/umRNnrUegj9xYsVyK9BV8Azhk=
github.com/absmach/fluxmq v0.40.0/go.mod h1:oVbq3VlkD0vPKc45gkcxXpD2tbweyF+CDw1x72AS+PA=
github.com/absmach/fluxmq v0.30.0 h1:U1DCgfg+aimz3B5L5Q2T9BGlIE7SDWH4hR+L2gv1vQQ=
github.com/absmach/fluxmq v0.30.0/go.mod h1:o8RKhK8AseC9MdWphJkAuLBh3mWXzpUpU5i5CGzS7s0=
github.com/absmach/senml v1.0.8 h1:+opem/r4g6c6eA/JLyCIuksyEhj7eBdysY3pEmy1mqo=
github.com/absmach/senml v1.0.8/go.mod h1:DRhzHLgvQoIUHroBgpFrSWso+bJZO9E96RlHAHy+VRI=
github.com/alecthomas/template v0.0.0-20160405071501-a0175ee3bccc/go.mod h1:LOuyumcjzFXgccqObfd/Ljyb9UuFJ6TxHnclSeseNhc=
@@ -82,8 +82,8 @@ github.com/davecgh/go-spew v1.1.2-0.20180830191138-d8f796af33cc h1:U9qPSI2PIWSS1
github.com/davecgh/go-spew v1.1.2-0.20180830191138-d8f796af33cc/go.mod h1:J7Y8YcW2NihsgmVo/mv3lAwl/skON4iLHjSsI+c5H38=
github.com/decred/dcrd/dcrec/secp256k1/v4 v4.4.1 h1:5RVFMOWjMyRy8cARdy79nAmgYw3hK/4HUq48LQ6Wwqo=
github.com/decred/dcrd/dcrec/secp256k1/v4 v4.4.1/go.mod h1:ZXNYxsqcloTdSy/rNShjYzMhyjf0LaoftYK0p+A3h40=
github.com/dgraph-io/ristretto/v2 v2.4.2 h1:x0cvjmUKxt764Yxdk2nr94we1AvPPAMh1rh5TQ+Jo80=
github.com/dgraph-io/ristretto/v2 v2.4.2/go.mod h1:0KsrXtXvnv0EqnzyowllbVJB8yBonswa2lTCK2gGo9E=
github.com/dgraph-io/ristretto/v2 v2.4.0 h1:I/w09yLjhdcVD2QV192UJcq8dPBaAJb9pOuMyNy0XlU=
github.com/dgraph-io/ristretto/v2 v2.4.0/go.mod h1:0KsrXtXvnv0EqnzyowllbVJB8yBonswa2lTCK2gGo9E=
github.com/dgryski/go-farm v0.0.0-20240924180020-3414d57e47da h1:aIftn67I1fkbMa512G+w+Pxci9hJPB8oMnkcP3iZF38=
github.com/dgryski/go-farm v0.0.0-20240924180020-3414d57e47da/go.mod h1:SqUrOPUnsFjfmXRMNPybcSiG0BgUW2AuFH8PAnS2iTw=
github.com/distribution/reference v0.6.0 h1:0IXCQ5g4/QMHHkarYzh5l+u8T3t73zM5QvfrDyIgxBk=
@@ -115,8 +115,8 @@ github.com/fsnotify/fsnotify v1.10.1 h1:b0/UzAf9yR5rhf3RPm9gf3ehBPpf0oZKIjtpKrx5
github.com/fsnotify/fsnotify v1.10.1/go.mod h1:TLheqan6HD6GBK6PrDWyDPBaEV8LspOxvPSjC+bVfgo=
github.com/fxamacker/cbor/v2 v2.9.2 h1:X4Ksno9+x3cz0TZv69ec1hxP/+tymuR8PXQJyDwfh78=
github.com/fxamacker/cbor/v2 v2.9.2/go.mod h1:vM4b+DJCtHn+zz7h3FFp/hDAI9WNWCsZj23V5ytsSxQ=
github.com/go-chi/chi/v5 v5.3.1 h1:3j4HZLGZQ3JpMCrPJF/Jl3mYJfWLKBfNJ6quurUGCf8=
github.com/go-chi/chi/v5 v5.3.1/go.mod h1:R+tYY2hNuVUUjxoPtqUdgBqevM9s9njzkTLutVsOCto=
github.com/go-chi/chi/v5 v5.3.0 h1:halUjDxhshgXHMrao5bB8eNBXo/rnzwr8m5m36glehM=
github.com/go-chi/chi/v5 v5.3.0/go.mod h1:R+tYY2hNuVUUjxoPtqUdgBqevM9s9njzkTLutVsOCto=
github.com/go-gorp/gorp/v3 v3.1.0 h1:ItKF/Vbuj31dmV4jxA1qblpSwkl9g1typ24xoe70IGs=
github.com/go-gorp/gorp/v3 v3.1.0/go.mod h1:dLEjIyyRNiXvNZ8PSmzpt1GsWAUK8kjVhEpjH8TixEw=
github.com/go-jose/go-jose/v4 v4.1.4 h1:moDMcTHmvE6Groj34emNPLs/qtYXRVcd6S7NHbHz3kA=
@@ -311,8 +311,8 @@ github.com/lestrrat-go/httprc v1.0.6 h1:qgmgIRhpvBqexMJjA/PmwSvhNk679oqD1RbovdCG
github.com/lestrrat-go/httprc v1.0.6/go.mod h1:mwwz3JMTPBjHUkkDv/IGJ39aALInZLrhBp0X7KGUZlo=
github.com/lestrrat-go/iter v1.0.2 h1:gMXo1q4c2pHmC3dn8LzRhJfP1ceCbgSiT9lUydIzltI=
github.com/lestrrat-go/iter v1.0.2/go.mod h1:Momfcq3AnRlRjI5b5O8/G5/BvpzrhoFTZcn06fEOPt4=
github.com/lestrrat-go/jwx/v2 v2.1.7 h1:bnYeET+S8IOyAw6W4LTc6SEeK7Xs58SKKZkR7scb3Ko=
github.com/lestrrat-go/jwx/v2 v2.1.7/go.mod h1:exQ9ZBuN1cMLYmxwhTlHUru08ykONG0z+HbLEeDG9qo=
github.com/lestrrat-go/jwx/v2 v2.1.6 h1:hxM1gfDILk/l5ylers6BX/Eq1m/pnxe9NBwW6lVfecA=
github.com/lestrrat-go/jwx/v2 v2.1.6/go.mod h1:Y722kU5r/8mV7fYDifjug0r8FK8mZdw0K0GpJw/l8pU=
github.com/lestrrat-go/option v1.0.1 h1:oAzP2fvZGQKWkvHa1/SAcFolBEca1oN+mQ7eooNBEYU=
github.com/lestrrat-go/option v1.0.1/go.mod h1:5ZHFbivi4xwXxhxY9XHDe2FHo6/Z7WWmtT7T5nBBp3I=
github.com/lib/pq v1.0.0/go.mod h1:5WUZQaWbwv1U+lTReE5YruASi9Al49XbQIvNi/34Woo=
@@ -576,8 +576,8 @@ golang.org/x/crypto v0.0.0-20210711020723-a769d52b0f97/go.mod h1:GvvjBRRGRdwPK5y
golang.org/x/crypto v0.0.0-20210921155107-089bfa567519/go.mod h1:GvvjBRRGRdwPK5ydBHafDWAxML/pGHZbMvKqRZ5+Abc=
golang.org/x/crypto v0.19.0/go.mod h1:Iy9bg/ha4yyC70EfRS8jz+B6ybOBKMaSxLj6P6oBDfU=
golang.org/x/crypto v0.20.0/go.mod h1:Xwo95rrVNIoSMx9wa1JroENMToLWn3RNVrTBpLHgZPQ=
golang.org/x/crypto v0.54.0 h1:YLIA59K4fiNzHzjnZt2tUJQjQtUWfWbeHBqKtk3eScw=
golang.org/x/crypto v0.54.0/go.mod h1:KWL8ny2AZdGR2cWmzeHrp2azQPGogOv+HeQaVEXC2dk=
golang.org/x/crypto v0.53.0 h1:QZ4Muo8THX6CizN2vPPd5fBGHyogrdK9fG4wLPFUsto=
golang.org/x/crypto v0.53.0/go.mod h1:DNLU434OwVakk9PzuwV8w62mAJpRJL3vsgcfp4Qnsio=
golang.org/x/exp v0.0.0-20260611194520-c48552f49976 h1:X8Hz2ImujgbmetVuW+w2YkyZChE3cBpZi2P158rTG9M=
golang.org/x/exp v0.0.0-20260611194520-c48552f49976/go.mod h1:vnf4pv9iKZXY58sQE1L86zmNWJ4159e1RkcWiLCkeEY=
golang.org/x/lint v0.0.0-20190930215403-16217165b5de/go.mod h1:6SW0HCj/g11FgYtHlgUYUwCkIfeOF89ocIRzGO/8vkc=
@@ -600,8 +600,8 @@ golang.org/x/net v0.0.0-20220722155237-a158d28d115b/go.mod h1:XRhObCWvk6IyKnWLug
golang.org/x/net v0.6.0/go.mod h1:2Tu9+aMcznHK/AK1HMvgo6xiTLG5rD5rZLDS+rp2Bjs=
golang.org/x/net v0.10.0/go.mod h1:0qNGK6F8kojg2nk9dLZ2mShWaEBan6FAoqfSigmmuDg=
golang.org/x/net v0.21.0/go.mod h1:bIjVDfnllIU7BJ2DNgfnXvpSvtn8VRwhlsaeUTyUS44=
golang.org/x/net v0.57.0 h1:K5+3DljvIuDG9/Jv9rvyMywYNFCQ9RSUY6OOTTkT+tE=
golang.org/x/net v0.57.0/go.mod h1:KpXc8iv+r3XplLAG/f7Jsf9RPszJzdR0f58q9vGOuEU=
golang.org/x/net v0.56.0 h1:Rw8j/hFzGvJUZwNBXnAtf5sVDVt+65SK2C7IxCxZt5o=
golang.org/x/net v0.56.0/go.mod h1:D3Ku6r+V6JROoZK144D2XfMHFcMq/0zSfLelVTCFKec=
golang.org/x/oauth2 v0.0.0-20190226205417-e64efc72b421/go.mod h1:gOpvHmFTYa4IltrdGE7lF6nIHvwfUNPOp7c8zoXwtLw=
golang.org/x/sync v0.0.0-20181108010431-42b317875d0f/go.mod h1:RxMgew5VJxzue5/jJTE5uejpjVlOe/izrB70Jof72aM=
golang.org/x/sync v0.0.0-20181221193216-37e7f081c4d4/go.mod h1:RxMgew5VJxzue5/jJTE5uejpjVlOe/izrB70Jof72aM=
@@ -610,8 +610,8 @@ golang.org/x/sync v0.0.0-20190911185100-cd5d95a43a6e/go.mod h1:RxMgew5VJxzue5/jJ
golang.org/x/sync v0.0.0-20201207232520-09787c993a3a/go.mod h1:RxMgew5VJxzue5/jJTE5uejpjVlOe/izrB70Jof72aM=
golang.org/x/sync v0.0.0-20220722155255-886fb9371eb4/go.mod h1:RxMgew5VJxzue5/jJTE5uejpjVlOe/izrB70Jof72aM=
golang.org/x/sync v0.1.0/go.mod h1:RxMgew5VJxzue5/jJTE5uejpjVlOe/izrB70Jof72aM=
golang.org/x/sync v0.22.0 h1:SZjpbeLmrCk4xhRSZFNZW5gFUeCeFgjekvI/+gfScek=
golang.org/x/sync v0.22.0/go.mod h1:9xrNwdLfx4jkKbNva9FpL6vEN7evnE43NNNJQ2LF3+0=
golang.org/x/sync v0.21.0 h1:HLII4xRRTtCRkxYp4HNFF0Js/Og6q2i++KXbg0gHCwM=
golang.org/x/sync v0.21.0/go.mod h1:9xrNwdLfx4jkKbNva9FpL6vEN7evnE43NNNJQ2LF3+0=
golang.org/x/sys v0.0.0-20180905080454-ebe1bf3edb33/go.mod h1:STP8DvDyc/dI5b8T5hshtkjS+E42TnysNCUPdjciGhY=
golang.org/x/sys v0.0.0-20181116152217-5ac8a444bdc5/go.mod h1:STP8DvDyc/dI5b8T5hshtkjS+E42TnysNCUPdjciGhY=
golang.org/x/sys v0.0.0-20190204203706-41f3e6584952/go.mod h1:STP8DvDyc/dI5b8T5hshtkjS+E42TnysNCUPdjciGhY=
@@ -640,8 +640,8 @@ golang.org/x/sys v0.0.0-20220722155257-8c9f86f7a55f/go.mod h1:oPkhp1MJrh7nUepCBc
golang.org/x/sys v0.5.0/go.mod h1:oPkhp1MJrh7nUepCBck5+mAzfO9JrbApNNgaTdGDITg=
golang.org/x/sys v0.8.0/go.mod h1:oPkhp1MJrh7nUepCBck5+mAzfO9JrbApNNgaTdGDITg=
golang.org/x/sys v0.17.0/go.mod h1:/VUhepiaJMQUp4+oa/7Zr1D23ma6VTLIYjOOTFZPUcA=
golang.org/x/sys v0.47.0 h1:o7XGOvZQCADBQQ4Y7VNq2dRWQR7JmOUW8Kxx4ZsNgWs=
golang.org/x/sys v0.47.0/go.mod h1:4GL1E5IUh+htKOUEOaiffhrAeqysfVGipDYzABqnCmw=
golang.org/x/sys v0.46.0 h1:noSf2Fq6F8DBgS+LysIkx7rIExoNHJsxOAtPp4rthXw=
golang.org/x/sys v0.46.0/go.mod h1:4GL1E5IUh+htKOUEOaiffhrAeqysfVGipDYzABqnCmw=
golang.org/x/term v0.0.0-20201117132131-f5c789dd3221/go.mod h1:Nr5EML6q2oocZ2LXRh80K7BxOlk5/8JxuGnuhpl+muw=
golang.org/x/term v0.0.0-20201126162022-7de9c90e9dd1/go.mod h1:bj7SfCRtBDWHUb9snDiAeCFNEtKQo2Wmx5Cou7ajbmo=
golang.org/x/term v0.0.0-20210927222741-03fcf44c2211/go.mod h1:jbD1KX2456YbFQfuXm/mYQcufACuNUgVhRMnK/tPxf8=
@@ -657,8 +657,8 @@ golang.org/x/text v0.3.7/go.mod h1:u+2+/6zg+i71rQMx5EYifcz6MCKuco9NR6JIITiCfzQ=
golang.org/x/text v0.7.0/go.mod h1:mrYo+phRRbMaCq/xk9113O4dZlRixOauAjOtrjsXDZ8=
golang.org/x/text v0.9.0/go.mod h1:e1OnstbJyHTd6l/uOt8jFFHp6TRDWZR/bV3emEE/zU8=
golang.org/x/text v0.14.0/go.mod h1:18ZOQIKpY8NJVqYksKHtTdi31H5itFRjB5/qKTNYzSU=
golang.org/x/text v0.40.0 h1:Ub2Z6/xjgF1WrYQz2nuITOEegKFtiIy+rieRJ5lHZKs=
golang.org/x/text v0.40.0/go.mod h1:hpnzDAfGV753zIKo+wk3u1bVKCGPbrnF7+7LBF/UHVY=
golang.org/x/text v0.38.0 h1:sXmwo9DwP3OK9EZ7PqAdaooSGozfl/3a6/xJcbzPRhE=
golang.org/x/text v0.38.0/go.mod h1:YXZt3QhHUKYT53r2lLKFIVi6Ao1jdzrTR/KQ09qyxF4=
golang.org/x/time v0.0.0-20210220033141-f8bda1e9f3ba/go.mod h1:tRJNPiyCQ0inRvYxbN9jk5I+vvW/OXSQhTDSoE431IQ=
golang.org/x/time v0.15.0 h1:bbrp8t3bGUeFOx08pvsMYRTCVSMk89u4tKbNOZbp88U=
golang.org/x/time v0.15.0/go.mod h1:Y4YMaQmXwGQZoFaVFk4YpCt4FLQMYKZe9oeV/f4MSno=
@@ -686,8 +686,8 @@ google.golang.org/genproto/googleapis/api v0.0.0-20260610212136-7ab31c22f7ad h1:
google.golang.org/genproto/googleapis/api v0.0.0-20260610212136-7ab31c22f7ad/go.mod h1:KdNqO+rCIWgFumrNBSEDlDNrkrQnpkax7Tv1WxNY8V4=
google.golang.org/genproto/googleapis/rpc v0.0.0-20260610212136-7ab31c22f7ad h1:45WmJvIV6C2+O/jjLkPUH+F3aOj/1miDoU2DD0+NWbg=
google.golang.org/genproto/googleapis/rpc v0.0.0-20260610212136-7ab31c22f7ad/go.mod h1:4Hqkh8ycfw05ld/3BWL7rJOSfebL2Q+DVDeRgYgxUU8=
google.golang.org/grpc v1.82.1 h1:NnAxzGRA0677vCa4BUkOAnO5+FfQqVl9iUXeD0IqcGE=
google.golang.org/grpc v1.82.1/go.mod h1:yzTZ1TB1Z3SG+LIYaI+WiE8D5+PZ3ArnrSp8zF3+/ZA=
google.golang.org/grpc v1.81.1 h1:VnnIIZ88UzOOKLukQi+ImGz8O1Wdp8nAGGnvOfEIWQQ=
google.golang.org/grpc v1.81.1/go.mod h1:xGH9GfzOyMTGIOXBJmXt+BX/V0kcdQbdcuwQ/zNw42I=
google.golang.org/protobuf v0.0.0-20200109180630-ec00e32a8dfd/go.mod h1:DFci5gLYBciE7Vtevhsrf46CRTquxDuWsQurQQe4oz8=
google.golang.org/protobuf v0.0.0-20200221191635-4d8936d0db64/go.mod h1:kwYJMbMJ01Woi6D6+Kah6886xMZcty6N08ah7+eCXa0=
google.golang.org/protobuf v0.0.0-20200228230310-ab0ca4ff8a60/go.mod h1:cfTl7dwQJ+fmap5saPgwCLgHXTUD7jkjRqWcaiX5VyM=
-12
View File
@@ -58,18 +58,6 @@ func CapabilityName(action string) string {
normalized == "admin_permission",
strings.Contains(normalized, "manage_role"):
return atomActionManage
case strings.Contains(normalized, atomActionAlarmRead):
return atomActionAlarmRead
case strings.Contains(normalized, atomActionAlarmUpdate):
return atomActionAlarmUpdate
case strings.Contains(normalized, atomActionAlarmDelete):
return atomActionAlarmDelete
case strings.Contains(normalized, atomActionAlarmAssign):
return atomActionAlarmAssign
case strings.Contains(normalized, atomActionAlarmAcknowledge):
return atomActionAlarmAcknowledge
case strings.Contains(normalized, atomActionAlarmResolve):
return atomActionAlarmResolve
case normalized == policies.ViewPermission,
normalized == atomActionRead,
strings.Contains(normalized, "read"),
+6 -19
View File
@@ -17,23 +17,10 @@ var magistralaActionDescriptions = map[string]string{
atomActionSubscribe: "Subscribe to channel messages",
atomActionExecute: "Execute a command or action",
atomActionList: "List objects",
atomActionAlarmRead: "Read alarms in a tenant or rule scope",
atomActionAlarmUpdate: "Update alarms in a tenant or rule scope",
atomActionAlarmDelete: "Delete alarms in a tenant or rule scope",
atomActionAlarmAssign: "Assign alarms in a tenant or rule scope",
atomActionAlarmAcknowledge: "Acknowledge alarms in a tenant or rule scope",
atomActionAlarmResolve: "Resolve alarms in a tenant or rule scope",
}
var magistralaActionApplicability = []CapabilityApplicabilitySpec{
{ActionName: atomActionWrite, ObjectKind: atomObjectKindTenant},
{ActionName: atomActionAlarmRead, ObjectKind: atomObjectKindTenant},
{ActionName: atomActionAlarmUpdate, ObjectKind: atomObjectKindTenant},
{ActionName: atomActionAlarmDelete, ObjectKind: atomObjectKindTenant},
{ActionName: atomActionAlarmAssign, ObjectKind: atomObjectKindTenant},
{ActionName: atomActionAlarmAcknowledge, ObjectKind: atomObjectKindTenant},
{ActionName: atomActionAlarmResolve, ObjectKind: atomObjectKindTenant},
{ActionName: atomActionRead, ObjectKind: atomObjectKindGroup},
{ActionName: atomActionWrite, ObjectKind: atomObjectKindGroup},
@@ -54,12 +41,6 @@ var magistralaActionApplicability = []CapabilityApplicabilitySpec{
{ActionName: atomActionManage, ObjectKind: atomObjectKindResource, ObjectType: atomObjectTypeResourceRule},
{ActionName: atomActionExecute, ObjectKind: atomObjectKindResource, ObjectType: atomObjectTypeResourceRule},
{ActionName: atomActionList, ObjectKind: atomObjectKindResource, ObjectType: atomObjectTypeResourceRule},
{ActionName: atomActionAlarmRead, ObjectKind: atomObjectKindResource, ObjectType: atomObjectTypeResourceRule},
{ActionName: atomActionAlarmUpdate, ObjectKind: atomObjectKindResource, ObjectType: atomObjectTypeResourceRule},
{ActionName: atomActionAlarmDelete, ObjectKind: atomObjectKindResource, ObjectType: atomObjectTypeResourceRule},
{ActionName: atomActionAlarmAssign, ObjectKind: atomObjectKindResource, ObjectType: atomObjectTypeResourceRule},
{ActionName: atomActionAlarmAcknowledge, ObjectKind: atomObjectKindResource, ObjectType: atomObjectTypeResourceRule},
{ActionName: atomActionAlarmResolve, ObjectKind: atomObjectKindResource, ObjectType: atomObjectTypeResourceRule},
{ActionName: atomActionRead, ObjectKind: atomObjectKindResource, ObjectType: atomObjectTypeResourceReport},
{ActionName: atomActionWrite, ObjectKind: atomObjectKindResource, ObjectType: atomObjectTypeResourceReport},
@@ -67,6 +48,12 @@ var magistralaActionApplicability = []CapabilityApplicabilitySpec{
{ActionName: atomActionManage, ObjectKind: atomObjectKindResource, ObjectType: atomObjectTypeResourceReport},
{ActionName: atomActionExecute, ObjectKind: atomObjectKindResource, ObjectType: atomObjectTypeResourceReport},
{ActionName: atomActionList, ObjectKind: atomObjectKindResource, ObjectType: atomObjectTypeResourceReport},
{ActionName: atomActionRead, ObjectKind: atomObjectKindResource, ObjectType: atomObjectTypeResourceAlarm},
{ActionName: atomActionWrite, ObjectKind: atomObjectKindResource, ObjectType: atomObjectTypeResourceAlarm},
{ActionName: atomActionDelete, ObjectKind: atomObjectKindResource, ObjectType: atomObjectTypeResourceAlarm},
{ActionName: atomActionManage, ObjectKind: atomObjectKindResource, ObjectType: atomObjectTypeResourceAlarm},
{ActionName: atomActionList, ObjectKind: atomObjectKindResource, ObjectType: atomObjectTypeResourceAlarm},
}
var magistralaActionAssignmentRules = []ActionAssignmentRuleSpec{
+3 -13
View File
@@ -103,7 +103,7 @@ func TestBootstrapMagistralaActionsCreatesMissingActionsAndApplicability(t *test
t.Fatalf("bootstrap failed: %v", err)
}
for _, name := range []string{atomActionRead, atomActionWrite, atomActionDelete, atomActionManage, atomActionPublish, atomActionSubscribe, atomActionExecute, atomActionList, atomActionAlarmRead, atomActionAlarmUpdate, atomActionAlarmDelete, atomActionAlarmAssign, atomActionAlarmAcknowledge, atomActionAlarmResolve} {
for _, name := range []string{atomActionRead, atomActionWrite, atomActionDelete, atomActionManage, atomActionPublish, atomActionSubscribe, atomActionExecute, atomActionList} {
if _, ok := actions[name]; !ok {
t.Fatalf("action %q was not ensured", name)
}
@@ -112,12 +112,6 @@ func TestBootstrapMagistralaActionsCreatesMissingActionsAndApplicability(t *test
t.Fatalf("unexpected applicability count: got %d want %d", len(applicability), len(magistralaActionApplicability))
}
assertApplicability(t, applicability, "write-id", atomObjectKindTenant, "")
assertApplicability(t, applicability, "alarm_read-id", atomObjectKindTenant, "")
assertApplicability(t, applicability, "alarm_update-id", atomObjectKindTenant, "")
assertApplicability(t, applicability, "alarm_delete-id", atomObjectKindTenant, "")
assertApplicability(t, applicability, "alarm_assign-id", atomObjectKindTenant, "")
assertApplicability(t, applicability, "alarm_acknowledge-id", atomObjectKindTenant, "")
assertApplicability(t, applicability, "alarm_resolve-id", atomObjectKindTenant, "")
assertApplicability(t, applicability, "read-id", atomObjectKindGroup, "")
assertApplicability(t, applicability, "write-id", atomObjectKindGroup, "")
assertApplicability(t, applicability, "delete-id", atomObjectKindGroup, "")
@@ -126,14 +120,10 @@ func TestBootstrapMagistralaActionsCreatesMissingActionsAndApplicability(t *test
assertApplicability(t, applicability, "publish-id", atomObjectKindResource, "resource:channel")
assertApplicability(t, applicability, "execute-id", atomObjectKindResource, "resource:rule")
assertApplicability(t, applicability, "list-id", atomObjectKindResource, "resource:rule")
assertApplicability(t, applicability, "alarm_read-id", atomObjectKindResource, "resource:rule")
assertApplicability(t, applicability, "alarm_update-id", atomObjectKindResource, "resource:rule")
assertApplicability(t, applicability, "alarm_delete-id", atomObjectKindResource, "resource:rule")
assertApplicability(t, applicability, "alarm_assign-id", atomObjectKindResource, "resource:rule")
assertApplicability(t, applicability, "alarm_acknowledge-id", atomObjectKindResource, "resource:rule")
assertApplicability(t, applicability, "alarm_resolve-id", atomObjectKindResource, "resource:rule")
assertApplicability(t, applicability, "execute-id", atomObjectKindResource, "resource:report")
assertApplicability(t, applicability, "list-id", atomObjectKindResource, "resource:report")
assertApplicability(t, applicability, "manage-id", atomObjectKindResource, "resource:alarm")
assertApplicability(t, applicability, "list-id", atomObjectKindResource, "resource:alarm")
if len(assignmentRules) != len(magistralaActionAssignmentRules) {
t.Fatalf("unexpected assignment guardrail count: got %d want %d", len(assignmentRules), len(magistralaActionAssignmentRules))
}
+1 -7
View File
@@ -12,13 +12,6 @@ const (
atomActionSubscribe = "subscribe"
atomActionExecute = "execute"
atomActionList = "list"
atomActionAlarmRead = "alarm_read"
atomActionAlarmUpdate = "alarm_update"
atomActionAlarmDelete = "alarm_delete"
atomActionAlarmAssign = "alarm_assign"
atomActionAlarmAcknowledge = "alarm_acknowledge"
atomActionAlarmResolve = "alarm_resolve"
)
const (
@@ -50,6 +43,7 @@ const (
atomObjectTypeResourceChannel = "resource:channel"
atomObjectTypeResourceRule = "resource:rule"
atomObjectTypeResourceReport = "resource:report"
atomObjectTypeResourceAlarm = "resource:alarm"
)
const atomDecisionAllow = "allow"
-1
View File
@@ -34,7 +34,6 @@ message PageMetadata {
string format = 17;
string order = 18;
string dir = 19;
repeated string publishers = 20;
}
message ReadMessagesRes {
-2
View File
@@ -63,7 +63,6 @@ func (client readersGrpcClient) ReadMessages(ctx context.Context, in *grpcReader
Interval: in.GetPageMetadata().GetInterval(),
Subtopic: in.GetPageMetadata().GetSubtopic(),
Publisher: in.GetPageMetadata().GetPublisher(),
Publishers: in.GetPageMetadata().GetPublishers(),
Protocol: in.GetPageMetadata().GetProtocol(),
Name: in.GetPageMetadata().GetName(),
Value: in.GetPageMetadata().GetValue(),
@@ -117,7 +116,6 @@ func encodeReadMessagesRequest(_ context.Context, grpcReq any) (any, error) {
Interval: req.pageMeta.Interval,
Subtopic: req.pageMeta.Subtopic,
Publisher: req.pageMeta.Publisher,
Publishers: req.pageMeta.Publishers,
Protocol: req.pageMeta.Protocol,
Name: req.pageMeta.Name,
Value: req.pageMeta.Value,
-1
View File
@@ -46,7 +46,6 @@ func decodeReadMessagesRequest(_ context.Context, grpcReq any) (any, error) {
Interval: req.GetPageMetadata().GetInterval(),
Subtopic: req.GetPageMetadata().GetSubtopic(),
Publisher: req.GetPageMetadata().GetPublisher(),
Publishers: req.GetPageMetadata().GetPublishers(),
Protocol: req.GetPageMetadata().GetProtocol(),
Name: req.GetPageMetadata().GetName(),
Value: req.GetPageMetadata().GetValue(),
+18 -19
View File
@@ -41,25 +41,24 @@ type MessagesPage struct {
// PageMetadata represents the parameters used to create database queries.
type PageMetadata struct {
Offset uint64 `json:"offset"`
Limit uint64 `json:"limit"`
Order string `json:"order,omitempty"`
Dir string `json:"dir,omitempty"`
Subtopic string `json:"subtopic,omitempty"`
Publisher string `json:"publisher,omitempty"`
Publishers []string `json:"publishers,omitempty"`
Protocol string `json:"protocol,omitempty"`
Name string `json:"name,omitempty"`
Value float64 `json:"v,omitempty"`
Comparator string `json:"comparator,omitempty"`
BoolValue bool `json:"vb,omitempty"`
StringValue string `json:"vs,omitempty"`
DataValue string `json:"vd,omitempty"`
From float64 `json:"from,omitempty"`
To float64 `json:"to,omitempty"`
Format string `json:"format,omitempty"`
Aggregation string `json:"aggregation,omitempty"`
Interval string `json:"interval,omitempty"`
Offset uint64 `json:"offset"`
Limit uint64 `json:"limit"`
Order string `json:"order,omitempty"`
Dir string `json:"dir,omitempty"`
Subtopic string `json:"subtopic,omitempty"`
Publisher string `json:"publisher,omitempty"`
Protocol string `json:"protocol,omitempty"`
Name string `json:"name,omitempty"`
Value float64 `json:"v,omitempty"`
Comparator string `json:"comparator,omitempty"`
BoolValue bool `json:"vb,omitempty"`
StringValue string `json:"vs,omitempty"`
DataValue string `json:"vd,omitempty"`
From float64 `json:"from,omitempty"`
To float64 `json:"to,omitempty"`
Format string `json:"format,omitempty"`
Aggregation string `json:"aggregation,omitempty"`
Interval string `json:"interval,omitempty"`
}
// ParseValueComparator convert comparison operator keys into mathematic anotation.
+20 -30
View File
@@ -18,13 +18,12 @@ import (
var _ readers.MessageRepository = (*postgresRepository)(nil)
const (
messageFieldChannel = "channel"
messageFieldName = "name"
messageFieldProtocol = "protocol"
messageFieldPublisher = "publisher"
messageFieldPublishers = "publishers"
messageFieldSubtopic = "subtopic"
messageFieldValue = "value"
messageFieldChannel = "channel"
messageFieldName = "name"
messageFieldProtocol = "protocol"
messageFieldPublisher = "publisher"
messageFieldSubtopic = "subtopic"
messageFieldValue = "value"
)
type postgresRepository struct {
@@ -53,20 +52,19 @@ func (tr postgresRepository) ReadAll(chanID string, rpm readers.PageMetadata) (r
LIMIT :limit OFFSET :offset;`, format, cond, order)
params := map[string]any{
messageFieldChannel: chanID,
"limit": rpm.Limit,
"offset": rpm.Offset,
messageFieldSubtopic: rpm.Subtopic,
messageFieldPublisher: rpm.Publisher,
messageFieldPublishers: rpm.Publishers,
messageFieldName: rpm.Name,
messageFieldProtocol: rpm.Protocol,
messageFieldValue: rpm.Value,
"bool_value": rpm.BoolValue,
"string_value": rpm.StringValue,
"data_value": rpm.DataValue,
"from": rpm.From,
"to": rpm.To,
messageFieldChannel: chanID,
"limit": rpm.Limit,
"offset": rpm.Offset,
messageFieldSubtopic: rpm.Subtopic,
messageFieldPublisher: rpm.Publisher,
messageFieldName: rpm.Name,
messageFieldProtocol: rpm.Protocol,
messageFieldValue: rpm.Value,
"bool_value": rpm.BoolValue,
"string_value": rpm.StringValue,
"data_value": rpm.DataValue,
"from": rpm.From,
"to": rpm.To,
}
rows, err := tr.db.NamedQuery(q, params)
if err != nil {
@@ -140,19 +138,11 @@ func fmtCondition(chanID string, rpm readers.PageMetadata) string {
return condition
}
_, hasPublishers := query[messageFieldPublishers]
for name := range query {
switch name {
case messageFieldPublisher:
if hasPublishers {
continue
}
condition = fmt.Sprintf(`%s AND %s = :%s`, condition, name, name)
case messageFieldPublishers:
condition = fmt.Sprintf(`%s AND %s = ANY(:%s)`, condition, messageFieldPublisher, messageFieldPublishers)
case
messageFieldSubtopic,
messageFieldPublisher,
messageFieldName,
messageFieldProtocol:
condition = fmt.Sprintf(`%s AND %s = :%s`, condition, name, name)
+20 -24
View File
@@ -27,13 +27,12 @@ const (
var _ readers.MessageRepository = (*timescaleRepository)(nil)
const (
messageFieldChannel = "channel"
messageFieldName = "name"
messageFieldProtocol = "protocol"
messageFieldPublisher = "publisher"
messageFieldPublishers = "publishers"
messageFieldSubtopic = "subtopic"
messageFieldValue = "value"
messageFieldChannel = "channel"
messageFieldName = "name"
messageFieldProtocol = "protocol"
messageFieldPublisher = "publisher"
messageFieldSubtopic = "subtopic"
messageFieldValue = "value"
)
type timescaleRepository struct {
@@ -113,20 +112,19 @@ func (tr timescaleRepository) ReadAll(chanID string, rpm readers.PageMetadata) (
}
params := map[string]any{
messageFieldChannel: chanID,
"limit": rpm.Limit,
"offset": rpm.Offset,
messageFieldSubtopic: rpm.Subtopic,
messageFieldPublisher: rpm.Publisher,
messageFieldPublishers: rpm.Publishers,
messageFieldName: rpm.Name,
messageFieldProtocol: rpm.Protocol,
messageFieldValue: rpm.Value,
"bool_value": rpm.BoolValue,
"string_value": rpm.StringValue,
"data_value": rpm.DataValue,
"from": rpm.From,
"to": rpm.To,
messageFieldChannel: chanID,
"limit": rpm.Limit,
"offset": rpm.Offset,
messageFieldSubtopic: rpm.Subtopic,
messageFieldPublisher: rpm.Publisher,
messageFieldName: rpm.Name,
messageFieldProtocol: rpm.Protocol,
messageFieldValue: rpm.Value,
"bool_value": rpm.BoolValue,
"string_value": rpm.StringValue,
"data_value": rpm.DataValue,
"from": rpm.From,
"to": rpm.To,
}
rows, err := tr.db.NamedQuery(q, params)
@@ -208,9 +206,7 @@ func fmtCondition(rpm readers.PageMetadata) string {
conditions = append(conditions, " subtopic = :subtopic ")
}
if _, ok := query[messageFieldPublishers]; ok {
conditions = append(conditions, " publisher = ANY(:publishers) ")
} else if _, ok := query[messageFieldPublisher]; ok {
if _, ok := query[messageFieldPublisher]; ok {
conditions = append(conditions, " publisher = :publisher ")
}