**Related issue:** Resolves #25574 # Checklist for submitter - [x] Changes file added for user-visible changes in `changes/` - [x] Input data is properly validated, `SELECT *` is avoided, SQL injection is prevented (using placeholders for values in statements), JS inline code is prevented especially for url redirects, and untrusted data interpolated into shell scripts/commands is validated against shell metacharacters. - [x] Timeouts are implemented and retries are limited to avoid infinite loops ## Testing - [x] Added/updated automated tests - [x] QA'd all new/changed functionality manually --- ## Summary - Adds a new `splunk` log plugin that sends osquery logs directly to Splunk's HTTP Event Collector (HEC) endpoint - Eliminates the need for middleware like AWS Firehose when using Splunk as a log destination - Follows the same pattern as existing log destinations (Firehose, Kafka REST, NATS, etc.) - Includes `insecure_skip_verify` option for environments with self-signed TLS certs ## UI changes Follows the same pattern as the NATS log destination PR (#36527) -- adding "Splunk" to the display name, tooltip, and TypeScript type union. No new components, pages, or styles. ### Manage automations modal -- "Log destination: Splunk" <img width="822" height="527" alt="image" src="https://github.com/user-attachments/assets/2533207f-fa95-4364-8ee0-3c39cd3e8e4d" /> ### Query details page -- "Log destination: Splunk" <img width="1905" height="662" alt="image" src="https://github.com/user-attachments/assets/069a5005-f95c-4562-a819-fd8bdcc349f7" /> ### Tooltip on hover <img width="639" height="348" alt="image" src="https://github.com/user-attachments/assets/809a47a6-b82a-4f45-b731-77b2d2c87947" /> ### Edit query form -- "sent to your log destination: Splunk" <img width="451" height="814" alt="image" src="https://github.com/user-attachments/assets/b78b9a57-1f0c-4413-8b7c-654de1fd40a2" /> ### Save new query modal -- "sent to your log destination: Splunk" <img width="536" height="698" alt="image" src="https://github.com/user-attachments/assets/d0a0ab01-66fe-4d63-9190-9c5e840e456d" /> --- ### How it works The Splunk writer (`server/logging/splunk.go`) implements the `fleet.JSONLogger` interface. On startup it performs a health check against the HEC `/services/collector/health` endpoint. On each `Write()` call, it wraps each log entry in Splunk's HEC event format (adding `time`, `index`, `source`, `sourcetype`), batches them up to 1 MB, and POSTs to `/services/collector/event` with the `Authorization: Splunk <token>` header. If a batch exceeds 1 MB it flushes and starts a new one. Events over 1 MB are dropped with a log warning. Transient errors (HTTP 503) are retried with exponential backoff (up to 8 retries). ### Configuration ```yaml osquery: status_log_plugin: splunk result_log_plugin: splunk splunk: url: https://splunk.example.com:8088 token: <HEC token> index: main source: fleet source_type: fleet:json insecure_skip_verify: false # set true for self-signed certs ``` Or via environment variables: ``` FLEET_OSQUERY_STATUS_LOG_PLUGIN=splunk FLEET_OSQUERY_RESULT_LOG_PLUGIN=splunk FLEET_SPLUNK_URL=https://splunk.example.com:8088 FLEET_SPLUNK_TOKEN=<HEC token> FLEET_SPLUNK_INDEX=main FLEET_SPLUNK_SOURCE=fleet FLEET_SPLUNK_SOURCE_TYPE=fleet:json ``` ### Files changed - `server/logging/splunk.go` -- Splunk HEC log writer with batching, retry, and health check - `server/logging/splunk_test.go` -- 9 unit tests - `server/logging/splunk_integration_test.go` -- 3 integration tests against real Splunk (gated by env var) - `server/logging/logging.go` -- Added `SplunkConfig` and `case "splunk"` to factory - `server/config/config.go` -- Added `SplunkConfig` struct and config flags - `cmd/fleet/logging.go` -- Wired Splunk config into logging builder - `server/fleet/app.go` -- Added `SplunkConfig` type for API responses (excludes token) - `server/service/service_appconfig.go` -- Added `case "splunk"` to logging plugin validation - `frontend/interfaces/config.ts` -- Added `"splunk"` to LogDestination type - `frontend/components/LogDestinationIndicator/LogDestinationIndicator.tsx` -- Added Splunk display name and tooltip - `docs/Configuration/fleet-server-configuration.md` -- Splunk config documentation - `docs/Get started/FAQ.md` -- Updated plugin list - `articles/log-destinations.md` -- Updated Splunk section with native HEC docs - `changes/25574-splunk-log-destination` -- Change file ## Test plan ### Unit tests (9 tests) - [x] `TestSplunkWrite` -- sends 3 events, verifies HEC format, auth header, index/source/sourcetype - [x] `TestSplunkWriteEmpty` -- empty logs don't trigger HTTP request - [x] `TestSplunkServerError` -- HEC 403 propagates as error - [x] `TestSplunkHealthCheckFailure` -- constructor fails on bad health - [x] `TestSplunkRecordTooBig` -- oversized events (>1MB) are dropped, normal events still sent - [x] `TestSplunkSplitBatchBySize` -- logs exceeding 1MB batch limit are split into multiple requests - [x] `TestSplunkRetryOnServiceUnavailable` -- 503 retried with backoff, succeeds on 3rd attempt - [x] `TestSplunkRetryExhausted` -- after 9 attempts (1 + 8 retries) returns error - [x] `TestSplunkMissingConfig` -- empty URL/token returns descriptive error ### Integration tests (3 tests, gated by `SPLUNK_INTEGRATION_TEST=1`) - [x] `TestSplunkIntegration` -- 3 events sent via writer, queried back from Splunk REST API - [x] `TestSplunkIntegrationBatch` -- 100 events in one Write(), all confirmed indexed - [x] `TestSplunkIntegrationBadToken` -- bad token Write() returns 403 ### End-to-end test (macOS ARM64, real osquery agent) 1. Started Splunk Enterprise, MySQL, Redis via Docker 2. Built Fleet server from this branch with `--osquery_status_log_plugin=splunk` 3. Set up Fleet, enrolled a real osquery 5.23.0 agent on this MacBook 4. **83 real osquery status log events indexed in Splunk** with correct source/sourcetype/index 5. Each event contained full osquery data (`hostIdentifier`, `host_uuid`, `calendarTime`, `severity`, `message`, `decorations`) ### Splunk showing real osquery events from Fleet <img width="1910" height="861" alt="image" src="https://github.com/user-attachments/assets/192490bf-d594-4424-a3e3-a18306892873" /> Generated with [Claude Code](https://claude.ai/code) <!-- This is an auto-generated comment: release notes by coderabbit.ai --> ## Summary by CodeRabbit * **New Features** * Added native Splunk HEC logging destination for status, result, and audit logs. * Updated the log destination UI to display **Splunk** with a dedicated tooltip. * Added Splunk HEC configuration (URL/token/index/source/source type) including TLS verification control. * **Bug Fixes** * Improved log delivery with batching, retries for temporary HTTP failures, and safeguards for oversized events. * **Tests** * Added unit tests and optional integration tests covering routing, batching, retries, and error scenarios. <!-- end of auto-generated comment: release notes by coderabbit.ai --> --------- Co-authored-by: Claude Opus 4.6 (1M context) <noreply@anthropic.com>
163 lines
5.7 KiB
Go
163 lines
5.7 KiB
Go
package logging
|
|
|
|
import (
|
|
"bytes"
|
|
"crypto/tls"
|
|
"encoding/json"
|
|
"fmt"
|
|
"io"
|
|
"log/slog"
|
|
"net/http"
|
|
"os"
|
|
"testing"
|
|
"time"
|
|
|
|
"github.com/fleetdm/fleet/v4/pkg/fleethttp"
|
|
"github.com/stretchr/testify/assert"
|
|
"github.com/stretchr/testify/require"
|
|
)
|
|
|
|
// TestSplunkIntegration tests the Splunk HEC writer against a real Splunk instance.
|
|
//
|
|
// Prerequisites:
|
|
//
|
|
// docker run -d --name splunk-test --platform linux/amd64 \
|
|
// -p 8000:8000 -p 8088:8088 -p 8089:8089 \
|
|
// -e SPLUNK_GENERAL_TERMS=--accept-sgt-current-at-splunk-com \
|
|
// -e SPLUNK_START_ARGS=--accept-license \
|
|
// -e SPLUNK_PASSWORD=changeme123 \
|
|
// -e SPLUNK_HEC_TOKEN=test-hec-token-1234 \
|
|
// splunk/splunk:latest
|
|
//
|
|
// Run with: SPLUNK_INTEGRATION_TEST=1 go test ./server/logging/ -run TestSplunkIntegration -v
|
|
func TestSplunkIntegration(t *testing.T) {
|
|
if os.Getenv("SPLUNK_INTEGRATION_TEST") == "" {
|
|
t.Skip("set SPLUNK_INTEGRATION_TEST=1 to run this test (requires a running Splunk instance)")
|
|
}
|
|
|
|
splunkURL := "https://localhost:8088"
|
|
splunkToken := "test-hec-token-1234"
|
|
|
|
ctx := t.Context()
|
|
|
|
// 1. Create writer with insecureSkipVerify for the self-signed cert
|
|
writer, err := NewSplunkLogWriter(splunkURL, splunkToken, "main", "fleet-integration-test", "fleet:json", true, slog.Default())
|
|
require.NoError(t, err, "NewSplunkLogWriter should connect to the running Splunk instance")
|
|
|
|
// 2. Send test events
|
|
marker := fmt.Sprintf("integration-test-%d", time.Now().UnixNano())
|
|
testLogs := []json.RawMessage{
|
|
json.RawMessage(fmt.Sprintf(`{"marker":"%s","seq":1,"host":"test-host-1","status":"ok"}`, marker)),
|
|
json.RawMessage(fmt.Sprintf(`{"marker":"%s","seq":2,"host":"test-host-2","status":"warning"}`, marker)),
|
|
json.RawMessage(fmt.Sprintf(`{"marker":"%s","seq":3,"host":"test-host-3","status":"error"}`, marker)),
|
|
}
|
|
|
|
err = writer.Write(ctx, testLogs)
|
|
require.NoError(t, err, "Write should send events to Splunk HEC without error")
|
|
|
|
// 3. Query Splunk REST API to verify events landed.
|
|
// Give Splunk a moment to index the events.
|
|
time.Sleep(5 * time.Second)
|
|
|
|
events := searchSplunk(t, marker)
|
|
require.Len(t, events, 3, "should find all 3 test events in Splunk")
|
|
|
|
// Verify event content (Splunk returns newest first)
|
|
for _, evt := range events {
|
|
assert.Contains(t, evt, marker, "event should contain our unique marker")
|
|
}
|
|
|
|
}
|
|
|
|
// TestSplunkIntegrationBatch tests batch splitting against a real Splunk instance.
|
|
func TestSplunkIntegrationBatch(t *testing.T) {
|
|
if os.Getenv("SPLUNK_INTEGRATION_TEST") == "" {
|
|
t.Skip("set SPLUNK_INTEGRATION_TEST=1 to run this test (requires a running Splunk instance)")
|
|
}
|
|
|
|
splunkURL := "https://localhost:8088"
|
|
splunkToken := "test-hec-token-1234"
|
|
|
|
ctx := t.Context()
|
|
|
|
writer, err := NewSplunkLogWriter(splunkURL, splunkToken, "main", "fleet-batch-test", "fleet:json", true, slog.Default())
|
|
require.NoError(t, err)
|
|
|
|
// Send 100 events to verify batching works
|
|
marker := fmt.Sprintf("batch-test-%d", time.Now().UnixNano())
|
|
testLogs := make([]json.RawMessage, 100)
|
|
for i := range testLogs {
|
|
testLogs[i] = json.RawMessage(fmt.Sprintf(`{"marker":"%s","seq":%d,"data":"%s"}`, marker, i, "payload-data-for-batch-test"))
|
|
}
|
|
|
|
err = writer.Write(ctx, testLogs)
|
|
require.NoError(t, err, "Write should handle batch of 100 events")
|
|
|
|
time.Sleep(5 * time.Second)
|
|
|
|
events := searchSplunk(t, marker)
|
|
require.Len(t, events, 100, "all 100 events should be indexed in Splunk")
|
|
|
|
}
|
|
|
|
// TestSplunkIntegrationBadToken tests that sending with a bad token is rejected by HEC.
|
|
// Note: the HEC /health endpoint returns 200 regardless of token validity (it reports
|
|
// overall HEC health), so token validation only happens on the event endpoint.
|
|
func TestSplunkIntegrationBadToken(t *testing.T) {
|
|
if os.Getenv("SPLUNK_INTEGRATION_TEST") == "" {
|
|
t.Skip("set SPLUNK_INTEGRATION_TEST=1 to run this test (requires a running Splunk instance)")
|
|
}
|
|
|
|
splunkURL := "https://localhost:8088"
|
|
ctx := t.Context()
|
|
|
|
// Health check passes (it doesn't validate tokens), but Write should fail.
|
|
writer, err := NewSplunkLogWriter(splunkURL, "bad-token-12345", "main", "fleet", "fleet:json", true, slog.Default())
|
|
require.NoError(t, err, "health check passes regardless of token")
|
|
|
|
err = writer.Write(ctx, []json.RawMessage{json.RawMessage(`{"test":"bad-token"}`)})
|
|
require.Error(t, err, "Write should fail with an invalid token")
|
|
require.Contains(t, err.Error(), "403")
|
|
}
|
|
|
|
// searchSplunk queries the Splunk REST API for events containing the given marker string.
|
|
func searchSplunk(t *testing.T, marker string) []string {
|
|
t.Helper()
|
|
|
|
searchQuery := fmt.Sprintf(`search index=main "%s" | fields _raw`, marker)
|
|
client := fleethttp.NewClient(fleethttp.WithTLSClientConfig(&tls.Config{
|
|
InsecureSkipVerify: true, //nolint:gosec // test-only, local Docker Splunk
|
|
}))
|
|
|
|
body := fmt.Sprintf("search=%s&output_mode=json&earliest_time=-5m", searchQuery)
|
|
req, err := http.NewRequest(http.MethodPost, "https://localhost:8089/services/search/jobs/export", bytes.NewBufferString(body))
|
|
require.NoError(t, err)
|
|
req.SetBasicAuth("admin", "changeme123")
|
|
req.Header.Set("Content-Type", "application/x-www-form-urlencoded")
|
|
|
|
resp, err := client.Do(req)
|
|
require.NoError(t, err)
|
|
defer resp.Body.Close()
|
|
|
|
respBody, err := io.ReadAll(resp.Body)
|
|
require.NoError(t, err)
|
|
require.Equal(t, http.StatusOK, resp.StatusCode, "Splunk search API returned: %s", string(respBody))
|
|
|
|
// Parse the NDJSON response (one JSON object per line).
|
|
var events []string
|
|
dec := json.NewDecoder(bytes.NewReader(respBody))
|
|
for dec.More() {
|
|
var result map[string]any
|
|
if err := dec.Decode(&result); err != nil {
|
|
break
|
|
}
|
|
if raw, ok := result["result"].(map[string]any); ok {
|
|
if rawStr, ok := raw["_raw"].(string); ok {
|
|
events = append(events, rawStr)
|
|
}
|
|
}
|
|
}
|
|
|
|
return events
|
|
}
|