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