**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>
259 lines
6.0 KiB
Go
259 lines
6.0 KiB
Go
// Package logging provides logger "plugins" for various destinations.
|
|
package logging
|
|
|
|
import (
|
|
"context"
|
|
"fmt"
|
|
"log/slog"
|
|
"time"
|
|
|
|
"github.com/fleetdm/fleet/v4/server/fleet"
|
|
)
|
|
|
|
type FilesystemConfig struct {
|
|
LogFile string
|
|
|
|
EnableLogRotation bool
|
|
EnableLogCompression bool
|
|
MaxSize int
|
|
MaxAge int
|
|
MaxBackups int
|
|
}
|
|
|
|
type WebhookConfig struct {
|
|
URL string
|
|
}
|
|
|
|
type FirehoseConfig struct {
|
|
StreamName string
|
|
|
|
Region string
|
|
EndpointURL string
|
|
AccessKeyID string
|
|
SecretAccessKey string
|
|
StsAssumeRoleArn string
|
|
StsExternalID string
|
|
}
|
|
|
|
type KinesisConfig struct {
|
|
StreamName string
|
|
|
|
Region string
|
|
EndpointURL string
|
|
AccessKeyID string
|
|
SecretAccessKey string
|
|
StsAssumeRoleArn string
|
|
StsExternalID string
|
|
}
|
|
|
|
type LambdaConfig struct {
|
|
Function string
|
|
|
|
Region string
|
|
AccessKeyID string
|
|
SecretAccessKey string
|
|
StsAssumeRoleArn string
|
|
StsExternalID string
|
|
}
|
|
|
|
type PubSubConfig struct {
|
|
Topic string
|
|
|
|
Project string
|
|
AddAttributes bool
|
|
}
|
|
|
|
type KafkaRESTConfig struct {
|
|
Topic string
|
|
|
|
ProxyHost string
|
|
ContentTypeValue string
|
|
Timeout int
|
|
}
|
|
|
|
type NatsConfig struct {
|
|
Server string
|
|
Subject string
|
|
|
|
CredFile string
|
|
NKeyFile string
|
|
|
|
TLSClientCertFile string
|
|
TLSClientKeyFile string
|
|
CACertFile string
|
|
|
|
Compression string
|
|
JetStream bool
|
|
|
|
Timeout time.Duration
|
|
}
|
|
|
|
type SplunkConfig struct {
|
|
URL string
|
|
Token string
|
|
Index string
|
|
Source string
|
|
SourceType string
|
|
InsecureSkipVerify bool
|
|
}
|
|
|
|
type Config struct {
|
|
Plugin string
|
|
|
|
Filesystem FilesystemConfig
|
|
Webhook WebhookConfig
|
|
Firehose FirehoseConfig
|
|
Kinesis KinesisConfig
|
|
Lambda LambdaConfig
|
|
PubSub PubSubConfig
|
|
KafkaREST KafkaRESTConfig
|
|
Nats NatsConfig
|
|
Splunk SplunkConfig
|
|
}
|
|
|
|
func NewJSONLogger(ctx context.Context, name string, config Config, logger *slog.Logger) (fleet.JSONLogger, error) {
|
|
switch config.Plugin {
|
|
case "":
|
|
// Allow "" to mean filesystem for backwards compatibility
|
|
logger.InfoContext(ctx, fmt.Sprintf("plugin for %s not explicitly specified. Assuming 'filesystem'", name))
|
|
fallthrough
|
|
case "filesystem":
|
|
writer, err := NewFilesystemLogWriter(
|
|
ctx,
|
|
config.Filesystem.LogFile,
|
|
logger,
|
|
config.Filesystem.EnableLogRotation,
|
|
config.Filesystem.EnableLogCompression,
|
|
config.Filesystem.MaxSize,
|
|
config.Filesystem.MaxAge,
|
|
config.Filesystem.MaxBackups,
|
|
)
|
|
if err != nil {
|
|
return nil, fmt.Errorf("create filesystem %s logger: %w", name, err)
|
|
}
|
|
return fleet.JSONLogger(writer), nil
|
|
case "webhook":
|
|
writer, err := NewWebhookLogWriter(config.Webhook.URL, logger)
|
|
if err != nil {
|
|
return nil, fmt.Errorf("create webhook %s logger: %w", name, err)
|
|
}
|
|
return fleet.JSONLogger(writer), nil
|
|
case "firehose":
|
|
writer, err := NewFirehoseLogWriter(
|
|
config.Firehose.Region,
|
|
config.Firehose.EndpointURL,
|
|
config.Firehose.AccessKeyID,
|
|
config.Firehose.SecretAccessKey,
|
|
config.Firehose.StsAssumeRoleArn,
|
|
config.Firehose.StsExternalID,
|
|
config.Firehose.StreamName,
|
|
logger,
|
|
)
|
|
if err != nil {
|
|
return nil, fmt.Errorf("create firehose %s logger: %w", name, err)
|
|
}
|
|
return fleet.JSONLogger(writer), nil
|
|
case "kinesis":
|
|
writer, err := NewKinesisLogWriter(
|
|
config.Kinesis.Region,
|
|
config.Kinesis.EndpointURL,
|
|
config.Kinesis.AccessKeyID,
|
|
config.Kinesis.SecretAccessKey,
|
|
config.Kinesis.StsAssumeRoleArn,
|
|
config.Kinesis.StsExternalID,
|
|
config.Kinesis.StreamName,
|
|
logger,
|
|
)
|
|
if err != nil {
|
|
return nil, fmt.Errorf("create kinesis %s logger: %w", name, err)
|
|
}
|
|
return fleet.JSONLogger(writer), nil
|
|
case "lambda":
|
|
writer, err := NewLambdaLogWriter(
|
|
config.Lambda.Region,
|
|
config.Lambda.AccessKeyID,
|
|
config.Lambda.SecretAccessKey,
|
|
config.Lambda.StsAssumeRoleArn,
|
|
config.Lambda.StsExternalID,
|
|
config.Lambda.Function,
|
|
logger,
|
|
)
|
|
if err != nil {
|
|
return nil, fmt.Errorf("create lambda %s logger: %w", name, err)
|
|
}
|
|
return fleet.JSONLogger(writer), nil
|
|
case "pubsub":
|
|
writer, err := NewPubSubLogWriter(
|
|
ctx,
|
|
config.PubSub.Project,
|
|
config.PubSub.Topic,
|
|
config.PubSub.AddAttributes,
|
|
logger,
|
|
)
|
|
if err != nil {
|
|
return nil, fmt.Errorf("create pubsub %s logger: %w", name, err)
|
|
}
|
|
return fleet.JSONLogger(writer), nil
|
|
case "stdout":
|
|
writer, err := NewStdoutLogWriter()
|
|
if err != nil {
|
|
return nil, fmt.Errorf("create stdout %s logger: %w", name, err)
|
|
}
|
|
return fleet.JSONLogger(writer), nil
|
|
case "kafkarest":
|
|
writer, err := NewKafkaRESTWriter(&KafkaRESTParams{
|
|
KafkaProxyHost: config.KafkaREST.ProxyHost,
|
|
KafkaTopic: config.KafkaREST.Topic,
|
|
KafkaContentTypeValue: config.KafkaREST.ContentTypeValue,
|
|
KafkaTimeout: config.KafkaREST.Timeout,
|
|
})
|
|
if err != nil {
|
|
return nil, fmt.Errorf("create kafka rest %s logger: %w", name, err)
|
|
}
|
|
return fleet.JSONLogger(writer), nil
|
|
case "nats":
|
|
writer, err := NewNatsLogWriter(
|
|
ctx,
|
|
config.Nats.Server,
|
|
config.Nats.Subject,
|
|
config.Nats.CredFile,
|
|
config.Nats.NKeyFile,
|
|
config.Nats.TLSClientCertFile,
|
|
config.Nats.TLSClientKeyFile,
|
|
config.Nats.CACertFile,
|
|
config.Nats.Compression,
|
|
config.Nats.JetStream,
|
|
config.Nats.Timeout,
|
|
logger,
|
|
)
|
|
if err != nil {
|
|
return nil, fmt.Errorf("create nats %s logger: %w", name, err)
|
|
}
|
|
return fleet.JSONLogger(writer), nil
|
|
case "splunk":
|
|
if config.Splunk.URL == "" {
|
|
return nil, fmt.Errorf("splunk %s logger: URL must not be empty", name)
|
|
}
|
|
if config.Splunk.Token == "" {
|
|
return nil, fmt.Errorf("splunk %s logger: HEC token must not be empty", name)
|
|
}
|
|
writer, err := NewSplunkLogWriter(
|
|
config.Splunk.URL,
|
|
config.Splunk.Token,
|
|
config.Splunk.Index,
|
|
config.Splunk.Source,
|
|
config.Splunk.SourceType,
|
|
config.Splunk.InsecureSkipVerify,
|
|
logger,
|
|
)
|
|
if err != nil {
|
|
return nil, fmt.Errorf("create splunk %s logger: %w", name, err)
|
|
}
|
|
return fleet.JSONLogger(writer), nil
|
|
default:
|
|
return nil, fmt.Errorf(
|
|
"unknown %s log plugin: %s", name, config.Plugin,
|
|
)
|
|
}
|
|
}
|