mirror of
https://github.com/agent-substrate/substrate.git
synced 2026-10-02 03:24:42 +08:00
otlprelay: forward log records (#1800)
Step 1 of #1748. The relay forwards `LogsService` the way it forwards traces and metrics: source gate, verbatim payload, atelet's own `OTEL_EXPORTER_OTLP_LOGS_HEADERS` upstream. No content check on the records. Identity spoofing through the relay is handled in #1748 by the consumer rule (a lifecycle record counts only under ateapi's resource, which the source gate already enforces) and per-pod sockets (#741), not by a denylist here. **Endpoint and compression resolution** (second commit, from review). The relay read the generic, traces, and metrics OTLP variables and ignored the logs ones, so a logs-specific collector or compression setting was silently overridden. All three signals now go through one resolver: each falls back to the generic, and every resolved value must agree, since the relay has one upstream connection. The conflict error names the variable each value came from; a signal that falls back names the generic variable, not its own unset one. One behavior change: `traces=A, metrics=A, generic=B` with logs unset used to pass and now fails at atelet startup, because logs falls back to `B` and that is a two-collector configuration. Tests: an ateom log batch arrives with the ateom's resource; the source gate, mixed-batch, header-isolation, and per-signal header tests cover logs; resolver tables gained logs rows including the fallback case. Docs say the relay carries logs, traces, and metrics.
This commit is contained in:
@@ -440,8 +440,8 @@ ateapi, through `serverboot.InitLogging`. They are off unless
|
||||
|
||||
Everything else is stdout. `serverboot.InitLogger` writes structured JSON there,
|
||||
and `ateom` wraps actor container output with the `ate.*` metadata labels
|
||||
described in [Actor Observability](../../observability.md). No worker pod
|
||||
exports a log record at all: the ateom relay carries traces and metrics only.
|
||||
described in [Actor Observability](../../observability.md). The ateom relay
|
||||
carries logs, traces, and metrics; worker pods do not emit log records yet.
|
||||
|
||||
Those labels sit in a nested group (`labels`, or `logging.googleapis.com/labels`
|
||||
on GKE, where the key promotes the group into `LogEntry.labels`). A filelog
|
||||
|
||||
@@ -178,7 +178,7 @@ The counter carries the same reason but no actor identity, so this record is the
|
||||
|
||||
#### The same records over OTLP
|
||||
|
||||
Both records also go out as OTLP log events, so a collector reads them without knowing substrate's stdout envelope. Set `OTEL_LOGS_EXPORTER=otlp` to turn it on; unset means `none`, which is what every environment but kind uses today. Only ateapi has a LoggerProvider — a worker pod cannot export a log record at all yet, because [the ateom relay](#the-ateom-otlp-relay) carries traces and metrics only.
|
||||
Both records also go out as OTLP log events, so a collector reads them without knowing substrate's stdout envelope. Set `OTEL_LOGS_EXPORTER=otlp` to turn it on; unset means `none`, which is what every environment but kind uses today. Only ateapi has a LoggerProvider today; [the ateom relay](#the-ateom-otlp-relay) carries logs, traces, and metrics, so an ateom exports log records the same way once it has one.
|
||||
|
||||
Two `event.name` values, which is the OTLP LogRecord's own field rather than an attribute:
|
||||
|
||||
@@ -410,7 +410,7 @@ Telemetry is emitted the same way everywhere; only the backend differs between a
|
||||
|
||||
### The ateom OTLP relay
|
||||
|
||||
ateom is the one component that does not talk to the collector directly. It exports over a unix socket at `/var/lib/ateom-gvisor/atelet-otlp.sock`, which `atelet` serves and forwards to the collector on the node's network ([`internal/otlprelay`](../internal/otlprelay)):
|
||||
ateom is the one component that does not talk to the collector directly. It exports logs, traces, and metrics over a unix socket at `/var/lib/ateom-gvisor/atelet-otlp.sock`, which `atelet` serves and forwards to the collector on the node's network ([`internal/otlprelay`](../internal/otlprelay)):
|
||||
|
||||
```
|
||||
ateom ──OTLP/gRPC over unix socket──► atelet relay ──OTLP/gRPC──► collector
|
||||
|
||||
+74
-54
@@ -14,7 +14,7 @@
|
||||
|
||||
// Package otlprelay carries ateom's OTLP telemetry to the collector over a unix
|
||||
// socket served by atelet, so a worker pod needs no network path of its own to
|
||||
// export spans and metrics.
|
||||
// export spans, metrics, and log records.
|
||||
//
|
||||
// Motivation. ateom runs inside the worker pod that hosts the actor, and until
|
||||
// now exported OTLP straight to the collector over the pod's network (the
|
||||
@@ -71,6 +71,7 @@ import (
|
||||
"k8s.io/utils/lru"
|
||||
|
||||
semconv "go.opentelemetry.io/otel/semconv/v1.40.0"
|
||||
collogspb "go.opentelemetry.io/proto/otlp/collector/logs/v1"
|
||||
colmetricspb "go.opentelemetry.io/proto/otlp/collector/metrics/v1"
|
||||
coltracepb "go.opentelemetry.io/proto/otlp/collector/trace/v1"
|
||||
resourcepb "go.opentelemetry.io/proto/otlp/resource/v1"
|
||||
@@ -83,20 +84,23 @@ const (
|
||||
endpointEnv = "OTEL_EXPORTER_OTLP_ENDPOINT"
|
||||
tracesEndpointEnv = "OTEL_EXPORTER_OTLP_TRACES_ENDPOINT"
|
||||
metricsEndpointEnv = "OTEL_EXPORTER_OTLP_METRICS_ENDPOINT"
|
||||
logsEndpointEnv = "OTEL_EXPORTER_OTLP_LOGS_ENDPOINT"
|
||||
|
||||
// compressionEnv and its signal-specific overrides configure upstream
|
||||
// gRPC compression (gzip or none).
|
||||
compressionEnv = "OTEL_EXPORTER_OTLP_COMPRESSION"
|
||||
tracesCompressionEnv = "OTEL_EXPORTER_OTLP_TRACES_COMPRESSION"
|
||||
metricsCompressionEnv = "OTEL_EXPORTER_OTLP_METRICS_COMPRESSION"
|
||||
logsCompressionEnv = "OTEL_EXPORTER_OTLP_LOGS_COMPRESSION"
|
||||
|
||||
// headersEnv and its signal-specific overrides carry the headers the
|
||||
// collector expects (an API key, a tenant id). Unlike the endpoint and the
|
||||
// compression, these are per-call metadata rather than per-connection, so
|
||||
// traces and metrics may legitimately differ and are resolved separately.
|
||||
// traces, metrics, and logs may legitimately differ and are resolved separately.
|
||||
headersEnv = "OTEL_EXPORTER_OTLP_HEADERS"
|
||||
tracesHeadersEnv = "OTEL_EXPORTER_OTLP_TRACES_HEADERS"
|
||||
metricsHeadersEnv = "OTEL_EXPORTER_OTLP_METRICS_HEADERS"
|
||||
logsHeadersEnv = "OTEL_EXPORTER_OTLP_LOGS_HEADERS"
|
||||
|
||||
// otlpDefaultPort matches atenet's normalizeOtlpCollector.
|
||||
otlpDefaultPort = "4317"
|
||||
@@ -286,10 +290,10 @@ func upstreamHeaders(signalEnv string) (metadata.MD, error) {
|
||||
return md, nil
|
||||
}
|
||||
|
||||
// The two OTLP services both declare a method named Export, with different
|
||||
// request types, so one type cannot implement both: the embedded Unimplemented
|
||||
// structs would give Server an ambiguous promoted Export and satisfy neither
|
||||
// interface. Each service gets its own tiny forwarder instead.
|
||||
// The OTLP services all declare a method named Export, with different request
|
||||
// types, so one type cannot implement more than one: the embedded Unimplemented
|
||||
// structs would give Server an ambiguous promoted Export and satisfy none of the
|
||||
// interfaces. Each service gets its own tiny forwarder instead.
|
||||
|
||||
type traceRelay struct {
|
||||
coltracepb.UnimplementedTraceServiceServer
|
||||
@@ -297,8 +301,8 @@ type traceRelay struct {
|
||||
// headers atelet presents to the collector; see upstreamContext. Resolved
|
||||
// once at construction: they come from atelet's environment, not the call.
|
||||
headers metadata.MD
|
||||
// gate is shared with metricRelay so a misnamed ateom is reported once, not
|
||||
// once per signal.
|
||||
// gate is shared with metricRelay and logRelay so a misnamed ateom is
|
||||
// reported once, not once per signal.
|
||||
gate *sourceGate
|
||||
}
|
||||
|
||||
@@ -335,6 +339,23 @@ func (m *metricRelay) Export(ctx context.Context, req *colmetricspb.ExportMetric
|
||||
return m.upstream.Export(upstreamContext(ctx, m.headers), req)
|
||||
}
|
||||
|
||||
type logRelay struct {
|
||||
collogspb.UnimplementedLogsServiceServer
|
||||
upstream collogspb.LogsServiceClient
|
||||
headers metadata.MD
|
||||
gate *sourceGate
|
||||
}
|
||||
|
||||
// Export forwards a batch of log records to the collector unchanged.
|
||||
func (l *logRelay) Export(ctx context.Context, req *collogspb.ExportLogsServiceRequest) (*collogspb.ExportLogsServiceResponse, error) {
|
||||
for _, rl := range req.GetResourceLogs() {
|
||||
if err := l.gate.check(ctx, rl.GetResource()); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
}
|
||||
return l.upstream.Export(upstreamContext(ctx, l.headers), req)
|
||||
}
|
||||
|
||||
// validateSocketPath rejects a relative path, which gRPC does not resolve:
|
||||
// "unix://foo/r.sock" parses as authority "foo", path "/r.sock". grpc.NewClient
|
||||
// being lazy, that wrong target is accepted at startup and fails per export
|
||||
@@ -385,6 +406,10 @@ func NewServer(ctx context.Context, sockPath string) (*Server, error) {
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
logHeaders, err := upstreamHeaders(logsHeadersEnv)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
dialOpts := []grpc.DialOption{
|
||||
// Plaintext by design today; TLS support for the upstream leg will be added
|
||||
@@ -418,12 +443,18 @@ func NewServer(ctx context.Context, sockPath string) (*Server, error) {
|
||||
headers: metricHeaders,
|
||||
gate: gate,
|
||||
})
|
||||
collogspb.RegisterLogsServiceServer(s.grpc, &logRelay{
|
||||
upstream: collogspb.NewLogsServiceClient(upstream),
|
||||
headers: logHeaders,
|
||||
gate: gate,
|
||||
})
|
||||
// Header names only: the values are credentials.
|
||||
slog.InfoContext(ctx, "OTLP relay forwarding to collector",
|
||||
slog.String("collector", target),
|
||||
slog.String("compression", comp),
|
||||
slog.Any("traceHeaders", headerNames(traceHeaders)),
|
||||
slog.Any("metricHeaders", headerNames(metricHeaders)))
|
||||
slog.Any("metricHeaders", headerNames(metricHeaders)),
|
||||
slog.Any("logHeaders", headerNames(logHeaders)))
|
||||
return s, nil
|
||||
}
|
||||
|
||||
@@ -477,27 +508,9 @@ func headerNames(md metadata.MD) []string {
|
||||
// upstreamCompression resolves the compression algorithm (gzip or none) to use
|
||||
// for upstream export.
|
||||
func upstreamCompression() (string, error) {
|
||||
generic := strings.TrimSpace(os.Getenv(compressionEnv))
|
||||
traces := strings.TrimSpace(os.Getenv(tracesCompressionEnv))
|
||||
metrics := strings.TrimSpace(os.Getenv(metricsCompressionEnv))
|
||||
|
||||
traceComp := generic
|
||||
if traces != "" {
|
||||
traceComp = traces
|
||||
}
|
||||
metricComp := generic
|
||||
if metrics != "" {
|
||||
metricComp = metrics
|
||||
}
|
||||
|
||||
if traceComp != "" && metricComp != "" && traceComp != metricComp {
|
||||
return "", fmt.Errorf("signal-specific compression settings conflict (%q for traces vs %q for metrics); the relay carries both signals over one connection",
|
||||
traceComp, metricComp)
|
||||
}
|
||||
|
||||
resolved := traceComp
|
||||
if resolved == "" {
|
||||
resolved = metricComp
|
||||
resolved, err := resolvePerSignal("compression settings", compressionEnv, tracesCompressionEnv, metricsCompressionEnv, logsCompressionEnv)
|
||||
if err != nil {
|
||||
return "", err
|
||||
}
|
||||
switch resolved {
|
||||
case "", "none":
|
||||
@@ -512,32 +525,14 @@ func upstreamCompression() (string, error) {
|
||||
// upstreamTarget resolves the collector address the relay forwards to, from the
|
||||
// standard OTLP endpoint variables, into the bare host:port grpc.NewClient wants.
|
||||
//
|
||||
// The signal-specific variables must agree: the relay carries traces and metrics
|
||||
// over one connection, so it cannot honor two different collectors. Configuring
|
||||
// both differently is a misconfiguration rather than something to silently pick
|
||||
// The signal-specific variables must agree: the relay carries every signal over
|
||||
// one connection, so it cannot honor two different collectors. Configuring
|
||||
// them differently is a misconfiguration rather than something to silently pick
|
||||
// a winner for.
|
||||
func upstreamTarget() (string, error) {
|
||||
generic := strings.TrimSpace(os.Getenv(endpointEnv))
|
||||
traces := strings.TrimSpace(os.Getenv(tracesEndpointEnv))
|
||||
metrics := strings.TrimSpace(os.Getenv(metricsEndpointEnv))
|
||||
|
||||
traceTarget := generic
|
||||
if traces != "" {
|
||||
traceTarget = traces
|
||||
}
|
||||
metricTarget := generic
|
||||
if metrics != "" {
|
||||
metricTarget = metrics
|
||||
}
|
||||
|
||||
if traceTarget != "" && metricTarget != "" && traceTarget != metricTarget {
|
||||
return "", fmt.Errorf("signal-specific endpoints conflict (%q for traces vs %q for metrics); the relay carries both signals over one connection",
|
||||
traceTarget, metricTarget)
|
||||
}
|
||||
|
||||
resolved := traceTarget
|
||||
if resolved == "" {
|
||||
resolved = metricTarget
|
||||
resolved, err := resolvePerSignal("endpoints", endpointEnv, tracesEndpointEnv, metricsEndpointEnv, logsEndpointEnv)
|
||||
if err != nil {
|
||||
return "", err
|
||||
}
|
||||
if resolved == "" {
|
||||
return "", nil
|
||||
@@ -545,6 +540,31 @@ func upstreamTarget() (string, error) {
|
||||
return normalizeEndpoint(resolved)
|
||||
}
|
||||
|
||||
// resolvePerSignal resolves each signal to its own variable, falling back to
|
||||
// the generic one, and requires every resolved value to agree: the relay
|
||||
// carries every signal over one connection. The error names the variables the
|
||||
// two values came from, since a fallback pulls the generic in under a signal
|
||||
// that is not itself set.
|
||||
func resolvePerSignal(what, genericEnv string, signalEnvs ...string) (string, error) {
|
||||
generic := strings.TrimSpace(os.Getenv(genericEnv))
|
||||
resolved, resolvedEnv := "", ""
|
||||
for _, env := range signalEnvs {
|
||||
v, src := strings.TrimSpace(os.Getenv(env)), env
|
||||
if v == "" {
|
||||
v, src = generic, genericEnv
|
||||
}
|
||||
if v == "" {
|
||||
continue
|
||||
}
|
||||
if resolved != "" && v != resolved {
|
||||
return "", fmt.Errorf("signal-specific %s conflict: %s=%q vs %s=%q; the relay carries every signal over one connection",
|
||||
what, resolvedEnv, resolved, src, v)
|
||||
}
|
||||
resolved, resolvedEnv = v, src
|
||||
}
|
||||
return resolved, nil
|
||||
}
|
||||
|
||||
// normalizeEndpoint accepts both a bare "host:port" and the URL form the OTLP
|
||||
// environment variables carry, and returns the host:port grpc.NewClient dials.
|
||||
//
|
||||
|
||||
@@ -36,9 +36,11 @@ import (
|
||||
"google.golang.org/grpc/status"
|
||||
"google.golang.org/protobuf/testing/protocmp"
|
||||
|
||||
collogspb "go.opentelemetry.io/proto/otlp/collector/logs/v1"
|
||||
colmetricspb "go.opentelemetry.io/proto/otlp/collector/metrics/v1"
|
||||
coltracepb "go.opentelemetry.io/proto/otlp/collector/trace/v1"
|
||||
commonpb "go.opentelemetry.io/proto/otlp/common/v1"
|
||||
logspb "go.opentelemetry.io/proto/otlp/logs/v1"
|
||||
metricspb "go.opentelemetry.io/proto/otlp/metrics/v1"
|
||||
resourcepb "go.opentelemetry.io/proto/otlp/resource/v1"
|
||||
tracepb "go.opentelemetry.io/proto/otlp/trace/v1"
|
||||
@@ -52,8 +54,10 @@ type fakeCollector struct {
|
||||
mu sync.Mutex
|
||||
traces []*coltracepb.ExportTraceServiceRequest
|
||||
metrics []*colmetricspb.ExportMetricsServiceRequest
|
||||
logs []*collogspb.ExportLogsServiceRequest
|
||||
traceMD []metadata.MD
|
||||
metricMD []metadata.MD
|
||||
logMD []metadata.MD
|
||||
got chan struct{}
|
||||
}
|
||||
|
||||
@@ -86,6 +90,22 @@ func (m *metricsSink) Export(ctx context.Context, req *colmetricspb.ExportMetric
|
||||
return &colmetricspb.ExportMetricsServiceResponse{}, nil
|
||||
}
|
||||
|
||||
type logsSink struct {
|
||||
collogspb.UnimplementedLogsServiceServer
|
||||
parent *fakeCollector
|
||||
}
|
||||
|
||||
func (l *logsSink) Export(ctx context.Context, req *collogspb.ExportLogsServiceRequest) (*collogspb.ExportLogsServiceResponse, error) {
|
||||
l.parent.mu.Lock()
|
||||
l.parent.logs = append(l.parent.logs, req)
|
||||
if md, ok := metadata.FromIncomingContext(ctx); ok {
|
||||
l.parent.logMD = append(l.parent.logMD, md.Copy())
|
||||
}
|
||||
l.parent.mu.Unlock()
|
||||
l.parent.got <- struct{}{}
|
||||
return &collogspb.ExportLogsServiceResponse{}, nil
|
||||
}
|
||||
|
||||
// startFakeCollector serves the OTLP collector services on a loopback TCP port
|
||||
// (the shape the relay forwards to) and returns the sink and its host:port.
|
||||
func startFakeCollector(t *testing.T) (*fakeCollector, string) {
|
||||
@@ -98,6 +118,7 @@ func startFakeCollector(t *testing.T) (*fakeCollector, string) {
|
||||
srv := grpc.NewServer()
|
||||
coltracepb.RegisterTraceServiceServer(srv, sink)
|
||||
colmetricspb.RegisterMetricsServiceServer(srv, &metricsSink{parent: sink})
|
||||
collogspb.RegisterLogsServiceServer(srv, &logsSink{parent: sink})
|
||||
go func() { _ = srv.Serve(lis) }()
|
||||
t.Cleanup(srv.Stop)
|
||||
return sink, lis.Addr().String()
|
||||
@@ -248,6 +269,49 @@ func TestRelayForwardsMetrics(t *testing.T) {
|
||||
}
|
||||
}
|
||||
|
||||
// usageLogs is an ateom log batch.
|
||||
func usageLogs() *collogspb.ExportLogsServiceRequest {
|
||||
return &collogspb.ExportLogsServiceRequest{
|
||||
ResourceLogs: []*logspb.ResourceLogs{{
|
||||
Resource: serviceResource("ateom-gvisor"),
|
||||
ScopeLogs: []*logspb.ScopeLogs{{
|
||||
LogRecords: []*logspb.LogRecord{{EventName: "ate.actor.usage_sampled"}},
|
||||
}},
|
||||
}},
|
||||
}
|
||||
}
|
||||
|
||||
func TestRelayForwardsLogsVerbatim(t *testing.T) {
|
||||
sink, collector := startFakeCollector(t)
|
||||
sock := startRelay(t, collector)
|
||||
|
||||
conn, err := Dial(context.Background(), sock)
|
||||
if err != nil {
|
||||
t.Fatalf("Dial: %v", err)
|
||||
}
|
||||
defer conn.Close()
|
||||
|
||||
req := usageLogs()
|
||||
if _, err := collogspb.NewLogsServiceClient(conn).Export(context.Background(), req); err != nil {
|
||||
t.Fatalf("Export through the relay: %v", err)
|
||||
}
|
||||
|
||||
select {
|
||||
case <-sink.got:
|
||||
case <-time.After(5 * time.Second):
|
||||
t.Fatal("collector never received the forwarded log export")
|
||||
}
|
||||
|
||||
sink.mu.Lock()
|
||||
defer sink.mu.Unlock()
|
||||
if len(sink.logs) != 1 {
|
||||
t.Fatalf("collector got %d log exports, want 1", len(sink.logs))
|
||||
}
|
||||
if diff := cmp.Diff(req, sink.logs[0], protocmp.Transform()); diff != "" {
|
||||
t.Errorf("forwarded request differs from what was sent (-sent +received):\n%s", diff)
|
||||
}
|
||||
}
|
||||
|
||||
// TestStopRemovesSocket matters for the restart path: a leftover socket makes
|
||||
// the next atelet's Listen fail with EADDRINUSE, and in the meantime makes
|
||||
// every ateom on the node believe a relay is there.
|
||||
@@ -404,13 +468,20 @@ func TestRelayRefusesNonAteomSource(t *testing.T) {
|
||||
if got := status.Code(err); got != codes.PermissionDenied {
|
||||
t.Errorf("metric Export from service.name %q = code %v (%v), want %v", tc.service, got, err, codes.PermissionDenied)
|
||||
}
|
||||
|
||||
_, err = collogspb.NewLogsServiceClient(conn).Export(context.Background(), &collogspb.ExportLogsServiceRequest{
|
||||
ResourceLogs: []*logspb.ResourceLogs{{Resource: serviceResource(tc.service)}},
|
||||
})
|
||||
if got := status.Code(err); got != codes.PermissionDenied {
|
||||
t.Errorf("log Export from service.name %q = code %v (%v), want %v", tc.service, got, err, codes.PermissionDenied)
|
||||
}
|
||||
})
|
||||
}
|
||||
|
||||
sink.mu.Lock()
|
||||
defer sink.mu.Unlock()
|
||||
if len(sink.traces) != 0 || len(sink.metrics) != 0 {
|
||||
t.Errorf("collector received %d traces and %d metrics from refused sources, want none to be forwarded", len(sink.traces), len(sink.metrics))
|
||||
if len(sink.traces) != 0 || len(sink.metrics) != 0 || len(sink.logs) != 0 {
|
||||
t.Errorf("collector received %d traces, %d metrics, %d logs from refused sources, want none to be forwarded", len(sink.traces), len(sink.metrics), len(sink.logs))
|
||||
}
|
||||
}
|
||||
|
||||
@@ -436,10 +507,20 @@ func TestRelayRefusesMixedBatch(t *testing.T) {
|
||||
t.Errorf("Export of a mixed batch = code %v (%v), want %v", got, err, codes.PermissionDenied)
|
||||
}
|
||||
|
||||
_, err = collogspb.NewLogsServiceClient(conn).Export(context.Background(), &collogspb.ExportLogsServiceRequest{
|
||||
ResourceLogs: []*logspb.ResourceLogs{
|
||||
{Resource: serviceResource("ateom-gvisor")},
|
||||
{Resource: serviceResource("actor")},
|
||||
},
|
||||
})
|
||||
if got := status.Code(err); got != codes.PermissionDenied {
|
||||
t.Errorf("log Export of a mixed batch = code %v (%v), want %v", got, err, codes.PermissionDenied)
|
||||
}
|
||||
|
||||
sink.mu.Lock()
|
||||
defer sink.mu.Unlock()
|
||||
if len(sink.traces) != 0 {
|
||||
t.Errorf("collector received %d exports from a mixed batch, want the batch refused whole", len(sink.traces))
|
||||
if len(sink.traces) != 0 || len(sink.logs) != 0 {
|
||||
t.Errorf("collector received %d traces and %d logs from mixed batches, want the batches refused whole", len(sink.traces), len(sink.logs))
|
||||
}
|
||||
}
|
||||
|
||||
@@ -459,13 +540,25 @@ func TestRelayAcceptsEveryAteomService(t *testing.T) {
|
||||
if _, err := coltracepb.NewTraceServiceClient(conn).Export(context.Background(), &coltracepb.ExportTraceServiceRequest{
|
||||
ResourceSpans: []*tracepb.ResourceSpans{{Resource: serviceResource(service)}},
|
||||
}); err != nil {
|
||||
t.Errorf("Export from allowlisted service %q: %v", service, err)
|
||||
t.Errorf("trace Export from allowlisted service %q: %v", service, err)
|
||||
continue
|
||||
}
|
||||
select {
|
||||
case <-sink.got:
|
||||
case <-time.After(5 * time.Second):
|
||||
t.Errorf("collector never received the export from %q", service)
|
||||
t.Errorf("collector never received the trace export from %q", service)
|
||||
}
|
||||
|
||||
if _, err := collogspb.NewLogsServiceClient(conn).Export(context.Background(), &collogspb.ExportLogsServiceRequest{
|
||||
ResourceLogs: []*logspb.ResourceLogs{{Resource: serviceResource(service)}},
|
||||
}); err != nil {
|
||||
t.Errorf("log Export from allowlisted service %q: %v", service, err)
|
||||
continue
|
||||
}
|
||||
select {
|
||||
case <-sink.got:
|
||||
case <-time.After(5 * time.Second):
|
||||
t.Errorf("collector never received the log export from %q", service)
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -603,12 +696,23 @@ func TestUpstreamTarget(t *testing.T) {
|
||||
generic string
|
||||
traces string
|
||||
metrics string
|
||||
logs string
|
||||
want string
|
||||
wantErr bool
|
||||
// errHas / errLacks, when set, pin which variable the error names.
|
||||
errHas string
|
||||
errLacks string
|
||||
}{
|
||||
{name: "unset", want: ""},
|
||||
{name: "generic only", generic: "collector:4317", want: "collector:4317"},
|
||||
{name: "signal specific overrides generic", generic: "generic:4317", traces: "specific:4317", metrics: "specific:4317", want: "specific:4317"},
|
||||
{name: "logs specific agrees", generic: "collector:4317", logs: "collector:4317", want: "collector:4317"},
|
||||
{name: "logs specific conflicts with generic", generic: "generic:4317", logs: "logs-only:4317", wantErr: true},
|
||||
{name: "logs specific conflicts with traces", traces: "a:4317", logs: "b:4317", wantErr: true},
|
||||
{name: "signal specific overrides generic", generic: "generic:4317", traces: "specific:4317", metrics: "specific:4317", logs: "specific:4317", want: "specific:4317"},
|
||||
// Logs left unset falls back to the generic, which then disagrees with
|
||||
// the two overrides: a real two-collector configuration, refused. The
|
||||
// error must name the generic variable, since LOGS_ENDPOINT is unset.
|
||||
{name: "two overrides and an unset signal conflict with generic", generic: "generic:4317", traces: "specific:4317", metrics: "specific:4317", wantErr: true, errHas: endpointEnv, errLacks: logsEndpointEnv},
|
||||
{name: "generic conflicts with different traces specific", generic: "generic:4317", traces: "traces-only:4317", wantErr: true},
|
||||
{name: "generic conflicts with different metrics specific", generic: "generic:4317", metrics: "metrics-only:4317", wantErr: true},
|
||||
{name: "matching generic and signal specific", generic: "collector:4317", traces: "collector:4317", metrics: "collector:4317", want: "collector:4317"},
|
||||
@@ -621,11 +725,18 @@ func TestUpstreamTarget(t *testing.T) {
|
||||
t.Setenv(endpointEnv, tc.generic)
|
||||
t.Setenv(tracesEndpointEnv, tc.traces)
|
||||
t.Setenv(metricsEndpointEnv, tc.metrics)
|
||||
t.Setenv(logsEndpointEnv, tc.logs)
|
||||
got, err := upstreamTarget()
|
||||
if tc.wantErr {
|
||||
if err == nil {
|
||||
t.Fatalf("upstreamTarget() = %q, want an error", got)
|
||||
}
|
||||
if tc.errHas != "" && !strings.Contains(err.Error(), tc.errHas) {
|
||||
t.Errorf("error %q does not name %s", err, tc.errHas)
|
||||
}
|
||||
if tc.errLacks != "" && strings.Contains(err.Error(), tc.errLacks) {
|
||||
t.Errorf("error %q names %s, which is not set", err, tc.errLacks)
|
||||
}
|
||||
return
|
||||
}
|
||||
if err != nil {
|
||||
@@ -644,10 +755,13 @@ func TestUpstreamCompression(t *testing.T) {
|
||||
generic string
|
||||
traces string
|
||||
metrics string
|
||||
logs string
|
||||
want string
|
||||
wantErr bool
|
||||
}{
|
||||
{name: "unset", want: "none"},
|
||||
{name: "logs specific gzip", logs: "gzip", want: "gzip"},
|
||||
{name: "logs conflicts with metrics", metrics: "gzip", logs: "none", wantErr: true},
|
||||
{name: "generic gzip", generic: "gzip", want: "gzip"},
|
||||
{name: "generic none", generic: "none", want: "none"},
|
||||
{name: "traces specific gzip", traces: "gzip", want: "gzip"},
|
||||
@@ -661,6 +775,7 @@ func TestUpstreamCompression(t *testing.T) {
|
||||
t.Setenv(compressionEnv, tc.generic)
|
||||
t.Setenv(tracesCompressionEnv, tc.traces)
|
||||
t.Setenv(metricsCompressionEnv, tc.metrics)
|
||||
t.Setenv(logsCompressionEnv, tc.logs)
|
||||
got, err := upstreamCompression()
|
||||
if tc.wantErr {
|
||||
if err == nil {
|
||||
@@ -678,9 +793,8 @@ func TestUpstreamCompression(t *testing.T) {
|
||||
}
|
||||
}
|
||||
|
||||
// exportBoth sends one empty trace batch and one empty metric batch from
|
||||
// service through the relay, and returns the metadata each arrived with.
|
||||
func exportBoth(t *testing.T, sink *fakeCollector, sock, service string, ctx context.Context) (traceMD, metricMD metadata.MD) {
|
||||
// exportAll sends one empty batch per signal from service through the relay, and returns the metadata each arrived with.
|
||||
func exportAll(t *testing.T, sink *fakeCollector, sock, service string, ctx context.Context) (traceMD, metricMD, logMD metadata.MD) {
|
||||
t.Helper()
|
||||
conn, err := Dial(context.Background(), sock)
|
||||
if err != nil {
|
||||
@@ -698,16 +812,18 @@ func exportBoth(t *testing.T, sink *fakeCollector, sock, service string, ctx con
|
||||
}); err != nil {
|
||||
t.Fatalf("metricClient.Export: %v", err)
|
||||
}
|
||||
if _, err := collogspb.NewLogsServiceClient(conn).Export(ctx, &collogspb.ExportLogsServiceRequest{
|
||||
ResourceLogs: []*logspb.ResourceLogs{{Resource: serviceResource(service)}},
|
||||
}); err != nil {
|
||||
t.Fatalf("logClient.Export: %v", err)
|
||||
}
|
||||
|
||||
sink.mu.Lock()
|
||||
defer sink.mu.Unlock()
|
||||
if len(sink.traceMD) == 0 {
|
||||
t.Fatal("collector received no metadata for trace export")
|
||||
if len(sink.traceMD) == 0 || len(sink.metricMD) == 0 || len(sink.logMD) == 0 {
|
||||
t.Fatalf("collector received metadata for %d trace, %d metric, %d log exports, want one each", len(sink.traceMD), len(sink.metricMD), len(sink.logMD))
|
||||
}
|
||||
if len(sink.metricMD) == 0 {
|
||||
t.Fatal("collector received no metadata for metrics export")
|
||||
}
|
||||
return sink.traceMD[0], sink.metricMD[0]
|
||||
return sink.traceMD[0], sink.metricMD[0], sink.logMD[0]
|
||||
}
|
||||
|
||||
// The upstream leg is atelet's connection, so its credentials are atelet's. An
|
||||
@@ -721,12 +837,12 @@ func TestExportDropsClientMetadata(t *testing.T) {
|
||||
"authorization", "Bearer client-token",
|
||||
"custom-header", "custom-value",
|
||||
)
|
||||
traceMD, metricMD := exportBoth(t, sink, sock, "ateom-gvisor", ctx)
|
||||
traceMD, metricMD, logMD := exportAll(t, sink, sock, "ateom-gvisor", ctx)
|
||||
|
||||
for _, tc := range []struct {
|
||||
signal string
|
||||
md metadata.MD
|
||||
}{{"trace", traceMD}, {"metric", metricMD}} {
|
||||
}{{"trace", traceMD}, {"metric", metricMD}, {"log", logMD}} {
|
||||
for _, key := range []string{"authorization", "custom-header"} {
|
||||
if got := tc.md.Get(key); len(got) != 0 {
|
||||
t.Errorf("%s export reached the collector with the client's %s = %v, want it dropped", tc.signal, key, got)
|
||||
@@ -742,12 +858,13 @@ func TestExportAttachesAteletHeaders(t *testing.T) {
|
||||
// Set before startRelay: NewServer resolves headers once, at construction.
|
||||
t.Setenv(headersEnv, "authorization=Bearer atelet-token,x-tenant=substrate")
|
||||
t.Setenv(metricsHeadersEnv, "authorization=Bearer metrics-token")
|
||||
t.Setenv(logsHeadersEnv, "authorization=Bearer logs-token")
|
||||
sock := startRelay(t, collector)
|
||||
|
||||
// The client sends its own, which must lose to atelet's rather than
|
||||
// appending a second value the collector might pick either way.
|
||||
ctx := metadata.AppendToOutgoingContext(context.Background(), "authorization", "Bearer client-token")
|
||||
traceMD, metricMD := exportBoth(t, sink, sock, "ateom-gvisor", ctx)
|
||||
traceMD, metricMD, logMD := exportAll(t, sink, sock, "ateom-gvisor", ctx)
|
||||
|
||||
if got := traceMD.Get("authorization"); len(got) != 1 || got[0] != "Bearer atelet-token" {
|
||||
t.Errorf("trace export authorization = %v, want exactly [Bearer atelet-token]", got)
|
||||
@@ -763,6 +880,12 @@ func TestExportAttachesAteletHeaders(t *testing.T) {
|
||||
if got := metricMD.Get("x-tenant"); len(got) != 0 {
|
||||
t.Errorf("metric export x-tenant = %v, want absent: %s replaces %s rather than merging", got, metricsHeadersEnv, headersEnv)
|
||||
}
|
||||
if got := logMD.Get("authorization"); len(got) != 1 || got[0] != "Bearer logs-token" {
|
||||
t.Errorf("log export authorization = %v, want exactly [Bearer logs-token]", got)
|
||||
}
|
||||
if got := logMD.Get("x-tenant"); len(got) != 0 {
|
||||
t.Errorf("log export x-tenant = %v, want absent: %s replaces %s rather than merging", got, logsHeadersEnv, headersEnv)
|
||||
}
|
||||
}
|
||||
|
||||
func TestParseHeaders(t *testing.T) {
|
||||
|
||||
Reference in New Issue
Block a user