diff --git a/services/postprocessing/pkg/metrics/metrics.go b/services/postprocessing/pkg/metrics/metrics.go new file mode 100644 index 000000000..077a04018 --- /dev/null +++ b/services/postprocessing/pkg/metrics/metrics.go @@ -0,0 +1,69 @@ +package metrics + +import ( + "github.com/prometheus/client_golang/prometheus" + "github.com/prometheus/client_golang/prometheus/promauto" +) + +var ( + // Namespace defines the namespace for the defines metrics. + Namespace = "opencloud" + + // Subsystem defines the subsystem for the defines metrics. + Subsystem = "postprocessing" +) + +// Metrics defines the available metrics of this service. +type Metrics struct { + // Counter *prometheus.CounterVec + BuildInfo *prometheus.GaugeVec + EventsOutstandingAcks prometheus.Gauge + EventsUnprocessed prometheus.Gauge + EventsRedelivered prometheus.Gauge + InProgress prometheus.Gauge + Finished *prometheus.CounterVec +} + +// New initializes the available metrics. +func New() *Metrics { + m := &Metrics{ + BuildInfo: promauto.NewGaugeVec(prometheus.GaugeOpts{ + Namespace: Namespace, + Subsystem: Subsystem, + Name: "build_info", + Help: "Build information", + }, []string{"version"}), + EventsOutstandingAcks: promauto.NewGauge(prometheus.GaugeOpts{ + Namespace: Namespace, + Subsystem: Subsystem, + Name: "events_outstanding_acks", + Help: "Number of outstanding acks for events", + }), + EventsUnprocessed: promauto.NewGauge(prometheus.GaugeOpts{ + Namespace: Namespace, + Subsystem: Subsystem, + Name: "events_unprocessed", + Help: "Number of unprocessed events", + }), + EventsRedelivered: promauto.NewGauge(prometheus.GaugeOpts{ + Namespace: Namespace, + Subsystem: Subsystem, + Name: "events_redelivered", + Help: "Number of redelivered events", + }), + InProgress: promauto.NewGauge(prometheus.GaugeOpts{ + Namespace: Namespace, + Subsystem: Subsystem, + Name: "in_progress", + Help: "Number of postprocessing events in progress", + }), + Finished: promauto.NewCounterVec(prometheus.CounterOpts{ + Namespace: Namespace, + Subsystem: Subsystem, + Name: "finished", + Help: "Number of finished postprocessing events", + }, []string{"status"}), + } + + return m +} diff --git a/services/postprocessing/pkg/service/service.go b/services/postprocessing/pkg/service/service.go index 2b41b2888..c9ca40d3f 100644 --- a/services/postprocessing/pkg/service/service.go +++ b/services/postprocessing/pkg/service/service.go @@ -9,7 +9,9 @@ import ( "time" "github.com/opencloud-eu/opencloud/pkg/log" + "github.com/opencloud-eu/opencloud/pkg/version" "github.com/opencloud-eu/opencloud/services/postprocessing/pkg/config" + "github.com/opencloud-eu/opencloud/services/postprocessing/pkg/metrics" "github.com/opencloud-eu/opencloud/services/postprocessing/pkg/postprocessing" ctxpkg "github.com/opencloud-eu/reva/v2/pkg/ctx" "github.com/opencloud-eu/reva/v2/pkg/events" @@ -22,14 +24,15 @@ import ( // PostprocessingService is an instance of the service handling postprocessing of files type PostprocessingService struct { - ctx context.Context - log log.Logger - events <-chan raw.Event - pub events.Publisher - steps []events.Postprocessingstep - store store.Store - c config.Postprocessing - tp trace.TracerProvider + ctx context.Context + log log.Logger + events <-chan raw.Event + pub events.Publisher + steps []events.Postprocessingstep + store store.Store + c config.Postprocessing + tp trace.TracerProvider + metrics *metrics.Metrics } var ( @@ -78,15 +81,20 @@ func NewPostprocessingService(ctx context.Context, logger log.Logger, sto store. return nil, err } + m := metrics.New() + m.BuildInfo.WithLabelValues(version.GetString()).Set(1) + monitorMetrics(raw, "postprocessing-pull", m, logger) + return &PostprocessingService{ - ctx: ctx, - log: logger, - events: evs, - pub: pub, - steps: getSteps(cfg.Postprocessing), - store: sto, - c: cfg.Postprocessing, - tp: tp, + ctx: ctx, + log: logger, + events: evs, + pub: pub, + steps: getSteps(cfg.Postprocessing), + store: sto, + c: cfg.Postprocessing, + tp: tp, + metrics: m, }, nil } @@ -150,6 +158,7 @@ func (pps *PostprocessingService) processEvent(e raw.Event) error { InitiatorID: e.InitiatorID, ImpersonatingUser: ev.ImpersonatingUser, } + pps.metrics.InProgress.Inc() next = pp.Init(ev) case events.PostprocessingStepFinished: if ev.UploadID == "" { @@ -200,7 +209,9 @@ func (pps *PostprocessingService) processEvent(e raw.Event) error { } }) case events.UploadReady: + pps.metrics.InProgress.Dec() if ev.Failed { + pps.metrics.Finished.WithLabelValues("failed", string(pp.Status.Outcome)).Inc() // the upload failed - let's keep it around for a while - but mark it as finished pp, err = pps.getPP(pps.store, ev.UploadID) if err != nil { @@ -211,6 +222,7 @@ func (pps *PostprocessingService) processEvent(e raw.Event) error { return storePP(pps.store, pp) } + pps.metrics.Finished.WithLabelValues("succeeded").Inc() // the storage provider thinks the upload is done - so no need to keep it any more if err := pps.store.Delete(ev.UploadID); err != nil { pps.log.Error().Str("uploadID", ev.UploadID).Err(err).Msg("cannot delete upload") @@ -360,3 +372,25 @@ func (pps *PostprocessingService) findUploadsByStep(step events.Postprocessingst return ids } + +func monitorMetrics(stream raw.Stream, name string, m *metrics.Metrics, logger log.Logger) { + ctx := context.Background() + consumer, err := stream.JetStream().Consumer(ctx, name) + if err != nil { + logger.Error().Err(err).Msg("failed to get consumer") + } + ticker := time.NewTicker(5 * time.Second) + go func() { + for range ticker.C { + info, err := consumer.Info(ctx) + if err != nil { + logger.Error().Err(err).Msg("failed to get consumer") + } + + m.EventsOutstandingAcks.Set(float64(info.NumAckPending)) + m.EventsUnprocessed.Set(float64(info.NumPending)) + m.EventsRedelivered.Set(float64(info.NumRedelivered)) + logger.Trace().Msg("updated postprocessing event metrics") + } + }() +} diff --git a/services/search/pkg/search/events.go b/services/search/pkg/search/events.go index eca39f94f..a02286d9c 100644 --- a/services/search/pkg/search/events.go +++ b/services/search/pkg/search/events.go @@ -127,7 +127,7 @@ func monitorMetrics(stream raw.Stream, name string, m *metrics.Metrics, logger l m.EventsOutstandingAcks.Set(float64(info.NumAckPending)) m.EventsUnprocessed.Set(float64(info.NumPending)) m.EventsRedelivered.Set(float64(info.NumRedelivered)) - logger.Debug().Msg("updated event metrics") + logger.Trace().Msg("updated search event metrics") } }() }