| 1 | // SPDX-License-Identifier: GPL-3.0-or-later |
| 2 | |
| 3 | package functions |
| 4 | |
| 5 | import ( |
| 6 | "context" |
| 7 | "sync" |
| 8 | "testing" |
| 9 | "time" |
| 10 | |
| 11 | "github.com/netdata/netdata/go/plugins/pkg/metrix" |
| 12 | "github.com/netdata/netdata/go/plugins/plugin/framework/runtimecomp" |
| 13 | "github.com/stretchr/testify/assert" |
| 14 | "github.com/stretchr/testify/require" |
| 15 | ) |
| 16 | |
| 17 | type runtimeServiceMock struct { |
| 18 | mu sync.Mutex |
| 19 | |
| 20 | registered []runtimecomp.ComponentConfig |
| 21 | unregistered []string |
| 22 | } |
| 23 | |
| 24 | func (m *runtimeServiceMock) RegisterComponent(cfg runtimecomp.ComponentConfig) error { |
| 25 | m.mu.Lock() |
| 26 | defer m.mu.Unlock() |
| 27 | m.registered = append(m.registered, cfg) |
| 28 | return nil |
| 29 | } |
| 30 | |
| 31 | func (m *runtimeServiceMock) UnregisterComponent(name string) { |
| 32 | m.mu.Lock() |
| 33 | defer m.mu.Unlock() |
| 34 | m.unregistered = append(m.unregistered, name) |
| 35 | } |
| 36 | |
| 37 | func (m *runtimeServiceMock) RegisterProducer(string, func() error) error { return nil } |
| 38 | func (m *runtimeServiceMock) UnregisterProducer(string) {} |
| 39 | |
| 40 | func (m *runtimeServiceMock) snapshot() ([]runtimecomp.ComponentConfig, []string) { |
| 41 | m.mu.Lock() |
| 42 | defer m.mu.Unlock() |
| 43 | |
| 44 | registered := append([]runtimecomp.ComponentConfig(nil), m.registered...) |
| 45 | unregistered := append([]string(nil), m.unregistered...) |
| 46 | return registered, unregistered |
| 47 | } |
| 48 | |
| 49 | func runtimeMetricValue(t *testing.T, store metrix.RuntimeStore, name string, labels metrix.Labels) float64 { |
| 50 | t.Helper() |
| 51 | require.NotNil(t, store) |
| 52 | |
| 53 | reader := store.Read(metrix.ReadRaw()) |
| 54 | v, ok := reader.Value(name, labels) |
| 55 | require.Truef(t, ok, "metric %q not found (labels=%v)", name, labels) |
| 56 | return v |
| 57 | } |
| 58 | |
| 59 | func TestManager_RuntimeMetricsScenarios(t *testing.T) { |
| 60 | tests := map[string]struct { |
| 61 | run func(t *testing.T, mgr *Manager, in *chanInput, out *safeBuffer) |
| 62 | }{ |
| 63 | "registers and unregisters runtime component around Run lifecycle": { |
| 64 | run: func(t *testing.T, mgr *Manager, in *chanInput, _ *safeBuffer) { |
| 65 | mockSvc := &runtimeServiceMock{} |
| 66 | mgr.SetRuntimeService(mockSvc) |
| 67 | close(in.ch) |
| 68 | |
| 69 | ctx, cancel := context.WithTimeout(context.Background(), time.Second) |
| 70 | defer cancel() |
| 71 | mgr.Run(ctx, nil) |
| 72 | |
| 73 | registered, unregistered := mockSvc.snapshot() |
| 74 | require.Len(t, registered, 1) |
| 75 | assert.Equal(t, functionsRuntimeComponentName, registered[0].Name) |
| 76 | assert.Equal(t, mgr.runtimeStore, registered[0].Store) |
| 77 | assert.True(t, registered[0].Autogen.Enabled) |
| 78 | assert.Equal(t, "functions", registered[0].Module) |
| 79 | assert.Equal(t, "manager", registered[0].JobName) |
| 80 | |
| 81 | require.Len(t, unregistered, 1) |
| 82 | assert.Equal(t, functionsRuntimeComponentName, unregistered[0]) |
| 83 | }, |
| 84 | }, |
| 85 | "pathology counters and gauges are updated": { |
| 86 | run: func(t *testing.T, mgr *Manager, in *chanInput, _ *safeBuffer) { |
| 87 | mgr.workerCount = 1 |
| 88 | mgr.queueSize = 1 |
| 89 | mgr.cancelFallbackDelay = 50 * time.Millisecond |
| 90 | |
| 91 | started := make(chan struct{}, 1) |
| 92 | release := make(chan struct{}) |
| 93 | |
| 94 | mgr.Register("fn", func(fn Function) { |
| 95 | if fn.UID == "tx1" { |
| 96 | started <- struct{}{} |
| 97 | <-release |
| 98 | mgr.respUID(fn.UID, 200, "late") |
| 99 | return |
| 100 | } |
| 101 | mgr.respUID(fn.UID, 200, "ok") |
| 102 | }) |
| 103 | |
| 104 | cancel, done := startFlowManager(t, mgr) |
| 105 | defer cancel() |
| 106 | |
| 107 | in.ch <- functionLine("tx1", "fn") |
| 108 | <-started |
| 109 | in.ch <- functionLine("tx2", "fn") |
| 110 | // tx3 (former queue-full step) removed: the scheduler now blocks |
| 111 | // on full instead of rejecting, so the reader would deadlock here |
| 112 | // and the FUNCTION_CANCEL below would never be processed. |
| 113 | in.ch <- functionLine("tx1", "fn") // duplicate-active |
| 114 | in.ch <- "FUNCTION_CANCEL tx1" // fallback->499 |
| 115 | |
| 116 | waitForCondition(t, time.Second, func() bool { |
| 117 | reader := mgr.runtimeStore.Read(metrix.ReadRaw()) |
| 118 | v, ok := reader.Value(functionsRuntimeMetricPrefix+".cancel_fallback_total", nil) |
| 119 | return ok && v >= 1 |
| 120 | }, "cancel fallback metric increments") |
| 121 | |
| 122 | in.ch <- functionLine("tx1", "fn") // duplicate-tombstone |
| 123 | close(release) |
| 124 | close(in.ch) |
| 125 | waitForDone(t, done) |
| 126 | waitForCondition(t, time.Second, func() bool { |
| 127 | reader := mgr.runtimeStore.Read(metrix.ReadRaw()) |
| 128 | active, ok := reader.Value(functionsRuntimeMetricPrefix+".invocations_active", nil) |
| 129 | if !ok || active != 0 { |
| 130 | return false |
| 131 | } |
| 132 | awaiting, ok := reader.Value(functionsRuntimeMetricPrefix+".invocations_awaiting_result", nil) |
| 133 | if !ok || awaiting != 0 { |
| 134 | return false |
| 135 | } |
| 136 | pending, ok := reader.Value(functionsRuntimeMetricPrefix+".scheduler_pending", nil) |
| 137 | return ok && pending == 0 |
| 138 | }, "runtime gauges settle after shutdown") |
| 139 | |
| 140 | assert.GreaterOrEqual(t, runtimeMetricValue(t, mgr.runtimeStore, functionsRuntimeMetricPrefix+".cancel_fallback_total", nil), float64(1)) |
| 141 | assert.GreaterOrEqual(t, runtimeMetricValue(t, mgr.runtimeStore, functionsRuntimeMetricPrefix+".late_terminal_dropped_total", nil), float64(1)) |
| 142 | assert.GreaterOrEqual(t, runtimeMetricValue(t, mgr.runtimeStore, functionsRuntimeMetricPrefix+".duplicate_uid_ignored_total", nil), float64(2)) |
| 143 | assert.Equal(t, float64(4), runtimeMetricValue(t, mgr.runtimeStore, functionsRuntimeMetricPrefix+".calls_total", nil)) |
| 144 | |
| 145 | assert.Equal(t, float64(0), runtimeMetricValue(t, mgr.runtimeStore, functionsRuntimeMetricPrefix+".invocations_active", nil)) |
| 146 | assert.Equal(t, float64(0), runtimeMetricValue(t, mgr.runtimeStore, functionsRuntimeMetricPrefix+".invocations_awaiting_result", nil)) |
| 147 | assert.Equal(t, float64(0), runtimeMetricValue(t, mgr.runtimeStore, functionsRuntimeMetricPrefix+".scheduler_pending", nil)) |
| 148 | }, |
| 149 | }, |
| 150 | } |
| 151 | |
| 152 | for name, tc := range tests { |
| 153 | t.Run(name, func(t *testing.T) { |
| 154 | mgr, out := newFlowManager() |
| 155 | in := &chanInput{ch: make(chan string, 32)} |
| 156 | mgr.input = in |
| 157 | tc.run(t, mgr, in, out) |
| 158 | }) |
| 159 | } |
| 160 | } |