feat: add vector-cron-sink for scheduled cron task output

Scheduled cron task output previously reached only the `dokku` user's cron mail, and could not be redirected because `app.json` rejects bare shell operators in a cron `command`. Setting `vector-cron-sink` on an app or globally routes that output to a dedicated sink instead, on both the `docker-local` and `k3s` schedulers, which keeps log destinations under operator control rather than in a deployed repository. Cron events carry `dokku_app` and `dokku_cron_id` fields so a sink can give each task its own destination. This also fixes a `k3s` bug where configuring a global `vector-sink` silently removed the vector prometheus exporter sink.
This commit is contained in:
Jose Diaz-Gonzalez
2026-08-09 00:54:09 -04:00
parent 047be485a2
commit 52b26a3760
14 changed files with 1346 additions and 112 deletions

View File

@@ -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)
}
}

View File

@@ -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

View 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")
}
}

View File

@@ -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
View 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"])
}
}

View File

@@ -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")
}

View File

@@ -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
}

View File

@@ -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 {

View File

@@ -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 {

View 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")
}
}