From e91e3ccc1b885c640af10a36964faacaf99f13f3 Mon Sep 17 00:00:00 2001 From: Javier Marcos <1271349+javuto@users.noreply.github.com> Date: Thu, 17 Sep 2026 22:54:29 +0200 Subject: [PATCH] Added a generic HTTP logsink --- .../src/features/log-sinks/LogSinksPage.tsx | 7 + pkg/config/loggers.go | 21 ++ pkg/config/types.go | 2 + pkg/config/validation.go | 1 + pkg/config/validation_test.go | 11 + pkg/logging/exporter_adapters.go | 13 + pkg/logging/exporter_factory.go | 7 + pkg/logging/http.go | 241 ++++++++++++++++ pkg/logging/http_test.go | 269 ++++++++++++++++++ pkg/logsinks/logsinks.go | 78 +++++ pkg/logsinks/logsinks_test.go | 132 ++++++++- 11 files changed, 781 insertions(+), 1 deletion(-) create mode 100644 pkg/logging/http.go create mode 100644 pkg/logging/http_test.go diff --git a/frontend/src/features/log-sinks/LogSinksPage.tsx b/frontend/src/features/log-sinks/LogSinksPage.tsx index a76c0aed4..3c21c4f58 100644 --- a/frontend/src/features/log-sinks/LogSinksPage.tsx +++ b/frontend/src/features/log-sinks/LogSinksPage.tsx @@ -115,6 +115,13 @@ const SINK_TYPE_ICONS: Record = { ), + // http — arrow upload / send + http: ( + + + + + ), }; /** Returns the icon for a sink type, or a default document icon. */ diff --git a/pkg/config/loggers.go b/pkg/config/loggers.go index 7e6027a3a..01bd4dd80 100644 --- a/pkg/config/loggers.go +++ b/pkg/config/loggers.go @@ -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"` diff --git a/pkg/config/types.go b/pkg/config/types.go index 6942eca4b..112f4c8b2 100644 --- a/pkg/config/types.go +++ b/pkg/config/types.go @@ -48,6 +48,7 @@ const ( LoggingS3 string = "s3" LoggingKafka string = "kafka" LoggingElastic string = "elastic" + LoggingHTTP string = "http" ) // Types of carver @@ -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 diff --git a/pkg/config/validation.go b/pkg/config/validation.go index 5e20919f1..57fec0380 100644 --- a/pkg/config/validation.go +++ b/pkg/config/validation.go @@ -23,6 +23,7 @@ var validLogging = map[string]bool{ LoggingS3: true, LoggingKafka: true, LoggingElastic: true, + LoggingHTTP: true, } // Valid values for carver in configuration diff --git a/pkg/config/validation_test.go b/pkg/config/validation_test.go index b68cc90fe..f848330b4 100644 --- a/pkg/config/validation_test.go +++ b/pkg/config/validation_test.go @@ -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}, diff --git a/pkg/logging/exporter_adapters.go b/pkg/logging/exporter_adapters.go index 759073792..c8a6f4efc 100644 --- a/pkg/logging/exporter_adapters.go +++ b/pkg/logging/exporter_adapters.go @@ -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 +} diff --git a/pkg/logging/exporter_factory.go b/pkg/logging/exporter_factory.go index b193a2f47..5609db51a 100644 --- a/pkg/logging/exporter_factory.go +++ b/pkg/logging/exporter_factory.go @@ -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) } diff --git a/pkg/logging/http.go b/pkg/logging/http.go new file mode 100644 index 000000000..08d51ce6f --- /dev/null +++ b/pkg/logging/http.go @@ -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 +} diff --git a/pkg/logging/http_test.go b/pkg/logging/http_test.go new file mode 100644 index 000000000..d3a66b778 --- /dev/null +++ b/pkg/logging/http_test.go @@ -0,0 +1,269 @@ +package logging + +import ( + "encoding/json" + "io" + "net/http" + "net/http/httptest" + "testing" + + "github.com/jmpsec/osctrl/pkg/config" + "github.com/jmpsec/osctrl/pkg/types" +) + +func TestCreateLoggerHTTPDefaults(t *testing.T) { + l, err := CreateLoggerHTTP(&config.HTTPLogger{URL: "http://x"}) + if err != nil { + t.Fatalf("create: %v", err) + } + if l.Configuration.Method != http.MethodPost { + t.Errorf("method default: got %q want POST", l.Configuration.Method) + } + if l.Configuration.Format != HTTPFormatJSON { + t.Errorf("format default: got %q want json", l.Configuration.Format) + } + if l.headers["Content-Type"] != "application/json" { + t.Errorf("content-type default: got %q", l.headers["Content-Type"]) + } + if !l.IsEnabled() || l.Name() != config.LoggingHTTP { + t.Errorf("exporter name/enabled: %q %v", l.Name(), l.IsEnabled()) + } +} + +func TestCreateLoggerHTTPNilConfigDoesNotPanic(t *testing.T) { + l, err := CreateLoggerHTTP(nil) + if err != nil { + t.Fatalf("create nil: %v", err) + } + // Enabled defaults to true; Send is a no-op when URL is empty. + l.Send("status", []byte(`{}`), "env", "uuid", false) +} + +func TestLoggerHTTPSendJSON(t *testing.T) { + var gotBody []byte + var gotCT string + srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + gotBody, _ = io.ReadAll(r.Body) + gotCT = r.Header.Get("Content-Type") + w.WriteHeader(http.StatusOK) + })) + defer srv.Close() + + l, err := CreateLoggerHTTP(&config.HTTPLogger{ + URL: srv.URL, + Method: "POST", + Format: HTTPFormatJSON, + }) + if err != nil { + t.Fatalf("create: %v", err) + } + payload := []byte(`[{"name":"a"},{"name":"b"}]`) + l.Send("result", payload, "prod", "node-1", false) + + if string(gotBody) != string(payload) { + t.Errorf("json body: got %q want %q", gotBody, payload) + } + if gotCT != "application/json" { + t.Errorf("content-type: got %q want application/json", gotCT) + } +} + +func TestLoggerHTTPSendNDJSON(t *testing.T) { + var gotBody []byte + srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + gotBody, _ = io.ReadAll(r.Body) + w.WriteHeader(http.StatusOK) + })) + defer srv.Close() + + l, _ := CreateLoggerHTTP(&config.HTTPLogger{ + URL: srv.URL, + Format: HTTPFormatNDJSON, + }) + l.Send("result", []byte(`[{"a":1},{"b":2}]`), "env", "uuid", false) + + lines := splitLines(string(gotBody)) + if len(lines) != 2 { + t.Fatalf("ndjson lines: got %d want 2", len(lines)) + } + if lines[0] != `{"a":1}` || lines[1] != `{"b":2}` { + t.Errorf("ndjson body: got %q", gotBody) + } +} + +func TestLoggerHTTPSendMetadataEnvelope(t *testing.T) { + var gotBody []byte + srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + gotBody, _ = io.ReadAll(r.Body) + w.WriteHeader(http.StatusOK) + })) + defer srv.Close() + + l, _ := CreateLoggerHTTP(&config.HTTPLogger{ + URL: srv.URL, + IncludeMetadata: true, + }) + l.Send("status", []byte(`[{"msg":"hi"}]`), "prod", "node-1", false) + + var envelopes []httpEnvelope + if err := json.Unmarshal(gotBody, &envelopes); err != nil { + t.Fatalf("unmarshal envelopes: %v (body=%q)", err, gotBody) + } + if len(envelopes) != 1 { + t.Fatalf("envelope count: got %d want 1", len(envelopes)) + } + if envelopes[0].Environment != "prod" || envelopes[0].UUID != "node-1" || envelopes[0].LogType != "status" { + t.Errorf("envelope metadata: %+v", envelopes[0]) + } + if string(envelopes[0].Event) != `{"msg":"hi"}` { + t.Errorf("envelope event: got %q", envelopes[0].Event) + } +} + +func TestLoggerHTTPSendMetadataNDJSON(t *testing.T) { + var gotBody []byte + srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + gotBody, _ = io.ReadAll(r.Body) + w.WriteHeader(http.StatusOK) + })) + defer srv.Close() + + l, _ := CreateLoggerHTTP(&config.HTTPLogger{ + URL: srv.URL, + Format: HTTPFormatNDJSON, + IncludeMetadata: true, + }) + l.Send("result", []byte(`[{"x":1},{"y":2}]`), "env", "node", false) + + lines := splitLines(string(gotBody)) + if len(lines) != 2 { + t.Fatalf("ndjson envelope lines: got %d want 2", len(lines)) + } + var first httpEnvelope + if err := json.Unmarshal([]byte(lines[0]), &first); err != nil { + t.Fatalf("unmarshal first line: %v", err) + } + if first.LogType != "result" { + t.Errorf("first envelope logType: got %q want result", first.LogType) + } +} + +func TestLoggerHTTPSendRaw(t *testing.T) { + var gotBody []byte + var gotCT string + srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + gotBody, _ = io.ReadAll(r.Body) + gotCT = r.Header.Get("Content-Type") + w.WriteHeader(http.StatusOK) + })) + defer srv.Close() + + l, _ := CreateLoggerHTTP(&config.HTTPLogger{ + URL: srv.URL, + Format: HTTPFormatRaw, + ContentType: "text/plain", + }) + payload := []byte(`raw bytes here`) + l.Send("status", payload, "env", "uuid", false) + + if string(gotBody) != string(payload) { + t.Errorf("raw body: got %q want %q", gotBody, payload) + } + if gotCT != "text/plain" { + t.Errorf("raw content-type: got %q want text/plain", gotCT) + } +} + +func TestLoggerHTTPCustomHeaders(t *testing.T) { + var gotAuth string + srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + gotAuth = r.Header.Get("Authorization") + w.WriteHeader(http.StatusOK) + })) + defer srv.Close() + + l, _ := CreateLoggerHTTP(&config.HTTPLogger{ + URL: srv.URL, + Headers: map[string]string{"Authorization": "Bearer xyz"}, + }) + l.Send("status", []byte(`{}`), "env", "uuid", false) + + if gotAuth != "Bearer xyz" { + t.Errorf("custom header: got %q want Bearer xyz", gotAuth) + } +} + +func TestLoggerHTTPExportAdapter(t *testing.T) { + var gotBody []byte + srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + gotBody, _ = io.ReadAll(r.Body) + w.WriteHeader(http.StatusOK) + })) + defer srv.Close() + + l, _ := CreateLoggerHTTP(&config.HTTPLogger{URL: srv.URL}) + err := l.Export(types.QueryLog, []byte(`{"q":"test"}`), ExportParams{ + Environment: "env", + UUID: "node", + QueryName: "q1", + Status: 0, + Debug: false, + }) + if err != nil { + t.Fatalf("export: %v", err) + } + if string(gotBody) != `{"q":"test"}` { + t.Errorf("export body: got %q", gotBody) + } +} + +func TestLoggerHTTPCloseIsIdempotent(t *testing.T) { + l, _ := CreateLoggerHTTP(&config.HTTPLogger{URL: "http://x"}) + if err := l.Close(); err != nil { + t.Fatalf("close: %v", err) + } + if err := l.Close(); err != nil { + t.Fatalf("close again: %v", err) + } +} + +func TestLoggerHTTPNon2xxDoesNotPanic(t *testing.T) { + srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + w.WriteHeader(http.StatusInternalServerError) + })) + defer srv.Close() + + l, _ := CreateLoggerHTTP(&config.HTTPLogger{URL: srv.URL}) + l.Send("status", []byte(`{}`), "env", "uuid", false) +} + +func TestLoggerHTTPMethodPUT(t *testing.T) { + var gotMethod string + srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + gotMethod = r.Method + w.WriteHeader(http.StatusOK) + })) + defer srv.Close() + + l, _ := CreateLoggerHTTP(&config.HTTPLogger{URL: srv.URL, Method: "PUT"}) + l.Send("status", []byte(`{}`), "env", "uuid", false) + + if gotMethod != http.MethodPut { + t.Errorf("method: got %q want PUT", gotMethod) + } +} + +func splitLines(s string) []string { + var out []string + start := 0 + for i := 0; i < len(s); i++ { + if s[i] == '\n' { + out = append(out, s[start:i]) + start = i + 1 + } + } + if start < len(s) { + out = append(out, s[start:]) + } + return out +} diff --git a/pkg/logsinks/logsinks.go b/pkg/logsinks/logsinks.go index 6044e9c41..fa34601aa 100644 --- a/pkg/logsinks/logsinks.go +++ b/pkg/logsinks/logsinks.go @@ -390,6 +390,29 @@ var Registry = map[string]SinkSpec{ return e, nil }, }, + config.LoggingHTTP: { + Type: config.LoggingHTTP, + Description: "Generic HTTP/HTTPS log sink. Forwards raw osquery payloads to any web endpoint.", + HasSecret: false, + Fields: []FieldSpec{ + {Name: "url", Label: "URL", Type: FieldString, Required: true, Placeholder: "https://logs.example.com/ingest", Help: "HTTP/HTTPS endpoint that will receive log payloads."}, + {Name: "method", Label: "Method", Type: FieldSelect, Options: []string{"POST", "PUT"}, Default: "POST", Help: "HTTP method used for each request."}, + {Name: "format", Label: "Format", Type: FieldSelect, Options: []string{"json", "ndjson", "raw"}, Default: "json", Help: "json: forward as JSON array. ndjson: one JSON object per line. raw: bytes verbatim."}, + {Name: "contentType", Label: "Content-Type", Type: FieldString, Placeholder: "application/json", Help: "Overrides the Content-Type header derived from format when set."}, + {Name: "headers", Label: "Custom headers", Type: FieldString, Placeholder: `{"X-Token":"abc"}`, Help: "JSON object of extra HTTP headers. Header names are case-sensitive."}, + {Name: "includeMetadata", Label: "Include metadata", Type: FieldBoolean, Default: false, Help: "Wrap each event with environment, uuid, log_type, and timestamp fields."}, + {Name: "timeoutSeconds", Label: "Timeout (seconds)", Type: FieldInteger, Default: 30, Help: "Per-request timeout. 0 uses the 30s default."}, + }, + Decode: decodeHTTPConfig, + Build: func(cfg any, smgr *settings.Settings) (logging.DataExporter, error) { + h, err := logging.CreateLoggerHTTP(cfg.(*config.HTTPLogger)) + if err != nil { + return nil, err + } + h.Settings(smgr) + return h, nil + }, + }, } // SupportedTypes returns the Registry keys in a stable, sorted order, for @@ -430,6 +453,56 @@ func decodeTyped[T any]() func(json.RawMessage) (any, error) { } } +// decodeHTTPConfig unmarshals an HTTP sink config. The "headers" field +// can arrive as a JSON object (from the API/YAML) or as a JSON string +// (from the frontend single-line string input, e.g. `{"X-Token":"abc"}`). +// This decoder accepts both shapes and normalizes to map[string]string. +func decodeHTTPConfig(raw json.RawMessage) (any, error) { + var cfg config.HTTPLogger + if len(raw) == 0 || string(raw) == "null" { + return &cfg, nil + } + // First try a direct unmarshal — works when headers is a JSON object. + if err := json.Unmarshal(raw, &cfg); err == nil { + return &cfg, nil + } + // Fall back: decode into a map and coerce headers from string to + // map[string]string if necessary. + var generic map[string]any + if err := json.Unmarshal(raw, &generic); err != nil { + return nil, fmt.Errorf("decode HTTPLogger: %w", err) + } + if v, ok := generic["headers"]; ok { + switch hv := v.(type) { + case string: + var m map[string]string + if hv != "" { + if err := json.Unmarshal([]byte(hv), &m); err != nil { + return nil, fmt.Errorf("decode HTTPLogger headers string: %w", err) + } + } + generic["headers"] = m + case map[string]any: + m := make(map[string]string, len(hv)) + for k, val := range hv { + m[k] = fmt.Sprintf("%v", val) + } + generic["headers"] = m + case nil: + generic["headers"] = map[string]string{} + } + } + // Re-marshal and unmarshal into the typed struct. + normalized, err := json.Marshal(generic) + if err != nil { + return nil, fmt.Errorf("decode HTTPLogger re-marshal: %w", err) + } + if err := json.Unmarshal(normalized, &cfg); err != nil { + return nil, fmt.Errorf("decode HTTPLogger: %w", err) + } + return &cfg, nil +} + // ErrSinkNotFound is returned when a sink row is not found. var ErrSinkNotFound = errors.New("log sink not found") @@ -881,6 +954,11 @@ func (m *LogSinksManager) configRowForType(cfg *config.YAMLConfigurationLogger, return nil, nil } return m.marshalConfigRow(name, typ, order, envID, cfg.Elastic) + case config.LoggingHTTP: + if cfg.HTTP == nil { + return nil, nil + } + return m.marshalConfigRow(name, typ, order, envID, cfg.HTTP) default: return nil, fmt.Errorf("%w: %q", ErrInvalidSinkType, typ) } diff --git a/pkg/logsinks/logsinks_test.go b/pkg/logsinks/logsinks_test.go index 4fe1fad0a..d7dec429a 100644 --- a/pkg/logsinks/logsinks_test.go +++ b/pkg/logsinks/logsinks_test.go @@ -33,7 +33,7 @@ func TestRegistryCoversAllYAMLTypes(t *testing.T) { config.LoggingNone, config.LoggingStdout, config.LoggingFile, config.LoggingDB, config.LoggingGraylog, config.LoggingSplunk, config.LoggingLogstash, config.LoggingKinesis, config.LoggingS3, - config.LoggingKafka, config.LoggingElastic, + config.LoggingKafka, config.LoggingElastic, config.LoggingHTTP, } for _, w := range want { if _, ok := Registry[w]; !ok { @@ -798,3 +798,133 @@ func TestDisabledSinkNotCounted(t *testing.T) { t.Errorf("stats sink ID: got %d want 2", stats[0].SinkID) } } + +func TestRegistryHTTPSink(t *testing.T) { + spec, ok := Registry[config.LoggingHTTP] + if !ok { + t.Fatalf("Registry missing %q", config.LoggingHTTP) + } + if spec.Type != config.LoggingHTTP { + t.Errorf("spec type: got %q want %q", spec.Type, config.LoggingHTTP) + } + if spec.HasSecret { + t.Error("HTTP sink should not have secrets") + } + // Verify the Decode func is set. + if spec.Decode == nil { + t.Error("Decode func is nil") + } + // Verify the Build func is set. + if spec.Build == nil { + t.Error("Build func is nil") + } +} + +func TestHTTPDecodeHeadersAsObject(t *testing.T) { + cfgJSON := `{"url":"http://x","method":"POST","headers":{"Authorization":"Bearer t"}}` + decoded, err := decodeHTTPConfig(json.RawMessage(cfgJSON)) + if err != nil { + t.Fatalf("decode: %v", err) + } + cfg, ok := decoded.(*config.HTTPLogger) + if !ok { + t.Fatalf("decoded type: got %T want *config.HTTPLogger", decoded) + } + if cfg.URL != "http://x" { + t.Errorf("url: got %q", cfg.URL) + } + if cfg.Headers["Authorization"] != "Bearer t" { + t.Errorf("headers: got %v", cfg.Headers) + } +} + +func TestHTTPDecodeHeadersAsString(t *testing.T) { + // The frontend sends headers as a JSON string since the field is + // rendered as a single-line text input. + cfgJSON := `{"url":"http://x","headers":"{\"X-Token\":\"abc\"}"}` + decoded, err := decodeHTTPConfig(json.RawMessage(cfgJSON)) + if err != nil { + t.Fatalf("decode: %v", err) + } + cfg, ok := decoded.(*config.HTTPLogger) + if !ok { + t.Fatalf("decoded type: got %T want *config.HTTPLogger", decoded) + } + if cfg.Headers["X-Token"] != "abc" { + t.Errorf("headers from string: got %v", cfg.Headers) + } +} + +func TestHTTPDecodeEmptyHeaders(t *testing.T) { + cfgJSON := `{"url":"http://x"}` + decoded, err := decodeHTTPConfig(json.RawMessage(cfgJSON)) + if err != nil { + t.Fatalf("decode: %v", err) + } + cfg := decoded.(*config.HTTPLogger) + if len(cfg.Headers) != 0 { + t.Errorf("headers should be nil/empty, got %v", cfg.Headers) + } +} + +func TestHTTPDecodeDefaultsWhenEmpty(t *testing.T) { + decoded, err := decodeHTTPConfig(json.RawMessage(`null`)) + if err != nil { + t.Fatalf("decode null: %v", err) + } + cfg := decoded.(*config.HTTPLogger) + if cfg.URL != "" { + t.Errorf("null config url: got %q want empty", cfg.URL) + } +} + +func TestHTTPBuildExporter(t *testing.T) { + decoded, _ := decodeHTTPConfig(json.RawMessage(`{"url":"http://localhost:9999"}`)) + exp, err := Registry[config.LoggingHTTP].Build(decoded, nil) + if err != nil { + t.Fatalf("build: %v", err) + } + if exp.Name() != config.LoggingHTTP { + t.Errorf("exporter name: got %q want %q", exp.Name(), config.LoggingHTTP) + } + if !exp.IsEnabled() { + t.Error("exporter should be enabled") + } + if err := exp.Close(); err != nil { + t.Errorf("close: %v", err) + } +} + +func TestSeedHTTPSink(t *testing.T) { + m := newTestManager(t) + params := &config.ServiceParameters{ + Logger: &config.YAMLConfigurationLogger{ + Types: []string{config.LoggingHTTP}, + HTTP: &config.HTTPLogger{ + URL: "https://logs.example.com/ingest", + Method: "POST", + }, + }, + DB: &config.YAMLConfigurationDB{Type: "sqlite", FilePath: t.TempDir() + "/p.db"}, + } + if err := m.Seed(params, 0); err != nil { + t.Fatalf("seed: %v", err) + } + rows, err := m.ListByEnvironment(0) + if err != nil { + t.Fatal(err) + } + if len(rows) != 1 { + t.Fatalf("seed count: got %d want 1", len(rows)) + } + if rows[0].Type != config.LoggingHTTP { + t.Errorf("seed type: got %q want %q", rows[0].Type, config.LoggingHTTP) + } + var cfg config.HTTPLogger + if err := json.Unmarshal([]byte(rows[0].Config), &cfg); err != nil { + t.Fatalf("unmarshal: %v", err) + } + if cfg.URL != "https://logs.example.com/ingest" { + t.Errorf("seed config url: got %q", cfg.URL) + } +}