From 8bba825706c2300e179d9bbec09751c6b37e8a47 Mon Sep 17 00:00:00 2001 From: Peter Grlica Date: Mon, 12 Apr 2021 10:27:54 +0200 Subject: [PATCH] Added rest api --- app/boot_levels.go | 3 +- pkg/messagebus/corteza.go | 2 +- pkg/messagebus/handler.go | 18 +- pkg/messagebus/messagebus.go | 27 +- pkg/messagebus/redis.go | 2 +- pkg/messagebus/settings.go | 21 +- pkg/messagebus/sql.go | 2 +- pkg/options/messagebus.gen.go | 37 ++ pkg/options/messagebus.yaml | 15 + store/messagebus_queuesettings.gen.go | 12 + store/messagebus_queuesettings.yaml | 9 +- store/rdbms/messagebus_queuesettings.gen.go | 19 +- store/rdbms/messagebus_queuesettings.go | 6 +- system/rest.yaml | 94 ++++ system/rest/handlers/queues.go | 133 +++++ system/rest/queues.go | 125 +++++ system/rest/request/queues.go | 437 +++++++++++++++ system/rest/router.go | 1 + system/service/queue.go | 210 +++++++ system/service/queue_actions.gen.go | 592 ++++++++++++++++++++ system/service/queue_actions.yaml | 89 +++ system/service/service.go | 2 + tests/messagebus/main_test.go | 6 +- 23 files changed, 1830 insertions(+), 32 deletions(-) create mode 100644 pkg/options/messagebus.gen.go create mode 100644 pkg/options/messagebus.yaml create mode 100644 system/rest/handlers/queues.go create mode 100644 system/rest/queues.go create mode 100644 system/rest/request/queues.go create mode 100644 system/service/queue.go create mode 100644 system/service/queue_actions.gen.go create mode 100644 system/service/queue_actions.yaml diff --git a/app/boot_levels.go b/app/boot_levels.go index b8d9edb64..6b7adece8 100644 --- a/app/boot_levels.go +++ b/app/boot_levels.go @@ -23,6 +23,7 @@ import ( "github.com/cortezaproject/corteza-server/pkg/mail" "github.com/cortezaproject/corteza-server/pkg/messagebus" "github.com/cortezaproject/corteza-server/pkg/monitor" + "github.com/cortezaproject/corteza-server/pkg/options" "github.com/cortezaproject/corteza-server/pkg/provision" "github.com/cortezaproject/corteza-server/pkg/rbac" "github.com/cortezaproject/corteza-server/pkg/scheduler" @@ -134,7 +135,7 @@ func (app *CortezaApp) Setup() (err error) { return err } - messagebus.Setup(app.Log, eventbus.Service()) + messagebus.Setup(options.Messagebus(), app.Log, eventbus.Service()) app.lvl = bootLevelSetup return diff --git a/pkg/messagebus/corteza.go b/pkg/messagebus/corteza.go index 3a298b372..7c08a4d15 100644 --- a/pkg/messagebus/corteza.go +++ b/pkg/messagebus/corteza.go @@ -18,7 +18,7 @@ type ( } CortezaQueueHandler struct { - handle handler + handle HandlerType messages CortezaMessageStore poll *time.Ticker } diff --git a/pkg/messagebus/handler.go b/pkg/messagebus/handler.go index 464e75989..7412bf671 100644 --- a/pkg/messagebus/handler.go +++ b/pkg/messagebus/handler.go @@ -6,14 +6,14 @@ import ( ) const ( - HandlerCorteza handler = "corteza" - HandlerNoop handler = "noop" - HandlerRedis handler = "redis" - HandlerSql handler = "sql" + HandlerCorteza HandlerType = "corteza" + HandlerNoop HandlerType = "noop" + HandlerRedis HandlerType = "redis" + HandlerSql HandlerType = "sql" ) type ( - handler string + HandlerType string Handler interface { ReadHandler @@ -53,3 +53,11 @@ type ( Process(context.Context, QueueMessage) error } ) + +func HandlerTypes() []HandlerType { + return []HandlerType{ + HandlerCorteza, + HandlerRedis, + HandlerSql, + } +} diff --git a/pkg/messagebus/messagebus.go b/pkg/messagebus/messagebus.go index 9fe740258..4bab85a7b 100644 --- a/pkg/messagebus/messagebus.go +++ b/pkg/messagebus/messagebus.go @@ -9,6 +9,7 @@ import ( "github.com/cortezaproject/corteza-server/pkg/eventbus" "github.com/cortezaproject/corteza-server/pkg/expr" + "github.com/cortezaproject/corteza-server/pkg/options" "github.com/cortezaproject/corteza-server/system/service/event" "go.uber.org/zap" ) @@ -20,6 +21,7 @@ var ( type ( messageBus struct { + opts *options.MessagebusOpt logger *zap.Logger eventbus dispatcher mutex sync.Mutex @@ -50,17 +52,19 @@ func Set(m *messageBus) { } // Setup handles the singleton service -func Setup(logger *zap.Logger, d dispatcher) { +func Setup(opts *options.MessagebusOpt, log *zap.Logger, d dispatcher) { if gMbus != nil { return } - gMbus = New(logger, d) + gMbus = New(opts, log, d) } -func New(logger *zap.Logger, d dispatcher) *messageBus { +func New(opts *options.MessagebusOpt, logger *zap.Logger, d dispatcher) *messageBus { + // force logger debug return &messageBus{ eventbus: d, + opts: opts, logger: logger, queues: QueueSet{}, w: make(chan bool), @@ -163,6 +167,7 @@ func (mb *messageBus) dispatchEvents(ctx context.Context, queues ...*Queue) { select { case p := <-q.dispatch: if !q.settings.CanDispatch() { + mb.logger.Debug("NOT dispatching queue with payload, dispatch disabled in queue settings", zap.String("queue", q.settings.Queue)) break } @@ -194,7 +199,7 @@ func (mb *messageBus) updateProcessedMessages(ctx context.Context, queues ...*Qu err := q.handler.Process(ctx, *message) if err != nil { - mb.logger.Debug("could not mark message as processed", zap.String("queue", message.Queue), zap.Error(err)) + mb.logger.Warn("could not mark message as processed", zap.String("queue", message.Queue), zap.Error(err)) } else { mb.logger.Debug("message processed", zap.String("queue", message.Queue)) } @@ -262,18 +267,18 @@ func (mb *messageBus) writeToQueue(ctx context.Context, q *Queue) { err := q.handler.Write(ctx, p) if err != nil { - mb.logger.Info("could not add message to queue", zap.String("queue", q.settings.Queue), zap.Error(err)) + mb.logger.Warn("could not add message to queue", zap.String("queue", q.settings.Queue), zap.Error(err)) break } mb.logger.Debug("wrote payload to queue", zap.String("queue", q.settings.Queue)) case <-q.err: - mb.logger.Info("closing message queue listener", zap.String("queue", q.settings.Queue), zap.Int("processed", len(q.processed))) + mb.logger.Debug("closing message queue listener", zap.String("queue", q.settings.Queue), zap.Int("processed", len(q.processed))) return case <-ctx.Done(): - mb.logger.Info("closing message queue listener", zap.String("queue", q.settings.Queue)) + mb.logger.Debug("closing message queue listener", zap.String("queue", q.settings.Queue)) return } } @@ -291,7 +296,7 @@ func (mb *messageBus) readFromQueue(ctx context.Context, q *Queue) { list, err := q.handler.Read(ctx) if err != nil { - mb.logger.Info("could not read message from queue", zap.String("queue", q.settings.Queue), zap.Error(err)) + mb.logger.Warn("could not read message from queue", zap.String("queue", q.settings.Queue), zap.Error(err)) break } @@ -311,11 +316,11 @@ func (mb *messageBus) readFromQueue(ctx context.Context, q *Queue) { } case <-q.err: - mb.logger.Info("closing message queue listener", zap.String("queue", q.settings.Queue), zap.Int("processed", len(q.processed))) + mb.logger.Debug("closing message queue listener", zap.String("queue", q.settings.Queue), zap.Int("processed", len(q.processed))) return case <-ctx.Done(): - mb.logger.Info("closing message queue listener", zap.String("queue", q.settings.Queue), zap.Int("processed", len(q.processed))) + mb.logger.Debug("closing message queue listener", zap.String("queue", q.settings.Queue), zap.Int("processed", len(q.processed))) return } } @@ -337,7 +342,7 @@ func (mb *messageBus) initQueues(ctx context.Context, storer QueueStorer) error mb.logger.Debug("initializing queue", zap.String("queue", qs.Queue)) if err != nil { - mb.logger.Info("could not init handler for queue", zap.String("queue", qs.Queue), zap.Error(err)) + mb.logger.Warn("could not init handler for queue", zap.String("queue", qs.Queue), zap.Error(err)) continue } diff --git a/pkg/messagebus/redis.go b/pkg/messagebus/redis.go index 7f5bb83a4..af354c6d7 100644 --- a/pkg/messagebus/redis.go +++ b/pkg/messagebus/redis.go @@ -10,7 +10,7 @@ import ( type ( RedisQueueHandler struct { queue string - handle handler + handle HandlerType client *redis.Client poll *time.Ticker } diff --git a/pkg/messagebus/settings.go b/pkg/messagebus/settings.go index f9ba02584..464d54428 100644 --- a/pkg/messagebus/settings.go +++ b/pkg/messagebus/settings.go @@ -18,10 +18,10 @@ type ( } QueueSettings struct { - ID uint64 - Handler string - Queue string - Meta QueueSettingsMeta + ID uint64 `json:"queueID,string"` + Handler string `json:"handler"` + Queue string `json:"queue"` + Meta QueueSettingsMeta `json:"meta"` CreatedAt time.Time `json:"createdAt,omitempty"` CreatedBy uint64 `json:"createdBy,string" ` @@ -32,6 +32,11 @@ type ( } QueueSettingsFilter struct { + Queue string `json:"queue"` + Handler HandlerType `json:"handler"` + + Deleted filter.State `json:"deleted"` + filter.Sorting filter.Paging } @@ -66,11 +71,17 @@ func (m QueueSettingsMeta) Value() (driver.Value, error) { } func (m QueueSettingsMeta) MarshalJSON() ([]byte, error) { + + pollDelay := "" + if m.PollDelay != nil { + pollDelay = m.PollDelay.String() + } + return json.Marshal(struct { PollDelay string `json:"poll_delay"` DispatchEvents bool `json:"dispatch_events,omitempty"` }{ - PollDelay: m.PollDelay.String(), + PollDelay: pollDelay, DispatchEvents: m.DispatchEvents, }) } diff --git a/pkg/messagebus/sql.go b/pkg/messagebus/sql.go index efb6f96dc..4041c6bbe 100644 --- a/pkg/messagebus/sql.go +++ b/pkg/messagebus/sql.go @@ -8,7 +8,7 @@ import ( type ( SqlQueueHandler struct { queue string - handle handler + handle HandlerType client SqlClient poll *time.Ticker } diff --git a/pkg/options/messagebus.gen.go b/pkg/options/messagebus.gen.go new file mode 100644 index 000000000..5510fbe96 --- /dev/null +++ b/pkg/options/messagebus.gen.go @@ -0,0 +1,37 @@ +package options + +// This file is auto-generated. +// +// Changes to this file may cause incorrect behavior and will be lost if +// the code is regenerated. +// +// Definitions file that controls how this file is generated: +// pkg/options/messagebus.yaml + +type ( + MessagebusOpt struct { + Enabled bool `env:"MESSAGEBUS_ENABLED"` + LogEnabled bool `env:"MESSAGEBUS_LOG_ENABLED"` + } +) + +// Messagebus initializes and returns a MessagebusOpt with default values +func Messagebus() (o *MessagebusOpt) { + o = &MessagebusOpt{ + Enabled: true, + LogEnabled: false, + } + + fill(o) + + // Function that allows access to custom logic inside the parent function. + // The custom logic in the other file should be like: + // func (o *Messagebus) Defaults() {...} + func(o interface{}) { + if def, ok := o.(interface{ Defaults() }); ok { + def.Defaults() + } + }(o) + + return +} diff --git a/pkg/options/messagebus.yaml b/pkg/options/messagebus.yaml new file mode 100644 index 000000000..89ff0232b --- /dev/null +++ b/pkg/options/messagebus.yaml @@ -0,0 +1,15 @@ +docs: + title: Messaging queue + +props: + - name: Enabled + type: bool + default: true + description: |- + Enable messagebus + + - name: logEnabled + type: bool + default: false + description: |- + Enable extra logging for messagebus watchers diff --git a/store/messagebus_queuesettings.gen.go b/store/messagebus_queuesettings.gen.go index 9c03c1364..00932d558 100644 --- a/store/messagebus_queuesettings.gen.go +++ b/store/messagebus_queuesettings.gen.go @@ -16,6 +16,8 @@ import ( type ( MessagebusQueuesettings interface { SearchMessagebusQueuesettings(ctx context.Context, f messagebus.QueueSettingsFilter) (messagebus.QueueSettingsSet, messagebus.QueueSettingsFilter, error) + LookupMessagebusQueuesettingByID(ctx context.Context, id uint64) (*messagebus.QueueSettings, error) + LookupMessagebusQueuesettingByQueue(ctx context.Context, queue string) (*messagebus.QueueSettings, error) CreateMessagebusQueuesetting(ctx context.Context, rr ...*messagebus.QueueSettings) error @@ -38,6 +40,16 @@ func SearchMessagebusQueuesettings(ctx context.Context, s MessagebusQueuesetting return s.SearchMessagebusQueuesettings(ctx, f) } +// LookupMessagebusQueuesettingByID searches for queue by ID +func LookupMessagebusQueuesettingByID(ctx context.Context, s MessagebusQueuesettings, id uint64) (*messagebus.QueueSettings, error) { + return s.LookupMessagebusQueuesettingByID(ctx, id) +} + +// LookupMessagebusQueuesettingByQueue searches for queue by queue name +func LookupMessagebusQueuesettingByQueue(ctx context.Context, s MessagebusQueuesettings, queue string) (*messagebus.QueueSettings, error) { + return s.LookupMessagebusQueuesettingByQueue(ctx, queue) +} + // CreateMessagebusQueuesetting creates one or more MessagebusQueuesettings in store func CreateMessagebusQueuesetting(ctx context.Context, s MessagebusQueuesettings, rr ...*messagebus.QueueSettings) error { return s.CreateMessagebusQueuesetting(ctx, rr...) diff --git a/store/messagebus_queuesettings.yaml b/store/messagebus_queuesettings.yaml index 16ced1129..16771e6dd 100644 --- a/store/messagebus_queuesettings.yaml +++ b/store/messagebus_queuesettings.yaml @@ -18,11 +18,18 @@ fields: - { field: UpdatedAt } - { field: DeletedAt } +lookups: + - fields: [ ID ] + description: |- + searches for queue by ID + - fields: [ Queue ] + description: |- + searches for queue by queue name rdbms: alias: mqs table: queue_settings - customFilterConverter: false + customFilterConverter: true search: enablePaging: true diff --git a/store/rdbms/messagebus_queuesettings.gen.go b/store/rdbms/messagebus_queuesettings.gen.go index 07fecbe87..7b0afd497 100644 --- a/store/rdbms/messagebus_queuesettings.gen.go +++ b/store/rdbms/messagebus_queuesettings.gen.go @@ -33,7 +33,10 @@ func (s Store) SearchMessagebusQueuesettings(ctx context.Context, f messagebus.Q ) return set, f, func() error { - q = s.messagebusQueuesettingsSelectBuilder() + q, err = s.convertMessagebusQueuesettingFilter(f) + if err != nil { + return err + } // Paging enabled // {search: {enablePaging:true}} @@ -257,6 +260,20 @@ func (s Store) QueryMessagebusQueuesettings( return set, rows.Err() } +// LookupMessagebusQueuesettingByID searches for queue by ID +func (s Store) LookupMessagebusQueuesettingByID(ctx context.Context, id uint64) (*messagebus.QueueSettings, error) { + return s.execLookupMessagebusQueuesetting(ctx, squirrel.Eq{ + s.preprocessColumn("mqs.id", ""): store.PreprocessValue(id, ""), + }) +} + +// LookupMessagebusQueuesettingByQueue searches for queue by queue name +func (s Store) LookupMessagebusQueuesettingByQueue(ctx context.Context, queue string) (*messagebus.QueueSettings, error) { + return s.execLookupMessagebusQueuesetting(ctx, squirrel.Eq{ + s.preprocessColumn("mqs.queue", ""): store.PreprocessValue(queue, ""), + }) +} + // CreateMessagebusQueuesetting creates one or more rows in queue_settings table func (s Store) CreateMessagebusQueuesetting(ctx context.Context, rr ...*messagebus.QueueSettings) (err error) { for _, res := range rr { diff --git a/store/rdbms/messagebus_queuesettings.go b/store/rdbms/messagebus_queuesettings.go index 8204315b2..b96e6e7bd 100644 --- a/store/rdbms/messagebus_queuesettings.go +++ b/store/rdbms/messagebus_queuesettings.go @@ -2,10 +2,12 @@ package rdbms import ( "github.com/Masterminds/squirrel" - "github.com/cortezaproject/corteza-server/federation/types" + "github.com/cortezaproject/corteza-server/pkg/filter" + "github.com/cortezaproject/corteza-server/pkg/messagebus" ) -func (s Store) convertMessagebusQueuesettingsFilter(f types.ExposedModuleFilter) (query squirrel.SelectBuilder, err error) { +func (s Store) convertMessagebusQueuesettingFilter(f messagebus.QueueSettingsFilter) (query squirrel.SelectBuilder, err error) { query = s.messagebusQueuesettingsSelectBuilder() + query = filter.StateCondition(query, "mqs.deleted_at", f.Deleted) return } diff --git a/system/rest.yaml b/system/rest.yaml index ca323a44e..b4990bd58 100644 --- a/system/rest.yaml +++ b/system/rest.yaml @@ -1474,3 +1474,97 @@ endpoints: - type: uint name: limit title: Limit + +- title: Messaging queues + entrypoint: queues + path: "/queues" + imports: + - sqlxTypes github.com/jmoiron/sqlx/types + apis: + - name: list + method: GET + title: Messaging queues + path: "/" + parameters: + get: + - type: string + name: query + required: false + title: Search query + - type: uint + name: limit + required: false + title: Limit + - type: string + name: pageCursor + required: false + title: Page cursor + - type: string + name: sort + required: false + title: Sort items + - name: deleted + required: false + title: Exclude (0, default), include (1) or return only (2) deleted queues + type: uint + - name: create + method: PUT + title: Create messaging queue + path: "" + parameters: + post: + - type: string + name: queue + required: true + title: Name of queue + - type: string + name: handler + required: true + title: Queue handler + - type: sqlxTypes.JSONText + name: meta + required: false + title: Meta data for queue + - name: read + method: GET + title: Messaging queue details + path: "/{queueID}" + parameters: + path: + - type: uint64 + name: queueID + required: true + title: Queue ID + - name: update + method: POST + title: Update role details + path: "/{queueID}" + parameters: + path: + - type: uint64 + name: queueID + required: true + title: Queue ID + post: + - type: string + name: queue + required: true + title: Name of queue + - type: string + name: handler + required: true + title: Queue handler + - type: sqlxTypes.JSONText + name: meta + required: false + title: Meta data for queue + - name: delete + method: DELETE + title: Messaging queue delete + path: "/{queueID}" + parameters: + path: + - type: uint64 + name: queueID + required: true + title: Queue ID diff --git a/system/rest/handlers/queues.go b/system/rest/handlers/queues.go new file mode 100644 index 000000000..6ba4b4de1 --- /dev/null +++ b/system/rest/handlers/queues.go @@ -0,0 +1,133 @@ +package handlers + +// This file is auto-generated. +// +// Changes to this file may cause incorrect behavior and will be lost if +// the code is regenerated. +// +// Definitions file that controls how this file is generated: +// + +import ( + "context" + "github.com/cortezaproject/corteza-server/pkg/api" + "github.com/cortezaproject/corteza-server/system/rest/request" + "github.com/go-chi/chi" + "net/http" +) + +type ( + // Internal API interface + QueuesAPI interface { + List(context.Context, *request.QueuesList) (interface{}, error) + Create(context.Context, *request.QueuesCreate) (interface{}, error) + Read(context.Context, *request.QueuesRead) (interface{}, error) + Update(context.Context, *request.QueuesUpdate) (interface{}, error) + Delete(context.Context, *request.QueuesDelete) (interface{}, error) + } + + // HTTP API interface + Queues struct { + List func(http.ResponseWriter, *http.Request) + Create func(http.ResponseWriter, *http.Request) + Read func(http.ResponseWriter, *http.Request) + Update func(http.ResponseWriter, *http.Request) + Delete func(http.ResponseWriter, *http.Request) + } +) + +func NewQueues(h QueuesAPI) *Queues { + return &Queues{ + List: func(w http.ResponseWriter, r *http.Request) { + defer r.Body.Close() + params := request.NewQueuesList() + if err := params.Fill(r); err != nil { + api.Send(w, r, err) + return + } + + value, err := h.List(r.Context(), params) + if err != nil { + api.Send(w, r, err) + return + } + + api.Send(w, r, value) + }, + Create: func(w http.ResponseWriter, r *http.Request) { + defer r.Body.Close() + params := request.NewQueuesCreate() + if err := params.Fill(r); err != nil { + api.Send(w, r, err) + return + } + + value, err := h.Create(r.Context(), params) + if err != nil { + api.Send(w, r, err) + return + } + + api.Send(w, r, value) + }, + Read: func(w http.ResponseWriter, r *http.Request) { + defer r.Body.Close() + params := request.NewQueuesRead() + if err := params.Fill(r); err != nil { + api.Send(w, r, err) + return + } + + value, err := h.Read(r.Context(), params) + if err != nil { + api.Send(w, r, err) + return + } + + api.Send(w, r, value) + }, + Update: func(w http.ResponseWriter, r *http.Request) { + defer r.Body.Close() + params := request.NewQueuesUpdate() + if err := params.Fill(r); err != nil { + api.Send(w, r, err) + return + } + + value, err := h.Update(r.Context(), params) + if err != nil { + api.Send(w, r, err) + return + } + + api.Send(w, r, value) + }, + Delete: func(w http.ResponseWriter, r *http.Request) { + defer r.Body.Close() + params := request.NewQueuesDelete() + if err := params.Fill(r); err != nil { + api.Send(w, r, err) + return + } + + value, err := h.Delete(r.Context(), params) + if err != nil { + api.Send(w, r, err) + return + } + + api.Send(w, r, value) + }, + } +} + +func (h Queues) MountRoutes(r chi.Router, middlewares ...func(http.Handler) http.Handler) { + r.Group(func(r chi.Router) { + r.Use(middlewares...) + r.Get("/queues/", h.List) + r.Put("/queues", h.Create) + r.Get("/queues/{queueID}", h.Read) + r.Post("/queues/{queueID}", h.Update) + r.Delete("/queues/{queueID}", h.Delete) + }) +} diff --git a/system/rest/queues.go b/system/rest/queues.go new file mode 100644 index 000000000..a67fe30a9 --- /dev/null +++ b/system/rest/queues.go @@ -0,0 +1,125 @@ +package rest + +import ( + "context" + + "github.com/cortezaproject/corteza-server/pkg/api" + "github.com/cortezaproject/corteza-server/pkg/filter" + "github.com/cortezaproject/corteza-server/pkg/messagebus" + "github.com/cortezaproject/corteza-server/system/rest/request" + "github.com/cortezaproject/corteza-server/system/service" + "github.com/davecgh/go-spew/spew" +) + +type ( + Queue struct { + svc service.QueueService + ac templateAccessController + } + + queuePayload struct { + *messagebus.QueueSettings + } + + queueSetPayload struct { + Filter messagebus.QueueSettingsFilter `json:"filter"` + Set []*queuePayload `json:"set"` + } +) + +func (Queue) New() *Queue { + return &Queue{ + svc: service.DefaultQueue, + ac: service.DefaultAccessControl, + } +} + +func (ctrl *Queue) List(ctx context.Context, r *request.QueuesList) (interface{}, error) { + var ( + err error + f = messagebus.QueueSettingsFilter{ + Deleted: filter.State(r.Deleted), + } + ) + + if f.Paging, err = filter.NewPaging(r.Limit, r.PageCursor); err != nil { + return nil, err + } + + if f.Sorting, err = filter.NewSorting(r.Sort); err != nil { + return nil, err + } + + set, filter, err := ctrl.svc.Search(ctx, f) + return ctrl.makeFilterPayload(ctx, set, filter, err) +} + +func (ctrl *Queue) Create(ctx context.Context, r *request.QueuesCreate) (interface{}, error) { + var ( + err error + q = &messagebus.QueueSettings{ + Handler: r.Handler, + Queue: r.Queue, + // Meta: r.Meta, + } + ) + + q, err = ctrl.svc.Create(ctx, q) + + return ctrl.makePayload(ctx, q, err) +} + +func (ctrl *Queue) Read(ctx context.Context, r *request.QueuesRead) (interface{}, error) { + return ctrl.svc.FindByID(ctx, r.QueueID) +} + +func (ctrl *Queue) Update(ctx context.Context, r *request.QueuesUpdate) (interface{}, error) { + var ( + err error + q = &messagebus.QueueSettings{ + ID: r.QueueID, + Handler: r.Handler, + Queue: r.Queue, + // Meta: r.Meta, + } + ) + + q, err = ctrl.svc.Update(ctx, q) + + return ctrl.makePayload(ctx, q, err) +} + +func (ctrl *Queue) Delete(ctx context.Context, r *request.QueuesDelete) (interface{}, error) { + spew.Dump("DELETE", r) + return api.OK(), ctrl.svc.DeleteByID(ctx, r.QueueID) +} + +func (ctrl *Queue) Undelete(ctx context.Context, r *request.TemplateUndelete) (interface{}, error) { + return api.OK(), ctrl.svc.UndeleteByID(ctx, r.TemplateID) +} + +func (ctrl *Queue) makePayload(ctx context.Context, q *messagebus.QueueSettings, err error) (*queuePayload, error) { + if err != nil || q == nil { + return nil, err + } + + qq := &queuePayload{ + QueueSettings: q, + } + + return qq, nil +} + +func (ctrl *Queue) makeFilterPayload(ctx context.Context, nn messagebus.QueueSettingsSet, f messagebus.QueueSettingsFilter, err error) (*queueSetPayload, error) { + if err != nil { + return nil, err + } + + msp := &queueSetPayload{Filter: f, Set: make([]*queuePayload, len(nn))} + + for i := range nn { + msp.Set[i], _ = ctrl.makePayload(ctx, nn[i], nil) + } + + return msp, nil +} diff --git a/system/rest/request/queues.go b/system/rest/request/queues.go new file mode 100644 index 000000000..346ada426 --- /dev/null +++ b/system/rest/request/queues.go @@ -0,0 +1,437 @@ +package request + +// This file is auto-generated. +// +// Changes to this file may cause incorrect behavior and will be lost if +// the code is regenerated. +// +// Definitions file that controls how this file is generated: +// + +import ( + "encoding/json" + "fmt" + "github.com/cortezaproject/corteza-server/pkg/payload" + "github.com/go-chi/chi" + sqlxTypes "github.com/jmoiron/sqlx/types" + "io" + "mime/multipart" + "net/http" + "strings" +) + +// dummy vars to prevent +// unused imports complain +var ( + _ = chi.URLParam + _ = multipart.ErrMessageTooLarge + _ = payload.ParseUint64s + _ = strings.ToLower + _ = io.EOF + _ = fmt.Errorf + _ = json.NewEncoder +) + +type ( + // Internal API interface + QueuesList struct { + // Query GET parameter + // + // Search query + Query string + + // Limit GET parameter + // + // Limit + Limit uint + + // PageCursor GET parameter + // + // Page cursor + PageCursor string + + // Sort GET parameter + // + // Sort items + Sort string + + // Deleted GET parameter + // + // Exclude (0, default), include (1) or return only (2) deleted queues + Deleted uint + } + + QueuesCreate struct { + // Queue POST parameter + // + // Name of queue + Queue string + + // Handler POST parameter + // + // Queue handler + Handler string + + // Meta POST parameter + // + // Meta data for queue + Meta sqlxTypes.JSONText + } + + QueuesRead struct { + // QueueID PATH parameter + // + // Queue ID + QueueID uint64 `json:",string"` + } + + QueuesUpdate struct { + // QueueID PATH parameter + // + // Queue ID + QueueID uint64 `json:",string"` + + // Queue POST parameter + // + // Name of queue + Queue string + + // Handler POST parameter + // + // Queue handler + Handler string + + // Meta POST parameter + // + // Meta data for queue + Meta sqlxTypes.JSONText + } + + QueuesDelete struct { + // QueueID PATH parameter + // + // Queue ID + QueueID uint64 `json:",string"` + } +) + +// NewQueuesList request +func NewQueuesList() *QueuesList { + return &QueuesList{} +} + +// Auditable returns all auditable/loggable parameters +func (r QueuesList) Auditable() map[string]interface{} { + return map[string]interface{}{ + "query": r.Query, + "limit": r.Limit, + "pageCursor": r.PageCursor, + "sort": r.Sort, + "deleted": r.Deleted, + } +} + +// Auditable returns all auditable/loggable parameters +func (r QueuesList) GetQuery() string { + return r.Query +} + +// Auditable returns all auditable/loggable parameters +func (r QueuesList) GetLimit() uint { + return r.Limit +} + +// Auditable returns all auditable/loggable parameters +func (r QueuesList) GetPageCursor() string { + return r.PageCursor +} + +// Auditable returns all auditable/loggable parameters +func (r QueuesList) GetSort() string { + return r.Sort +} + +// Auditable returns all auditable/loggable parameters +func (r QueuesList) GetDeleted() uint { + return r.Deleted +} + +// Fill processes request and fills internal variables +func (r *QueuesList) Fill(req *http.Request) (err error) { + + { + // GET params + tmp := req.URL.Query() + + if val, ok := tmp["query"]; ok && len(val) > 0 { + r.Query, err = val[0], nil + if err != nil { + return err + } + } + if val, ok := tmp["limit"]; ok && len(val) > 0 { + r.Limit, err = payload.ParseUint(val[0]), nil + if err != nil { + return err + } + } + if val, ok := tmp["pageCursor"]; ok && len(val) > 0 { + r.PageCursor, err = val[0], nil + if err != nil { + return err + } + } + if val, ok := tmp["sort"]; ok && len(val) > 0 { + r.Sort, err = val[0], nil + if err != nil { + return err + } + } + if val, ok := tmp["deleted"]; ok && len(val) > 0 { + r.Deleted, err = payload.ParseUint(val[0]), nil + if err != nil { + return err + } + } + } + + return err +} + +// NewQueuesCreate request +func NewQueuesCreate() *QueuesCreate { + return &QueuesCreate{} +} + +// Auditable returns all auditable/loggable parameters +func (r QueuesCreate) Auditable() map[string]interface{} { + return map[string]interface{}{ + "queue": r.Queue, + "handler": r.Handler, + "meta": r.Meta, + } +} + +// Auditable returns all auditable/loggable parameters +func (r QueuesCreate) GetQueue() string { + return r.Queue +} + +// Auditable returns all auditable/loggable parameters +func (r QueuesCreate) GetHandler() string { + return r.Handler +} + +// Auditable returns all auditable/loggable parameters +func (r QueuesCreate) GetMeta() sqlxTypes.JSONText { + return r.Meta +} + +// Fill processes request and fills internal variables +func (r *QueuesCreate) Fill(req *http.Request) (err error) { + + if strings.ToLower(req.Header.Get("content-type")) == "application/json" { + err = json.NewDecoder(req.Body).Decode(r) + + switch { + case err == io.EOF: + err = nil + case err != nil: + return fmt.Errorf("error parsing http request body: %w", err) + } + } + + { + if err = req.ParseForm(); err != nil { + return err + } + + // POST params + + if val, ok := req.Form["queue"]; ok && len(val) > 0 { + r.Queue, err = val[0], nil + if err != nil { + return err + } + } + + if val, ok := req.Form["handler"]; ok && len(val) > 0 { + r.Handler, err = val[0], nil + if err != nil { + return err + } + } + + if val, ok := req.Form["meta"]; ok && len(val) > 0 { + r.Meta, err = payload.ParseJSONTextWithErr(val[0]) + if err != nil { + return err + } + } + } + + return err +} + +// NewQueuesRead request +func NewQueuesRead() *QueuesRead { + return &QueuesRead{} +} + +// Auditable returns all auditable/loggable parameters +func (r QueuesRead) Auditable() map[string]interface{} { + return map[string]interface{}{ + "queueID": r.QueueID, + } +} + +// Auditable returns all auditable/loggable parameters +func (r QueuesRead) GetQueueID() uint64 { + return r.QueueID +} + +// Fill processes request and fills internal variables +func (r *QueuesRead) Fill(req *http.Request) (err error) { + + { + var val string + // path params + + val = chi.URLParam(req, "queueID") + r.QueueID, err = payload.ParseUint64(val), nil + if err != nil { + return err + } + + } + + return err +} + +// NewQueuesUpdate request +func NewQueuesUpdate() *QueuesUpdate { + return &QueuesUpdate{} +} + +// Auditable returns all auditable/loggable parameters +func (r QueuesUpdate) Auditable() map[string]interface{} { + return map[string]interface{}{ + "queueID": r.QueueID, + "queue": r.Queue, + "handler": r.Handler, + "meta": r.Meta, + } +} + +// Auditable returns all auditable/loggable parameters +func (r QueuesUpdate) GetQueueID() uint64 { + return r.QueueID +} + +// Auditable returns all auditable/loggable parameters +func (r QueuesUpdate) GetQueue() string { + return r.Queue +} + +// Auditable returns all auditable/loggable parameters +func (r QueuesUpdate) GetHandler() string { + return r.Handler +} + +// Auditable returns all auditable/loggable parameters +func (r QueuesUpdate) GetMeta() sqlxTypes.JSONText { + return r.Meta +} + +// Fill processes request and fills internal variables +func (r *QueuesUpdate) Fill(req *http.Request) (err error) { + + if strings.ToLower(req.Header.Get("content-type")) == "application/json" { + err = json.NewDecoder(req.Body).Decode(r) + + switch { + case err == io.EOF: + err = nil + case err != nil: + return fmt.Errorf("error parsing http request body: %w", err) + } + } + + { + if err = req.ParseForm(); err != nil { + return err + } + + // POST params + + if val, ok := req.Form["queue"]; ok && len(val) > 0 { + r.Queue, err = val[0], nil + if err != nil { + return err + } + } + + if val, ok := req.Form["handler"]; ok && len(val) > 0 { + r.Handler, err = val[0], nil + if err != nil { + return err + } + } + + if val, ok := req.Form["meta"]; ok && len(val) > 0 { + r.Meta, err = payload.ParseJSONTextWithErr(val[0]) + if err != nil { + return err + } + } + } + + { + var val string + // path params + + val = chi.URLParam(req, "queueID") + r.QueueID, err = payload.ParseUint64(val), nil + if err != nil { + return err + } + + } + + return err +} + +// NewQueuesDelete request +func NewQueuesDelete() *QueuesDelete { + return &QueuesDelete{} +} + +// Auditable returns all auditable/loggable parameters +func (r QueuesDelete) Auditable() map[string]interface{} { + return map[string]interface{}{ + "queueID": r.QueueID, + } +} + +// Auditable returns all auditable/loggable parameters +func (r QueuesDelete) GetQueueID() uint64 { + return r.QueueID +} + +// Fill processes request and fills internal variables +func (r *QueuesDelete) Fill(req *http.Request) (err error) { + + { + var val string + // path params + + val = chi.URLParam(req, "queueID") + r.QueueID, err = payload.ParseUint64(val), nil + if err != nil { + return err + } + + } + + return err +} diff --git a/system/rest/router.go b/system/rest/router.go index 13403c419..2b70cd976 100644 --- a/system/rest/router.go +++ b/system/rest/router.go @@ -36,5 +36,6 @@ func MountRoutes(r chi.Router) { handlers.NewStats(Stats{}.New()).MountRoutes(r) handlers.NewReminder(Reminder{}.New()).MountRoutes(r) handlers.NewActionlog(Actionlog{}.New()).MountRoutes(r) + handlers.NewQueues(Queue{}.New()).MountRoutes(r) }) } diff --git a/system/service/queue.go b/system/service/queue.go new file mode 100644 index 000000000..24483c20f --- /dev/null +++ b/system/service/queue.go @@ -0,0 +1,210 @@ +package service + +import ( + "context" + + "github.com/cortezaproject/corteza-server/pkg/actionlog" + "github.com/cortezaproject/corteza-server/pkg/messagebus" + "github.com/cortezaproject/corteza-server/store" +) + +type ( + queue struct { + actionlog actionlog.Recorder + store store.Storer + ac templateAccessController + } + + QueueService interface { + FindByID(ctx context.Context, ID uint64) (*messagebus.QueueSettings, error) + Search(context.Context, messagebus.QueueSettingsFilter) (messagebus.QueueSettingsSet, messagebus.QueueSettingsFilter, error) + + Create(ctx context.Context, q *messagebus.QueueSettings) (*messagebus.QueueSettings, error) + Update(ctx context.Context, q *messagebus.QueueSettings) (*messagebus.QueueSettings, error) + + DeleteByID(ctx context.Context, ID uint64) error + UndeleteByID(ctx context.Context, ID uint64) error + } +) + +func Queue() QueueService { + return (&queue{ + ac: DefaultAccessControl, + actionlog: DefaultActionlog, + store: DefaultStore, + }) +} + +func (svc *queue) FindByID(ctx context.Context, ID uint64) (q *messagebus.QueueSettings, err error) { + var ( + qProps = &queueActionProps{} + ) + + err = func() error { + if ID == 0 { + return QueueErrInvalidID() + } + + if q, err = store.LookupMessagebusQueuesettingByID(ctx, svc.store, ID); err != nil { + return TemplateErrInvalidID().Wrap(err) + } + + qProps.setQueue(q) + + // if !svc.ac.CanReadTemplate(ctx, tpl) { + // return TemplateErrNotAllowedToRead() + // } + + return nil + }() + + return q, svc.recordAction(ctx, qProps, QueueActionLookup, err) +} + +func (svc *queue) Create(ctx context.Context, new *messagebus.QueueSettings) (q *messagebus.QueueSettings, err error) { + var ( + qProps = &queueActionProps{new: new} + ) + + err = func() (err error) { + if !svc.isValidHandler(messagebus.HandlerType(new.Handler)) { + return QueueErrInvalidHandler(qProps) + } + + // Set new values after beforeCreate events are emitted + new.ID = nextID() + new.CreatedAt = *now() + + if err = store.CreateMessagebusQueuesetting(ctx, svc.store, new); err != nil { + return + } + + q = new + + return nil + }() + + return q, svc.recordAction(ctx, qProps, QueueActionCreate, err) +} + +func (svc *queue) Update(ctx context.Context, upd *messagebus.QueueSettings) (q *messagebus.QueueSettings, err error) { + var ( + qProps = &queueActionProps{update: upd} + ) + + err = func() (err error) { + if _, e := store.LookupMessagebusQueuesettingByID(ctx, svc.store, upd.ID); e != nil { + return QueueErrNotFound(qProps) + } + + if qq, e := store.LookupMessagebusQueuesettingByQueue(ctx, svc.store, upd.Queue); e == nil && qq != nil && qq.ID != upd.ID { + return QueueErrAlreadyExists(qProps) + } + + if !svc.isValidHandler(messagebus.HandlerType(upd.Handler)) { + return QueueErrInvalidHandler(qProps) + } + + // Set new values after beforeCreate events are emitted + upd.UpdatedAt = now() + + if err = store.UpdateMessagebusQueuesetting(ctx, svc.store, upd); err != nil { + return + } + + q = upd + + return nil + }() + + return q, svc.recordAction(ctx, qProps, QueueActionUpdate, err) +} + +func (svc *queue) DeleteByID(ctx context.Context, ID uint64) (err error) { + var ( + qProps = &queueActionProps{} + q *messagebus.QueueSettings + ) + + err = func() (err error) { + if ID == 0 { + return QueueErrInvalidID() + } + + if q, err = store.LookupMessagebusQueuesettingByID(ctx, svc.store, ID); err != nil { + return + } + + qProps.setQueue(q) + + // if !svc.ac.CanDeleteTemplate(ctx, tpl) { + // return TemplateErrNotAllowedToDelete() + // } + + q.DeletedAt = now() + if err = store.UpdateMessagebusQueuesetting(ctx, svc.store, q); err != nil { + return + } + + return nil + }() + + return svc.recordAction(ctx, qProps, QueueActionDelete, err) +} + +func (svc *queue) UndeleteByID(ctx context.Context, ID uint64) (err error) { + var ( + qProps = &queueActionProps{} + q *messagebus.QueueSettings + ) + + err = func() (err error) { + if ID == 0 { + return QueueErrInvalidID() + } + + if q, err = store.LookupMessagebusQueuesettingByID(ctx, svc.store, ID); err != nil { + return + } + + qProps.setQueue(q) + + // if !svc.ac.CanDeleteTemplate(ctx, tpl) { + // return TemplateErrNotAllowedToDelete() + // } + + q.DeletedAt = nil + if err = store.UpdateMessagebusQueuesetting(ctx, svc.store, q); err != nil { + return + } + + return nil + }() + + return svc.recordAction(ctx, qProps, QueueActionDelete, err) +} + +func (svc *queue) Search(ctx context.Context, filter messagebus.QueueSettingsFilter) (q messagebus.QueueSettingsSet, f messagebus.QueueSettingsFilter, err error) { + var ( + aProps = &queueActionProps{search: &filter} + ) + + err = func() error { + if q, f, err = store.SearchMessagebusQueuesettings(ctx, svc.store, filter); err != nil { + return err + } + + return nil + }() + + return q, f, svc.recordAction(ctx, aProps, QueueActionSearch, err) +} + +func (svc *queue) isValidHandler(h messagebus.HandlerType) bool { + for _, hh := range messagebus.HandlerTypes() { + if h == hh { + return true + } + } + return false +} diff --git a/system/service/queue_actions.gen.go b/system/service/queue_actions.gen.go new file mode 100644 index 000000000..9f8651c28 --- /dev/null +++ b/system/service/queue_actions.gen.go @@ -0,0 +1,592 @@ +package service + +// This file is auto-generated. +// +// Changes to this file may cause incorrect behavior and will be lost if +// the code is regenerated. +// +// Definitions file that controls how this file is generated: +// system/service/queue_actions.yaml + +import ( + "context" + "fmt" + "github.com/cortezaproject/corteza-server/pkg/actionlog" + "github.com/cortezaproject/corteza-server/pkg/errors" + "github.com/cortezaproject/corteza-server/pkg/messagebus" + "strings" + "time" +) + +type ( + queueActionProps struct { + queue *messagebus.QueueSettings + new *messagebus.QueueSettings + update *messagebus.QueueSettings + search *messagebus.QueueSettingsFilter + } + + queueAction struct { + timestamp time.Time + resource string + action string + log string + severity actionlog.Severity + + // prefix for error when action fails + errorMessage string + + props *queueActionProps + } + + queueLogMetaKey struct{} + queuePropsMetaKey struct{} +) + +var ( + // just a placeholder to cover template cases w/o fmt package use + _ = fmt.Println +) + +// ********************************************************************************************************************* +// ********************************************************************************************************************* +// Props methods +// setQueue updates queueActionProps's queue +// +// Allows method chaining +// +// This function is auto-generated. +// +func (p *queueActionProps) setQueue(queue *messagebus.QueueSettings) *queueActionProps { + p.queue = queue + return p +} + +// setNew updates queueActionProps's new +// +// Allows method chaining +// +// This function is auto-generated. +// +func (p *queueActionProps) setNew(new *messagebus.QueueSettings) *queueActionProps { + p.new = new + return p +} + +// setUpdate updates queueActionProps's update +// +// Allows method chaining +// +// This function is auto-generated. +// +func (p *queueActionProps) setUpdate(update *messagebus.QueueSettings) *queueActionProps { + p.update = update + return p +} + +// setSearch updates queueActionProps's search +// +// Allows method chaining +// +// This function is auto-generated. +// +func (p *queueActionProps) setSearch(search *messagebus.QueueSettingsFilter) *queueActionProps { + p.search = search + return p +} + +// Serialize converts queueActionProps to actionlog.Meta +// +// This function is auto-generated. +// +func (p queueActionProps) Serialize() actionlog.Meta { + var ( + m = make(actionlog.Meta) + ) + + if p.queue != nil { + m.Set("queue.queue", p.queue.Queue, true) + m.Set("queue.ID", p.queue.ID, true) + } + if p.new != nil { + m.Set("new.queue", p.new.Queue, true) + m.Set("new.handler", p.new.Handler, true) + m.Set("new.ID", p.new.ID, true) + } + if p.update != nil { + m.Set("update.queue", p.update.Queue, true) + m.Set("update.handler", p.update.Handler, true) + m.Set("update.ID", p.update.ID, true) + } + if p.search != nil { + } + + return m +} + +// tr translates string and replaces meta value placeholder with values +// +// This function is auto-generated. +// +func (p queueActionProps) Format(in string, err error) string { + var ( + pairs = []string{"{err}"} + // first non-empty string + fns = func(ii ...interface{}) string { + for _, i := range ii { + if s := fmt.Sprintf("%v", i); len(s) > 0 { + return s + } + } + + return "" + } + ) + + if err != nil { + pairs = append(pairs, err.Error()) + } else { + pairs = append(pairs, "nil") + } + + if p.queue != nil { + // replacement for "{queue}" (in order how fields are defined) + pairs = append( + pairs, + "{queue}", + fns( + p.queue.Queue, + p.queue.ID, + ), + ) + pairs = append(pairs, "{queue.queue}", fns(p.queue.Queue)) + pairs = append(pairs, "{queue.ID}", fns(p.queue.ID)) + } + + if p.new != nil { + // replacement for "{new}" (in order how fields are defined) + pairs = append( + pairs, + "{new}", + fns( + p.new.Queue, + p.new.Handler, + p.new.ID, + ), + ) + pairs = append(pairs, "{new.queue}", fns(p.new.Queue)) + pairs = append(pairs, "{new.handler}", fns(p.new.Handler)) + pairs = append(pairs, "{new.ID}", fns(p.new.ID)) + } + + if p.update != nil { + // replacement for "{update}" (in order how fields are defined) + pairs = append( + pairs, + "{update}", + fns( + p.update.Queue, + p.update.Handler, + p.update.ID, + ), + ) + pairs = append(pairs, "{update.queue}", fns(p.update.Queue)) + pairs = append(pairs, "{update.handler}", fns(p.update.Handler)) + pairs = append(pairs, "{update.ID}", fns(p.update.ID)) + } + + if p.search != nil { + // replacement for "{search}" (in order how fields are defined) + pairs = append( + pairs, + "{search}", + fns(), + ) + } + return strings.NewReplacer(pairs...).Replace(in) +} + +// ********************************************************************************************************************* +// ********************************************************************************************************************* +// Action methods + +// String returns loggable description as string +// +// This function is auto-generated. +// +func (a *queueAction) String() string { + var props = &queueActionProps{} + + if a.props != nil { + props = a.props + } + + return props.Format(a.log, nil) +} + +func (e *queueAction) ToAction() *actionlog.Action { + return &actionlog.Action{ + Resource: e.resource, + Action: e.action, + Severity: e.severity, + Description: e.String(), + Meta: e.props.Serialize(), + } +} + +// ********************************************************************************************************************* +// ********************************************************************************************************************* +// Action constructors + +// QueueActionSearch returns "system:queue.search" action +// +// This function is auto-generated. +// +func QueueActionSearch(props ...*queueActionProps) *queueAction { + a := &queueAction{ + timestamp: time.Now(), + resource: "system:queue", + action: "search", + log: "searched for queues", + severity: actionlog.Info, + } + + if len(props) > 0 { + a.props = props[0] + } + + return a +} + +// QueueActionLookup returns "system:queue.lookup" action +// +// This function is auto-generated. +// +func QueueActionLookup(props ...*queueActionProps) *queueAction { + a := &queueAction{ + timestamp: time.Now(), + resource: "system:queue", + action: "lookup", + log: "looked-up for a {queue}", + severity: actionlog.Info, + } + + if len(props) > 0 { + a.props = props[0] + } + + return a +} + +// QueueActionCreate returns "system:queue.create" action +// +// This function is auto-generated. +// +func QueueActionCreate(props ...*queueActionProps) *queueAction { + a := &queueAction{ + timestamp: time.Now(), + resource: "system:queue", + action: "create", + log: "created {queue}", + severity: actionlog.Notice, + } + + if len(props) > 0 { + a.props = props[0] + } + + return a +} + +// QueueActionUpdate returns "system:queue.update" action +// +// This function is auto-generated. +// +func QueueActionUpdate(props ...*queueActionProps) *queueAction { + a := &queueAction{ + timestamp: time.Now(), + resource: "system:queue", + action: "update", + log: "updated {queue}", + severity: actionlog.Notice, + } + + if len(props) > 0 { + a.props = props[0] + } + + return a +} + +// QueueActionDelete returns "system:queue.delete" action +// +// This function is auto-generated. +// +func QueueActionDelete(props ...*queueActionProps) *queueAction { + a := &queueAction{ + timestamp: time.Now(), + resource: "system:queue", + action: "delete", + log: "deleted {queue}", + severity: actionlog.Notice, + } + + if len(props) > 0 { + a.props = props[0] + } + + return a +} + +// QueueActionUndelete returns "system:queue.undelete" action +// +// This function is auto-generated. +// +func QueueActionUndelete(props ...*queueActionProps) *queueAction { + a := &queueAction{ + timestamp: time.Now(), + resource: "system:queue", + action: "undelete", + log: "undeleted {queue}", + severity: actionlog.Notice, + } + + if len(props) > 0 { + a.props = props[0] + } + + return a +} + +// ********************************************************************************************************************* +// ********************************************************************************************************************* +// Error constructors + +// QueueErrGeneric returns "system:queue.generic" as *errors.Error +// +// +// This function is auto-generated. +// +func QueueErrGeneric(mm ...*queueActionProps) *errors.Error { + var p = &queueActionProps{} + if len(mm) > 0 { + p = mm[0] + } + + var e = errors.New( + errors.KindInternal, + + p.Format("failed to complete request due to internal error", nil), + + errors.Meta("type", "generic"), + errors.Meta("resource", "system:queue"), + + // action log entry; no formatting, it will be applied inside recordAction fn. + errors.Meta(queueLogMetaKey{}, "{err}"), + errors.Meta(queuePropsMetaKey{}, p), + + errors.StackSkip(1), + ) + + if len(mm) > 0 { + } + + return e +} + +// QueueErrNotFound returns "system:queue.notFound" as *errors.Error +// +// +// This function is auto-generated. +// +func QueueErrNotFound(mm ...*queueActionProps) *errors.Error { + var p = &queueActionProps{} + if len(mm) > 0 { + p = mm[0] + } + + var e = errors.New( + errors.KindInternal, + + p.Format("queue not found", nil), + + errors.Meta("type", "notFound"), + errors.Meta("resource", "system:queue"), + + errors.Meta(queuePropsMetaKey{}, p), + + errors.StackSkip(1), + ) + + if len(mm) > 0 { + } + + return e +} + +// QueueErrInvalidID returns "system:queue.invalidID" as *errors.Error +// +// +// This function is auto-generated. +// +func QueueErrInvalidID(mm ...*queueActionProps) *errors.Error { + var p = &queueActionProps{} + if len(mm) > 0 { + p = mm[0] + } + + var e = errors.New( + errors.KindInternal, + + p.Format("invalid ID", nil), + + errors.Meta("type", "invalidID"), + errors.Meta("resource", "system:queue"), + + errors.Meta(queuePropsMetaKey{}, p), + + errors.StackSkip(1), + ) + + if len(mm) > 0 { + } + + return e +} + +// QueueErrInvalidHandler returns "system:queue.invalidHandler" as *errors.Error +// +// +// This function is auto-generated. +// +func QueueErrInvalidHandler(mm ...*queueActionProps) *errors.Error { + var p = &queueActionProps{} + if len(mm) > 0 { + p = mm[0] + } + + var e = errors.New( + errors.KindInternal, + + p.Format("invalid handler", nil), + + errors.Meta("type", "invalidHandler"), + errors.Meta("resource", "system:queue"), + + errors.Meta(queuePropsMetaKey{}, p), + + errors.StackSkip(1), + ) + + if len(mm) > 0 { + } + + return e +} + +// QueueErrAlreadyExists returns "system:queue.alreadyExists" as *errors.Error +// +// +// This function is auto-generated. +// +func QueueErrAlreadyExists(mm ...*queueActionProps) *errors.Error { + var p = &queueActionProps{} + if len(mm) > 0 { + p = mm[0] + } + + var e = errors.New( + errors.KindInternal, + + p.Format("queue by that name already exists", nil), + + errors.Meta("type", "alreadyExists"), + errors.Meta("resource", "system:queue"), + + errors.Meta(queuePropsMetaKey{}, p), + + errors.StackSkip(1), + ) + + if len(mm) > 0 { + } + + return e +} + +// QueueErrNotAllowedToRead returns "system:queue.notAllowedToRead" as *errors.Error +// +// +// This function is auto-generated. +// +func QueueErrNotAllowedToRead(mm ...*queueActionProps) *errors.Error { + var p = &queueActionProps{} + if len(mm) > 0 { + p = mm[0] + } + + var e = errors.New( + errors.KindInternal, + + p.Format("not allowed to read this queue", nil), + + errors.Meta("type", "notAllowedToRead"), + errors.Meta("resource", "system:queue"), + + // action log entry; no formatting, it will be applied inside recordAction fn. + errors.Meta(queueLogMetaKey{}, "failed to read {queue.queue}; insufficient permissions"), + errors.Meta(queuePropsMetaKey{}, p), + + errors.StackSkip(1), + ) + + if len(mm) > 0 { + } + + return e +} + +// ********************************************************************************************************************* +// ********************************************************************************************************************* + +// recordAction is a service helper function wraps function that can return error +// +// It will wrap unrecognized/internal errors with generic errors. +// +// This function is auto-generated. +// +func (svc queue) recordAction(ctx context.Context, props *queueActionProps, actionFn func(...*queueActionProps) *queueAction, err error) error { + if svc.actionlog == nil || actionFn == nil { + // action log disabled or no action fn passed, return error as-is + return err + } else if err == nil { + // action completed w/o error, record it + svc.actionlog.Record(ctx, actionFn(props).ToAction()) + return nil + } + + a := actionFn(props).ToAction() + + // Extracting error information and recording it as action + a.Error = err.Error() + + switch c := err.(type) { + case *errors.Error: + m := c.Meta() + + a.Error = err.Error() + a.Severity = actionlog.Severity(m.AsInt("severity")) + a.Description = props.Format(m.AsString(queueLogMetaKey{}), err) + + if p, has := m[queuePropsMetaKey{}]; has { + a.Meta = p.(*queueActionProps).Serialize() + } + + svc.actionlog.Record(ctx, a) + default: + svc.actionlog.Record(ctx, a) + } + + // Original error is passed on + return err +} diff --git a/system/service/queue_actions.yaml b/system/service/queue_actions.yaml new file mode 100644 index 000000000..05aac5c5d --- /dev/null +++ b/system/service/queue_actions.yaml @@ -0,0 +1,89 @@ +# List of loggable service actions + +resource: system:queue +service: queue + +# Default sensitivity for actions +defaultActionSeverity: notice + +# default severity for errors +defaultErrorSeverity: error + +import: + - github.com/cortezaproject/corteza-server/pkg/messagebus + +props: + - name: queue + type: "*messagebus.QueueSettings" + fields: [ queue, ID ] + - name: new + type: "*messagebus.QueueSettings" + fields: [ queue, handler, ID ] + - name: update + type: "*messagebus.QueueSettings" + fields: [ queue, handler, ID ] + - name: search + type: "*messagebus.QueueSettingsFilter" + fields: [ ] + +actions: + - action: search + log: "searched for queues" + severity: info + + - action: lookup + log: "looked-up for a {queue}" + severity: info + + - action: create + log: "created {queue}" + + - action: update + log: "updated {queue}" + + - action: delete + log: "deleted {queue}" + + - action: undelete + log: "undeleted {queue}" + +errors: + - error: notFound + message: "queue not found" + severity: warning + + - error: invalidID + message: "invalid ID" + severity: warning + + - error: invalidHandler + message: "invalid handler" + severity: warning + + - error: alreadyExists + message: "queue by that name already exists" + severity: warning + + - error: notAllowedToRead + message: "not allowed to read this queue" + log: "failed to read {queue.queue}; insufficient permissions" + + # - error: notAllowedToCreate + # message: "not allowed to create templates" + # log: "failed to create template; insufficient permissions" + + # - error: notAllowedToUpdate + # message: "not allowed to update this template" + # log: "failed to update {template.handle}; insufficient permissions" + + # - error: notAllowedToDelete + # message: "not allowed to delete this template" + # log: "failed to delete {template.handle}; insufficient permissions" + + # - error: notAllowedToUndelete + # message: "not allowed to undelete this template" + # log: "failed to undelete {template.handle}; insufficient permissions" + + # - error: notAllowedToRender + # message: "not allowed to render this template" + # log: "failed to render {template.handle}; insufficient permissions" diff --git a/system/service/service.go b/system/service/service.go index 30c355626..4af9d1e64 100644 --- a/system/service/service.go +++ b/system/service/service.go @@ -75,6 +75,7 @@ var ( DefaultReminder ReminderService DefaultAttachment AttachmentService DefaultRenderer TemplateService + DefaultQueue QueueService DefaultStatistics *statistics @@ -167,6 +168,7 @@ func Initialize(ctx context.Context, log *zap.Logger, s store.Storer, c Config) DefaultSink = Sink() DefaultStatistics = Statistics() DefaultAttachment = Attachment(DefaultObjectStore) + DefaultQueue = Queue() automationService.DefaultUser = DefaultUser diff --git a/tests/messagebus/main_test.go b/tests/messagebus/main_test.go index 0b10806ea..26aa28386 100644 --- a/tests/messagebus/main_test.go +++ b/tests/messagebus/main_test.go @@ -16,6 +16,7 @@ import ( "github.com/cortezaproject/corteza-server/pkg/id" "github.com/cortezaproject/corteza-server/pkg/logger" "github.com/cortezaproject/corteza-server/pkg/messagebus" + "github.com/cortezaproject/corteza-server/pkg/options" "github.com/cortezaproject/corteza-server/pkg/rand" "github.com/cortezaproject/corteza-server/pkg/rbac" "github.com/cortezaproject/corteza-server/store/sqlite3" @@ -24,7 +25,6 @@ import ( "github.com/go-chi/chi" _ "github.com/joho/godotenv/autoload" "github.com/stretchr/testify/require" - "go.uber.org/zap" ) type ( @@ -73,8 +73,8 @@ func InitTestApp() { eventbus.Set(eventBus) - messageBus := messagebus.New(zap.NewNop(), eventbus.Service()) - // messageBus := messagebus.New(logger.Default(), eventbus.Service()) + // messageBus := messagebus.New(&options.MessagebusOpt{LogEnabled: false}, zap.NewNop(), eventbus.Service()) + messageBus := messagebus.New(&options.MessagebusOpt{LogEnabled: false}, logger.Default(), eventbus.Service()) messagebus.Set(messageBus) return nil