Produce OpenTelemetry metrics (#185)

* .golangci.yml: remove mentions of deprecated linters

* Fix "staticcheck" linter error by using grpc.NewClient

* Configure OpenTelemetry

Metrics only for now.

* Produce OpenTelemetry metrics

* Update DeploymentGuide.md

Co-authored-by: Fedor Korotkov <fedor.korotkov@gmail.com>

* Update DeploymentGuide.md

Co-authored-by: Fedor Korotkov <fedor.korotkov@gmail.com>

* Introduce "org.cirruslabs.orchard.controller.worker_status"

---------

Co-authored-by: Fedor Korotkov <fedor.korotkov@gmail.com>
This commit is contained in:
Nikolay Edigaryev
2024-06-24 18:19:51 +04:00
committed by GitHub
co-authored by Fedor Korotkov
parent 1f5bb26d20
commit ff0497b1d8
11 changed files with 338 additions and 51 deletions
+128 -1
View File
@@ -12,8 +12,12 @@ import (
storepkg "github.com/cirruslabs/orchard/internal/controller/store"
"github.com/cirruslabs/orchard/internal/controller/store/badger"
"github.com/cirruslabs/orchard/internal/netconstants"
"github.com/cirruslabs/orchard/internal/opentelemetry"
v1 "github.com/cirruslabs/orchard/pkg/resource/v1"
"github.com/cirruslabs/orchard/rpc"
"github.com/samber/lo"
"go.opentelemetry.io/otel/attribute"
"go.opentelemetry.io/otel/metric"
"go.uber.org/zap"
"golang.org/x/crypto/ssh"
"golang.org/x/net/http2"
@@ -111,8 +115,11 @@ func New(opts ...Option) (*Controller, error) {
controller.workerNotifier = notifier.NewNotifier(controller.logger.With("component", "rpc"))
// Instantiate the scheduler
controller.scheduler = scheduler.NewScheduler(store, controller.workerNotifier,
controller.scheduler, err = scheduler.NewScheduler(store, controller.workerNotifier,
controller.workerOfflineTimeout, controller.logger)
if err != nil {
return nil, err
}
// Instantiate the SSH server (if configured)
if controller.sshListenAddr != "" && controller.sshSigner != nil {
@@ -169,6 +176,11 @@ func New(opts ...Option) (*Controller, error) {
ErrInitFailed, err)
}
// Metrics
if err := controller.initializeMetrics(); err != nil {
return nil, err
}
return controller, nil
}
@@ -239,3 +251,118 @@ func (controller *Controller) SSHAddress() (string, bool) {
return controller.sshServer.Address(), true
}
//nolint:gocognit // looks OK for now
func (controller *Controller) initializeMetrics() error {
_, err := opentelemetry.DefaultMeter.Int64ObservableGauge("org.cirruslabs.orchard.controller.vm_status",
metric.WithInt64Callback(func(ctx context.Context, observer metric.Int64Observer) error {
return controller.store.View(func(txn storepkg.Transaction) error {
vms, err := txn.ListVMs()
if err != nil {
return err
}
type Key struct {
Worker string
Status v1.VMStatus
}
groups := lo.CountValuesBy(vms, func(vm v1.VM) Key {
return Key{
Worker: vm.Worker,
Status: vm.Status,
}
})
for key, count := range groups {
observer.Observe(int64(count), metric.WithAttributes(
attribute.String("worker", key.Worker),
attribute.String("status", key.Status.String()),
))
}
return nil
})
}),
)
if err != nil {
return err
}
_, err = opentelemetry.DefaultMeter.Int64ObservableGauge("org.cirruslabs.orchard.controller.worker_status",
metric.WithInt64Callback(func(ctx context.Context, observer metric.Int64Observer) error {
return controller.store.View(func(txn storepkg.Transaction) error {
workers, err := txn.ListWorkers()
if err != nil {
return err
}
groups := lo.CountValuesBy(workers, func(worker v1.Worker) string {
if worker.Offline(time.Minute) {
return "offline"
}
return "online"
})
for status, count := range groups {
observer.Observe(int64(count), metric.WithAttributes(
attribute.String("status", status),
))
}
return nil
})
}),
)
if err != nil {
return err
}
_, err = opentelemetry.DefaultMeter.Int64ObservableGauge("org.cirruslabs.orchard.controller.worker_resource",
metric.WithInt64Callback(func(ctx context.Context, observer metric.Int64Observer) error {
return controller.store.View(func(txn storepkg.Transaction) error {
workers, err := txn.ListWorkers()
if err != nil {
return err
}
vms, err := txn.ListVMs()
if err != nil {
return err
}
_, workerToResources := scheduler.ProcessVMs(vms)
for _, worker := range workers {
resourcesUsed := workerToResources.Get(worker.Name)
for key, value := range resourcesUsed {
observer.Observe(int64(value), metric.WithAttributes(
attribute.String("worker", worker.Name),
attribute.String("resource", key),
attribute.String("type", "used"),
))
}
resourcesAvailable := worker.Resources.Subtracted(resourcesUsed)
for key, value := range resourcesAvailable {
observer.Observe(int64(value), metric.WithAttributes(
attribute.String("worker", worker.Name),
attribute.String("resource", key),
attribute.String("type", "available"),
))
}
}
return nil
})
}),
)
if err != nil {
return err
}
return nil
}
+23 -4
View File
@@ -4,10 +4,12 @@ import (
"context"
"github.com/cirruslabs/orchard/internal/controller/notifier"
storepkg "github.com/cirruslabs/orchard/internal/controller/store"
"github.com/cirruslabs/orchard/internal/opentelemetry"
"github.com/cirruslabs/orchard/pkg/resource/v1"
"github.com/cirruslabs/orchard/rpc"
"github.com/prometheus/client_golang/prometheus"
"github.com/prometheus/client_golang/prometheus/promauto"
"go.opentelemetry.io/otel/metric"
"go.uber.org/zap"
"sort"
"time"
@@ -37,6 +39,8 @@ type Scheduler struct {
workerOfflineTimeout time.Duration
logger *zap.SugaredLogger
schedulingRequested chan bool
schedulingTimeHistogram metric.Float64Histogram
}
func NewScheduler(
@@ -44,14 +48,25 @@ func NewScheduler(
notifier *notifier.Notifier,
workerOfflineTimeout time.Duration,
logger *zap.SugaredLogger,
) *Scheduler {
return &Scheduler{
) (*Scheduler, error) {
scheduler := &Scheduler{
store: store,
notifier: notifier,
workerOfflineTimeout: workerOfflineTimeout,
logger: logger,
schedulingRequested: make(chan bool, 1),
}
// Metrics
var err error
scheduler.schedulingTimeHistogram, err = opentelemetry.DefaultMeter.
Float64Histogram("org.cirruslabs.orchard.controller.scheduling_time")
if err != nil {
return nil, err
}
return scheduler, nil
}
func (scheduler *Scheduler) Run() {
@@ -103,7 +118,7 @@ func (scheduler *Scheduler) schedulingLoopIteration() error {
if err != nil {
return err
}
unscheduledVMs, workerToResources := processVMs(vms)
unscheduledVMs, workerToResources := ProcessVMs(vms)
workers, err := txn.ListWorkers()
if err != nil {
@@ -119,6 +134,10 @@ func (scheduler *Scheduler) schedulingLoopIteration() error {
if resourcesRemaining.CanFit(unscheduledVM.Resources) &&
!worker.Offline(scheduler.workerOfflineTimeout) &&
!worker.SchedulingPaused {
// Metrics
scheduler.schedulingTimeHistogram.Record(context.Background(),
time.Since(unscheduledVM.CreatedAt).Seconds())
unscheduledVM.Worker = worker.Name
if err := txn.SetVM(unscheduledVM); err != nil {
@@ -148,7 +167,7 @@ func (scheduler *Scheduler) schedulingLoopIteration() error {
return err
}
func processVMs(vms []v1.VM) ([]v1.VM, WorkerToResources) {
func ProcessVMs(vms []v1.VM) ([]v1.VM, WorkerToResources) {
var unscheduledVMs []v1.VM
workerToResources := make(WorkerToResources)