Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
7 changes: 7 additions & 0 deletions frontend/src/features/log-sinks/LogSinksPage.tsx
Original file line number Diff line number Diff line change
Expand Up @@ -115,6 +115,13 @@ const SINK_TYPE_ICONS: Record<string, ReactNode> = {
<path d="M3 12h4l2-5 3 10 2-5h7" />
</svg>
),
// http — arrow upload / send
http: (
<svg viewBox="0 0 24 24" fill="none" stroke="currentColor" strokeWidth="1.5">
<path d="M12 19V5M5 12l7-7 7 7" />
<path d="M5 21h14" />
</svg>
),
};

/** Returns the icon for a sink type, or a default document icon. */
Expand Down
21 changes: 21 additions & 0 deletions pkg/config/loggers.go
Original file line number Diff line number Diff line change
Expand Up @@ -85,6 +85,27 @@ type LogstashLogger struct {
Path string `yaml:"path"`
}

// HTTPLogger holds configuration for the generic HTTP log sink. It sends
// the raw osquery log payload (optionally wrapped with metadata) to an
// arbitrary HTTP/HTTPS endpoint with configurable method, headers, and
// serialization format.
type HTTPLogger struct {
URL string `yaml:"url" json:"url"`
Method string `yaml:"method" json:"method"`
Headers map[string]string `yaml:"headers" json:"headers"`
// Format controls the wire encoding. "json" (the default) sends the
// payload as-is with a JSON content type; "ndjson" splits array
// payloads into newline-delimited JSON objects; "raw" sends the
// bytes verbatim with the configured content type.
Format string `yaml:"format" json:"format"`
ContentType string `yaml:"contentType" json:"contentType"`
TimeoutSeconds int `yaml:"timeoutSeconds" json:"timeoutSeconds"`
// IncludeMetadata, when true, wraps each event in an envelope that
// adds environment, uuid, logType, and timestamp fields. When false
// the raw osquery payload is forwarded unchanged.
IncludeMetadata bool `yaml:"includeMetadata" json:"includeMetadata"`
}

// LocalLogger to hold all local logger configuration values
type LocalLogger struct {
FilePath string `yaml:"filePath"`
Expand Down
2 changes: 2 additions & 0 deletions pkg/config/types.go
Original file line number Diff line number Diff line change
Expand Up @@ -48,6 +48,7 @@ const (
LoggingS3 string = "s3"
LoggingKafka string = "kafka"
LoggingElastic string = "elastic"
LoggingHTTP string = "http"
)

// Types of carver
Expand Down Expand Up @@ -296,6 +297,7 @@ type YAMLConfigurationLogger struct {
Kinesis *KinesisLogger `mapstructure:"kinesis"`
Kafka *KafkaLogger `mapstructure:"kafka"`
Local *LocalLogger `mapstructure:"local"`
HTTP *HTTPLogger `mapstructure:"http"`
}

// YAMLConfigurationCarver to hold the carver configuration values
Expand Down
1 change: 1 addition & 0 deletions pkg/config/validation.go
Original file line number Diff line number Diff line change
Expand Up @@ -23,6 +23,7 @@ var validLogging = map[string]bool{
LoggingS3: true,
LoggingKafka: true,
LoggingElastic: true,
LoggingHTTP: true,
}

// Valid values for carver in configuration
Expand Down
11 changes: 11 additions & 0 deletions pkg/config/validation_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -17,6 +17,17 @@ func TestValidateTLSConfigValuesAcceptsKafkaAndMultipleLoggers(t *testing.T) {
}
}

func TestValidateTLSConfigValuesAcceptsHTTP(t *testing.T) {
err := ValidateTLSConfigValues(TLSConfiguration{
Service: YAMLConfigurationService{Auth: AuthNone},
Logger: &YAMLConfigurationLogger{Types: []string{LoggingHTTP}},
Carver: &YAMLConfigurationCarver{Type: CarverDB},
})
if err != nil {
t.Fatalf("expected http logger type to validate: %v", err)
}
}

func TestValidateTLSConfigValuesRejectsInvalidMultipleLogger(t *testing.T) {
err := ValidateTLSConfigValues(TLSConfiguration{
Service: YAMLConfigurationService{Auth: AuthNone},
Expand Down
13 changes: 13 additions & 0 deletions pkg/logging/exporter_adapters.go
Original file line number Diff line number Diff line change
Expand Up @@ -177,3 +177,16 @@ func (logE *LoggerElastic) Export(logType string, data []byte, params ExportPara
logE.Send(logType, data, params.Environment, params.UUID, params.Debug)
return nil
}

func (l *LoggerHTTP) Name() string {
return config.LoggingHTTP
}

func (l *LoggerHTTP) IsEnabled() bool {
return l != nil && l.Enabled
}

func (l *LoggerHTTP) Export(logType string, data []byte, params ExportParams) error {
l.Send(logType, data, params.Environment, params.UUID, params.Debug)
return nil
}
7 changes: 7 additions & 0 deletions pkg/logging/exporter_factory.go
Original file line number Diff line number Diff line change
Expand Up @@ -141,6 +141,13 @@ func CreateExporter(exporterType string, cfg config.ServiceParameters, mgr *sett
}
e.Settings(mgr)
return e, false, nil
case config.LoggingHTTP:
h, err := CreateLoggerHTTP(cfg.Logger.HTTP)
if err != nil {
return nil, false, err
}
h.Settings(mgr)
return h, false, nil
default:
return nil, false, fmt.Errorf("unknown exporter type: %s", exporterType)
}
Expand Down
241 changes: 241 additions & 0 deletions pkg/logging/http.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,241 @@
package logging

import (
"bytes"
"encoding/json"
"fmt"
"io"
"net/http"
"strings"
"time"

"github.com/jmpsec/osctrl/pkg/config"
"github.com/jmpsec/osctrl/pkg/settings"
"github.com/jmpsec/osctrl/pkg/types"
"github.com/rs/zerolog/log"
)

const (
// HTTPFormatJSON sends the payload as-is with a JSON content type.
HTTPFormatJSON = "json"
// HTTPFormatNDJSON splits array payloads into newline-delimited JSON
// objects, one per line. Non-array payloads are sent as a single line.
HTTPFormatNDJSON = "ndjson"
// HTTPFormatRaw sends the bytes verbatim using the configured
// content type, with no re-serialization.
HTTPFormatRaw = "raw"
)

// httpDefaultTimeout caps a single HTTP request so an unresponsive
// endpoint cannot stall the log ingestion path indefinitely.
const httpDefaultTimeout = 30 * time.Second

// LoggerHTTP is a generic HTTP log sink. It forwards osquery log
// payloads to an arbitrary HTTP/HTTPS endpoint with configurable
// method, headers, and serialization format. Unlike Splunk or Graylog,
// it does not impose a vendor-specific envelope; operators choose
// whether metadata is wrapped via IncludeMetadata.
type LoggerHTTP struct {
Configuration config.HTTPLogger
client *http.Client
headers map[string]string
Enabled bool
}

// CreateLoggerHTTP initializes a generic HTTP logger from the given
// configuration. The HTTP client is reused across calls so connection
// pooling is effective.
func CreateLoggerHTTP(cfg *config.HTTPLogger) (*LoggerHTTP, error) {
if cfg == nil {
cfg = &config.HTTPLogger{}
}
method := strings.ToUpper(strings.TrimSpace(cfg.Method))
if method == "" {
method = http.MethodPost
}
format := strings.ToLower(strings.TrimSpace(cfg.Format))
if format == "" {
format = HTTPFormatJSON
}
contentType := strings.TrimSpace(cfg.ContentType)
if contentType == "" {
switch format {
case HTTPFormatNDJSON:
contentType = "application/x-ndjson"
case HTTPFormatRaw:
contentType = "application/octet-stream"
default:
contentType = "application/json"
}
}
timeout := time.Duration(cfg.TimeoutSeconds) * time.Second
if timeout <= 0 {
timeout = httpDefaultTimeout
}
headers := make(map[string]string, len(cfg.Headers)+1)
for k, v := range cfg.Headers {
headers[k] = v
}
if _, ok := headers["Content-Type"]; !ok {
headers["Content-Type"] = contentType
}
l := &LoggerHTTP{
Configuration: config.HTTPLogger{
URL: cfg.URL,
Method: method,
Headers: cfg.Headers,
Format: format,
ContentType: contentType,
TimeoutSeconds: cfg.TimeoutSeconds,
IncludeMetadata: cfg.IncludeMetadata,
},
client: &http.Client{Timeout: timeout},
headers: headers,
Enabled: true,
}
return l, nil
}

// Settings is a no-op for the HTTP logger; all configuration is carried
// in the config struct. The method exists to satisfy the exporter
// convention used by the registry builder.
func (l *LoggerHTTP) Settings(mgr *settings.Settings) {
log.Info().Msg("Setting HTTP logging settings")
}

// Close releases the HTTP client transport so pooled connections are
// not leaked across hot reloads.
func (l *LoggerHTTP) Close() error {
if l.client != nil {
l.client.CloseIdleConnections()
}
return nil
}

// httpEnvelope wraps a raw osquery event with osctrl metadata when
// IncludeMetadata is enabled.
type httpEnvelope struct {
Environment string `json:"environment"`
UUID string `json:"uuid"`
LogType string `json:"log_type"`
Timestamp int64 `json:"timestamp"`
Event json.RawMessage `json:"event"`
}

// Send serializes and transmits the osquery payload to the configured
// HTTP endpoint. The format and metadata wrapping are controlled by
// the configuration.
func (l *LoggerHTTP) Send(logType string, data []byte, environment, uuid string, debug bool) {
if l.Configuration.URL == "" {
return
}
body, err := l.encode(logType, data, environment, uuid)
if err != nil {
log.Err(err).Str("type", logType).Msg("http sink: encode error")
return
}
if debug {
log.Debug().Msgf("http sink: sending %d bytes (%s) to %s", body.Len(), logType, l.Configuration.URL)
}
req, err := http.NewRequest(l.Configuration.Method, l.Configuration.URL, body)
if err != nil {
log.Err(err).Msg("http sink: build request")
return
}
for k, v := range l.headers {
req.Header.Set(k, v)
}
resp, err := l.client.Do(req)
if err != nil {
log.Err(err).Msg("http sink: send request")
return
}
defer resp.Body.Close()
if resp.StatusCode >= 400 {
respBody, _ := io.ReadAll(io.LimitReader(resp.Body, 512))
log.Warn().
Int("status", resp.StatusCode).
Str("type", logType).
Bytes("body", respBody).
Msg("http sink: non-2xx response")
}
if debug {
log.Debug().Msgf("http sink: response %d", resp.StatusCode)
}
}

// encode produces the request body according to the configured format
// and metadata preference.
func (l *LoggerHTTP) encode(logType string, data []byte, environment, uuid string) (*bytes.Buffer, error) {
if !l.Configuration.IncludeMetadata {
if l.Configuration.Format == HTTPFormatNDJSON {
return l.encodeNDJSON(data, nil)
}
return bytes.NewBuffer(data), nil
}
// Metadata wrapping: split array payloads into individual events,
// envelope each one, then re-serialize according to format.
var events []json.RawMessage
if logType == types.QueryLog {
events = []json.RawMessage{data}
} else {
if err := json.Unmarshal(data, &events); err != nil {
events = []json.RawMessage{data}
}
}
now := time.Now().Unix()
envelopes := make([]httpEnvelope, 0, len(events))
for _, ev := range events {
envelopes = append(envelopes, httpEnvelope{
Environment: environment,
UUID: uuid,
LogType: logType,
Timestamp: now,
Event: ev,
})
}
switch l.Configuration.Format {
case HTTPFormatNDJSON:
return l.encodeNDJSON(nil, envelopes)
case HTTPFormatRaw:
return bytes.NewBuffer(data), nil
default:
out, err := json.Marshal(envelopes)
if err != nil {
return nil, fmt.Errorf("marshal http envelopes: %w", err)
}
return bytes.NewBuffer(out), nil
}
}

// encodeNDJSON writes events as newline-delimited JSON. When raw is
// non-nil and envelopes is nil, it splits the raw array payload into
// individual lines. When envelopes is non-nil, each envelope is written
// as one line.
func (l *LoggerHTTP) encodeNDJSON(raw []byte, envelopes []httpEnvelope) (*bytes.Buffer, error) {
buf := bytes.NewBuffer(nil)
if envelopes != nil {
for _, e := range envelopes {
line, err := json.Marshal(e)
if err != nil {
return nil, fmt.Errorf("marshal ndjson line: %w", err)
}
buf.Write(line)
buf.WriteByte('\n')
}
return buf, nil
}
var items []json.RawMessage
if err := json.Unmarshal(raw, &items); err != nil {
buf.Write(raw)
if !bytes.HasSuffix(raw, []byte{'\n'}) {
buf.WriteByte('\n')
}
return buf, nil
}
for _, item := range items {
buf.Write(item)
buf.WriteByte('\n')
}
return buf, nil
}
Loading
Loading