From 996c66c10cb8d64a0f66c8735f7967ad92614a8b Mon Sep 17 00:00:00 2001 From: Stelios Fragkakis <52996999+stelfrag@users.noreply.github.com> Date: Thu, 18 Jun 2026 16:06:25 +0300 Subject: [PATCH 1/9] Track and account memory allocation size in replication queries (#22756) - Added `alloc_size` field to `struct replication_query` to store allocation size. - Modified `replication_buffers_allocated` updates to use `alloc_size` for both increments and decrements. - Introduced an additional buffer deallocation safeguard during request cancellation. --- src/streaming/stream-replication-sender.c | 10 +++++++--- 1 file changed, 7 insertions(+), 3 deletions(-) diff --git a/src/streaming/stream-replication-sender.c b/src/streaming/stream-replication-sender.c index e15378df96730c..0391bc213bb028 100644 --- a/src/streaming/stream-replication-sender.c +++ b/src/streaming/stream-replication-sender.c @@ -98,6 +98,7 @@ struct replication_query { STORAGE_ENGINE_BACKEND backend; struct replication_request *rq; + size_t alloc_size; // bytes accounted in replication_buffers_allocated for this query size_t dimensions; struct replication_dimension data[]; }; @@ -118,9 +119,11 @@ static struct replication_query *replication_query_prepare( bool synchronous ) { size_t dimensions = rrdset_number_of_dimensions(st); - struct replication_query *q = callocz(1, sizeof(struct replication_query) + dimensions * sizeof(struct replication_dimension)); - __atomic_add_fetch(&replication_buffers_allocated, sizeof(struct replication_query) + dimensions * sizeof(struct replication_dimension), __ATOMIC_RELAXED); + size_t alloc_size = sizeof(struct replication_query) + dimensions * sizeof(struct replication_dimension); + struct replication_query *q = callocz(1, alloc_size); + __atomic_add_fetch(&replication_buffers_allocated, alloc_size, __ATOMIC_RELAXED); + q->alloc_size = alloc_size; q->dimensions = dimensions; q->st = st; @@ -301,7 +304,7 @@ static void replication_query_finalize(BUFFER *wb, struct replication_query *q, spinlock_unlock(&replication_queries.spinlock); } - __atomic_sub_fetch(&replication_buffers_allocated, sizeof(struct replication_query) + dimensions * sizeof(struct replication_dimension), __ATOMIC_RELAXED); + __atomic_sub_fetch(&replication_buffers_allocated, q->alloc_size, __ATOMIC_RELAXED); freez(q); } @@ -1574,6 +1577,7 @@ static void replication_pipeline_cancel_and_cleanup(void) { internal_error(true, "REPLICATION: cancelled %zu inflight queries", cancelled); + __atomic_sub_fetch(&replication_buffers_allocated, rtp.max_requests_ahead * sizeof(struct replication_request), __ATOMIC_RELAXED); freez(rtp.rqs); rtp.rqs = NULL; rtp.max_requests_ahead = 0; From 8daf86c9fb9697301fb84a29451e4309b39e02c1 Mon Sep 17 00:00:00 2001 From: "nedi-app[bot]" <267954999+nedi-app[bot]@users.noreply.github.com> Date: Thu, 18 Jun 2026 13:21:22 +0000 Subject: [PATCH 2/9] docs: Clarify StatsD [app] section is a namespace requiring chart (#22673) * docs: Added a prominent note in the StatsD README 'Application Section Options' section clarifying that the [app] section is a namespace/container that does not create dashboard charts by itself, and that users must add chart definition sections to see synthetic charts. * docs: clarify statsd [app] section is a namespace requiring chart * Delete .nedi-audit/state.json * docs: Convert the troubleshooting/explanatory content about why StatsD [app] * review --------- Co-authored-by: nedi-app[bot] Co-authored-by: Fotis Voutsas --- src/collectors/statsd.plugin/README.md | 14 ++++++++++---- 1 file changed, 10 insertions(+), 4 deletions(-) diff --git a/src/collectors/statsd.plugin/README.md b/src/collectors/statsd.plugin/README.md index 1d1019a655e613..5eed2e576f8f70 100644 --- a/src/collectors/statsd.plugin/README.md +++ b/src/collectors/statsd.plugin/README.md @@ -517,14 +517,22 @@ Example of a synthetic chart combining multiple metrics: The `[app]` section defines the application and has these options: +:::warning + +The `[app]` section is a **namespace/container** — it groups metrics and sets defaults, but does **not** create any dashboard charts by itself. To see synthetic charts on the dashboard, you **must** add one or more chart definition sections (e.g., `[mychart]`) below the `[app]` section. If you only define an `[app]` section without chart definitions, the only visible charts will be private charts for individual metrics (if `private charts = yes` or the global default is enabled). + +Settings like `private charts`, `gaps when not collected`, and `history` configure how the app's metrics and charts behave — they are not chart-level settings. The `memory mode` setting under `[app]` is currently ignored. See [Chart Definitions](#chart-definitions) below for how to create charts. + +::: + :::note - **name** - Defines the application name - **metrics** - [Simple pattern](https://github.com/netdata/netdata/blob/master/src/libnetdata/simple_pattern/README.md) matching all metrics for this app - **private charts** - Enable/disable private charts for matched metrics (yes|no) - **gaps when not collected** - Show gaps when no metrics are collected (yes|no) -- **memory mode** - Sets memory mode for application charts (optional, default is global Netdata setting) -- **history** - Size of round-robin database (optional, only relevant with `memory mode = save`) +- **memory mode** - Ignored in the `[app]` section; application charts use the host's default memory mode +- **history** - Size of round-robin database for application charts (optional, minimum 5) ::: @@ -919,7 +927,6 @@ Start with this basic configuration: metrics = k6* private charts = yes gaps when not collected = no - memory mode = dbengine ``` @@ -957,7 +964,6 @@ Here's a complete configuration for k6: metrics = k6* private charts = yes gaps when not collected = no - memory mode = dbengine [dictionary] http_req_blocked = Blocked HTTP Requests From 58f8b9191e4ab9d9e76339101b10f8548cdb37e3 Mon Sep 17 00:00:00 2001 From: Ilya Mashchenko Date: Thu, 18 Jun 2026 16:51:41 +0300 Subject: [PATCH 3/9] refactor(go.d): bind single-instance functions to runtime job (#22767) --- .../SKILL.md | 7 + .../agent/jobmgr/funcctl/controller_test.go | 188 +++++++++++++++++- .../plugin/agent/jobmgr/funcctl/dispatch.go | 24 +++ .../plugin/framework/collectorapi/registry.go | 8 +- .../go.d/collector/snmp_topology/collector.go | 22 +- .../snmp_topology/collector_refresh_test.go | 16 +- .../collector/snmp_topology/func_topology.go | 4 +- .../snmp_topology/func_topology_handler.go | 12 +- .../func_topology_presentation_test.go | 17 +- .../snmp_topology/func_topology_test.go | 32 +-- .../snmp_topology/topology_registry.go | 2 - .../snmp_topology/topology_trap_enrich.go | 37 +++- .../topology_trap_enrich_test.go | 84 +++++++- 13 files changed, 378 insertions(+), 75 deletions(-) diff --git a/.agents/skills/project-writing-go-modules-framework-v2/SKILL.md b/.agents/skills/project-writing-go-modules-framework-v2/SKILL.md index f47de2eed98434..5fc31239d3b300 100644 --- a/.agents/skills/project-writing-go-modules-framework-v2/SKILL.md +++ b/.agents/skills/project-writing-go-modules-framework-v2/SKILL.md @@ -67,6 +67,13 @@ source files for evidence. - If Functions exist, isolate them in a `func/` subpackage with a narrow `Deps` interface declared there. The Function package MUST NOT import the collector package or hold `*Collector`. +- If a single-instance collector exposes `AgentWide` module `Methods`, its + `MethodHandler(job)` receives the running canonical runtime job. Use + `job.Collector()` to bind the Function handler to collector-owned state; do + not add a `__job` parameter or introduce a package-global registry to bridge + Function dispatch. During method execution, when the singleton job is not + running, the framework returns an unavailable response before calling + `MethodHandler`. - `collectorapi.Creator.InstancePolicy` defaults to `InstancePolicyPerJob`. Use `InstancePolicySingle` only for collectors that are intentionally one canonical job per agent. Single-instance configs MUST diff --git a/src/go/plugin/agent/jobmgr/funcctl/controller_test.go b/src/go/plugin/agent/jobmgr/funcctl/controller_test.go index 3b7c72084800dc..87865dfbdd652a 100644 --- a/src/go/plugin/agent/jobmgr/funcctl/controller_test.go +++ b/src/go/plugin/agent/jobmgr/funcctl/controller_test.go @@ -938,6 +938,188 @@ func TestControllerRawAgentWideModuleMethodDoesNotRequireRunningJob(t *testing.T assert.Nil(t, gotJob) } +func TestControllerRawSingleInstanceAgentWideModuleMethodUsesRunningJob(t *testing.T) { + var gotCode int + var gotResp map[string]any + var gotJob collectorapi.RuntimeJob + reg := newTestFunctionRegistry() + controller := New(Options{ + FnReg: reg, + JSONWriter: func(data []byte, code int) { + gotCode = code + require.NoError(t, json.Unmarshal(data, &gotResp)) + }, + }) + controller.RegisterModules(collectorapi.Registry{ + "mod": collectorapi.Creator{ + InstancePolicy: collectorapi.InstancePolicySingle, + Methods: func() []funcapi.MethodConfig { + return []funcapi.MethodConfig{{ + ID: "logs", + RawRequest: true, + AgentWide: true, + }} + }, + MethodHandler: func(job collectorapi.RuntimeJob) funcapi.MethodHandler { + gotJob = job + return &rawTestHandler{ + raw: func(_ context.Context, req funcapi.RawMethodRequest) *funcapi.FunctionResponse { + assert.Equal(t, "logs", req.Method) + return funcapi.RawResponse(map[string]any{ + "status": 200, + "type": "table", + }) + }, + } + }, + }, + }) + job := newTestRuntimeJob("mod", "mod", true) + controller.OnJobStart(job) + + reg.call("mod:logs", context.Background(), functions.Function{ + UID: "raw-single-agent-wide", + Timeout: time.Second, + }) + + assert.Equal(t, 200, gotCode) + assert.Equal(t, float64(200), gotResp["status"]) + assert.Same(t, job, gotJob) +} + +func TestControllerSingleInstanceAgentWideModuleMethodUsesRunningJob(t *testing.T) { + var gotCode int + var gotResp map[string]any + var gotJob collectorapi.RuntimeJob + reg := newTestFunctionRegistry() + controller := New(Options{ + FnReg: reg, + JSONWriter: func(data []byte, code int) { + gotCode = code + require.NoError(t, json.Unmarshal(data, &gotResp)) + }, + }) + controller.RegisterModules(collectorapi.Registry{ + "mod": collectorapi.Creator{ + InstancePolicy: collectorapi.InstancePolicySingle, + Methods: func() []funcapi.MethodConfig { + return []funcapi.MethodConfig{{ + ID: "status", + AgentWide: true, + }} + }, + MethodHandler: func(job collectorapi.RuntimeJob) funcapi.MethodHandler { + gotJob = job + return &rawTestHandler{ + params: func(context.Context, string) ([]funcapi.ParamConfig, error) { + return []funcapi.ParamConfig{{ + ID: "scope", + Name: "Scope", + Selection: funcapi.ParamSelect, + Options: []funcapi.ParamOption{{ + ID: "all", + Name: "All", + Default: true, + }}, + }}, nil + }, + handle: func(_ context.Context, method string, params funcapi.ResolvedParams) *funcapi.FunctionResponse { + assert.Equal(t, "status", method) + assert.Equal(t, "all", params.GetOne("scope")) + return &funcapi.FunctionResponse{ + Status: 200, + ResponseType: "table", + Help: "status", + } + }, + } + }, + }, + }) + job := newTestRuntimeJob("mod", "mod", true) + controller.OnJobStart(job) + + reg.call("mod:status", context.Background(), functions.Function{ + UID: "single-agent-wide", + Timeout: time.Second, + Payload: []byte(`{"scope":"all"}`), + }) + + assert.Equal(t, 200, gotCode) + assert.Equal(t, float64(200), gotResp["status"]) + assert.Same(t, job, gotJob) + assert.Equal(t, []any{"scope"}, gotResp["accepted_params"]) +} + +func TestControllerSingleInstanceAgentWideModuleMethodRequiresRunningJob(t *testing.T) { + tests := map[string]struct { + setup func(*Controller) + message string + }{ + "before start": { + message: "module 'mod' is not running", + }, + "after stop": { + setup: func(controller *Controller) { + job := newTestRuntimeJob("mod", "mod", true) + controller.OnJobStart(job) + controller.OnJobStop(job) + }, + message: "module 'mod' is not running", + }, + "registered but not running": { + setup: func(controller *Controller) { + controller.OnJobStart(newTestRuntimeJob("mod", "mod", false)) + }, + message: "job 'mod' is no longer running", + }, + } + + for name, tc := range tests { + t.Run(name, func(t *testing.T) { + var gotCode int + var gotResp map[string]any + var gotHandler bool + reg := newTestFunctionRegistry() + controller := New(Options{ + FnReg: reg, + JSONWriter: func(data []byte, code int) { + gotCode = code + require.NoError(t, json.Unmarshal(data, &gotResp)) + }, + }) + controller.RegisterModules(collectorapi.Registry{ + "mod": collectorapi.Creator{ + InstancePolicy: collectorapi.InstancePolicySingle, + Methods: func() []funcapi.MethodConfig { + return []funcapi.MethodConfig{{ + ID: "status", + AgentWide: true, + }} + }, + MethodHandler: func(collectorapi.RuntimeJob) funcapi.MethodHandler { + gotHandler = true + return &tableTestHandler{} + }, + }, + }) + if tc.setup != nil { + tc.setup(controller) + } + + reg.call("mod:status", context.Background(), functions.Function{ + UID: "single-agent-wide-missing-job", + Timeout: time.Second, + }) + + assert.Equal(t, 503, gotCode) + assert.Equal(t, float64(503), gotResp["status"]) + assert.Equal(t, tc.message, gotResp["errorMessage"]) + assert.False(t, gotHandler) + }) + } +} + func TestControllerModuleMethodRequestContextCancellation(t *testing.T) { ctx, cancel := context.WithCancel(context.Background()) cancel() @@ -1119,13 +1301,17 @@ func (j *testRuntimeJob) IsRunning() bool { return j.running } func (j *testRuntimeJob) Collector() any { return nil } type rawTestHandler struct { + params func(context.Context, string) ([]funcapi.ParamConfig, error) handle func(context.Context, string, funcapi.ResolvedParams) *funcapi.FunctionResponse raw func(context.Context, funcapi.RawMethodRequest) *funcapi.FunctionResponse } var _ funcapi.RawMethodHandler = (*rawTestHandler)(nil) -func (h *rawTestHandler) MethodParams(context.Context, string) ([]funcapi.ParamConfig, error) { +func (h *rawTestHandler) MethodParams(ctx context.Context, method string) ([]funcapi.ParamConfig, error) { + if h.params != nil { + return h.params(ctx, method) + } return nil, nil } diff --git a/src/go/plugin/agent/jobmgr/funcctl/dispatch.go b/src/go/plugin/agent/jobmgr/funcctl/dispatch.go index a38e5d5f14577e..bd645fce147805 100644 --- a/src/go/plugin/agent/jobmgr/funcctl/dispatch.go +++ b/src/go/plugin/agent/jobmgr/funcctl/dispatch.go @@ -164,12 +164,20 @@ func (c *Controller) makeMethodFuncHandler(moduleName, methodID string) function includeJobParam := methodRequiresJobParam(methodCfg) if !includeJobParam { + job, jobName, jobGen, ok := c.resolveAgentWideMethodJob(fn, moduleName) + if !ok { + return + } + c.executeMethodRequest(ctx, methodExecutionInput{ fn: fn, moduleName: moduleName, + jobName: jobName, jobLabel: moduleName, methodID: methodID, methodCfg: methodCfg, + job: job, + jobGen: jobGen, info: info, payload: payload, argValues: argValues, @@ -239,6 +247,22 @@ func (c *Controller) makeMethodFuncHandler(moduleName, methodID string) function } } +func (c *Controller) resolveAgentWideMethodJob(fn functions.Function, moduleName string) (collectorapi.RuntimeJob, string, uint64, bool) { + creator, ok := c.registry.getCreator(moduleName) + if !ok || creator.InstancePolicy != collectorapi.InstancePolicySingle { + return nil, "", 0, true + } + + jobName := moduleName + job, jobGen := c.registry.getJobWithGeneration(moduleName, jobName) + if job == nil { + c.respondError(fn, 503, "module '%s' is not running", moduleName) + return nil, "", 0, false + } + + return job, jobName, jobGen, true +} + func (c *Controller) handleMethodFuncInfo(moduleName, methodID string, fn functions.Function) { methodCfg, ok := c.registry.getMethod(moduleName, methodID) if !ok { diff --git a/src/go/plugin/framework/collectorapi/registry.go b/src/go/plugin/framework/collectorapi/registry.go index 7268c4ff8c77da..10d93d0c5507b2 100644 --- a/src/go/plugin/framework/collectorapi/registry.go +++ b/src/go/plugin/framework/collectorapi/registry.go @@ -67,8 +67,12 @@ type ( Methods func() []funcapi.MethodConfig // Optional: MethodHandler returns a handler for method requests. - // AgentWide module methods are dispatched with nil job. Job-bound module - // methods and JobMethods are dispatched with the selected running job. + // AgentWide module methods are dispatched with nil job, except that + // single-instance collectors receive their running canonical job. + // Job-bound module methods and JobMethods are dispatched with the + // selected running job. + // When the canonical single-instance job is not running, dispatch returns + // unavailable before calling MethodHandler. // The handler implements funcapi.MethodHandler interface with: // - MethodParams(ctx, method) for dynamic params // - Handle(ctx, method, params) for request handling diff --git a/src/go/plugin/go.d/collector/snmp_topology/collector.go b/src/go/plugin/go.d/collector/snmp_topology/collector.go index 5725401c26f4bd..58e670a660da3d 100644 --- a/src/go/plugin/go.d/collector/snmp_topology/collector.go +++ b/src/go/plugin/go.d/collector/snmp_topology/collector.go @@ -44,6 +44,7 @@ func New() *Collector { return &Collector{ deviceCaches: make(map[string]*topologyCache), deviceLastCollected: make(map[string]time.Time), + topologyRegistry: newTopologyRegistry(), registeredDevices: ddsnmp.DeviceRegistry.Devices, newSnmpClient: gosnmp.NewHandler, newDdSnmpColl: func(cfg ddsnmpcollector.Config) ddCollector { @@ -62,6 +63,7 @@ type ( deviceCaches map[string]*topologyCache // one cache per SNMP device deviceLastCollected map[string]time.Time // last collection time per device topologyCache *topologyCache // current device cache (set during refreshDeviceTopology) + topologyRegistry *topologyRegistry refreshMu sync.Mutex statsMu sync.RWMutex @@ -103,6 +105,9 @@ func (c *Collector) Run(ctx context.Context) error { if err := ctx.Err(); err != nil { return nil } + c.publishTrapTopologyEnrichment() + defer c.unpublishTrapTopologyEnrichment() + c.refreshTopologyRecovering(ctx) ticker := time.NewTicker(c.deviceCheckEvery()) @@ -138,11 +143,13 @@ func (c *Collector) refreshEvery() time.Duration { } func (c *Collector) Cleanup(context.Context) { + c.unpublishTrapTopologyEnrichment() + c.refreshMu.Lock() defer c.refreshMu.Unlock() for key, cache := range c.deviceCaches { - snmpTopologyRegistry.unregister(cache) + c.topologyRegistry.unregister(cache) delete(c.deviceCaches, key) delete(c.deviceLastCollected, key) } @@ -328,7 +335,7 @@ func (c *Collector) getOrCreateDeviceCache(key string) *topologyCache { if !ok { cache = newTopologyCache() c.deviceCaches[key] = cache - snmpTopologyRegistry.register(cache) + c.topologyRegistry.register(cache) } return cache } @@ -345,7 +352,7 @@ func (c *Collector) newDeviceCollectionCache(dev ddsnmp.DeviceConnectionInfo) *t func (c *Collector) pruneStaleDeviceCaches(seen map[string]bool) { for key, cache := range c.deviceCaches { if !seen[key] { - snmpTopologyRegistry.unregister(cache) + c.topologyRegistry.unregister(cache) delete(c.deviceCaches, key) delete(c.deviceLastCollected, key) } @@ -428,5 +435,12 @@ func topologyMethods() []funcapi.MethodConfig { } func topologyFunctionHandler(job collectorapi.RuntimeJob) funcapi.MethodHandler { - return &funcTopology{} + if job == nil { + return nil + } + coll, ok := job.Collector().(*Collector) + if !ok || coll == nil { + return nil + } + return &funcTopology{registry: coll.topologyRegistry} } diff --git a/src/go/plugin/go.d/collector/snmp_topology/collector_refresh_test.go b/src/go/plugin/go.d/collector/snmp_topology/collector_refresh_test.go index 95ac1cc78598fd..c588bd0f7ada94 100644 --- a/src/go/plugin/go.d/collector/snmp_topology/collector_refresh_test.go +++ b/src/go/plugin/go.d/collector/snmp_topology/collector_refresh_test.go @@ -125,23 +125,18 @@ func TestCollectorRunDoesNotPollWhenContextAlreadyCanceled(t *testing.T) { } func TestCollectorPruneStaleDeviceCachesRemovesLastDeviceCache(t *testing.T) { - previousRegistry := snmpTopologyRegistry - registry := newTopologyRegistry() - snmpTopologyRegistry = registry - t.Cleanup(func() { snmpTopologyRegistry = previousRegistry }) - coll := New() cache := newTopologyCache() coll.deviceCaches["gone:161"] = cache coll.deviceLastCollected["gone:161"] = time.Now() - registry.register(cache) + coll.topologyRegistry.register(cache) coll.registeredDevices = func() []ddsnmp.DeviceConnectionInfo { return nil } coll.refreshTopology(context.Background()) require.Empty(t, coll.deviceCaches) require.Empty(t, coll.deviceLastCollected) - require.False(t, topologyRegistryHasCache(registry, cache)) + require.False(t, topologyRegistryHasCache(coll.topologyRegistry, cache)) } func TestCollectorRefreshTopologyRecoveringHandlesPanic(t *testing.T) { @@ -305,11 +300,6 @@ func TestCollectorNewDeviceCollectionCacheUsesEffectiveDeviceCheckEvery(t *testi } func TestCollector_RefreshKeepsPublishedSnapshotWhileCollectionRuns(t *testing.T) { - previousRegistry := snmpTopologyRegistry - registry := newTopologyRegistry() - snmpTopologyRegistry = registry - t.Cleanup(func() { snmpTopologyRegistry = previousRegistry }) - ctrl := gomock.NewController(t) defer ctrl.Finish() @@ -324,7 +314,6 @@ func TestCollector_RefreshKeepsPublishedSnapshotWhileCollectionRuns(t *testing.T key := "10.0.0.10:161" published := newTopologyCache() seedPublishedEndpointSnapshot(published) - registry.register(published) started := make(chan struct{}) release := make(chan struct{}) @@ -332,6 +321,7 @@ func TestCollector_RefreshKeepsPublishedSnapshotWhileCollectionRuns(t *testing.T coll := New() coll.deviceCaches[key] = published + coll.topologyRegistry.register(published) coll.newSnmpClient = func() gosnmp.Handler { return mockHandler } coll.newDdSnmpColl = func(ddsnmpcollector.Config) ddCollector { return &blockingTopologyCollector{ diff --git a/src/go/plugin/go.d/collector/snmp_topology/func_topology.go b/src/go/plugin/go.d/collector/snmp_topology/func_topology.go index 22671414d84c28..43e5a3f8b26b2a 100644 --- a/src/go/plugin/go.d/collector/snmp_topology/func_topology.go +++ b/src/go/plugin/go.d/collector/snmp_topology/func_topology.go @@ -7,7 +7,9 @@ import "github.com/netdata/netdata/go/plugins/pkg/funcapi" // Compile-time interface check. var _ funcapi.MethodHandler = (*funcTopology)(nil) -type funcTopology struct{} +type funcTopology struct { + registry *topologyRegistry +} const ( topologyFunctionName = "snmp:topology:snmp" diff --git a/src/go/plugin/go.d/collector/snmp_topology/func_topology_handler.go b/src/go/plugin/go.d/collector/snmp_topology/func_topology_handler.go index 5134d16e2a32a0..ce975fdcc5e38f 100644 --- a/src/go/plugin/go.d/collector/snmp_topology/func_topology_handler.go +++ b/src/go/plugin/go.d/collector/snmp_topology/func_topology_handler.go @@ -18,7 +18,7 @@ func (f *funcTopology) MethodParams(_ context.Context, method string) ([]funcapi topologyNodesIdentityParamConfig(), topologyMapTypeParamConfig(), topologyInferenceStrategyParamConfig(), - topologyManagedFocusParamConfig(topologyManagedFocusParamOptions()), + topologyManagedFocusParamConfig(topologyManagedFocusParamOptions(f.registry)), topologyDepthParamConfig(), }, nil } @@ -30,13 +30,13 @@ func (f *funcTopology) Handle(_ context.Context, method string, params funcapi.R return funcapi.NotFoundResponse(method) } - if snmpTopologyRegistry == nil { + if f.registry == nil { return funcapi.UnavailableResponse("topology data not available yet, please retry after topology refresh") } options := resolveTopologyQueryOptions(params) options.ResolveDNSName = resolveTopologyReverseDNSNameCached // never block on network I/O - data, ok := snmpTopologyRegistry.snapshotWithOptions(options) + data, ok := f.registry.snapshotWithOptions(options) if !ok { return funcapi.UnavailableResponse("topology data not available yet, please retry after topology refresh") } @@ -53,13 +53,13 @@ func (f *funcTopology) Handle(_ context.Context, method string, params funcapi.R } } -func topologyManagedFocusParamOptions() []funcapi.ParamOption { - if snmpTopologyRegistry == nil { +func topologyManagedFocusParamOptions(registry *topologyRegistry) []funcapi.ParamOption { + if registry == nil { return nil } options := make([]funcapi.ParamOption, 0) - for _, target := range snmpTopologyRegistry.managedDeviceFocusTargets() { + for _, target := range registry.managedDeviceFocusTargets() { if strings.TrimSpace(target.Value) == "" { continue } diff --git a/src/go/plugin/go.d/collector/snmp_topology/func_topology_presentation_test.go b/src/go/plugin/go.d/collector/snmp_topology/func_topology_presentation_test.go index 0e41515dd8580b..8a6d630de84da0 100644 --- a/src/go/plugin/go.d/collector/snmp_topology/func_topology_presentation_test.go +++ b/src/go/plugin/go.d/collector/snmp_topology/func_topology_presentation_test.go @@ -35,5 +35,20 @@ func TestSNMPTopologyCreatorOwnsTopologyFunction(t *testing.T) { require.Equal(t, topologyMethodID, methods[0].ID) require.Equal(t, topologyFunctionName, methods[0].FunctionName) require.True(t, methods[0].AgentWide) - require.IsType(t, &funcTopology{}, creator.MethodHandler(nil)) + + coll := New() + handler := creator.MethodHandler(&topologyRuntimeJobForTest{collector: coll}) + require.IsType(t, &funcTopology{}, handler) + require.Same(t, coll.topologyRegistry, handler.(*funcTopology).registry) + require.Nil(t, creator.MethodHandler(nil)) +} + +type topologyRuntimeJobForTest struct { + collector *Collector } + +func (j *topologyRuntimeJobForTest) FullName() string { return "snmp_topology" } +func (j *topologyRuntimeJobForTest) ModuleName() string { return "snmp_topology" } +func (j *topologyRuntimeJobForTest) Name() string { return "snmp_topology" } +func (j *topologyRuntimeJobForTest) IsRunning() bool { return true } +func (j *topologyRuntimeJobForTest) Collector() any { return j.collector } diff --git a/src/go/plugin/go.d/collector/snmp_topology/func_topology_test.go b/src/go/plugin/go.d/collector/snmp_topology/func_topology_test.go index e0dc21d2a47270..61957e4af669f3 100644 --- a/src/go/plugin/go.d/collector/snmp_topology/func_topology_test.go +++ b/src/go/plugin/go.d/collector/snmp_topology/func_topology_test.go @@ -70,13 +70,7 @@ func TestTopologyMethodConfigIncludesSelectors(t *testing.T) { } func TestFuncTopology_MethodParams(t *testing.T) { - prev := snmpTopologyRegistry - t.Cleanup(func() { - snmpTopologyRegistry = prev - }) - registry := newTopologyRegistry() - snmpTopologyRegistry = registry registry.register(newTestTopologyCacheLLDP( "agent-test", time.Now().UTC(), @@ -90,7 +84,7 @@ func TestFuncTopology_MethodParams(t *testing.T) { "Gi0/2", )) - f := &funcTopology{} + f := &funcTopology{registry: registry} params, err := f.MethodParams(context.Background(), topologyMethodID) require.NoError(t, err) @@ -110,13 +104,7 @@ func TestFuncTopology_MethodParams(t *testing.T) { } func TestFuncTopology_Handle_DefaultStrictL2(t *testing.T) { - prev := snmpTopologyRegistry - t.Cleanup(func() { - snmpTopologyRegistry = prev - }) - registry := newTopologyRegistry() - snmpTopologyRegistry = registry registry.register(newTestTopologyCacheLLDP( "agent-test", time.Now().UTC(), @@ -130,7 +118,7 @@ func TestFuncTopology_Handle_DefaultStrictL2(t *testing.T) { "Gi0/2", )) - f := &funcTopology{} + f := &funcTopology{registry: registry} resp := f.Handle(context.Background(), topologyMethodID, nil) require.NotNil(t, resp) assert.Equal(t, 200, resp.Status) @@ -150,13 +138,7 @@ func TestFuncTopology_Handle_DefaultStrictL2(t *testing.T) { } func TestFuncTopology_Handle_AcceptsSelectorParams(t *testing.T) { - prev := snmpTopologyRegistry - t.Cleanup(func() { - snmpTopologyRegistry = prev - }) - registry := newTopologyRegistry() - snmpTopologyRegistry = registry registry.register(newTestTopologyCacheLLDP( "agent-test", time.Now().UTC(), @@ -170,7 +152,7 @@ func TestFuncTopology_Handle_AcceptsSelectorParams(t *testing.T) { "Gi0/2", )) - f := &funcTopology{} + f := &funcTopology{registry: registry} cfg := []funcapi.ParamConfig{ topologyNodesIdentityParamConfig(), topologyMapTypeParamConfig(), @@ -197,13 +179,7 @@ func TestFuncTopology_Handle_AcceptsSelectorParams(t *testing.T) { } func TestFuncTopology_Handle_UnknownSelectorsFallbackToDefaults(t *testing.T) { - prev := snmpTopologyRegistry - t.Cleanup(func() { - snmpTopologyRegistry = prev - }) - registry := newTopologyRegistry() - snmpTopologyRegistry = registry registry.register(newTestTopologyCacheLLDP( "agent-test", time.Now().UTC(), @@ -217,7 +193,7 @@ func TestFuncTopology_Handle_UnknownSelectorsFallbackToDefaults(t *testing.T) { "Gi0/2", )) - f := &funcTopology{} + f := &funcTopology{registry: registry} cfg := []funcapi.ParamConfig{ topologyNodesIdentityParamConfig(), topologyMapTypeParamConfig(), diff --git a/src/go/plugin/go.d/collector/snmp_topology/topology_registry.go b/src/go/plugin/go.d/collector/snmp_topology/topology_registry.go index 6a1161320addf6..3599283d23e640 100644 --- a/src/go/plugin/go.d/collector/snmp_topology/topology_registry.go +++ b/src/go/plugin/go.d/collector/snmp_topology/topology_registry.go @@ -32,8 +32,6 @@ func newTopologyRegistry() *topologyRegistry { } } -var snmpTopologyRegistry = newTopologyRegistry() - func (r *topologyRegistry) register(cache *topologyCache) { if r == nil || cache == nil { return diff --git a/src/go/plugin/go.d/collector/snmp_topology/topology_trap_enrich.go b/src/go/plugin/go.d/collector/snmp_topology/topology_trap_enrich.go index 50871141454c8c..d31df678ba9f16 100644 --- a/src/go/plugin/go.d/collector/snmp_topology/topology_trap_enrich.go +++ b/src/go/plugin/go.d/collector/snmp_topology/topology_trap_enrich.go @@ -5,6 +5,7 @@ package snmptopology import ( "sort" "strings" + "sync/atomic" ) type TrapTopologyEnrichment struct { @@ -22,27 +23,45 @@ type TrapTopologyEnrichment struct { Neighbors []string } -// TrapEnrichmentForIP returns source-device topology enrichment for a trap -// received from the given source IP. It intentionally does not infer a trap -// interface from the source IP. -func TrapEnrichmentForIP(ip string) *TrapTopologyEnrichment { - return TrapEnrichmentForSource(ip, "") +var activeTrapTopologyRegistry atomic.Pointer[topologyRegistry] + +func (c *Collector) publishTrapTopologyEnrichment() { + if c.topologyRegistry != nil { + activeTrapTopologyRegistry.Store(c.topologyRegistry) + } +} + +func (c *Collector) unpublishTrapTopologyEnrichment() { + if c.topologyRegistry != nil { + activeTrapTopologyRegistry.CompareAndSwap(c.topologyRegistry, nil) + } } // TrapEnrichmentForSource returns topology enrichment data for a trap received // from the given source IP and, when available, the trap subject ifIndex. // Interface and neighbor enrichment only use the trap ifIndex after the source // IP matches exactly one local topology cache. -// -// It copies active cache pointers under the registry lock, reads each cache -// under its own lock, and never blocks on I/O. func TrapEnrichmentForSource(ip, trapIfIndex string) *TrapTopologyEnrichment { + registry := activeTrapTopologyRegistry.Load() + if registry == nil { + return nil + } + return registry.trapEnrichmentForSource(ip, trapIfIndex) +} + +// trapEnrichmentForSource copies active cache pointers under the registry lock, +// reads each cache under its own lock, and never blocks on I/O. +func (r *topologyRegistry) trapEnrichmentForSource(ip, trapIfIndex string) *TrapTopologyEnrichment { + if r == nil { + return nil + } + ip = normalizeIPAddress(ip) if ip == "" { return nil } - caches := snmpTopologyRegistry.activeCaches() + caches := r.activeCaches() if len(caches) == 0 { return nil } diff --git a/src/go/plugin/go.d/collector/snmp_topology/topology_trap_enrich_test.go b/src/go/plugin/go.d/collector/snmp_topology/topology_trap_enrich_test.go index 3ffc3b80971808..e88029f95f5e0f 100644 --- a/src/go/plugin/go.d/collector/snmp_topology/topology_trap_enrich_test.go +++ b/src/go/plugin/go.d/collector/snmp_topology/topology_trap_enrich_test.go @@ -3,7 +3,10 @@ package snmptopology import ( + "context" + "errors" "testing" + "time" "github.com/stretchr/testify/require" ) @@ -86,14 +89,16 @@ func TestTopologyCacheTrapEnrichmentForSourceIncludesLocalDeviceIdentity(t *test require.Equal(t, "vnode-node-id", enrich.SourceVnodeID) } -func TestTrapEnrichmentForSourceUsesGlobalRegistry(t *testing.T) { +func TestTrapEnrichmentForSourceUsesActiveRegistry(t *testing.T) { + registry := newTopologyRegistry() + publishTrapTopologyRegistryForTest(t, registry) + cache := newTopologyCache() cache.localDevice.ManagementIP = "192.0.2.20" cache.ifNamesByIndex["11"] = "Gi0/11" cache.lldpRemotes["11:1"] = &lldpRemote{sysName: "dist-c"} - snmpTopologyRegistry.register(cache) - defer snmpTopologyRegistry.unregister(cache) + registry.register(cache) enrich := TrapEnrichmentForSource("192.0.2.20", "11") require.NotNil(t, enrich) @@ -106,7 +111,10 @@ func TestTrapEnrichmentForSourceUsesGlobalRegistry(t *testing.T) { require.Equal(t, "Gi0/11", mapped.Interface) } -func TestTrapEnrichmentForSourceAmbiguousGlobalRegistryMatchDoesNotEnrich(t *testing.T) { +func TestTrapEnrichmentForSourceAmbiguousActiveRegistryMatchDoesNotEnrich(t *testing.T) { + registry := newTopologyRegistry() + publishTrapTopologyRegistryForTest(t, registry) + cacheA := newTopologyCache() cacheA.localDevice.ManagementIP = "192.0.2.20" cacheA.ifNamesByIndex["11"] = "Gi0/11" @@ -114,10 +122,8 @@ func TestTrapEnrichmentForSourceAmbiguousGlobalRegistryMatchDoesNotEnrich(t *tes cacheB.localDevice.ManagementIP = "192.0.2.20" cacheB.ifNamesByIndex["11"] = "Gi0/11" - snmpTopologyRegistry.register(cacheA) - snmpTopologyRegistry.register(cacheB) - defer snmpTopologyRegistry.unregister(cacheA) - defer snmpTopologyRegistry.unregister(cacheB) + registry.register(cacheA) + registry.register(cacheB) enrich := TrapEnrichmentForSource("192.0.2.20", "11") require.NotNil(t, enrich) @@ -126,3 +132,65 @@ func TestTrapEnrichmentForSourceAmbiguousGlobalRegistryMatchDoesNotEnrich(t *tes require.Empty(t, enrich.Interface) require.Empty(t, enrich.Neighbors) } + +func TestCollectorRunPublishesAndClearsTrapTopologyRegistry(t *testing.T) { + previous := activeTrapTopologyRegistry.Swap(nil) + t.Cleanup(func() { activeTrapTopologyRegistry.Store(previous) }) + + coll := New() + coll.UpdateEvery = 3600 + coll.registeredDevices = nil + + ctx, cancel := context.WithCancel(context.Background()) + errCh := make(chan error, 1) + go func() { + errCh <- coll.Run(ctx) + }() + + stopped := false + stopRunner := func() error { + if stopped { + return nil + } + stopped = true + cancel() + select { + case err := <-errCh: + return err + case <-time.After(time.Second): + return errors.New("runner did not stop") + } + } + defer func() { + require.NoError(t, stopRunner()) + }() + + require.Eventually(t, func() bool { + return activeTrapTopologyRegistry.Load() == coll.topologyRegistry + }, time.Second, 10*time.Millisecond) + + require.NoError(t, stopRunner()) + require.Nil(t, activeTrapTopologyRegistry.Load()) +} + +func TestCollectorCleanupDoesNotClearNewerTrapTopologyRegistry(t *testing.T) { + previous := activeTrapTopologyRegistry.Swap(nil) + t.Cleanup(func() { activeTrapTopologyRegistry.Store(previous) }) + + oldColl := New() + newColl := New() + activeTrapTopologyRegistry.Store(newColl.topologyRegistry) + + oldColl.Cleanup(context.Background()) + + require.Same(t, newColl.topologyRegistry, activeTrapTopologyRegistry.Load()) +} + +func publishTrapTopologyRegistryForTest(t *testing.T, registry *topologyRegistry) { + t.Helper() + + previous := activeTrapTopologyRegistry.Swap(registry) + t.Cleanup(func() { + activeTrapTopologyRegistry.CompareAndSwap(registry, previous) + }) +} From e0d3673475512f95c16e45a962c7597069b5ae53 Mon Sep 17 00:00:00 2001 From: Stelios Fragkakis <52996999+stelfrag@users.noreply.github.com> Date: Thu, 18 Jun 2026 18:06:44 +0300 Subject: [PATCH 4/9] Add chart to monitor in-flight ACLK messages (#22757) * feat(aclk): read MQTT 5.0 Receive Maximum and chart in-flight QoS1 Extract the broker's Receive Maximum from CONNACK, expose it in aclk state (text/JSON), and add an always-on netdata.aclk_mqtt_inflight pulse chart (in flight vs receive maximum). Observation only, no pacing. * fix(aclk): use atomic access for rx_maximum Read in mqtt_ng_get_stats (pulse thread) races the CONNACK write (ACLK main thread). Make both stores and the load __ATOMIC_RELAXED, matching the other stats fields. --- src/aclk/aclk.c | 8 ++++-- src/aclk/mqtt_websockets/common_public.h | 2 ++ src/aclk/mqtt_websockets/mqtt_ng.c | 12 +++++++++ src/daemon/pulse/pulse-network.c | 34 ++++++++++++++++++++++++ 4 files changed, 54 insertions(+), 2 deletions(-) diff --git a/src/aclk/aclk.c b/src/aclk/aclk.c index 42342199155825..1daabe04fb347e 100644 --- a/src/aclk/aclk.c +++ b/src/aclk/aclk.c @@ -1128,8 +1128,9 @@ char *aclk_state(void) } if (aclk_is_online) { - buffer_sprintf(wb, "Received Cloud MQTT Messages: %d\nMQTT Messages Confirmed by Remote Broker (PUBACKs): %d\nPending PUBACKS: %d\n", - aclk_rcvd_cloud_msgs, aclk_pubacks_per_conn, aclk_stats.mqtt.packets_waiting_puback); + buffer_sprintf(wb, "Received Cloud MQTT Messages: %d\nMQTT Messages Confirmed by Remote Broker (PUBACKs): %d\nPending PUBACKS: %d\nServer Receive Maximum: %u\n", + aclk_rcvd_cloud_msgs, aclk_pubacks_per_conn, aclk_stats.mqtt.packets_waiting_puback, + (unsigned)aclk_stats.mqtt.rx_maximum); RRDHOST *host; rrd_rdlock(); @@ -1269,6 +1270,9 @@ char *aclk_state_json(void) tmp = json_object_new_int((int32_t) aclk_stats.mqtt.packets_waiting_puback); json_object_object_add(msg, "pending-mqtt-pubacks", tmp); + tmp = json_object_new_int((int32_t) aclk_stats.mqtt.rx_maximum); + json_object_object_add(msg, "server-receive-maximum", tmp); + tmp = json_object_new_int(aclk_connection_counter > 0 ? (aclk_connection_counter - 1) : 0); json_object_object_add(msg, "reconnect-count", tmp); diff --git a/src/aclk/mqtt_websockets/common_public.h b/src/aclk/mqtt_websockets/common_public.h index 31ea4400128462..ef36ff8347168f 100644 --- a/src/aclk/mqtt_websockets/common_public.h +++ b/src/aclk/mqtt_websockets/common_public.h @@ -27,6 +27,8 @@ struct mqtt_ng_stats { int tx_messages_sent; int rx_messages_rcvd; int packets_waiting_puback; + // MQTT 5.0 server Receive Maximum from CONNACK (max concurrent unacked QoS1/2); 65535 if unset + uint16_t rx_maximum; size_t tx_buffer_used; size_t tx_buffer_free; size_t tx_buffer_size; diff --git a/src/aclk/mqtt_websockets/mqtt_ng.c b/src/aclk/mqtt_websockets/mqtt_ng.c index 88c690cc0a4593..c7672268153b4f 100644 --- a/src/aclk/mqtt_websockets/mqtt_ng.c +++ b/src/aclk/mqtt_websockets/mqtt_ng.c @@ -250,6 +250,9 @@ struct mqtt_ng_client { c_rhash rx_aliases; size_t max_msg_size; + + // MQTT 5.0 server Receive Maximum (CONNACK prop 0x21); absent => 65535 default [MQTT-3.2.2.3.3] + uint16_t rx_maximum; }; usec_t publish_latency; @@ -630,6 +633,9 @@ struct mqtt_ng_client *mqtt_ng_init(struct mqtt_ng_init *settings) client->tx_topic_aliases.stoi_dict = TX_ALIASES_INITIALIZE(); client->tx_topic_aliases.idx_max = UINT16_MAX; + // MQTT 5.0 default Receive Maximum when the server omits the property [MQTT-3.2.2.3.3] + __atomic_store_n(&client->rx_maximum, UINT16_MAX, __ATOMIC_RELAXED); + // TODO just embed the struct into mqtt_ng_client client->parser.received_data = settings->data_in; client->send_fnc_ptr = settings->data_out_fnc; @@ -2144,6 +2150,11 @@ int handle_incoming_traffic(struct mqtt_ng_client *client) client->max_msg_size = prop->data.uint32; } + if ((prop = get_property_by_id(client->parser.properties_parser.head, MQTT_PROP_RECEIVE_MAX)) != NULL) { + nd_log(NDLS_DAEMON, NDLP_INFO, "ACLK: MQTT server receive maximum is %" PRIu16, prop->data.uint16); + __atomic_store_n(&client->rx_maximum, prop->data.uint16, __ATOMIC_RELAXED); + } + if (client->connack_callback) client->connack_callback(client->user_ctx, client->parser.mqtt_packet.connack.reason_code); if (!client->parser.mqtt_packet.connack.reason_code) { @@ -2300,6 +2311,7 @@ void mqtt_ng_get_stats(struct mqtt_ng_client *client, struct mqtt_ng_stats *stat stats->tx_messages_sent = __atomic_load_n(&client->stats.tx_messages_sent, __ATOMIC_RELAXED); stats->rx_messages_rcvd = __atomic_load_n(&client->stats.rx_messages_rcvd, __ATOMIC_RELAXED); stats->packets_waiting_puback = __atomic_load_n(&client->stats.packets_waiting_puback, __ATOMIC_RELAXED); + stats->rx_maximum = __atomic_load_n(&client->rx_maximum, __ATOMIC_RELAXED); stats->tx_bytes_queued = 0; stats->tx_buffer_reclaimable = 0; diff --git a/src/daemon/pulse/pulse-network.c b/src/daemon/pulse/pulse-network.c index ee5024cb1b2f5a..e173d45c3c46ea 100644 --- a/src/daemon/pulse/pulse-network.c +++ b/src/daemon/pulse/pulse-network.c @@ -335,6 +335,40 @@ void pulse_network_do(bool extended __maybe_unused) { pulse_aclk_time_heatmap(); + { + // In-flight QoS1 messages vs the broker's MQTT 5.0 Receive Maximum. + // When "in flight" approaches "receive maximum" the agent is at the + // broker's window limit. Always available (not extended) since it is + // the primary signal for the MQTT 5.0 Receive Maximum behavior. + static RRDSET *st_aclk_inflight = NULL; + static RRDDIM *rd_in_flight = NULL, *rd_receive_max = NULL; + + if (unlikely(!st_aclk_inflight)) { + st_aclk_inflight = rrdset_create_localhost( + "netdata", + "aclk_mqtt_inflight", + NULL, + PULSE_NETWORK_CHART_FAMILY, + "netdata.aclk_mqtt_inflight", + "Netdata ACLK MQTT In-Flight QoS1 Window", + "messages", + "netdata", + "pulse", + PULSE_NETWORK_CHART_PRIORITY + 2, + localhost->rrd_update_every, + RRDSET_TYPE_LINE); + + rrdlabels_add(st_aclk_inflight->rrdlabels, "endpoint", "aclk", RRDLABEL_SRC_AUTO); + + rd_in_flight = rrddim_add(st_aclk_inflight, "in flight", NULL, 1, 1, RRD_ALGORITHM_ABSOLUTE); + rd_receive_max = rrddim_add(st_aclk_inflight, "receive maximum", NULL, 1, 1, RRD_ALGORITHM_ABSOLUTE); + } + + rrddim_set_by_pointer(st_aclk_inflight, rd_in_flight, (collected_number)t.mqtt.packets_waiting_puback); + rrddim_set_by_pointer(st_aclk_inflight, rd_receive_max, (collected_number)t.mqtt.rx_maximum); + rrdset_done(st_aclk_inflight); + } + if(extended) { static RRDSET *st_aclk_queue_size = NULL; static RRDDIM *rd_messages = NULL; From d04bce705c8bfac7963992d685f532e105db3f28 Mon Sep 17 00:00:00 2001 From: Costa Tsaousis Date: Thu, 18 Jun 2026 18:08:13 +0300 Subject: [PATCH 5/9] streaming: add receiver disconnect diagnostics and per-child join keys (#22769) --- src/streaming/stream-receiver-internals.h | 1 + src/streaming/stream-receiver.c | 36 +++++++++++++++++-- .../functions/function-netdata-streaming.c | 23 ++++++++++++ 3 files changed, 57 insertions(+), 3 deletions(-) diff --git a/src/streaming/stream-receiver-internals.h b/src/streaming/stream-receiver-internals.h index 856601b687418a..c68a98d7b0f87e 100644 --- a/src/streaming/stream-receiver-internals.h +++ b/src/streaming/stream-receiver-internals.h @@ -74,6 +74,7 @@ struct receiver_state { nd_poll_event_t wanted; usec_t last_traffic_ut; + size_t bytes_received; // raw socket bytes received on this connection (diagnostics) struct pollfd_meta meta; } thread; diff --git a/src/streaming/stream-receiver.c b/src/streaming/stream-receiver.c index 624c3e0360e87b..ca9396eb399ad3 100644 --- a/src/streaming/stream-receiver.c +++ b/src/streaming/stream-receiver.c @@ -153,6 +153,7 @@ static ssize_t receiver_read_uncompressed(struct receiver_state *r) { r->thread.uncompressed.read_len += bytes; r->thread.uncompressed.read_buffer[r->thread.uncompressed.read_len] = '\0'; + r->thread.bytes_received += bytes; pulse_stream_received_bytes(bytes); } @@ -271,6 +272,7 @@ static ssize_t receiver_read_compressed(struct receiver_state *r) { if(bytes > 0) { r->thread.compressed.used += bytes; + r->thread.bytes_received += bytes; worker_set_metric(WORKER_RECEIVER_JOB_BYTES_READ, (NETDATA_DOUBLE)bytes); pulse_stream_received_bytes(bytes); } @@ -423,6 +425,11 @@ void stream_receiver_move_to_running_unsafe(struct stream_thread *sth, struct re rpt->thread.compressed.enabled = stream_decompression_initialize(rpt); buffered_reader_init(&rpt->thread.uncompressed); + // start fresh at admission: the no-traffic timeout must be measured from when we start + // reading (now), not from when the connection was accepted and queued. + rpt->thread.bytes_received = 0; + rpt->thread.last_traffic_ut = now_monotonic_usec(); + rpt->thread.line_buffer = buffer_create(sizeof(rpt->thread.uncompressed.read_buffer), NULL); // help preferred_sender_buffer() select the right buffer @@ -516,16 +523,39 @@ static void stream_receiver_remove_internal(struct stream_thread *sth, struct re if(parser) count = parser->user.data_collections_count; + // gather diagnostics for a single, uniform disconnect line (same fields for every reason): + // bytes_in distinguishes "the child sent nothing" (network/silent) from "sent but unparsed"; + // iface exposes the child's connection type (e.g. ppp0 cellular vs eth0) for log-side pivots. + size_t bytes_out = 0; + spinlock_lock(&rpt->thread.send_to_child.spinlock); + if(rpt->thread.send_to_child.scb) + bytes_out = stream_circular_buffer_stats_unsafe(rpt->thread.send_to_child.scb)->bytes_sent; + spinlock_unlock(&rpt->thread.send_to_child.spinlock); + + char iface[64] = ""; + if(rpt->host && rpt->host->rrdlabels) + rrdlabels_get_value_strcpyz(rpt->host->rrdlabels, iface, sizeof(iface), "_net_default_iface"); + + time_t connected_s = rpt->connected_since_s ? (now_realtime_sec() - rpt->connected_since_s) : 0; + long long idle_s = (long long)((now_monotonic_usec() - rpt->thread.last_traffic_ut) / USEC_PER_SEC); + double repl_pct = rpt->host ? rpt->host->stream.rcv.status.replication.percent : 0.0; + errno_clear(); nd_log(NDLS_DAEMON, NDLP_ERR, - "STREAM RCV[%zu] '%s' [from [%s]:%s]: " - "receiver disconnected (after %zu received messages): %s" + "STREAM RCV[%zu] '%s' [from [%s]:%s]: receiver disconnected: " + "reason=\"%s\" msgs=%zu bytes_in=%zu bytes_out=%zu connected=%llds idle=%llds repl=%.0f%% iface=%s" , sth->id , rpt->hostname ? rpt->hostname : "-" , rpt->remote_ip ? rpt->remote_ip : "-" , rpt->remote_port ? rpt->remote_port : "-" + , stream_handshake_error_to_string(reason) , count - , stream_handshake_error_to_string(reason)); + , rpt->thread.bytes_received + , bytes_out + , (long long)connected_s + , idle_s + , repl_pct + , iface[0] ? iface : "-"); internal_fatal(META_GET(&sth->run.meta, (Word_t)&rpt->thread.meta) == NULL, "Receiver to be removed is not found in the list of receivers"); diff --git a/src/web/api/functions/function-netdata-streaming.c b/src/web/api/functions/function-netdata-streaming.c index f138de69a5204b..913d17c00d0504 100644 --- a/src/web/api/functions/function-netdata-streaming.c +++ b/src/web/api/functions/function-netdata-streaming.c @@ -309,6 +309,16 @@ int function_netdata_streaming(BUFFER *wb, const char *function __maybe_unused, buffer_json_add_array_item_string(wb, NULL); // MlSilenced } + // MachineGUID + NodeID — hidden columns, to join streaming rows to the node inventory + buffer_json_add_array_item_string(wb, host->machine_guid); // MachineGUID + if(!UUIDiszero(host->node_id)) { + char node_id_str[UUID_STR_LEN]; + uuid_unparse_lower(host->node_id.uuid, node_id_str); + buffer_json_add_array_item_string(wb, node_id_str); // NodeID + } + else + buffer_json_add_array_item_string(wb, NULL); // NodeID + // close buffer_json_array_close(wb); } @@ -879,6 +889,19 @@ int function_netdata_streaming(BUFFER *wb, const char *function __maybe_unused, RRDF_FIELD_SORT_DESCENDING, NULL, RRDF_FIELD_SUMMARY_SUM, RRDF_FIELD_FILTER_RANGE, RRDF_FIELD_OPTS_NONE, NULL); + + // MachineGUID + NodeID — hidden, used to join streaming rows to the node inventory + buffer_rrdf_table_add_field(wb, field_id++, "MachineGUID", "Machine GUID", + RRDF_FIELD_TYPE_STRING, RRDF_FIELD_VISUAL_VALUE, RRDF_FIELD_TRANSFORM_NONE, + 0, NULL, NAN, RRDF_FIELD_SORT_ASCENDING, NULL, + RRDF_FIELD_SUMMARY_COUNT, RRDF_FIELD_FILTER_MULTISELECT, + RRDF_FIELD_OPTS_NONE, NULL); + + buffer_rrdf_table_add_field(wb, field_id++, "NodeID", "Cloud Node ID", + RRDF_FIELD_TYPE_STRING, RRDF_FIELD_VISUAL_VALUE, RRDF_FIELD_TRANSFORM_NONE, + 0, NULL, NAN, RRDF_FIELD_SORT_ASCENDING, NULL, + RRDF_FIELD_SUMMARY_COUNT, RRDF_FIELD_FILTER_MULTISELECT, + RRDF_FIELD_OPTS_NONE, NULL); } buffer_json_object_close(wb); // columns buffer_json_member_add_string(wb, "default_sort_column", "Node"); From ba4b2b2c8aee194fa924ae9cd2cf984674fb2e67 Mon Sep 17 00:00:00 2001 From: Stelios Fragkakis <52996999+stelfrag@users.noreply.github.com> Date: Thu, 18 Jun 2026 19:30:54 +0300 Subject: [PATCH 6/9] Adjust ML model deletion query (#22763) Fix pruning logic in ML database to ensure proper row deletion behavior without relying on DELETE with LIMIT --- src/ml/ml.cc | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/src/ml/ml.cc b/src/ml/ml.cc index 24258b48f735ea..cc4181c23e2cf0 100644 --- a/src/ml/ml.cc +++ b/src/ml/ml.cc @@ -346,7 +346,7 @@ const char *db_models_delete = const char *db_models_prune = "DELETE FROM models " - "WHERE after < @after LIMIT @n;"; + "WHERE rowid IN (SELECT rowid FROM models WHERE after < @after LIMIT @n);"; static int ml_dimension_add_model(const nd_uuid_t *metric_uuid, const ml_kmeans_inlined_t *inlined_km) From 38fed237553984e2d590b92ba27a7ab3aee175f0 Mon Sep 17 00:00:00 2001 From: Ilya Mashchenko Date: Thu, 18 Jun 2026 19:33:30 +0300 Subject: [PATCH 7/9] refactor(go.d/snmp): inject shared SNMP state (#22770) * refactor(go.d/snmp): inject shared SNMP state * refactor(go.d/snmp): tighten shared state constructors --- .agents/sow/specs/snmp-traps/netdata.md | 2 +- src/go/plugin/go.d/collector/init.go | 19 ++- src/go/plugin/go.d/collector/init_test.go | 87 ++++++++++++ .../collector/snmp/bgp_typed_metrics_test.go | 8 +- .../plugin/go.d/collector/snmp/charts_test.go | 4 +- src/go/plugin/go.d/collector/snmp/collect.go | 2 +- .../plugin/go.d/collector/snmp/collector.go | 28 +++- .../snmp/collector_licensing_edge_test.go | 2 +- .../go.d/collector/snmp/collector_test.go | 99 +++++++++++--- .../snmp/ddsnmp/device_registry_test.go | 83 ------------ .../{device_registry.go => device_store.go} | 124 ++++++++---------- .../snmp/ddsnmp/device_store_test.go | 100 ++++++++++++++ ...ogy_device_registry.go => device_state.go} | 13 +- .../collector/snmp/func_bgp_peers_test.go | 4 +- .../go.d/collector/snmp/test_helpers_test.go | 9 ++ .../collector/snmp_topology/charts_test.go | 2 +- .../go.d/collector/snmp_topology/collector.go | 53 +++++--- .../snmp_topology/collector_refresh_test.go | 87 ++++++------ .../go.d/collector/snmp_topology/config.go | 2 +- .../func_topology_presentation_test.go | 24 +++- .../collector/snmp_topology/metrix_test.go | 19 ++- .../topology_integration_test.go | 6 +- .../topology_test_helpers_test.go | 22 ++++ .../snmp_topology/topology_trap_enrich.go | 27 ++-- .../topology_trap_enrich_test.go | 59 ++++----- .../go.d/collector/snmp_traps/collector.go | 38 +++++- .../snmp_traps/collector_e2e_test.go | 2 +- .../go.d/collector/snmp_traps/enrich.go | 24 +++- .../go.d/collector/snmp_traps/enrich_test.go | 106 +++++++-------- .../collector/snmp_traps/func_logs_test.go | 5 +- .../go.d/collector/snmp_traps/init_test.go | 74 +++++++---- .../collector/snmp_traps/listener_test.go | 2 +- .../collector/snmp_traps/pipeline_test.go | 12 +- .../go.d/collector/snmp_traps/profile_test.go | 10 +- .../collector/snmp_traps/test_helpers_test.go | 20 +++ 35 files changed, 762 insertions(+), 416 deletions(-) create mode 100644 src/go/plugin/go.d/collector/init_test.go delete mode 100644 src/go/plugin/go.d/collector/snmp/ddsnmp/device_registry_test.go rename src/go/plugin/go.d/collector/snmp/ddsnmp/{device_registry.go => device_store.go} (53%) create mode 100644 src/go/plugin/go.d/collector/snmp/ddsnmp/device_store_test.go rename src/go/plugin/go.d/collector/snmp/{topology_device_registry.go => device_state.go} (83%) create mode 100644 src/go/plugin/go.d/collector/snmp/test_helpers_test.go diff --git a/.agents/sow/specs/snmp-traps/netdata.md b/.agents/sow/specs/snmp-traps/netdata.md index 78d7fe5757cfd7..05feb2568906af 100644 --- a/.agents/sow/specs/snmp-traps/netdata.md +++ b/.agents/sow/specs/snmp-traps/netdata.md @@ -784,7 +784,7 @@ _HOSTNAME= 0 } -func init() { - collectorapi.Register("snmp_traps", collectorapi.Creator{ +// Register registers the SNMP traps collector with shared SNMP-family enrichment state. +func Register(deviceStore *ddsnmp.DeviceStore, topologyEnricher *snmptopology.TrapEnrichmentHandle) { + collectorapi.Register("snmp_traps", newCreator(deviceStore, topologyEnricher)) +} + +func newCreator(deviceStore *ddsnmp.DeviceStore, topologyEnricher *snmptopology.TrapEnrichmentHandle) collectorapi.Creator { + if deviceStore == nil { + panic("snmp_traps Register requires a non-nil device store") + } + if topologyEnricher == nil { + panic("snmp_traps Register requires a non-nil trap enrichment handle") + } + return collectorapi.Creator{ JobConfigSchema: configSchema, Defaults: collectorapi.Defaults{ UpdateEvery: 1, }, - CreateV2: func() collectorapi.CollectorV2 { return New() }, + CreateV2: func() collectorapi.CollectorV2 { return New(deviceStore, topologyEnricher) }, Config: func() any { return &Config{} }, Methods: snmpTrapsMethods, MethodHandler: snmpTrapsMethodHandler, - }) + } } -func New() *Collector { +// New returns an SNMP traps collector using the provided SNMP-family enrichment state. +func New(deviceStore *ddsnmp.DeviceStore, topologyEnricher *snmptopology.TrapEnrichmentHandle) *Collector { + if deviceStore == nil { + panic("snmp_traps New requires a non-nil device store") + } + if topologyEnricher == nil { + panic("snmp_traps New requires a non-nil trap enrichment handle") + } store := metrix.NewCollectorStore() return &Collector{ @@ -63,7 +83,9 @@ func New() *Collector { ReceiveBuffer: defaultListenerReceiveBuffer, }, }, - store: store, + store: store, + deviceLookup: deviceStore, + topologyEnricher: topologyEnricher, } } @@ -75,6 +97,8 @@ type Collector struct { trapWriter TrapWriter journalDir string store metrix.CollectorStore + deviceLookup deviceLookup + topologyEnricher trapTopologyEnricher jobName string vnode string versions map[SnmpVersion]struct{} @@ -636,7 +660,7 @@ func (c *Collector) handlePacket(data []byte, peerIP net.IP, conn *net.UDPConn, entry := trapEntryFromPDU(c.jobName, pdu, td, time.Now().UnixMicro(), monotonicUsec()) entry.PacketSequence = packetSequence - enrichTrapEntry(entry, c.reverseDNSEnabled, c.reverseDNS) + c.enrichTrapEntry(entry, c.reverseDNSEnabled, c.reverseDNS) renderTrapEntryTemplates(entry, td) if unknownOID { c.incTrapError("unknown_oid") diff --git a/src/go/plugin/go.d/collector/snmp_traps/collector_e2e_test.go b/src/go/plugin/go.d/collector/snmp_traps/collector_e2e_test.go index ee738679215996..a6a63aa32a42c9 100644 --- a/src/go/plugin/go.d/collector/snmp_traps/collector_e2e_test.go +++ b/src/go/plugin/go.d/collector/snmp_traps/collector_e2e_test.go @@ -17,7 +17,7 @@ func TestCollectorReplayPcapThroughListenerToJournal(t *testing.T) { withTestCacheDir(t) port := freeUDPPort(t) - c := New() + c := newTestSNMPTrapsCollector() c.SetJobName("e2e") c.Listen.Endpoints = []EndpointConfig{{Protocol: "udp", Address: "127.0.0.1", Port: port}} c.Versions = []string{"v2c"} diff --git a/src/go/plugin/go.d/collector/snmp_traps/enrich.go b/src/go/plugin/go.d/collector/snmp_traps/enrich.go index 990b2963476caa..e738b1bd3992d4 100644 --- a/src/go/plugin/go.d/collector/snmp_traps/enrich.go +++ b/src/go/plugin/go.d/collector/snmp_traps/enrich.go @@ -278,10 +278,19 @@ type deviceEnrichment struct { matches int } -var trapTopologyEnrichmentForSource = snmptopology.TrapEnrichmentForSource +type deviceLookup interface { + DevicesByHostname(hostname string) []ddsnmp.DeviceConnectionInfo +} + +type trapTopologyEnricher interface { + EnrichmentForSource(ip, trapIfIndex string) *snmptopology.TrapTopologyEnrichment +} -func resolveDeviceEnrichment(sourceIP string) deviceEnrichment { - devices := ddsnmp.DeviceRegistry.DevicesByHostname(sourceIP) +func (c *Collector) resolveDeviceEnrichment(sourceIP string) deviceEnrichment { + if c == nil || c.deviceLookup == nil { + return deviceEnrichment{} + } + devices := c.deviceLookup.DevicesByHostname(sourceIP) enrich := deviceEnrichment{matches: len(devices)} if len(devices) != 1 { return enrich @@ -305,7 +314,7 @@ func resolveDeviceEnrichment(sourceIP string) deviceEnrichment { return enrich } -func enrichTrapEntry(entry *TrapEntry, useReverseDNS bool, dns *reverseDNSResolver) { +func (c *Collector) enrichTrapEntry(entry *TrapEntry, useReverseDNS bool, dns *reverseDNSResolver) { if entry == nil { return } @@ -324,7 +333,7 @@ func enrichTrapEntry(entry *TrapEntry, useReverseDNS bool, dns *reverseDNSResolv audit.Source = &TrapSourceAudit{Selected: sourceIP, Method: "entry_source"} } - enrich := resolveDeviceEnrichment(sourceIP) + enrich := c.resolveDeviceEnrichment(sourceIP) audit.Registry = &TrapEnrichmentLookup{ Key: sourceIP, Status: lookupStatus(enrich.matches), @@ -364,7 +373,10 @@ func enrichTrapEntry(entry *TrapEntry, useReverseDNS bool, dns *reverseDNSResolv addTrapEnrichmentApplied(audit, "TRAP_INTERFACE", iface) } - topo := trapTopologyEnrichmentForSource(sourceIP, trapIfIndex) + var topo *snmptopology.TrapTopologyEnrichment + if c != nil && c.topologyEnricher != nil { + topo = c.topologyEnricher.EnrichmentForSource(sourceIP, trapIfIndex) + } topologyTrusted := topo != nil && topo.DeviceStatus == "matched" if topo != nil { audit.Topology = &TrapEnrichmentLookup{ diff --git a/src/go/plugin/go.d/collector/snmp_traps/enrich_test.go b/src/go/plugin/go.d/collector/snmp_traps/enrich_test.go index 90922ace5a1dbe..e0cb2b0c9a2594 100644 --- a/src/go/plugin/go.d/collector/snmp_traps/enrich_test.go +++ b/src/go/plugin/go.d/collector/snmp_traps/enrich_test.go @@ -12,15 +12,15 @@ import ( ) func TestEnrichTrapEntryHostnamePriority(t *testing.T) { + c, store := newTestTrapEnrichmentCollector(nil) regKey := "key:10.1.2.3:162" - ddsnmp.DeviceRegistry.Register(regKey, ddsnmp.DeviceConnectionInfo{ + store.Register(regKey, ddsnmp.DeviceConnectionInfo{ Hostname: "10.1.2.3", SysName: "core-sw-01", VnodeHostname: "core-sw.mydc.example.com", Vendor: "cisco", VnodeGUID: "8f72c1e2-3a4b-5c6d-7e8f-9a0b1c2d3e4f", }) - defer ddsnmp.DeviceRegistry.Unregister(regKey) dns := newReverseDNSResolver() @@ -49,7 +49,7 @@ func TestEnrichTrapEntryHostnamePriority(t *testing.T) { entry := &TrapEntry{ SourceIP: tc.sourceIP, } - enrichTrapEntry(entry, tc.useReverseDNS, dns) + c.enrichTrapEntry(entry, tc.useReverseDNS, dns) if entry.DeviceHostname != tc.wantHostname { t.Errorf("DeviceHostname = %q, want %q", entry.DeviceHostname, tc.wantHostname) @@ -101,8 +101,7 @@ func TestEnrichTrapEntryRegistryHostnameWinsOverTopologyAndReverseDNS(t *testing }, } - prev := trapTopologyEnrichmentForSource - trapTopologyEnrichmentForSource = func(ip, ifIndex string) *snmptopology.TrapTopologyEnrichment { + topologyEnricher := testTrapTopologyEnricher(func(ip, ifIndex string) *snmptopology.TrapTopologyEnrichment { vnodeID := "topology-vnode-id" if ip == "10.1.2.6" { vnodeID = "registry-vnode-id" @@ -120,14 +119,13 @@ func TestEnrichTrapEntryRegistryHostnameWinsOverTopologyAndReverseDNS(t *testing NeighborStatus: "matched", Neighbors: []string{"topo-neighbor"}, } - } - t.Cleanup(func() { trapTopologyEnrichmentForSource = prev }) + }) for tcName, tc := range tests { t.Run(tcName, func(t *testing.T) { + c, store := newTestTrapEnrichmentCollector(topologyEnricher) regKey := "key:" + tc.info.Hostname + ":162" - ddsnmp.DeviceRegistry.Register(regKey, tc.info) - defer ddsnmp.DeviceRegistry.Unregister(regKey) + store.Register(regKey, tc.info) dns := newReverseDNSResolver() dns.cache[tc.info.Hostname] = reverseDNSCacheEntry{ @@ -142,7 +140,7 @@ func TestEnrichTrapEntryRegistryHostnameWinsOverTopologyAndReverseDNS(t *testing {Name: "ifIndex", OID: ifIndexOIDPrefix + ".1", Type: "InterfaceIndex", Value: int64(1)}, }, } - enrichTrapEntry(entry, true, dns) + c.enrichTrapEntry(entry, true, dns) if entry.DeviceHostname != tc.wantHost { t.Errorf("DeviceHostname = %q, want %q", entry.DeviceHostname, tc.wantHost) @@ -164,16 +162,16 @@ func TestEnrichTrapEntryRegistryHostnameWinsOverTopologyAndReverseDNS(t *testing } func TestEnrichTrapEntrySysNameOverVnodeUnknown(t *testing.T) { + c, store := newTestTrapEnrichmentCollector(nil) regKey := "key:10.1.2.4:162" - ddsnmp.DeviceRegistry.Register(regKey, ddsnmp.DeviceConnectionInfo{ + store.Register(regKey, ddsnmp.DeviceConnectionInfo{ Hostname: "10.1.2.4", SysName: "real-switch", VnodeHostname: "unknown", }) - defer ddsnmp.DeviceRegistry.Unregister(regKey) entry := &TrapEntry{SourceIP: "10.1.2.4"} - enrichTrapEntry(entry, false, nil) + c.enrichTrapEntry(entry, false, nil) if entry.DeviceHostname != "real-switch" { t.Errorf("DeviceHostname = %q, want real-switch (unknown vnode hostname treated as unresolved)", entry.DeviceHostname) @@ -181,24 +179,25 @@ func TestEnrichTrapEntrySysNameOverVnodeUnknown(t *testing.T) { } func TestEnrichTrapEntryEmptySysNameSkipped(t *testing.T) { + c, store := newTestTrapEnrichmentCollector(nil) regKey := "key:10.1.2.5:162" - ddsnmp.DeviceRegistry.Register(regKey, ddsnmp.DeviceConnectionInfo{ + store.Register(regKey, ddsnmp.DeviceConnectionInfo{ Hostname: "10.1.2.5", SysName: "", }) - defer ddsnmp.DeviceRegistry.Unregister(regKey) entry := &TrapEntry{SourceIP: "10.1.2.5"} - enrichTrapEntry(entry, false, nil) + c.enrichTrapEntry(entry, false, nil) if entry.DeviceHostname != "" { t.Errorf("DeviceHostname = %q, want empty (empty sysName treated as unresolved)", entry.DeviceHostname) } } -func TestEnrichTrapEntryNoDeviceRegistryMatch(t *testing.T) { +func TestEnrichTrapEntryNoDeviceStoreMatch(t *testing.T) { + c, _ := newTestTrapEnrichmentCollector(nil) entry := &TrapEntry{SourceIP: "172.16.0.99"} - enrichTrapEntry(entry, false, nil) + c.enrichTrapEntry(entry, false, nil) if entry.DeviceHostname != "" { t.Errorf("DeviceHostname = %q, want empty for unknown device", entry.DeviceHostname) @@ -211,22 +210,21 @@ func TestEnrichTrapEntryNoDeviceRegistryMatch(t *testing.T) { } } -func TestEnrichTrapEntryAmbiguousDeviceRegistryMatchDoesNotEnrich(t *testing.T) { - ddsnmp.DeviceRegistry.Register("job-a:10.9.9.1:162", ddsnmp.DeviceConnectionInfo{ +func TestEnrichTrapEntryAmbiguousDeviceStoreMatchDoesNotEnrich(t *testing.T) { + c, store := newTestTrapEnrichmentCollector(nil) + store.Register("job-a:10.9.9.1:162", ddsnmp.DeviceConnectionInfo{ Hostname: "10.9.9.1", SysName: "switch-a", Vendor: "vendor-a", }) - defer ddsnmp.DeviceRegistry.Unregister("job-a:10.9.9.1:162") - ddsnmp.DeviceRegistry.Register("job-b:10.9.9.1:162", ddsnmp.DeviceConnectionInfo{ + store.Register("job-b:10.9.9.1:162", ddsnmp.DeviceConnectionInfo{ Hostname: "10.9.9.1", SysName: "switch-b", Vendor: "vendor-b", }) - defer ddsnmp.DeviceRegistry.Unregister("job-b:10.9.9.1:162") entry := &TrapEntry{SourceIP: "10.9.9.1"} - enrichTrapEntry(entry, false, nil) + c.enrichTrapEntry(entry, false, nil) if entry.DeviceHostname != "" { t.Errorf("DeviceHostname = %q, want empty for ambiguous registry source", entry.DeviceHostname) @@ -243,8 +241,7 @@ func TestEnrichTrapEntryAmbiguousDeviceRegistryMatchDoesNotEnrich(t *testing.T) } func TestEnrichTrapEntryDoesNotUseTopologyOnVnodeConflict(t *testing.T) { - prev := trapTopologyEnrichmentForSource - trapTopologyEnrichmentForSource = func(_, ifIndex string) *snmptopology.TrapTopologyEnrichment { + topologyEnricher := testTrapTopologyEnricher(func(_, ifIndex string) *snmptopology.TrapTopologyEnrichment { return &snmptopology.TrapTopologyEnrichment{ DeviceStatus: "matched", DeviceMethod: "management_ip", @@ -258,15 +255,14 @@ func TestEnrichTrapEntryDoesNotUseTopologyOnVnodeConflict(t *testing.T) { NeighborStatus: "matched", Neighbors: []string{"dist-a"}, } - } - t.Cleanup(func() { trapTopologyEnrichmentForSource = prev }) + }) - ddsnmp.DeviceRegistry.Register("job-a:10.9.9.2:162", ddsnmp.DeviceConnectionInfo{ + c, store := newTestTrapEnrichmentCollector(topologyEnricher) + store.Register("job-a:10.9.9.2:162", ddsnmp.DeviceConnectionInfo{ Hostname: "10.9.9.2", SysName: "registry-switch", VnodeGUID: "registry-vnode-id", }) - defer ddsnmp.DeviceRegistry.Unregister("job-a:10.9.9.2:162") entry := &TrapEntry{ SourceIP: "10.9.9.2", @@ -274,7 +270,7 @@ func TestEnrichTrapEntryDoesNotUseTopologyOnVnodeConflict(t *testing.T) { {Name: "ifIndex", OID: ifIndexOIDPrefix + ".1", Type: "InterfaceIndex", Value: int64(1)}, }, } - enrichTrapEntry(entry, false, nil) + c.enrichTrapEntry(entry, false, nil) if entry.DeviceHostname != "registry-switch" { t.Errorf("DeviceHostname = %q, want registry-switch", entry.DeviceHostname) @@ -294,11 +290,10 @@ func TestEnrichTrapEntryDoesNotUseTopologyOnVnodeConflict(t *testing.T) { } func TestEnrichTrapEntryUsesTrapVarbindInterfaceWithoutTopology(t *testing.T) { - prev := trapTopologyEnrichmentForSource - trapTopologyEnrichmentForSource = func(_, _ string) *snmptopology.TrapTopologyEnrichment { + topologyEnricher := testTrapTopologyEnricher(func(_, _ string) *snmptopology.TrapTopologyEnrichment { return nil - } - t.Cleanup(func() { trapTopologyEnrichmentForSource = prev }) + }) + c, _ := newTestTrapEnrichmentCollector(topologyEnricher) entry := &TrapEntry{ SourceIP: "10.9.9.3", @@ -307,7 +302,7 @@ func TestEnrichTrapEntryUsesTrapVarbindInterfaceWithoutTopology(t *testing.T) { {Name: "ifName", OID: ifNameOIDPrefix + ".29", Type: "OctetString", Value: "uplink-29"}, }, } - enrichTrapEntry(entry, false, nil) + c.enrichTrapEntry(entry, false, nil) if entry.TopologyInterface != "uplink-29" { t.Errorf("TopologyInterface = %q, want uplink-29", entry.TopologyInterface) @@ -327,8 +322,9 @@ func TestEnrichTrapEntryUsesTrapVarbindInterfaceWithoutTopology(t *testing.T) { } func TestEnrichTrapEntrySourceUDPPeerFallback(t *testing.T) { + c, _ := newTestTrapEnrichmentCollector(nil) entry := &TrapEntry{SourceUDPPeer: "192.168.1.1"} - enrichTrapEntry(entry, false, nil) + c.enrichTrapEntry(entry, false, nil) if entry.DeviceHostname != "" { t.Errorf("DeviceHostname = %q, want empty (no device match)", entry.DeviceHostname) @@ -336,12 +332,14 @@ func TestEnrichTrapEntrySourceUDPPeerFallback(t *testing.T) { } func TestEnrichTrapEntryNilEntry(t *testing.T) { - enrichTrapEntry(nil, false, nil) + c, _ := newTestTrapEnrichmentCollector(nil) + c.enrichTrapEntry(nil, false, nil) } func TestEnrichTrapEntryNoSource(t *testing.T) { + c, _ := newTestTrapEnrichmentCollector(nil) entry := &TrapEntry{} - enrichTrapEntry(entry, false, nil) + c.enrichTrapEntry(entry, false, nil) if entry.DeviceHostname != "" { t.Errorf("DeviceHostname = %q, want empty", entry.DeviceHostname) @@ -349,11 +347,11 @@ func TestEnrichTrapEntryNoSource(t *testing.T) { } func TestEnrichTrapEntryReverseDNSDefaultOff(t *testing.T) { + c, store := newTestTrapEnrichmentCollector(nil) regKey := "key:10.5.5.1:162" - ddsnmp.DeviceRegistry.Register(regKey, ddsnmp.DeviceConnectionInfo{ + store.Register(regKey, ddsnmp.DeviceConnectionInfo{ Hostname: "10.5.5.1", }) - defer ddsnmp.DeviceRegistry.Unregister(regKey) dns := newReverseDNSResolver() dns.cache["10.5.5.1"] = reverseDNSCacheEntry{ @@ -362,7 +360,7 @@ func TestEnrichTrapEntryReverseDNSDefaultOff(t *testing.T) { } entry := &TrapEntry{SourceIP: "10.5.5.1"} - enrichTrapEntry(entry, false, dns) + c.enrichTrapEntry(entry, false, dns) if entry.DeviceHostname != "" { t.Errorf("DeviceHostname = %q, want empty (reverse DNS disabled, no vnode/sysName)", entry.DeviceHostname) @@ -370,6 +368,7 @@ func TestEnrichTrapEntryReverseDNSDefaultOff(t *testing.T) { } func TestEnrichTrapEntryReverseDNSEnabledNoSNMPState(t *testing.T) { + c, _ := newTestTrapEnrichmentCollector(nil) dns := newReverseDNSResolver() dns.cache["10.6.6.1"] = reverseDNSCacheEntry{ name: "peer.mydc.example.com", @@ -377,7 +376,7 @@ func TestEnrichTrapEntryReverseDNSEnabledNoSNMPState(t *testing.T) { } entry := &TrapEntry{SourceIP: "10.6.6.1"} - enrichTrapEntry(entry, true, dns) + c.enrichTrapEntry(entry, true, dns) if entry.DeviceHostname != "" { t.Errorf("DeviceHostname = %q, want empty because reverse DNS is not authoritative identity", entry.DeviceHostname) @@ -394,11 +393,11 @@ func TestEnrichTrapEntryReverseDNSEnabledNoSNMPState(t *testing.T) { } func TestEnrichTrapEntryReverseDNSDisabledNoCacheUse(t *testing.T) { + c, store := newTestTrapEnrichmentCollector(nil) regKey := "key:10.7.7.1:162" - ddsnmp.DeviceRegistry.Register(regKey, ddsnmp.DeviceConnectionInfo{ + store.Register(regKey, ddsnmp.DeviceConnectionInfo{ Hostname: "10.7.7.1", }) - defer ddsnmp.DeviceRegistry.Unregister(regKey) dns := newReverseDNSResolver() dns.cache["10.7.7.1"] = reverseDNSCacheEntry{ @@ -407,7 +406,7 @@ func TestEnrichTrapEntryReverseDNSDisabledNoCacheUse(t *testing.T) { } entry := &TrapEntry{SourceIP: "10.7.7.1"} - enrichTrapEntry(entry, false, dns) + c.enrichTrapEntry(entry, false, dns) if entry.DeviceHostname != "" { t.Errorf("DeviceHostname = %q, want empty (reverse DNS disabled, no SNMP state)", entry.DeviceHostname) @@ -415,12 +414,12 @@ func TestEnrichTrapEntryReverseDNSDisabledNoCacheUse(t *testing.T) { } func TestEnrichTrapEntryReverseDNSDoesNotReplaceKnownHostname(t *testing.T) { + c, store := newTestTrapEnrichmentCollector(nil) regKey := "key:10.7.7.2:162" - ddsnmp.DeviceRegistry.Register(regKey, ddsnmp.DeviceConnectionInfo{ + store.Register(regKey, ddsnmp.DeviceConnectionInfo{ Hostname: "10.7.7.2", SysName: "known-switch", }) - defer ddsnmp.DeviceRegistry.Unregister(regKey) dns := newReverseDNSResolver() dns.cache["10.7.7.2"] = reverseDNSCacheEntry{ @@ -428,7 +427,7 @@ func TestEnrichTrapEntryReverseDNSDoesNotReplaceKnownHostname(t *testing.T) { expiresAt: farFuture(), } entry := &TrapEntry{SourceIP: "10.7.7.2"} - enrichTrapEntry(entry, true, dns) + c.enrichTrapEntry(entry, true, dns) if entry.DeviceHostname != "known-switch" { t.Errorf("DeviceHostname = %q, want known-switch", entry.DeviceHostname) @@ -439,6 +438,7 @@ func TestEnrichTrapEntryReverseDNSDoesNotReplaceKnownHostname(t *testing.T) { } func TestEnrichTrapEntryReverseDNSEnabledSchedulesAsyncLookup(t *testing.T) { + c, _ := newTestTrapEnrichmentCollector(nil) dns := newReverseDNSResolver() defer dns.Close() @@ -455,7 +455,7 @@ func TestEnrichTrapEntryReverseDNSEnabledSchedulesAsyncLookup(t *testing.T) { } entry := &TrapEntry{SourceIP: "203.0.113.10"} - enrichTrapEntry(entry, true, dns) + c.enrichTrapEntry(entry, true, dns) select { case <-started: @@ -515,17 +515,17 @@ func TestEnrichTrapEntryVendorAndVnodeEnrichment(t *testing.T) { for tcName, tc := range tests { t.Run(tcName, func(t *testing.T) { + c, store := newTestTrapEnrichmentCollector(nil) regKey := "key:" + tc.hostname + ":162" - ddsnmp.DeviceRegistry.Register(regKey, ddsnmp.DeviceConnectionInfo{ + store.Register(regKey, ddsnmp.DeviceConnectionInfo{ Hostname: tc.hostname, SysName: tc.sysName, Vendor: tc.vendor, VnodeGUID: tc.vnodeGUID, }) - defer ddsnmp.DeviceRegistry.Unregister(regKey) entry := &TrapEntry{SourceIP: tc.hostname} - enrichTrapEntry(entry, false, nil) + c.enrichTrapEntry(entry, false, nil) if entry.DeviceVendor != tc.wantVendor { t.Errorf("DeviceVendor = %q, want %q", entry.DeviceVendor, tc.wantVendor) diff --git a/src/go/plugin/go.d/collector/snmp_traps/func_logs_test.go b/src/go/plugin/go.d/collector/snmp_traps/func_logs_test.go index a219b1564dd656..4f242ed945967c 100644 --- a/src/go/plugin/go.d/collector/snmp_traps/func_logs_test.go +++ b/src/go/plugin/go.d/collector/snmp_traps/func_logs_test.go @@ -14,6 +14,8 @@ import ( "github.com/netdata/netdata/go/plugins/plugin/agent/jobmgr/funcctl" "github.com/netdata/netdata/go/plugins/plugin/framework/collectorapi" "github.com/netdata/netdata/go/plugins/plugin/framework/functions" + "github.com/netdata/netdata/go/plugins/plugin/go.d/collector/snmp/ddsnmp" + snmptopology "github.com/netdata/netdata/go/plugins/plugin/go.d/collector/snmp_topology" "github.com/netdata/netdata/go/plugins/plugin/go.d/collector/snmp_traps/snmptrapsfunc" "github.com/stretchr/testify/assert" "github.com/stretchr/testify/require" @@ -147,8 +149,7 @@ func TestSNMPTrapsLogsDispatchDoesNotRequireRunningJob(t *testing.T) { t.Cleanup(func() { activeDirectJournalJobs.Store(startJournalJobs) }) t.Setenv(netdataLogDirEnv, filepath.Join(t.TempDir(), "logs")) - creator, ok := collectorapi.DefaultRegistry.Lookup("snmp_traps") - require.True(t, ok) + creator := newCreator(ddsnmp.NewDeviceStore(), snmptopology.NewTrapEnrichmentHandle()) reg := newSNMPTrapsTestFunctionRegistry() var gotCode int diff --git a/src/go/plugin/go.d/collector/snmp_traps/init_test.go b/src/go/plugin/go.d/collector/snmp_traps/init_test.go index 731b05240cd4a8..426e5574e0db16 100644 --- a/src/go/plugin/go.d/collector/snmp_traps/init_test.go +++ b/src/go/plugin/go.d/collector/snmp_traps/init_test.go @@ -13,18 +13,19 @@ import ( "github.com/netdata/netdata/go/plugins/plugin/framework/chartengine" "github.com/netdata/netdata/go/plugins/plugin/framework/charttpl" - "github.com/netdata/netdata/go/plugins/plugin/framework/collectorapi" + "github.com/netdata/netdata/go/plugins/plugin/go.d/collector/snmp/ddsnmp" + snmptopology "github.com/netdata/netdata/go/plugins/plugin/go.d/collector/snmp_topology" "github.com/netdata/netdata/go/plugins/plugin/go.d/pkg/collecttest" "github.com/stretchr/testify/assert" "github.com/stretchr/testify/require" ) func TestCollectorChartTemplateYAML(t *testing.T) { - collecttest.AssertChartTemplateSchema(t, New().ChartTemplateYAML()) + collecttest.AssertChartTemplateSchema(t, newTestSNMPTrapsCollector().ChartTemplateYAML()) } func TestCollectorChartTemplateYAMLChartsDeclareAlgorithms(t *testing.T) { - charts := chartTemplatesByIDFromYAML(t, New().ChartTemplateYAML()) + charts := chartTemplatesByIDFromYAML(t, newTestSNMPTrapsCollector().ChartTemplateYAML()) assertAllChartTemplatesDeclareAlgorithm(t, charts) for _, id := range []string{ @@ -68,7 +69,7 @@ func TestCollectorChartTemplateYAMLIncludesProfileMetricCharts(t *testing.T) { require.NotNil(t, rt) require.NotEmpty(t, tmpl) - c := New() + c := newTestSNMPTrapsCollector() c.profileMetrics = rt c.dynamicChartYAML = tmpl @@ -90,12 +91,29 @@ func TestCollectorChartTemplateYAMLIncludesProfileMetricCharts(t *testing.T) { assert.Contains(t, contexts, "snmp.trap.cisco.config.changes") } -func TestCollectorRegistrationAvailableByDefault(t *testing.T) { - creator, ok := collectorapi.DefaultRegistry.Lookup("snmp_traps") - require.True(t, ok) +func TestCollectorCreatorDefaults(t *testing.T) { + creator := newCreator(ddsnmp.NewDeviceStore(), snmptopology.NewTrapEnrichmentHandle()) assert.False(t, creator.Defaults.Disabled) } +func TestCollectorCreatorRequiresSharedDependencies(t *testing.T) { + require.PanicsWithValue(t, "snmp_traps Register requires a non-nil device store", func() { + _ = newCreator(nil, snmptopology.NewTrapEnrichmentHandle()) + }) + require.PanicsWithValue(t, "snmp_traps Register requires a non-nil trap enrichment handle", func() { + _ = newCreator(ddsnmp.NewDeviceStore(), nil) + }) +} + +func TestCollectorNewRequiresSharedDependencies(t *testing.T) { + require.PanicsWithValue(t, "snmp_traps New requires a non-nil device store", func() { + _ = New(nil, snmptopology.NewTrapEnrichmentHandle()) + }) + require.PanicsWithValue(t, "snmp_traps New requires a non-nil trap enrichment handle", func() { + _ = New(ddsnmp.NewDeviceStore(), nil) + }) +} + func chartTemplatesByIDFromYAML(t *testing.T, raw string) map[string]charttpl.Chart { t.Helper() @@ -215,7 +233,7 @@ func TestConfigSchemaDynCfgRetentionDefaultDisablesTimeRotation(t *testing.T) { } func TestCollectorDefaultListenReceiveBuffer(t *testing.T) { - assert.Equal(t, defaultListenerReceiveBuffer, New().Listen.ReceiveBuffer) + assert.Equal(t, defaultListenerReceiveBuffer, newTestSNMPTrapsCollector().Listen.ReceiveBuffer) } func TestConfigSchemaDynCfgTabsRenderAllTopLevelFieldsOnce(t *testing.T) { @@ -469,7 +487,7 @@ func TestCollectorInit_BindsEndpointsAndCheckIsNoop(t *testing.T) { withTestCacheDir(t) port := freeUDPPort(t) - c := New() + c := newTestSNMPTrapsCollector() c.SetJobName("local") c.Listen.Endpoints = []EndpointConfig{{Protocol: "udp", Address: "127.0.0.1", Port: port}} @@ -493,7 +511,7 @@ func TestCollectorInit_IdempotentDoubleInit(t *testing.T) { withTestCacheDir(t) port := freeUDPPort(t) - c := New() + c := newTestSNMPTrapsCollector() c.SetJobName("local") c.Listen.Endpoints = []EndpointConfig{{Protocol: "udp", Address: "127.0.0.1", Port: port}} @@ -510,7 +528,7 @@ func TestCollectorInit_IdempotentDoubleInit(t *testing.T) { func TestCollectorInit_InvalidJobNameIsCodedError(t *testing.T) { withTestCacheDir(t) - c := New() + c := newTestSNMPTrapsCollector() c.SetJobName("../bad") c.Listen.Endpoints = []EndpointConfig{{Protocol: "udp", Address: "127.0.0.1", Port: 162}} @@ -528,7 +546,7 @@ func TestCollectorInit_InvalidJobNameIsCodedError(t *testing.T) { func TestCollectorInit_InvalidEndpointsIsCodedError(t *testing.T) { withTestCacheDir(t) - c := New() + c := newTestSNMPTrapsCollector() c.SetJobName("local") c.Listen.Endpoints = []EndpointConfig{{Protocol: "tcp", Address: "127.0.0.1", Port: 162}} @@ -543,7 +561,7 @@ func TestCollectorInit_InvalidEndpointsIsCodedError(t *testing.T) { func TestCollectorInit_InvalidReceiveBufferIsCodedError(t *testing.T) { withTestCacheDir(t) - c := New() + c := newTestSNMPTrapsCollector() c.SetJobName("local") c.Listen.Endpoints = []EndpointConfig{{Protocol: "udp", Address: "127.0.0.1", Port: freeUDPPort(t)}} c.Listen.ReceiveBuffer = -1 @@ -560,7 +578,7 @@ func TestCollectorInit_InvalidReceiveBufferIsCodedError(t *testing.T) { func TestCollectorInit_TooLargeReceiveBufferIsCodedError(t *testing.T) { withTestCacheDir(t) - c := New() + c := newTestSNMPTrapsCollector() c.SetJobName("local") c.Listen.Endpoints = []EndpointConfig{{Protocol: "udp", Address: "127.0.0.1", Port: freeUDPPort(t)}} c.Listen.ReceiveBuffer = maxListenerReceiveBuffer + 1 @@ -576,7 +594,7 @@ func TestCollectorInit_TooLargeReceiveBufferIsCodedError(t *testing.T) { func TestCollectorInit_NoOutputBackendIsCodedError(t *testing.T) { disabled := false - c := New() + c := newTestSNMPTrapsCollector() c.SetJobName("local") c.Listen.Endpoints = []EndpointConfig{{Protocol: "udp", Address: "127.0.0.1", Port: freeUDPPort(t)}} c.Journal.Enabled = &disabled @@ -595,7 +613,7 @@ func TestCollectorInit_MissingNetdataLogRootIsRetryableCodedError(t *testing.T) root := filepath.Join(t.TempDir(), "missing") withNetdataLogDir(t, root) - c := New() + c := newTestSNMPTrapsCollector() c.SetJobName("local") c.Listen.Endpoints = []EndpointConfig{{Protocol: "udp", Address: "127.0.0.1", Port: freeUDPPort(t)}} @@ -623,7 +641,7 @@ func TestCollectorInit_OTELOnlySkipsJournalCreation(t *testing.T) { srv := startOTLPFixture(t, nil) const jobName = "otel-only" - c := New() + c := newTestSNMPTrapsCollector() c.SetJobName(jobName) c.Listen.Endpoints = []EndpointConfig{{Protocol: "udp", Address: "127.0.0.1", Port: freeUDPPort(t)}} c.Journal.Enabled = &disabled @@ -657,7 +675,7 @@ func TestCollectorInit_OTLPPreflightFailureIsRetryableCodedError(t *testing.T) { endpoint := "http://" + ln.Addr().String() require.NoError(t, ln.Close()) - c := New() + c := newTestSNMPTrapsCollector() c.SetJobName("otlp-preflight") c.Listen.Endpoints = []EndpointConfig{{Protocol: "udp", Address: "127.0.0.1", Port: freeUDPPort(t)}} c.Journal.Enabled = &disabled @@ -685,7 +703,7 @@ func TestCollectorInit_BindsMultipleEndpoints(t *testing.T) { firstPort := freeUDPPort(t) secondPort := freeUDPPort(t) - c := New() + c := newTestSNMPTrapsCollector() c.SetJobName("local") c.Listen.Endpoints = []EndpointConfig{ {Protocol: "udp", Address: "127.0.0.1", Port: firstPort}, @@ -717,7 +735,7 @@ func TestCollectorInit_BindFailureIsRetryableCodedError(t *testing.T) { port := conn.LocalAddr().(*net.UDPAddr).Port - c := New() + c := newTestSNMPTrapsCollector() c.SetJobName("local") c.Listen.Endpoints = []EndpointConfig{{Protocol: "udp", Address: "127.0.0.1", Port: port}} @@ -742,7 +760,7 @@ func TestCollectorInit_ReceiveBufferFailureIsRetryableCodedError(t *testing.T) { return errors.New("set buffer failed") } - c := New() + c := newTestSNMPTrapsCollector() c.SetJobName("local") c.Listen.Endpoints = []EndpointConfig{{Protocol: "udp", Address: "127.0.0.1", Port: freeUDPPort(t)}} @@ -761,7 +779,7 @@ func TestCollectorInit_ReceiveBufferFailureIsRetryableCodedError(t *testing.T) { func TestCollectorInit_InvalidVersionIsCodedError(t *testing.T) { withTestCacheDir(t) - c := New() + c := newTestSNMPTrapsCollector() c.SetJobName("local") c.Listen.Endpoints = []EndpointConfig{{Protocol: "udp", Address: "127.0.0.1", Port: 162}} c.Versions = []string{"v5"} @@ -779,7 +797,7 @@ func TestCollectorInit_ProfileLoadFailureIsCodedError(t *testing.T) { resetProfileCacheForTest() withTestCacheDir(t) - c := New() + c := newTestSNMPTrapsCollector() c.SetJobName("local") c.Listen.Endpoints = []EndpointConfig{{Protocol: "udp", Address: "127.0.0.1", Port: freeUDPPort(t)}} c.Versions = []string{" V1 ", "V2C"} @@ -802,7 +820,7 @@ func TestCollectorInit_PartialBindFailureClosesPriorSockets(t *testing.T) { defer secondConn.Close() secondPort := secondConn.LocalAddr().(*net.UDPAddr).Port - c := New() + c := newTestSNMPTrapsCollector() c.SetJobName("local") c.Listen.Endpoints = []EndpointConfig{ {Protocol: "udp", Address: "127.0.0.1", Port: firstPort}, @@ -836,7 +854,7 @@ func TestCollectorInit_EngineStateStatErrorIsRetryableCodedError(t *testing.T) { const jobName = "engine-state-stat-error" require.NoError(t, os.WriteFile(engineBootsDir(jobName), []byte("not a directory"), 0644)) - c := New() + c := newTestSNMPTrapsCollector() c.SetJobName(jobName) c.Listen.Endpoints = []EndpointConfig{{Protocol: "udp", Address: "127.0.0.1", Port: freeUDPPort(t)}} c.Versions = []string{"v3"} @@ -870,7 +888,7 @@ func TestCollectorInit_CleansCreatedV3StateOnEngineBootsFailure(t *testing.T) { const jobName = "cleanup-v3-state" require.NoError(t, os.MkdirAll(engineBootsPath(jobName), 0750)) - c := New() + c := newTestSNMPTrapsCollector() c.SetJobName(jobName) c.Listen.Endpoints = []EndpointConfig{{Protocol: "udp", Address: "127.0.0.1", Port: freeUDPPort(t)}} c.Versions = []string{"v3"} @@ -896,7 +914,7 @@ func TestCollectorCleanupIsIdempotent(t *testing.T) { withTestCacheDir(t) port := freeUDPPort(t) - c := New() + c := newTestSNMPTrapsCollector() c.SetJobName("local") c.Listen.Endpoints = []EndpointConfig{{Protocol: "udp", Address: "127.0.0.1", Port: port}} @@ -909,7 +927,7 @@ func TestCollectorCleanupIsIdempotent(t *testing.T) { } func TestCollectorCollectRequiresStartedListener(t *testing.T) { - c := New() + c := newTestSNMPTrapsCollector() err := c.Collect(context.Background()) require.Error(t, err) assert.Contains(t, err.Error(), "listener not started") diff --git a/src/go/plugin/go.d/collector/snmp_traps/listener_test.go b/src/go/plugin/go.d/collector/snmp_traps/listener_test.go index 8fad389483bd7d..fb93c506bd4667 100644 --- a/src/go/plugin/go.d/collector/snmp_traps/listener_test.go +++ b/src/go/plugin/go.d/collector/snmp_traps/listener_test.go @@ -81,7 +81,7 @@ func TestListenerReadLoopDoesNotReportReadErrorDuringClose(t *testing.T) { func TestCollectorLogListenerReadErrorIsRateLimited(t *testing.T) { var buf bytes.Buffer - c := New() + c := newTestSNMPTrapsCollector() c.Logger = logger.NewWithWriter(&buf) ep := EndpointConfig{ Protocol: "udp4", diff --git a/src/go/plugin/go.d/collector/snmp_traps/pipeline_test.go b/src/go/plugin/go.d/collector/snmp_traps/pipeline_test.go index 82fea57ff8332f..7ba7d01a9d523c 100644 --- a/src/go/plugin/go.d/collector/snmp_traps/pipeline_test.go +++ b/src/go/plugin/go.d/collector/snmp_traps/pipeline_test.go @@ -119,18 +119,19 @@ func TestCollectorHandlePacketRecoversFromPanic(t *testing.T) { func TestCollectorHandlePacketRendersTemplatesAfterEnrichment(t *testing.T) { packet := readColdStartUDPPacket(t) + deviceStore := ddsnmp.NewDeviceStore() regKey := "test:198.51.100.10:162" - ddsnmp.DeviceRegistry.Register(regKey, ddsnmp.DeviceConnectionInfo{ + deviceStore.Register(regKey, ddsnmp.DeviceConnectionInfo{ Hostname: "198.51.100.10", SysName: "core-sw-01", Vendor: "cisco", }) - defer ddsnmp.DeviceRegistry.Unregister(regKey) trap := testColdStartTrap("security", "warning", "security coldStart on {_HOSTNAME} from {TRAP_DEVICE_VENDOR}") setSingleTestTrap(t, trap) writer := &mockTrapWriter{} c := newDefaultTestV2Collector(writer) + c.deviceLookup = deviceStore c.handlePacket(packet.payload, packet.peer, nil, nil) @@ -162,8 +163,7 @@ func TestCollectorHandlePacketDoesNotUseListenerVnodeAsSourceNode(t *testing.T) func TestCollectorHandlePacketRendersTopologyEnrichmentBeforeReverseDNS(t *testing.T) { packet := readColdStartUDPPacket(t) - prev := trapTopologyEnrichmentForSource - trapTopologyEnrichmentForSource = func(ip, ifIndex string) *snmptopology.TrapTopologyEnrichment { + topologyEnricher := testTrapTopologyEnricher(func(ip, ifIndex string) *snmptopology.TrapTopologyEnrichment { if ip != "198.51.100.10" { t.Fatalf("topology enrichment looked up IP %q, want 198.51.100.10", ip) } @@ -180,8 +180,7 @@ func TestCollectorHandlePacketRendersTopologyEnrichmentBeforeReverseDNS(t *testi InterfaceStatus: "skipped", NeighborStatus: "skipped", } - } - t.Cleanup(func() { trapTopologyEnrichmentForSource = prev }) + }) dns := newReverseDNSResolver() dns.cache["198.51.100.10"] = reverseDNSCacheEntry{ @@ -198,6 +197,7 @@ func TestCollectorHandlePacketRendersTopologyEnrichmentBeforeReverseDNS(t *testi setSingleTestTrap(t, trap) writer := &mockTrapWriter{} c := newDefaultTestV2Collector(writer) + c.topologyEnricher = topologyEnricher c.reverseDNSEnabled = true c.reverseDNS = dns diff --git a/src/go/plugin/go.d/collector/snmp_traps/profile_test.go b/src/go/plugin/go.d/collector/snmp_traps/profile_test.go index 2978eb74fc234c..ac5f6e2bc80d10 100644 --- a/src/go/plugin/go.d/collector/snmp_traps/profile_test.go +++ b/src/go/plugin/go.d/collector/snmp_traps/profile_test.go @@ -90,7 +90,7 @@ func TestCollectorInitAcquiresProfileCache(t *testing.T) { port := freeUDPPort(t) - c := New() + c := newTestSNMPTrapsCollector() c.SetJobName("local") c.Listen.Endpoints = []EndpointConfig{{Protocol: "udp", Address: "127.0.0.1", Port: port}} @@ -109,11 +109,11 @@ func TestMultipleCollectorsShareSameCache(t *testing.T) { port1 := freeUDPPort(t) port2 := freeUDPPort(t) - c1 := New() + c1 := newTestSNMPTrapsCollector() c1.SetJobName("job1") c1.Listen.Endpoints = []EndpointConfig{{Protocol: "udp", Address: "127.0.0.1", Port: port1}} - c2 := New() + c2 := newTestSNMPTrapsCollector() c2.SetJobName("job2") c2.Listen.Endpoints = []EndpointConfig{{Protocol: "udp", Address: "127.0.0.1", Port: port2}} @@ -142,7 +142,7 @@ func TestInitBindFailureReleasesProfileRef(t *testing.T) { require.NoError(t, err) defer conn.Close() - c := New() + c := newTestSNMPTrapsCollector() c.SetJobName("local") c.Listen.Endpoints = []EndpointConfig{{Protocol: "udp", Address: "127.0.0.1", Port: conn.LocalAddr().(*net.UDPAddr).Port}} @@ -444,7 +444,7 @@ traps: require.Equal(t, "diagnostic", td.Category) require.Equal(t, "warning", td.Severity) - c := New() + c := newTestSNMPTrapsCollector() c.overrides = buildOverrideMap([]OverrideConfig{ { OID: oid, diff --git a/src/go/plugin/go.d/collector/snmp_traps/test_helpers_test.go b/src/go/plugin/go.d/collector/snmp_traps/test_helpers_test.go index a41b0ba0a6d556..6bc7baa67e1d8e 100644 --- a/src/go/plugin/go.d/collector/snmp_traps/test_helpers_test.go +++ b/src/go/plugin/go.d/collector/snmp_traps/test_helpers_test.go @@ -10,6 +10,8 @@ import ( "github.com/gosnmp/gosnmp" "github.com/netdata/netdata/go/plugins/pkg/metrix" + "github.com/netdata/netdata/go/plugins/plugin/go.d/collector/snmp/ddsnmp" + snmptopology "github.com/netdata/netdata/go/plugins/plugin/go.d/collector/snmp_topology" ) func readSinglePcapUDPPacket(t *testing.T, fixture string) pcapUDPPacket { @@ -55,6 +57,24 @@ func newDefaultTestV2Collector(writer TrapWriter) *Collector { return newTestV2Collector("test", writer, nil, []string{"public"}) } +func newTestSNMPTrapsCollector() *Collector { + return New(ddsnmp.NewDeviceStore(), snmptopology.NewTrapEnrichmentHandle()) +} + +type testTrapTopologyEnricher func(ip, trapIfIndex string) *snmptopology.TrapTopologyEnrichment + +func (f testTrapTopologyEnricher) EnrichmentForSource(ip, trapIfIndex string) *snmptopology.TrapTopologyEnrichment { + return f(ip, trapIfIndex) +} + +func newTestTrapEnrichmentCollector(topologyEnricher trapTopologyEnricher) (*Collector, *ddsnmp.DeviceStore) { + store := ddsnmp.NewDeviceStore() + return &Collector{ + deviceLookup: store, + topologyEnricher: topologyEnricher, + }, store +} + func withCleanJobMetrics(t *testing.T, jobName string) *perJobMetrics { t.Helper() removeJobMetrics(jobName) From 412b9556cebacd78bd7613bdb008c9c5debec174 Mon Sep 17 00:00:00 2001 From: Stelios Fragkakis <52996999+stelfrag@users.noreply.github.com> Date: Thu, 18 Jun 2026 19:48:13 +0300 Subject: [PATCH 8/9] Sonar cleanup (#22760) * claim/ui: remove unused static hToken/hRoom (sonar c:S1481) Co-Authored-By: Claude Opus 4.8 (1M context) * c_rhash/tests: remove unused 'val' in test_uint64_ptr_incremental (sonar c:S1481) Co-Authored-By: Claude Opus 4.8 (1M context) * freebsd.plugin: remove unused st_intr/rd_intr statics (sonar c:S1481) Co-Authored-By: Claude Opus 4.8 (1M context) * systemd-cat-native: drop dead store to remaining_dst_len before break (sonar c:S1854) Co-Authored-By: Claude Opus 4.8 (1M context) * websocket-send: drop dead "" init of disconnect_msg (sonar c:S1854) All paths to abnormal_disconnect set disconnect_msg first; the success path returns before the label, so the initializer was never read. Co-Authored-By: Claude Opus 4.8 (1M context) * api_v2_claim: drop dead store to can_be_claimed (sonar c:S1854) Only read at the guarding 'if(can_be_claimed && key)'; never read after the assignment, so setting it to false had no effect. Co-Authored-By: Claude Opus 4.8 (1M context) * inlined/murmur64: use ULL suffix for 64-bit constants (sonar c:S6996) The constants are 64-bit and used in uint64_t math; UL is only 32-bit on LLP64 (Windows). ULL is the portable suffix matching intent. Co-Authored-By: Claude Opus 4.8 (1M context) * netipc/uds: uppercase ull -> ULL literal suffix (sonar c:S818) Co-Authored-By: Claude Opus 4.8 (1M context) * rrdhost-system-info: document empty coverity_remove_taint marker (sonar c:S1186) Co-Authored-By: Claude Opus 4.8 (1M context) * nd_log: document empty debug_dummy no-op (sonar c:S1186) Co-Authored-By: Claude Opus 4.8 (1M context) * cgroups/discovery: use while loop instead of empty-clause for (sonar c:S1264) Co-Authored-By: Claude Opus 4.8 (1M context) * stream-connector: merge adjacent string literals in thread tag (sonar c:S3728) Co-Authored-By: Claude Opus 4.8 (1M context) --------- Co-authored-by: Claude Opus 4.8 (1M context) --- src/claim/ui.c | 2 -- src/collectors/cgroups.plugin/cgroup-discovery.c | 2 +- src/collectors/freebsd.plugin/freebsd_sysctl.c | 3 --- src/database/rrdhost-system-info.c | 5 ++++- src/libnetdata/c_rhash/tests.c | 1 - src/libnetdata/inlined.h | 4 ++-- src/libnetdata/log/nd_log.h | 2 +- src/libnetdata/log/systemd-cat-native.c | 1 - .../netipc/src/transport/posix/netipc_uds_receive.c | 2 +- src/streaming/stream-connector.c | 2 +- src/web/api/v2/api_v2_claim.c | 1 - src/web/websocket/websocket-send.c | 2 +- 12 files changed, 11 insertions(+), 16 deletions(-) diff --git a/src/claim/ui.c b/src/claim/ui.c index 92805d777be382..986577349a7cea 100644 --- a/src/claim/ui.c +++ b/src/claim/ui.c @@ -10,8 +10,6 @@ static LPCTSTR szWindowClass = _T("DesktopApp"); static HINSTANCE hInst; -static HWND hToken; -static HWND hRoom; LRESULT CALLBACK WndProc(HWND hNetdatawnd, UINT message, WPARAM wParam, LPARAM lParam) { diff --git a/src/collectors/cgroups.plugin/cgroup-discovery.c b/src/collectors/cgroups.plugin/cgroup-discovery.c index a895b8ae53bd99..483dbebd659d2f 100644 --- a/src/collectors/cgroups.plugin/cgroup-discovery.c +++ b/src/collectors/cgroups.plugin/cgroup-discovery.c @@ -794,7 +794,7 @@ static inline void discovery_update_filenames_all_cgroups() { static inline void discovery_cleanup_all_cgroups() { struct cgroup *cg = discovered_cgroup_root, *last = NULL; - for(; cg ;) { + while(cg) { if(!cg->available) { // enable the first duplicate cgroup { diff --git a/src/collectors/freebsd.plugin/freebsd_sysctl.c b/src/collectors/freebsd.plugin/freebsd_sysctl.c index 7dea4a71b61b5f..475ba1b098d0d7 100644 --- a/src/collectors/freebsd.plugin/freebsd_sysctl.c +++ b/src/collectors/freebsd.plugin/freebsd_sysctl.c @@ -574,9 +574,6 @@ int do_hw_intcnt(int update_every, usec_t dt) { for (i = 0; i < nintr; i++) totalintr += intrcnt[i]; - static RRDSET *st_intr = NULL; - static RRDDIM *rd_intr = NULL; - common_interrupts(totalintr, update_every, "hw.intrcnt"); size_t size; diff --git a/src/database/rrdhost-system-info.c b/src/database/rrdhost-system-info.c index c59a25b98c661b..6705f24f99bc2c 100644 --- a/src/database/rrdhost-system-info.c +++ b/src/database/rrdhost-system-info.c @@ -6,7 +6,10 @@ #include "daemon/win_system-info.h" // coverity[ +tainted_string_sanitize_content : arg-0 ] -static inline void coverity_remove_taint(char *s __maybe_unused) { } +static inline void coverity_remove_taint(char *s __maybe_unused) { + // intentionally empty: only a marker for the Coverity taint sanitizer + // (see the annotation above); it has no runtime effect. +} void rrdhost_system_info_swap(struct rrdhost_system_info *a, struct rrdhost_system_info *b) { if(a && b) diff --git a/src/libnetdata/c_rhash/tests.c b/src/libnetdata/c_rhash/tests.c index 3caa7d003662d6..061f4e5ee51a75 100644 --- a/src/libnetdata/c_rhash/tests.c +++ b/src/libnetdata/c_rhash/tests.c @@ -105,7 +105,6 @@ int test_uint64_ptr() { #define UINT64_PTR_INC_ITERATION_COUNT 5000 int test_uint64_ptr_incremental() { c_rhash hash = c_rhash_new(100); - void *val; TEST_START(); diff --git a/src/libnetdata/inlined.h b/src/libnetdata/inlined.h index ee5d90051ab2d4..aa8ee11b9530d7 100644 --- a/src/libnetdata/inlined.h +++ b/src/libnetdata/inlined.h @@ -108,9 +108,9 @@ static inline uint32_t murmur32(uint32_t k) { static uint64_t murmur64(uint64_t k) __attribute__((const)); static inline uint64_t murmur64(uint64_t k) { k ^= k >> 33; - k *= 0xff51afd7ed558ccdUL; + k *= 0xff51afd7ed558ccdULL; k ^= k >> 33; - k *= 0xc4ceb9fe1a85ec53UL; + k *= 0xc4ceb9fe1a85ec53ULL; k ^= k >> 33; return k; diff --git a/src/libnetdata/log/nd_log.h b/src/libnetdata/log/nd_log.h index 211540f3c80471..b2fe65b8540e51 100644 --- a/src/libnetdata/log/nd_log.h +++ b/src/libnetdata/log/nd_log.h @@ -117,7 +117,7 @@ extern int aclklog_enabled; #define LOG_DATE_LENGTH 26 void log_date(char *buffer, size_t len, time_t now); -static inline void debug_dummy(void) {} +static inline void debug_dummy(void) { /* no-op: target for debug/timing macros when those are compiled out */ } void nd_log_limits_reset(void); void nd_log_limits_unlimited(void); diff --git a/src/libnetdata/log/systemd-cat-native.c b/src/libnetdata/log/systemd-cat-native.c index 7d2a9b05ef9c19..381cf6b6b7a291 100644 --- a/src/libnetdata/log/systemd-cat-native.c +++ b/src/libnetdata/log/systemd-cat-native.c @@ -84,7 +84,6 @@ static inline size_t copy_replacing_newlines(char *dst, size_t dst_len, const ch memcpy(current_dst, current_src, copy_len); current_dst += copy_len; - remaining_dst_len -= copy_len; bytes_copied += copy_len; break; } diff --git a/src/libnetdata/netipc/src/transport/posix/netipc_uds_receive.c b/src/libnetdata/netipc/src/transport/posix/netipc_uds_receive.c index 50da34e5879888..b528bfc57b7a3a 100644 --- a/src/libnetdata/netipc/src/transport/posix/netipc_uds_receive.c +++ b/src/libnetdata/netipc/src/transport/posix/netipc_uds_receive.c @@ -17,7 +17,7 @@ static uint64_t monotonic_ms(void) { struct timespec ts; clock_gettime(CLOCK_MONOTONIC, &ts); - return (uint64_t)ts.tv_sec * 1000ull + (uint64_t)ts.tv_nsec / 1000000ull; + return (uint64_t)ts.tv_sec * 1000ULL + (uint64_t)ts.tv_nsec / 1000000ULL; } static receive_wait_t receive_wait_init(uint32_t timeout_ms, int abort_fd) diff --git a/src/streaming/stream-connector.c b/src/streaming/stream-connector.c index cfe02a1568f7d0..a94c6314857b1d 100644 --- a/src/streaming/stream-connector.c +++ b/src/streaming/stream-connector.c @@ -693,7 +693,7 @@ bool stream_connector_init(struct sender_state *s) { completion_init(&sc->completion); char tag[NETDATA_THREAD_TAG_MAX + 1]; - snprintfz(tag, NETDATA_THREAD_TAG_MAX, THREAD_TAG_STREAM_SENDER "-CN" "[%d]", + snprintfz(tag, NETDATA_THREAD_TAG_MAX, THREAD_TAG_STREAM_SENDER "-CN[%d]", sc->id); sc->thread = nd_thread_create(tag, NETDATA_THREAD_OPTION_DEFAULT, stream_connector_thread, sc); diff --git a/src/web/api/v2/api_v2_claim.c b/src/web/api/v2/api_v2_claim.c index 423a1d4098b4dc..a3ea2141f43a4e 100644 --- a/src/web/api/v2/api_v2_claim.c +++ b/src/web/api/v2/api_v2_claim.c @@ -219,7 +219,6 @@ static int api_claim(uint8_t version, struct web_client *w, char *url) { if(claim_agent(base_url, token, rooms, cloud_config_proxy_get(), cloud_config_insecure_get())) { msg = "ok"; - can_be_claimed = false; claim_reload_and_wait_online(); response = CLAIM_RESP_ACTION_OK; } diff --git a/src/web/websocket/websocket-send.c b/src/web/websocket/websocket-send.c index a836d2c061a39a..0d9000d539c610 100644 --- a/src/web/websocket/websocket-send.c +++ b/src/web/websocket/websocket-send.c @@ -78,7 +78,7 @@ static int websocket_protocol_send_frame( if(!wsc) return -1; - const char *disconnect_msg = ""; + const char *disconnect_msg; if (wsc->sock.fd < 0) { disconnect_msg = "Client not connected"; From 803c8239002f743acda4759c37f5f37080ea535f Mon Sep 17 00:00:00 2001 From: Ilya Mashchenko Date: Thu, 18 Jun 2026 21:35:36 +0300 Subject: [PATCH 9/9] chore(go.d/snmp_topology): remove dead cache snapshot path (#22773) * refactor(go.d/snmp_topology): remove dead cache snapshot path * refactor(go.d/snmp_topology): clarify dns resolver hook --- .../snmp_topology/func_topology_handler.go | 2 +- .../snmp_topology/topology_cache_test.go | 156 +++++++----------- .../collector/snmp_topology/topology_dns.go | 138 +--------------- .../topology_integration_test.go | 7 +- .../topology_observation_remote.go | 109 ------------ .../topology_observation_remote_cdp.go | 74 --------- .../topology_observation_remote_identity.go | 55 ------ .../topology_observation_remote_lldp.go | 85 ---------- .../topology_snapshot_builder.go | 45 +---- .../snmp_topology/topology_snmprec_test.go | 7 +- .../topology_test_helpers_test.go | 26 ++- 11 files changed, 92 insertions(+), 612 deletions(-) delete mode 100644 src/go/plugin/go.d/collector/snmp_topology/topology_observation_remote.go delete mode 100644 src/go/plugin/go.d/collector/snmp_topology/topology_observation_remote_cdp.go delete mode 100644 src/go/plugin/go.d/collector/snmp_topology/topology_observation_remote_identity.go delete mode 100644 src/go/plugin/go.d/collector/snmp_topology/topology_observation_remote_lldp.go diff --git a/src/go/plugin/go.d/collector/snmp_topology/func_topology_handler.go b/src/go/plugin/go.d/collector/snmp_topology/func_topology_handler.go index ce975fdcc5e38f..d9ad19bf3fc55f 100644 --- a/src/go/plugin/go.d/collector/snmp_topology/func_topology_handler.go +++ b/src/go/plugin/go.d/collector/snmp_topology/func_topology_handler.go @@ -35,7 +35,7 @@ func (f *funcTopology) Handle(_ context.Context, method string, params funcapi.R } options := resolveTopologyQueryOptions(params) - options.ResolveDNSName = resolveTopologyReverseDNSNameCached // never block on network I/O + options.ResolveDNSName = resolveTopologyReverseDNSNameNoop // never block on network I/O data, ok := f.registry.snapshotWithOptions(options) if !ok { return funcapi.UnavailableResponse("topology data not available yet, please retry after topology refresh") diff --git a/src/go/plugin/go.d/collector/snmp_topology/topology_cache_test.go b/src/go/plugin/go.d/collector/snmp_topology/topology_cache_test.go index b6eced3b6374fd..36555ff0cfae77 100644 --- a/src/go/plugin/go.d/collector/snmp_topology/topology_cache_test.go +++ b/src/go/plugin/go.d/collector/snmp_topology/topology_cache_test.go @@ -90,14 +90,13 @@ func TestTopologyCache_LldpSnapshot(t *testing.T) { tagLldpRemPortIDSubtype: "5", tagLldpRemPortDesc: "downlink", tagLldpRemSysName: "sw2", + tagLldpRemMgmtAddr: "10.0.0.2", }, }) coll.finalizeTopologyCache() - coll.topologyCache.mu.RLock() - data, ok := coll.topologyCache.snapshot() - coll.topologyCache.mu.RUnlock() + data, ok := snapshotTopologyCacheForTest(coll.topologyCache) require.True(t, ok) require.Len(t, data.Actors, 2) @@ -105,7 +104,7 @@ func TestTopologyCache_LldpSnapshot(t *testing.T) { link := data.Links[0] assert.Equal(t, "lldp", link.Protocol) - assert.Equal(t, "bidirectional", link.Direction) + assert.Equal(t, "unidirectional", link.Direction) assert.Equal(t, "Gi0/1", link.Src.Attributes["port_id"]) assert.Equal(t, "Gi0/2", link.Dst.Attributes["port_id"]) assert.Equal(t, "sw2", link.Dst.Attributes["sys_name"]) @@ -130,15 +129,13 @@ func TestTopologyCache_CdpSnapshot(t *testing.T) { address: "10.0.0.3", } - cache.mu.RLock() - data, ok := cache.snapshot() - cache.mu.RUnlock() + data, ok := snapshotTopologyCacheForTest(cache) require.True(t, ok) require.Len(t, data.Actors, 2) require.Len(t, data.Links, 1) assert.Equal(t, "cdp", data.Links[0].Protocol) - assert.Equal(t, "bidirectional", data.Links[0].Direction) + assert.Equal(t, "unidirectional", data.Links[0].Direction) assert.Equal(t, "Gi0/2", data.Links[0].Src.Attributes["if_name"]) assert.Equal(t, "Gi0/3", data.Links[0].Dst.Attributes["port_id"]) } @@ -291,14 +288,12 @@ func TestTopologyCache_CdpSnapshotHexAddress(t *testing.T) { address: "0a000003", } - cache.mu.RLock() - data, ok := cache.snapshot() - cache.mu.RUnlock() + data, ok := snapshotTopologyCacheForTest(cache) require.True(t, ok) require.Len(t, data.Links, 1) assert.Equal(t, "cdp", data.Links[0].Protocol) - assert.Equal(t, "bidirectional", data.Links[0].Direction) + assert.Equal(t, "unidirectional", data.Links[0].Direction) assert.True(t, linkHasRawAddressMetric(data.Links[0], "0a000003")) remote := findDeviceActorBySysName(data, "sw3") @@ -338,14 +333,14 @@ func TestTopologyCache_CdpSnapshotRawAddressWithoutIP(t *testing.T) { address: "edge-sw3.mgmt.local", } - cache.mu.RLock() - data, ok := cache.snapshot() - cache.mu.RUnlock() + options := defaultTopologyQueryOptionsForTest() + options.EliminateNonIPInferred = false + data, ok := snapshotTopologyCacheForTestWithOptions(cache, options) require.True(t, ok) require.Len(t, data.Links, 1) assert.Equal(t, "cdp", data.Links[0].Protocol) - assert.Equal(t, "bidirectional", data.Links[0].Direction) + assert.Equal(t, "unidirectional", data.Links[0].Direction) assert.True(t, linkHasRawAddressMetric(data.Links[0], "edge-sw3.mgmt.local")) } @@ -378,9 +373,38 @@ func TestTopologyCache_SnapshotBidirectionalPairMetadata(t *testing.T) { managementAddr: "10.0.0.2", } - cache.mu.RLock() - data, ok := cache.snapshot() - cache.mu.RUnlock() + remoteCache := newTopologyCache() + remoteCache.updateTime = cache.updateTime + remoteCache.lastUpdate = cache.lastUpdate + remoteCache.agentID = "agent2" + remoteCache.localDevice = topologyDevice{ + ChassisID: "aa:bb:cc:dd:ee:ff", + ChassisIDType: "macAddress", + SysName: "sw2", + ManagementIP: "10.0.0.2", + } + remoteCache.lldpLocPorts["2"] = &lldpLocPort{ + portNum: "2", + portID: "Gi0/2", + portIDSubtype: "interfaceName", + portDesc: "downlink", + } + remoteCache.lldpRemotes["2:1"] = &lldpRemote{ + localPortNum: "2", + remIndex: "1", + chassisID: "00:11:22:33:44:55", + chassisIDSubtype: "macAddress", + portID: "Gi0/1", + portIDSubtype: "interfaceName", + portDesc: "uplink", + sysName: "sw1", + managementAddr: "10.0.0.1", + } + + registry := newTopologyRegistry() + registry.register(cache) + registry.register(remoteCache) + data, ok := snapshotTopologyRegistryForTest(registry) require.True(t, ok) require.Len(t, data.Links, 1) @@ -427,9 +451,7 @@ func TestTopologyCache_SnapshotMergesRemoteIdentityAcrossProtocols(t *testing.T) address: "10.0.0.2", } - cache.mu.RLock() - data, ok := cache.snapshot() - cache.mu.RUnlock() + data, ok := snapshotTopologyCacheForTest(cache) require.True(t, ok) require.Equal(t, 2, countDeviceActors(data)) @@ -525,9 +547,7 @@ func TestTopologyCache_LLDPManagementAddressesAndCaps(t *testing.T) { coll.finalizeTopologyCache() - coll.topologyCache.mu.RLock() - data, ok := coll.topologyCache.snapshot() - coll.topologyCache.mu.RUnlock() + data, ok := snapshotTopologyCacheForTest(coll.topologyCache) require.True(t, ok) require.Greater(t, len(data.Actors), 1) @@ -562,9 +582,7 @@ func TestTopologyCache_CDPManagementAddresses(t *testing.T) { tagCdpSecondaryMgmtAddr: "0a000004", }) - cache.mu.RLock() - data, ok := cache.snapshot() - cache.mu.RUnlock() + data, ok := snapshotTopologyCacheForTest(cache) require.True(t, ok) require.True(t, containsMgmtAddr(data, map[string]struct{}{"10.0.0.3": {}, "10.0.0.4": {}})) @@ -602,9 +620,9 @@ func TestTopologyCache_FDBAndARPEnrichment(t *testing.T) { tagArpState: "reachable", }) - cache.mu.RLock() - data, ok := cache.snapshot() - cache.mu.RUnlock() + options := defaultTopologyQueryOptionsForTest() + options.MapType = topologyMapTypeAllDevicesLowConfidence + data, ok := snapshotTopologyCacheForTestWithOptions(cache, options) require.True(t, ok) require.GreaterOrEqual(t, len(data.Actors), 2) @@ -1034,9 +1052,9 @@ func TestTopologyCache_SnapshotDeterministicEndpointIPSelection(t *testing.T) { expectedIPs := []string{"10.20.4.205", "10.20.4.60"} for range 25 { - cache.mu.RLock() - data, ok := cache.snapshot() - cache.mu.RUnlock() + options := defaultTopologyQueryOptionsForTest() + options.MapType = topologyMapTypeAllDevicesLowConfidence + data, ok := snapshotTopologyCacheForTestWithOptions(cache, options) require.True(t, ok) ep := findActorByMAC(data, "d8:5e:d3:0e:c5:e6") @@ -1079,9 +1097,7 @@ func TestTopologyCache_SnapshotDeterministicOrdering(t *testing.T) { address: "10.0.0.3", } - cache.mu.RLock() - data, ok := cache.snapshot() - cache.mu.RUnlock() + data, ok := snapshotTopologyCacheForTest(cache) require.True(t, ok) require.NotEmpty(t, data.Actors) @@ -1104,69 +1120,6 @@ func TestTopologyCache_SnapshotDeterministicOrdering(t *testing.T) { assert.Equal(t, expectedLinkOrder, linkOrder) } -func TestTopologyCache_BuildEngineObservations_SeparatesProtocolSpecificRemoteObservations(t *testing.T) { - cache := newTopologyCache() - cache.localDevice = topologyDevice{ - ChassisID: "00:11:22:33:44:55", - ChassisIDType: "macAddress", - SysName: "sw-a", - ManagementIP: "10.0.0.1", - } - cache.lldpLocPorts["1"] = &lldpLocPort{ - portNum: "1", - portID: "Gi0/1", - portIDSubtype: "interfaceName", - portDesc: "uplink", - } - cache.lldpRemotes["1:1"] = &lldpRemote{ - localPortNum: "1", - remIndex: "1", - chassisID: "aa:bb:cc:dd:ee:ff", - chassisIDSubtype: "macAddress", - portID: "Gi0/2", - portIDSubtype: "interfaceName", - portDesc: "downlink", - sysName: "sw-b", - managementAddr: "10.0.0.2", - } - cache.cdpRemotes["1:1"] = &cdpRemote{ - ifIndex: "1", - ifName: "Gi0/1", - deviceID: "sw-b", - sysName: "switch-b", - devicePort: "Gi0/2", - address: "10.0.0.2", - } - - observations, localDeviceID := cache.buildEngineObservations(cache.localDevice) - require.Equal(t, "macAddress:00:11:22:33:44:55", localDeviceID) - require.Len(t, observations, 3) - require.Equal(t, localDeviceID, observations[0].DeviceID) - - var lldpObservation *topologyengine.L2Observation - var cdpObservation *topologyengine.L2Observation - for i := 1; i < len(observations); i++ { - observation := &observations[i] - switch { - case len(observation.LLDPRemotes) > 0: - lldpObservation = observation - case len(observation.CDPRemotes) > 0: - cdpObservation = observation - } - } - - require.NotNil(t, lldpObservation) - require.NotNil(t, cdpObservation) - require.Equal(t, lldpObservation.DeviceID, cdpObservation.DeviceID) - require.Equal(t, "macAddress:aa:bb:cc:dd:ee:ff", lldpObservation.DeviceID) - require.Equal(t, "10.0.0.2", lldpObservation.ManagementIP) - require.Equal(t, "10.0.0.2", cdpObservation.ManagementIP) - require.Equal(t, "sw-b", lldpObservation.Hostname) - require.Equal(t, "switch-b", cdpObservation.Hostname) - require.Len(t, lldpObservation.LLDPRemotes, 1) - require.Len(t, cdpObservation.CDPRemotes, 1) -} - func TestTopologyObservationIdentityResolver_ReusesStableRemoteIdentityAcrossSignals(t *testing.T) { resolver := newTopologyObservationIdentityResolver(topologyengine.L2Observation{ DeviceID: "macAddress:00:11:22:33:44:55", @@ -1626,6 +1579,9 @@ func linkHasRawAddressMetric(link topologyLink, raw string) bool { if raw == "" || len(link.Metrics) == 0 { return false } + if value, ok := link.Metrics["remote_address_raw"].(string); ok && value == raw { + return true + } srcRaw, srcOK := link.Metrics["src_remote_address_raw"].(string) dstRaw, dstOK := link.Metrics["dst_remote_address_raw"].(string) return (srcOK && srcRaw == raw) || (dstOK && dstRaw == raw) diff --git a/src/go/plugin/go.d/collector/snmp_topology/topology_dns.go b/src/go/plugin/go.d/collector/snmp_topology/topology_dns.go index 059da2e0e82f43..936d20ebd7e739 100644 --- a/src/go/plugin/go.d/collector/snmp_topology/topology_dns.go +++ b/src/go/plugin/go.d/collector/snmp_topology/topology_dns.go @@ -2,140 +2,8 @@ package snmptopology -import ( - "context" - "net" - "net/netip" - "sort" - "strings" - "sync" - "time" -) - -const ( - topologyReverseDNSTimeout = 50 * time.Millisecond - topologyReverseDNSCacheTTL = 10 * time.Minute - topologyReverseDNSNegTTL = 30 * time.Second -) - -type topologyReverseDNSCacheEntry struct { - name string - expiresAt time.Time -} - -type topologyReverseDNSResolver struct { - mu sync.RWMutex - timeout time.Duration - ttl time.Duration - cache map[string]topologyReverseDNSCacheEntry -} - -func newTopologyReverseDNSResolver(timeout, ttl time.Duration) *topologyReverseDNSResolver { - return &topologyReverseDNSResolver{ - timeout: timeout, - ttl: ttl, - cache: make(map[string]topologyReverseDNSCacheEntry), - } -} - -// lookupCached returns the cached result for ip without performing any network I/O. -// Returns "" when the IP has never been resolved or its cache entry has expired. -func (r *topologyReverseDNSResolver) lookupCached(ip string) string { - if r == nil { - return "" - } - addr, err := netip.ParseAddr(strings.TrimSpace(ip)) - if err != nil || !addr.IsValid() { - return "" - } - ip = addr.Unmap().String() - - r.mu.RLock() - entry, ok := r.cache[ip] - r.mu.RUnlock() - if ok && time.Now().Before(entry.expiresAt) { - return entry.name - } +// resolveTopologyReverseDNSNameNoop is the non-blocking DNS resolver hook used +// during function responses. SNMP topology has no reverse-DNS warmer today. +func resolveTopologyReverseDNSNameNoop(_ string) string { return "" } - -func (r *topologyReverseDNSResolver) lookup(ip string) string { - if r == nil { - return "" - } - addr, err := netip.ParseAddr(strings.TrimSpace(ip)) - if err != nil || !addr.IsValid() { - return "" - } - ip = addr.Unmap().String() - now := time.Now() - - r.mu.RLock() - entry, ok := r.cache[ip] - r.mu.RUnlock() - if ok && now.Before(entry.expiresAt) { - return entry.name - } - - ctx, cancel := context.WithTimeout(context.Background(), r.timeout) - defer cancel() - names, err := net.DefaultResolver.LookupAddr(ctx, ip) - resolved := "" - if err == nil { - resolved = topologyNormalizeReverseDNSName(names) - } - ttl := r.ttl - if resolved == "" && topologyReverseDNSNegTTL > 0 { - ttl = topologyReverseDNSNegTTL - } - - r.mu.Lock() - r.cache[ip] = topologyReverseDNSCacheEntry{ - name: resolved, - expiresAt: now.Add(ttl), - } - r.mu.Unlock() - - return resolved -} - -func topologyNormalizeReverseDNSName(names []string) string { - if len(names) == 0 { - return "" - } - seen := make(map[string]struct{}, len(names)) - out := make([]string, 0, len(names)) - for _, name := range names { - name = strings.TrimSpace(name) - name = strings.TrimSuffix(name, ".") - name = strings.ToLower(name) - if name == "" { - continue - } - if _, ok := seen[name]; ok { - continue - } - seen[name] = struct{}{} - out = append(out, name) - } - if len(out) == 0 { - return "" - } - sort.Strings(out) - return out[0] -} - -var defaultTopologyReverseDNSResolver = newTopologyReverseDNSResolver(topologyReverseDNSTimeout, topologyReverseDNSCacheTTL) - -// resolveTopologyReverseDNSName performs a live DNS lookup (with cache). -// Used while building topology snapshots to warm the cache. -func resolveTopologyReverseDNSName(ip string) string { - return defaultTopologyReverseDNSResolver.lookup(ip) -} - -// resolveTopologyReverseDNSNameCached returns a cached DNS name if available, -// or an empty string if the IP has not been resolved yet. Never blocks on network I/O. -// Used during function responses to avoid external calls. -func resolveTopologyReverseDNSNameCached(ip string) string { - return defaultTopologyReverseDNSResolver.lookupCached(ip) -} diff --git a/src/go/plugin/go.d/collector/snmp_topology/topology_integration_test.go b/src/go/plugin/go.d/collector/snmp_topology/topology_integration_test.go index 24e93359f8f09e..fd3264df2bdac0 100644 --- a/src/go/plugin/go.d/collector/snmp_topology/topology_integration_test.go +++ b/src/go/plugin/go.d/collector/snmp_topology/topology_integration_test.go @@ -111,11 +111,12 @@ func collectTopologySnapshotFromDevice(t *testing.T, dev ddsnmp.DeviceConnection if cache == nil { return false } - cache.mu.RLock() - defer cache.mu.RUnlock() var ok bool - snapshot, ok = cache.snapshot() + options := defaultTopologyQueryOptionsForTest() + options.CollapseActorsByIP = false + options.EliminateNonIPInferred = false + snapshot, ok = snapshotTopologyCacheForTestWithOptions(cache, options) return ok }, 5*time.Second, 100*time.Millisecond, "topology snapshot did not become available for %q", dev.SysName) diff --git a/src/go/plugin/go.d/collector/snmp_topology/topology_observation_remote.go b/src/go/plugin/go.d/collector/snmp_topology/topology_observation_remote.go deleted file mode 100644 index e730efec5bf682..00000000000000 --- a/src/go/plugin/go.d/collector/snmp_topology/topology_observation_remote.go +++ /dev/null @@ -1,109 +0,0 @@ -// SPDX-License-Identifier: GPL-3.0-or-later - -package snmptopology - -import ( - "sort" - "strings" - - topologyengine "github.com/netdata/netdata/go/plugins/pkg/l2topology" -) - -type topologyRemoteObservationBuilder struct { - cache *topologyCache - local topologyDevice - localObservation topologyengine.L2Observation - localManagementIP string - localSysName string - localGlobalID string - - resolver *topologyObservationIdentityResolver - remoteObservations map[string]*topologyengine.L2Observation - remoteOrder []string - remoteManagementByID map[string]string - remoteChassisByID map[string]string -} - -func newTopologyRemoteObservationBuilder(cache *topologyCache, local topologyDevice, localObservation topologyengine.L2Observation) *topologyRemoteObservationBuilder { - localManagementIP := normalizeIPAddress(local.ManagementIP) - if localManagementIP == "" { - localManagementIP = pickManagementIP(local.ManagementAddresses) - } - - localGlobalID := strings.TrimSpace(localObservation.Hostname) - if localGlobalID == "" { - localGlobalID = localObservation.DeviceID - } - - return &topologyRemoteObservationBuilder{ - cache: cache, - local: local, - localObservation: localObservation, - localManagementIP: localManagementIP, - localSysName: strings.TrimSpace(local.SysName), - localGlobalID: localGlobalID, - resolver: newTopologyObservationIdentityResolver(localObservation), - remoteObservations: make(map[string]*topologyengine.L2Observation), - remoteOrder: make([]string, 0, len(cache.lldpRemotes)+len(cache.cdpRemotes)), - remoteManagementByID: make(map[string]string), - remoteChassisByID: make(map[string]string), - } -} - -func (c *topologyCache) buildEngineObservations(local topologyDevice) ([]topologyengine.L2Observation, string) { - localObservation := c.buildEngineObservation(local) - localObservation.DeviceID = strings.TrimSpace(localObservation.DeviceID) - if localObservation.DeviceID == "" { - return nil, "" - } - - builder := newTopologyRemoteObservationBuilder(c, local, localObservation) - builder.collectLLDPRemoteObservations() - builder.collectCDPRemoteObservations() - - return builder.observations(), localObservation.DeviceID -} - -func (b *topologyRemoteObservationBuilder) observations() []topologyengine.L2Observation { - observations := make([]topologyengine.L2Observation, 0, 1+len(b.remoteObservations)) - observations = append(observations, b.localObservation) - - sort.Strings(b.remoteOrder) - for _, key := range b.remoteOrder { - entry := b.remoteObservations[key] - if entry == nil { - continue - } - if entry.ManagementIP == "" { - entry.ManagementIP = b.remoteManagementByID[entry.DeviceID] - } - if entry.ChassisID == "" { - entry.ChassisID = b.remoteChassisByID[entry.DeviceID] - } - if len(entry.LLDPRemotes) == 0 && len(entry.CDPRemotes) == 0 { - continue - } - if entry.Hostname == "" { - entry.Hostname = entry.DeviceID - } - observations = append(observations, *entry) - } - - return observations -} - -func selectTopologyRemoteHostname(current, candidate, deviceID string) string { - current = strings.TrimSpace(current) - candidate = strings.TrimSpace(candidate) - deviceID = strings.TrimSpace(deviceID) - if candidate == "" { - if current != "" { - return current - } - return deviceID - } - if current == "" || current == deviceID { - return candidate - } - return current -} diff --git a/src/go/plugin/go.d/collector/snmp_topology/topology_observation_remote_cdp.go b/src/go/plugin/go.d/collector/snmp_topology/topology_observation_remote_cdp.go deleted file mode 100644 index 114d63bb110115..00000000000000 --- a/src/go/plugin/go.d/collector/snmp_topology/topology_observation_remote_cdp.go +++ /dev/null @@ -1,74 +0,0 @@ -// SPDX-License-Identifier: GPL-3.0-or-later - -package snmptopology - -import ( - "sort" - "strings" - - topologyengine "github.com/netdata/netdata/go/plugins/pkg/l2topology" -) - -func (b *topologyRemoteObservationBuilder) collectCDPRemoteObservations() { - keys := make([]string, 0, len(b.cache.cdpRemotes)) - for key := range b.cache.cdpRemotes { - keys = append(keys, key) - } - sort.Strings(keys) - - for _, key := range keys { - remote := b.cache.cdpRemotes[key] - if remote == nil { - continue - } - - remoteDeviceToken := strings.TrimSpace(remote.deviceID) - remoteSysName := strings.TrimSpace(remote.sysName) - remoteManagementIP := normalizeIPAddress(remote.address) - if remoteManagementIP == "" { - remoteManagementIP = pickManagementIP(remote.managementAddrs) - } - if remoteManagementIP == "" && remoteDeviceToken == "" && remoteSysName == "" { - continue - } - - remoteDeviceID := b.resolver.resolve( - []string{remoteDeviceToken, remoteSysName}, - "", - "", - remoteManagementIP, - ) - if remoteDeviceID == "" || remoteDeviceID == b.localObservation.DeviceID { - continue - } - b.updateRemoteIdentity(remoteDeviceID, remoteManagementIP, "") - - remoteIfName := strings.TrimSpace(remote.devicePort) - localIfName := strings.TrimSpace(remote.ifName) - if localIfName == "" && strings.TrimSpace(remote.ifIndex) != "" { - localIfName = strings.TrimSpace(b.cache.ifNamesByIndex[remote.ifIndex]) - } - if remoteIfName == "" || localIfName == "" { - continue - } - - remoteObservation := b.ensureRemoteObservation( - "cdp", - remoteDeviceID, - firstNonEmpty(remoteSysName, remoteDeviceToken, remoteDeviceID), - remoteManagementIP, - "", - ) - if remoteObservation == nil { - continue - } - - remoteObservation.CDPRemotes = append(remoteObservation.CDPRemotes, topologyengine.CDPRemoteObservation{ - LocalIfName: remoteIfName, - DeviceID: b.localGlobalID, - SysName: b.localSysName, - DevicePort: localIfName, - Address: b.localManagementIP, - }) - } -} diff --git a/src/go/plugin/go.d/collector/snmp_topology/topology_observation_remote_identity.go b/src/go/plugin/go.d/collector/snmp_topology/topology_observation_remote_identity.go deleted file mode 100644 index bfa439191a1938..00000000000000 --- a/src/go/plugin/go.d/collector/snmp_topology/topology_observation_remote_identity.go +++ /dev/null @@ -1,55 +0,0 @@ -// SPDX-License-Identifier: GPL-3.0-or-later - -package snmptopology - -import ( - "strings" - - topologyengine "github.com/netdata/netdata/go/plugins/pkg/l2topology" -) - -func (b *topologyRemoteObservationBuilder) updateRemoteIdentity(deviceID, managementIP, chassisID string) { - deviceID = strings.TrimSpace(deviceID) - if deviceID == "" { - return - } - if managementIP = canonicalObservationIP(managementIP); managementIP != "" { - if _, ok := b.remoteManagementByID[deviceID]; !ok { - b.remoteManagementByID[deviceID] = managementIP - } - } - if chassisID = strings.TrimSpace(chassisID); chassisID != "" { - if _, ok := b.remoteChassisByID[deviceID]; !ok { - b.remoteChassisByID[deviceID] = chassisID - } - } -} - -func (b *topologyRemoteObservationBuilder) ensureRemoteObservation(protocol, deviceID, hostname, managementIP, chassisID string) *topologyengine.L2Observation { - deviceID = strings.TrimSpace(deviceID) - if deviceID == "" { - return nil - } - - key := protocol + "|" + deviceID - entry := b.remoteObservations[key] - if entry == nil { - entry = &topologyengine.L2Observation{ - DeviceID: deviceID, - Inferred: true, - } - b.remoteObservations[key] = entry - b.remoteOrder = append(b.remoteOrder, key) - } - - entry.Hostname = selectTopologyRemoteHostname(entry.Hostname, hostname, deviceID) - b.updateRemoteIdentity(deviceID, managementIP, chassisID) - if entry.ManagementIP == "" { - entry.ManagementIP = b.remoteManagementByID[deviceID] - } - if entry.ChassisID == "" { - entry.ChassisID = b.remoteChassisByID[deviceID] - } - b.resolver.register(deviceID, []string{entry.Hostname}, entry.ChassisID, entry.ManagementIP) - return entry -} diff --git a/src/go/plugin/go.d/collector/snmp_topology/topology_observation_remote_lldp.go b/src/go/plugin/go.d/collector/snmp_topology/topology_observation_remote_lldp.go deleted file mode 100644 index bea84a06e5a196..00000000000000 --- a/src/go/plugin/go.d/collector/snmp_topology/topology_observation_remote_lldp.go +++ /dev/null @@ -1,85 +0,0 @@ -// SPDX-License-Identifier: GPL-3.0-or-later - -package snmptopology - -import ( - "sort" - "strings" - - topologyengine "github.com/netdata/netdata/go/plugins/pkg/l2topology" -) - -func (b *topologyRemoteObservationBuilder) collectLLDPRemoteObservations() { - keys := make([]string, 0, len(b.cache.lldpRemotes)) - for key := range b.cache.lldpRemotes { - keys = append(keys, key) - } - sort.Strings(keys) - - for _, key := range keys { - remote := b.cache.lldpRemotes[key] - if remote == nil { - continue - } - - remoteSysName := strings.TrimSpace(remote.sysName) - remoteChassisID := strings.TrimSpace(remote.chassisID) - remoteManagementIP := normalizeIPAddress(remote.managementAddr) - if remoteManagementIP == "" { - remoteManagementIP = pickManagementIP(remote.managementAddrs) - } - - remoteDeviceID := b.resolver.resolve( - []string{remoteSysName}, - remoteChassisID, - strings.TrimSpace(remote.chassisIDSubtype), - remoteManagementIP, - ) - if remoteDeviceID == "" || remoteDeviceID == b.localObservation.DeviceID { - continue - } - b.updateRemoteIdentity(remoteDeviceID, remoteManagementIP, remoteChassisID) - - remoteObservation := b.ensureRemoteObservation( - "lldp", - remoteDeviceID, - firstNonEmpty(remoteSysName, remoteDeviceID), - remoteManagementIP, - remoteChassisID, - ) - if remoteObservation == nil { - continue - } - - localPort := b.cache.lldpLocPorts[remote.localPortNum] - localPortID := "" - localPortIDSubtype := "" - localPortDesc := "" - if localPort != nil { - localPortID = strings.TrimSpace(localPort.portID) - localPortIDSubtype = strings.TrimSpace(localPort.portIDSubtype) - localPortDesc = strings.TrimSpace(localPort.portDesc) - } - - if strings.TrimSpace(remote.portID) == "" && - strings.TrimSpace(remote.portDesc) == "" && - localPortID == "" && - localPortDesc == "" { - continue - } - - remoteObservation.LLDPRemotes = append(remoteObservation.LLDPRemotes, topologyengine.LLDPRemoteObservation{ - LocalPortNum: strings.TrimSpace(remote.remIndex), - RemoteIndex: strings.TrimSpace(remote.localPortNum), - LocalPortID: strings.TrimSpace(remote.portID), - LocalPortIDSubtype: strings.TrimSpace(remote.portIDSubtype), - LocalPortDesc: strings.TrimSpace(remote.portDesc), - ChassisID: strings.TrimSpace(b.local.ChassisID), - SysName: b.localSysName, - PortID: localPortID, - PortIDSubtype: localPortIDSubtype, - PortDesc: localPortDesc, - ManagementIP: b.localManagementIP, - }) - } -} diff --git a/src/go/plugin/go.d/collector/snmp_topology/topology_snapshot_builder.go b/src/go/plugin/go.d/collector/snmp_topology/topology_snapshot_builder.go index eae52866e52168..018936f5e234ff 100644 --- a/src/go/plugin/go.d/collector/snmp_topology/topology_snapshot_builder.go +++ b/src/go/plugin/go.d/collector/snmp_topology/topology_snapshot_builder.go @@ -2,12 +2,7 @@ package snmptopology -import ( - "time" - - topologyengine "github.com/netdata/netdata/go/plugins/pkg/l2topology" - "github.com/netdata/netdata/go/plugins/plugin/go.d/collector/snmp/ddsnmp" -) +import "github.com/netdata/netdata/go/plugins/plugin/go.d/collector/snmp/ddsnmp" func buildLocalTopologyDevice(dev ddsnmp.DeviceConnectionInfo) topologyDevice { device := topologyDevice{ @@ -72,41 +67,3 @@ func buildLocalTopologyDevice(dev ddsnmp.DeviceConnectionInfo) topologyDevice { return device } - -func (c *topologyCache) snapshot() (topologyData, bool) { - if !c.hasFreshSnapshotAt(time.Now()) { - return topologyData{}, false - } - - local := c.localDevice - local = normalizeTopologyDevice(local) - - observations, localDeviceID := c.buildEngineObservations(local) - if len(observations) == 0 { - return topologyData{}, false - } - - result, err := topologyengine.BuildL2ResultFromObservations(observations, topologyengine.DiscoverOptions{ - EnableLLDP: true, - EnableCDP: true, - EnableBridge: true, - EnableARP: true, - }) - if err != nil { - return topologyData{}, false - } - - data := topologyengine.ToGraph(result, topologyengine.GraphOptions{ - SchemaVersion: topologySchemaVersion, - Source: "snmp", - Layer: "2", - View: "summary", - AgentID: c.agentID, - LocalDeviceID: localDeviceID, - CollectedAt: c.lastUpdate, - ResolveDNSName: resolveTopologyReverseDNSName, - }) - - augmentLocalActorFromCache(&data, local) - return data, true -} diff --git a/src/go/plugin/go.d/collector/snmp_topology/topology_snmprec_test.go b/src/go/plugin/go.d/collector/snmp_topology/topology_snmprec_test.go index e2d116f4a203cb..89130d4717bbe9 100644 --- a/src/go/plugin/go.d/collector/snmp_topology/topology_snmprec_test.go +++ b/src/go/plugin/go.d/collector/snmp_topology/topology_snmprec_test.go @@ -67,9 +67,10 @@ func TestTopologyCache_RealSnmprecFixtures(t *testing.T) { } coll.finalizeTopologyCache() - coll.topologyCache.mu.RLock() - snapshot, ok := coll.topologyCache.snapshot() - coll.topologyCache.mu.RUnlock() + options := defaultTopologyQueryOptionsForTest() + options.CollapseActorsByIP = false + options.EliminateNonIPInferred = false + snapshot, ok := snapshotTopologyCacheForTestWithOptions(coll.topologyCache, options) require.True(t, ok) require.GreaterOrEqual(t, len(snapshot.Actors), 1) diff --git a/src/go/plugin/go.d/collector/snmp_topology/topology_test_helpers_test.go b/src/go/plugin/go.d/collector/snmp_topology/topology_test_helpers_test.go index 9f892f8c2352dd..38b244de67f668 100644 --- a/src/go/plugin/go.d/collector/snmp_topology/topology_test_helpers_test.go +++ b/src/go/plugin/go.d/collector/snmp_topology/topology_test_helpers_test.go @@ -25,15 +25,35 @@ func registerTestDeviceState(store *ddsnmp.DeviceStore, devices ...ddsnmp.Device } func snapshotTopologyRegistryForTest(registry *topologyRegistry) (topologyData, bool) { - return registry.snapshotWithOptions(topologyQueryOptions{ + return snapshotTopologyRegistryForTestWithOptions(registry, defaultTopologyQueryOptionsForTest()) +} + +func snapshotTopologyRegistryForTestWithOptions(registry *topologyRegistry, options topologyQueryOptions) (topologyData, bool) { + if options.ResolveDNSName == nil { + options.ResolveDNSName = resolveTopologyReverseDNSNameNoop + } + return registry.snapshotWithOptions(options) +} + +func snapshotTopologyCacheForTest(cache *topologyCache) (topologyData, bool) { + return snapshotTopologyCacheForTestWithOptions(cache, defaultTopologyQueryOptionsForTest()) +} + +func snapshotTopologyCacheForTestWithOptions(cache *topologyCache, options topologyQueryOptions) (topologyData, bool) { + registry := newTopologyRegistry() + registry.register(cache) + return snapshotTopologyRegistryForTestWithOptions(registry, options) +} + +func defaultTopologyQueryOptionsForTest() topologyQueryOptions { + return topologyQueryOptions{ CollapseActorsByIP: true, EliminateNonIPInferred: true, MapType: topologyMapTypeLLDPCDPManaged, InferenceStrategy: topologyInferenceStrategyFDBMinimumKnowledge, ManagedDeviceFocus: topologyManagedFocusAllDevices, Depth: topologyDepthAllInternal, - ResolveDNSName: resolveTopologyReverseDNSName, - }) + } } func containsMgmtAddr(snapshot topologyData, addrs map[string]struct{}) bool {