mirror of
https://github.com/dokku/dokku.git
synced 2026-08-29 10:08:53 +02:00
Merge pull request #8914 from dokku/feat/vector-cron-sink
Add vector-cron-sink for scheduled cron task output
This commit is contained in:
@@ -100,3 +100,26 @@ func TestDokkuRunCommandAltCommandWithLogFile(t *testing.T) {
|
||||
t.Errorf("DokkuRunCommand() = %q, want %q", got, want)
|
||||
}
|
||||
}
|
||||
|
||||
// TestDokkuRunCommandAppTaskIgnoresLogFile pins that LogFile is honored only
|
||||
// for internally injected tasks from the cron-entries trigger. App tasks never
|
||||
// interpolate a path into the crontab line - their output is shipped by the
|
||||
// vector integration instead.
|
||||
func TestDokkuRunCommandAppTaskIgnoresLogFile(t *testing.T) {
|
||||
task := CronTask{
|
||||
App: "myapp",
|
||||
ID: "abc123",
|
||||
Command: "npm run send-email",
|
||||
LogFile: "/var/log/dokku/should-not-appear.log",
|
||||
}
|
||||
|
||||
got := task.DokkuRunCommand()
|
||||
want := "dokku cron:run myapp abc123"
|
||||
if got != want {
|
||||
t.Errorf("DokkuRunCommand() = %q, want %q", got, want)
|
||||
}
|
||||
|
||||
if strings.Contains(got, ">>") {
|
||||
t.Errorf("DokkuRunCommand() interpolated a redirect into an app task line: %q", got)
|
||||
}
|
||||
}
|
||||
|
||||
@@ -15,7 +15,13 @@ import (
|
||||
|
||||
type vectorConfig struct {
|
||||
Sources map[string]vectorSource `json:"sources"`
|
||||
Sinks map[string]VectorSink `json:"sinks"`
|
||||
|
||||
// Transforms is left nil unless a cron sink is configured, so that configs
|
||||
// without one marshal identically to those generated before cron routing
|
||||
// existed
|
||||
Transforms map[string]any `json:"transforms,omitempty"`
|
||||
|
||||
Sinks map[string]VectorSink `json:"sinks"`
|
||||
}
|
||||
|
||||
type vectorSource struct {
|
||||
@@ -23,6 +29,24 @@ type vectorSource struct {
|
||||
IncludeLabels []string `json:"include_labels,omitempty"`
|
||||
}
|
||||
|
||||
type vectorRouteTransform struct {
|
||||
Type string `json:"type"`
|
||||
Inputs []string `json:"inputs"`
|
||||
RerouteUnmatched bool `json:"reroute_unmatched"`
|
||||
Route map[string]vectorCondition `json:"route"`
|
||||
}
|
||||
|
||||
type vectorCondition struct {
|
||||
Type string `json:"type"`
|
||||
Source string `json:"source"`
|
||||
}
|
||||
|
||||
type vectorRemapTransform struct {
|
||||
Type string `json:"type"`
|
||||
Inputs []string `json:"inputs"`
|
||||
Source string `json:"source"`
|
||||
}
|
||||
|
||||
type vectorTemplateData struct {
|
||||
DokkuLibRoot string
|
||||
DokkuLogsDir string
|
||||
@@ -187,45 +211,121 @@ func stopVectorContainer() error {
|
||||
})
|
||||
}
|
||||
|
||||
func writeVectorConfig() error {
|
||||
apps, _ := common.UnfilteredDokkuApps()
|
||||
// vectorAppSinks holds the resolved sink configuration for a single app or for
|
||||
// the global scope
|
||||
type vectorAppSinks struct {
|
||||
// SourceID is the vector source component id
|
||||
SourceID string
|
||||
|
||||
// IncludeLabels is the docker_logs label filter for the source
|
||||
IncludeLabels []string
|
||||
|
||||
// SinkID is the component id for the sink receiving non-cron logs
|
||||
SinkID string
|
||||
|
||||
// CronSinkID is the component id for the sink receiving cron task logs
|
||||
CronSinkID string
|
||||
|
||||
// RouterID is the component id for the route transform splitting the source
|
||||
RouterID string
|
||||
|
||||
// CronRemapID is the component id for the remap transform on the cron branch
|
||||
CronRemapID string
|
||||
|
||||
// Sink is the DSN for non-cron logs, empty when unset
|
||||
Sink string
|
||||
|
||||
// CronSink is the DSN for cron task logs, empty when unset
|
||||
CronSink string
|
||||
}
|
||||
|
||||
// cronRouteTransforms returns the route and remap pair that splits a source
|
||||
// into cron and non-cron branches. The remap flattens the cron labels into
|
||||
// top-level fields because vector drops any event whose sink template
|
||||
// references a missing field, and a nested quoted path is awkward to template.
|
||||
//
|
||||
// The route carries a single condition, so an event either matches it or falls
|
||||
// through to the reserved _unmatched output. A second route added later would
|
||||
// need a mutually exclusive condition, since route fans out to every match.
|
||||
func cronRouteTransforms(routerID string, remapID string, sourceID string, hasSink bool) map[string]any {
|
||||
return map[string]any{
|
||||
routerID: vectorRouteTransform{
|
||||
Type: "route",
|
||||
Inputs: []string{sourceID},
|
||||
RerouteUnmatched: hasSink,
|
||||
Route: map[string]vectorCondition{
|
||||
CronRouteName: {
|
||||
Type: "vrl",
|
||||
Source: fmt.Sprintf("%s == %q", vrlLabelPath(ContainerTypeLabel), CronContainerType),
|
||||
},
|
||||
},
|
||||
},
|
||||
remapID: vectorRemapTransform{
|
||||
Type: "remap",
|
||||
Inputs: []string{fmt.Sprintf("%s.%s", routerID, CronRouteName)},
|
||||
Source: fmt.Sprintf(".dokku_app = to_string(%s) ?? \"\"\n.dokku_cron_id = to_string(%s) ?? \"\"",
|
||||
vrlLabelPath(AppLabelAlias), vrlLabelPath(CronIDLabel)),
|
||||
},
|
||||
}
|
||||
}
|
||||
|
||||
// vrlLabelPath renders a docker label lookup as a VRL path. Segments holding
|
||||
// characters outside [A-Za-z0-9_] must be double quoted.
|
||||
func vrlLabelPath(label string) string {
|
||||
return fmt.Sprintf(".label.%q", label)
|
||||
}
|
||||
|
||||
// buildVectorConfig assembles the vector configuration for the supplied scopes.
|
||||
// It performs no IO so that the generated shape can be asserted directly.
|
||||
func buildVectorConfig(scopes []vectorAppSinks) (vectorConfig, error) {
|
||||
data := vectorConfig{
|
||||
Sources: map[string]vectorSource{},
|
||||
Sinks: map[string]VectorSink{},
|
||||
}
|
||||
for _, appName := range apps {
|
||||
value := common.PropertyGet("logs", appName, "vector-sink")
|
||||
if value == "" {
|
||||
|
||||
for _, scope := range scopes {
|
||||
if scope.Sink == "" && scope.CronSink == "" {
|
||||
continue
|
||||
}
|
||||
|
||||
inflectedAppName := strings.ReplaceAll(appName, ".", "-")
|
||||
sink, err := SinkValueToConfig(inflectedAppName, value)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
data.Sources[fmt.Sprintf("docker-source:%s", inflectedAppName)] = vectorSource{
|
||||
data.Sources[scope.SourceID] = vectorSource{
|
||||
Type: "docker_logs",
|
||||
IncludeLabels: []string{fmt.Sprintf("%s=%s", reportComputedAppLabelAlias(appName), appName)},
|
||||
IncludeLabels: scope.IncludeLabels,
|
||||
}
|
||||
|
||||
data.Sinks[fmt.Sprintf("docker-sink:%s", inflectedAppName)] = sink
|
||||
}
|
||||
sinkInputs := []string{scope.SourceID}
|
||||
if scope.CronSink != "" {
|
||||
if data.Transforms == nil {
|
||||
data.Transforms = map[string]any{}
|
||||
}
|
||||
for id, transform := range cronRouteTransforms(scope.RouterID, scope.CronRemapID, scope.SourceID, scope.Sink != "") {
|
||||
data.Transforms[id] = transform
|
||||
}
|
||||
|
||||
value := common.PropertyGet("logs", "--global", "vector-sink")
|
||||
if value != "" {
|
||||
sink, err := SinkValueToConfig("--global", value)
|
||||
if err != nil {
|
||||
return err
|
||||
sinkInputs = []string{fmt.Sprintf("%s._unmatched", scope.RouterID)}
|
||||
|
||||
cronSink, err := SinkValueToConfig(SinkValueToConfigInput{
|
||||
SinkValue: scope.CronSink,
|
||||
Inputs: []string{scope.CronRemapID},
|
||||
})
|
||||
if err != nil {
|
||||
return data, err
|
||||
}
|
||||
|
||||
data.Sinks[scope.CronSinkID] = cronSink
|
||||
}
|
||||
|
||||
data.Sources["docker-global-source"] = vectorSource{
|
||||
Type: "docker_logs",
|
||||
IncludeLabels: []string{reportComputedAppLabelAlias("global")},
|
||||
}
|
||||
if scope.Sink != "" {
|
||||
sink, err := SinkValueToConfig(SinkValueToConfigInput{
|
||||
SinkValue: scope.Sink,
|
||||
Inputs: sinkInputs,
|
||||
})
|
||||
if err != nil {
|
||||
return data, err
|
||||
}
|
||||
|
||||
data.Sinks["docker-global-sink"] = sink
|
||||
data.Sinks[scope.SinkID] = sink
|
||||
}
|
||||
}
|
||||
|
||||
if len(data.Sources) == 0 {
|
||||
@@ -238,14 +338,56 @@ func writeVectorConfig() error {
|
||||
|
||||
if len(data.Sinks) == 0 {
|
||||
// write logs to a blackhole
|
||||
sink, err := SinkValueToConfig("--null", VectorDefaultSink)
|
||||
sink, err := SinkValueToConfig(SinkValueToConfigInput{
|
||||
SinkValue: VectorDefaultSink,
|
||||
Inputs: []string{"docker-null-source"},
|
||||
})
|
||||
if err != nil {
|
||||
return err
|
||||
return data, err
|
||||
}
|
||||
|
||||
data.Sinks["docker-null-sink"] = sink
|
||||
}
|
||||
|
||||
return data, nil
|
||||
}
|
||||
|
||||
// vectorScopes collects the sink configuration for every app plus the global scope
|
||||
func vectorScopes() []vectorAppSinks {
|
||||
apps, _ := common.UnfilteredDokkuApps()
|
||||
scopes := []vectorAppSinks{}
|
||||
for _, appName := range apps {
|
||||
inflectedAppName := strings.ReplaceAll(appName, ".", "-")
|
||||
scopes = append(scopes, vectorAppSinks{
|
||||
SourceID: fmt.Sprintf("docker-source:%s", inflectedAppName),
|
||||
IncludeLabels: []string{fmt.Sprintf("%s=%s", reportComputedAppLabelAlias(appName), appName)},
|
||||
SinkID: fmt.Sprintf("docker-sink:%s", inflectedAppName),
|
||||
CronSinkID: fmt.Sprintf("docker-cron-sink:%s", inflectedAppName),
|
||||
RouterID: fmt.Sprintf("docker-router:%s", inflectedAppName),
|
||||
CronRemapID: fmt.Sprintf("docker-cron-remap:%s", inflectedAppName),
|
||||
Sink: common.PropertyGet("logs", appName, "vector-sink"),
|
||||
CronSink: common.PropertyGet("logs", appName, "vector-cron-sink"),
|
||||
})
|
||||
}
|
||||
|
||||
return append(scopes, vectorAppSinks{
|
||||
SourceID: "docker-global-source",
|
||||
IncludeLabels: []string{reportComputedAppLabelAlias("global")},
|
||||
SinkID: "docker-global-sink",
|
||||
CronSinkID: "docker-global-cron-sink",
|
||||
RouterID: "docker-global-router",
|
||||
CronRemapID: "docker-global-cron-remap",
|
||||
Sink: common.PropertyGet("logs", "--global", "vector-sink"),
|
||||
CronSink: common.PropertyGet("logs", "--global", "vector-cron-sink"),
|
||||
})
|
||||
}
|
||||
|
||||
func writeVectorConfig() error {
|
||||
data, err := buildVectorConfig(vectorScopes())
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
b, err := json.MarshalIndent(data, "", " ")
|
||||
if err != nil {
|
||||
return err
|
||||
|
||||
174
plugins/logs/functions_test.go
Normal file
174
plugins/logs/functions_test.go
Normal file
@@ -0,0 +1,174 @@
|
||||
package logs
|
||||
|
||||
import (
|
||||
"encoding/json"
|
||||
"strings"
|
||||
"testing"
|
||||
)
|
||||
|
||||
func appScope(sink string, cronSink string) vectorAppSinks {
|
||||
return vectorAppSinks{
|
||||
SourceID: "docker-source:myapp",
|
||||
IncludeLabels: []string{"com.dokku.app-name=myapp"},
|
||||
SinkID: "docker-sink:myapp",
|
||||
CronSinkID: "docker-cron-sink:myapp",
|
||||
RouterID: "docker-router:myapp",
|
||||
CronRemapID: "docker-cron-remap:myapp",
|
||||
Sink: sink,
|
||||
CronSink: cronSink,
|
||||
}
|
||||
}
|
||||
|
||||
func marshalConfig(t *testing.T, scopes []vectorAppSinks) (string, map[string]interface{}) {
|
||||
t.Helper()
|
||||
|
||||
config, err := buildVectorConfig(scopes)
|
||||
if err != nil {
|
||||
t.Fatalf("buildVectorConfig() error = %v", err)
|
||||
}
|
||||
|
||||
b, err := json.MarshalIndent(config, "", " ")
|
||||
if err != nil {
|
||||
t.Fatalf("MarshalIndent() error = %v", err)
|
||||
}
|
||||
|
||||
var decoded map[string]interface{}
|
||||
if err := json.Unmarshal(b, &decoded); err != nil {
|
||||
t.Fatalf("Unmarshal() error = %v", err)
|
||||
}
|
||||
|
||||
return string(b), decoded
|
||||
}
|
||||
|
||||
func lookup(t *testing.T, decoded map[string]interface{}, path ...string) interface{} {
|
||||
t.Helper()
|
||||
|
||||
var current interface{} = decoded
|
||||
for _, key := range path {
|
||||
asMap, ok := current.(map[string]interface{})
|
||||
if !ok {
|
||||
t.Fatalf("path %v: %q is not a map", path, key)
|
||||
}
|
||||
current, ok = asMap[key]
|
||||
if !ok {
|
||||
t.Fatalf("path %v: missing key %q", path, key)
|
||||
}
|
||||
}
|
||||
|
||||
return current
|
||||
}
|
||||
|
||||
// TestBuildVectorConfigOmitsTransforms is the backwards compatibility guard:
|
||||
// a config without a cron sink must marshal exactly as it did before cron
|
||||
// routing existed, which means no transforms key at all.
|
||||
func TestBuildVectorConfigOmitsTransforms(t *testing.T) {
|
||||
raw, decoded := marshalConfig(t, []vectorAppSinks{appScope("console://?encoding[codec]=json", "")})
|
||||
|
||||
if strings.Contains(raw, "transforms") {
|
||||
t.Errorf("config should not contain a transforms key:\n%s", raw)
|
||||
}
|
||||
|
||||
inputs := lookup(t, decoded, "sinks", "docker-sink:myapp", "inputs")
|
||||
if got := inputs.([]interface{})[0]; got != "docker-source:myapp" {
|
||||
t.Errorf("sink inputs[0] = %v, want docker-source:myapp", got)
|
||||
}
|
||||
}
|
||||
|
||||
func TestBuildVectorConfigCronSinkOnly(t *testing.T) {
|
||||
_, decoded := marshalConfig(t, []vectorAppSinks{appScope("", "console://?encoding[codec]=json")})
|
||||
|
||||
if got := lookup(t, decoded, "transforms", "docker-router:myapp", "type"); got != "route" {
|
||||
t.Errorf("router type = %v, want route", got)
|
||||
}
|
||||
|
||||
// nothing consumes the unmatched output when there is no plain sink
|
||||
if got := lookup(t, decoded, "transforms", "docker-router:myapp", "reroute_unmatched"); got != false {
|
||||
t.Errorf("reroute_unmatched = %v, want false", got)
|
||||
}
|
||||
|
||||
condition := lookup(t, decoded, "transforms", "docker-router:myapp", "route", "cron", "source")
|
||||
want := `.label."com.dokku.container-type" == "cron"`
|
||||
if condition != want {
|
||||
t.Errorf("route condition = %v, want %v", condition, want)
|
||||
}
|
||||
|
||||
remapSource := lookup(t, decoded, "transforms", "docker-cron-remap:myapp", "source").(string)
|
||||
for _, fragment := range []string{".dokku_app", ".dokku_cron_id", `.label."com.dokku.cron-id"`} {
|
||||
if !strings.Contains(remapSource, fragment) {
|
||||
t.Errorf("remap source %q missing %q", remapSource, fragment)
|
||||
}
|
||||
}
|
||||
|
||||
inputs := lookup(t, decoded, "sinks", "docker-cron-sink:myapp", "inputs")
|
||||
if got := inputs.([]interface{})[0]; got != "docker-cron-remap:myapp" {
|
||||
t.Errorf("cron sink inputs[0] = %v, want docker-cron-remap:myapp", got)
|
||||
}
|
||||
|
||||
sinks := lookup(t, decoded, "sinks").(map[string]interface{})
|
||||
if _, ok := sinks["docker-sink:myapp"]; ok {
|
||||
t.Error("plain sink should not exist when only a cron sink is set")
|
||||
}
|
||||
}
|
||||
|
||||
func TestBuildVectorConfigBothSinks(t *testing.T) {
|
||||
_, decoded := marshalConfig(t, []vectorAppSinks{
|
||||
appScope("console://?encoding[codec]=json", "console://?encoding[codec]=text"),
|
||||
})
|
||||
|
||||
if got := lookup(t, decoded, "transforms", "docker-router:myapp", "reroute_unmatched"); got != true {
|
||||
t.Errorf("reroute_unmatched = %v, want true", got)
|
||||
}
|
||||
|
||||
inputs := lookup(t, decoded, "sinks", "docker-sink:myapp", "inputs")
|
||||
if got := inputs.([]interface{})[0]; got != "docker-router:myapp._unmatched" {
|
||||
t.Errorf("plain sink inputs[0] = %v, want docker-router:myapp._unmatched", got)
|
||||
}
|
||||
|
||||
cronInputs := lookup(t, decoded, "sinks", "docker-cron-sink:myapp", "inputs")
|
||||
if got := cronInputs.([]interface{})[0]; got != "docker-cron-remap:myapp" {
|
||||
t.Errorf("cron sink inputs[0] = %v, want docker-cron-remap:myapp", got)
|
||||
}
|
||||
}
|
||||
|
||||
func TestBuildVectorConfigGlobalScope(t *testing.T) {
|
||||
_, decoded := marshalConfig(t, []vectorAppSinks{{
|
||||
SourceID: "docker-global-source",
|
||||
IncludeLabels: []string{"com.dokku.app-name"},
|
||||
SinkID: "docker-global-sink",
|
||||
CronSinkID: "docker-global-cron-sink",
|
||||
RouterID: "docker-global-router",
|
||||
CronRemapID: "docker-global-cron-remap",
|
||||
Sink: "console://?encoding[codec]=json",
|
||||
CronSink: "console://?encoding[codec]=text",
|
||||
}})
|
||||
|
||||
lookup(t, decoded, "transforms", "docker-global-router")
|
||||
lookup(t, decoded, "transforms", "docker-global-cron-remap")
|
||||
lookup(t, decoded, "sinks", "docker-global-cron-sink")
|
||||
|
||||
inputs := lookup(t, decoded, "sinks", "docker-global-sink", "inputs")
|
||||
if got := inputs.([]interface{})[0]; got != "docker-global-router._unmatched" {
|
||||
t.Errorf("global sink inputs[0] = %v, want docker-global-router._unmatched", got)
|
||||
}
|
||||
}
|
||||
|
||||
func TestBuildVectorConfigNoSinks(t *testing.T) {
|
||||
raw, decoded := marshalConfig(t, []vectorAppSinks{appScope("", "")})
|
||||
|
||||
if strings.Contains(raw, "transforms") {
|
||||
t.Errorf("config should not contain a transforms key:\n%s", raw)
|
||||
}
|
||||
|
||||
lookup(t, decoded, "sources", "docker-null-source")
|
||||
|
||||
inputs := lookup(t, decoded, "sinks", "docker-null-sink", "inputs")
|
||||
if got := inputs.([]interface{})[0]; got != "docker-null-source" {
|
||||
t.Errorf("null sink inputs[0] = %v, want docker-null-source", got)
|
||||
}
|
||||
}
|
||||
|
||||
func TestBuildVectorConfigInvalidSink(t *testing.T) {
|
||||
if _, err := buildVectorConfig([]vectorAppSinks{appScope("console://?sinks=nope", "")}); err == nil {
|
||||
t.Fatal("buildVectorConfig() expected an error for an invalid sink DSN")
|
||||
}
|
||||
}
|
||||
@@ -22,21 +22,35 @@ const MaxSize = "10m"
|
||||
// AppLabelAlias is the property key for the app label alias
|
||||
const AppLabelAlias = "com.dokku.app-name"
|
||||
|
||||
// ContainerTypeLabel is the docker label holding the type of a dokku container
|
||||
const ContainerTypeLabel = "com.dokku.container-type"
|
||||
|
||||
// CronContainerType is the ContainerTypeLabel value used for cron task containers
|
||||
const CronContainerType = "cron"
|
||||
|
||||
// CronIDLabel is the docker label holding the cron task id
|
||||
const CronIDLabel = "com.dokku.cron-id"
|
||||
|
||||
// CronRouteName is the vector route output carrying cron task logs
|
||||
const CronRouteName = "cron"
|
||||
|
||||
var (
|
||||
// DefaultProperties is a map of all valid logs properties with corresponding default property values
|
||||
DefaultProperties = map[string]string{
|
||||
"app-label-alias": AppLabelAlias,
|
||||
"max-size": MaxSize,
|
||||
"vector-sink": "",
|
||||
"app-label-alias": AppLabelAlias,
|
||||
"max-size": MaxSize,
|
||||
"vector-cron-sink": "",
|
||||
"vector-sink": "",
|
||||
}
|
||||
|
||||
// GlobalProperties is a map of all valid global logs properties
|
||||
GlobalProperties = map[string]bool{
|
||||
"app-label-alias": true,
|
||||
"max-size": true,
|
||||
"vector-image": true,
|
||||
"vector-networks": true,
|
||||
"vector-sink": true,
|
||||
"app-label-alias": true,
|
||||
"max-size": true,
|
||||
"vector-cron-sink": true,
|
||||
"vector-image": true,
|
||||
"vector-networks": true,
|
||||
"vector-sink": true,
|
||||
}
|
||||
)
|
||||
|
||||
@@ -64,9 +78,21 @@ func GetFailedLogs(appName string) error {
|
||||
return err
|
||||
}
|
||||
|
||||
// SinkValueToConfigInput is the input for the SinkValueToConfig function
|
||||
type SinkValueToConfigInput struct {
|
||||
// SinkValue is the sink DSN to convert
|
||||
SinkValue string
|
||||
|
||||
// Inputs are the vector component ids feeding the sink. When empty, the
|
||||
// inputs key is omitted, which is appropriate for callers that only need
|
||||
// the parsed sink for validation or redaction.
|
||||
Inputs []string
|
||||
}
|
||||
|
||||
// SinkValueToConfig converts a sink DSN value to a VectorSink
|
||||
func SinkValueToConfig(appName string, sinkValue string) (VectorSink, error) {
|
||||
func SinkValueToConfig(input SinkValueToConfigInput) (VectorSink, error) {
|
||||
var data VectorSink
|
||||
sinkValue := input.SinkValue
|
||||
if strings.Contains(sinkValue, "://") {
|
||||
parts := strings.SplitN(sinkValue, "://", 2)
|
||||
parts[0] = strings.ReplaceAll(parts[0], "_", "-")
|
||||
@@ -96,12 +122,8 @@ func SinkValueToConfig(appName string, sinkValue string) (VectorSink, error) {
|
||||
}
|
||||
|
||||
data["type"] = u.Scheme
|
||||
data["inputs"] = []string{"docker-source:" + appName}
|
||||
if appName == "--global" {
|
||||
data["inputs"] = []string{"docker-global-source"}
|
||||
}
|
||||
if appName == "--null" {
|
||||
data["inputs"] = []string{"docker-null-source"}
|
||||
if len(input.Inputs) > 0 {
|
||||
data["inputs"] = input.Inputs
|
||||
}
|
||||
|
||||
// add special support for `base64enc:VAL` fields
|
||||
|
||||
108
plugins/logs/logs_test.go
Normal file
108
plugins/logs/logs_test.go
Normal file
@@ -0,0 +1,108 @@
|
||||
package logs
|
||||
|
||||
import (
|
||||
"encoding/base64"
|
||||
"reflect"
|
||||
"testing"
|
||||
)
|
||||
|
||||
func TestSinkValueToConfigInputs(t *testing.T) {
|
||||
sink, err := SinkValueToConfig(SinkValueToConfigInput{
|
||||
SinkValue: "console://?encoding[codec]=json",
|
||||
Inputs: []string{"docker-source:myapp"},
|
||||
})
|
||||
if err != nil {
|
||||
t.Fatalf("SinkValueToConfig() error = %v", err)
|
||||
}
|
||||
|
||||
if sink["type"] != "console" {
|
||||
t.Errorf("type = %v, want console", sink["type"])
|
||||
}
|
||||
|
||||
want := []string{"docker-source:myapp"}
|
||||
if !reflect.DeepEqual(sink["inputs"], want) {
|
||||
t.Errorf("inputs = %v, want %v", sink["inputs"], want)
|
||||
}
|
||||
}
|
||||
|
||||
func TestSinkValueToConfigOmitsEmptyInputs(t *testing.T) {
|
||||
sink, err := SinkValueToConfig(SinkValueToConfigInput{
|
||||
SinkValue: "console://?encoding[codec]=json",
|
||||
})
|
||||
if err != nil {
|
||||
t.Fatalf("SinkValueToConfig() error = %v", err)
|
||||
}
|
||||
|
||||
if _, ok := sink["inputs"]; ok {
|
||||
t.Errorf("inputs should be omitted when no inputs are supplied, got %v", sink["inputs"])
|
||||
}
|
||||
}
|
||||
|
||||
func TestSinkValueToConfigRejectsSinksOption(t *testing.T) {
|
||||
_, err := SinkValueToConfig(SinkValueToConfigInput{
|
||||
SinkValue: "console://?sinks=nope",
|
||||
})
|
||||
if err == nil {
|
||||
t.Fatal("SinkValueToConfig() expected an error for the sinks option")
|
||||
}
|
||||
}
|
||||
|
||||
func TestSinkValueToConfigSchemeUnderscores(t *testing.T) {
|
||||
for _, scheme := range []string{"datadog_logs", "aws_cloudwatch_logs"} {
|
||||
sink, err := SinkValueToConfig(SinkValueToConfigInput{
|
||||
SinkValue: scheme + "://?api_key=abc123",
|
||||
})
|
||||
if err != nil {
|
||||
t.Fatalf("SinkValueToConfig(%s) error = %v", scheme, err)
|
||||
}
|
||||
|
||||
if sink["type"] != scheme {
|
||||
t.Errorf("type = %v, want %s", sink["type"], scheme)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
func TestSinkValueToConfigBase64Enc(t *testing.T) {
|
||||
encoded := base64.StdEncoding.EncodeToString([]byte("{{ pod }}"))
|
||||
sink, err := SinkValueToConfig(SinkValueToConfigInput{
|
||||
SinkValue: "http://?process=base64enc%3A" + encoded,
|
||||
})
|
||||
if err != nil {
|
||||
t.Fatalf("SinkValueToConfig() error = %v", err)
|
||||
}
|
||||
|
||||
if sink["process"] != "{{ pod }}" {
|
||||
t.Errorf("process = %v, want {{ pod }}", sink["process"])
|
||||
}
|
||||
}
|
||||
|
||||
// TestSinkValueToConfigTemplatedFilePath guards the DSN handling that the
|
||||
// documented cron file sink depends on: a vector template in a query-string
|
||||
// value must survive url.Parse and the qson decode with its braces and
|
||||
// interior spaces intact.
|
||||
func TestSinkValueToConfigTemplatedFilePath(t *testing.T) {
|
||||
sink, err := SinkValueToConfig(SinkValueToConfigInput{
|
||||
SinkValue: "file://?path=/var/log/dokku/apps/myapp/cron-{{ dokku_cron_id }}.log&encoding[codec]=text",
|
||||
Inputs: []string{"docker-cron-remap:myapp"},
|
||||
})
|
||||
if err != nil {
|
||||
t.Fatalf("SinkValueToConfig() error = %v", err)
|
||||
}
|
||||
|
||||
if sink["type"] != "file" {
|
||||
t.Errorf("type = %v, want file", sink["type"])
|
||||
}
|
||||
|
||||
wantPath := "/var/log/dokku/apps/myapp/cron-{{ dokku_cron_id }}.log"
|
||||
if sink["path"] != wantPath {
|
||||
t.Errorf("path = %v, want %v", sink["path"], wantPath)
|
||||
}
|
||||
|
||||
encoding, ok := sink["encoding"].(map[string]interface{})
|
||||
if !ok {
|
||||
t.Fatalf("encoding = %v, want a map", sink["encoding"])
|
||||
}
|
||||
if encoding["codec"] != "text" {
|
||||
t.Errorf("encoding.codec = %v, want text", encoding["codec"])
|
||||
}
|
||||
}
|
||||
@@ -20,32 +20,37 @@ func ReportSingleApp(appName string, format string, infoFlag string) error {
|
||||
var flags map[string]common.ReportFunc
|
||||
if appName == "--global" {
|
||||
flags = map[string]common.ReportFunc{
|
||||
"--logs-computed-app-label-alias": reportComputedAppLabelAlias,
|
||||
"--logs-computed-max-size": reportComputedMaxSize,
|
||||
"--logs-computed-vector-image": reportComputedVectorImage,
|
||||
"--logs-computed-vector-networks": reportComputedVectorNetworks,
|
||||
"--logs-computed-vector-sink": reportComputedVectorSink,
|
||||
"--logs-global-app-label-alias": reportGlobalAppLabelAlias,
|
||||
"--logs-global-max-size": reportGlobalMaxSize,
|
||||
"--logs-global-vector-image": reportGlobalVectorImage,
|
||||
"--logs-global-vector-networks": reportGlobalVectorNetworks,
|
||||
"--logs-global-vector-sink": reportGlobalVectorSink,
|
||||
"--logs-computed-app-label-alias": reportComputedAppLabelAlias,
|
||||
"--logs-computed-max-size": reportComputedMaxSize,
|
||||
"--logs-computed-vector-cron-sink": reportComputedVectorCronSink,
|
||||
"--logs-computed-vector-image": reportComputedVectorImage,
|
||||
"--logs-computed-vector-networks": reportComputedVectorNetworks,
|
||||
"--logs-computed-vector-sink": reportComputedVectorSink,
|
||||
"--logs-global-app-label-alias": reportGlobalAppLabelAlias,
|
||||
"--logs-global-max-size": reportGlobalMaxSize,
|
||||
"--logs-global-vector-cron-sink": reportGlobalVectorCronSink,
|
||||
"--logs-global-vector-image": reportGlobalVectorImage,
|
||||
"--logs-global-vector-networks": reportGlobalVectorNetworks,
|
||||
"--logs-global-vector-sink": reportGlobalVectorSink,
|
||||
}
|
||||
} else {
|
||||
flags = map[string]common.ReportFunc{
|
||||
"--logs-app-label-alias": reportAppLabelAlias,
|
||||
"--logs-computed-app-label-alias": reportComputedAppLabelAlias,
|
||||
"--logs-computed-max-size": reportComputedMaxSize,
|
||||
"--logs-computed-vector-image": reportComputedVectorImage,
|
||||
"--logs-computed-vector-networks": reportComputedVectorNetworks,
|
||||
"--logs-computed-vector-sink": reportComputedVectorSink,
|
||||
"--logs-global-app-label-alias": reportGlobalAppLabelAlias,
|
||||
"--logs-global-max-size": reportGlobalMaxSize,
|
||||
"--logs-global-vector-image": reportGlobalVectorImage,
|
||||
"--logs-global-vector-networks": reportGlobalVectorNetworks,
|
||||
"--logs-global-vector-sink": reportGlobalVectorSink,
|
||||
"--logs-max-size": reportMaxSize,
|
||||
"--logs-vector-sink": reportVectorSink,
|
||||
"--logs-app-label-alias": reportAppLabelAlias,
|
||||
"--logs-computed-app-label-alias": reportComputedAppLabelAlias,
|
||||
"--logs-computed-max-size": reportComputedMaxSize,
|
||||
"--logs-computed-vector-cron-sink": reportComputedVectorCronSink,
|
||||
"--logs-computed-vector-image": reportComputedVectorImage,
|
||||
"--logs-computed-vector-networks": reportComputedVectorNetworks,
|
||||
"--logs-computed-vector-sink": reportComputedVectorSink,
|
||||
"--logs-global-app-label-alias": reportGlobalAppLabelAlias,
|
||||
"--logs-global-max-size": reportGlobalMaxSize,
|
||||
"--logs-global-vector-cron-sink": reportGlobalVectorCronSink,
|
||||
"--logs-global-vector-image": reportGlobalVectorImage,
|
||||
"--logs-global-vector-networks": reportGlobalVectorNetworks,
|
||||
"--logs-global-vector-sink": reportGlobalVectorSink,
|
||||
"--logs-max-size": reportMaxSize,
|
||||
"--logs-vector-cron-sink": reportVectorCronSink,
|
||||
"--logs-vector-sink": reportVectorSink,
|
||||
}
|
||||
}
|
||||
|
||||
@@ -129,7 +134,29 @@ func reportComputedVectorSink(appName string) string {
|
||||
}
|
||||
|
||||
func reportGlobalVectorSink(appName string) string {
|
||||
value := common.PropertyGet("logs", "--global", "vector-sink")
|
||||
return redactedSink(common.PropertyGet("logs", "--global", "vector-sink"), "--logs-global-vector-sink")
|
||||
}
|
||||
|
||||
func reportComputedVectorCronSink(appName string) string {
|
||||
value := reportVectorCronSink(appName)
|
||||
if value == "" {
|
||||
value = reportGlobalVectorCronSink(appName)
|
||||
}
|
||||
return value
|
||||
}
|
||||
|
||||
func reportGlobalVectorCronSink(appName string) string {
|
||||
return redactedSink(common.PropertyGet("logs", "--global", "vector-cron-sink"), "--logs-global-vector-cron-sink")
|
||||
}
|
||||
|
||||
func reportVectorCronSink(appName string) string {
|
||||
return redactedSink(common.PropertyGet("logs", appName, "vector-cron-sink"), "--logs-vector-cron-sink")
|
||||
}
|
||||
|
||||
// redactedSink returns the sink value as-is for json reports or when the exact
|
||||
// flag was requested, and otherwise reduces it to its scheme so that
|
||||
// credentials embedded in the DSN are not printed in a general report
|
||||
func redactedSink(value string, exactFlag string) string {
|
||||
if value == "" {
|
||||
return value
|
||||
}
|
||||
@@ -138,12 +165,11 @@ func reportGlobalVectorSink(appName string) string {
|
||||
return value
|
||||
}
|
||||
|
||||
if os.Getenv("DOKKU_REPORT_FLAG") == "--logs-global-vector-sink" {
|
||||
if os.Getenv("DOKKU_REPORT_FLAG") == exactFlag {
|
||||
return value
|
||||
}
|
||||
|
||||
// only show the schema and sanitize the rest
|
||||
sink, err := SinkValueToConfig("--global", value)
|
||||
sink, err := SinkValueToConfig(SinkValueToConfigInput{SinkValue: value})
|
||||
if err != nil {
|
||||
return ""
|
||||
}
|
||||
@@ -156,24 +182,5 @@ func reportMaxSize(appName string) string {
|
||||
}
|
||||
|
||||
func reportVectorSink(appName string) string {
|
||||
value := common.PropertyGet("logs", appName, "vector-sink")
|
||||
if value == "" {
|
||||
return value
|
||||
}
|
||||
|
||||
if os.Getenv("DOKKU_REPORT_FORMAT") != "stdout" {
|
||||
return value
|
||||
}
|
||||
|
||||
if os.Getenv("DOKKU_REPORT_FLAG") == "--logs-vector-sink" {
|
||||
return value
|
||||
}
|
||||
|
||||
// only show the schema and sanitize the rest
|
||||
sink, err := SinkValueToConfig(appName, value)
|
||||
if err != nil {
|
||||
return ""
|
||||
}
|
||||
|
||||
return fmt.Sprintf("%s://redacted", sink["type"])
|
||||
return redactedSink(common.PropertyGet("logs", appName, "vector-sink"), "--logs-vector-sink")
|
||||
}
|
||||
|
||||
@@ -22,7 +22,7 @@ func validateSetValue(appName string, key string, value string) error {
|
||||
return validateVectorNetworks(appName, value)
|
||||
}
|
||||
|
||||
if key == "vector-sink" {
|
||||
if key == "vector-sink" || key == "vector-cron-sink" {
|
||||
return validateVectorSink(appName, value)
|
||||
}
|
||||
|
||||
@@ -60,7 +60,7 @@ func validateVectorSink(appName string, value string) error {
|
||||
return nil
|
||||
}
|
||||
|
||||
_, err := SinkValueToConfig(appName, value)
|
||||
_, err := SinkValueToConfig(SinkValueToConfigInput{SinkValue: value})
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
@@ -78,8 +78,9 @@ func CommandSet(appName string, property string, value string) error {
|
||||
common.CommandPropertySet("logs", appName, property, value, DefaultProperties, GlobalProperties)
|
||||
|
||||
vectorProperties := map[string]bool{
|
||||
"app-label-alias": true,
|
||||
"vector-sink": true,
|
||||
"app-label-alias": true,
|
||||
"vector-cron-sink": true,
|
||||
"vector-sink": true,
|
||||
}
|
||||
|
||||
if _, ok := vectorProperties[property]; ok {
|
||||
|
||||
@@ -2175,7 +2175,10 @@ func installHelmCharts(ctx context.Context, clientset KubernetesClient, shouldIn
|
||||
}
|
||||
|
||||
if chart.ReleaseName == "vector" && chart.Namespace == "vector" {
|
||||
values = updateVectorValues(values)
|
||||
values, err = updateVectorValues(values)
|
||||
if err != nil {
|
||||
return fmt.Errorf("Error updating vector values: %w", err)
|
||||
}
|
||||
}
|
||||
|
||||
chartProperties, err := common.PropertyMapGet("scheduler-k3s", "--global", "chart-overrides."+chart.ReleaseName)
|
||||
@@ -2256,26 +2259,129 @@ func installHelmCharts(ctx context.Context, clientset KubernetesClient, shouldIn
|
||||
return nil
|
||||
}
|
||||
|
||||
func updateVectorValues(values map[string]interface{}) map[string]interface{} {
|
||||
value := common.PropertyGet("logs", "--global", "vector-sink")
|
||||
if value == "" {
|
||||
return values
|
||||
const (
|
||||
// kubernetesLogsTransform is the base transform enriching every container log event
|
||||
kubernetesLogsTransform = "kubernetes_container_logs"
|
||||
|
||||
// kubernetesRouterTransform splits container logs into cron and non-cron branches
|
||||
kubernetesRouterTransform = "kubernetes_router"
|
||||
|
||||
// kubernetesCronRemapTransform flattens cron metadata onto the event
|
||||
kubernetesCronRemapTransform = "kubernetes_cron_remap"
|
||||
|
||||
// kubernetesDefaultSink is the console sink shipped in the base values file
|
||||
kubernetesDefaultSink = "default_global_sink"
|
||||
|
||||
// kubernetesGlobalSink receives non-cron logs when a vector-sink is configured
|
||||
kubernetesGlobalSink = "kubernetes_global_sink"
|
||||
|
||||
// kubernetesCronSink receives cron task logs when a vector-cron-sink is configured
|
||||
kubernetesCronSink = "kubernetes_cron_sink"
|
||||
)
|
||||
|
||||
// kubernetesCronRouteTransforms returns the route and remap pair that splits
|
||||
// container logs into cron and non-cron branches.
|
||||
//
|
||||
// Cron pods are labelled app.kubernetes.io/name=cron. The raw cron id exceeds
|
||||
// the Kubernetes label length cap, so it lives in an annotation and is
|
||||
// flattened onto the event here - vector drops any event whose sink template
|
||||
// references a missing field.
|
||||
func kubernetesCronRouteTransforms() map[string]interface{} {
|
||||
return map[string]interface{}{
|
||||
kubernetesRouterTransform: map[string]interface{}{
|
||||
"type": "route",
|
||||
"inputs": []string{kubernetesLogsTransform},
|
||||
"reroute_unmatched": true,
|
||||
"route": map[string]interface{}{
|
||||
logs.CronRouteName: map[string]interface{}{
|
||||
"type": "vrl",
|
||||
"source": `.kubernetes.pod_labels."app.kubernetes.io/name" == "cron"`,
|
||||
},
|
||||
},
|
||||
},
|
||||
kubernetesCronRemapTransform: map[string]interface{}{
|
||||
"type": "remap",
|
||||
"inputs": []string{kubernetesRouterTransform + "." + logs.CronRouteName},
|
||||
"source": ".dokku_app = to_string(.kubernetes.pod_labels.\"app.kubernetes.io/part-of\") ?? \"\"\n" +
|
||||
".dokku_cron_id = to_string(.kubernetes.pod_annotations.\"dokku.com/cron-id\") ?? \"\"",
|
||||
},
|
||||
}
|
||||
}
|
||||
|
||||
// updateVectorValues layers the configured log sinks onto the vector chart
|
||||
// values. Sinks and transforms from the base values file are preserved unless
|
||||
// they are explicitly superseded, so unrelated components such as the
|
||||
// prometheus exporter keep working.
|
||||
func updateVectorValues(values map[string]interface{}) (map[string]interface{}, error) {
|
||||
sinkValue := common.PropertyGet("logs", "--global", "vector-sink")
|
||||
cronSinkValue := common.PropertyGet("logs", "--global", "vector-cron-sink")
|
||||
if sinkValue == "" && cronSinkValue == "" {
|
||||
return values, nil
|
||||
}
|
||||
|
||||
sink, err := logs.SinkValueToConfig("--global", value)
|
||||
if err != nil {
|
||||
return nil
|
||||
customConfig, ok := values["customConfig"].(map[string]interface{})
|
||||
if !ok {
|
||||
return values, errors.New("Missing or invalid customConfig in vector chart values")
|
||||
}
|
||||
|
||||
sink["inputs"] = []string{"kubernetes_container_logs"}
|
||||
|
||||
sinkMap := map[string]interface{}{
|
||||
"kubernetes_global_sink": sink,
|
||||
sinks, ok := customConfig["sinks"].(map[string]interface{})
|
||||
if !ok {
|
||||
sinks = map[string]interface{}{}
|
||||
}
|
||||
|
||||
values["customConfig"].(map[string]interface{})["sinks"] = sinkMap
|
||||
// a configured sink supersedes the console sink from the base values;
|
||||
// without one, the console sink remains the destination for non-cron logs
|
||||
nonCronSinkID := kubernetesDefaultSink
|
||||
if sinkValue != "" {
|
||||
nonCronSinkID = kubernetesGlobalSink
|
||||
delete(sinks, kubernetesDefaultSink)
|
||||
|
||||
return values
|
||||
sink, err := logs.SinkValueToConfig(logs.SinkValueToConfigInput{
|
||||
SinkValue: sinkValue,
|
||||
Inputs: []string{kubernetesLogsTransform},
|
||||
})
|
||||
if err != nil {
|
||||
return values, fmt.Errorf("Error parsing vector-sink: %w", err)
|
||||
}
|
||||
|
||||
// stored as a plain map so that later lookups, and the yaml encoder,
|
||||
// see the same shape as the sinks parsed from the base values file
|
||||
sinks[kubernetesGlobalSink] = map[string]interface{}(sink)
|
||||
}
|
||||
|
||||
if cronSinkValue != "" {
|
||||
transforms, ok := customConfig["transforms"].(map[string]interface{})
|
||||
if !ok {
|
||||
transforms = map[string]interface{}{}
|
||||
}
|
||||
for id, transform := range kubernetesCronRouteTransforms() {
|
||||
transforms[id] = transform
|
||||
}
|
||||
customConfig["transforms"] = transforms
|
||||
|
||||
cronSink, err := logs.SinkValueToConfig(logs.SinkValueToConfigInput{
|
||||
SinkValue: cronSinkValue,
|
||||
Inputs: []string{kubernetesCronRemapTransform},
|
||||
})
|
||||
if err != nil {
|
||||
return values, fmt.Errorf("Error parsing vector-cron-sink: %w", err)
|
||||
}
|
||||
sinks[kubernetesCronSink] = map[string]interface{}(cronSink)
|
||||
|
||||
// cron logs are routed away from the non-cron sink so that each event
|
||||
// lands in exactly one destination
|
||||
nonCronSink, ok := sinks[nonCronSinkID].(map[string]interface{})
|
||||
if !ok {
|
||||
return values, fmt.Errorf("Missing or invalid %s sink in vector chart values", nonCronSinkID)
|
||||
}
|
||||
nonCronSink["inputs"] = []string{kubernetesRouterTransform + "._unmatched"}
|
||||
sinks[nonCronSinkID] = nonCronSink
|
||||
}
|
||||
|
||||
customConfig["sinks"] = sinks
|
||||
values["customConfig"] = customConfig
|
||||
|
||||
return values, nil
|
||||
}
|
||||
|
||||
func installHelperCommands(ctx context.Context) error {
|
||||
|
||||
208
plugins/scheduler-k3s/vector_values_test.go
Normal file
208
plugins/scheduler-k3s/vector_values_test.go
Normal file
@@ -0,0 +1,208 @@
|
||||
package scheduler_k3s
|
||||
|
||||
import (
|
||||
"testing"
|
||||
|
||||
"github.com/dokku/dokku/plugins/common"
|
||||
"gopkg.in/yaml.v3"
|
||||
)
|
||||
|
||||
// baseVectorValues mirrors the shipped templates/helm-config/vector.yaml
|
||||
// closely enough to exercise the sink and transform merging
|
||||
const baseVectorValues = `
|
||||
customConfig:
|
||||
sources:
|
||||
kubernetes_logs:
|
||||
type: kubernetes_logs
|
||||
internal_metrics:
|
||||
type: internal_metrics
|
||||
transforms:
|
||||
kubernetes_container_logs:
|
||||
type: remap
|
||||
inputs:
|
||||
- kubernetes_logs
|
||||
source: |
|
||||
.container = .kubernetes.container_name
|
||||
sinks:
|
||||
default_global_sink:
|
||||
type: console
|
||||
inputs:
|
||||
- kubernetes_container_logs
|
||||
encoding:
|
||||
codec: json
|
||||
prom_exporter:
|
||||
type: prometheus_exporter
|
||||
inputs:
|
||||
- internal_metrics
|
||||
address: 0.0.0.0:9090
|
||||
`
|
||||
|
||||
func setupVectorValuesTest(t *testing.T, sink string, cronSink string) map[string]interface{} {
|
||||
t.Helper()
|
||||
t.Setenv("PLUGIN_PATH", "/var/lib/dokku/plugins")
|
||||
t.Setenv("PLUGIN_ENABLED_PATH", "/var/lib/dokku/plugins/enabled")
|
||||
t.Setenv("DOKKU_LIB_ROOT", t.TempDir())
|
||||
t.Setenv("DOKKU_SYSTEM_USER", "root")
|
||||
t.Setenv("DOKKU_SYSTEM_GROUP", "root")
|
||||
|
||||
if err := common.PropertySetup("logs"); err != nil {
|
||||
t.Fatalf("PropertySetup: %v", err)
|
||||
}
|
||||
if sink != "" {
|
||||
if err := common.PropertyWrite("logs", "--global", "vector-sink", sink); err != nil {
|
||||
t.Fatalf("PropertyWrite vector-sink: %v", err)
|
||||
}
|
||||
}
|
||||
if cronSink != "" {
|
||||
if err := common.PropertyWrite("logs", "--global", "vector-cron-sink", cronSink); err != nil {
|
||||
t.Fatalf("PropertyWrite vector-cron-sink: %v", err)
|
||||
}
|
||||
}
|
||||
|
||||
values := map[string]interface{}{}
|
||||
if err := yaml.Unmarshal([]byte(baseVectorValues), &values); err != nil {
|
||||
t.Fatalf("Unmarshal base values: %v", err)
|
||||
}
|
||||
|
||||
return values
|
||||
}
|
||||
|
||||
func vectorSinks(t *testing.T, values map[string]interface{}) map[string]interface{} {
|
||||
t.Helper()
|
||||
customConfig, ok := values["customConfig"].(map[string]interface{})
|
||||
if !ok {
|
||||
t.Fatal("customConfig is not a map")
|
||||
}
|
||||
sinks, ok := customConfig["sinks"].(map[string]interface{})
|
||||
if !ok {
|
||||
t.Fatal("sinks is not a map")
|
||||
}
|
||||
return sinks
|
||||
}
|
||||
|
||||
func vectorTransforms(t *testing.T, values map[string]interface{}) map[string]interface{} {
|
||||
t.Helper()
|
||||
customConfig, ok := values["customConfig"].(map[string]interface{})
|
||||
if !ok {
|
||||
t.Fatal("customConfig is not a map")
|
||||
}
|
||||
transforms, ok := customConfig["transforms"].(map[string]interface{})
|
||||
if !ok {
|
||||
t.Fatal("transforms is not a map")
|
||||
}
|
||||
return transforms
|
||||
}
|
||||
|
||||
func sinkInput(t *testing.T, sinks map[string]interface{}, sinkID string) string {
|
||||
t.Helper()
|
||||
sink, ok := sinks[sinkID].(map[string]interface{})
|
||||
if !ok {
|
||||
t.Fatalf("sink %q is missing or not a map", sinkID)
|
||||
}
|
||||
|
||||
switch inputs := sink["inputs"].(type) {
|
||||
case []string:
|
||||
return inputs[0]
|
||||
case []interface{}:
|
||||
return inputs[0].(string)
|
||||
default:
|
||||
t.Fatalf("sink %q has unexpected inputs type %T", sinkID, sink["inputs"])
|
||||
return ""
|
||||
}
|
||||
}
|
||||
|
||||
func TestUpdateVectorValuesNoSinks(t *testing.T) {
|
||||
values := setupVectorValuesTest(t, "", "")
|
||||
|
||||
updated, err := updateVectorValues(values)
|
||||
if err != nil {
|
||||
t.Fatalf("updateVectorValues() error = %v", err)
|
||||
}
|
||||
|
||||
sinks := vectorSinks(t, updated)
|
||||
if _, ok := sinks[kubernetesDefaultSink]; !ok {
|
||||
t.Error("default_global_sink should remain when no sink is configured")
|
||||
}
|
||||
if _, ok := sinks[kubernetesGlobalSink]; ok {
|
||||
t.Error("kubernetes_global_sink should not exist when no sink is configured")
|
||||
}
|
||||
}
|
||||
|
||||
// TestUpdateVectorValuesPreservesPromExporter guards against the sinks map
|
||||
// being replaced wholesale, which previously dropped the prometheus exporter
|
||||
// that the chart still exposes on port 9090
|
||||
func TestUpdateVectorValuesPreservesPromExporter(t *testing.T) {
|
||||
values := setupVectorValuesTest(t, "console://?encoding[codec]=json", "")
|
||||
|
||||
updated, err := updateVectorValues(values)
|
||||
if err != nil {
|
||||
t.Fatalf("updateVectorValues() error = %v", err)
|
||||
}
|
||||
|
||||
sinks := vectorSinks(t, updated)
|
||||
if _, ok := sinks["prom_exporter"]; !ok {
|
||||
t.Error("prom_exporter should survive configuring a vector-sink")
|
||||
}
|
||||
if _, ok := sinks[kubernetesDefaultSink]; ok {
|
||||
t.Error("default_global_sink should be superseded by the configured sink")
|
||||
}
|
||||
if got := sinkInput(t, sinks, kubernetesGlobalSink); got != kubernetesLogsTransform {
|
||||
t.Errorf("global sink input = %v, want %v", got, kubernetesLogsTransform)
|
||||
}
|
||||
}
|
||||
|
||||
func TestUpdateVectorValuesCronSinkOnly(t *testing.T) {
|
||||
values := setupVectorValuesTest(t, "", "console://?encoding[codec]=text")
|
||||
|
||||
updated, err := updateVectorValues(values)
|
||||
if err != nil {
|
||||
t.Fatalf("updateVectorValues() error = %v", err)
|
||||
}
|
||||
|
||||
transforms := vectorTransforms(t, updated)
|
||||
if _, ok := transforms[kubernetesRouterTransform]; !ok {
|
||||
t.Error("router transform should be added")
|
||||
}
|
||||
if _, ok := transforms[kubernetesLogsTransform]; !ok {
|
||||
t.Error("base container logs transform should be preserved")
|
||||
}
|
||||
|
||||
sinks := vectorSinks(t, updated)
|
||||
if got := sinkInput(t, sinks, kubernetesCronSink); got != kubernetesCronRemapTransform {
|
||||
t.Errorf("cron sink input = %v, want %v", got, kubernetesCronRemapTransform)
|
||||
}
|
||||
|
||||
// without a configured plain sink, the console default takes the unmatched branch
|
||||
if got := sinkInput(t, sinks, kubernetesDefaultSink); got != kubernetesRouterTransform+"._unmatched" {
|
||||
t.Errorf("default sink input = %v, want %v._unmatched", got, kubernetesRouterTransform)
|
||||
}
|
||||
}
|
||||
|
||||
func TestUpdateVectorValuesBothSinks(t *testing.T) {
|
||||
values := setupVectorValuesTest(t, "console://?encoding[codec]=json", "console://?encoding[codec]=text")
|
||||
|
||||
updated, err := updateVectorValues(values)
|
||||
if err != nil {
|
||||
t.Fatalf("updateVectorValues() error = %v", err)
|
||||
}
|
||||
|
||||
sinks := vectorSinks(t, updated)
|
||||
if got := sinkInput(t, sinks, kubernetesGlobalSink); got != kubernetesRouterTransform+"._unmatched" {
|
||||
t.Errorf("global sink input = %v, want %v._unmatched", got, kubernetesRouterTransform)
|
||||
}
|
||||
if got := sinkInput(t, sinks, kubernetesCronSink); got != kubernetesCronRemapTransform {
|
||||
t.Errorf("cron sink input = %v, want %v", got, kubernetesCronRemapTransform)
|
||||
}
|
||||
if _, ok := sinks["prom_exporter"]; !ok {
|
||||
t.Error("prom_exporter should survive configuring both sinks")
|
||||
}
|
||||
}
|
||||
|
||||
func TestUpdateVectorValuesInvalidCustomConfig(t *testing.T) {
|
||||
values := setupVectorValuesTest(t, "console://", "")
|
||||
delete(values, "customConfig")
|
||||
|
||||
if _, err := updateVectorValues(values); err == nil {
|
||||
t.Fatal("updateVectorValues() expected an error for missing customConfig")
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user