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

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
The table of contents is too big for display.
Diff view
Diff view
  •  
  •  
  •  
The diff you're trying to view is too large. We only load the first 3000 changed files.
160 changes: 159 additions & 1 deletion confgenerator/agentmetrics.go
Original file line number Diff line number Diff line change
Expand Up @@ -22,9 +22,22 @@ import (

"github.com/GoogleCloudPlatform/ops-agent/confgenerator/otel"
"github.com/GoogleCloudPlatform/ops-agent/confgenerator/otel/ottl"
"github.com/GoogleCloudPlatform/ops-agent/internal/version"
)

// AgentSelfMetrics provides the agent.googleapis.com/agent/ metrics.
var (
agentKind string = "ops-agent"
schemaVersion string = "v1"
)

const (
healthLogsTag string = "ops-agent-health"
agentVersionKey string = "agent.googleapis.com/health/agentVersion"
agentKindKey string = "agent.googleapis.com/health/agentKind"
schemaVersionKey string = "agent.googleapis.com/health/schemaVersion"
)

// AgentSelfMetrics provides the agent.googleapis.com/agent/ metric and self logs.
// It is never referenced in the config file, and instead is forcibly added in confgenerator.go.
// Therefore, it does not need to implement any interfaces.
type AgentSelfMetrics struct {
Expand All @@ -34,6 +47,7 @@ type AgentSelfMetrics struct {
OtelPort int
OtelRuntimeDir string
OtlpExporterEnabled bool
LogsDir string
}

// Following reference : https://github.com/googleapis/googleapis/blob/master/google/rpc/code.proto
Expand Down Expand Up @@ -112,6 +126,18 @@ func (r AgentSelfMetrics) AddSelfMetricsPipelines(receiverPipelines map[string]o
Type: "metrics",
ReceiverPipelineName: "ops_agent",
}

receiverPipelines["logging_ping"] = r.LoggingPingPipeline(ctx)
pipelines["loggingping"] = otel.Pipeline{
Type: "logs",
ReceiverPipelineName: "logging_ping",
}

receiverPipelines["health_checks"] = r.HealthChecksPipeline(ctx)
pipelines["healthchecks"] = otel.Pipeline{
Type: "logs",
ReceiverPipelineName: "health_checks",
}
}

func (r AgentSelfMetrics) PrometheusMetricsPipeline(ctx context.Context) otel.ReceiverPipeline {
Expand Down Expand Up @@ -439,4 +465,136 @@ func (r AgentSelfMetrics) OpsAgentPipeline(ctx context.Context) otel.ReceiverPip
}, ctx)
}

// This method creates a component that enforces the `Structured Health Logs` format to
// all `ops-agent-health` logs. It sets `agentKind`, `agentVersion` and `schemaVersion`.
// This method also processes all self logs to set the severity field correctly.
func generateStructuredHealthLogsOtelComponents(ctx context.Context) []otel.Component {
components, err := LoggingProcessorModifyFields{
Fields: map[string]*ModifyField{
fmt.Sprintf(`labels."%s"`, agentKindKey): {
StaticValue: &agentKind,
},
fmt.Sprintf(`labels."%s"`, agentVersionKey): {
StaticValue: &version.Version,
},
fmt.Sprintf(`labels."%s"`, schemaVersionKey): {
StaticValue: &schemaVersion,
},
"severity": {
MoveFrom: "jsonPayload.severity",
MapValues: map[string]string{
"error": "ERROR",
"warn": "WARNING",
"info": "INFO",
"debug": "DEBUG",
},
MapValuesExclusive: false,
},
},
}.Processors(ctx)
if err != nil {
// We're generating a hard-coded config, so this should never fail.
panic(err)
}
return components
}

// generateHealthLogsParsingComponents creates OTel processors that parse health check logs from JSON
// and filter out any lines that do not have a valid severity.
func generateHealthLogsParsingComponents(ctx context.Context) []otel.Component {
components := []otel.Component{}
parseJsonProccesor, err := LoggingProcessorParseJson{
ParserShared: ParserShared{
TimeKey: "time",
TimeFormat: "%Y-%m-%dT%H:%M:%S%z",
},
}.Processors(ctx)
if err != nil {
// We're generating a hard-coded config, so this should never fail.
panic(err)
}
components = append(components, parseJsonProccesor...)

// This is used to exclude any previous content of the `health-checks.log` file that does not contain
// the `jsonPayload.severity` field.
body := ottl.LValue{"body"}
bodySeverity := ottl.LValue{"body", "severity"}
excludeLogFilter := otel.Filter("logs", "log_record",
[]ottl.Value{
ottl.IsNil(body),
ottl.And(body.IsPresent(), ottl.IsNil(bodySeverity)),
ottl.And(bodySeverity.IsPresent(), ottl.Not(ottl.IsMatch(bodySeverity, "INFO|ERROR|WARNING|DEBUG|info|error|warning|debug"))),
})
components = append(components, excludeLogFilter)

return components
}

func (r AgentSelfMetrics) LoggingPingPipeline(ctx context.Context) otel.ReceiverPipeline {
logProccesors := []otel.Component{
otel.Transform("log", "log",
[]ottl.Statement{
"set(observed_time, Now())",
"set(time, Now())",
},
),
}
logProccesors = append(logProccesors, generateStructuredHealthLogsOtelComponents(ctx)...)
logProccesors = append(logProccesors, otelSetLogNameComponents(ctx, healthLogsTag)...)

return ConvertGCMSystemExporterToOtlpExporter(otel.ReceiverPipeline{
Receiver: otel.Component{
Type: "otlpjsonfile",
Config: map[string]any{
"include": []string{
filepath.Join(r.OtelRuntimeDir, "logging_ping_otlp.json"),
},
"replay_file": true,
"poll_interval": time.Duration(600 * time.Second).String(),
"start_at": "beginning",
},
},
Processors: map[string][]otel.Component{
"logs": logProccesors,
},
ExporterTypes: map[string]otel.ExporterType{
"logs": otel.Logging,
},
}, ctx)
}

func (r AgentSelfMetrics) HealthChecksPipeline(ctx context.Context) otel.ReceiverPipeline {
healthChecksPath := filepath.Join(r.LogsDir, "health-checks.log")

logProccesors := generateHealthLogsParsingComponents(ctx)
logProccesors = append(logProccesors, generateStructuredHealthLogsOtelComponents(ctx)...)
logProccesors = append(logProccesors, otelSetLogNameComponents(ctx, healthLogsTag)...)

return ConvertGCMSystemExporterToOtlpExporter(otel.ReceiverPipeline{
Receiver: otel.Component{
Type: "file_log",
Config: map[string]any{
"include": []string{healthChecksPath},
"start_at": "beginning",
"storage": fileStorageExtensionType,
"operators": []map[string]any{
{
"id": "body",
"type": "move",
"from": "body",
"to": "body.message",
},
},
},
},
Processors: map[string][]otel.Component{
"logs": logProccesors,
},
ExporterTypes: map[string]otel.ExporterType{
"logs": otel.Logging,
},
UsedExtensions: []string{fileStorageExtensionType},
}, ctx)
}

// intentionally not registered as a component because this is not created by users
3 changes: 2 additions & 1 deletion confgenerator/confgenerator.go
Original file line number Diff line number Diff line change
Expand Up @@ -241,7 +241,7 @@ func fileStorageExtension(stateDir string) otel.Component {
}
}

func (uc *UnifiedConfig) GenerateOtelConfig(ctx context.Context, outDir, stateDir string) (string, error) {
func (uc *UnifiedConfig) GenerateOtelConfig(ctx context.Context, outDir, stateDir, logsDir string) (string, error) {
p := platform.FromContext(ctx)

userAgent, _ := p.UserAgent("Google-Cloud-Ops-Agent-Metrics")
Expand All @@ -260,6 +260,7 @@ func (uc *UnifiedConfig) GenerateOtelConfig(ctx context.Context, outDir, stateDi
OtelPort: int(uc.GetOtelMetricsPort()),
OtelRuntimeDir: outDir,
OtlpExporterEnabled: uc.Global.GetOtlpExporter(),
LogsDir: logsDir,
}
agentSelfMetrics.AddSelfMetricsPipelines(receiverPipelines, pipelines, ctx)

Expand Down
10 changes: 8 additions & 2 deletions confgenerator/confgenerator_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -281,7 +281,7 @@ func generateConfigs(pc platformConfig, testDir string) (got map[string]string,
}

// Otel configs
otelGeneratedConfig, err := mergedUc.GenerateOtelConfig(ctx, "", "")
otelGeneratedConfig, err := mergedUc.GenerateOtelConfig(ctx, "", "", "")
if err != nil {
return
}
Expand Down Expand Up @@ -333,6 +333,12 @@ func generateConfigs(pc platformConfig, testDir string) (got map[string]string,
}
got["enabled_receivers_otlp.json"] = string(generatedEnabledReceiversOTLPJSON)

generatedLoggingPingOTLPJSON, err := self_metrics.CollectLoggingPingToOTLPJSON()
if err != nil {
return
}
got["logging_ping_otlp.json"] = string(generatedLoggingPingOTLPJSON)

// If the confgenerator test is designed to test the otel_logging experiment, generate an OTEL config with both otlp_exporter and otel_logging enabled.
if len(enabledExperiments) == 1 && enabledExperiments["otel_logging"] {
generateOtelConfigWithOtlpExporterEnabled(got, pc, testDir, otelGeneratedConfig)
Expand All @@ -358,7 +364,7 @@ func generateOtelConfigWithOtlpExporterEnabled(got map[string]string, pc platfor
enabled := true
mergedUcOtlp.Global.OtlpExporter = &enabled

otelGeneratedConfigOtlp, err := mergedUcOtlp.GenerateOtelConfig(ctxOtlp, "", "")
otelGeneratedConfigOtlp, err := mergedUcOtlp.GenerateOtelConfig(ctxOtlp, "", "", "")
if err == nil {
got["otel_otlp_exporter.yaml"] = otelGeneratedConfigOtlp
}
Expand Down
2 changes: 1 addition & 1 deletion confgenerator/config.go
Original file line number Diff line number Diff line change
Expand Up @@ -1239,7 +1239,7 @@ func (uc *UnifiedConfig) OTelLoggingSupported(ctx context.Context) bool {
}
t := true
ucLoggingCopy.Logging.Service.OTelLogging = &t
_, err = ucLoggingCopy.GenerateOtelConfig(ctx, "", "")
_, err = ucLoggingCopy.GenerateOtelConfig(ctx, "", "", "")
return err == nil
}

Expand Down
2 changes: 1 addition & 1 deletion confgenerator/files.go
Original file line number Diff line number Diff line change
Expand Up @@ -59,7 +59,7 @@ func (uc *UnifiedConfig) GenerateFilesFromConfig(ctx context.Context, service, l
}
}
case "otel":
otelConfig, err := uc.GenerateOtelConfig(ctx, outDir, stateDir)
otelConfig, err := uc.GenerateOtelConfig(ctx, outDir, stateDir, logsDir)
if err != nil {
return fmt.Errorf("can't parse configuration: %w", err)
}
Expand Down
4 changes: 4 additions & 0 deletions confgenerator/otel/ottl/ottl.go
Original file line number Diff line number Diff line change
Expand Up @@ -247,6 +247,10 @@ func Or(conditions ...Value) Value {
return valuef(`(%s)`, strings.Join(out, " or "))
}

func IsNil(a Value) Value {
return valuef(`%s == nil`, a)
}

func IsNotNil(a Value) Value {
return valuef(`%s != nil`, a)
}
Expand Down
63 changes: 0 additions & 63 deletions confgenerator/self_logs.go
Original file line number Diff line number Diff line change
Expand Up @@ -22,23 +22,13 @@ import (

"github.com/GoogleCloudPlatform/ops-agent/confgenerator/fluentbit"
"github.com/GoogleCloudPlatform/ops-agent/internal/healthchecks"
"github.com/GoogleCloudPlatform/ops-agent/internal/logs"
"github.com/GoogleCloudPlatform/ops-agent/internal/platform"
"github.com/GoogleCloudPlatform/ops-agent/internal/version"
)

var (
agentKind string = "ops-agent"
schemaVersion string = "v1"
)

const (
opsAgentLogsMatch string = "ops-agent-*"
fluentBitSelfLogsTag string = "ops-agent-fluent-bit"
healthLogsTag string = "ops-agent-health"
agentVersionKey string = "agent.googleapis.com/health/agentVersion"
agentKindKey string = "agent.googleapis.com/health/agentKind"
schemaVersionKey string = "agent.googleapis.com/health/schemaVersion"
)

func fluentbitSelfLogsPath(p platform.Platform) string {
Expand All @@ -49,57 +39,6 @@ func fluentbitSelfLogsPath(p platform.Platform) string {
return path.Join("${logs_dir}", "subagents", loggingModule)
}

func healthChecksLogsPath() string {
return path.Join("${logs_dir}", "health-checks.log")
}

func generateInputHealthLoggingPingComponent(ctx context.Context) []fluentbit.Component {
return []fluentbit.Component{
{
Kind: "INPUT",
Config: map[string]string{
"Name": "dummy",
"Tag": healthLogsTag,
"Dummy": `{"code": "LogPingOpsAgent", "severity": "DEBUG"}`,
"Interval_Sec": "600",
"Interval_NSec": "0",
},
},
}
}

// This method creates a file input for the `health-checks.log` file, a json parser for the
// structured logs and a grep filter to avoid ingesting previous content of the file.
func generateInputHealthChecksLogsComponents(ctx context.Context) []fluentbit.Component {
out := make([]fluentbit.Component, 0)
out = append(out, LoggingReceiverFilesMixin{
IncludePaths: []string{healthChecksLogsPath()},
BufferInMemory: true,
}.Components(ctx, healthLogsTag)...)
out = append(out, LoggingProcessorParseJson{
// TODO(b/282754149): Remove TimeKey and TimeFormat when feature gets implemented.
ParserShared: ParserShared{
TimeKey: logs.TimeZapKey,
TimeFormat: "%Y-%m-%dT%H:%M:%S%z",
},
}.Components(ctx, healthLogsTag, "health-checks-json")...)
out = append(out, []fluentbit.Component{
// This is used to exclude any previous content of the `health-checks.log` file that does not contain
// the `jsonPayload.severity` field. Due to `https://github.com/fluent/fluent-bit/issues/7092` the
// filtering can't be done directly to the `logging.googleapis.com/severity` field.
// We cannot use `LoggingProcessorExcludeLogs` here since it doesn't exclude when the field is missing.
{
Kind: "FILTER",
Config: map[string]string{
"Name": "grep",
"Match": healthLogsTag,
"Regex": fmt.Sprintf("%s INFO|ERROR|WARNING|DEBUG|info|error|warning|debug", logs.SeverityZapKey),
},
},
}...)
return out
}

// This method creates a file input for the `logging-module.log` file, a regex parser for the
// fluent-bit self logs and a translator of severity to the logging api format.
func generateInputFluentBitSelfLogsComponents(ctx context.Context, logLevel string) []fluentbit.Component {
Expand Down Expand Up @@ -216,9 +155,7 @@ func generateOutputSelfLogsComponent(ctx context.Context, userAgent string, inge

func (uc *UnifiedConfig) generateSelfLogsComponents(ctx context.Context, userAgent string) []fluentbit.Component {
out := make([]fluentbit.Component, 0)
out = append(out, generateInputHealthLoggingPingComponent(ctx)...)
out = append(out, generateInputFluentBitSelfLogsComponents(ctx, uc.Logging.Service.LogLevel)...)
out = append(out, generateInputHealthChecksLogsComponents(ctx)...)
out = append(out, generateFilterSelfLogsSamplingComponents(ctx)...)
out = append(out, generateFilterStructuredHealthLogsComponents(ctx)...)
out = append(out, generateFilterMapSeverityFieldComponent(ctx)...)
Expand Down
Loading
Loading