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 dc63fc020fe2eb..f47de2eed98434 100644 --- a/.agents/skills/project-writing-go-modules-framework-v2/SKILL.md +++ b/.agents/skills/project-writing-go-modules-framework-v2/SKILL.md @@ -50,6 +50,11 @@ source files for evidence. embedded `charts.yaml` is RECOMMENDED. - `Collect(ctx)` MUST return `error` and write metrics to `metrix`; it MUST NOT return a V1 `map[string]int64`. +- Long-running side-effect loops that must start only with the running job MAY + implement optional `collectorapi.CollectorV2Runner`. `Run(ctx)` MUST return + promptly after cancellation. Do not start operational polling from `Init()` or + `Check()`, because DynCfg `test` and autodetection use those methods without + starting the runtime job. - Collector `Cleanup(ctx)` MUST be idempotent. The framework may call it more than once, including after partial `Init` / `Check` setup. - Files SHOULD stay boring: public lifecycle methods in `collector.go`, setup diff --git a/docs/netdata-cloud/authentication-and-authorization/role-based-access-model.md b/docs/netdata-cloud/authentication-and-authorization/role-based-access-model.md index a5dc9ab04fce89..a95336d08faa73 100644 --- a/docs/netdata-cloud/authentication-and-authorization/role-based-access-model.md +++ b/docs/netdata-cloud/authentication-and-authorization/role-based-access-model.md @@ -26,6 +26,10 @@ You can control what functionalities users can access in Netdata Cloud through t | **View-only access** - monitor specific systems without making changes | **Observer** | | **Billing management** - handle invoices and payments without system access | **Billing** | +## Role Change Propagation + +Role changes take effect immediately. When an Admin or Manager changes a user's role, the updated permissions are applied right away by the Netdata Cloud backend. + ## Quick Reference
diff --git a/src/collectors/statsd.plugin/README.md b/src/collectors/statsd.plugin/README.md index 5581c41613d739..1d1019a655e613 100644 --- a/src/collectors/statsd.plugin/README.md +++ b/src/collectors/statsd.plugin/README.md @@ -496,7 +496,13 @@ For example, to monitor the application `myapp` using StatsD and Netdata, create Using this configuration, `myapp` gets its own dashboard section with one chart containing two [dimensions](https://learn.netdata.cloud/docs/developer-and-contributor-corner/glossary#d). -When you send metrics like `foo:10|g` and `bar:20|g`, you'll see both private charts and your synthetic chart. +When you send metrics like `myapp.metric1:10|g` and `myapp.metric2:20|g`, you'll see both private charts and your synthetic chart. These metric names must match the pattern defined in the `[app]` section (e.g., `myapp.*`) for them to appear in your synthetic charts. + +:::note + +**Synthetic chart appears empty or is missing?** This happens when the metric names you send don't match the `metrics` pattern in your `[app]` section. StatsD matches incoming metric names against the `metrics` pattern using Netdata's [simple pattern](/src/libnetdata/simple_pattern/README.md) syntax — if a metric name doesn't match, it is never linked to the app's synthetic charts. For example, with `metrics = myapp.*`, sending bare names like `foo:10|g` creates a private chart for `foo` but never feeds the synthetic chart. To fix this, send metric names that include the prefix matching the pattern (e.g., `myapp.foo:10|g`). + +:::
Synthetic Chart Example diff --git a/src/database/contexts/rrdcontext.c b/src/database/contexts/rrdcontext.c index 766dd1c96c75d9..a73ccc93971963 100644 --- a/src/database/contexts/rrdcontext.c +++ b/src/database/contexts/rrdcontext.c @@ -324,7 +324,7 @@ static void rrdcontext_checkpoint_execute(RRDHOST *host, const char *claim_id, c return; if(rrdhost_flag_check(host, RRDHOST_FLAG_ACLK_STREAM_CONTEXTS)) { - nd_log(NDLS_DAEMON, NDLP_NOTICE, + nd_log(NDLS_DAEMON, NDLP_DEBUG, "RRDCONTEXT: checkpoint for claim id '%s', node id '%s', " "while node '%s' has an active context streaming.", claim_id, node_id, rrdhost_hostname(host)); diff --git a/src/go/pkg/funcapi/response.go b/src/go/pkg/funcapi/response.go index 3efdf94d49fcb1..02c3c6da4e5692 100644 --- a/src/go/pkg/funcapi/response.go +++ b/src/go/pkg/funcapi/response.go @@ -23,9 +23,8 @@ type MethodConfig struct { // RawRequest routes the complete Function request to a RawMethodHandler. // Use this for Function APIs that need raw payloads, args, or full response envelopes. RawRequest bool - // FIXME: AgentWide currently removes __job from the public API, but funcctl still - // dispatches through the first running job for the module instead of a true - // agent-level execution path. + // AgentWide marks a module/static method as agent-level: funcctl omits + // __job from the public API and dispatches the method without a RuntimeJob. AgentWide bool // Method is agent-wide (does not require __job selector) RequiredParams []ParamConfig // Required parameters for this method (including __sort if used) // FIXME: Presentation is intentionally untyped here, while the shared UI schema diff --git a/src/go/plugin/agent/jobmgr/dyncfg_collector_test.go b/src/go/plugin/agent/jobmgr/dyncfg_collector_test.go index afd463a768d635..a6e419931aaf4c 100644 --- a/src/go/plugin/agent/jobmgr/dyncfg_collector_test.go +++ b/src/go/plugin/agent/jobmgr/dyncfg_collector_test.go @@ -9,6 +9,7 @@ import ( "errors" "os" "strings" + "sync/atomic" "testing" "time" @@ -53,6 +54,11 @@ type jobNameRequiredV2Collector struct { store metrix.CollectorStore } +type runnerProbeV2Collector struct { + *jobNameRequiredV2Collector + runCalled atomic.Bool +} + func newJobNameRequiredV2Collector() *jobNameRequiredV2Collector { return &jobNameRequiredV2Collector{ store: metrix.NewCollectorStore(), @@ -79,6 +85,11 @@ func (c *jobNameRequiredV2Collector) MetricStore() metrix.CollectorStore { } func (c *jobNameRequiredV2Collector) ChartTemplateYAML() string { return "" } +func (c *runnerProbeV2Collector) Run(context.Context) error { + c.runCalled.Store(true) + return nil +} + func TestDyncfgConfigUserconfig_InvalidPayload_Returns400Only(t *testing.T) { tests := map[string]struct { contentType string @@ -210,6 +221,64 @@ func TestDyncfgCmdTest_PassesJobNameToV2Collector(t *testing.T) { assert.Equal(t, float64(200), resp["status"]) } +func TestDyncfgCmdTest_DoesNotRunV2CollectorRunner(t *testing.T) { + var buf bytes.Buffer + probe := &runnerProbeV2Collector{jobNameRequiredV2Collector: newJobNameRequiredV2Collector()} + + mgr := newCollectorTestManager() + mgr.SetDyncfgResponder(dyncfg.NewResponder(netdataapi.New(safewriter.New(&buf)))) + mgr.modules.Register("runnerprobe", collectorapi.Creator{ + CreateV2: func() collectorapi.CollectorV2 { + return probe + }, + }) + + fn := dyncfg.NewFunction(functions.Function{ + UID: "test-v2-runner", + ContentType: "application/json", + Payload: mustMarshalCollectorConfigPayload(t, prepareDyncfgCfg("runnerprobe", "payload-name")), + Args: []string{mgr.dyncfgModID("runnerprobe"), string(dyncfg.CommandTest), "tested-job"}, + }) + + mgr.dyncfgCmdTest(fn) + mgr.cmdTestWG.Wait() + + var resp map[string]any + mustDecodeFunctionPayload(t, buf.String(), "test-v2-runner", &resp) + assert.Equal(t, float64(200), resp["status"]) + require.Never(t, probe.runCalled.Load, 100*time.Millisecond, 10*time.Millisecond, "runner started during dyncfg test") +} + +func TestDyncfgCmdTest_SingleInstanceV2RunnerUsesCanonicalName(t *testing.T) { + var buf bytes.Buffer + probe := &runnerProbeV2Collector{jobNameRequiredV2Collector: newJobNameRequiredV2Collector()} + + mgr := newCollectorTestManager() + mgr.SetDyncfgResponder(dyncfg.NewResponder(netdataapi.New(safewriter.New(&buf)))) + mgr.modules.Register("singlev2runner", collectorapi.Creator{ + InstancePolicy: collectorapi.InstancePolicySingle, + CreateV2: func() collectorapi.CollectorV2 { + return probe + }, + }) + + fn := dyncfg.NewFunction(functions.Function{ + UID: "single-v2-runner-test", + ContentType: "application/json", + Payload: mustMarshalCollectorConfigPayload(t, prepareDyncfgCfg("singlev2runner", "payload-name")), + Args: collectorTestArgs(mgr, "singlev2runner", string(dyncfg.CommandTest)), + }) + + mgr.dyncfgCmdTest(fn) + mgr.cmdTestWG.Wait() + + var resp map[string]any + mustDecodeFunctionPayload(t, buf.String(), "single-v2-runner-test", &resp) + assert.Equal(t, float64(200), resp["status"]) + assert.Equal(t, "singlev2runner", probe.jobName) + require.Never(t, probe.runCalled.Load, 100*time.Millisecond, 10*time.Millisecond, "runner started during single-instance dyncfg test") +} + func TestDyncfgCmdTest_SingleInstancePolicyUsesCanonicalName(t *testing.T) { var buf bytes.Buffer mod := &namedTestV1Module{} diff --git a/src/go/plugin/agent/jobmgr/funcctl/controller.go b/src/go/plugin/agent/jobmgr/funcctl/controller.go index bb7f73eb1c14e6..477a4612a9bbcb 100644 --- a/src/go/plugin/agent/jobmgr/funcctl/controller.go +++ b/src/go/plugin/agent/jobmgr/funcctl/controller.go @@ -86,6 +86,7 @@ func (c *Controller) RegisterModules(modules collectorapi.Registry) { continue } c.registry.registerModuleWithMethods(name, creator, methods) + c.registerAvailableModuleMethods(name, methods, true) } } @@ -167,7 +168,14 @@ func (c *Controller) registerModuleMethodsOnJobStart(moduleName string) { return } - for _, method := range creator.Methods() { + c.registerAvailableModuleMethods(moduleName, creator.Methods(), false) +} + +func (c *Controller) registerAvailableModuleMethods(moduleName string, methods []funcapi.MethodConfig, agentWideOnly bool) { + for _, method := range methods { + if agentWideOnly && !method.AgentWide { + continue + } if !methodAvailable(method) { continue } diff --git a/src/go/plugin/agent/jobmgr/funcctl/controller_test.go b/src/go/plugin/agent/jobmgr/funcctl/controller_test.go index 26847b4d8f63b8..3b7c72084800dc 100644 --- a/src/go/plugin/agent/jobmgr/funcctl/controller_test.go +++ b/src/go/plugin/agent/jobmgr/funcctl/controller_test.go @@ -411,17 +411,19 @@ func TestParseArgsParams(t *testing.T) { func TestControllerLifecycleHooks(t *testing.T) { tests := map[string]struct{}{ - "register modules does not register static methods yet": {}, - "first job start registers static methods once": {}, - "availability-gated static method registers when available": {}, - "public method name collision skips colliding module": {}, - "rejected module does not poison planned public names": {}, - "topology methods register direct alias": {}, - "job stop unregisters job methods": {}, - "cleanup unregisters static methods": {}, - "cleanup ignores unavailable static methods": {}, - "cleanup with api configured still unregisters static methods": {}, - "api registration honors method tags": {}, + "register modules does not register static methods yet": {}, + "register modules registers available agent-wide methods": {}, + "first job start registers static methods once": {}, + "availability-gated static method registers when available": {}, + "availability-gated agent-wide method registers when available": {}, + "public method name collision skips colliding module": {}, + "rejected module does not poison planned public names": {}, + "topology methods register direct alias": {}, + "job stop unregisters job methods": {}, + "cleanup unregisters static methods": {}, + "cleanup ignores unavailable static methods": {}, + "cleanup with api configured still unregisters static methods": {}, + "api registration honors method tags": {}, } for name := range tests { @@ -439,6 +441,17 @@ func TestControllerLifecycleHooks(t *testing.T) { assert.Empty(t, reg.registeredNames()) + case "register modules registers available agent-wide methods": + controller.RegisterModules(collectorapi.Registry{ + "mod": collectorapi.Creator{ + Methods: func() []funcapi.MethodConfig { + return []funcapi.MethodConfig{{ID: "logs", AgentWide: true}} + }, + }, + }) + + assert.Equal(t, []string{"mod:logs"}, reg.registeredNames()) + case "first job start registers static methods once": controller.RegisterModules(collectorapi.Registry{ "mod": collectorapi.Creator{ @@ -473,6 +486,28 @@ func TestControllerLifecycleHooks(t *testing.T) { assert.Equal(t, []string{"mod:logs"}, reg.registeredNames()) + case "availability-gated agent-wide method registers when available": + available := false + controller.RegisterModules(collectorapi.Registry{ + "mod": collectorapi.Creator{ + Methods: func() []funcapi.MethodConfig { + return []funcapi.MethodConfig{{ + ID: "logs", + AgentWide: true, + Available: func() bool { return available }, + }} + }, + }, + }) + + assert.Empty(t, reg.registeredNames()) + + available = true + controller.OnJobStart(newTestRuntimeJob("mod", "job1", true)) + controller.OnJobStart(newTestRuntimeJob("mod", "job2", true)) + + assert.Equal(t, []string{"mod:logs"}, reg.registeredNames()) + case "public method name collision skips colliding module": controller.RegisterModules(collectorapi.Registry{ "aaa": collectorapi.Creator{ @@ -520,15 +555,18 @@ func TestControllerLifecycleHooks(t *testing.T) { case "topology methods register direct alias": controller.RegisterModules(collectorapi.Registry{ - "snmp": collectorapi.Creator{ + "snmp_topology": collectorapi.Creator{ Methods: func() []funcapi.MethodConfig { - return []funcapi.MethodConfig{{ID: "topology:snmp", Aliases: []string{"topology:snmp"}}} + return []funcapi.MethodConfig{{ + ID: "topology:snmp", + FunctionName: "snmp:topology:snmp", + Aliases: []string{"topology:snmp"}, + AgentWide: true, + }} }, }, }) - controller.OnJobStart(newTestRuntimeJob("snmp", "edge-router", true)) - assert.ElementsMatch(t, []string{"snmp:topology:snmp", "topology:snmp"}, reg.registeredNames()) controller.Cleanup() @@ -853,6 +891,53 @@ func TestControllerRawModuleMethodRequest(t *testing.T) { } } +func TestControllerRawAgentWideModuleMethodDoesNotRequireRunningJob(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{ + 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", + }) + }, + } + }, + }, + }) + + reg.call("mod:logs", context.Background(), functions.Function{ + UID: "raw-agent-wide", + Timeout: time.Second, + }) + + assert.Equal(t, 200, gotCode) + assert.Equal(t, float64(200), gotResp["status"]) + assert.Nil(t, gotJob) +} + func TestControllerModuleMethodRequestContextCancellation(t *testing.T) { ctx, cancel := context.WithCancel(context.Background()) cancel() diff --git a/src/go/plugin/agent/jobmgr/funcctl/dispatch.go b/src/go/plugin/agent/jobmgr/funcctl/dispatch.go index 0858c166270dbe..a38e5d5f14577e 100644 --- a/src/go/plugin/agent/jobmgr/funcctl/dispatch.go +++ b/src/go/plugin/agent/jobmgr/funcctl/dispatch.go @@ -18,7 +18,7 @@ const ( paramJob = "__job" ) -type methodParamResolver func(ctx context.Context, methodCfg *funcapi.MethodConfig, handler funcapi.MethodHandler, methodID string) ([]funcapi.ParamConfig, bool, error) +type methodParamResolver func(ctx context.Context, methodCfg *funcapi.MethodConfig, handler funcapi.MethodHandler, methodID string) ([]funcapi.ParamConfig, error) type methodExecutionInput struct { fn functions.Function @@ -59,7 +59,7 @@ func (c *Controller) executeMethodRequest(parent context.Context, in methodExecu ctx, cancel := context.WithTimeout(parent, in.fn.Timeout) defer cancel() - if !in.job.IsRunning() { + if in.job != nil && !in.job.IsRunning() { c.respondError(in.fn, 503, "job '%s' is no longer running", in.jobLabel) return } @@ -72,6 +72,10 @@ func (c *Controller) executeMethodRequest(parent context.Context, in methodExecu handler := creator.MethodHandler(in.job) if handler == nil { + if in.job == nil { + c.respondError(in.fn, 500, "module '%s' returned nil handler for agent-wide method '%s'", in.moduleName, in.methodID) + return + } c.respondError(in.fn, 500, "module '%s' returned nil handler for job '%s'", in.moduleName, in.jobName) return } @@ -96,7 +100,7 @@ func (c *Controller) executeMethodRequest(parent context.Context, in methodExecu Source: in.fn.Source, }) - if !c.registry.verifyJobGeneration(in.moduleName, in.jobName, in.jobGen) { + if in.job != nil && !c.registry.verifyJobGeneration(in.moduleName, in.jobName, in.jobGen) { c.respondError(in.fn, 503, "job '%s' was replaced during request, please retry", in.jobLabel) return } @@ -106,17 +110,19 @@ func (c *Controller) executeMethodRequest(parent context.Context, in methodExecu return } - methodParams, paramsFromJob, err := in.resolveParams(ctx, in.methodCfg, handler, in.methodID) + methodParams, err := in.resolveParams(ctx, in.methodCfg, handler, in.methodID) if err != nil { + if in.job == nil { + c.respondError(in.fn, 503, "module '%s' method '%s' cannot provide parameters: %v", in.moduleName, in.methodID, err) + return + } c.respondError(in.fn, 503, "job '%s' cannot provide parameters: %v", in.jobLabel, err) return } - if paramsFromJob { - if err := validateParamValues(methodParams, in.argValues, in.payload, in.jobName); err != nil { - c.respondError(in.fn, 400, "%v", err) - return - } + if err := validateParamValues(methodParams, in.argValues, in.payload, in.jobName); err != nil { + c.respondError(in.fn, 400, "%v", err) + return } methodParamValues := make(map[string][]string, len(methodParams)) @@ -130,7 +136,7 @@ func (c *Controller) executeMethodRequest(parent context.Context, in methodExecu dataResp := handler.Handle(ctx, in.methodID, resolvedParams) - if !c.registry.verifyJobGeneration(in.moduleName, in.jobName, in.jobGen) { + if in.job != nil && !c.registry.verifyJobGeneration(in.moduleName, in.jobName, in.jobGen) { c.respondError(in.fn, 503, "job '%s' was replaced during request, please retry", in.jobLabel) return } @@ -157,34 +163,49 @@ func (c *Controller) makeMethodFuncHandler(moduleName, methodID string) function argValues := parseArgsParams(fn.Args) includeJobParam := methodRequiresJobParam(methodCfg) + if !includeJobParam { + c.executeMethodRequest(ctx, methodExecutionInput{ + fn: fn, + moduleName: moduleName, + jobLabel: moduleName, + methodID: methodID, + methodCfg: methodCfg, + info: info, + payload: payload, + argValues: argValues, + resolveParams: func(ctx context.Context, methodCfg *funcapi.MethodConfig, handler funcapi.MethodHandler, methodID string) ([]funcapi.ParamConfig, error) { + return c.resolveModuleMethodParams(ctx, methodCfg, handler, methodID) + }, + respond: func(dataResp *funcapi.FunctionResponse, methodParams []funcapi.ParamConfig, updateEvery int) { + c.respondWithParams(fn, moduleName, dataResp, methodParams, updateEvery, methodCfg.ResponseType, false) + }, + }) + return + } + jobs := c.registry.getJobNames(moduleName) if len(jobs) == 0 { c.respondError(fn, 422, "no %s instances configured", moduleName) return } - // FIXME: AgentWide currently means "omit __job from the public API" rather - // than "dispatch without a job"; we still route through the first running - // job for the module. jobName := jobs[0] var resolvedJob funcapi.ResolvedParam - if includeJobParam { - jobParam := buildJobParamConfig(jobs) - jobValues := paramValues(argValues, payload, paramJob) - if len(jobValues) > 1 { - c.respondError(fn, 400, "parameter '%s' expects a single value", paramJob) - return - } - resolvedJob = funcapi.ResolveParam(jobParam, jobValues) - jobName = resolvedJob.GetOne() - if len(jobValues) > 0 && jobValues[0] != jobName { - c.respondError(fn, 404, "unknown job '%s', available: %v", jobValues[0], jobs) - return - } - if jobName == "" { - c.respondError(fn, 404, "no %s instances configured", moduleName) - return - } + jobParam := buildJobParamConfig(jobs) + jobValues := paramValues(argValues, payload, paramJob) + if len(jobValues) > 1 { + c.respondError(fn, 400, "parameter '%s' expects a single value", paramJob) + return + } + resolvedJob = funcapi.ResolveParam(jobParam, jobValues) + jobName = resolvedJob.GetOne() + if len(jobValues) > 0 && jobValues[0] != jobName { + c.respondError(fn, 404, "unknown job '%s', available: %v", jobValues[0], jobs) + return + } + if jobName == "" { + c.respondError(fn, 404, "no %s instances configured", moduleName) + return } job, jobGen := c.registry.getJobWithGeneration(moduleName, jobName) @@ -205,8 +226,8 @@ func (c *Controller) makeMethodFuncHandler(moduleName, methodID string) function info: info, payload: payload, argValues: argValues, - resolveParams: func(ctx context.Context, methodCfg *funcapi.MethodConfig, handler funcapi.MethodHandler, methodID string) ([]funcapi.ParamConfig, bool, error) { - return c.resolveMethodParamsForJob(ctx, methodCfg, job, handler, methodID) + resolveParams: func(ctx context.Context, methodCfg *funcapi.MethodConfig, handler funcapi.MethodHandler, methodID string) ([]funcapi.ParamConfig, error) { + return c.resolveModuleMethodParams(ctx, methodCfg, handler, methodID) }, augmentParams: func(resolvedParams funcapi.ResolvedParams) { resolvedParams[paramJob] = resolvedJob @@ -267,18 +288,18 @@ func (c *Controller) buildRequiredParams(moduleName string, methodParams []funca return required } -func (c *Controller) resolveMethodParamsForJob(ctx context.Context, methodCfg *funcapi.MethodConfig, job collectorapi.RuntimeJob, handler funcapi.MethodHandler, methodID string) ([]funcapi.ParamConfig, bool, error) { +func (c *Controller) resolveModuleMethodParams(ctx context.Context, methodCfg *funcapi.MethodConfig, handler funcapi.MethodHandler, methodID string) ([]funcapi.ParamConfig, error) { methodParams := methodCfg.RequiredParams - jobParams, err := handler.MethodParams(ctx, methodID) + handlerParams, err := handler.MethodParams(ctx, methodID) if err != nil { - return nil, false, err + return nil, err } - if len(jobParams) == 0 { - return methodParams, true, nil + if len(handlerParams) == 0 { + return methodParams, nil } - return funcapi.MergeParamConfigs(methodParams, jobParams), true, nil + return funcapi.MergeParamConfigs(methodParams, handlerParams), nil } func validateParamValues(methodParams []funcapi.ParamConfig, argValues map[string][]string, payload map[string]any, jobName string) error { @@ -288,11 +309,17 @@ func validateParamValues(methodParams []funcapi.ParamConfig, argValues map[strin continue } if cfg.Selection == funcapi.ParamSelect && len(values) > 1 { + if jobName == "" { + return fmt.Errorf("parameter '%s' expects a single value", cfg.ID) + } return fmt.Errorf("parameter '%s' expects a single value for job '%s'", cfg.ID, jobName) } allowed := allowedOptions(cfg.Options) for _, value := range values { if !allowed[value] { + if jobName == "" { + return fmt.Errorf("parameter '%s' option '%s' is not supported", cfg.ID, value) + } return fmt.Errorf("parameter '%s' option '%s' is not supported by job '%s'", cfg.ID, value, jobName) } } @@ -496,7 +523,7 @@ func (c *Controller) makeJobMethodFuncHandler(moduleName, jobName, methodID stri info: info, payload: payload, argValues: argValues, - resolveParams: func(ctx context.Context, methodCfg *funcapi.MethodConfig, handler funcapi.MethodHandler, methodID string) ([]funcapi.ParamConfig, bool, error) { + resolveParams: func(ctx context.Context, methodCfg *funcapi.MethodConfig, handler funcapi.MethodHandler, methodID string) ([]funcapi.ParamConfig, error) { return c.resolveJobMethodParams(ctx, methodCfg, handler, methodID) }, respond: func(dataResp *funcapi.FunctionResponse, methodParams []funcapi.ParamConfig, updateEvery int) { @@ -539,18 +566,18 @@ func (c *Controller) handleJobMethodFuncInfo(moduleName, jobName, methodID strin c.respondJSON(fn, resp) } -func (c *Controller) resolveJobMethodParams(ctx context.Context, methodCfg *funcapi.MethodConfig, handler funcapi.MethodHandler, methodID string) ([]funcapi.ParamConfig, bool, error) { +func (c *Controller) resolveJobMethodParams(ctx context.Context, methodCfg *funcapi.MethodConfig, handler funcapi.MethodHandler, methodID string) ([]funcapi.ParamConfig, error) { methodParams := methodCfg.RequiredParams jobParams, err := handler.MethodParams(ctx, methodID) if err != nil { - return nil, false, err + return nil, err } if len(jobParams) == 0 { - return methodParams, true, nil + return methodParams, nil } - return funcapi.MergeParamConfigs(methodParams, jobParams), true, nil + return funcapi.MergeParamConfigs(methodParams, jobParams), nil } func buildJobMethodAcceptedParams(methodParams []funcapi.ParamConfig) []string { accepted := make([]string, 0, len(methodParams)) diff --git a/src/go/plugin/agent/jobmgr/funcdispatch_test.go b/src/go/plugin/agent/jobmgr/funcdispatch_test.go index 9c48449fe5b060..d57947ccbd4285 100644 --- a/src/go/plugin/agent/jobmgr/funcdispatch_test.go +++ b/src/go/plugin/agent/jobmgr/funcdispatch_test.go @@ -203,6 +203,116 @@ func TestExecuteFunction_ModuleMethodPublicFunctionName(t *testing.T) { assert.Equal(t, "logs", gotMethod) } +func TestExecuteFunction_AgentWideModuleMethodDoesNotRequireRunningJob(t *testing.T) { + tests := map[string]struct { + functionName string + }{ + "public function name": { + functionName: "snmp:topology:snmp", + }, + "alias": { + functionName: "topology:snmp", + }, + } + + for name, tc := range tests { + t.Run(name, func(t *testing.T) { + writer := &jsonWriteCapture{} + var gotMethod string + var gotJob collectorapi.RuntimeJob + + mgr := New(Config{ + PluginName: testPluginName, + FunctionJSONWriter: writer.write, + }) + mgr.modules = collectorapi.Registry{ + "snmp_topology": collectorapi.Creator{ + Methods: func() []funcapi.MethodConfig { + return []funcapi.MethodConfig{{ + ID: "topology:snmp", + FunctionName: "snmp:topology:snmp", + Aliases: []string{"topology:snmp"}, + AgentWide: true, + RequiredParams: []funcapi.ParamConfig{{ + ID: "scope", + Name: "Scope", + Selection: funcapi.ParamSelect, + Options: []funcapi.ParamOption{{ID: "all", Name: "All"}}, + }}, + }} + }, + MethodHandler: func(job collectorapi.RuntimeJob) funcapi.MethodHandler { + gotJob = job + return &mockMethodHandler{ + handleFunc: func(ctx context.Context, method string, params funcapi.ResolvedParams) *funcapi.FunctionResponse { + gotMethod = method + assert.Equal(t, "all", params.GetOne("scope")) + assert.Empty(t, params.GetOne("__job")) + return &funcapi.FunctionResponse{Status: 200, Help: "topology"} + }, + } + }, + }, + } + mgr.funcCtl.RegisterModules(mgr.modules) + + mgr.ExecuteFunction(tc.functionName, functions.Function{ + UID: "agent-wide-public-name", + Timeout: time.Second, + Args: []string{"scope:all"}, + }) + + resp := writer.requireResponse(t) + assert.Equal(t, float64(200), resp["status"]) + assert.Equal(t, "topology", resp["help"]) + assert.Equal(t, "topology:snmp", gotMethod) + assert.Nil(t, gotJob) + }) + } +} + +func TestExecuteFunction_AgentWideModuleMethodValidationErrorOmitsJob(t *testing.T) { + writer := &jsonWriteCapture{} + + mgr := New(Config{ + PluginName: testPluginName, + FunctionJSONWriter: writer.write, + }) + mgr.modules = collectorapi.Registry{ + "mod": collectorapi.Creator{ + Methods: func() []funcapi.MethodConfig { + return []funcapi.MethodConfig{{ + ID: "logs", + AgentWide: true, + RequiredParams: []funcapi.ParamConfig{{ + ID: "scope", + Name: "Scope", + Selection: funcapi.ParamSelect, + Options: []funcapi.ParamOption{ + {ID: "a", Name: "A"}, + {ID: "b", Name: "B"}, + }, + }}, + }} + }, + MethodHandler: func(collectorapi.RuntimeJob) funcapi.MethodHandler { + return &mockMethodHandler{} + }, + }, + } + mgr.funcCtl.RegisterModules(mgr.modules) + + mgr.ExecuteFunction("mod:logs", functions.Function{ + UID: "agent-wide-validation", + Timeout: time.Second, + Args: []string{"scope:a,b"}, + }) + + resp := writer.requireResponse(t) + assert.Equal(t, float64(400), resp["status"]) + assert.Contains(t, resp["errorMessage"], "parameter 'scope' expects a single value") +} + func TestExecuteFunction_ContextBehavior(t *testing.T) { tests := map[string]struct { managerCtx context.Context diff --git a/src/go/plugin/agent/jobmgr/manager_process_test.go b/src/go/plugin/agent/jobmgr/manager_process_test.go index 118fab45b3a679..159391c03f9e15 100644 --- a/src/go/plugin/agent/jobmgr/manager_process_test.go +++ b/src/go/plugin/agent/jobmgr/manager_process_test.go @@ -203,7 +203,7 @@ func TestManagerAddConfigSingleInstancePolicyPublishesDyncfgSingle(t *testing.T) assert.NotContains(t, out, " remove") } -func TestRun_DoesNotRegisterModuleMethodsBeforeAnyJobStarts(t *testing.T) { +func TestRun_DoesNotRegisterJobBoundModuleMethodsBeforeAnyJobStarts(t *testing.T) { fnReg := &recordingFunctionRegistry{} mgr := New(Config{PluginName: testPluginName, FnReg: fnReg}) @@ -241,6 +241,44 @@ func TestRun_DoesNotRegisterModuleMethodsBeforeAnyJobStarts(t *testing.T) { assert.Empty(t, fnReg.registeredNames(), "static methods must not be registered before first started job") } +func TestRun_RegistersAgentWideModuleMethodsBeforeAnyJobStarts(t *testing.T) { + fnReg := &recordingFunctionRegistry{} + mgr := New(Config{PluginName: testPluginName, FnReg: fnReg}) + + mgr.modules = collectorapi.Registry{ + "mod": collectorapi.Creator{ + Methods: func() []funcapi.MethodConfig { + return []funcapi.MethodConfig{{ID: "a", AgentWide: true}} + }, + }, + } + + ctx, cancel := context.WithCancel(context.Background()) + defer cancel() + + in := make(chan []*confgroup.Group) + done := make(chan struct{}) + go func() { + mgr.Run(ctx, in) + close(done) + }() + + waitCtx, waitCancel := context.WithTimeout(context.Background(), time.Second) + defer waitCancel() + require.True(t, mgr.WaitStarted(waitCtx), "manager did not report started") + + cancel() + close(in) + + select { + case <-done: + case <-time.After(2 * time.Second): + t.Fatal("manager did not stop after cancel") + } + + assert.Equal(t, []string{"mod:a"}, fnReg.registeredNames()) +} + func TestStartRunningJob_RegistersModuleMethodsOnFirstStartedJob(t *testing.T) { fnReg := &recordingFunctionRegistry{} mgr := New(Config{PluginName: testPluginName, FnReg: fnReg}) diff --git a/src/go/plugin/framework/collectorapi/collector.go b/src/go/plugin/framework/collectorapi/collector.go index 0b308b59513457..062284f82b8505 100644 --- a/src/go/plugin/framework/collectorapi/collector.go +++ b/src/go/plugin/framework/collectorapi/collector.go @@ -56,6 +56,17 @@ type CollectorV2 interface { ChartTemplateYAML() string } +// CollectorV2Runner is an optional long-running V2 collector hook. +// +// Run is called only after the runtime job starts. It is not called during +// config validation or autodetection. Implementations should block until ctx is +// canceled, and must return promptly after ctx.Done(). Returning before +// cancellation is treated as an unexpected runner stop and is logged by the job +// runtime. +type CollectorV2Runner interface { + Run(context.Context) error +} + // CollectorV2EnginePolicy allows a V2 collector to provide chartengine policy // (series selector + autogen behavior). type CollectorV2EnginePolicy interface { diff --git a/src/go/plugin/framework/collectorapi/registry.go b/src/go/plugin/framework/collectorapi/registry.go index e8bb375e9a67ea..7268c4ff8c77da 100644 --- a/src/go/plugin/framework/collectorapi/registry.go +++ b/src/go/plugin/framework/collectorapi/registry.go @@ -66,7 +66,9 @@ type ( // If Methods is non-nil, this module provides functions Methods func() []funcapi.MethodConfig - // Optional: MethodHandler returns a handler for method requests on a specific job. + // 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. // 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/framework/jobruntime/job_common.go b/src/go/plugin/framework/jobruntime/job_common.go index 09301d2522e082..db423f1ce274d9 100644 --- a/src/go/plugin/framework/jobruntime/job_common.go +++ b/src/go/plugin/framework/jobruntime/job_common.go @@ -49,6 +49,15 @@ func (c *stopController) requestStop() { c.stopOnce.Do(func() { close(c.stopCh) }) } +func (c *stopController) stopRequested() bool { + select { + case <-c.stopCh: + return true + default: + return false + } +} + func (c *stopController) stopAndWait() { c.requestStop() if !c.started.Load() { diff --git a/src/go/plugin/framework/jobruntime/job_v2.go b/src/go/plugin/framework/jobruntime/job_v2.go index 54424a64609901..b17a7f532400e5 100644 --- a/src/go/plugin/framework/jobruntime/job_v2.go +++ b/src/go/plugin/framework/jobruntime/job_v2.go @@ -295,6 +295,10 @@ func (j *JobV2) Start() { j.running.Store(true) runCtx, cancel := context.WithCancel(context.Background()) j.setRunContext(runCtx, cancel) + var runnerDone <-chan error + if !j.stopCtrl.stopRequested() { + runnerDone = j.startCollectorRunner(runCtx) + } if j.functionOnly { j.Info("started in function-only mode") } else { @@ -312,6 +316,9 @@ LOOP: select { case <-j.stopCtrl.stopCh: break LOOP + case err := <-runnerDone: + runnerDone = nil + j.handleCollectorRunnerExit(runCtx, err) case t := <-j.tick: if !j.functionOnly && j.shouldCollect(t) { markRunStartWithResumeLog(&j.skipTracker, j.Logger) @@ -320,12 +327,62 @@ LOOP: } } } + cancel() + j.waitCollectorRunner(runCtx, runnerDone) // Mark not-running before cleanup so external function dispatch can reject requests // while module resources are being torn down. j.running.Store(false) j.Cleanup() } +func (j *JobV2) startCollectorRunner(ctx context.Context) <-chan error { + runner, ok := j.module.(collectorapi.CollectorV2Runner) + if !ok { + return nil + } + done := make(chan error, 1) + go func() { + done <- j.runCollectorRunner(ctx, runner) + }() + return done +} + +func (j *JobV2) runCollectorRunner(ctx context.Context, runner collectorapi.CollectorV2Runner) (err error) { + defer func() { + if r := recover(); r != nil { + j.panicked.Store(true) + err = fmt.Errorf("panic %v", r) + j.Errorf("PANIC: %v", r) + if logger.Level.Enabled(slog.LevelDebug) { + j.Errorf("STACK: %s", debug.Stack()) + } + } + }() + + err = runner.Run(ctx) + if err != nil && ctx.Err() != nil && errors.Is(err, ctx.Err()) { + return nil + } + return err +} + +func (j *JobV2) waitCollectorRunner(ctx context.Context, done <-chan error) { + if done == nil { + return + } + j.handleCollectorRunnerExit(ctx, <-done) +} + +func (j *JobV2) handleCollectorRunnerExit(ctx context.Context, err error) { + if err != nil { + j.Errorf("collector runner failed: %v", err) + return + } + if ctx.Err() == nil { + j.Warningf("collector runner stopped before job stop") + } +} + func (j *JobV2) Stop() { j.cancelRunContext() j.stopCtrl.stopAndWait() diff --git a/src/go/plugin/framework/jobruntime/job_v2_test.go b/src/go/plugin/framework/jobruntime/job_v2_test.go index 65510d4242ee3e..9a09aa52c5d467 100644 --- a/src/go/plugin/framework/jobruntime/job_v2_test.go +++ b/src/go/plugin/framework/jobruntime/job_v2_test.go @@ -39,6 +39,11 @@ type mockModuleV2 struct { vnode *vnodes.VirtualNode } +type mockRunnerModuleV2 struct { + *mockModuleV2 + runFunc func(context.Context) error +} + type mockRuntimeComponentService struct { registerErr error registered []runtimecomp.ComponentConfig @@ -105,6 +110,13 @@ func (m *mockModuleV2) ChartTemplateYAML() string { return m.template } +func (m *mockRunnerModuleV2) Run(ctx context.Context) error { + if m.runFunc == nil { + return nil + } + return m.runFunc(ctx) +} + func newTestJobV2(mod collectorapi.CollectorV2, out *bytes.Buffer) *JobV2 { return NewJobV2(JobV2Config{ PluginName: pluginName, @@ -230,6 +242,175 @@ groups: ` } +func TestJobV2RunnerDoesNotRunDuringAutoDetection(t *testing.T) { + started := make(chan struct{}) + mod := &mockRunnerModuleV2{ + mockModuleV2: &mockModuleV2{ + store: metrix.NewCollectorStore(), + template: chartTemplateV2(), + }, + runFunc: func(context.Context) error { + close(started) + return nil + }, + } + job := newTestJobV2(mod, &bytes.Buffer{}) + + require.NoError(t, job.AutoDetection()) + + require.Never(t, func() bool { + select { + case <-started: + return true + default: + return false + } + }, 100*time.Millisecond, 10*time.Millisecond, "runner started during autodetection") +} + +func TestJobV2RunnerDoesNotStartAfterPreStartStop(t *testing.T) { + started := make(chan struct{}) + stopped := make(chan struct{}) + mod := &mockRunnerModuleV2{ + mockModuleV2: &mockModuleV2{ + store: metrix.NewCollectorStore(), + template: chartTemplateV2(), + }, + runFunc: func(context.Context) error { + close(started) + return nil + }, + } + job := newTestJobV2(mod, &bytes.Buffer{}) + require.NoError(t, job.AutoDetection()) + + job.Stop() + go func() { + job.Start() + close(stopped) + }() + + select { + case <-stopped: + case <-time.After(time.Second): + t.Fatal("pre-stopped job did not exit") + } + select { + case <-started: + t.Fatal("runner started after pre-start stop") + default: + } +} + +func TestJobV2RunnerStopWaitsBeforeCleanup(t *testing.T) { + started := make(chan struct{}) + ctxCanceled := make(chan struct{}) + releaseRunner := make(chan struct{}) + cleanupCalled := make(chan struct{}) + stopped := make(chan struct{}) + + mod := &mockRunnerModuleV2{ + mockModuleV2: &mockModuleV2{ + store: metrix.NewCollectorStore(), + template: chartTemplateV2(), + cleanupFunc: func(context.Context) { + close(cleanupCalled) + }, + }, + runFunc: func(ctx context.Context) error { + close(started) + <-ctx.Done() + close(ctxCanceled) + <-releaseRunner + return nil + }, + } + job := newTestJobV2(mod, &bytes.Buffer{}) + require.NoError(t, job.AutoDetection()) + + go job.Start() + select { + case <-started: + case <-time.After(time.Second): + t.Fatal("runner did not start") + } + + go func() { + job.Stop() + close(stopped) + }() + + select { + case <-ctxCanceled: + case <-time.After(time.Second): + t.Fatal("runner context was not canceled") + } + + select { + case <-cleanupCalled: + t.Fatal("cleanup ran before runner returned") + default: + } + select { + case <-stopped: + t.Fatal("Stop returned before runner returned") + default: + } + + close(releaseRunner) + require.Eventually(t, func() bool { + select { + case <-cleanupCalled: + return true + default: + return false + } + }, time.Second, 10*time.Millisecond) + require.Eventually(t, func() bool { + select { + case <-stopped: + return true + default: + return false + } + }, time.Second, 10*time.Millisecond) +} + +func TestJobV2RunnerPanicRecovered(t *testing.T) { + started := make(chan struct{}) + stopped := make(chan struct{}) + mod := &mockRunnerModuleV2{ + mockModuleV2: &mockModuleV2{ + store: metrix.NewCollectorStore(), + template: chartTemplateV2(), + }, + runFunc: func(context.Context) error { + close(started) + panic("runner boom") + }, + } + job := newTestJobV2(mod, &bytes.Buffer{}) + require.NoError(t, job.AutoDetection()) + + go func() { + job.Start() + close(stopped) + }() + select { + case <-started: + case <-time.After(time.Second): + t.Fatal("runner did not start") + } + require.Eventually(t, job.Panicked, time.Second, 10*time.Millisecond) + + job.Stop() + select { + case <-stopped: + case <-time.After(time.Second): + t.Fatal("job did not stop") + } +} + func TestJobV2Scenarios(t *testing.T) { tests := map[string]struct { run func(t *testing.T) diff --git a/src/go/plugin/go.d/collector/snmp/ddsnmp/topology_provider.go b/src/go/plugin/go.d/collector/snmp/ddsnmp/topology_provider.go deleted file mode 100644 index d920e6211859bf..00000000000000 --- a/src/go/plugin/go.d/collector/snmp/ddsnmp/topology_provider.go +++ /dev/null @@ -1,15 +0,0 @@ -// SPDX-License-Identifier: GPL-3.0-or-later - -package ddsnmp - -import "github.com/netdata/netdata/go/plugins/pkg/funcapi" - -// TopologyHandler is set by the snmp_topology module at init time. -// The snmp module's function router delegates topology:snmp requests to it. -// This avoids circular imports: snmp -> ddsnmp <- snmp_topology. -var TopologyHandler funcapi.MethodHandler - -// TopologyMethodConfig is set by the snmp_topology module at init time. -// The snmp module includes it in its method list so the function appears -// as snmp:topology:snmp. -var TopologyMethodConfig *funcapi.MethodConfig diff --git a/src/go/plugin/go.d/collector/snmp/func_interfaces_test.go b/src/go/plugin/go.d/collector/snmp/func_interfaces_test.go index b45fb69704b674..c875dfa7884d38 100644 --- a/src/go/plugin/go.d/collector/snmp/func_interfaces_test.go +++ b/src/go/plugin/go.d/collector/snmp/func_interfaces_test.go @@ -23,7 +23,9 @@ func TestSnmpMethods(t *testing.T) { var ifacesMethod *funcapi.MethodConfig var licensesMethod *funcapi.MethodConfig var bgpMethod *funcapi.MethodConfig + methodIDs := make(map[string]bool, len(methods)) for i := range methods { + methodIDs[methods[i].ID] = true switch methods[i].ID { case "interfaces": ifacesMethod = &methods[i] @@ -44,6 +46,7 @@ func TestSnmpMethods(t *testing.T) { require.NotNil(t, bgpMethod) assert.Equal(t, "BGP Peers", bgpMethod.Name) require.NotEmpty(t, bgpMethod.RequiredParams) + assert.False(t, methodIDs["topology:snmp"], "snmp_topology owns the topology Function") // Verify type group param exists var typeGroupParam *funcapi.ParamConfig diff --git a/src/go/plugin/go.d/collector/snmp/func_router.go b/src/go/plugin/go.d/collector/snmp/func_router.go index b7d3f233d1c314..a5ec6c73f8c6a4 100644 --- a/src/go/plugin/go.d/collector/snmp/func_router.go +++ b/src/go/plugin/go.d/collector/snmp/func_router.go @@ -40,7 +40,6 @@ func newFuncRouter(ifaceCache *ifaceCache, extraHandlers ...registeredSNMPFuncti r.registerHandler(h.methodID, h.handler) } } - addTopologyFunctionHandler(r.handlers) return r } @@ -100,7 +99,7 @@ func snmpBaseMethods() []funcapi.MethodConfig { ifacesMethodConfig(), } methods = append(methods, collectorSpecificMethodConfigs()...) - return appendTopologyMethodConfig(methods) + return methods } func snmpFunctionHandler(job collectorapi.RuntimeJob) funcapi.MethodHandler { diff --git a/src/go/plugin/go.d/collector/snmp/topology_func_router.go b/src/go/plugin/go.d/collector/snmp/topology_func_router.go deleted file mode 100644 index 36ca5a93df2033..00000000000000 --- a/src/go/plugin/go.d/collector/snmp/topology_func_router.go +++ /dev/null @@ -1,22 +0,0 @@ -// SPDX-License-Identifier: GPL-3.0-or-later - -package snmp - -import ( - "github.com/netdata/netdata/go/plugins/pkg/funcapi" - "github.com/netdata/netdata/go/plugins/plugin/go.d/collector/snmp/ddsnmp" -) - -func addTopologyFunctionHandler(handlers map[string]funcapi.MethodHandler) { - if ddsnmp.TopologyHandler == nil || ddsnmp.TopologyMethodConfig == nil { - return - } - handlers[ddsnmp.TopologyMethodConfig.ID] = ddsnmp.TopologyHandler -} - -func appendTopologyMethodConfig(methods []funcapi.MethodConfig) []funcapi.MethodConfig { - if ddsnmp.TopologyMethodConfig == nil { - return methods - } - return append(methods, *ddsnmp.TopologyMethodConfig) -} diff --git a/src/go/plugin/go.d/collector/snmp_topology/charts.go b/src/go/plugin/go.d/collector/snmp_topology/charts.go index 9cc2062d75d00b..7e686017970247 100644 --- a/src/go/plugin/go.d/collector/snmp_topology/charts.go +++ b/src/go/plugin/go.d/collector/snmp_topology/charts.go @@ -2,44 +2,9 @@ package snmptopology -import "github.com/netdata/netdata/go/plugins/plugin/framework/collectorapi" +import _ "embed" -var ( - topologyDevicesChart = collectorapi.Chart{ - ID: "topology_devices", - Title: "Topology devices", - Units: "devices", - Fam: "Topology", - Ctx: "snmp_topology.devices", - Priority: 39200, - Dims: collectorapi.Dims{ - {ID: "snmp_topology_devices_total", Name: "total"}, - {ID: "snmp_topology_devices_discovered", Name: "discovered"}, - }, - } - topologyLinksChart = collectorapi.Chart{ - ID: "topology_links", - Title: "Topology links", - Units: "links", - Fam: "Topology", - Ctx: "snmp_topology.links", - Priority: 39201, - Dims: collectorapi.Dims{ - {ID: "snmp_topology_links_total", Name: "total"}, - {ID: "snmp_topology_links_lldp", Name: "lldp"}, - {ID: "snmp_topology_links_cdp", Name: "cdp"}, - {ID: "snmp_topology_links_stp", Name: "stp"}, - }, - } - topologyCharts = collectorapi.Charts{ - topologyDevicesChart.Copy(), - topologyLinksChart.Copy(), - } -) +//go:embed charts.yaml +var chartTemplateYAML string -func (c *Collector) addTopologyCharts() { - charts := topologyCharts.Copy() - if err := c.Charts().Add(*charts...); err != nil { - c.Warningf("failed to add topology charts: %v", err) - } -} +func (c *Collector) ChartTemplateYAML() string { return chartTemplateYAML } diff --git a/src/go/plugin/go.d/collector/snmp_topology/charts.yaml b/src/go/plugin/go.d/collector/snmp_topology/charts.yaml new file mode 100644 index 00000000000000..d4c11ce9519629 --- /dev/null +++ b/src/go/plugin/go.d/collector/snmp_topology/charts.yaml @@ -0,0 +1,44 @@ +version: v1 +groups: + - family: Internal + metrics: + - netdata.go.plugin.collector.snmp_topology.devices_registered + - netdata.go.plugin.collector.snmp_topology.devices_cached + - netdata.go.plugin.collector.snmp_topology.last_refresh_age_seconds + - netdata.go.plugin.collector.snmp_topology.last_refresh_duration_seconds + - netdata.go.plugin.collector.snmp_topology.refresh_runs_total + - netdata.go.plugin.collector.snmp_topology.refresh_errors_total + charts: + - id: devices + title: SNMP topology internal devices + context: netdata.go.plugin.collector.snmp_topology.devices + units: devices + algorithm: absolute + type: line + dimensions: + - selector: netdata.go.plugin.collector.snmp_topology.devices_registered + name: registered + - selector: netdata.go.plugin.collector.snmp_topology.devices_cached + name: cached + - id: last_refresh + title: SNMP topology internal last refresh + context: netdata.go.plugin.collector.snmp_topology.last_refresh + units: seconds + algorithm: absolute + type: line + dimensions: + - selector: netdata.go.plugin.collector.snmp_topology.last_refresh_age_seconds + name: age + - selector: netdata.go.plugin.collector.snmp_topology.last_refresh_duration_seconds + name: duration + - id: refreshes + title: SNMP topology internal refreshes + context: netdata.go.plugin.collector.snmp_topology.refreshes + units: events/s + algorithm: incremental + type: line + dimensions: + - selector: netdata.go.plugin.collector.snmp_topology.refresh_runs_total + name: runs + - selector: netdata.go.plugin.collector.snmp_topology.refresh_errors_total + name: errors diff --git a/src/go/plugin/go.d/collector/snmp_topology/charts_test.go b/src/go/plugin/go.d/collector/snmp_topology/charts_test.go new file mode 100644 index 00000000000000..e89373b9293bfe --- /dev/null +++ b/src/go/plugin/go.d/collector/snmp_topology/charts_test.go @@ -0,0 +1,43 @@ +// SPDX-License-Identifier: GPL-3.0-or-later + +package snmptopology + +import ( + "testing" + + "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/go.d/pkg/collecttest" + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" +) + +func TestCollector_ChartTemplateYAML(t *testing.T) { + raw := New().ChartTemplateYAML() + collecttest.AssertChartTemplateSchema(t, raw) + + spec, err := charttpl.DecodeYAML([]byte(raw)) + require.NoError(t, err) + _, err = chartengine.Compile(spec, 1) + require.NoError(t, err) + + contexts := chartTemplateContexts(spec.Groups) + assert.Contains(t, contexts, "netdata.go.plugin.collector.snmp_topology.devices") + assert.Contains(t, contexts, "netdata.go.plugin.collector.snmp_topology.last_refresh") + assert.Contains(t, contexts, "netdata.go.plugin.collector.snmp_topology.refreshes") + assert.NotContains(t, contexts, "snmp_topology.devices") + assert.NotContains(t, contexts, "snmp_topology.links") +} + +func chartTemplateContexts(groups []charttpl.Group) map[string]struct{} { + contexts := make(map[string]struct{}) + for _, group := range groups { + for _, chart := range group.Charts { + contexts[chart.Context] = struct{}{} + } + for ctx := range chartTemplateContexts(group.Groups) { + contexts[ctx] = struct{}{} + } + } + return contexts +} 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 30cbae74fae511..5725401c26f4bd 100644 --- a/src/go/plugin/go.d/collector/snmp_topology/collector.go +++ b/src/go/plugin/go.d/collector/snmp_topology/collector.go @@ -6,13 +6,16 @@ import ( "context" _ "embed" "fmt" + "log/slog" + "runtime/debug" + "sync" "time" "github.com/gosnmp/gosnmp" - topologyengine "github.com/netdata/netdata/go/plugins/pkg/l2topology" - + "github.com/netdata/netdata/go/plugins/logger" "github.com/netdata/netdata/go/plugins/pkg/funcapi" + "github.com/netdata/netdata/go/plugins/pkg/metrix" "github.com/netdata/netdata/go/plugins/plugin/framework/collectorapi" "github.com/netdata/netdata/go/plugins/plugin/go.d/collector/snmp/ddsnmp" "github.com/netdata/netdata/go/plugins/plugin/go.d/collector/snmp/ddsnmp/ddsnmpcollector" @@ -28,25 +31,26 @@ func init() { Defaults: collectorapi.Defaults{ UpdateEvery: 60, }, - Create: func() collectorapi.CollectorV1 { return New() }, - Config: func() any { return &Config{} }, + CreateV2: func() collectorapi.CollectorV2 { return New() }, + Config: func() any { return &Config{} }, + InstancePolicy: collectorapi.InstancePolicySingle, + Methods: topologyMethods, + MethodHandler: topologyFunctionHandler, }) - - // Register the topology function handler and method config so the snmp module - // can serve topology:snmp requests under the snmp:topology:snmp function name. - ddsnmp.TopologyHandler = &funcTopology{} - cfg := topologyMethodConfig() - ddsnmp.TopologyMethodConfig = &cfg } func New() *Collector { + store := metrix.NewCollectorStore() return &Collector{ deviceCaches: make(map[string]*topologyCache), deviceLastCollected: make(map[string]time.Time), + registeredDevices: ddsnmp.DeviceRegistry.Devices, newSnmpClient: gosnmp.NewHandler, newDdSnmpColl: func(cfg ddsnmpcollector.Config) ddCollector { return ddsnmpcollector.New(cfg) }, + store: store, + metrics: newCollectorMetrics(store), } } @@ -55,14 +59,21 @@ type ( collectorapi.Base `yaml:",inline"` Config `yaml:",inline"` - charts *collectorapi.Charts 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) - topologyChartsAdded bool - newSnmpClient func() gosnmp.Handler - newDdSnmpColl func(ddsnmpcollector.Config) ddCollector + refreshMu sync.Mutex + statsMu sync.RWMutex + stats collectorRuntimeStats + + store metrix.CollectorStore + metrics *collectorMetrics + + registeredDevices func() []ddsnmp.DeviceConnectionInfo + topologyProfiles func(ddsnmp.DeviceConnectionInfo) []*ddsnmp.Profile + newSnmpClient func() gosnmp.Handler + newDdSnmpColl func(ddsnmpcollector.Config) ddCollector } ddCollector interface { Collect() ([]*ddsnmp.ProfileMetrics, error) @@ -81,42 +92,43 @@ func (c *Collector) Check(context.Context) error { return nil } -func (c *Collector) Charts() *collectorapi.Charts { - if c.charts == nil { - c.charts = &collectorapi.Charts{} - } - return c.charts +func (c *Collector) Collect(context.Context) error { + c.writeInternalMetrics(time.Now()) + return nil } -func (c *Collector) Collect(context.Context) map[string]int64 { - if devices := ddsnmp.DeviceRegistry.Devices(); len(devices) > 0 { - refreshEvery := c.refreshEvery() - now := time.Now() - seen := make(map[string]bool, len(devices)) +func (c *Collector) MetricStore() metrix.CollectorStore { return c.store } - for _, dev := range devices { - key := fmt.Sprintf("%s:%d", dev.Hostname, dev.Port) - seen[key] = true +func (c *Collector) Run(ctx context.Context) error { + if err := ctx.Err(); err != nil { + return nil + } + c.refreshTopologyRecovering(ctx) - lastCollected, exists := c.deviceLastCollected[key] - isNew := !exists - isStale := exists && now.Sub(lastCollected) >= refreshEvery + ticker := time.NewTicker(c.deviceCheckEvery()) + defer ticker.Stop() - if isNew || isStale { - c.refreshDeviceTopology(key, dev) - c.deviceLastCollected[key] = now - } + for { + select { + case <-ctx.Done(): + return nil + case <-ticker.C: + c.refreshTopologyRecovering(ctx) } - - c.pruneStaleDeviceCaches(seen) } - - mx := make(map[string]int64) - c.collectTopologyMetrics(mx) - return mx } -const defaultRefreshEvery = 30 * time.Minute +const ( + defaultDeviceCheckEvery = time.Minute + defaultRefreshEvery = 30 * time.Minute +) + +func (c *Collector) deviceCheckEvery() time.Duration { + if c.UpdateEvery > 0 { + return time.Duration(c.UpdateEvery) * time.Second + } + return defaultDeviceCheckEvery +} func (c *Collector) refreshEvery() time.Duration { if d := c.RefreshEvery.Duration(); d > 0 { @@ -126,31 +138,125 @@ func (c *Collector) refreshEvery() time.Duration { } func (c *Collector) Cleanup(context.Context) { + c.refreshMu.Lock() + defer c.refreshMu.Unlock() + for key, cache := range c.deviceCaches { snmpTopologyRegistry.unregister(cache) delete(c.deviceCaches, key) + delete(c.deviceLastCollected, key) + } + c.recordCleanupStats() +} + +func (c *Collector) refreshTopology(ctx context.Context) refreshStats { + start := time.Now() + c.refreshMu.Lock() + defer c.refreshMu.Unlock() + + devices := c.getRegisteredDevices() + refreshEvery := c.refreshEvery() + now := time.Now() + seen := make(map[string]bool, len(devices)) + stats := refreshStats{ + hasDeviceCounts: true, + registeredDevices: len(devices), + } + + for _, dev := range devices { + if ctx.Err() != nil { + break + } + + key := fmt.Sprintf("%s:%d", dev.Hostname, dev.Port) + seen[key] = true + + lastCollected, exists := c.deviceLastCollected[key] + isNew := !exists + isStale := exists && now.Sub(lastCollected) >= refreshEvery + + if isNew || isStale { + if !c.refreshDeviceTopology(ctx, key, dev) { + if ctx.Err() != nil { + break + } + stats.errors++ + } + c.deviceLastCollected[key] = now + } } + + if ctx.Err() == nil { + c.pruneStaleDeviceCaches(seen) + } + stats.cachedDevices = len(c.deviceCaches) + stats.completedAt = time.Now() + stats.duration = stats.completedAt.Sub(start) + return stats +} + +func (c *Collector) getRegisteredDevices() []ddsnmp.DeviceConnectionInfo { + if c.registeredDevices == nil { + return nil + } + return c.registeredDevices() +} + +func (c *Collector) refreshTopologyRecovering(ctx context.Context) { + start := time.Now() + defer func() { + if r := recover(); r != nil { + c.recordRefreshStats(refreshStats{ + errors: 1, + completedAt: time.Now(), + duration: time.Since(start), + }) + c.Errorf("PANIC: %v", r) + if logger.Level.Enabled(slog.LevelDebug) { + c.Errorf("STACK: %s", debug.Stack()) + } + } + }() + + c.recordRefreshStats(c.refreshTopology(ctx)) } // refreshDeviceTopology collects topology data for a single device into its own cache. -func (c *Collector) refreshDeviceTopology(key string, dev ddsnmp.DeviceConnectionInfo) { +func (c *Collector) refreshDeviceTopology(ctx context.Context, key string, dev ddsnmp.DeviceConnectionInfo) bool { + if ctx.Err() != nil { + return false + } + snmpClient, err := newSNMPClientFromDeviceInfo(c.newSnmpClient, dev) if err != nil { c.Warningf("device '%s': failed to create SNMP client: %v", dev.Hostname, err) - return + return false } if dev.MaxRepetitions != 0 { snmpClient.SetMaxRepetitions(dev.MaxRepetitions) } if err := snmpClient.Connect(); err != nil { + if ctx.Err() != nil { + return false + } c.Warningf("device '%s': failed to connect: %v", dev.Hostname, err) - return + return false } + stopContextClose := closeSNMPClientOnContextCancel(ctx, snmpClient) + defer stopContextClose() defer func() { _ = snmpClient.Close() }() - profiles := c.findTopologyProfiles(dev) + if ctx.Err() != nil { + return false + } + + profiles := c.getTopologyProfiles(dev) if len(profiles) == 0 { - return + return true + } + + if ctx.Err() != nil { + return false } coll := c.newDdSnmpColl(ddsnmpcollector.Config{ @@ -163,15 +269,26 @@ func (c *Collector) refreshDeviceTopology(key string, dev ddsnmp.DeviceConnectio pms, err := coll.Collect() if err != nil { + if ctx.Err() != nil { + return false + } c.Warningf("device '%s': topology collection failed: %v", dev.Hostname, err) - return + return false + } + + if ctx.Err() != nil { + return false } sysUptime, err := snmputils.GetSysUptime(snmpClient) - if err != nil { + if err != nil && ctx.Err() == nil { c.Debugf("device '%s': failed to query system uptime: %v", dev.Hostname, err) } + if ctx.Err() != nil { + return false + } + // Build the next snapshot off-registry. Function readers keep seeing the // previous complete snapshot until this collection is fully ingested. next := c.newDeviceCollectionCache(dev) @@ -181,13 +298,29 @@ func (c *Collector) refreshDeviceTopology(key string, dev ddsnmp.DeviceConnectio c.updateTopologySysUptime(sysUptime) c.updateTopologyProfileTags(pms) c.ingestTopologyProfileMetrics(pms) - c.collectTopologyVTPVLANContexts(dev) + c.collectTopologyVTPVLANContexts(ctx, dev) + if ctx.Err() != nil { + return false + } c.finalizeTopologyCache() cache := c.getOrCreateDeviceCache(key) cache.mu.Lock() cache.replaceWith(next) cache.mu.Unlock() + return true +} + +func closeSNMPClientOnContextCancel(ctx context.Context, client gosnmp.Handler) func() { + done := make(chan struct{}) + go func() { + select { + case <-ctx.Done(): + _ = client.Close() + case <-done: + } + }() + return func() { close(done) } } func (c *Collector) getOrCreateDeviceCache(key string) *topologyCache { @@ -203,7 +336,7 @@ func (c *Collector) getOrCreateDeviceCache(key string) *topologyCache { func (c *Collector) newDeviceCollectionCache(dev ddsnmp.DeviceConnectionInfo) *topologyCache { cache := newTopologyCache() cache.updateTime = time.Now() - cache.staleAfter = c.refreshEvery() + time.Duration(c.UpdateEvery*2)*time.Second + cache.staleAfter = c.refreshEvery() + 2*c.deviceCheckEvery() cache.agentID = dev.Hostname cache.localDevice = buildLocalTopologyDevice(dev) return cache @@ -228,6 +361,13 @@ func (c *Collector) findTopologyProfiles(dev ddsnmp.DeviceConnectionInfo) []*dds }).Project(ddsnmp.ConsumerTopology).Profiles() } +func (c *Collector) getTopologyProfiles(dev ddsnmp.DeviceConnectionInfo) []*ddsnmp.Profile { + if c.topologyProfiles != nil { + return c.topologyProfiles(dev) + } + return c.findTopologyProfiles(dev) +} + func (c *Collector) ingestTopologyProfileMetrics(pms []*ddsnmp.ProfileMetrics) { for _, pm := range pms { c.ingestTopologyMetricSet(pm.TopologyMetrics) @@ -240,51 +380,6 @@ func (c *Collector) ingestTopologyMetricSet(metrics []ddsnmp.Metric) { } } -// collectTopologyMetrics reads the aggregated topology from the global registry. -func (c *Collector) collectTopologyMetrics(mx map[string]int64) { - if !c.topologyChartsAdded { - c.addTopologyCharts() - c.topologyChartsAdded = true - } - - data, ok := snmpTopologyRegistry.snapshot() - if !ok { - mx["snmp_topology_devices_total"] = 0 - mx["snmp_topology_devices_discovered"] = 0 - mx["snmp_topology_links_total"] = 0 - mx["snmp_topology_links_lldp"] = 0 - mx["snmp_topology_links_cdp"] = 0 - mx["snmp_topology_links_stp"] = 0 - return - } - - totalDevices := 0 - for _, actor := range data.Actors { - if topologyengine.IsDeviceActorType(actor.ActorType) { - totalDevices++ - } - } - - var lldpLinks, cdpLinks, stpLinks int64 - for _, link := range data.Links { - switch link.Protocol { - case "lldp": - lldpLinks++ - case "cdp": - cdpLinks++ - case "stp": - stpLinks++ - } - } - - mx["snmp_topology_devices_total"] = int64(totalDevices) - mx["snmp_topology_devices_discovered"] = int64(maxInt(totalDevices-1, 0)) - mx["snmp_topology_links_total"] = int64(len(data.Links)) - mx["snmp_topology_links_lldp"] = lldpLinks - mx["snmp_topology_links_cdp"] = cdpLinks - mx["snmp_topology_links_stp"] = stpLinks -} - func newSNMPClientFromDeviceInfo(newClient func() gosnmp.Handler, dev ddsnmp.DeviceConnectionInfo) (gosnmp.Handler, error) { client := newClient() 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 84eb8df06b3e91..95ac1cc78598fd 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 @@ -3,6 +3,8 @@ package snmptopology import ( + "context" + "sync" "testing" "time" @@ -16,6 +18,292 @@ import ( "github.com/netdata/netdata/go/plugins/plugin/go.d/pkg/snmputils" ) +func TestCollectorValidationLifecycleDoesNotStartPolling(t *testing.T) { + coll := New() + coll.UpdateEvery = 3600 + coll.registeredDevices = func() []ddsnmp.DeviceConnectionInfo { + return []ddsnmp.DeviceConnectionInfo{{ + Hostname: "192.0.2.10", + Port: 161, + }} + } + coll.newSnmpClient = func() gosnmp.Handler { + t.Fatal("validation lifecycle must not start topology polling") + return nil + } + + require.NoError(t, coll.Init(context.Background())) + require.NoError(t, coll.Check(context.Background())) + require.NoError(t, coll.Check(context.Background())) + coll.Cleanup(context.Background()) +} + +func TestCollectorRunRefreshesImmediatelyBeforeUpdateEvery(t *testing.T) { + coll := New() + coll.UpdateEvery = 3600 + coll.registeredDevices = func() []ddsnmp.DeviceConnectionInfo { + return []ddsnmp.DeviceConnectionInfo{{ + Hostname: "192.0.2.10", + Port: 161, + }} + } + + refreshed := make(chan struct{}, 1) + coll.newSnmpClient = func() gosnmp.Handler { + select { + case refreshed <- struct{}{}: + default: + } + panic("stop") + } + + ctx, cancel := context.WithCancel(context.Background()) + done := make(chan struct{}) + go func() { + defer close(done) + require.NoError(t, coll.Run(ctx)) + }() + + seen := false + require.Eventually(t, func() bool { + if seen { + return true + } + select { + case <-refreshed: + seen = true + return true + default: + return false + } + }, 2500*time.Millisecond, 10*time.Millisecond) + + cancel() + select { + case <-done: + case <-time.After(time.Second): + require.Fail(t, "runner did not stop") + } +} + +func TestCollectorRunStopsOnContextCancel(t *testing.T) { + coll := New() + coll.UpdateEvery = 3600 + + ctx, cancel := context.WithCancel(context.Background()) + done := make(chan struct{}) + go func() { + defer close(done) + require.NoError(t, coll.Run(ctx)) + }() + + cancel() + select { + case <-done: + case <-time.After(time.Second): + require.Fail(t, "runner did not stop") + } +} + +func TestCollectorRunDoesNotPollWhenContextAlreadyCanceled(t *testing.T) { + coll := New() + coll.registeredDevices = func() []ddsnmp.DeviceConnectionInfo { + return []ddsnmp.DeviceConnectionInfo{{ + Hostname: "192.0.2.10", + Port: 161, + }} + } + coll.newSnmpClient = func() gosnmp.Handler { + t.Fatal("Run must not poll with an already canceled context") + return nil + } + + ctx, cancel := context.WithCancel(context.Background()) + cancel() + + require.NoError(t, coll.Run(ctx)) +} + +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.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)) +} + +func TestCollectorRefreshTopologyRecoveringHandlesPanic(t *testing.T) { + coll := New() + coll.registeredDevices = func() []ddsnmp.DeviceConnectionInfo { + return []ddsnmp.DeviceConnectionInfo{{ + Hostname: "192.0.2.10", + Port: 161, + }} + } + coll.newSnmpClient = func() gosnmp.Handler { + panic("boom") + } + + require.NotPanics(t, func() { coll.refreshTopologyRecovering(context.Background()) }) + require.NotPanics(t, func() { coll.refreshTopologyRecovering(context.Background()) }) +} + +func TestCollectorRunCancelsInFlightRefresh(t *testing.T) { + ctrl := gomock.NewController(t) + defer ctrl.Finish() + + dev := ddsnmp.DeviceConnectionInfo{ + Hostname: "192.0.2.10", + Port: 161, + SysObjectID: "1.3.6.1.4.1.9.1.1", + } + mockHandler := snmpmock.NewMockHandler(ctrl) + mockHandler.EXPECT().SetTarget(dev.Hostname) + mockHandler.EXPECT().SetPort(uint16(dev.Port)) + mockHandler.EXPECT().SetRetries(dev.Retries) + mockHandler.EXPECT().SetTimeout(time.Duration(dev.Timeout) * time.Second) + mockHandler.EXPECT().SetMaxOids(dev.MaxOIDs) + mockHandler.EXPECT().SetMaxRepetitions(uint32(dev.MaxRepetitions)) + mockHandler.EXPECT().SetCommunity(dev.Community) + mockHandler.EXPECT().SetVersion(gosnmp.Version2c) + mockHandler.EXPECT().Connect().Return(nil) + + getStarted := make(chan struct{}) + closeCalled := make(chan struct{}) + var closeOnce sync.Once + mockHandler.EXPECT().Get(gomock.InAnyOrder([]string{ + snmputils.OidSnmpEngineTime, + snmputils.OidHrSystemUptime, + snmputils.OidSysUpTime, + })).DoAndReturn(func([]string) (*gosnmp.SnmpPacket, error) { + close(getStarted) + <-closeCalled + return nil, context.Canceled + }) + mockHandler.EXPECT().Close().DoAndReturn(func() error { + closeOnce.Do(func() { close(closeCalled) }) + return nil + }).AnyTimes() + + coll := New() + coll.UpdateEvery = 3600 + coll.registeredDevices = func() []ddsnmp.DeviceConnectionInfo { + return []ddsnmp.DeviceConnectionInfo{dev} + } + coll.newSnmpClient = func() gosnmp.Handler { return mockHandler } + coll.topologyProfiles = func(ddsnmp.DeviceConnectionInfo) []*ddsnmp.Profile { + return []*ddsnmp.Profile{{}} + } + coll.newDdSnmpColl = func(ddsnmpcollector.Config) ddCollector { + return ddCollectorFunc(func() ([]*ddsnmp.ProfileMetrics, error) { + return replacementEndpointProfileMetrics(), nil + }) + } + + ctx, cancel := context.WithCancel(context.Background()) + errCh := make(chan error, 1) + go func() { + errCh <- coll.Run(ctx) + }() + + select { + case <-getStarted: + case <-time.After(5 * time.Second): + require.Fail(t, "refresh did not start") + } + + cancel() + select { + case err := <-errCh: + require.NoError(t, err) + case <-time.After(time.Second): + require.Fail(t, "runner did not stop after context cancellation") + } +} + +func TestCollectorCancelsInFlightVLANContextRefresh(t *testing.T) { + ctrl := gomock.NewController(t) + defer ctrl.Finish() + + dev := ddsnmp.DeviceConnectionInfo{ + Hostname: "192.0.2.10", + Port: 161, + Community: "public", + } + mockHandler := snmpmock.NewMockHandler(ctrl) + mockHandler.EXPECT().SetTarget(dev.Hostname) + mockHandler.EXPECT().SetPort(uint16(dev.Port)) + mockHandler.EXPECT().SetRetries(dev.Retries) + mockHandler.EXPECT().SetTimeout(time.Duration(dev.Timeout) * time.Second) + mockHandler.EXPECT().SetMaxOids(dev.MaxOIDs) + mockHandler.EXPECT().SetMaxRepetitions(uint32(dev.MaxRepetitions)) + mockHandler.EXPECT().SetCommunity(dev.Community) + mockHandler.EXPECT().SetVersion(gosnmp.Version2c) + mockHandler.EXPECT().Version().Return(gosnmp.Version2c) + mockHandler.EXPECT().Community().Return(dev.Community) + mockHandler.EXPECT().SetCommunity(dev.Community + "@100") + + closeCalled := make(chan struct{}) + var closeOnce sync.Once + mockHandler.EXPECT().Connect().Return(nil) + mockHandler.EXPECT().Close().DoAndReturn(func() error { + closeOnce.Do(func() { close(closeCalled) }) + return nil + }).AnyTimes() + + coll := New() + coll.newSnmpClient = func() gosnmp.Handler { return mockHandler } + collectStarted := make(chan struct{}) + coll.newDdSnmpColl = func(ddsnmpcollector.Config) ddCollector { + return ddCollectorFunc(func() ([]*ddsnmp.ProfileMetrics, error) { + close(collectStarted) + <-closeCalled + return nil, context.Canceled + }) + } + + ctx, cancel := context.WithCancel(context.Background()) + errCh := make(chan error, 1) + go func() { + _, err := collectTopologyVLANContext(ctx, coll, dev, "100", nil) + errCh <- err + }() + + select { + case <-collectStarted: + case <-time.After(time.Second): + require.Fail(t, "vlan-context refresh did not start") + } + + cancel() + select { + case err := <-errCh: + require.ErrorIs(t, err, context.Canceled) + case <-time.After(time.Second): + require.Fail(t, "vlan-context refresh did not stop after context cancellation") + } +} + +func TestCollectorNewDeviceCollectionCacheUsesEffectiveDeviceCheckEvery(t *testing.T) { + coll := New() + + cache := coll.newDeviceCollectionCache(ddsnmp.DeviceConnectionInfo{Hostname: "switch-a"}) + + require.Equal(t, defaultRefreshEvery+2*defaultDeviceCheckEvery, cache.staleAfter) +} + func TestCollector_RefreshKeepsPublishedSnapshotWhileCollectionRuns(t *testing.T) { previousRegistry := snmpTopologyRegistry registry := newTopologyRegistry() @@ -55,7 +343,7 @@ func TestCollector_RefreshKeepsPublishedSnapshotWhileCollectionRuns(t *testing.T go func() { defer close(done) - coll.refreshDeviceTopology(key, dev) + coll.refreshDeviceTopology(context.Background(), key, dev) }() <-started @@ -76,6 +364,10 @@ type blockingTopologyCollector struct { result []*ddsnmp.ProfileMetrics } +type ddCollectorFunc func() ([]*ddsnmp.ProfileMetrics, error) + +func (f ddCollectorFunc) Collect() ([]*ddsnmp.ProfileMetrics, error) { return f() } + func (c *blockingTopologyCollector) Collect() ([]*ddsnmp.ProfileMetrics, error) { close(c.started) <-c.release @@ -160,3 +452,10 @@ func replacementEndpointProfileMetrics() []*ddsnmp.ProfileMetrics { }, }} } + +func topologyRegistryHasCache(registry *topologyRegistry, cache *topologyCache) bool { + registry.mu.RLock() + defer registry.mu.RUnlock() + _, ok := registry.caches[cache] + return ok +} diff --git a/src/go/plugin/go.d/collector/snmp_topology/config.go b/src/go/plugin/go.d/collector/snmp_topology/config.go index 42394a58d01ec1..7109f1beda564b 100644 --- a/src/go/plugin/go.d/collector/snmp_topology/config.go +++ b/src/go/plugin/go.d/collector/snmp_topology/config.go @@ -5,7 +5,7 @@ package snmptopology import "github.com/netdata/netdata/go/plugins/pkg/confopt" // Config for the snmp_topology module. -// This module has a single global job — device list comes from the SNMP device registry. +// This module has a single global collector instance; devices come from the SNMP device registry. type Config struct { UpdateEvery int `yaml:"update_every,omitempty" json:"update_every"` RefreshEvery confopt.LongDuration `yaml:"refresh_every,omitempty" json:"refresh_every,omitempty"` 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 dcf9afed42d2f4..22671414d84c28 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 @@ -9,7 +9,10 @@ var _ funcapi.MethodHandler = (*funcTopology)(nil) type funcTopology struct{} -const topologyMethodID = "topology:snmp" +const ( + topologyFunctionName = "snmp:topology:snmp" + topologyMethodID = "topology:snmp" +) const ( topologyParamNodesIdentity = "nodes_identity" diff --git a/src/go/plugin/go.d/collector/snmp_topology/func_topology_presentation.go b/src/go/plugin/go.d/collector/snmp_topology/func_topology_presentation.go index 60176077872c4e..caeb56dab284df 100644 --- a/src/go/plugin/go.d/collector/snmp_topology/func_topology_presentation.go +++ b/src/go/plugin/go.d/collector/snmp_topology/func_topology_presentation.go @@ -7,6 +7,7 @@ import "github.com/netdata/netdata/go/plugins/pkg/funcapi" func topologyMethodConfig() funcapi.MethodConfig { return funcapi.MethodConfig{ ID: topologyMethodID, + FunctionName: topologyFunctionName, Aliases: []string{topologyMethodID}, Name: "Topology (SNMP)", UpdateEvery: 10, 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 498b975376f65f..0e41515dd8580b 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 @@ -5,6 +5,7 @@ package snmptopology import ( "testing" + "github.com/netdata/netdata/go/plugins/plugin/framework/collectorapi" "github.com/stretchr/testify/require" ) @@ -12,6 +13,27 @@ func TestSNMPTopologyMethodConfigDoesNotUseLegacyPresentation(t *testing.T) { method := topologyMethodConfig() require.Equal(t, topologyMethodID, method.ID) + require.Equal(t, topologyFunctionName, method.FunctionName) + require.Equal(t, []string{topologyMethodID}, method.Aliases) require.Equal(t, "topology", method.ResponseType) require.Nil(t, method.Presentation()) } + +func TestSNMPTopologyCreatorOwnsTopologyFunction(t *testing.T) { + creator, ok := collectorapi.DefaultRegistry.Lookup("snmp_topology") + require.True(t, ok) + require.Nil(t, creator.Create) + require.NotNil(t, creator.CreateV2) + require.Equal(t, collectorapi.InstancePolicySingle, creator.InstancePolicy) + require.False(t, creator.FunctionOnly) + require.NotNil(t, creator.Methods) + require.NotNil(t, creator.MethodHandler) + require.Implements(t, (*collectorapi.CollectorV2Runner)(nil), creator.CreateV2()) + + methods := creator.Methods() + require.Len(t, methods, 1) + 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)) +} diff --git a/src/go/plugin/go.d/collector/snmp_topology/metrix.go b/src/go/plugin/go.d/collector/snmp_topology/metrix.go new file mode 100644 index 00000000000000..d13999f3e3cdcf --- /dev/null +++ b/src/go/plugin/go.d/collector/snmp_topology/metrix.go @@ -0,0 +1,114 @@ +// SPDX-License-Identifier: GPL-3.0-or-later + +package snmptopology + +import ( + "time" + + "github.com/netdata/netdata/go/plugins/pkg/metrix" +) + +const internalMetricPrefix = "netdata.go.plugin.collector.snmp_topology" + +type collectorMetrics struct { + devicesRegistered metrix.SnapshotGauge + devicesCached metrix.SnapshotGauge + lastRefreshAgeSeconds metrix.SnapshotGauge + lastRefreshDurationSecs metrix.SnapshotGauge + refreshRuns metrix.SnapshotCounter + refreshErrors metrix.SnapshotCounter +} + +type collectorRuntimeStats struct { + registeredDevices int + cachedDevices int + lastRefresh time.Time + lastDuration time.Duration + refreshRuns uint64 + refreshErrors uint64 +} + +type refreshStats struct { + hasDeviceCounts bool + registeredDevices int + cachedDevices int + errors int + completedAt time.Time + duration time.Duration +} + +func newCollectorMetrics(store metrix.CollectorStore) *collectorMetrics { + meter := store.Write().SnapshotMeter(internalMetricPrefix) + return &collectorMetrics{ + devicesRegistered: meter.Gauge("devices_registered"), + devicesCached: meter.Gauge("devices_cached"), + lastRefreshAgeSeconds: meter.Gauge("last_refresh_age_seconds"), + lastRefreshDurationSecs: meter.Gauge("last_refresh_duration_seconds"), + refreshRuns: meter.Counter("refresh_runs_total"), + refreshErrors: meter.Counter("refresh_errors_total"), + } +} + +func (c *Collector) recordRefreshStats(stats refreshStats) { + c.statsMu.Lock() + defer c.statsMu.Unlock() + + if stats.hasDeviceCounts { + c.stats.registeredDevices = stats.registeredDevices + c.stats.cachedDevices = stats.cachedDevices + } + c.stats.lastRefresh = stats.completedAt + c.stats.lastDuration = stats.duration + c.stats.refreshRuns++ + c.stats.refreshErrors += uint64(stats.errors) +} + +func (c *Collector) recordCleanupStats() { + c.statsMu.Lock() + defer c.statsMu.Unlock() + + c.stats.registeredDevices = 0 + c.stats.cachedDevices = 0 +} + +func (c *Collector) writeInternalMetrics(now time.Time) { + stats := c.runtimeStatsSnapshot(now) + + c.metrics.devicesRegistered.Observe(float64(stats.registeredDevices)) + c.metrics.devicesCached.Observe(float64(stats.cachedDevices)) + c.metrics.lastRefreshAgeSeconds.Observe(stats.lastRefreshAgeSeconds) + c.metrics.lastRefreshDurationSecs.Observe(stats.lastRefreshDurationSeconds) + c.metrics.refreshRuns.ObserveTotal(float64(stats.refreshRuns)) + c.metrics.refreshErrors.ObserveTotal(float64(stats.refreshErrors)) +} + +type runtimeStatsSnapshot struct { + registeredDevices int + cachedDevices int + lastRefreshAgeSeconds float64 + lastRefreshDurationSeconds float64 + refreshRuns uint64 + refreshErrors uint64 +} + +func (c *Collector) runtimeStatsSnapshot(now time.Time) runtimeStatsSnapshot { + c.statsMu.RLock() + defer c.statsMu.RUnlock() + + var ageSeconds float64 + if !c.stats.lastRefresh.IsZero() { + ageSeconds = now.Sub(c.stats.lastRefresh).Seconds() + if ageSeconds < 0 { + ageSeconds = 0 + } + } + + return runtimeStatsSnapshot{ + registeredDevices: c.stats.registeredDevices, + cachedDevices: c.stats.cachedDevices, + lastRefreshAgeSeconds: ageSeconds, + lastRefreshDurationSeconds: c.stats.lastDuration.Seconds(), + refreshRuns: c.stats.refreshRuns, + refreshErrors: c.stats.refreshErrors, + } +} diff --git a/src/go/plugin/go.d/collector/snmp_topology/metrix_test.go b/src/go/plugin/go.d/collector/snmp_topology/metrix_test.go new file mode 100644 index 00000000000000..213c497fd166a3 --- /dev/null +++ b/src/go/plugin/go.d/collector/snmp_topology/metrix_test.go @@ -0,0 +1,81 @@ +// SPDX-License-Identifier: GPL-3.0-or-later + +package snmptopology + +import ( + "context" + "testing" + "time" + + "github.com/stretchr/testify/require" + + "github.com/netdata/netdata/go/plugins/pkg/metrix" + "github.com/netdata/netdata/go/plugins/plugin/go.d/collector/snmp/ddsnmp" + "github.com/netdata/netdata/go/plugins/plugin/go.d/pkg/collecttest" +) + +func TestCollector_WriteInternalMetrics(t *testing.T) { + coll := New() + now := time.Unix(100, 0) + coll.recordRefreshStats(refreshStats{ + hasDeviceCounts: true, + registeredDevices: 2, + cachedDevices: 1, + errors: 3, + completedAt: now.Add(-5 * time.Second), + duration: 1500 * time.Millisecond, + }) + + managed, ok := metrix.AsCycleManagedStore(coll.MetricStore()) + require.True(t, ok) + managed.CycleController().BeginCycle() + coll.writeInternalMetrics(now) + require.NoError(t, managed.CycleController().CommitCycleSuccess()) + + reader := coll.MetricStore().Read(metrix.ReadRaw()) + requireMetricValue(t, reader, internalMetricPrefix+".devices_registered", 2) + requireMetricValue(t, reader, internalMetricPrefix+".devices_cached", 1) + requireMetricValue(t, reader, internalMetricPrefix+".last_refresh_age_seconds", 5) + requireMetricValue(t, reader, internalMetricPrefix+".last_refresh_duration_seconds", 1.5) + requireMetricValue(t, reader, internalMetricPrefix+".refresh_runs_total", 1) + requireMetricValue(t, reader, internalMetricPrefix+".refresh_errors_total", 3) +} + +func TestCollectorCollectWritesInternalMetrics(t *testing.T) { + coll := New() + coll.registeredDevices = func() []ddsnmp.DeviceConnectionInfo { + t.Fatal("Collect must not poll SNMP devices") + return nil + } + coll.recordRefreshStats(refreshStats{ + hasDeviceCounts: true, + registeredDevices: 1, + cachedDevices: 1, + completedAt: time.Now(), + }) + + managed, ok := metrix.AsCycleManagedStore(coll.MetricStore()) + require.True(t, ok) + managed.CycleController().BeginCycle() + require.NoError(t, coll.Collect(context.Background())) + require.NoError(t, managed.CycleController().CommitCycleSuccess()) + + reader := coll.MetricStore().Read(metrix.ReadRaw()) + requireMetricValue(t, reader, internalMetricPrefix+".devices_registered", 1) + requireMetricValue(t, reader, internalMetricPrefix+".devices_cached", 1) + requireMetricValue(t, reader, internalMetricPrefix+".refresh_runs_total", 1) + collecttest.AssertChartCoverage(t, coll, collecttest.ChartCoverageExpectation{ + RequiredContexts: map[string][]string{ + "netdata.go.plugin.collector.snmp_topology.devices": {"cached", "registered"}, + "netdata.go.plugin.collector.snmp_topology.last_refresh": {"age", "duration"}, + "netdata.go.plugin.collector.snmp_topology.refreshes": {"errors", "runs"}, + }, + }) +} + +func requireMetricValue(t *testing.T, reader metrix.Reader, name string, want metrix.SampleValue) { + t.Helper() + got, ok := reader.Value(name, nil) + require.True(t, ok, "metric %s not found", name) + require.Equal(t, want, got) +} 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 1c7ac5868b09df..059da2e0e82f43 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 @@ -128,7 +128,7 @@ func topologyNormalizeReverseDNSName(names []string) string { var defaultTopologyReverseDNSResolver = newTopologyReverseDNSResolver(topologyReverseDNSTimeout, topologyReverseDNSCacheTTL) // resolveTopologyReverseDNSName performs a live DNS lookup (with cache). -// Used during the collector's Collect() cycle to warm the cache. +// Used while building topology snapshots to warm the cache. func resolveTopologyReverseDNSName(ip string) string { return defaultTopologyReverseDNSResolver.lookup(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 0839b407807b0a..77ed7445487639 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 @@ -97,12 +97,12 @@ func collectTopologySnapshotFromDevice(t *testing.T, dev ddsnmp.DeviceConnection defer ddsnmp.DeviceRegistry.Unregister(deviceKey) coll := New() - coll.Config = Config{UpdateEvery: 1} + coll.Config = Config{UpdateEvery: 3600} require.NoError(t, coll.Init(context.Background())) defer coll.Cleanup(context.Background()) require.NoError(t, coll.Check(context.Background())) - _ = coll.Collect(context.Background()) + coll.refreshTopology(context.Background()) var snapshot topologyData cacheKey := dev.Hostname + ":" + strconv.Itoa(dev.Port) diff --git a/src/go/plugin/go.d/collector/snmp_topology/topology_metrics_test.go b/src/go/plugin/go.d/collector/snmp_topology/topology_metrics_test.go deleted file mode 100644 index 3995a6a863111a..00000000000000 --- a/src/go/plugin/go.d/collector/snmp_topology/topology_metrics_test.go +++ /dev/null @@ -1,48 +0,0 @@ -// SPDX-License-Identifier: GPL-3.0-or-later - -package snmptopology - -import ( - "testing" - "time" - - "github.com/stretchr/testify/assert" -) - -func TestCollectTopologyMetrics(t *testing.T) { - cache := newTopologyCache() - cache.lastUpdate = time.Now() - cache.localDevice = topologyDevice{ - ManagementIP: "192.0.2.1", - ChassisID: "aa:bb:cc:dd:ee:ff", - ChassisIDType: "macAddress", - } - cache.lldpLocPorts["1"] = &lldpLocPort{ - portNum: "1", - portID: "Gi0/1", - portIDSubtype: "interfaceName", - } - cache.lldpRemotes["1:1"] = &lldpRemote{ - localPortNum: "1", - remIndex: "1", - chassisID: "11:22:33:44:55:66", - chassisIDSubtype: "macAddress", - portID: "Gi0/2", - portIDSubtype: "interfaceName", - } - - snmpTopologyRegistry.register(cache) - defer snmpTopologyRegistry.unregister(cache) - - c := New() - - mx := make(map[string]int64) - c.collectTopologyMetrics(mx) - - // The topology engine pipeline processes raw cache data into actors and links. - // With minimal test data (one local + one LLDP remote), the engine produces - // at least the local device and the remote device as actors. - assert.GreaterOrEqual(t, mx["snmp_topology_devices_total"], int64(1)) - assert.GreaterOrEqual(t, mx["snmp_topology_links_total"], int64(0)) - assert.True(t, c.topologyChartsAdded) -} 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 39cabff6d95cf4..6a1161320addf6 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 @@ -52,18 +52,6 @@ func (r *topologyRegistry) unregister(cache *topologyCache) { r.mu.Unlock() } -func (r *topologyRegistry) snapshot() (topologyData, bool) { - return r.snapshotWithOptions(topologyQueryOptions{ - CollapseActorsByIP: true, - EliminateNonIPInferred: true, - MapType: topologyMapTypeLLDPCDPManaged, - InferenceStrategy: topologyInferenceStrategyFDBMinimumKnowledge, - ManagedDeviceFocus: topologyManagedFocusAllDevices, - Depth: topologyDepthAllInternal, - ResolveDNSName: resolveTopologyReverseDNSName, // live resolver — warms the cache - }) -} - func (r *topologyRegistry) snapshotWithOptions(options topologyQueryOptions) (topologyData, bool) { if r == nil { return topologyData{}, false diff --git a/src/go/plugin/go.d/collector/snmp_topology/topology_registry_test.go b/src/go/plugin/go.d/collector/snmp_topology/topology_registry_test.go index 8d425cb356c5fd..727d98317be21b 100644 --- a/src/go/plugin/go.d/collector/snmp_topology/topology_registry_test.go +++ b/src/go/plugin/go.d/collector/snmp_topology/topology_registry_test.go @@ -68,7 +68,7 @@ func TestTopologyRegistry_SnapshotAggregatesAcrossCaches(t *testing.T) { registry.register(cacheA) registry.register(cacheB) - data, ok := registry.snapshot() + data, ok := snapshotTopologyRegistryForTest(registry) require.True(t, ok) require.Equal(t, "2", data.Layer) require.Equal(t, "snmp", data.Source) @@ -110,7 +110,7 @@ func TestTopologyRegistry_SnapshotSingleCacheKeepsLLDPUnidirectional(t *testing. registry.register(cache) - data, ok := registry.snapshot() + data, ok := snapshotTopologyRegistryForTest(registry) require.True(t, ok) require.Len(t, data.Links, 1) require.Equal(t, "lldp", data.Links[0].Protocol) @@ -298,7 +298,7 @@ func TestTopologyRegistry_SnapshotReturnsFalseWithoutCollectedCaches(t *testing. cache := newTopologyCache() registry.register(cache) - _, ok := registry.snapshot() + _, ok := snapshotTopologyRegistryForTest(registry) require.False(t, ok) } @@ -360,13 +360,13 @@ func TestTopologyRegistry_SnapshotDeterministicAcrossRepeatedCalls(t *testing.T) registry.register(cacheA) registry.register(cacheB) - baseline, ok := registry.snapshot() + baseline, ok := snapshotTopologyRegistryForTest(registry) require.True(t, ok) require.NotEmpty(t, baseline.Actors) require.NotEmpty(t, baseline.Links) for range 10 { - next, ok := registry.snapshot() + next, ok := snapshotTopologyRegistryForTest(registry) require.True(t, ok) require.Equal(t, baseline, next) } @@ -412,7 +412,7 @@ func TestTopologyRegistry_SnapshotDeduplicatesDuplicateDeviceObservations(t *tes registry.register(cacheA) registry.register(cacheB) - data, ok := registry.snapshot() + data, ok := snapshotTopologyRegistryForTest(registry) require.True(t, ok) require.Len(t, data.Links, 1) diff --git a/src/go/plugin/go.d/collector/snmp_topology/topology_snapshot_race_test.go b/src/go/plugin/go.d/collector/snmp_topology/topology_snapshot_race_test.go index 9ec23b0014dd97..c31552ff05c41c 100644 --- a/src/go/plugin/go.d/collector/snmp_topology/topology_snapshot_race_test.go +++ b/src/go/plugin/go.d/collector/snmp_topology/topology_snapshot_race_test.go @@ -12,8 +12,7 @@ import ( // concurrent-map-write crash where the snapshot read path mutated a cache-shared // map. snapshotEngineObservations holds only an RLock and passes c.localDevice // (by value, so Labels still aliases the cache map) to normalizeTopologyDevice, -// which writes Labels. Two concurrent readers (Collect's snapshot() and a -// Function's snapshotWithOptions()) then write the same map. +// which writes Labels. Two concurrent snapshot readers then write the same map. // // Run with -race. Before the fix this reports a data race / fatal concurrent map // writes; after the fix (normalizeTopologyDevice clones Labels) it passes. @@ -50,13 +49,12 @@ func TestTopologyRegistry_ConcurrentSnapshotsDoNotRaceOnDeviceLabels(t *testing. var wg sync.WaitGroup wg.Add(goroutines) for i := range goroutines { - // Alternate the two production read paths: Collect (snapshot) and the - // Function handler (snapshotWithOptions). + // Alternate the default and option-aware registry read paths. collectPath := i%2 == 0 go func() { defer wg.Done() if collectPath { - _, _ = registry.snapshot() + _, _ = snapshotTopologyRegistryForTest(registry) } else { _, _ = registry.snapshotWithOptions(topologyQueryOptions{}) } diff --git a/src/go/plugin/go.d/collector/snmp_topology/topology_snmp_value_helpers.go b/src/go/plugin/go.d/collector/snmp_topology/topology_snmp_value_helpers.go index 1790023dd39a03..53e00d574424db 100644 --- a/src/go/plugin/go.d/collector/snmp_topology/topology_snmp_value_helpers.go +++ b/src/go/plugin/go.d/collector/snmp_topology/topology_snmp_value_helpers.go @@ -29,10 +29,3 @@ func parseIndex(value string) int { } return v } - -func maxInt(a, b int) int { - if a > b { - return a - } - return b -} 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 af8c3691dafe89..f1b09429a5be68 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 @@ -2,6 +2,18 @@ package snmptopology +func snapshotTopologyRegistryForTest(registry *topologyRegistry) (topologyData, bool) { + return registry.snapshotWithOptions(topologyQueryOptions{ + CollapseActorsByIP: true, + EliminateNonIPInferred: true, + MapType: topologyMapTypeLLDPCDPManaged, + InferenceStrategy: topologyInferenceStrategyFDBMinimumKnowledge, + ManagedDeviceFocus: topologyManagedFocusAllDevices, + Depth: topologyDepthAllInternal, + ResolveDNSName: resolveTopologyReverseDNSName, + }) +} + func containsMgmtAddr(snapshot topologyData, addrs map[string]struct{}) bool { for _, actor := range snapshot.Actors { for _, ip := range actor.Match.IPAddresses { diff --git a/src/go/plugin/go.d/collector/snmp_topology/topology_vlan_context.go b/src/go/plugin/go.d/collector/snmp_topology/topology_vlan_context.go index ed8ba68acd7384..851421c157fb74 100644 --- a/src/go/plugin/go.d/collector/snmp_topology/topology_vlan_context.go +++ b/src/go/plugin/go.d/collector/snmp_topology/topology_vlan_context.go @@ -3,13 +3,18 @@ package snmptopology import ( + "context" + "github.com/netdata/netdata/go/plugins/plugin/go.d/collector/snmp/ddsnmp" ) -func (c *Collector) collectTopologyVTPVLANContexts(dev ddsnmp.DeviceConnectionInfo) { +func (c *Collector) collectTopologyVTPVLANContexts(ctx context.Context, dev ddsnmp.DeviceConnectionInfo) { if c.topologyCache == nil { return } + if ctx.Err() != nil { + return + } contexts := c.topologyCache.vtpVLANContexts() if len(contexts) == 0 { @@ -23,8 +28,15 @@ func (c *Collector) collectTopologyVTPVLANContexts(dev ddsnmp.DeviceConnectionIn } for _, context := range contexts { - pms, err := collectTopologyVLANContext(c, dev, context.vlanID, profiles) + if ctx.Err() != nil { + return + } + + pms, err := collectTopologyVLANContext(ctx, c, dev, context.vlanID, profiles) if err != nil { + if ctx.Err() != nil { + return + } c.Warningf("device '%s': topology vlan-context polling failed for vlan %s: %v", dev.Hostname, context.vlanID, err) continue } diff --git a/src/go/plugin/go.d/collector/snmp_topology/topology_vlan_context_collect.go b/src/go/plugin/go.d/collector/snmp_topology/topology_vlan_context_collect.go index 5ef6a841a2a1ce..191f9edd10ce99 100644 --- a/src/go/plugin/go.d/collector/snmp_topology/topology_vlan_context_collect.go +++ b/src/go/plugin/go.d/collector/snmp_topology/topology_vlan_context_collect.go @@ -3,6 +3,7 @@ package snmptopology import ( + "context" "fmt" "strconv" "strings" @@ -22,22 +23,30 @@ func loadTopologyVLANContextProfiles(dev ddsnmp.DeviceConnectionInfo) ([]*ddsnmp }).Project(ddsnmp.ConsumerTopology).FilterByKind(vlanScopableKinds).Profiles(), nil } -func collectTopologyVLANContext(c *Collector, dev ddsnmp.DeviceConnectionInfo, vlanID string, profiles []*ddsnmp.Profile) ([]*ddsnmp.ProfileMetrics, error) { +func collectTopologyVLANContext(ctx context.Context, c *Collector, dev ddsnmp.DeviceConnectionInfo, vlanID string, profiles []*ddsnmp.Profile) ([]*ddsnmp.ProfileMetrics, error) { if strings.TrimSpace(vlanID) == "" { return nil, fmt.Errorf("empty vlan id") } if _, err := strconv.Atoi(vlanID); err != nil { return nil, fmt.Errorf("invalid vlan id '%s': %w", vlanID, err) } + if err := ctx.Err(); err != nil { + return nil, err + } - snmpClient, err := initTopologyVLANClient(c, dev, vlanID) + snmpClient, stopContextClose, err := initTopologyVLANClient(ctx, c, dev, vlanID) if err != nil { return nil, err } + defer stopContextClose() defer func() { _ = snmpClient.Close() }() + if err := ctx.Err(); err != nil { + return nil, err + } + vlanCollector := c.newDdSnmpColl(ddsnmpcollector.Config{ SnmpClient: snmpClient, Profiles: profiles, @@ -46,13 +55,24 @@ func collectTopologyVLANContext(c *Collector, dev ddsnmp.DeviceConnectionInfo, v DisableBulkWalk: dev.DisableBulkWalk, }) - return vlanCollector.Collect() + pms, err := vlanCollector.Collect() + if err != nil { + return nil, err + } + if err := ctx.Err(); err != nil { + return nil, err + } + return pms, nil } -func initTopologyVLANClient(c *Collector, dev ddsnmp.DeviceConnectionInfo, vlanID string) (gosnmp.Handler, error) { +func initTopologyVLANClient(ctx context.Context, c *Collector, dev ddsnmp.DeviceConnectionInfo, vlanID string) (gosnmp.Handler, func(), error) { + if err := ctx.Err(); err != nil { + return nil, func() {}, err + } + client, err := newSNMPClientFromDeviceInfo(c.newSnmpClient, dev) if err != nil { - return nil, err + return nil, func() {}, err } switch client.Version() { @@ -71,8 +91,9 @@ func initTopologyVLANClient(c *Collector, dev ddsnmp.DeviceConnectionInfo, vlanI } if err := client.Connect(); err != nil { - return nil, err + return nil, func() {}, err } - return client, nil + stopContextClose := closeSNMPClientOnContextCancel(ctx, client) + return client, stopContextClose, nil } 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 6a1d781b7c206a..a219b1564dd656 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 @@ -11,6 +11,9 @@ import ( "time" "github.com/netdata/netdata/go/plugins/pkg/funcapi" + "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_traps/snmptrapsfunc" "github.com/stretchr/testify/assert" "github.com/stretchr/testify/require" @@ -138,6 +141,44 @@ func TestSNMPTrapsLogsFunctionUnavailableWithoutJournal(t *testing.T) { assert.Contains(t, resp.Message, "direct journal output has no sources") } +func TestSNMPTrapsLogsDispatchDoesNotRequireRunningJob(t *testing.T) { + startJournalJobs := activeDirectJournalJobs.Load() + activeDirectJournalJobs.Store(1) + 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) + + reg := newSNMPTrapsTestFunctionRegistry() + var gotCode int + var gotResp map[string]any + controller := funcctl.New(funcctl.Options{ + FnReg: reg, + JSONWriter: func(data []byte, code int) { + gotCode = code + require.NoError(t, json.Unmarshal(data, &gotResp)) + }, + }) + controller.RegisterModules(collectorapi.Registry{"snmp_traps": creator}) + + require.Contains(t, reg.registeredNames(), snmpTrapsFunctionName) + + reg.call(snmpTrapsFunctionName, context.Background(), functions.Function{ + UID: "snmp-traps-jobless-dispatch", + Timeout: time.Second, + Payload: []byte(`{}`), + }) + + require.NotNil(t, gotResp) + assert.Equal(t, 503, gotCode) + assert.Equal(t, float64(503), gotResp["status"]) + errorMessage := fmt.Sprint(gotResp["errorMessage"]) + assert.Contains(t, errorMessage, "direct journal output has no sources") + assert.NotContains(t, errorMessage, "unknown job") + assert.NotContains(t, errorMessage, "requires raw request handling") +} + func TestSNMPTrapsLogsFunctionRejectsUnknownMethod(t *testing.T) { handler := snmptrapsfunc.NewHandler(t.TempDir()) @@ -231,6 +272,55 @@ func funcapiRawRequest(method string, info bool, payload []byte) funcapi.RawMeth } } +type snmpTrapsTestFunctionRegistry struct { + handlers map[string]func(functions.Function) + contextHandlers map[string]functions.Handler +} + +func newSNMPTrapsTestFunctionRegistry() *snmpTrapsTestFunctionRegistry { + return &snmpTrapsTestFunctionRegistry{ + handlers: make(map[string]func(functions.Function)), + contextHandlers: make(map[string]functions.Handler), + } +} + +func (r *snmpTrapsTestFunctionRegistry) Register(name string, fn func(functions.Function)) { + r.handlers[name] = fn +} + +func (r *snmpTrapsTestFunctionRegistry) RegisterWithContext(name string, fn functions.Handler) { + r.handlers[name] = func(f functions.Function) { + fn(context.Background(), f) + } + r.contextHandlers[name] = fn +} + +func (r *snmpTrapsTestFunctionRegistry) Unregister(name string) { + delete(r.handlers, name) + delete(r.contextHandlers, name) +} + +func (r *snmpTrapsTestFunctionRegistry) RegisterPrefix(string, string, func(functions.Function)) {} +func (r *snmpTrapsTestFunctionRegistry) UnregisterPrefix(string, string) {} +func (r *snmpTrapsTestFunctionRegistry) RegisterPrefixWithContext(string, string, functions.Handler) { +} + +func (r *snmpTrapsTestFunctionRegistry) call(name string, ctx context.Context, fn functions.Function) { + if handler := r.contextHandlers[name]; handler != nil { + handler(ctx, fn) + return + } + r.handlers[name](fn) +} + +func (r *snmpTrapsTestFunctionRegistry) registeredNames() []string { + out := make([]string, 0, len(r.handlers)) + for name := range r.handlers { + out = append(out, name) + } + return out +} + func assertResponseHistogramID(t *testing.T, response map[string]any, want string) { t.Helper() histogram, ok := response["histogram"].(map[string]any) diff --git a/src/go/plugin/go.d/docs/how-to-write-a-collector.md b/src/go/plugin/go.d/docs/how-to-write-a-collector.md index a6dc0bb5dca93b..5accd1947b4f0d 100644 --- a/src/go/plugin/go.d/docs/how-to-write-a-collector.md +++ b/src/go/plugin/go.d/docs/how-to-write-a-collector.md @@ -161,6 +161,16 @@ Public lifecycle and framework-contract methods MUST stay in `collector.go`: - `MetricStore() metrix.CollectorStore` - `ChartTemplateYAML() string` +Collectors that need a long-running side-effect loop MAY additionally implement +`collectorapi.CollectorV2Runner` with `Run(context.Context) error`. Use this +only when work must start with the running job lifecycle but must not wait for +the next globally aligned `Collect()` tick, such as an agent-wide Function state +refresh. The runtime starts `Run()` only after the job starts, never during +autodetection or DynCfg `test`, cancels it on stop, and waits for it before +`Cleanup()`. The implementation MUST return promptly after `ctx.Done()` and +SHOULD make in-flight I/O cancellation-aware where the underlying library allows +it. + `Init()` validates config, prepares matchers/clients, and initializes persistent state. Explicit setup details SHOULD live in helper methods, preferably in `init.go`, so the public method reads as the lifecycle sequence. `Check()` MUST @@ -414,6 +424,8 @@ exactly what was run. - New collector using `Collect() map[string]int64`. - Full live collection from `Check()`. +- Starting operational background polling from `Check()` or `Init()`; use the + optional V2 runner hook when polling must be tied to the running job lifecycle. - Public config knobs for internal implementation details. - Custom selector or retry framework when existing package/framework behavior is enough.