Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 1 addition & 1 deletion apricot/protos/apricot.pb.go

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

2 changes: 1 addition & 1 deletion apricot/protos/apricot_grpc.pb.go

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

1,499 changes: 756 additions & 743 deletions coconut/protos/o2control.pb.go

Large diffs are not rendered by default.

16 changes: 15 additions & 1 deletion coconut/protos/o2control_grpc.pb.go

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

6 changes: 2 additions & 4 deletions common/event/writer.go
Original file line number Diff line number Diff line change
Expand Up @@ -110,8 +110,7 @@ func NewWriterWithTopic(topic topic.Topic) *KafkaWriter {
metric.SetFieldUInt64("messages_failed", 0)

metricDuration := writer.newMetric(KAFKAWRITER)
defer monitoring.SendHistogrammable(&metricDuration)
defer monitoring.TimerNS(&metricDuration)()
defer monitoring.TimerSendHist(&metricDuration, monitoring.Nanosecond)()

if err := writer.WriteMessages(context.Background(), messages...); err != nil {
metric.SetFieldUInt64("messages_failed", uint64(len(messages)))
Expand Down Expand Up @@ -250,8 +249,7 @@ func (w *KafkaWriter) WriteEventWithTimestamp(e interface{}, timestamp time.Time

metric := w.newMetric(KAFKAPREPARE)

defer monitoring.SendHistogrammable(&metric)
defer monitoring.TimerNS(&metric)()
defer monitoring.TimerSendHist(&metric, monitoring.Nanosecond)()

wrappedEvent, key, err := internalEventToKafkaEvent(e, timestamp)
if err != nil {
Expand Down
79 changes: 79 additions & 0 deletions common/monitoring/grpcinterceptor.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,79 @@
/*
* === This file is part of ALICE O² ===
*
* Copyright 2025 CERN and copyright holders of ALICE O².
* Author: Michal Tichak <michal.tichak@cern.ch>
*
* This program is free software: you can redistribute it and/or modify
* it under the terms of the GNU General Public License as published by
* the Free Software Foundation, either version 3 of the License, or
* (at your option) any later version.
*
* This program is distributed in the hope that it will be useful,
* but WITHOUT ANY WARRANTY; without even the implied warranty of
* MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the
* GNU General Public License for more details.
*
* You should have received a copy of the GNU General Public License
* along with this program. If not, see <http://www.gnu.org/licenses/>.
*
* In applying this license CERN does not waive the privileges and
* immunities granted to it by virtue of its status as an
* Intergovernmental Organization or submit itself to any jurisdiction.
*/

package monitoring

import (
"context"

"google.golang.org/grpc"
)

type measuredClientStream struct {
grpc.ClientStream
method string
metricName string
}

func (t *measuredClientStream) RecvMsg(m interface{}) error {
metric := NewMetric(t.metricName)
metric.AddTag("method", t.method)
defer TimerSendSingle(&metric, Millisecond)()

err := t.ClientStream.RecvMsg(m)
return err
}

type NameConvertType func(string) string

func SetupStreamClientInterceptor(metricName string, convert NameConvertType) grpc.StreamClientInterceptor {
return func(
ctx context.Context,
desc *grpc.StreamDesc,
cc *grpc.ClientConn,
method string,
streamer grpc.Streamer,
opts ...grpc.CallOption,
) (grpc.ClientStream, error) {
clientStream, err := streamer(ctx, desc, cc, method, opts...)
if err != nil {
return nil, err
}

return &measuredClientStream{
ClientStream: clientStream,
method: convert(method),
metricName: metricName,
}, nil
}
}

func SetupUnaryClientInterceptor(name string, convert NameConvertType) grpc.UnaryClientInterceptor {
return func(ctx context.Context, method string, req, reply any, cc *grpc.ClientConn, invoker grpc.UnaryInvoker, opts ...grpc.CallOption) error {
metric := NewMetric(name)
metric.AddTag("method", convert(method))
defer TimerSendSingle(&metric, Millisecond)()
return invoker(ctx, method, req, reply, cc, opts...)
}
}
5 changes: 5 additions & 0 deletions common/monitoring/metric.go
Original file line number Diff line number Diff line change
Expand Up @@ -27,6 +27,7 @@ package monitoring
import (
"fmt"
"io"
"sort"
"time"

lp "github.com/influxdata/line-protocol/v2/lineprotocol"
Expand Down Expand Up @@ -110,6 +111,10 @@ func Format(writer io.Writer, metrics []Metric) error {
var enc lp.Encoder

for _, metric := range metrics {
// AddTag requires tags sorted lexicografically
sort.Slice(metric.tags, func(i int, j int) bool {
return metric.tags[i].name < metric.tags[j].name
})
enc.StartLine(metric.name)
for _, tag := range metric.tags {
enc.AddTag(tag.name, tag.value)
Expand Down
4 changes: 2 additions & 2 deletions common/monitoring/monitoring_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -546,8 +546,8 @@ func TestMetricsHistogramObject(t *testing.T) {
}

func measureFunc(metric *Metric) {
defer TimerMS(metric)()
defer TimerNS(metric)()
defer Timer(metric, Millisecond)()
defer Timer(metric, Nanosecond)()
time.Sleep(100 * time.Millisecond)
}

Expand Down
60 changes: 50 additions & 10 deletions common/monitoring/timer.go
Original file line number Diff line number Diff line change
@@ -1,19 +1,59 @@
package monitoring

import "time"
import (
"time"

// Timer* functions are meant to be used with defer statement to measure runtime of given function:
// defer TimerNS(&metric)()
func TimerMS(metric *Metric) func() {
start := time.Now()
return func() {
metric.SetFieldInt64("execution_time_ms", time.Since(start).Milliseconds())
}
"github.com/AliceO2Group/Control/common/logger/infologger"
)

type TimeResolution int

const (
Millisecond TimeResolution = iota
Nanosecond
)

// Timer function is meant to be used with defer statement to measure runtime of given function:
// defer Timer(&metric, Milliseconds)()
func Timer(metric *Metric, unit TimeResolution) func() {
return timer(metric, unit, false, false)
}

func TimerNS(metric *Metric) func() {
// Timer function is meant to be used with defer statement to measure runtime of given function:
// defer Timer(&metric, Milliseconds)()
// sends measured value as Send(metric)
func TimerSendSingle(metric *Metric, unit TimeResolution) func() {
return timer(metric, unit, true, false)
}

// Timer function is meant to be used with defer statement to measure runtime of given function:
// defer Timer(&metric, Milliseconds)()
// sends measured value as SendHistogrammable(metric)
func TimerSendHist(metric *Metric, unit TimeResolution) func() {
return timer(metric, unit, true, true)
}

func timer(metric *Metric, unit TimeResolution, send bool, sendHistogrammable bool) func() {
start := time.Now()

return func() {
metric.SetFieldInt64("execution_time_ns", time.Since(start).Nanoseconds())
dur := time.Since(start)
// we are setting default value as Nanoseconds
switch unit {
case Millisecond:
metric.SetFieldInt64("execution_time_ms", dur.Milliseconds())
case Nanosecond:
metric.SetFieldInt64("execution_time_ns", dur.Nanoseconds())
default:
log.WithField("level", infologger.IL_Devel).Warnf("trying to use unknown time resolution in monitoring.timer function [%d], skipping", unit)
}

if send {
if sendHistogrammable {
SendHistogrammable(metric)
} else {
Send(metric)
}
}
}
}
4 changes: 2 additions & 2 deletions common/protos/common.pb.go

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

Loading
Loading