diff --git a/.chloggen/splunk-inputs-precedence.yaml b/.chloggen/splunk-inputs-precedence.yaml new file mode 100644 index 0000000..f6332c4 --- /dev/null +++ b/.chloggen/splunk-inputs-precedence.yaml @@ -0,0 +1,14 @@ +change_type: enhancement + +component: splunk_inputs + +note: "add Splunk btool config precedence across system and `splunk_ta_*` app directories" + +issues: [117] + +subtext: | + The `base_dir` config field (or `$SPLUNK_HOME` env var) now points to the Splunk + installation root. The receiver discovers all `splunk_ta_*` directories automatically + and merges `inputs.conf`, `transforms.conf`, and `props.conf` per TA using standard + btool precedence — later layers win per key. Unsupported input and output stanza + kinds are now skipped silently with an info log instead of returning an error. diff --git a/internal/collector/collector.go b/internal/collector/collector.go index 07b4ba4..5ef0dc7 100644 --- a/internal/collector/collector.go +++ b/internal/collector/collector.go @@ -36,15 +36,16 @@ func Run(baseDir string, cfg *config.Config) (func(), error) { if err != nil { return nil, err } - inputs, err := tabuilder.ReadInputs(baseDir) + dirs := tabuilder.ConfDirs(baseDir) + inputs, err := tabuilder.ReadInputs(dirs) if err != nil { return nil, err } - transforms, err := tabuilder.ReadTransforms(baseDir) + transforms, err := tabuilder.ReadTransforms(dirs) if err != nil { return nil, err } - props, err := tabuilder.ReadProps(baseDir) + props, err := tabuilder.ReadProps(dirs) if err != nil { return nil, err } diff --git a/internal/collector/collector_test.go b/internal/collector/collector_test.go index a62ea28..8e4fd90 100644 --- a/internal/collector/collector_test.go +++ b/internal/collector/collector_test.go @@ -37,7 +37,7 @@ func TestRunTA(t *testing.T) { defer func() { _ = rcvr.Shutdown(context.Background()) }() - cancel, err := Run(filepath.Join("testdata", "ta"), &config.Config{ + cancel, err := Run(filepath.Join("testdata", "ta_home"), &config.Config{ Type: "otlp_http", Endpoint: "http://localhost:1337", }) @@ -60,7 +60,7 @@ func TestRunPeriodic(t *testing.T) { defer func() { _ = rcvr.Shutdown(context.Background()) }() - cancel, err := Run(filepath.Join("testdata", "periodic"), &config.Config{ + cancel, err := Run(filepath.Join("testdata", "periodic_home"), &config.Config{ Type: "otlp_http", Endpoint: "http://localhost:1338", }) @@ -101,7 +101,7 @@ func TestRunDisabled(t *testing.T) { defer func() { _ = rcvr.Shutdown(context.Background()) }() - cancel, err := Run(filepath.Join("testdata", "disabled"), &config.Config{ + cancel, err := Run(filepath.Join("testdata", "disabled_home"), &config.Config{ Type: "otlp_http", Endpoint: "http://localhost:1339", }) @@ -123,7 +123,7 @@ func TestRunDisabledInterval(t *testing.T) { defer func() { _ = rcvr.Shutdown(context.Background()) }() - cancel, err := Run(filepath.Join("testdata", "disabled_interval"), &config.Config{ + cancel, err := Run(filepath.Join("testdata", "disabled_interval_home"), &config.Config{ Type: "otlp_http", Endpoint: "http://localhost:1340", }) @@ -144,7 +144,7 @@ func TestRunScriptedInputs(t *testing.T) { defer func() { _ = rcvr.Shutdown(context.Background()) }() - cancel, err := Run(filepath.Join("testdata", "script"), &config.Config{ + cancel, err := Run(filepath.Join("testdata", "script_home"), &config.Config{ Type: "otlp_http", Endpoint: "http://localhost:1341", }) @@ -157,7 +157,7 @@ func TestRunScriptedInputs(t *testing.T) { } func TestUseTCP(t *testing.T) { - rootDir := filepath.Join("testdata", "tcp") + rootDir := filepath.Join("testdata", "tcp_home") logsSink := &consumertest.LogsSink{} cfg := otlpreceiver.NewFactory().CreateDefaultConfig().(*otlpreceiver.Config) cfg.Protocols.HTTP.GetOrInsertDefault().ServerConfig.NetAddr.Endpoint = "localhost:1342" @@ -190,7 +190,7 @@ func TestUseTCP(t *testing.T) { } func TestUseUDP(t *testing.T) { - rootDir := filepath.Join("testdata", "udp") + rootDir := filepath.Join("testdata", "udp_home") logsSink := &consumertest.LogsSink{} cfg := otlpreceiver.NewFactory().CreateDefaultConfig().(*otlpreceiver.Config) cfg.Protocols.HTTP.GetOrInsertDefault().ServerConfig.NetAddr.Endpoint = "localhost:1343" @@ -232,7 +232,7 @@ func TestRunScriptedInputsWithHEC(t *testing.T) { defer func() { _ = rcvr.Shutdown(context.Background()) }() - cancel, err := Run(filepath.Join("testdata", "script"), &config.Config{ + cancel, err := Run(filepath.Join("testdata", "script_home"), &config.Config{ Endpoint: "http://localhost:1341", Token: "foo", }) diff --git a/internal/collector/testdata/disabled/bin/darwin_arm64/foo b/internal/collector/testdata/disabled_home/etc/apps/splunk_ta_disabled/bin/darwin_arm64/foo similarity index 100% rename from internal/collector/testdata/disabled/bin/darwin_arm64/foo rename to internal/collector/testdata/disabled_home/etc/apps/splunk_ta_disabled/bin/darwin_arm64/foo diff --git a/internal/collector/testdata/disabled/bin/linux_amd64/foo b/internal/collector/testdata/disabled_home/etc/apps/splunk_ta_disabled/bin/linux_amd64/foo similarity index 100% rename from internal/collector/testdata/disabled/bin/linux_amd64/foo rename to internal/collector/testdata/disabled_home/etc/apps/splunk_ta_disabled/bin/linux_amd64/foo diff --git a/internal/collector/testdata/disabled/default/inputs.conf b/internal/collector/testdata/disabled_home/etc/apps/splunk_ta_disabled/default/inputs.conf similarity index 100% rename from internal/collector/testdata/disabled/default/inputs.conf rename to internal/collector/testdata/disabled_home/etc/apps/splunk_ta_disabled/default/inputs.conf diff --git a/internal/collector/testdata/disabled_interval/bin/darwin_arm64/foo b/internal/collector/testdata/disabled_interval_home/etc/apps/splunk_ta_disabled_interval/bin/darwin_arm64/foo similarity index 100% rename from internal/collector/testdata/disabled_interval/bin/darwin_arm64/foo rename to internal/collector/testdata/disabled_interval_home/etc/apps/splunk_ta_disabled_interval/bin/darwin_arm64/foo diff --git a/internal/collector/testdata/disabled_interval/bin/linux_amd64/foo b/internal/collector/testdata/disabled_interval_home/etc/apps/splunk_ta_disabled_interval/bin/linux_amd64/foo similarity index 100% rename from internal/collector/testdata/disabled_interval/bin/linux_amd64/foo rename to internal/collector/testdata/disabled_interval_home/etc/apps/splunk_ta_disabled_interval/bin/linux_amd64/foo diff --git a/internal/collector/testdata/disabled_interval/default/inputs.conf b/internal/collector/testdata/disabled_interval_home/etc/apps/splunk_ta_disabled_interval/default/inputs.conf similarity index 100% rename from internal/collector/testdata/disabled_interval/default/inputs.conf rename to internal/collector/testdata/disabled_interval_home/etc/apps/splunk_ta_disabled_interval/default/inputs.conf diff --git a/internal/collector/testdata/periodic/bin/darwin_arm64/foo b/internal/collector/testdata/periodic_home/etc/apps/splunk_ta_periodic/bin/darwin_arm64/foo similarity index 100% rename from internal/collector/testdata/periodic/bin/darwin_arm64/foo rename to internal/collector/testdata/periodic_home/etc/apps/splunk_ta_periodic/bin/darwin_arm64/foo diff --git a/internal/collector/testdata/periodic/bin/linux_amd64/foo b/internal/collector/testdata/periodic_home/etc/apps/splunk_ta_periodic/bin/linux_amd64/foo similarity index 100% rename from internal/collector/testdata/periodic/bin/linux_amd64/foo rename to internal/collector/testdata/periodic_home/etc/apps/splunk_ta_periodic/bin/linux_amd64/foo diff --git a/internal/collector/testdata/periodic/default/inputs.conf b/internal/collector/testdata/periodic_home/etc/apps/splunk_ta_periodic/default/inputs.conf similarity index 100% rename from internal/collector/testdata/periodic/default/inputs.conf rename to internal/collector/testdata/periodic_home/etc/apps/splunk_ta_periodic/default/inputs.conf diff --git a/internal/collector/testdata/script/bin/foo.sh b/internal/collector/testdata/script_home/etc/apps/splunk_ta_script/bin/foo.sh similarity index 100% rename from internal/collector/testdata/script/bin/foo.sh rename to internal/collector/testdata/script_home/etc/apps/splunk_ta_script/bin/foo.sh diff --git a/internal/collector/testdata/script/default/inputs.conf b/internal/collector/testdata/script_home/etc/apps/splunk_ta_script/default/inputs.conf similarity index 100% rename from internal/collector/testdata/script/default/inputs.conf rename to internal/collector/testdata/script_home/etc/apps/splunk_ta_script/default/inputs.conf diff --git a/internal/collector/testdata/ta/bin/darwin_arm64/foo b/internal/collector/testdata/ta_home/etc/apps/splunk_ta_test/bin/darwin_arm64/foo similarity index 100% rename from internal/collector/testdata/ta/bin/darwin_arm64/foo rename to internal/collector/testdata/ta_home/etc/apps/splunk_ta_test/bin/darwin_arm64/foo diff --git a/internal/collector/testdata/ta/bin/linux_amd64/foo b/internal/collector/testdata/ta_home/etc/apps/splunk_ta_test/bin/linux_amd64/foo similarity index 100% rename from internal/collector/testdata/ta/bin/linux_amd64/foo rename to internal/collector/testdata/ta_home/etc/apps/splunk_ta_test/bin/linux_amd64/foo diff --git a/internal/collector/testdata/ta/default/inputs.conf b/internal/collector/testdata/ta_home/etc/apps/splunk_ta_test/default/inputs.conf similarity index 100% rename from internal/collector/testdata/ta/default/inputs.conf rename to internal/collector/testdata/ta_home/etc/apps/splunk_ta_test/default/inputs.conf diff --git a/internal/collector/testdata/tcp/default/inputs.conf b/internal/collector/testdata/tcp_home/etc/apps/splunk_ta_tcp/default/inputs.conf similarity index 100% rename from internal/collector/testdata/tcp/default/inputs.conf rename to internal/collector/testdata/tcp_home/etc/apps/splunk_ta_tcp/default/inputs.conf diff --git a/internal/collector/testdata/udp/default/inputs.conf b/internal/collector/testdata/udp_home/etc/apps/splunk_ta_udp/default/inputs.conf similarity index 100% rename from internal/collector/testdata/udp/default/inputs.conf rename to internal/collector/testdata/udp_home/etc/apps/splunk_ta_udp/default/inputs.conf diff --git a/internal/conf/inputs.go b/internal/conf/inputs.go index dde4ee6..9e3374e 100644 --- a/internal/conf/inputs.go +++ b/internal/conf/inputs.go @@ -20,10 +20,17 @@ type Input struct { ServerURI string `xml:"server_uri"` SessionKey string `xml:"session_key"` CheckpointDir string `xml:"checkpoint_dir"` + AppDir string `xml:"-"` Configuration Configuration `xml:"configuration"` } -func ReadInput(payload []byte) ([]Input, error) { +// IsDisabled reports whether the stanza has disabled=1. +func (s *Stanza) IsDisabled() bool { + p := s.Params.Get("disabled") + return p != nil && p.Value == "1" +} + +func ReadInput(payload []byte, appDir string) ([]Input, error) { f, err := ini.Load(payload) if err != nil { return nil, err @@ -35,6 +42,7 @@ func ReadInput(payload []byte) ([]Input, error) { continue // disregard default section. We need a stanza per input. } i := Input{ + AppDir: appDir, Configuration: Configuration{ Stanza: Stanza{ Name: section.Name(), @@ -58,6 +66,43 @@ func ReadInput(payload []byte) ([]Input, error) { return result, nil } +// MergeInputs merges layered inputs; later layers take precedence per param key. +func MergeInputs(layers [][]Input) []Input { + seen := make(map[string]int) + var result []Input + for _, layer := range layers { + for _, input := range layer { + name := input.Configuration.Stanza.Name + if idx, ok := seen[name]; ok { + result[idx] = mergeInput(result[idx], input) + } else { + seen[name] = len(result) + result = append(result, input) + } + } + } + return result +} + +func mergeInput(base, override Input) Input { + merged := base + params := make(map[string]int, len(base.Configuration.Stanza.Params)) + mergedParams := append([]Param{}, base.Configuration.Stanza.Params...) + for i, p := range mergedParams { + params[p.Name] = i + } + for _, p := range override.Configuration.Stanza.Params { + if idx, ok := params[p.Name]; ok { + mergedParams[idx] = p + } else { + params[p.Name] = len(mergedParams) + mergedParams = append(mergedParams, p) + } + } + merged.Configuration.Stanza.Params = mergedParams + return merged +} + func (i *Input) ToXML() ([]byte, error) { b, err := xml.MarshalIndent(i, "", " ") return append([]byte(xmlDeclaration), b...), err diff --git a/internal/conf/inputs_test.go b/internal/conf/inputs_test.go index 6f16bdb..81c7ef5 100644 --- a/internal/conf/inputs_test.go +++ b/internal/conf/inputs_test.go @@ -15,7 +15,7 @@ import ( func TestOneInput(t *testing.T) { b, err := os.ReadFile(filepath.Join("testdata", "oneinput.conf")) require.NoError(t, err) - res, err := ReadInput(b) + res, err := ReadInput(b, "") require.NoError(t, err) assert.Equal( t, @@ -50,7 +50,7 @@ func TestOneInput(t *testing.T) { func TestTwoInputs(t *testing.T) { b, err := os.ReadFile(filepath.Join("testdata", "twoinputs.conf")) require.NoError(t, err) - res, err := ReadInput(b) + res, err := ReadInput(b, "") require.NoError(t, err) assert.Equal(t, []Input{{ ServerHost: "", @@ -112,6 +112,38 @@ func TestTwoInputs(t *testing.T) { }}, res) } +func TestMergeInputsPartialOverride(t *testing.T) { + base := []Input{{Configuration: Configuration{Stanza: Stanza{ + Name: "monitor:///var/log/syslog", + Params: Params{{Name: "sourcetype", Value: "syslog"}, {Name: "index", Value: "main"}}, + }}}} + override := []Input{{Configuration: Configuration{Stanza: Stanza{ + Name: "monitor:///var/log/syslog", + Params: Params{{Name: "index", Value: "override"}}, + }}}} + + merged := MergeInputs([][]Input{base, override}) + require.Len(t, merged, 1) + assert.Equal(t, "syslog", merged[0].Configuration.Stanza.Params.Get("sourcetype").Value) + assert.Equal(t, "override", merged[0].Configuration.Stanza.Params.Get("index").Value) +} + +func TestMergeInputsFullOverride(t *testing.T) { + base := []Input{{Configuration: Configuration{Stanza: Stanza{ + Name: "monitor:///var/log/syslog", + Params: Params{{Name: "sourcetype", Value: "syslog"}, {Name: "index", Value: "main"}}, + }}}} + override := []Input{{Configuration: Configuration{Stanza: Stanza{ + Name: "monitor:///var/log/syslog", + Params: Params{{Name: "sourcetype", Value: "sourcetype_override"}, {Name: "index", Value: "index_override"}}, + }}}} + + merged := MergeInputs([][]Input{base, override}) + require.Len(t, merged, 1) + assert.Equal(t, "sourcetype_override", merged[0].Configuration.Stanza.Params.Get("sourcetype").Value) + assert.Equal(t, "index_override", merged[0].Configuration.Stanza.Params.Get("index").Value) +} + func TestToXML(t *testing.T) { testStr := ` @@ -133,7 +165,7 @@ func TestToXML(t *testing.T) { ` b, err := os.ReadFile(filepath.Join("testdata", "oneinput.conf")) require.NoError(t, err) - res, err := ReadInput(b) + res, err := ReadInput(b, "") require.NoError(t, err) assert.Len(t, res, 1) res[0].ServerHost = "773c28971b2a" diff --git a/internal/conf/props.go b/internal/conf/props.go index 4598ee2..acf81fd 100644 --- a/internal/conf/props.go +++ b/internal/conf/props.go @@ -60,6 +60,26 @@ func (p *Prop) Type() PropType { } } +// MergeProps merges multiple slices of props, with later slices taking +// precedence. Stanzas are keyed by name; the last definition wins. +// The merged result is re-ordered by specificity. +func MergeProps(layers [][]Prop) []Prop { + seen := make(map[string]int) + var result []Prop + for _, layer := range layers { + for _, p := range layer { + if idx, ok := seen[p.Name]; ok { + result[idx] = p + } else { + seen[p.Name] = len(result) + result = append(result, p) + } + } + } + orderProps(result) + return result +} + func ReadProps(payload []byte) ([]Prop, error) { f, err := ini.Load(payload) if err != nil { diff --git a/internal/conf/transforms.go b/internal/conf/transforms.go index eb10cd5..5d58369 100644 --- a/internal/conf/transforms.go +++ b/internal/conf/transforms.go @@ -11,6 +11,24 @@ type Transform struct { Format string } +// MergeTransforms merges multiple slices of transforms, with later slices +// taking precedence. Stanzas are keyed by name; the last definition wins. +func MergeTransforms(layers [][]Transform) []Transform { + seen := make(map[string]int) + var result []Transform + for _, layer := range layers { + for _, t := range layer { + if idx, ok := seen[t.Name]; ok { + result[idx] = t + } else { + seen[t.Name] = len(result) + result = append(result, t) + } + } + } + return result +} + func ReadTransforms(payload []byte) ([]Transform, error) { f, err := ini.Load(payload) if err != nil { diff --git a/internal/receiver/batchreceiver/config.go b/internal/receiver/batchreceiver/config.go index 159d714..b7ba1ee 100644 --- a/internal/receiver/batchreceiver/config.go +++ b/internal/receiver/batchreceiver/config.go @@ -8,9 +8,8 @@ import ( ) type Config struct { + Input conf.Input `mapstructure:"-"` + BaseDir string `mapstructure:"-"` Transforms []conf.Transform `mapstructure:"-"` Props []conf.Prop `mapstructure:"-"` - - BaseDir string `mapstructure:"-"` - Input conf.Input `mapstructure:"-"` } diff --git a/internal/receiver/monitorreceiver/config.go b/internal/receiver/monitorreceiver/config.go index 31d37a8..855f4ad 100644 --- a/internal/receiver/monitorreceiver/config.go +++ b/internal/receiver/monitorreceiver/config.go @@ -8,9 +8,8 @@ import ( ) type Config struct { + Input conf.Input `mapstructure:"-"` + BaseDir string `mapstructure:"-"` Transforms []conf.Transform `mapstructure:"-"` Props []conf.Prop `mapstructure:"-"` - - BaseDir string `mapstructure:"-"` - Input conf.Input `mapstructure:"-"` } diff --git a/internal/receiver/scriptreceiver/config.go b/internal/receiver/scriptreceiver/config.go index 38889cc..7ce0228 100644 --- a/internal/receiver/scriptreceiver/config.go +++ b/internal/receiver/scriptreceiver/config.go @@ -6,8 +6,8 @@ package scriptreceiver import "github.com/splunk/tarunner/internal/conf" type Config struct { + conf.Input `mapstructure:"-"` BaseDir string `mapstructure:"-"` Props []conf.Prop `mapstructure:"-"` Transforms []conf.Transform `mapstructure:"-"` - conf.Input `mapstructure:"-"` } diff --git a/internal/receiver/tcpreceiver/config.go b/internal/receiver/tcpreceiver/config.go index 4aa7592..08df79d 100644 --- a/internal/receiver/tcpreceiver/config.go +++ b/internal/receiver/tcpreceiver/config.go @@ -9,11 +9,10 @@ import ( ) type Config struct { + Input conf.Input `mapstructure:"-"` + BaseDir string `mapstructure:"-"` Transforms []conf.Transform `mapstructure:"-"` Props []conf.Prop `mapstructure:"-"` - - BaseDir string `mapstructure:"-"` - Input conf.Input `mapstructure:"-"` } func (cfg *Config) Validate() error { diff --git a/internal/receiver/udpreceiver/config.go b/internal/receiver/udpreceiver/config.go index ef33801..f180647 100644 --- a/internal/receiver/udpreceiver/config.go +++ b/internal/receiver/udpreceiver/config.go @@ -9,11 +9,10 @@ import ( ) type Config struct { + Input conf.Input `mapstructure:"-"` + BaseDir string `mapstructure:"-"` Transforms []conf.Transform `mapstructure:"-"` Props []conf.Prop `mapstructure:"-"` - - BaseDir string `mapstructure:"-"` - Input conf.Input `mapstructure:"-"` } func (cfg *Config) Validate() error { diff --git a/internal/receiver/wineventlogreceiver/config.go b/internal/receiver/wineventlogreceiver/config.go index 80a6fc1..6c6b15a 100644 --- a/internal/receiver/wineventlogreceiver/config.go +++ b/internal/receiver/wineventlogreceiver/config.go @@ -8,9 +8,8 @@ import ( ) type Config struct { + Input conf.Input `mapstructure:"-"` + BaseDir string `mapstructure:"-"` Transforms []conf.Transform `mapstructure:"-"` Props []conf.Prop `mapstructure:"-"` - - BaseDir string `mapstructure:"-"` - Input conf.Input `mapstructure:"-"` } diff --git a/internal/script/command.go b/internal/script/command.go index 2b85487..2161f62 100644 --- a/internal/script/command.go +++ b/internal/script/command.go @@ -18,13 +18,23 @@ func DetermineCommandName(baseDir string, input conf.Input) (string, error) { if err != nil { return "", err } + resolveDir := baseDir + if input.AppDir != "" { + resolveDir = input.AppDir + } switch parsed.Kind { case "monitor", "batch": return parsed.Target, nil case "script": - return GetPath(baseDir, parsed.Target) + if filepath.IsAbs(parsed.Target) { + return parsed.Target, nil + } + return GetPath(resolveDir, parsed.Target) case "": - return GetPath(baseDir, filepath.Join("bin", fmt.Sprintf("%s_%s", runtime.GOOS, runtime.GOARCH), parsed.Target)) + if filepath.IsAbs(parsed.Target) { + return parsed.Target, nil + } + return GetPath(resolveDir, filepath.Join("bin", fmt.Sprintf("%s_%s", runtime.GOOS, runtime.GOARCH), parsed.Target)) default: return "", fmt.Errorf("unknown scheme %q", parsed.Kind) } diff --git a/internal/tabuilder/tabuilder.go b/internal/tabuilder/tabuilder.go index bf73eec..c202c65 100644 --- a/internal/tabuilder/tabuilder.go +++ b/internal/tabuilder/tabuilder.go @@ -14,7 +14,7 @@ import ( "fmt" "os" "path/filepath" - "sort" + "strings" "github.com/open-telemetry/opentelemetry-collector-contrib/exporter/splunkhecexporter" "go.opentelemetry.io/collector/component" @@ -38,26 +38,40 @@ import ( "github.com/splunk/tarunner/internal/stanza" ) +// ResolveSplunkHome returns baseDir if set, otherwise falls back to $SPLUNK_HOME. +func ResolveSplunkHome(baseDir string) (string, error) { + if baseDir != "" { + return baseDir, nil + } + if home := os.Getenv("SPLUNK_HOME"); home != "" { + return home, nil + } + return "", fmt.Errorf("base_dir is not set and $SPLUNK_HOME is not defined") +} + // CreateReceivers builds a logs receiver for every enabled input stanza, -// dispatching by the stanza name's input kind. Stanzas with disabled=1 are -// skipped. +// dispatching by the stanza name's input kind. Stanzas with disabled=1 or +// unsupported kinds are skipped silently. func CreateReceivers(ctx context.Context, inputs []conf.Input, transforms []conf.Transform, props []conf.Prop, baseDir string, next consumer.Logs, telemetrySettings component.TelemetrySettings) ([]receiver.Logs, error) { var receivers []receiver.Logs for _, input := range inputs { - disabled := input.Configuration.Stanza.Params.Get("disabled") - if disabled != nil && disabled.Value == "1" { + if input.Configuration.Stanza.IsDisabled() { continue } l, err := CreateReceiver(ctx, baseDir, next, input, transforms, props, telemetrySettings) if err != nil { return nil, fmt.Errorf("failed to create receiver %q: %w", input.Configuration.Stanza.Name, err) } + if l == nil { + continue + } receivers = append(receivers, l) } return receivers, nil } // CreateReceiver builds a single logs receiver for an input stanza. +// Returns nil receiver for unsupported input kinds; the caller is responsible for logging. func CreateReceiver(ctx context.Context, baseDir string, next consumer.Logs, input conf.Input, transforms []conf.Transform, props []conf.Prop, telemetrySettings component.TelemetrySettings) (receiver.Logs, error) { parsed, err := stanza.ParseName(input.Configuration.Stanza.Name) if err != nil { @@ -113,7 +127,7 @@ func CreateReceiver(ctx context.Context, baseDir string, next consumer.Logs, inp Props: props, }, next) default: - return nil, fmt.Errorf("unsupported scheme %q", parsed.Kind) + return nil, nil } } @@ -132,24 +146,54 @@ func confFilePaths(dirs []string, filename string) []string { return paths } +// DiscoverTAs returns splunk_ta_* directories under splunkHome/etc/apps. +func DiscoverTAs(splunkHome string) ([]string, error) { + appsDir := filepath.Join(splunkHome, "etc", "apps") + entries, err := os.ReadDir(appsDir) + if err != nil { + if errors.Is(err, os.ErrNotExist) { + return nil, nil + } + return nil, fmt.Errorf("tabuilder: failed to scan %s: %w", appsDir, err) + } + var taDirs []string + for _, entry := range entries { + if !entry.IsDir() { + continue + } + if strings.HasPrefix(strings.ToLower(entry.Name()), "splunk_ta_") { + taDirs = append(taDirs, filepath.Join(appsDir, entry.Name())) + } + } + return taDirs, nil +} + func splunkHomeDirs(splunkHome string) []string { + taDirs, _ := DiscoverTAs(splunkHome) etcDir := filepath.Join(splunkHome, "etc") - appDirs, _ := filepath.Glob(filepath.Join(etcDir, "apps", "*")) - sort.Strings(appDirs) - dirs := []string{filepath.Join(etcDir, "system", "default")} - for _, app := range appDirs { - dirs = append(dirs, filepath.Join(app, "default")) + for _, ta := range taDirs { + dirs = append(dirs, filepath.Join(ta, "default")) } - for _, app := range appDirs { - dirs = append(dirs, filepath.Join(app, "local")) + for _, ta := range taDirs { + dirs = append(dirs, filepath.Join(ta, "local")) } dirs = append(dirs, filepath.Join(etcDir, "system", "local")) return dirs } +func taDirsWithSystem(splunkHome, taDir string) []string { + etcDir := filepath.Join(splunkHome, "etc") + return []string{ + filepath.Join(etcDir, "system", "default"), + filepath.Join(taDir, "default"), + filepath.Join(taDir, "local"), + filepath.Join(etcDir, "system", "local"), + } +} + func readConfFiles(paths []string) ([][]byte, error) { var payloads [][]byte for _, path := range paths { @@ -165,60 +209,79 @@ func readConfFiles(paths []string) ([][]byte, error) { return payloads, nil } -// ReadInputs reads inputs.conf, preferring local/ over default/. -func ReadInputs(baseDir string) ([]conf.Input, error) { - fileToRead := filepath.Join(baseDir, "local", "inputs.conf") - if _, err := os.Stat(fileToRead); errors.Is(err, os.ErrNotExist) { - fileToRead = filepath.Join(baseDir, "default", "inputs.conf") - if _, err := os.Stat(fileToRead); errors.Is(err, os.ErrNotExist) { +// ConfDirs returns the Splunk btool conf search path for splunkHome. +func ConfDirs(splunkHome string) []string { + return splunkHomeDirs(splunkHome) +} + +// ConfDirsWithSystem returns the Splunk btool conf search path for a single TA merged with system config. +func ConfDirsWithSystem(splunkHome, taDir string) []string { + return taDirsWithSystem(splunkHome, taDir) +} + +// ReadInputs discovers and merges inputs.conf files from the given search +// directories. Use ConfDirs or ConfDirsWithSystem to build the dirs slice. +// Returns nil (no error) when no inputs.conf files are found. +func ReadInputs(dirs []string) ([]conf.Input, error) { + var layers [][]conf.Input + for _, path := range confFilePaths(dirs, "inputs.conf") { + b, err := os.ReadFile(path) + if errors.Is(err, os.ErrNotExist) { + continue + } + if err != nil { return nil, err } + appDir := filepath.Dir(filepath.Dir(path)) // strip /default or /local + inputs, err := conf.ReadInput(b, appDir) + if err != nil { + return nil, fmt.Errorf("parse %s: %w", path, err) + } + layers = append(layers, inputs) } - b, err := os.ReadFile(fileToRead) - if err != nil { - return nil, err - } - return conf.ReadInput(b) + return conf.MergeInputs(layers), nil } -// ReadTransforms reads transforms.conf, preferring local/ over default/. It -// returns a nil slice (and no error) when the file is absent. -func ReadTransforms(baseDir string) ([]conf.Transform, error) { - fileToRead := filepath.Join(baseDir, "local", "transforms.conf") - if _, err := os.Stat(fileToRead); errors.Is(err, os.ErrNotExist) { - fileToRead = filepath.Join(baseDir, "default", "transforms.conf") - if _, err := os.Stat(fileToRead); errors.Is(err, os.ErrNotExist) { - return nil, nil - } - } - b, err := os.ReadFile(fileToRead) +// ReadTransforms discovers and merges transforms.conf files from the given +// search directories. Returns nil (no error) when absent. +func ReadTransforms(dirs []string) ([]conf.Transform, error) { + payloads, err := readConfFiles(confFilePaths(dirs, "transforms.conf")) if err != nil { return nil, err } - return conf.ReadTransforms(b) -} - -// ReadProps reads props.conf, preferring local/ over default/. It returns a nil -// slice (and no error) when the file is absent. -func ReadProps(baseDir string) ([]conf.Prop, error) { - fileToRead := filepath.Join(baseDir, "local", "props.conf") - if _, err := os.Stat(fileToRead); errors.Is(err, os.ErrNotExist) { - fileToRead = filepath.Join(baseDir, "default", "props.conf") - if _, err := os.Stat(fileToRead); errors.Is(err, os.ErrNotExist) { - return nil, nil + var layers [][]conf.Transform + for _, b := range payloads { + transforms, err := conf.ReadTransforms(b) + if err != nil { + return nil, err } + layers = append(layers, transforms) } - b, err := os.ReadFile(fileToRead) + return conf.MergeTransforms(layers), nil +} + +// ReadProps discovers and merges props.conf files from the given search +// directories. Returns nil (no error) when absent. +func ReadProps(dirs []string) ([]conf.Prop, error) { + payloads, err := readConfFiles(confFilePaths(dirs, "props.conf")) if err != nil { return nil, err } - return conf.ReadProps(b) + var layers [][]conf.Prop + for _, b := range payloads { + props, err := conf.ReadProps(b) + if err != nil { + return nil, err + } + layers = append(layers, props) + } + return conf.MergeProps(layers), nil } // ReadOutputs merges outputs.conf across $SPLUNK_HOME using standard Splunk // precedence. Use HTTPOut (or future TCPOut, etc.) to extract a specific type. func ReadOutputs(splunkHome string) (conf.ConfMap, error) { - payloads, err := readConfFiles(confFilePaths(splunkHomeDirs(splunkHome), "outputs.conf")) + payloads, err := readConfFiles(confFilePaths(ConfDirs(splunkHome), "outputs.conf")) if err != nil { return nil, err } @@ -249,13 +312,14 @@ func CreateExporter(merged conf.ConfMap, logger *zap.Logger, telemetrySettings c } // CreateOutputExporter builds a logs exporter from one built-in output stanza. +// Returns nil exporter for unsupported output kinds; the caller is responsible for logging. func CreateOutputExporter(output *conf.Output, logger *zap.Logger, telemetrySettings component.TelemetrySettings) (exporter.Logs, error) { parsed, err := stanza.ParseOutputName(output.Configuration.Stanza.Name) if err != nil { return nil, err } if parsed.Kind != "httpout" { - return nil, fmt.Errorf("unsupported scheme %q", parsed.Kind) + return nil, nil } return newHECExporter(output, logger, telemetrySettings) } diff --git a/internal/tabuilder/tabuilder_test.go b/internal/tabuilder/tabuilder_test.go index ef56bf0..8744170 100644 --- a/internal/tabuilder/tabuilder_test.go +++ b/internal/tabuilder/tabuilder_test.go @@ -92,33 +92,33 @@ func TestReadTransforms(t *testing.T) { rootDir := filepath.Join("testdata", "transforms") tests := []struct { name string - path string + splunkHome string expectedName string expectedRegex string }{ { name: "default", - path: filepath.Join(rootDir, "default"), + splunkHome: filepath.Join(rootDir, "splunk_default"), expectedName: "example_default", expectedRegex: "default", }, { name: "local", - path: filepath.Join(rootDir, "local"), + splunkHome: filepath.Join(rootDir, "splunk_local"), expectedName: "example_local", expectedRegex: "local", }, { name: "both", - path: filepath.Join(rootDir, "both"), - expectedName: "example_local2", + splunkHome: filepath.Join(rootDir, "splunk_both"), + expectedName: "example_transform", expectedRegex: "local", }, } for _, test := range tests { t.Run(test.name, func(t *testing.T) { - transforms, err := ReadTransforms(test.path) + transforms, err := ReadTransforms(ConfDirs(test.splunkHome)) require.NoError(t, err) require.Len(t, transforms, 1) require.Equal(t, test.expectedName, transforms[0].Name) diff --git a/internal/tabuilder/testdata/outputs/app_local_wins_over_app_default/etc/apps/my_ta/default/outputs.conf b/internal/tabuilder/testdata/outputs/app_local_wins_over_app_default/etc/apps/splunk_ta_test/default/outputs.conf similarity index 100% rename from internal/tabuilder/testdata/outputs/app_local_wins_over_app_default/etc/apps/my_ta/default/outputs.conf rename to internal/tabuilder/testdata/outputs/app_local_wins_over_app_default/etc/apps/splunk_ta_test/default/outputs.conf diff --git a/internal/tabuilder/testdata/outputs/app_local_wins_over_app_default/etc/apps/my_ta/local/outputs.conf b/internal/tabuilder/testdata/outputs/app_local_wins_over_app_default/etc/apps/splunk_ta_test/local/outputs.conf similarity index 100% rename from internal/tabuilder/testdata/outputs/app_local_wins_over_app_default/etc/apps/my_ta/local/outputs.conf rename to internal/tabuilder/testdata/outputs/app_local_wins_over_app_default/etc/apps/splunk_ta_test/local/outputs.conf diff --git a/internal/tabuilder/testdata/outputs/app_overrides/etc/apps/my_ta/default/outputs.conf b/internal/tabuilder/testdata/outputs/app_overrides/etc/apps/splunk_ta_test/default/outputs.conf similarity index 100% rename from internal/tabuilder/testdata/outputs/app_overrides/etc/apps/my_ta/default/outputs.conf rename to internal/tabuilder/testdata/outputs/app_overrides/etc/apps/splunk_ta_test/default/outputs.conf diff --git a/internal/tabuilder/testdata/outputs/app_overrides/etc/apps/my_ta/local/outputs.conf b/internal/tabuilder/testdata/outputs/app_overrides/etc/apps/splunk_ta_test/local/outputs.conf similarity index 100% rename from internal/tabuilder/testdata/outputs/app_overrides/etc/apps/my_ta/local/outputs.conf rename to internal/tabuilder/testdata/outputs/app_overrides/etc/apps/splunk_ta_test/local/outputs.conf diff --git a/internal/tabuilder/testdata/outputs/system_local_wins/etc/apps/my_ta/local/outputs.conf b/internal/tabuilder/testdata/outputs/system_local_wins/etc/apps/splunk_ta_test/local/outputs.conf similarity index 100% rename from internal/tabuilder/testdata/outputs/system_local_wins/etc/apps/my_ta/local/outputs.conf rename to internal/tabuilder/testdata/outputs/system_local_wins/etc/apps/splunk_ta_test/local/outputs.conf diff --git a/internal/tabuilder/testdata/transforms/splunk_both/etc/apps/splunk_ta_test/default/transforms.conf b/internal/tabuilder/testdata/transforms/splunk_both/etc/apps/splunk_ta_test/default/transforms.conf new file mode 100644 index 0000000..eb3772d --- /dev/null +++ b/internal/tabuilder/testdata/transforms/splunk_both/etc/apps/splunk_ta_test/default/transforms.conf @@ -0,0 +1,2 @@ +[example_transform] +REGEX=default diff --git a/internal/tabuilder/testdata/transforms/splunk_both/etc/apps/splunk_ta_test/local/transforms.conf b/internal/tabuilder/testdata/transforms/splunk_both/etc/apps/splunk_ta_test/local/transforms.conf new file mode 100644 index 0000000..7813599 --- /dev/null +++ b/internal/tabuilder/testdata/transforms/splunk_both/etc/apps/splunk_ta_test/local/transforms.conf @@ -0,0 +1,2 @@ +[example_transform] +REGEX=local \ No newline at end of file diff --git a/internal/tabuilder/testdata/transforms/splunk_default/etc/apps/splunk_ta_test/default/transforms.conf b/internal/tabuilder/testdata/transforms/splunk_default/etc/apps/splunk_ta_test/default/transforms.conf new file mode 100644 index 0000000..51c987f --- /dev/null +++ b/internal/tabuilder/testdata/transforms/splunk_default/etc/apps/splunk_ta_test/default/transforms.conf @@ -0,0 +1,2 @@ +[example_default] +REGEX=default \ No newline at end of file diff --git a/internal/tabuilder/testdata/transforms/splunk_local/etc/apps/splunk_ta_test/local/transforms.conf b/internal/tabuilder/testdata/transforms/splunk_local/etc/apps/splunk_ta_test/local/transforms.conf new file mode 100644 index 0000000..071044b --- /dev/null +++ b/internal/tabuilder/testdata/transforms/splunk_local/etc/apps/splunk_ta_test/local/transforms.conf @@ -0,0 +1,2 @@ +[example_local] +REGEX=local \ No newline at end of file diff --git a/pkg/splunkinputsreceiver/config.go b/pkg/splunkinputsreceiver/config.go index 50fe229..9d9b9cb 100644 --- a/pkg/splunkinputsreceiver/config.go +++ b/pkg/splunkinputsreceiver/config.go @@ -4,5 +4,9 @@ package splunkinputsreceiver type Config struct { - BaseDir string `mapstructure:"path"` + // BaseDir is the Splunk installation root ($SPLUNK_HOME). TAs are + // discovered under etc/apps/* and conf files are layered using standard + // Splunk precedence. Falls back to the $SPLUNK_HOME environment variable + // when not set. + BaseDir string `mapstructure:"base_dir"` } diff --git a/pkg/splunkinputsreceiver/factory_test.go b/pkg/splunkinputsreceiver/factory_test.go index 910f481..dc9e01c 100644 --- a/pkg/splunkinputsreceiver/factory_test.go +++ b/pkg/splunkinputsreceiver/factory_test.go @@ -14,79 +14,81 @@ import ( "go.opentelemetry.io/collector/consumer" "go.opentelemetry.io/collector/pdata/plog" "go.opentelemetry.io/collector/receiver" + "go.uber.org/zap" "github.com/splunk/tarunner/pkg/splunkinputsreceiver" ) func TestWithSubReceiverRegistersCustomScheme(t *testing.T) { - baseDir := writeTA(t, "[custom:///thing]\nsourcetype = custom\n") + splunkHome := writeTA(t, "[custom:///thing]\nsourcetype = custom\n") fake := &fakeSubReceiverFactory{scheme: "custom"} factory := splunkinputsreceiver.NewFactory(splunkinputsreceiver.WithSubReceiver(fake)) - rcvr, err := factory.CreateLogs(context.Background(), newReceiverSettings(), splunkinputsreceiver.Config{BaseDir: baseDir}, nopConsumer{}) + rcvr, err := factory.CreateLogs(context.Background(), newReceiverSettings(), splunkinputsreceiver.Config{BaseDir: splunkHome}, nopConsumer{}) require.NoError(t, err) require.NotNil(t, rcvr) require.True(t, fake.called) - require.Equal(t, baseDir, fake.request.BaseDir) + require.Equal(t, filepath.Join(splunkHome, "etc", "apps", "Splunk_TA_test"), fake.request.BaseDir) require.Equal(t, "/thing", fake.request.Path) require.Equal(t, "custom:///thing", fake.request.Input.Configuration.Stanza.Name) } func TestWithSubReceiverOverridesBuiltInAndHandlesEmptySchemeAsScript(t *testing.T) { - baseDir := writeTA(t, "[modinput]\ninterval = -1\n") + splunkHome := writeTA(t, "[modinput]\ninterval = -1\n") fake := &fakeSubReceiverFactory{scheme: "script"} factory := splunkinputsreceiver.NewFactory(splunkinputsreceiver.WithSubReceiver(fake)) - rcvr, err := factory.CreateLogs(context.Background(), newReceiverSettings(), splunkinputsreceiver.Config{BaseDir: baseDir}, nopConsumer{}) + rcvr, err := factory.CreateLogs(context.Background(), newReceiverSettings(), splunkinputsreceiver.Config{BaseDir: splunkHome}, nopConsumer{}) require.NoError(t, err) require.NotNil(t, rcvr) require.True(t, fake.called) - require.Equal(t, "modinput", fake.request.Path) + require.Equal(t, filepath.Join(splunkHome, "etc", "apps", "Splunk_TA_test"), fake.request.BaseDir) require.Equal(t, "modinput", fake.request.Input.Configuration.Stanza.Name) } func TestWithSubReceiverMatchesNormalizedScheme(t *testing.T) { - baseDir := writeTA(t, "[Custom://thing]\n") + splunkHome := writeTA(t, "[Custom://thing]\n") fake := &fakeSubReceiverFactory{scheme: "custom"} factory := splunkinputsreceiver.NewFactory(splunkinputsreceiver.WithSubReceiver(fake)) - rcvr, err := factory.CreateLogs(context.Background(), newReceiverSettings(), splunkinputsreceiver.Config{BaseDir: baseDir}, nopConsumer{}) + rcvr, err := factory.CreateLogs(context.Background(), newReceiverSettings(), splunkinputsreceiver.Config{BaseDir: splunkHome}, nopConsumer{}) require.NoError(t, err) require.NotNil(t, rcvr) require.True(t, fake.called) + require.Equal(t, filepath.Join(splunkHome, "etc", "apps", "Splunk_TA_test"), fake.request.BaseDir) require.Equal(t, "thing", fake.request.Path) require.Equal(t, "Custom://thing", fake.request.Input.Configuration.Stanza.Name) } func TestWithSubReceiverRejectsUnsupportedScheme(t *testing.T) { - baseDir := writeTA(t, "[unsupported://thing]\n") + splunkHome := writeTA(t, "[unsupported://thing]\n") factory := splunkinputsreceiver.NewFactory() - rcvr, err := factory.CreateLogs(context.Background(), newReceiverSettings(), splunkinputsreceiver.Config{BaseDir: baseDir}, nopConsumer{}) + rcvr, err := factory.CreateLogs(context.Background(), newReceiverSettings(), splunkinputsreceiver.Config{BaseDir: splunkHome}, nopConsumer{}) require.ErrorContains(t, err, `unsupported scheme "unsupported"`) require.Nil(t, rcvr) } func TestWithSubReceiverSkipsDisabledCustomStanza(t *testing.T) { - baseDir := writeTA(t, "[custom:///thing]\ndisabled = 1\n") + splunkHome := writeTA(t, "[custom:///thing]\ndisabled = 1\n") fake := &fakeSubReceiverFactory{scheme: "custom"} factory := splunkinputsreceiver.NewFactory(splunkinputsreceiver.WithSubReceiver(fake)) - rcvr, err := factory.CreateLogs(context.Background(), newReceiverSettings(), splunkinputsreceiver.Config{BaseDir: baseDir}, nopConsumer{}) + rcvr, err := factory.CreateLogs(context.Background(), newReceiverSettings(), splunkinputsreceiver.Config{BaseDir: splunkHome}, nopConsumer{}) require.NoError(t, err) require.NotNil(t, rcvr) require.False(t, fake.called) } func TestWithSubReceiverRequestIncludesPropsAndTransforms(t *testing.T) { - baseDir := writeTA(t, "[custom:///thing]\nsourcetype = custom\n") - defaultDir := filepath.Join(baseDir, "default") - require.NoError(t, os.WriteFile(filepath.Join(defaultDir, "props.conf"), []byte("[custom]\nTRANSFORMS-routing = route\n"), 0o600)) - require.NoError(t, os.WriteFile(filepath.Join(defaultDir, "transforms.conf"), []byte("[route]\nREGEX = ^(.*)$\nFORMAT = $1\n"), 0o600)) + splunkHome := writeTA(t, "[custom:///thing]\nsourcetype = custom\n") + taDefaultDir := filepath.Join(splunkHome, "etc", "apps", "Splunk_TA_test", "default") + require.NoError(t, os.WriteFile(filepath.Join(taDefaultDir, "props.conf"), []byte("[custom]\nTRANSFORMS-routing = route\n"), 0o600)) + require.NoError(t, os.WriteFile(filepath.Join(taDefaultDir, "transforms.conf"), []byte("[route]\nREGEX = ^(.*)$\nFORMAT = $1\n"), 0o600)) fake := &fakeSubReceiverFactory{scheme: "custom"} factory := splunkinputsreceiver.NewFactory(splunkinputsreceiver.WithSubReceiver(fake)) - rcvr, err := factory.CreateLogs(context.Background(), newReceiverSettings(), splunkinputsreceiver.Config{BaseDir: baseDir}, nopConsumer{}) + rcvr, err := factory.CreateLogs(context.Background(), newReceiverSettings(), splunkinputsreceiver.Config{BaseDir: splunkHome}, nopConsumer{}) require.NoError(t, err) require.NotNil(t, rcvr) require.True(t, fake.called) @@ -137,15 +139,18 @@ func (nopConsumer) ConsumeLogs(context.Context, plog.Logs) error { func newReceiverSettings() receiver.Settings { return receiver.Settings{ - ID: component.MustNewID("splunk_inputs"), + ID: component.MustNewID("splunk_inputs"), + TelemetrySettings: component.TelemetrySettings{Logger: zap.NewNop()}, } } +// writeTA creates a minimal Splunk home layout with one TA and returns the +// splunk home path. The TA is placed at /etc/apps/ta/. func writeTA(t *testing.T, inputsConf string) string { t.Helper() - baseDir := t.TempDir() - defaultDir := filepath.Join(baseDir, "default") - require.NoError(t, os.Mkdir(defaultDir, 0o755)) - require.NoError(t, os.WriteFile(filepath.Join(defaultDir, "inputs.conf"), []byte(inputsConf), 0o600)) - return baseDir + splunkHome := t.TempDir() + taDefaultDir := filepath.Join(splunkHome, "etc", "apps", "Splunk_TA_test", "default") + require.NoError(t, os.MkdirAll(taDefaultDir, 0o755)) + require.NoError(t, os.WriteFile(filepath.Join(taDefaultDir, "inputs.conf"), []byte(inputsConf), 0o600)) + return splunkHome } diff --git a/pkg/splunkinputsreceiver/go.mod b/pkg/splunkinputsreceiver/go.mod index 2b1fb9e..45d455b 100644 --- a/pkg/splunkinputsreceiver/go.mod +++ b/pkg/splunkinputsreceiver/go.mod @@ -9,6 +9,7 @@ require ( go.opentelemetry.io/collector/consumer v1.63.0 go.opentelemetry.io/collector/pdata v1.63.0 go.opentelemetry.io/collector/receiver v1.63.0 + go.uber.org/zap v1.28.0 ) require ( @@ -97,7 +98,6 @@ require ( go.opentelemetry.io/otel/metric v1.44.0 // indirect go.opentelemetry.io/otel/trace v1.44.0 // indirect go.uber.org/multierr v1.11.0 // indirect - go.uber.org/zap v1.28.0 // indirect go.yaml.in/yaml/v3 v3.0.4 // indirect golang.org/x/crypto v0.54.0 // indirect golang.org/x/net v0.57.0 // indirect diff --git a/pkg/splunkinputsreceiver/subreceiver.go b/pkg/splunkinputsreceiver/subreceiver.go index cc60a6d..b4e6bd4 100644 --- a/pkg/splunkinputsreceiver/subreceiver.go +++ b/pkg/splunkinputsreceiver/subreceiver.go @@ -11,6 +11,7 @@ import ( "go.opentelemetry.io/collector/component" "go.opentelemetry.io/collector/consumer" "go.opentelemetry.io/collector/receiver" + "go.uber.org/zap" "github.com/splunk/tarunner/internal/conf" "github.com/splunk/tarunner/internal/stanza" @@ -83,37 +84,56 @@ func newFactoryOptions(opts ...Option) factoryOptions { func (o factoryOptions) createLogsFunc(ctx context.Context, settings receiver.Settings, config component.Config, logs consumer.Logs) (receiver.Logs, error) { cfg := config.(Config) - baseDir := cfg.BaseDir - inputs, err := tabuilder.ReadInputs(baseDir) - if err != nil { - return nil, err - } - transforms, err := tabuilder.ReadTransforms(baseDir) + + splunkHome, err := tabuilder.ResolveSplunkHome(cfg.BaseDir) if err != nil { - return nil, err + return nil, fmt.Errorf("splunk_inputs: %w", err) } - props, err := tabuilder.ReadProps(baseDir) + + taDirs, err := tabuilder.DiscoverTAs(splunkHome) if err != nil { return nil, err } - receivers, err := o.createReceivers(ctx, inputs, transforms, props, baseDir, logs, settings) - if err != nil { - return nil, err + var allReceivers []receiver.Logs + for _, taDir := range taDirs { + dirs := tabuilder.ConfDirsWithSystem(splunkHome, taDir) + inputs, err := tabuilder.ReadInputs(dirs) + if err != nil { + return nil, err + } + transforms, err := tabuilder.ReadTransforms(dirs) + if err != nil { + return nil, err + } + props, err := tabuilder.ReadProps(dirs) + if err != nil { + return nil, err + } + receivers, err := o.createReceivers(ctx, inputs, transforms, props, taDir, logs, settings) + if err != nil { + return nil, err + } + allReceivers = append(allReceivers, receivers...) } - return packReceivers(receivers), nil + return packReceivers(allReceivers), nil } func (o factoryOptions) createReceivers(ctx context.Context, inputs []Input, transforms []Transform, props []Prop, baseDir string, next consumer.Logs, settings receiver.Settings) ([]receiver.Logs, error) { var receivers []receiver.Logs for _, input := range inputs { - disabled := input.Configuration.Stanza.Params.Get("disabled") - if disabled != nil && disabled.Value == "1" { + name := input.Configuration.Stanza.Name + if input.Configuration.Stanza.IsDisabled() { + settings.Logger.Info("splunk_inputs: skipping disabled stanza", zap.String("stanza", name)) continue } l, err := o.createReceiver(ctx, baseDir, next, input, transforms, props, settings) if err != nil { - return nil, fmt.Errorf("failed to create receiver %q: %w", input.Configuration.Stanza.Name, err) + return nil, fmt.Errorf("failed to create receiver %q: %w", name, err) + } + if l == nil { + settings.Logger.Info("splunk_inputs: skipping unsupported input stanza", zap.String("stanza", name)) + continue } receivers = append(receivers, l) } @@ -138,5 +158,9 @@ func (o factoryOptions) createReceiver(ctx context.Context, baseDir string, next Props: props, }, next) } - return tabuilder.CreateReceiver(ctx, baseDir, next, input, transforms, props, settings.TelemetrySettings) + l, err := tabuilder.CreateReceiver(ctx, baseDir, next, input, transforms, props, settings.TelemetrySettings) + if l == nil && err == nil { + return nil, fmt.Errorf("unsupported scheme %q", scheme) + } + return l, err } diff --git a/pkg/splunkoutputsexporter/factory_test.go b/pkg/splunkoutputsexporter/factory_test.go index 9abaf94..fbe7da9 100644 --- a/pkg/splunkoutputsexporter/factory_test.go +++ b/pkg/splunkoutputsexporter/factory_test.go @@ -14,6 +14,7 @@ import ( "go.opentelemetry.io/collector/consumer" "go.opentelemetry.io/collector/exporter" "go.opentelemetry.io/collector/pdata/plog" + "go.uber.org/zap" "github.com/splunk/tarunner/pkg/splunkoutputsexporter" ) @@ -71,13 +72,13 @@ func TestWithSubExporterRegistersCustomScheme(t *testing.T) { require.Equal(t, "splunk:9997", fake.requests[0].Output.Configuration.Stanza.Params.Get("server").Value) } -func TestWithSubExporterRejectsUnsupportedScheme(t *testing.T) { +func TestWithSubExporterSkipsUnsupportedScheme(t *testing.T) { baseDir := writeTA(t, "[tcpout:primary]\nserver = splunk:9997\n") factory := splunkoutputsexporter.NewFactory() exp, err := factory.CreateLogs(context.Background(), newExporterSettings(), splunkoutputsexporter.Config{BaseDir: baseDir}) - require.ErrorContains(t, err, `unsupported scheme "tcpout"`) - require.Nil(t, exp) + require.NoError(t, err) + require.NotNil(t, exp) } type fakeSubExporterFactory struct { @@ -114,7 +115,8 @@ func (fakeLogsExporter) ConsumeLogs(context.Context, plog.Logs) error { func newExporterSettings() exporter.Settings { return exporter.Settings{ - ID: component.MustNewID("splunk_outputs"), + ID: component.MustNewID("splunk_outputs"), + TelemetrySettings: component.TelemetrySettings{Logger: zap.NewNop()}, } } diff --git a/pkg/splunkoutputsexporter/subexporter.go b/pkg/splunkoutputsexporter/subexporter.go index 2fc2540..17639ae 100644 --- a/pkg/splunkoutputsexporter/subexporter.go +++ b/pkg/splunkoutputsexporter/subexporter.go @@ -6,11 +6,11 @@ package splunkoutputsexporter import ( "context" "fmt" - "os" "strings" "go.opentelemetry.io/collector/component" "go.opentelemetry.io/collector/exporter" + "go.uber.org/zap" "github.com/splunk/tarunner/internal/conf" "github.com/splunk/tarunner/internal/stanza" @@ -76,12 +76,9 @@ func newFactoryOptions(opts ...Option) factoryOptions { func (o factoryOptions) createLogsFunc(ctx context.Context, settings exporter.Settings, config component.Config) (exporter.Logs, error) { cfg := config.(Config) - splunkHome := cfg.BaseDir - if splunkHome == "" { - splunkHome = os.Getenv("SPLUNK_HOME") - } - if splunkHome == "" { - return nil, fmt.Errorf("splunk_outputs: base_dir is not set and SPLUNK_HOME is not defined") + splunkHome, err := tabuilder.ResolveSplunkHome(cfg.BaseDir) + if err != nil { + return nil, fmt.Errorf("splunk_outputs: %w", err) } outputs, err := tabuilder.ReadOutputGroups(splunkHome) @@ -103,6 +100,10 @@ func (o factoryOptions) createExporters(ctx context.Context, baseDir string, out if err != nil { return nil, fmt.Errorf("failed to create exporter %q: %w", output.Configuration.Stanza.Name, err) } + if e == nil { + settings.Logger.Info("splunk_outputs: skipping unsupported output stanza", zap.String("stanza", output.Configuration.Stanza.Name)) + continue + } exporters = append(exporters, e) } return exporters, nil