From 6de0558e51366353e99501d32d3a6556026e58b9 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Martin=20Gr=C3=B8nlien=20Pejcoch?= Date: Thu, 12 Aug 2021 15:44:14 +0200 Subject: [PATCH] Add product-drop-timeout parameter --- cmd/mmsd/main.go | 9 ++++++++- internal/server/api.go | 4 ++-- internal/server/cache_test.go | 2 +- internal/server/productstatus.go | 21 ++++++++++++++++----- internal/server/productstatus_test.go | 2 +- 5 files changed, 28 insertions(+), 10 deletions(-) diff --git a/cmd/mmsd/main.go b/cmd/mmsd/main.go index cf29d15..1026ddf 100644 --- a/cmd/mmsd/main.go +++ b/cmd/mmsd/main.go @@ -77,6 +77,11 @@ func main() { Usage: "Specify the port number for the API listening port.", Value: 8080, }), + altsrc.NewIntFlag(&cli.IntFlag{ + Name: "product-drop-timeout", + Usage: "Specify how many seconds to wait before a product not seen by the system is dropped from the overview. Default is 604800 (7d)", + Value: 604800, + }), altsrc.NewIntFlag(&cli.IntFlag{ Name: "nats-port", Usage: "Specify the port number for the NATS listening port.", @@ -166,7 +171,7 @@ func main() { } templates := server.CreateTemplates() - webService := server.NewService(templates, eventsDB, stateDB, natsURL) + webService := server.NewService(templates, eventsDB, stateDB, natsURL, ctx.Int("product-drop-timeout")) log.Println("Populating productstatus from the local events database ...") events, err := webService.GetAllEvents(context.Background()) @@ -354,6 +359,8 @@ func startEventLoop(webService *server.Service) { if err := webService.DeleteOldEvents(time.Now().AddDate(0, 0, -3)); err != nil { log.Printf("failed to delete old events from events db: %s", err) } + + webService.Productstatus.PurgeOldProducts(604800) } } }() diff --git a/internal/server/api.go b/internal/server/api.go index 04f3868..73e0ff2 100644 --- a/internal/server/api.go +++ b/internal/server/api.go @@ -56,7 +56,7 @@ type HTTPServerError struct { } // NewService creates a service struct, containing all that is needed for a mmsd server to run. -func NewService(templates *template.Template, eventsDB *sql.DB, stateDB *sql.DB, natsURL string) *Service { +func NewService(templates *template.Template, eventsDB *sql.DB, stateDB *sql.DB, natsURL string, productDropTimeout int) *Service { m := NewServiceMetrics(MetricsOpts{}) service := Service{ @@ -67,7 +67,7 @@ func NewService(templates *template.Template, eventsDB *sql.DB, stateDB *sql.DB, Router: mux.NewRouter(), NatsURL: natsURL, Metrics: m, - Productstatus: NewProductstatus(m), + Productstatus: NewProductstatus(m, productDropTimeout), } service.setRoutes() diff --git a/internal/server/cache_test.go b/internal/server/cache_test.go index 8448989..4b442e1 100644 --- a/internal/server/cache_test.go +++ b/internal/server/cache_test.go @@ -75,7 +75,7 @@ func NewMockService() (*Service, sqlmock.Sqlmock, error) { } templates := CreateTemplates() - webService := NewService(templates, eventsDB, nil, "") + webService := NewService(templates, eventsDB, nil, "", 60) return webService, mock, nil } diff --git a/internal/server/productstatus.go b/internal/server/productstatus.go index fc62c0e..6e1a23e 100644 --- a/internal/server/productstatus.go +++ b/internal/server/productstatus.go @@ -14,13 +14,15 @@ type Product struct { } type Productstatus struct { - Products map[string]Product - GaugeVec *prometheus.GaugeVec + Products map[string]Product + ProductDropTimeout int + GaugeVec *prometheus.GaugeVec } -func NewProductstatus(m *metrics) *Productstatus { +func NewProductstatus(m *metrics, productDropTimeout int) *Productstatus { productstatus := Productstatus{ - Products: make(map[string]Product), + Products: make(map[string]Product), + ProductDropTimeout: productDropTimeout, GaugeVec: prometheus.NewGaugeVec( prometheus.GaugeOpts{ Subsystem: "mmsd", @@ -60,7 +62,7 @@ func (p *Productstatus) GetProductDelays(t time.Time) { func (p *Productstatus) UpdateMetrics() { for k, v := range p.Products { - diff := time.Now().Sub(v.NextInstanceExpected) + diff := time.Since(v.NextInstanceExpected) p.GaugeVec.WithLabelValues(k).Set(diff.Seconds()) } } @@ -70,3 +72,12 @@ func (p *Productstatus) Populate(events []*mms.ProductEvent) { p.PushEvent(*event) } } + +func (p *Productstatus) PurgeOldProducts(secondsAgo int) { + for k, v := range p.Products { + diff := time.Since(v.NextInstanceExpected) + if diff.Seconds() <= -float64(secondsAgo) { + delete(p.Products, k) + } + } +} diff --git a/internal/server/productstatus_test.go b/internal/server/productstatus_test.go index 1a9d44e..79a7bb2 100644 --- a/internal/server/productstatus_test.go +++ b/internal/server/productstatus_test.go @@ -9,7 +9,7 @@ import ( func TestPushEvent(t *testing.T) { metrics := NewServiceMetrics(MetricsOpts{}) - ps := NewProductstatus(metrics) + ps := NewProductstatus(metrics, 60) var productEventList [3]mms.ProductEvent