upgrade otel log to 0.21 (#1958)

fixes https://github.com/agent-substrate/substrate/pull/1952

security bump stuck by breaking changes.

---------

Signed-off-by: dependabot[bot] <support@github.com>
Co-authored-by: dependabot[bot] <49699333+dependabot[bot]@users.noreply.github.com>
This commit is contained in:
Benjamin Elder
2026-09-28 22:26:40 +00:00
committed by GitHub
co-authored by dependabot[bot]
parent dd57a119ae
commit 51a71cb365
55 changed files with 1457 additions and 4687 deletions
+3 -2
View File
@@ -32,6 +32,7 @@ import (
"github.com/agent-substrate/substrate/internal/resources"
"github.com/agent-substrate/substrate/pkg/proto/ateapipb"
"github.com/google/go-cmp/cmp"
"go.opentelemetry.io/otel/attribute"
otellog "go.opentelemetry.io/otel/log"
"go.opentelemetry.io/otel/log/global"
sdklog "go.opentelemetry.io/otel/sdk/log"
@@ -633,8 +634,8 @@ func (otlpSinkExporter) Export(_ context.Context, records []sdklog.Record) error
timestamp: r.Timestamp(),
attrs: map[string]string{},
}
r.WalkAttributes(func(kv otellog.KeyValue) bool {
e.attrs[kv.Key] = kv.Value.String()
r.WalkAttributes(func(kv attribute.KeyValue) bool {
e.attrs[string(kv.Key)] = kv.Value.String()
return true
})
otlpSink = append(otlpSink, e)
+3 -3
View File
@@ -50,14 +50,14 @@ require (
go.opentelemetry.io/contrib/instrumentation/google.golang.org/grpc/otelgrpc v0.70.0
go.opentelemetry.io/contrib/instrumentation/net/http/otelhttp v0.70.0
go.opentelemetry.io/otel v1.45.0
go.opentelemetry.io/otel/exporters/otlp/otlplog/otlploggrpc v0.20.0
go.opentelemetry.io/otel/exporters/otlp/otlplog/otlploggrpc v0.21.0
go.opentelemetry.io/otel/exporters/otlp/otlpmetric/otlpmetricgrpc v1.43.0
go.opentelemetry.io/otel/exporters/otlp/otlptrace/otlptracegrpc v1.45.0
go.opentelemetry.io/otel/exporters/prometheus v0.65.0
go.opentelemetry.io/otel/log v0.20.0
go.opentelemetry.io/otel/log v0.21.0
go.opentelemetry.io/otel/metric v1.45.0
go.opentelemetry.io/otel/sdk v1.45.0
go.opentelemetry.io/otel/sdk/log v0.20.0
go.opentelemetry.io/otel/sdk/log v0.21.0
go.opentelemetry.io/otel/sdk/metric v1.45.0
go.opentelemetry.io/otel/trace v1.45.0
go.opentelemetry.io/proto/otlp v1.11.0
+8 -8
View File
@@ -520,8 +520,8 @@ go.opentelemetry.io/contrib/instrumentation/net/http/otelhttp v0.70.0 h1:LMuyCAy
go.opentelemetry.io/contrib/instrumentation/net/http/otelhttp v0.70.0/go.mod h1:085m8qbm4hgc8rZWGDEa4vmyyo2c3nPxUslYUKUIU04=
go.opentelemetry.io/otel v1.45.0 h1:pdrWmLHofpubmArBv1LgFSv1Z0Ie/ppdZzu+kUN5EeU=
go.opentelemetry.io/otel v1.45.0/go.mod h1:XZxIqPapzEYnhNSScF5DIqXhm/rYi0FzCe2XddAwZfQ=
go.opentelemetry.io/otel/exporters/otlp/otlplog/otlploggrpc v0.20.0 h1:rydZ9sxbcFdm/oWrVyfLTjHIygMgv0bEeMd+3B/BvoM=
go.opentelemetry.io/otel/exporters/otlp/otlplog/otlploggrpc v0.20.0/go.mod h1:earQ25dooT0Hhspq59DZ8YCC50jWfOlFEeWoxy/P444=
go.opentelemetry.io/otel/exporters/otlp/otlplog/otlploggrpc v0.21.0 h1:WseeVYf5dJZTsyPiyW5L14k5qsSibqXAMTSiFEDiWr0=
go.opentelemetry.io/otel/exporters/otlp/otlplog/otlploggrpc v0.21.0/go.mod h1:SiLZnQS6Qk2eCpvr2CH/XMAOa64TWGXxEZJZCpD2Lmc=
go.opentelemetry.io/otel/exporters/otlp/otlpmetric/otlpmetricgrpc v1.43.0 h1:8UQVDcZxOJLtX6gxtDt3vY2WTgvZqMQRzjsqiIHQdkc=
go.opentelemetry.io/otel/exporters/otlp/otlpmetric/otlpmetricgrpc v1.43.0/go.mod h1:2lmweYCiHYpEjQ/lSJBYhj9jP1zvCvQW4BqL9dnT7FQ=
go.opentelemetry.io/otel/exporters/otlp/otlptrace v1.45.0 h1:QRefszxJmfPdjXUUm3j6iDzY03mTPXMjqErFqQ67vUg=
@@ -532,18 +532,18 @@ go.opentelemetry.io/otel/exporters/prometheus v0.65.0 h1:jOveH/b4lU9HT7y+Gfamf18
go.opentelemetry.io/otel/exporters/prometheus v0.65.0/go.mod h1:i1P8pcumauPtUI4YNopea1dhzEMuEqWP1xoUZDylLHo=
go.opentelemetry.io/otel/exporters/stdout/stdoutmetric v1.43.0 h1:TC+BewnDpeiAmcscXbGMfxkO+mwYUwE/VySwvw88PfA=
go.opentelemetry.io/otel/exporters/stdout/stdoutmetric v1.43.0/go.mod h1:J/ZyF4vfPwsSr9xJSPyQ4LqtcTPULFR64KwTikGLe+A=
go.opentelemetry.io/otel/log v0.20.0 h1:/5i0vuHxCLWUfChWG41K9wkM0jafruPw9NU1/RCJirs=
go.opentelemetry.io/otel/log v0.20.0/go.mod h1:wOcMcjsZpG8x7Bak7IhSi/lg8wscV2C1VdrKCLPlt0E=
go.opentelemetry.io/otel/log v0.21.0 h1:SLsVDGmtyBrdw8/a2Z0bOIxou/+bN4z56GebH7T0LvA=
go.opentelemetry.io/otel/log v0.21.0/go.mod h1:iReetQrZL9Wyg84cCkOoCmqDHS5RCFfyxC7J+r8fn8g=
go.opentelemetry.io/otel/metric v1.45.0 h1:7Eg1uH7CJ5cXv9is6tnBe1FI6rj1nwUdbFypRm3br/M=
go.opentelemetry.io/otel/metric v1.45.0/go.mod h1:HAPbm1nd3p1PmFH7v2dR+6BjXxw+Lq4a2+pndMAm08s=
go.opentelemetry.io/otel/metric/x v0.67.0 h1:PcicCNZFkZ4bXfSooXdo3WN7RBOVOtjVdo1wD358Uns=
go.opentelemetry.io/otel/metric/x v0.67.0/go.mod h1:FBjCWZe6wgcqxcMtjdGiClDKXb2YxxXii0CXftE4QtI=
go.opentelemetry.io/otel/sdk v1.45.0 h1:4VVSMgQ83dUgW2aoX5f6JgLvHwIvzcuLnF9lUdCSpCw=
go.opentelemetry.io/otel/sdk v1.45.0/go.mod h1:Sr40LgXV7DsKMMJMKOhUWOgMWTfAaqvm2kF0g7ilwuA=
go.opentelemetry.io/otel/sdk/log v0.20.0 h1:vM3xI7TQgKPiSghe6urZtAkyFY7SodrSpC83CffDFuY=
go.opentelemetry.io/otel/sdk/log v0.20.0/go.mod h1:Knej2nmsTUzN79T2eeXdRsjjPcoxoq2pUyUHz9TFyyU=
go.opentelemetry.io/otel/sdk/log/logtest v0.20.0 h1:OqdRZ1guyzamK3M6LlRsmGqRrjkHWw6WZOKKli5ELpg=
go.opentelemetry.io/otel/sdk/log/logtest v0.20.0/go.mod h1:PuMIlm7zAt7c3z8zfOI5ox4iT1Z87We+PF6YoINux/M=
go.opentelemetry.io/otel/sdk/log v0.21.0 h1:QsE7XSR0ktQdKmRKGnR+f1ObGF32WG+7MER/P9KgmYc=
go.opentelemetry.io/otel/sdk/log v0.21.0/go.mod h1:m9mApjCoD2/1QuKCAptjv+BrG9WKOvQLVdNx+iBldTo=
go.opentelemetry.io/otel/sdk/log/logtest v0.21.0 h1:X+JBBgKlswCGYsmgL0CnoUUtlE//VB345c84jYAYkdQ=
go.opentelemetry.io/otel/sdk/log/logtest v0.21.0/go.mod h1:HD1575K8e6sIFBBDd5tZB3t9DlMytWXq9FuR+Y4rfjE=
go.opentelemetry.io/otel/sdk/metric v1.45.0 h1:oVFszMfyj1Am6s24Vtc7wBb8BKLcwepJjNEYILuiE3o=
go.opentelemetry.io/otel/sdk/metric v1.45.0/go.mod h1:vUWUxDZvu1WVRj8JA8S0AdhsPrZoDpA2DdZauIh4mDA=
go.opentelemetry.io/otel/trace v1.45.0 h1:l/mP6Uv7oNO7/TblbhpbgMidxhq1uO/rPsikOyVhxag=
+13 -12
View File
@@ -33,6 +33,7 @@ import (
"sync"
"time"
"go.opentelemetry.io/otel/attribute"
"go.opentelemetry.io/otel/log"
"go.opentelemetry.io/otel/log/global"
@@ -142,40 +143,40 @@ func BuildRecord(ev Event, t time.Time, attrs []slog.Attr) log.Record {
rec.SetEventName(ev.Name)
rec.SetTimestamp(t)
rec.SetSeverity(ev.Severity)
rec.SetBody(log.StringValue(ev.Body))
rec.SetBody(attribute.StringValue(ev.Body))
kvs := make([]log.KeyValue, 0, len(attrs))
kvs := make([]attribute.KeyValue, 0, len(attrs))
for _, a := range attrs {
kvs = append(kvs, log.KeyValue{Key: a.Key, Value: logValue(a.Value)})
kvs = append(kvs, attribute.KeyValue{Key: attribute.Key(a.Key), Value: logValue(a.Value)})
}
rec.AddAttributes(kvs...)
return rec
}
// logValue keeps the kind slog's JSON handler writes, so the two copies match.
func logValue(v slog.Value) log.Value {
func logValue(v slog.Value) attribute.Value {
switch v.Kind() {
case slog.KindString:
return log.StringValue(v.String())
return attribute.StringValue(v.String())
case slog.KindInt64:
return log.Int64Value(v.Int64())
return attribute.Int64Value(v.Int64())
case slog.KindUint64:
// OTel has no unsigned kind. Clamp rather than wrap, as the metric
// side does in addSat.
u := v.Uint64()
if u > math.MaxInt64 {
return log.Int64Value(math.MaxInt64)
return attribute.Int64Value(math.MaxInt64)
}
return log.Int64Value(int64(u))
return attribute.Int64Value(int64(u))
case slog.KindFloat64:
return log.Float64Value(v.Float64())
return attribute.Float64Value(v.Float64())
case slog.KindBool:
return log.BoolValue(v.Bool())
return attribute.BoolValue(v.Bool())
case slog.KindDuration:
// nanoseconds, not "1.5s"
return log.Int64Value(int64(v.Duration()))
return attribute.Int64Value(int64(v.Duration()))
default:
return log.StringValue(v.String())
return attribute.StringValue(v.String())
}
}
+19 -18
View File
@@ -23,6 +23,7 @@ import (
"testing"
"time"
"go.opentelemetry.io/otel/attribute"
"go.opentelemetry.io/otel/log"
sdklog "go.opentelemetry.io/otel/sdk/log"
sdktrace "go.opentelemetry.io/otel/sdk/trace"
@@ -90,8 +91,8 @@ func usagePendingAttrs() []slog.Attr {
func recordAttrs(rec log.Record) map[string]string {
got := make(map[string]string, rec.AttributesLen())
rec.WalkAttributes(func(kv log.KeyValue) bool {
got[kv.Key] = kv.Value.String()
rec.WalkAttributes(func(kv attribute.KeyValue) bool {
got[string(kv.Key)] = kv.Value.String()
return true
})
return got
@@ -346,8 +347,8 @@ func TestLogWritesBothCopies(t *testing.T) {
return true
})
otlpAttrs := map[string]string{}
otlpRec.WalkAttributes(func(kv log.KeyValue) bool {
otlpAttrs[kv.Key] = kv.Value.String()
otlpRec.WalkAttributes(func(kv attribute.KeyValue) bool {
otlpAttrs[string(kv.Key)] = kv.Value.String()
return true
})
if !maps.Equal(stdoutAttrs, otlpAttrs) {
@@ -387,7 +388,7 @@ func TestEmitCarriesTraceContext(t *testing.T) {
// Trace context belongs on the record's own fields. The stdout copy carries
// it as attributes; the OTLP copy must not, or it is there twice.
rec.WalkAttributes(func(kv log.KeyValue) bool {
rec.WalkAttributes(func(kv attribute.KeyValue) bool {
switch kv.Key {
case ateattr.LogTraceIDField, ateattr.LogSpanIDField, ateattr.LogTraceFlagsField:
t.Errorf("record carries trace context as the attribute %q", kv.Key)
@@ -414,16 +415,16 @@ func TestBuildRecordKeepsValueKinds(t *testing.T) {
tests := []struct {
name string
attr slog.Attr
want log.Value
want attribute.Value
}{
{"string", slog.String("k", "v"), log.StringValue("v")},
{"int", slog.Int64("k", 7), log.Int64Value(7)},
{"uint", slog.Uint64("k", 7), log.Int64Value(7)},
{"uint above int64 clamps", slog.Uint64("k", math.MaxUint64), log.Int64Value(math.MaxInt64)},
{"float", slog.Float64("k", 1.5), log.Float64Value(1.5)},
{"bool", slog.Bool("k", true), log.BoolValue(true)},
{"duration is nanoseconds, as in the stdout copy", slog.Duration("k", 1500*time.Millisecond), log.Int64Value(1_500_000_000)},
{"anything else falls back to its string form", slog.Any("k", struct{}{}), log.StringValue("{}")},
{"string", slog.String("k", "v"), attribute.StringValue("v")},
{"int", slog.Int64("k", 7), attribute.Int64Value(7)},
{"uint", slog.Uint64("k", 7), attribute.Int64Value(7)},
{"uint above int64 clamps", slog.Uint64("k", math.MaxUint64), attribute.Int64Value(math.MaxInt64)},
{"float", slog.Float64("k", 1.5), attribute.Float64Value(1.5)},
{"bool", slog.Bool("k", true), attribute.BoolValue(true)},
{"duration is nanoseconds, as in the stdout copy", slog.Duration("k", 1500*time.Millisecond), attribute.Int64Value(1_500_000_000)},
{"anything else falls back to its string form", slog.Any("k", struct{}{}), attribute.StringValue("{}")},
}
for _, tt := range tests {
@@ -431,13 +432,13 @@ func TestBuildRecordKeepsValueKinds(t *testing.T) {
t.Parallel()
rec := BuildRecord(StateChanged, time.Now(), []slog.Attr{tt.attr})
var got log.Value
rec.WalkAttributes(func(kv log.KeyValue) bool {
var got attribute.Value
rec.WalkAttributes(func(kv attribute.KeyValue) bool {
got = kv.Value
return false
})
if got.Kind() != tt.want.Kind() {
t.Fatalf("kind = %v, want %v", got.Kind(), tt.want.Kind())
if got.Type() != tt.want.Type() {
t.Fatalf("type = %v, want %v", got.Type(), tt.want.Type())
}
if got.String() != tt.want.String() {
t.Errorf("value = %v, want %v", got, tt.want)
@@ -1,7 +1,7 @@
// Copyright The OpenTelemetry Authors
// SPDX-License-Identifier: Apache-2.0
package otlploggrpc // import "go.opentelemetry.io/otel/exporters/otlp/otlplog/otlploggrpc"
package otlploggrpc
import (
"context"
@@ -103,10 +103,12 @@ func newGRPCDialOptions(cfg config) []grpc.DialOption {
if cfg.serviceConfig.Value != "" {
dialOpts = append(dialOpts, grpc.WithDefaultServiceConfig(cfg.serviceConfig.Value))
}
// Prioritize GRPCCredentials over Insecure (passing both is an error).
// Prioritize configured credentials over Insecure (passing both is an error).
switch {
case cfg.gRPCCredentials.Value != nil:
dialOpts = append(dialOpts, grpc.WithTransportCredentials(cfg.gRPCCredentials.Value))
case cfg.tlsCfg.Value != nil:
dialOpts = append(dialOpts, grpc.WithTransportCredentials(credentials.NewTLS(cfg.tlsCfg.Value)))
case cfg.insecure.Value:
dialOpts = append(dialOpts, grpc.WithTransportCredentials(insecure.NewCredentials()))
default:
@@ -1,7 +1,7 @@
// Copyright The OpenTelemetry Authors
// SPDX-License-Identifier: Apache-2.0
package otlploggrpc // import "go.opentelemetry.io/otel/exporters/otlp/otlplog/otlploggrpc"
package otlploggrpc
import (
"crypto/tls"
+1 -1
View File
@@ -60,4 +60,4 @@ The configuration can be overridden by [WithTLSCredentials], [WithGRPCConn] opti
[W3C Baggage HTTP Header Content Format]: https://www.w3.org/TR/baggage/#header-content
*/
package otlploggrpc // import "go.opentelemetry.io/otel/exporters/otlp/otlplog/otlploggrpc"
package otlploggrpc
@@ -1,7 +1,7 @@
// Copyright The OpenTelemetry Authors
// SPDX-License-Identifier: Apache-2.0
package otlploggrpc // import "go.opentelemetry.io/otel/exporters/otlp/otlplog/otlploggrpc"
package otlploggrpc
import (
"context"
@@ -3,9 +3,9 @@
// Package internal provides internal functionality for the otlploggrpc
// package.
package internal // import "go.opentelemetry.io/otel/exporters/otlp/otlplog/otlploggrpc/internal"
package internal
//go:generate gotmpl --body=../../../../../internal/shared/otlp/observ/target.go.tmpl "--data={ \"pkg\": \"observ\", \"pkg_path\": \"go.opentelemetry.io/otel/exporters/otlp/otlplog/otlploggrpc/internal/observ\" }" --out=observ/target.go
//go:generate gotmpl --body=../../../../../internal/shared/otlp/observ/target.go.tmpl "--data={ \"pkg\": \"observ\" }" --out=observ/target.go
//go:generate gotmpl --body=../../../../../internal/shared/otlp/observ/target_test.go.tmpl "--data={ \"pkg\": \"observ\" }" --out=observ/target_test.go
//go:generate gotmpl --body=../../../../../internal/shared/x/x.go.tmpl "--data={ \"pkg\": \"go.opentelemetry.io/otel/exporters/otlp/otlplog/otlploggrpc\" }" --out=x/x.go
@@ -3,7 +3,7 @@
// Package observ provides observability metrics for OTLP log exporters.
// This is an experimental feature controlled by the x.Observability feature flag.
package observ // import "go.opentelemetry.io/otel/exporters/otlp/otlplog/otlploggrpc/internal/observ"
package observ
import (
"context"
@@ -21,8 +21,8 @@ import (
"go.opentelemetry.io/otel/exporters/otlp/otlplog/otlploggrpc/internal/x"
"go.opentelemetry.io/otel/internal/global"
"go.opentelemetry.io/otel/metric"
semconv "go.opentelemetry.io/otel/semconv/v1.41.0"
"go.opentelemetry.io/otel/semconv/v1.41.0/otelconv"
semconv "go.opentelemetry.io/otel/semconv/v1.43.0"
"go.opentelemetry.io/otel/semconv/v1.43.0/otelconv"
)
const (
@@ -1,10 +1,10 @@
// Code generated by gotmpl. DO NOT MODIFY.
// source: internal/shared/otlp/observ/target.go.tmpl
// Copyright The OpenTelemetry Authors
// SPDX-License-Identifier: Apache-2.0
package observ // import "go.opentelemetry.io/otel/exporters/otlp/otlplog/otlploggrpc/internal/observ"
// DO NOT MODIFY. Generated by gotmpl.
// source: internal/shared/otlp/observ/target.go.tmpl
package observ
import (
"errors"
@@ -41,11 +41,10 @@ func ParseCanonicalTarget(target string) (string, int, error) {
const sep = "://"
// Find scheme. Do not allocate the string by using url.Parse.
idx := strings.Index(target, sep)
if idx == -1 {
scheme, endpoint, found := strings.Cut(target, sep)
if !found {
return "", -1, fmt.Errorf("invalid target %q: missing scheme", target)
}
scheme, endpoint := target[:idx], target[idx+len(sep):]
// Check for unix schemes.
if scheme == schemeUnix || scheme == schemeUnixAbstract {
@@ -1,7 +1,7 @@
// Copyright The OpenTelemetry Authors
// SPDX-License-Identifier: Apache-2.0
package internal // import "go.opentelemetry.io/otel/exporters/otlp/otlplog/otlploggrpc/internal"
package internal
import "fmt"
@@ -1,13 +1,13 @@
// Code generated by gotmpl. DO NOT MODIFY.
// source: internal/shared/otlp/retry/retry.go.tmpl
// Copyright The OpenTelemetry Authors
// SPDX-License-Identifier: Apache-2.0
// DO NOT MODIFY. Generated by gotmpl.
// source: internal/shared/otlp/retry/retry.go.tmpl
// Package retry provides request retry functionality that can perform
// configurable exponential backoff for transient errors and honor any
// explicit throttle responses received.
package retry // import "go.opentelemetry.io/otel/exporters/otlp/otlplog/otlploggrpc/internal/retry"
package retry
import (
"context"
@@ -1,12 +1,12 @@
// Code generated by gotmpl. DO NOT MODIFY.
// source: internal/shared/otlp/otlplog/transform/log.go.tmpl
// Copyright The OpenTelemetry Authors
// SPDX-License-Identifier: Apache-2.0
// DO NOT MODIFY. Generated by gotmpl.
// source: internal/shared/otlp/otlplog/transform/log.go.tmpl
// Package transform provides transformation functionality from the
// sdk/log data-types into OTLP data-types.
package transform // import "go.opentelemetry.io/otel/exporters/otlp/otlplog/otlploggrpc/internal/transform"
package transform
import (
"time"
@@ -94,13 +94,13 @@ func LogRecord(record log.Record) *lpb.LogRecord {
EventName: record.EventName(),
SeverityNumber: SeverityNumber(record.Severity()),
SeverityText: record.SeverityText(),
Body: LogAttrValue(record.Body()),
Body: AttrValue(record.Body()),
Attributes: make([]*cpb.KeyValue, 0, record.AttributesLen()),
Flags: uint32(record.TraceFlags()),
// TODO: DroppedAttributesCount: /* ... */,
}
record.WalkAttributes(func(kv api.KeyValue) bool {
r.Attributes = append(r.Attributes, LogAttr(kv))
record.WalkAttributes(func(kv attribute.KeyValue) bool {
r.Attributes = append(r.Attributes, Attr(kv))
return true
})
if tID := record.TraceID(); tID.IsValid() {
@@ -205,6 +205,12 @@ func AttrValue(v attribute.Value) *cpb.AnyValue {
Values: attrValues(v.AsSlice()),
},
}
case attribute.MAP:
av.Value = &cpb.AnyValue_KvlistValue{
KvlistValue: &cpb.KeyValueList{
Values: Attrs(v.AsMap()),
},
}
case attribute.STRINGSLICE:
av.Value = &cpb.AnyValue_ArrayValue{
ArrayValue: &cpb.ArrayValue{
@@ -276,85 +282,6 @@ func attrValues(vals []attribute.Value) []*cpb.AnyValue {
return converted
}
// LogAttrs transforms a slice of [api.KeyValue] into OTLP key-values.
func LogAttrs(attrs []api.KeyValue) []*cpb.KeyValue {
if len(attrs) == 0 {
return nil
}
out := make([]*cpb.KeyValue, 0, len(attrs))
for _, kv := range attrs {
out = append(out, LogAttr(kv))
}
return out
}
// LogAttr transforms an [api.KeyValue] into an OTLP key-value.
func LogAttr(attr api.KeyValue) *cpb.KeyValue {
return &cpb.KeyValue{
Key: attr.Key,
Value: LogAttrValue(attr.Value),
}
}
// LogAttrValues transforms a slice of [api.Value] into an OTLP []AnyValue.
func LogAttrValues(vals []api.Value) []*cpb.AnyValue {
if len(vals) == 0 {
return nil
}
out := make([]*cpb.AnyValue, 0, len(vals))
for _, v := range vals {
out = append(out, LogAttrValue(v))
}
return out
}
// LogAttrValue transforms an [api.Value] into an OTLP AnyValue.
func LogAttrValue(v api.Value) *cpb.AnyValue {
av := new(cpb.AnyValue)
switch v.Kind() {
case api.KindBool:
av.Value = &cpb.AnyValue_BoolValue{
BoolValue: v.AsBool(),
}
case api.KindInt64:
av.Value = &cpb.AnyValue_IntValue{
IntValue: v.AsInt64(),
}
case api.KindFloat64:
av.Value = &cpb.AnyValue_DoubleValue{
DoubleValue: v.AsFloat64(),
}
case api.KindString:
av.Value = &cpb.AnyValue_StringValue{
StringValue: v.AsString(),
}
case api.KindBytes:
av.Value = &cpb.AnyValue_BytesValue{
BytesValue: v.AsBytes(),
}
case api.KindSlice:
av.Value = &cpb.AnyValue_ArrayValue{
ArrayValue: &cpb.ArrayValue{
Values: LogAttrValues(v.AsSlice()),
},
}
case api.KindMap:
av.Value = &cpb.AnyValue_KvlistValue{
KvlistValue: &cpb.KeyValueList{
Values: LogAttrs(v.AsMap()),
},
}
case api.KindEmpty:
default:
av.Value = &cpb.AnyValue_StringValue{
StringValue: "INVALID",
}
}
return av
}
// SeverityNumber transforms a [log.Severity] into an OTLP SeverityNumber.
func SeverityNumber(s api.Severity) lpb.SeverityNumber {
switch s {
@@ -1,8 +1,8 @@
// Copyright The OpenTelemetry Authors
// SPDX-License-Identifier: Apache-2.0
package internal // import "go.opentelemetry.io/otel/exporters/otlp/otlplog/otlploggrpc/internal"
package internal
// Version is the current release version of the OpenTelemetry otlploggrpc
// exporter in use.
const Version = "0.20.0"
const Version = "0.21.0"
@@ -2,7 +2,7 @@
// SPDX-License-Identifier: Apache-2.0
// Package x documents experimental features for [go.opentelemetry.io/otel/exporters/otlp/otlplog/otlploggrpc].
package x // import "go.opentelemetry.io/otel/exporters/otlp/otlplog/otlploggrpc/internal/x"
package x
import "strings"
@@ -1,11 +1,11 @@
// Code generated by gotmpl. DO NOT MODIFY.
// source: internal/shared/x/x.go.tmpl
// Copyright The OpenTelemetry Authors
// SPDX-License-Identifier: Apache-2.0
// DO NOT MODIFY. Generated by gotmpl.
// source: internal/shared/x/x.go.tmpl
// Package x documents experimental features for [go.opentelemetry.io/otel/exporters/otlp/otlplog/otlploggrpc].
package x // import "go.opentelemetry.io/otel/exporters/otlp/otlplog/otlploggrpc/internal/x"
package x
import (
"os"
@@ -1,7 +1,7 @@
// Copyright The OpenTelemetry Authors
// SPDX-License-Identifier: Apache-2.0
package otlploggrpc // import "go.opentelemetry.io/otel/exporters/otlp/otlplog/otlploggrpc"
package otlploggrpc
// Version is the current release version of the OpenTelemetry OTLP over gRPC logs exporter in use.
func Version() string {
+41 -109
View File
@@ -62,10 +62,10 @@ Implementation requirements:
- The [specification requires](https://opentelemetry.io/docs/specs/otel/logs/api/#concurrency-requirements)
the method to be safe to be called concurrently.
- The method should use some default name if the passed name is empty
in order to meet the [specification's SDK requirement](https://opentelemetry.io/docs/specs/otel/logs/sdk/#logger-creation)
to return a working logger when an invalid name is passed
as well as to resemble the behavior of getting tracers and meters.
- If the passed name is empty, this should be reported as invalid as specified
by the [Logs SDK specification](https://opentelemetry.io/docs/specs/otel/logs/sdk/#logger-creation).
The method should still use it as the instrumentation scope name and return
a working logger.
`Logger` can be extended by adding new `LoggerOption` options
and adding new exported fields to the `LoggerConfig` struct.
@@ -148,16 +148,16 @@ func (r *Record) SetSeverityText(s string)
is accessed using following methods:
```go
func (r *Record) Body() Value
func (r *Record) SetBody(v Value)
func (r *Record) Body() attribute.Value
func (r *Record) SetBody(v attribute.Value)
```
[Log record attributes](https://opentelemetry.io/docs/specs/otel/logs/data-model/#field-attributes)
are accessed using following methods:
```go
func (r *Record) WalkAttributes(f func(KeyValue) bool)
func (r *Record) AddAttributes(attrs ...KeyValue)
func (r *Record) WalkAttributes(f func(attribute.KeyValue) bool)
func (r *Record) AddAttributes(attrs ...attribute.KeyValue)
```
`Record` has a `AttributesLen` method that returns
@@ -175,49 +175,20 @@ while keeping the API user friendly.
It relieves the user from making his own improvements
for reducing the number of allocations when passing attributes.
The abstractions described in
[the specification](https://opentelemetry.io/docs/specs/otel/logs/#new-first-party-application-logs)
are defined in [keyvalue.go](keyvalue.go).
Log body and attributes use the common
[`attribute.Value`](https://pkg.go.dev/go.opentelemetry.io/otel/attribute#Value)
and
[`attribute.KeyValue`](https://pkg.go.dev/go.opentelemetry.io/otel/attribute#KeyValue)
types.
These types cover the Logs Data Model `any` value shape, including empty
values, byte slices, generic slices, and maps.
Reusing the attribute package keeps the API consistent across signals and
avoids maintaining a second value model with conversion helpers.
`Value` is representing `any`.
`KeyValue` is representing a key(string)-value(`any`) pair.
`Kind` is an enumeration used for specifying the underlying value type.
`KindEmpty` is used for an empty (zero) value.
`KindBool` is used for boolean value.
`KindFloat64` is used for a double precision floating point (IEEE 754-1985) value.
`KindInt64` is used for a signed integer value.
`KindString` is used for a string value.
`KindBytes` is used for a slice of bytes (in spec: A byte array).
`KindSlice` is used for a slice of values (in spec: an array (a list) of any values).
`KindMap` is used for a slice of key-value pairs (in spec: `map<string, any>`).
These types are defined in `go.opentelemetry.io/otel/log` package
as they are tightly coupled with the API and different from common attributes.
The internal implementation of `Value` is based on
[`slog.Value`](https://pkg.go.dev/log/slog#Value)
and the API is mostly inspired by
[`attribute.Value`](https://pkg.go.dev/go.opentelemetry.io/otel/attribute#Value).
The benchmarks[^1] show that the implementation is more performant than
[`attribute.Value`](https://pkg.go.dev/go.opentelemetry.io/otel/attribute#Value).
The value accessors (`func (v Value) As[Kind]` methods) must not panic,
as it would violate the [specification](https://opentelemetry.io/docs/specs/otel/error-handling/):
> API methods MUST NOT throw unhandled exceptions when used incorrectly by end
> users. The API and SDK SHOULD provide safe defaults for missing or invalid
> arguments. [...] Whenever the library suppresses an error that would otherwise
> have been exposed to the user, the library SHOULD log the error using
> language-specific conventions.
Therefore, the value accessors should return a zero value
and log an error when a bad accessor is called.
The `Severity`, `Kind`, `Value`, `KeyValue` may implement
the [`fmt.Stringer`](https://pkg.go.dev/fmt#Stringer) interface.
However, it is not needed for the first stable release
and the `String` methods can be added later.
The zero value of `attribute.Value` represents an empty body.
Log maps use `attribute.MAP`, which may contain duplicate keys when callers
construct one that way. Duplicate-key normalization is an SDK policy, not an
API behavior.
The caller must not subsequently mutate the record passed to `Emit`.
This would allow the implementation to not clone the record,
@@ -252,12 +223,11 @@ Rejected alternatives:
- [Passing record as pointer to Logger.Emit](#passing-record-as-pointer-to-loggeremit)
- [Logger.WithAttributes](#loggerwithattributes)
- [Record attributes as slice](#record-attributes-as-slice)
- [Use any instead of defining Value](#use-any-instead-of-defining-value)
- [Use any instead of attribute.Value](#use-any-instead-of-attributevalue)
- [Severity type encapsulating number and text](#severity-type-encapsulating-number-and-text)
- [Reuse attribute package](#reuse-attribute-package)
- [Define log-specific value types](#define-log-specific-value-types)
- [Mix receiver types for Record](#mix-receiver-types-for-record)
- [Add XYZ method to Logger](#add-xyz-method-to-logger)
- [Rename KeyValue to Attr](#rename-keyvalue-to-attr)
### Logger.Enabled
@@ -476,30 +446,21 @@ less user friendly (users and bridges would use e.g. a `sync.Pool` to reduce
the number of heap allocation), less safe (more prone to use after free bugs
and race conditions), and the benchmark differences were not significant.
### Use any instead of defining Value
### Use any instead of attribute.Value
[Logs Data Model](https://opentelemetry.io/docs/specs/otel/logs/data-model/#field-body)
defines Body to be `any`.
One could propose to define `Body` (and attribute values) as `any`
instead of a defining a new type (`Value`).
instead of using `attribute.Value`.
First of all, [`any` type defined in the specification](https://opentelemetry.io/docs/specs/otel/logs/data-model/#type-any)
is not the same as `any` (`interface{}`) in Go.
Moreover, using `any` as a field would decrease the performance.[^7]
Notice it will be still possible to add following kind and factories
in a backwards compatible way:
```go
const KindMap Kind
func AnyValue(value any) KeyValue
func Any(key string, value any) KeyValue
```
However, currently, it would not be specification compliant.
Using `attribute.Value` preserves a typed, allocation-conscious representation
of the Logs Data Model `any` values while avoiding unconstrained `interface{}`
handling in bridge implementations.
### Severity type encapsulating number and text
@@ -518,20 +479,22 @@ It should be more user friendly to have them separated.
Especially when having getter and setter methods, setting one value
when the other is already set would be unpleasant.
### Reuse attribute package
### Define log-specific value types
It was tempting to reuse the existing
[https://pkg.go.dev/go.opentelemetry.io/otel/attribute] package
for defining log attributes and body.
The original design defined `Kind`, `Value`, and `KeyValue` in
`go.opentelemetry.io/otel/log`.
That avoided coupling log bodies to the common attribute package while logs were
still exploring structured value support.
However, this would be wrong because [the log attribute definition](https://opentelemetry.io/docs/specs/otel/logs/data-model/#field-attributes)
is different from [the common attribute definition](https://opentelemetry.io/docs/specs/otel/common/#attribute).
The Logs Data Model and common attribute value model now share the structured
`any` shapes needed by Go: empty, bool, int64, float64, string, byte slice,
homogeneous slices, generic slices, and maps.
The specification direction is to reuse these value shapes across signals.
Keeping log-specific types would duplicate API surface, require conversion
helpers, and make bridge code choose between two equivalent value models.
Moreover, it there is nothing telling that [the body definition](https://opentelemetry.io/docs/specs/otel/logs/data-model/#field-body)
has anything in common with a common attribute value.
Therefore, we define new types representing the abstract types defined
in the [Logs Data Model](https://opentelemetry.io/docs/specs/otel/logs/data-model/#definitions-used-in-this-document).
Therefore log records now use `attribute.Value` for body values and
`attribute.KeyValue` for attributes and map entries.
### Mix receiver types for Record
@@ -591,36 +554,6 @@ The `Logger` does not have methods like `SetSeverity`, etc.
as the Logs API needs to follow (be compliant with)
the [specification](https://opentelemetry.io/docs/specs/otel/logs/api/)
### Rename KeyValue to Attr
There was a proposal to rename `KeyValue` to `Attr` (or `Attribute`).[^11]
New developers may not intuitively know that `log.KeyValue` is an attribute in
the OpenTelemetry parlance.
During the discussion we agreed to keep the `KeyValue` name.
The type is used in multiple semantics:
- as a log attribute,
- as a map item,
- as a log record Body.
As for map item semantics, this type is a key-value pair, not an attribute.
Naming the type as `Attr` would convey semantical meaning
that would not be correct for a map.
We expect that most of the Logs API users will be OpenTelemetry contributors.
We plan to implement bridges for the most popular logging libraries ourselves.
Given we will all have the context needed to disambiguate these overlapping
names, developers' confusion should not be an issue.
For bridges not developed by us,
developers will likely look at our existing bridges for inspiration.
Our correct use of these types will be a reference to them.
At last, we provide `ValueFromAttribute` and `KeyValueFromAttribute`
to offer reuse of `attribute.Value` and `attribute.KeyValue`.
[^1]: [Handle structured body and attributes](https://github.com/pellared/opentelemetry-go/pull/7)
[^2]: Jonathan Amsterdam, [The Go Blog: Structured Logging with slog](https://go.dev/blog/slog)
[^3]: Jonathan Amsterdam, [GopherCon Europe 2023: A Fast Structured Logging Package](https://www.youtube.com/watch?v=tC4Jt3i62ns)
@@ -631,4 +564,3 @@ to offer reuse of `attribute.Value` and `attribute.KeyValue`.
[^8]: [log/slog: structured, leveled logging](https://github.com/golang/go/issues/56345#issuecomment-1302563756)
[^9]: [Record with pointer receivers only](https://github.com/pellared/opentelemetry-go/pull/8)
[^10]: [Go FAQ: Stack or heap](https://go.dev/doc/faq#stack_or_heap)
[^11]: [Rename KeyValue to Attr discussion](https://github.com/open-telemetry/opentelemetry-go/pull/4809#discussion_r1476080093)
+1 -1
View File
@@ -82,4 +82,4 @@ implement all the API interfaces when a user updates their API.
[registry]: https://opentelemetry.io/ecosystem/registry/?language=go&component=log-bridge
*/
package log // import "go.opentelemetry.io/otel/log"
package log
+1 -1
View File
@@ -11,7 +11,7 @@
// bump of the API package).
//
// [OpenTelemetry Logs API]: https://pkg.go.dev/go.opentelemetry.io/otel/log
package embedded // import "go.opentelemetry.io/otel/log/embedded"
package embedded
// LoggerProvider is embedded in the [Logs API LoggerProvider].
//
+1 -1
View File
@@ -9,7 +9,7 @@ This package is experimental. It will be deprecated and removed when the [log]
package becomes stable. Its functionality will be migrated to
go.opentelemetry.io/otel.
*/
package global // import "go.opentelemetry.io/otel/log/global"
package global
import (
"go.opentelemetry.io/otel/log"
+1 -1
View File
@@ -3,7 +3,7 @@
// Package global is the internal implementation of the OpenTelemetry global
// Logs API.
package global // import "go.opentelemetry.io/otel/log/internal/global"
package global
import (
"context"
+1 -1
View File
@@ -1,7 +1,7 @@
// Copyright The OpenTelemetry Authors
// SPDX-License-Identifier: Apache-2.0
package global // import "go.opentelemetry.io/otel/log/internal/global"
package global
import (
"errors"
-453
View File
@@ -1,453 +0,0 @@
// Copyright The OpenTelemetry Authors
// SPDX-License-Identifier: Apache-2.0
//go:generate stringer -type=Kind -trimprefix=Kind
package log // import "go.opentelemetry.io/otel/log"
import (
"bytes"
"cmp"
"errors"
"fmt"
"math"
"slices"
"strconv"
"unsafe"
"go.opentelemetry.io/otel/attribute"
"go.opentelemetry.io/otel/internal/global"
)
// errKind is logged when a Value is decoded to an incompatible type.
var errKind = errors.New("invalid Kind")
// Kind is the kind of a [Value].
type Kind int
// Kind values.
const (
KindEmpty Kind = iota
KindBool
KindFloat64
KindInt64
KindString
KindBytes
KindSlice
KindMap
)
// A Value represents a structured log value.
// A zero value is valid and represents an empty value.
type Value struct {
// Ensure forward compatibility by explicitly making this not comparable.
noCmp [0]func() //nolint: unused // This is indeed used.
// num holds the value for Int64, Float64, and Bool. It holds the length
// for String, Bytes, Slice, Map.
num uint64
// any holds either the KindBool, KindInt64, KindFloat64, stringptr,
// bytesptr, sliceptr, or mapptr. If KindBool, KindInt64, or KindFloat64
// then the value of Value is in num as described above. Otherwise, it
// contains the value wrapped in the appropriate type.
any any
}
type (
// sliceptr represents a value in Value.any for KindString Values.
stringptr *byte
// bytesptr represents a value in Value.any for KindBytes Values.
bytesptr *byte
// sliceptr represents a value in Value.any for KindSlice Values.
sliceptr *Value
// mapptr represents a value in Value.any for KindMap Values.
mapptr *KeyValue
)
// StringValue returns a new [Value] for a string.
func StringValue(v string) Value {
return Value{
num: uint64(len(v)),
any: stringptr(unsafe.StringData(v)),
}
}
// IntValue returns a [Value] for an int.
func IntValue(v int) Value { return Int64Value(int64(v)) }
// Int64Value returns a [Value] for an int64.
func Int64Value(v int64) Value {
// This can be later converted back to int64 (overflow not checked).
return Value{num: uint64(v), any: KindInt64} // nolint:gosec
}
// Float64Value returns a [Value] for a float64.
func Float64Value(v float64) Value {
return Value{num: math.Float64bits(v), any: KindFloat64}
}
// BoolValue returns a [Value] for a bool.
func BoolValue(v bool) Value { //nolint:revive // Not a control flag.
var n uint64
if v {
n = 1
}
return Value{num: n, any: KindBool}
}
// BytesValue returns a [Value] for a byte slice. The passed slice must not be
// changed after it is passed.
func BytesValue(v []byte) Value {
return Value{
num: uint64(len(v)),
any: bytesptr(unsafe.SliceData(v)),
}
}
// SliceValue returns a [Value] for a slice of [Value]. The passed slice must
// not be changed after it is passed.
func SliceValue(vs ...Value) Value {
return Value{
num: uint64(len(vs)),
any: sliceptr(unsafe.SliceData(vs)),
}
}
// MapValue returns a new [Value] for a slice of key-value pairs. The passed
// slice must not be changed after it is passed.
func MapValue(kvs ...KeyValue) Value {
return Value{
num: uint64(len(kvs)),
any: mapptr(unsafe.SliceData(kvs)),
}
}
// AsString returns the value held by v as a string.
func (v Value) AsString() string {
if sp, ok := v.any.(stringptr); ok {
return unsafe.String(sp, v.num)
}
global.Error(errKind, "AsString", "Kind", v.Kind())
return ""
}
// asString returns the value held by v as a string. It will panic if the Value
// is not KindString.
func (v Value) asString() string {
return unsafe.String(v.any.(stringptr), v.num)
}
// AsInt64 returns the value held by v as an int64.
func (v Value) AsInt64() int64 {
if v.Kind() != KindInt64 {
global.Error(errKind, "AsInt64", "Kind", v.Kind())
return 0
}
return v.asInt64()
}
// asInt64 returns the value held by v as an int64. If v is not of KindInt64,
// this will return garbage.
func (v Value) asInt64() int64 {
// Assumes v.num was a valid int64 (overflow not checked).
return int64(v.num) // nolint: gosec
}
// AsBool returns the value held by v as a bool.
func (v Value) AsBool() bool {
if v.Kind() != KindBool {
global.Error(errKind, "AsBool", "Kind", v.Kind())
return false
}
return v.asBool()
}
// asBool returns the value held by v as a bool. If v is not of KindBool, this
// will return garbage.
func (v Value) asBool() bool { return v.num == 1 }
// AsFloat64 returns the value held by v as a float64.
func (v Value) AsFloat64() float64 {
if v.Kind() != KindFloat64 {
global.Error(errKind, "AsFloat64", "Kind", v.Kind())
return 0
}
return v.asFloat64()
}
// asFloat64 returns the value held by v as a float64. If v is not of
// KindFloat64, this will return garbage.
func (v Value) asFloat64() float64 { return math.Float64frombits(v.num) }
// AsBytes returns the value held by v as a []byte.
func (v Value) AsBytes() []byte {
if sp, ok := v.any.(bytesptr); ok {
return unsafe.Slice((*byte)(sp), v.num)
}
global.Error(errKind, "AsBytes", "Kind", v.Kind())
return nil
}
// asBytes returns the value held by v as a []byte. It will panic if the Value
// is not KindBytes.
func (v Value) asBytes() []byte {
return unsafe.Slice((*byte)(v.any.(bytesptr)), v.num)
}
// AsSlice returns the value held by v as a []Value.
func (v Value) AsSlice() []Value {
if sp, ok := v.any.(sliceptr); ok {
return unsafe.Slice((*Value)(sp), v.num)
}
global.Error(errKind, "AsSlice", "Kind", v.Kind())
return nil
}
// asSlice returns the value held by v as a []Value. It will panic if the Value
// is not KindSlice.
func (v Value) asSlice() []Value {
return unsafe.Slice((*Value)(v.any.(sliceptr)), v.num)
}
// AsMap returns the value held by v as a []KeyValue.
func (v Value) AsMap() []KeyValue {
if sp, ok := v.any.(mapptr); ok {
return unsafe.Slice((*KeyValue)(sp), v.num)
}
global.Error(errKind, "AsMap", "Kind", v.Kind())
return nil
}
// asMap returns the value held by v as a []KeyValue. It will panic if the
// Value is not KindMap.
func (v Value) asMap() []KeyValue {
return unsafe.Slice((*KeyValue)(v.any.(mapptr)), v.num)
}
// Kind returns the Kind of v.
func (v Value) Kind() Kind {
switch x := v.any.(type) {
case Kind:
return x
case stringptr:
return KindString
case bytesptr:
return KindBytes
case sliceptr:
return KindSlice
case mapptr:
return KindMap
default:
return KindEmpty
}
}
// Empty reports whether v does not hold any value.
func (v Value) Empty() bool { return v.Kind() == KindEmpty }
// Equal reports whether v is equal to w.
func (v Value) Equal(w Value) bool {
k1 := v.Kind()
k2 := w.Kind()
if k1 != k2 {
return false
}
switch k1 {
case KindInt64, KindBool:
return v.num == w.num
case KindString:
return v.asString() == w.asString()
case KindFloat64:
return v.asFloat64() == w.asFloat64()
case KindSlice:
return slices.EqualFunc(v.asSlice(), w.asSlice(), Value.Equal)
case KindMap:
sv := sortMap(v.asMap())
sw := sortMap(w.asMap())
return slices.EqualFunc(sv, sw, KeyValue.Equal)
case KindBytes:
return bytes.Equal(v.asBytes(), w.asBytes())
case KindEmpty:
return true
default:
global.Error(errKind, "Equal", "Kind", k1)
return false
}
}
func sortMap(m []KeyValue) []KeyValue {
sm := make([]KeyValue, len(m))
copy(sm, m)
slices.SortFunc(sm, func(a, b KeyValue) int {
return cmp.Compare(a.Key, b.Key)
})
return sm
}
// String returns Value's value as a string, formatted like [fmt.Sprint].
//
// The returned string is meant for debugging;
// the string representation is not stable.
func (v Value) String() string {
switch v.Kind() {
case KindString:
return v.asString()
case KindInt64:
// Assumes v.num was a valid int64 (overflow not checked).
return strconv.FormatInt(int64(v.num), 10) // nolint: gosec
case KindFloat64:
return strconv.FormatFloat(v.asFloat64(), 'g', -1, 64)
case KindBool:
return strconv.FormatBool(v.asBool())
case KindBytes:
return fmt.Sprint(v.asBytes()) // nolint:staticcheck // Use fmt.Sprint to encode as slice.
case KindMap:
return fmt.Sprint(v.asMap())
case KindSlice:
return fmt.Sprint(v.asSlice())
case KindEmpty:
return "<nil>"
default:
// Try to handle this as gracefully as possible.
//
// Don't panic here. The goal here is to have developers find this
// first if a slog.Kind is is not handled. It is
// preferable to have user's open issue asking why their attributes
// have a "unhandled: " prefix than say that their code is panicking.
return fmt.Sprintf("<unhandled log.Kind: %s>", v.Kind())
}
}
// A KeyValue is a key-value pair used to represent a log attribute (a
// superset of [go.opentelemetry.io/otel/attribute.KeyValue]) and map item.
type KeyValue struct {
Key string
Value Value
}
// Equal reports whether a is equal to b.
func (a KeyValue) Equal(b KeyValue) bool {
return a.Key == b.Key && a.Value.Equal(b.Value)
}
// String returns a KeyValue for a string value.
func String(key, value string) KeyValue {
return KeyValue{key, StringValue(value)}
}
// Int64 returns a KeyValue for an int64 value.
func Int64(key string, value int64) KeyValue {
return KeyValue{key, Int64Value(value)}
}
// Int returns a KeyValue for an int value.
func Int(key string, value int) KeyValue {
return KeyValue{key, IntValue(value)}
}
// Float64 returns a KeyValue for a float64 value.
func Float64(key string, value float64) KeyValue {
return KeyValue{key, Float64Value(value)}
}
// Bool returns a KeyValue for a bool value.
func Bool(key string, value bool) KeyValue {
return KeyValue{key, BoolValue(value)}
}
// Bytes returns a KeyValue for a []byte value.
// The passed slice must not be changed after it is passed.
func Bytes(key string, value []byte) KeyValue {
return KeyValue{key, BytesValue(value)}
}
// Slice returns a KeyValue for a []Value value.
// The passed slice must not be changed after it is passed.
func Slice(key string, value ...Value) KeyValue {
return KeyValue{key, SliceValue(value...)}
}
// Map returns a KeyValue for a map value.
// The passed slice must not be changed after it is passed.
func Map(key string, value ...KeyValue) KeyValue {
return KeyValue{key, MapValue(value...)}
}
// Empty returns a KeyValue with an empty value.
func Empty(key string) KeyValue {
return KeyValue{key, Value{}}
}
// String returns key-value pair as a string, formatted like "key:value".
//
// The returned string is meant for debugging;
// the string representation is not stable.
func (a KeyValue) String() string {
return fmt.Sprintf("%s:%s", a.Key, a.Value)
}
// ValueFromAttribute converts [attribute.Value] to [Value].
func ValueFromAttribute(value attribute.Value) Value {
switch value.Type() {
case attribute.EMPTY:
return Value{}
case attribute.BOOL:
return BoolValue(value.AsBool())
case attribute.BOOLSLICE:
val := value.AsBoolSlice()
res := make([]Value, 0, len(val))
for _, v := range val {
res = append(res, BoolValue(v))
}
return SliceValue(res...)
case attribute.INT64:
return Int64Value(value.AsInt64())
case attribute.INT64SLICE:
val := value.AsInt64Slice()
res := make([]Value, 0, len(val))
for _, v := range val {
res = append(res, Int64Value(v))
}
return SliceValue(res...)
case attribute.FLOAT64:
return Float64Value(value.AsFloat64())
case attribute.FLOAT64SLICE:
val := value.AsFloat64Slice()
res := make([]Value, 0, len(val))
for _, v := range val {
res = append(res, Float64Value(v))
}
return SliceValue(res...)
case attribute.STRING:
return StringValue(value.AsString())
case attribute.STRINGSLICE:
val := value.AsStringSlice()
res := make([]Value, 0, len(val))
for _, v := range val {
res = append(res, StringValue(v))
}
return SliceValue(res...)
case attribute.BYTESLICE:
val := value.AsByteSlice()
return BytesValue(val)
case attribute.SLICE:
val := value.AsSlice()
res := make([]Value, 0, len(val))
for _, v := range val {
res = append(res, ValueFromAttribute(v))
}
return SliceValue(res...)
}
// This code should never be reached
// as log attributes are a superset of standard attributes.
panic("unknown attribute type")
}
// KeyValueFromAttribute converts [attribute.KeyValue] to [KeyValue].
func KeyValueFromAttribute(kv attribute.KeyValue) KeyValue {
return KeyValue{
Key: string(kv.Key),
Value: ValueFromAttribute(kv.Value),
}
}
-31
View File
@@ -1,31 +0,0 @@
// Code generated by "stringer -type=Kind -trimprefix=Kind"; DO NOT EDIT.
package log
import "strconv"
func _() {
// An "invalid array index" compiler error signifies that the constant values have changed.
// Re-run the stringer command to generate them again.
var x [1]struct{}
_ = x[KindEmpty-0]
_ = x[KindBool-1]
_ = x[KindFloat64-2]
_ = x[KindInt64-3]
_ = x[KindString-4]
_ = x[KindBytes-5]
_ = x[KindSlice-6]
_ = x[KindMap-7]
}
const _Kind_name = "EmptyBoolFloat64Int64StringBytesSliceMap"
var _Kind_index = [...]uint8{0, 5, 9, 16, 21, 27, 32, 37, 40}
func (i Kind) String() string {
idx := int(i) - 0
if i < 0 || idx >= len(_Kind_index)-1 {
return "Kind(" + strconv.FormatInt(int64(i), 10) + ")"
}
return _Kind_name[_Kind_index[idx]:_Kind_index[idx+1]]
}
+3 -2
View File
@@ -1,7 +1,7 @@
// Copyright The OpenTelemetry Authors
// SPDX-License-Identifier: Apache-2.0
package log // import "go.opentelemetry.io/otel/log"
package log
import (
"context"
@@ -36,7 +36,8 @@ type Logger interface {
//
// This is useful for users that want to know if a [Record]
// will be processed or dropped before they perform complex operations to
// construct the [Record].
// construct the [Record]. Callers should invoke Enabled before each call
// to [Logger.Emit] because the enabled state may change over time.
//
// The passed param is likely to be a partial record information being
// provided (e.g a param with only the Severity set).
+1 -1
View File
@@ -12,7 +12,7 @@
// defaults to no operation for methods it does not implement.
//
// [OpenTelemetry Logs API]: https://pkg.go.dev/go.opentelemetry.io/otel/log
package noop // import "go.opentelemetry.io/otel/log/noop"
package noop
import (
"context"
+3 -2
View File
@@ -1,7 +1,7 @@
// Copyright The OpenTelemetry Authors
// SPDX-License-Identifier: Apache-2.0
package log // import "go.opentelemetry.io/otel/log"
package log
import "go.opentelemetry.io/otel/log/embedded"
@@ -24,7 +24,8 @@ type LoggerProvider interface {
// commonly, this means a bridge will need to accept this value from its
// users.
//
// If name is empty, implementations need to provide a default name.
// An empty name is invalid. Implementations should retain the empty value as the
// instrumentation scope name, return a working Logger, and report the invalid value.
//
// The version of the packages using a bridge can be critical information
// to include when logging. The bridge should accept this version
+12 -9
View File
@@ -1,11 +1,13 @@
// Copyright The OpenTelemetry Authors
// SPDX-License-Identifier: Apache-2.0
package log // import "go.opentelemetry.io/otel/log"
package log
import (
"slices"
"time"
"go.opentelemetry.io/otel/attribute"
)
// attributesInlineCount is the number of attributes that are efficiently
@@ -25,7 +27,7 @@ type Record struct {
observedTimestamp time.Time
severity Severity
severityText string
body Value
body attribute.Value
err error
// The fields below are for optimizing the implementation of Attributes and
@@ -35,7 +37,7 @@ type Record struct {
// Allocation optimization: an inline array sized to hold
// the majority of log calls (based on examination of open-source
// code). It holds the start of the list of attributes.
front [attributesInlineCount]KeyValue
front [attributesInlineCount]attribute.KeyValue
// The number of attributes in front.
nFront int
@@ -44,7 +46,7 @@ type Record struct {
// Invariants:
// - len(back) > 0 if nFront == len(front)
// - Unused array elements are zero-ed. Used to detect mistakes.
back []KeyValue
back []attribute.KeyValue
}
// EventName returns the event name.
@@ -55,6 +57,7 @@ func (r *Record) EventName() string {
// SetEventName sets the event name.
// A log record with non-empty event name is interpreted as an event record.
// Event names should uniquely identify the event's attribute and body structure.
func (r *Record) SetEventName(s string) {
r.eventName = s
}
@@ -102,12 +105,12 @@ func (r *Record) SetSeverityText(text string) {
}
// Body returns the body of the log record.
func (r *Record) Body() Value {
func (r *Record) Body() attribute.Value {
return r.body
}
// SetBody sets the body of the log record.
func (r *Record) SetBody(v Value) {
func (r *Record) SetBody(v attribute.Value) {
r.body = v
}
@@ -122,8 +125,8 @@ func (r *Record) SetErr(err error) {
}
// WalkAttributes walks all attributes the log record holds by calling f for
// each on each [KeyValue] in the [Record]. Iteration stops if f returns false.
func (r *Record) WalkAttributes(f func(KeyValue) bool) {
// each on each [attribute.KeyValue] in the [Record]. Iteration stops if f returns false.
func (r *Record) WalkAttributes(f func(attribute.KeyValue) bool) {
for i := 0; i < r.nFront; i++ {
if !f(r.front[i]) {
return
@@ -137,7 +140,7 @@ func (r *Record) WalkAttributes(f func(KeyValue) bool) {
}
// AddAttributes adds attributes to the log record.
func (r *Record) AddAttributes(attrs ...KeyValue) {
func (r *Record) AddAttributes(attrs ...attribute.KeyValue) {
var i int
for i = 0; i < len(attrs) && r.nFront < len(r.front); i++ {
a := attrs[i]
+1 -1
View File
@@ -3,7 +3,7 @@
//go:generate stringer -type=Severity -linecomment
package log // import "go.opentelemetry.io/otel/log"
package log
// Severity represents a log record severity (also known as log level). Smaller
// numerical values correspond to less severe log records (such as debug
+30 -2
View File
@@ -80,7 +80,25 @@ is implemented as `SimpleProcessor` struct in [simple.go](simple.go).
The [Batching processor](https://opentelemetry.io/docs/specs/otel/logs/sdk/#batching-processor)
is implemented as `BatchProcessor` struct in [batch.go](batch.go).
The `Batcher` can be also configured using the `OTEL_BLRP_*` environment variables as
`OnEmit` clones accepted records into a bounded, drop-oldest queue. Once the
queue reaches the batch size, it attempts a non-blocking, coalesced worker
notification. It may briefly contend on the queue lock but never waits for
exporter I/O. A single worker goroutine owns dequeueing, scheduled exports, and
all exporter lifecycle calls, so exporter backpressure blocks the worker
instead of causing repeated polling.
`ForceFlush` is serialized through the worker. When the worker accepts a
request, it drains the records then queued in batches no larger than the
configured maximum and calls the exporter's `ForceFlush` while the request
context remains valid. `Shutdown` first closes queue admission, drains accepted
records, calls `ForceFlush` while its context remains valid, and always invokes
the exporter's `Shutdown`. Cancellation stops additional drain chunks subject
to the exporter honoring its context.
`WithMaxQueueSize` bounds the pending-record queue; the worker additionally
retains bounded batch scratch space and records currently being exported.
The `BatchProcessor` can also be configured using the `OTEL_BLRP_*` environment variables as
[defined by the specification](https://opentelemetry.io/docs/specs/otel/configuration/sdk-environment-variables/#batch-logrecord-processor).
### Exporter
@@ -106,8 +124,18 @@ is defined as `Record` struct in [record.go](record.go).
The `Record` is designed similarly to [`log.Record`](https://pkg.go.dev/go.opentelemetry.io/otel/log#Record)
in order to reduce the number of heap allocations when processing attributes.
The log body and attributes use `attribute.Value` and `attribute.KeyValue`
from `go.opentelemetry.io/otel/attribute`, matching the API representation and
avoiding an SDK-specific value tree.
The SDK does not have have an additional definition of
Top-level duplicate attributes are handled with last-value-wins semantics
unless `WithAllowKeyDuplication` is configured. Nested `attribute.MAP` values
use the same semantics for duplicate map keys, including the body value, so
exporters receive a canonical attribute tree by default. The attribute value
length limit recursively truncates string, string slice, byte slice, slice, and
map attribute values. The body is deduplicated but not truncated.
The SDK does not have an additional definition of
[ReadableLogRecord](https://opentelemetry.io/docs/specs/otel/logs/sdk/#readablelogrecord)
as the specification does not say that the exporters must not be able to modify
the log records. It simply requires them to be able to read the log records.
+248 -183
View File
@@ -1,17 +1,20 @@
// Copyright The OpenTelemetry Authors
// SPDX-License-Identifier: Apache-2.0
package log // import "go.opentelemetry.io/otel/sdk/log"
package log
import (
"context"
"errors"
"slices"
"math"
"sync"
"sync/atomic"
"time"
"go.opentelemetry.io/otel"
"go.opentelemetry.io/otel/internal/global"
"go.opentelemetry.io/otel/sdk/log/internal/counter"
"go.opentelemetry.io/otel/sdk/log/internal/observ"
)
const (
@@ -19,7 +22,6 @@ const (
dfltExpInterval = time.Second
dfltExpTimeout = 30 * time.Second
dfltExpMaxBatchSize = 512
dfltExpBufferSize = 1
envarMaxQSize = "OTEL_BLRP_MAX_QUEUE_SIZE"
envarExpInterval = "OTEL_BLRP_SCHEDULE_DELAY"
@@ -35,82 +37,83 @@ var _ Processor = (*BatchProcessor)(nil)
// Use [NewBatchProcessor] to create a BatchProcessor. An empty BatchProcessor
// is shut down by default, no records will be batched or exported.
type BatchProcessor struct {
// The BatchProcessor is designed to provide the highest throughput of
// log records possible while being compatible with OpenTelemetry. The
// entry point of log records is the OnEmit method. This method is designed
// to receive records as fast as possible while still honoring shutdown
// commands. All records received are enqueued to queue.
//
// In order to block OnEmit as little as possible, a separate "poll"
// goroutine is spawned at the creation of a BatchProcessor. This
// goroutine is responsible for batching the queue at regular polled
// intervals, or when it is directly signaled to.
//
// To keep the polling goroutine from backing up, all batches it makes are
// exported with a bufferedExporter. This exporter allows the poll
// goroutine to enqueue an export payload that will be handled in a
// separate goroutine dedicated to the export. This asynchronous behavior
// allows the poll goroutine to maintain accurate interval polling.
//
// ___BatchProcessor____ __Poll Goroutine__ __Export Goroutine__
// || || || || || ||
// || ********** || || || || ********** ||
// || Records=>* OnEmit * || || | - ticker || || * export * ||
// || ********** || || | - trigger || || ********** ||
// || || || || | || || || ||
// || || || || | || || || ||
// || __________\/___ || || |*********** || || ______/\_______ ||
// || (____queue______)>=||=||===|* batch *===||=||=>[_export_buffer_] ||
// || || || |*********** || || ||
// ||_____________________|| ||__________________|| ||____________________||
//
//
// The "release valve" in this processing is the record queue. This queue
// is a ring buffer. It will overwrite the oldest records first when writes
// to OnEmit are made faster than the queue can be flushed. If batches
// cannot be flushed to the export buffer, the records will remain in the
// queue.
// exporter is the bufferedExporter all batches are exported with.
exporter *bufferExporter
// A single goroutine owns dequeueing and all exporter calls. OnEmit only
// writes to the bounded queue and signals that goroutine. Consequently,
// exporter backpressure blocks the exporter goroutine instead of causing
// another goroutine to retry without making progress.
exporter Exporter
// q is the active queue of records that have not yet been exported.
q *queue
// batchSize is the minimum number of records needed before an export is
// triggered (unless the interval expires).
// batchSize is the maximum number of records in a scheduled export.
batchSize int
// pollTrigger triggers the poll goroutine to flush a batch from the queue.
// This is sent to when it is known that the queue contains at least one
// complete batch.
//
// When a send is made to the channel, the poll loop will be reset after
// the flush. If there is still enough records in the queue for another
// batch the reset of the poll loop will automatically re-trigger itself.
// There is no need for the original sender to monitor and resend.
pollTrigger chan struct{}
// pollKill kills the poll goroutine. This is only expected to be closed
// once by the Shutdown method.
pollKill chan struct{}
// pollDone signals the poll goroutine has completed.
pollDone chan struct{}
// exportTrigger is a coalesced signal that records are ready to export.
exportTrigger chan struct{}
// flush serializes ForceFlush requests through the worker.
flush chan batchProcessorRequest
// shutdown accepts the single Shutdown request. It is separate from flush
// so shutdown cannot be blocked behind concurrent ForceFlush callers.
shutdown chan batchProcessorRequest
// done is closed by the exporter goroutine after exporter shutdown.
done chan struct{}
// stopped holds the stopped state of the BatchProcessor.
stopped atomic.Bool
// inst is the instrumentation for observability (nil when disabled).
inst *observ.BLP
noCmp [0]func() //nolint: unused // This is indeed used.
}
type batchProcessorRequest struct {
ctx context.Context
resp chan<- error
}
func (r batchProcessorRequest) respond(err error) {
r.resp <- err
}
// NewBatchProcessor decorates the provided exporter
// so that the log records are batched before exporting.
//
// All of the exporter's methods are called synchronously.
// Calls to the exporter's Export, ForceFlush, and Shutdown methods are
// synchronized and never invoked concurrently.
func NewBatchProcessor(exporter Exporter, opts ...BatchProcessorOption) *BatchProcessor {
cfg := newBatchConfig(opts)
if exporter == nil {
// Do not panic on nil export.
exporter = defaultNoopExporter
}
b := &BatchProcessor{
q: newQueue(cfg.maxQSize.Value),
batchSize: cfg.expMaxBatchSize.Value,
exportTrigger: make(chan struct{}, 1),
flush: make(chan batchProcessorRequest),
shutdown: make(chan batchProcessorRequest, 1),
done: make(chan struct{}),
}
var err error
b.inst, err = observ.NewBLP(
counter.NextExporterID(),
func() int64 { return int64(b.q.Len()) },
int64(cfg.maxQSize.Value),
)
if err != nil {
otel.Handle(err)
}
// Wrap exporter with metrics recording if observability is enabled.
// This must be the innermost wrapper (closest to user exporter) to record
// metrics just before calling the actual exporter.
if b.inst != nil {
exporter = newMetricsExporter(exporter, b.inst)
}
// Order is important here. Wrap the timeoutExporter with the chunkExporter
// to ensure each export completes in timeout (instead of all chunked
// exports).
@@ -119,69 +122,144 @@ func NewBatchProcessor(exporter Exporter, opts ...BatchProcessorOption) *BatchPr
// appropriately on export.
exporter = newChunkExporter(exporter, cfg.expMaxBatchSize.Value)
b := &BatchProcessor{
exporter: newBufferExporter(exporter, cfg.expBufferSize.Value),
q: newQueue(cfg.maxQSize.Value),
batchSize: cfg.expMaxBatchSize.Value,
pollTrigger: make(chan struct{}, 1),
pollKill: make(chan struct{}),
}
b.pollDone = b.poll(cfg.expInterval.Value)
b.exporter = exporter
b.process(cfg.expInterval.Value)
return b
}
// poll spawns a goroutine to handle interval polling and batch exporting. The
// returned done chan is closed when the spawned goroutine completes.
func (b *BatchProcessor) poll(interval time.Duration) (done chan struct{}) {
done = make(chan struct{})
ticker := time.NewTicker(interval)
// TODO: investigate using a sync.Pool instead of cloning.
buf := make([]Record, b.batchSize)
// process starts the goroutine that owns dequeueing and all exporter calls.
func (b *BatchProcessor) process(interval time.Duration) {
go func() {
defer close(done)
defer ticker.Stop()
timer := time.NewTimer(interval)
defer timer.Stop()
// The worker owns and reuses buf. Exporters must not retain the slice
// passed to them, so it is safe to refill after Export returns.
buf := make([]Record, b.batchSize)
for {
// Probe shutdown by itself first. This makes an already queued terminal
// request win over every other ready case. Closing done before replying
// also means a successful Shutdown response observes a stopped worker.
select {
case <-ticker.C:
case <-b.pollTrigger:
ticker.Reset(interval)
case <-b.pollKill:
case req := <-b.shutdown:
err := b.shutdownExporter(req.ctx)
close(b.done)
req.respond(err)
return
default:
}
if d := b.q.Dropped(); d > 0 {
global.Warn("dropped log records", "dropped", d)
// With no queued shutdown, service a waiting ForceFlush before ordinary
// export wakes. The default keeps this priority check non-blocking.
// Shutdown remains selectable in case it arrived after the first probe.
select {
case req := <-b.shutdown:
err := b.shutdownExporter(req.ctx)
close(b.done)
req.respond(err)
return
case req := <-b.flush:
err := b.flushExporter(req.ctx)
req.respond(err)
continue
default:
}
var qLen int
// Don't copy data from queue unless exporter can accept more, it is very expensive.
if b.exporter.Ready() {
qLen = b.q.TryDequeue(buf, func(r []Record) bool {
ok := b.exporter.EnqueueExport(r)
if ok {
buf = slices.Clone(buf)
}
return ok
})
} else {
qLen = b.q.Len()
}
if qLen >= b.batchSize {
// There is another full batch ready. Immediately trigger
// another export attempt.
select {
case b.pollTrigger <- struct{}{}:
default:
// Another flush signal already received.
}
// No lifecycle request was waiting, so block on the complete event set.
// Both timer and size-triggered exports start a new interval window.
select {
case req := <-b.shutdown:
err := b.shutdownExporter(req.ctx)
close(b.done)
req.respond(err)
return
case req := <-b.flush:
err := b.flushExporter(req.ctx)
req.respond(err)
case <-timer.C:
resetTimer(timer, interval)
b.exportBatch(buf)
case <-b.exportTrigger:
resetTimer(timer, interval)
b.exportBatch(buf)
}
}
}()
return done
}
func resetTimer(timer *time.Timer, interval time.Duration) {
// Handle both GODEBUG=asynctimerchan=[0|1] properly.
if !timer.Stop() {
select {
case <-timer.C:
default:
}
}
timer.Reset(interval)
}
func (b *BatchProcessor) exportBatch(buf []Record) {
b.logDroppedRecords()
n, remaining := b.q.Dequeue(buf)
if n == 0 {
return
}
err := b.exporter.Export(context.Background(), buf[:n])
clear(buf[:n])
if err != nil {
otel.Handle(err)
}
if remaining >= b.batchSize {
b.triggerExport()
}
}
func (b *BatchProcessor) flushExporter(ctx context.Context) error {
if err := ctx.Err(); err != nil {
return err
}
b.logDroppedRecords()
records := b.q.Flush()
err := b.exporter.Export(ctx, records)
clear(records)
if ctxErr := ctx.Err(); ctxErr != nil {
return errors.Join(err, ctxErr)
}
return errors.Join(err, b.exporter.ForceFlush(ctx))
}
func (b *BatchProcessor) shutdownExporter(ctx context.Context) error {
b.logDroppedRecords()
records := b.q.Flush()
err := b.exporter.Export(ctx, records)
clear(records)
if ctxErr := ctx.Err(); ctxErr != nil {
err = errors.Join(err, ctxErr)
} else {
err = errors.Join(err, b.exporter.ForceFlush(ctx))
}
err = errors.Join(err, b.exporter.Shutdown(ctx))
if b.inst != nil {
err = errors.Join(err, b.inst.Shutdown())
}
return err
}
func (b *BatchProcessor) logDroppedRecords() {
if d := b.q.Dropped(); d > 0 {
if b.inst != nil {
b.inst.ProcessedQueueFull(context.Background(), int64(min(math.MaxInt64, d))) // nolint:gosec
}
global.Warn("dropped log records", "dropped", d)
}
}
func (b *BatchProcessor) triggerExport() {
select {
case b.exportTrigger <- struct{}{}:
default:
}
}
// Enabled returns true, indicating this Processor will process all records.
@@ -196,43 +274,31 @@ func (b *BatchProcessor) OnEmit(_ context.Context, r *Record) error {
}
// The record is cloned so that changes done by subsequent processors
// are not going to lead to a data race.
if n := b.q.Enqueue(r.Clone()); n >= b.batchSize {
select {
case b.pollTrigger <- struct{}{}:
default:
// Flush chan full. The poll goroutine will handle this by
// re-sending any trigger until the queue has less than batchSize
// records.
}
if n, accepted := b.q.Enqueue(r.Clone()); accepted && n >= b.batchSize {
b.triggerExport()
}
return nil
}
// Shutdown flushes queued log records and shuts down the decorated exporter.
// Shutdown flushes queued log records and the decorated exporter before
// shutting it down.
func (b *BatchProcessor) Shutdown(ctx context.Context) error {
if b.stopped.Swap(true) || b.q == nil {
return nil
}
// Stop the poll goroutine.
close(b.pollKill)
b.q.Close()
resp := make(chan error, 1)
b.shutdown <- batchProcessorRequest{ctx: ctx, resp: resp}
if err := ctx.Err(); err != nil {
return err
}
select {
case <-b.pollDone:
case err := <-resp:
return err
case <-ctx.Done():
// Out of time.
return errors.Join(ctx.Err(), b.exporter.Shutdown(ctx))
return ctx.Err()
}
// Flush remaining queued before exporter shutdown.
err := b.exporter.Export(ctx, b.q.Flush())
return errors.Join(err, b.exporter.Shutdown(ctx))
}
var errPartialFlush = errors.New("partial flush: export buffer full")
// Used for testing.
var ctxErr = func(ctx context.Context) error {
return ctx.Err()
}
// ForceFlush flushes queued log records and flushes the decorated exporter.
@@ -240,27 +306,26 @@ func (b *BatchProcessor) ForceFlush(ctx context.Context) error {
if b.stopped.Load() || b.q == nil {
return nil
}
if err := ctx.Err(); err != nil {
return err
}
buf := make([]Record, b.q.cap)
notFlushed := func() bool {
var flushed bool
_ = b.q.TryDequeue(buf, func(r []Record) bool {
flushed = b.exporter.EnqueueExport(r)
return flushed
})
return !flushed
resp := make(chan error, 1)
req := batchProcessorRequest{ctx: ctx, resp: resp}
select {
case b.flush <- req:
case <-b.done:
return nil
case <-ctx.Done():
return ctx.Err()
}
var err error
// For as long as ctx allows, try to make a single flush of the queue.
for notFlushed() {
// Use ctxErr instead of calling ctx.Err directly so we can test
// the partial error return.
if e := ctxErr(ctx); e != nil {
err = errors.Join(e, errPartialFlush)
break
}
select {
case err := <-resp:
return err
case <-ctx.Done():
return ctx.Err()
}
return errors.Join(err, b.exporter.ForceFlush(ctx))
}
// queue holds a queue of logging records.
@@ -273,6 +338,7 @@ type queue struct {
dropped atomic.Uint64
cap, len int
read, write *ring
closed bool
}
func newQueue(size int) *queue {
@@ -302,10 +368,14 @@ func (q *queue) Dropped() uint64 {
//
// If enqueueing r will exceed the capacity of q, the oldest Record held in q
// will be dropped and r retained.
func (q *queue) Enqueue(r Record) int {
func (q *queue) Enqueue(r Record) (int, bool) {
q.Lock()
defer q.Unlock()
if q.closed {
return q.len, false
}
q.write.Value = r
q.write = q.write.Next()
@@ -316,35 +386,23 @@ func (q *queue) Enqueue(r Record) int {
q.read = q.read.Next()
q.dropped.Add(1)
}
return q.len
return q.len, true
}
// TryDequeue attempts to dequeue up to len(buf) Records. The available Records
// will be assigned into buf and passed to write. If write fails, returning
// false, the Records will not be removed from the queue. If write succeeds,
// returning true, the dequeued Records are removed from the queue. The number
// of Records remaining in the queue are returned.
//
// When write is called the lock of q is held. The write function must not call
// other methods of this q that acquire the lock.
func (q *queue) TryDequeue(buf []Record, write func([]Record) bool) int {
// Dequeue removes up to len(buf) records from the queue and copies them into
// buf. The number copied and the number remaining are returned.
func (q *queue) Dequeue(buf []Record) (int, int) {
q.Lock()
defer q.Unlock()
origRead := q.read
n := min(len(buf), q.len)
for i := range n {
buf[i] = q.read.Value // nolint:gosec // n is bounded by len(buf)
q.read.Value = Record{}
q.read = q.read.Next()
}
if write(buf[:n]) {
q.len -= n
} else {
q.read = origRead
}
return q.len
q.len -= n
return n, q.len
}
// Flush returns all the Records held in the queue and resets it to be
@@ -353,9 +411,22 @@ func (q *queue) Flush() []Record {
q.Lock()
defer q.Unlock()
return q.flush()
}
// Close stops the queue from accepting records.
func (q *queue) Close() {
q.Lock()
defer q.Unlock()
q.closed = true
}
func (q *queue) flush() []Record {
out := make([]Record, q.len)
for i := range out {
out[i] = q.read.Value
q.read.Value = Record{}
q.read = q.read.Next()
}
q.len = 0
@@ -368,7 +439,6 @@ type batchConfig struct {
expInterval setting[time.Duration]
expTimeout setting[time.Duration]
expMaxBatchSize setting[int]
expBufferSize setting[int]
}
func newBatchConfig(options []BatchProcessorOption) batchConfig {
@@ -399,14 +469,9 @@ func newBatchConfig(options []BatchProcessorOption) batchConfig {
clearLessThanOne[int](),
getenv[int](envarExpMaxBatchSize),
clearLessThanOne[int](), // nolint:gocritic // the function argument is duplicated on purpose
clampMax[int](c.maxQSize.Value),
fallback[int](dfltExpMaxBatchSize),
clampMax[int](c.maxQSize.Value),
)
c.expBufferSize = c.expBufferSize.Resolve(
clearLessThanOne[int](),
fallback[int](dfltExpBufferSize),
)
return c
}
@@ -474,8 +539,9 @@ func WithExportTimeout(d time.Duration) BatchProcessorOption {
// and this option is not passed, that variable value will be used.
//
// By default, if an environment variable is not set, and this option is not
// passed, 512 will be used.
// passed, 512 or the maximum queue size, if smaller, will be used.
// The default value is also used when the provided value is less than one.
// The effective batch size will not exceed the configured maximum queue size.
func WithExportMaxBatchSize(size int) BatchProcessorOption {
return batchOptionFunc(func(cfg batchConfig) batchConfig {
cfg.expMaxBatchSize = newSetting(size)
@@ -483,14 +549,13 @@ func WithExportMaxBatchSize(size int) BatchProcessorOption {
})
}
// WithExportBufferSize sets the batch buffer size.
// Batches will be temporarily kept in a memory buffer until they are exported.
// WithExportBufferSize is retained for source compatibility and has no effect.
// The processor no longer maintains a separately configurable export-request
// buffer. [WithMaxQueueSize] bounds the pending-record queue.
//
// By default, a value of 1 will be used.
// The default value is also used when the provided value is less than one.
func WithExportBufferSize(size int) BatchProcessorOption {
// Deprecated: This option is no longer used.
func WithExportBufferSize(_ int) BatchProcessorOption {
return batchOptionFunc(func(cfg batchConfig) batchConfig {
cfg.expBufferSize = newSetting(size)
return cfg
})
}
+1 -1
View File
@@ -36,4 +36,4 @@ the experimental features.
See [go.opentelemetry.io/otel/log] for more information about
the OpenTelemetry Logs API.
*/
package log // import "go.opentelemetry.io/otel/sdk/log"
package log
+32 -197
View File
@@ -1,17 +1,14 @@
// Copyright The OpenTelemetry Authors
// SPDX-License-Identifier: Apache-2.0
package log // import "go.opentelemetry.io/otel/sdk/log"
package log
import (
"context"
"errors"
"fmt"
"sync"
"sync/atomic"
"time"
"go.opentelemetry.io/otel"
"go.opentelemetry.io/otel/sdk/log/internal/observ"
)
// Exporter handles the delivery of log records to external receivers.
@@ -67,6 +64,11 @@ func (noopExporter) Shutdown(context.Context) error { return nil }
func (noopExporter) ForceFlush(context.Context) error { return nil }
func shutdownExporter(ctx context.Context, exporter Exporter) error {
err := exporter.ForceFlush(ctx)
return errors.Join(err, exporter.Shutdown(ctx))
}
// chunkExporter wraps an Exporter's Export method so it is called with
// appropriately sized export payloads. Any payload larger than a defined size
// is chunked into smaller payloads and exported sequentially.
@@ -90,12 +92,19 @@ func newChunkExporter(exporter Exporter, size int) Exporter {
// Export exports records in chunks no larger than c.size.
func (c chunkExporter) Export(ctx context.Context, records []Record) error {
n := len(records)
var errs []error
for i, j := 0, min(c.size, n); i < n; i, j = i+c.size, min(j+c.size, n) {
if ctxErr := ctx.Err(); ctxErr != nil {
return errors.Join(append(errs, ctxErr)...)
}
if err := c.Exporter.Export(ctx, records[i:j]); err != nil {
return err
errs = append(errs, err)
}
if ctxErr := ctx.Err(); ctxErr != nil {
return errors.Join(append(errs, ctxErr)...)
}
}
return nil
return errors.Join(errs...)
}
// timeoutExporter wraps an Exporter and ensures any call to Export will have a
@@ -126,201 +135,27 @@ func (e *timeoutExporter) Export(ctx context.Context, records []Record) error {
return e.Exporter.Export(ctx, records)
}
// exportSync exports all data from input using exporter in a spawned
// goroutine. The returned chan will be closed when the spawned goroutine
// completes.
func exportSync(input <-chan exportData, exporter Exporter) (done chan struct{}) {
done = make(chan struct{})
go func() {
defer close(done)
for data := range input {
data.DoExport(exporter.Export)
}
}()
return done
}
// exportData is data related to an export.
type exportData struct {
ctx context.Context
records []Record
// respCh is the channel any error returned from the export will be sent
// on. If this is nil, and the export error is non-nil, the error will
// passed to the OTel error handler.
respCh chan<- error
}
// DoExport calls exportFn with the data contained in e. The error response
// will be returned on e's respCh if not nil. The error will be handled by the
// default OTel error handle if it is not nil and respCh is nil or full.
func (e exportData) DoExport(exportFn func(context.Context, []Record) error) {
if len(e.records) == 0 {
e.respond(nil)
return
}
e.respond(exportFn(e.ctx, e.records))
}
func (e exportData) respond(err error) {
select {
case e.respCh <- err:
default:
// e.respCh is nil or busy, default to otel.Handler.
if err != nil {
otel.Handle(err)
}
}
}
// bufferExporter provides asynchronous and synchronous export functionality by
// buffering export requests.
type bufferExporter struct {
// metricsExporter wraps an Exporter to record log processing metrics
// just before calling the wrapped exporter.
type metricsExporter struct {
Exporter
input chan exportData
inputMu sync.Mutex
done chan struct{}
stopped atomic.Bool
inst *observ.BLP
}
// newBufferExporter returns a new bufferExporter that wraps exporter. The
// returned bufferExporter will buffer at most size number of export requests.
// If size is less than 1, 1 will be used.
func newBufferExporter(exporter Exporter, size int) *bufferExporter {
if size < 1 {
size = 1
}
input := make(chan exportData, size)
return &bufferExporter{
// newMetricsExporter creates a metricsExporter that wraps the given exporter.
func newMetricsExporter(exporter Exporter, inst *observ.BLP) Exporter {
return &metricsExporter{
Exporter: exporter,
input: input,
done: exportSync(input, exporter),
inst: inst,
}
}
func (e *bufferExporter) Ready() bool {
return len(e.input) != cap(e.input)
}
var errStopped = errors.New("exporter stopped")
func (e *bufferExporter) enqueue(ctx context.Context, records []Record, rCh chan<- error) error {
data := exportData{ctx, records, rCh}
e.inputMu.Lock()
defer e.inputMu.Unlock()
// Check stopped before enqueueing now that e.inputMu is held. This
// prevents sends on a closed chan when Shutdown is called concurrently.
if e.stopped.Load() {
return errStopped
}
select {
case e.input <- data:
case <-ctx.Done():
return ctx.Err()
}
return nil
}
// EnqueueExport enqueues an export of records in the context of ctx to be
// performed asynchronously. This will return true if the records are
// successfully enqueued (or the bufferExporter is shut down), false otherwise.
//
// The passed records are held after this call returns.
func (e *bufferExporter) EnqueueExport(records []Record) bool {
if len(records) == 0 {
// Nothing to enqueue, do not waste input space.
return true
}
data := exportData{ctx: context.Background(), records: records}
e.inputMu.Lock()
defer e.inputMu.Unlock()
// Check stopped before enqueueing now that e.inputMu is held. This
// prevents sends on a closed chan when Shutdown is called concurrently.
if e.stopped.Load() {
return true
}
select {
case e.input <- data:
return true
default:
return false
}
}
// Export synchronously exports records in the context of ctx. This will not
// return until the export has been completed.
func (e *bufferExporter) Export(ctx context.Context, records []Record) error {
if len(records) == 0 {
return nil
}
resp := make(chan error, 1)
err := e.enqueue(ctx, records, resp)
if err != nil {
if errors.Is(err, errStopped) {
return nil
}
return fmt.Errorf("%w: dropping %d records", err, len(records))
}
select {
case err := <-resp:
return err
case <-ctx.Done():
return ctx.Err()
}
}
// ForceFlush flushes buffered exports. Any existing exports that is buffered
// is flushed before this returns.
func (e *bufferExporter) ForceFlush(ctx context.Context) error {
resp := make(chan error, 1)
err := e.enqueue(ctx, nil, resp)
if err != nil {
if errors.Is(err, errStopped) {
return nil
}
return err
}
select {
case <-resp:
case <-ctx.Done():
return ctx.Err()
}
return e.Exporter.ForceFlush(ctx)
}
// Shutdown shuts down e.
//
// Any buffered exports are flushed before this returns.
//
// All calls to EnqueueExport or Exporter will return nil without any export
// after this is called.
func (e *bufferExporter) Shutdown(ctx context.Context) error {
if e.stopped.Swap(true) {
return nil
}
e.inputMu.Lock()
defer e.inputMu.Unlock()
// No more sends will be made.
close(e.input)
select {
case <-e.done:
case <-ctx.Done():
return errors.Join(ctx.Err(), e.Exporter.Shutdown(ctx))
}
return e.Exporter.Shutdown(ctx)
// Export records the number of log records as a metric then forwards
// them to the wrapped Exporter. Error returned from wrapped exporter
// is not considered as per specification (to be measured by exporter).
func (e *metricsExporter) Export(ctx context.Context, records []Record) error {
if e.inst != nil {
e.inst.Processed(ctx, int64(len(records)))
}
return e.Exporter.Export(ctx, records)
}
+3 -3
View File
@@ -1,7 +1,7 @@
// Copyright The OpenTelemetry Authors
// SPDX-License-Identifier: Apache-2.0
package log // import "go.opentelemetry.io/otel/sdk/log"
package log
import (
"context"
@@ -11,8 +11,8 @@ import (
"go.opentelemetry.io/otel/metric"
"go.opentelemetry.io/otel/sdk"
"go.opentelemetry.io/otel/sdk/log/internal/x"
semconv "go.opentelemetry.io/otel/semconv/v1.41.0"
"go.opentelemetry.io/otel/semconv/v1.41.0/otelconv"
semconv "go.opentelemetry.io/otel/semconv/v1.43.0"
"go.opentelemetry.io/otel/semconv/v1.43.0/otelconv"
)
// newRecordCounterIncr returns a function that increments the log record
+300
View File
@@ -0,0 +1,300 @@
// Copyright The OpenTelemetry Authors
// SPDX-License-Identifier: Apache-2.0
// DO NOT MODIFY. Generated by gotmpl.
// source: internal/shared/attrnorm/dedup.go.tmpl
// Package attrnorm normalizes attribute values.
package attrnorm
import (
"reflect"
"unsafe"
"go.opentelemetry.io/otel/attribute"
)
var (
keyValueType = reflect.TypeFor[attribute.KeyValue]()
valueType = reflect.TypeFor[attribute.Value]()
)
// rawValue mirrors attribute.Value. It is used only to read immutable slice
// storage without calling AsMap or AsSlice on no-op paths.
type rawValue struct {
vtype attribute.Type
numeric uint64
stringly string
slice any
}
// Value returns value with all map values deduplicated and whether it changed.
//
// Duplicate map keys are resolved using last-value-wins semantics.
func Value(value attribute.Value) (attribute.Value, bool) {
switch value.Type() {
case attribute.SLICE:
return deduplicateSliceValue(value)
case attribute.MAP:
return deduplicateMapValue(value)
default:
return value, false
}
}
// KeyValue returns kv with all map values deduplicated and whether it changed.
func KeyValue(kv attribute.KeyValue) (attribute.KeyValue, bool) {
value, changed := Value(kv.Value)
if changed {
kv.Value = value
}
return kv, changed
}
// KeyValues returns kvs with all map values deduplicated and whether they changed.
//
// The returned slice is the original kvs slice if no value needs
// deduplication. Top-level keys in kvs are not deduplicated.
func KeyValues(kvs []attribute.KeyValue) ([]attribute.KeyValue, bool) {
// Preserve the caller's slice on the common no-op path. Once a changed
// value is found, copy the prior values exactly once and fill the rest in
// place as the scan continues.
var normalized []attribute.KeyValue
for i, kv := range kvs {
kv, changed := KeyValue(kv)
if normalized != nil {
normalized[i] = kv
continue
}
if !changed {
continue
}
normalized = make([]attribute.KeyValue, len(kvs))
copy(normalized, kvs[:i])
normalized[i] = kv
}
if normalized == nil {
return kvs, false
}
return normalized, true
}
// Set returns set with all map values deduplicated and whether it changed.
//
// The returned Set is the original set if no value needs deduplication.
// Top-level key uniqueness remains attribute.Set's responsibility; this only
// normalizes map attribute values.
func Set(set attribute.Set) (attribute.Set, bool) {
if set.Len() == 0 {
return set, false
}
// Most attribute sets contain no duplicate map keys. Delay allocation until
// the first changed value so the no-op path returns the original Set.
var normalized []attribute.KeyValue
for i := range set.Len() {
kv, _ := set.Get(i)
kv, changed := KeyValue(kv)
if normalized != nil {
normalized = append(normalized, kv)
continue
}
if !changed {
continue
}
normalized = make([]attribute.KeyValue, 0, set.Len())
for j := range i {
prior, _ := set.Get(j)
normalized = append(normalized, prior)
}
normalized = append(normalized, kv)
}
if normalized == nil {
return set, false
}
return attribute.NewSet(normalized...), true
}
func deduplicateSliceValue(value attribute.Value) (attribute.Value, bool) {
storage := valueStorage(value)
length := valueLen(storage)
// Slice values can contain map values, so recurse into each element while
// keeping the original attribute.Value when no element changes.
var normalized []attribute.Value
for i := range length {
elem := valueAt(storage, i)
elem, changed := Value(elem)
if normalized != nil {
normalized[i] = elem
continue
}
if !changed {
continue
}
normalized = make([]attribute.Value, length)
for j := range i {
normalized[j] = valueAt(storage, j)
}
normalized[i] = elem
}
if normalized == nil {
return value, false
}
return attribute.SliceValue(normalized...), true
}
func deduplicateMapValue(value attribute.Value) (attribute.Value, bool) {
storage := valueStorage(value)
length := keyValueLen(storage)
if length <= 1 {
// A single map entry cannot duplicate its own key, but its value might
// contain a map or slice that needs recursive normalization.
if length == 1 {
kv, changed := KeyValue(keyValueAt(storage, 0))
if changed {
return attribute.MapValue(kv), true
}
}
return value, false
}
var normalized []attribute.KeyValue
for i := 0; i < length; {
// attribute.MapValue stores key-values sorted by key using a stable
// sort. Equal keys therefore form a contiguous run, and the last
// element in that run is the last value provided by the caller.
first := keyValueAt(storage, i)
j := i + 1
for j < length && keyValueAt(storage, j).Key == first.Key {
j++
}
kv, nestedChanged := KeyValue(keyValueAt(storage, j-1))
// j-i > 1 means the current key run contained duplicates.
changed := nestedChanged || j-i > 1
if normalized != nil {
normalized = append(normalized, kv)
} else if changed {
normalized = make([]attribute.KeyValue, 0, length)
for k := range i {
normalized = append(normalized, keyValueAt(storage, k))
}
normalized = append(normalized, kv)
}
i = j
}
if normalized == nil {
return value, false
}
return attribute.MapValue(normalized...), true
}
func valueStorage(value attribute.Value) any {
// attribute.Value does not expose allocation-free map/slice iteration.
// The raw mirror lets us read the immutable backing array directly and
// reserve AsMap/AsSlice-style allocation for paths that actually change.
return (*rawValue)(
unsafe.Pointer(&value),
).slice //nolint:gosec // Read-only mirror of attribute.Value for allocation-free iteration.
}
func valueLen(storage any) int {
// attribute.Value stores small slices in fixed-size array values. Handle
// the common sizes directly and fall back to reflection for larger arrays.
switch storage.(type) {
case [0]attribute.Value:
return 0
case [1]attribute.Value:
return 1
case [2]attribute.Value:
return 2
case [3]attribute.Value:
return 3
case [4]attribute.Value:
return 4
case [5]attribute.Value:
return 5
default:
return arrayLen(storage, valueType)
}
}
func valueAt(storage any, i int) attribute.Value {
switch values := storage.(type) {
case [1]attribute.Value:
return values[i]
case [2]attribute.Value:
return values[i]
case [3]attribute.Value:
return values[i]
case [4]attribute.Value:
return values[i]
case [5]attribute.Value:
return values[i]
default:
return arrayAt[attribute.Value](storage, valueType, i)
}
}
func keyValueLen(storage any) int {
// attribute.Value stores small maps in fixed-size key-value arrays. Handle
// the common sizes directly and fall back to reflection for larger arrays.
switch storage.(type) {
case [0]attribute.KeyValue:
return 0
case [1]attribute.KeyValue:
return 1
case [2]attribute.KeyValue:
return 2
case [3]attribute.KeyValue:
return 3
case [4]attribute.KeyValue:
return 4
case [5]attribute.KeyValue:
return 5
default:
return arrayLen(storage, keyValueType)
}
}
func keyValueAt(storage any, i int) attribute.KeyValue {
switch kvs := storage.(type) {
case [1]attribute.KeyValue:
return kvs[i]
case [2]attribute.KeyValue:
return kvs[i]
case [3]attribute.KeyValue:
return kvs[i]
case [4]attribute.KeyValue:
return kvs[i]
case [5]attribute.KeyValue:
return kvs[i]
default:
return arrayAt[attribute.KeyValue](storage, keyValueType, i)
}
}
func arrayLen(storage any, elem reflect.Type) int {
// Be defensive around invalid or unexpected Value storage. Returning zero
// makes malformed storage a no-op instead of panicking in telemetry paths.
array := reflect.ValueOf(storage)
if array.Kind() != reflect.Array || array.Type().Elem() != elem {
return 0
}
return array.Len()
}
func arrayAt[T any](storage any, elem reflect.Type, i int) T {
// Match arrayLen's fail-closed behavior for unexpected storage.
array := reflect.ValueOf(storage)
if array.Kind() != reflect.Array || array.Type().Elem() != elem || i < 0 || i >= array.Len() {
var zero T
return zero
}
return array.Index(i).Interface().(T)
}
+232
View File
@@ -0,0 +1,232 @@
// Copyright The OpenTelemetry Authors
// SPDX-License-Identifier: Apache-2.0
// DO NOT MODIFY. Generated by gotmpl.
// source: internal/shared/attrnorm/truncate.go.tmpl
package attrnorm
import (
"slices"
"strings"
"unicode/utf8"
"go.opentelemetry.io/otel/attribute"
)
// Truncate returns a truncated version of attr. Only string, string slice,
// byte slice, slice, and map attribute values are truncated. String values are
// truncated to at most a length of limit. Each string slice value is truncated
// in this fashion (the slice length itself is unaffected), and byte slice
// values are truncated to at most limit bytes. For slice and map attribute
// values, the limit is applied recursively to contained values.
//
// No truncation is performed for a negative limit.
func Truncate(limit int, attr attribute.KeyValue) attribute.KeyValue {
if limit < 0 {
return attr
}
switch attr.Value.Type() {
case attribute.STRING:
v := attr.Value.AsString()
return attr.Key.String(truncate(limit, v))
case attribute.STRINGSLICE:
v := attr.Value.AsStringSlice()
for i := range v {
v[i] = truncate(limit, v[i])
}
return attr.Key.StringSlice(v)
case attribute.BYTESLICE:
v := attr.Value.AsString()
if len(v) > limit {
return attr.Key.ByteSlice([]byte(v[:limit]))
}
return attr
case attribute.SLICE:
v := attr.Value.AsSlice()
if !slices.ContainsFunc(v, func(e attribute.Value) bool { return needsTruncation(limit, e) }) {
return attr
}
newV := make([]attribute.Value, len(v))
for i, elem := range v {
newV[i] = TruncateValue(limit, elem)
}
return attr.Key.Slice(newV...)
case attribute.MAP:
v := attr.Value.AsMap()
if !slices.ContainsFunc(v, func(kv attribute.KeyValue) bool { return needsTruncation(limit, kv.Value) }) {
return attr
}
newV := make([]attribute.KeyValue, len(v))
for i, elem := range v {
elem.Value = TruncateValue(limit, elem.Value)
newV[i] = elem
}
return attr.Key.Map(newV...)
}
return attr
}
// TruncateValue returns a truncated version of v. Only string, string
// slice, byte slice, and (recursively) slice and map values are modified.
//
// No truncation is performed for a negative limit.
func TruncateValue(limit int, v attribute.Value) attribute.Value {
if limit < 0 {
return v
}
switch v.Type() {
case attribute.STRING:
return attribute.StringValue(truncate(limit, v.AsString()))
case attribute.STRINGSLICE:
ss := v.AsStringSlice()
for i := range ss {
ss[i] = truncate(limit, ss[i])
}
return attribute.StringSliceValue(ss)
case attribute.BYTESLICE:
// len(v.AsString()) is identical to len(v.AsByteSlice()) but
// avoids allocating the full slice before truncation.
s := v.AsString()
if limit >= 0 && len(s) > limit {
return attribute.ByteSliceValue([]byte(s[:limit]))
}
case attribute.SLICE:
sl := v.AsSlice()
if !slices.ContainsFunc(sl, func(e attribute.Value) bool { return needsTruncation(limit, e) }) {
return v
}
newSl := make([]attribute.Value, len(sl))
for i, elem := range sl {
newSl[i] = TruncateValue(limit, elem)
}
return attribute.SliceValue(newSl...)
case attribute.MAP:
m := v.AsMap()
if !slices.ContainsFunc(m, func(kv attribute.KeyValue) bool { return needsTruncation(limit, kv.Value) }) {
return v
}
newM := make([]attribute.KeyValue, len(m))
for i, elem := range m {
elem.Value = TruncateValue(limit, elem.Value)
newM[i] = elem
}
return attribute.MapValue(newM...)
}
return v
}
// stringNeedsTruncation reports whether s would be modified by truncate for the
// given limit.
func stringNeedsTruncation(limit int, s string) bool {
if limit < 0 || len(s) <= limit {
return false
}
return utf8.RuneCountInString(s) > limit || !utf8.ValidString(s)
}
// needsTruncation reports whether v would be modified by TruncateValue for the
// given limit.
func needsTruncation(limit int, v attribute.Value) bool {
switch v.Type() {
case attribute.STRING:
return stringNeedsTruncation(limit, v.AsString())
case attribute.BYTESLICE:
// len(v.AsString()) is identical to len(v.AsByteSlice()) but
// avoids memory allocation.
if limit >= 0 && len(v.AsString()) > limit {
return true
}
case attribute.STRINGSLICE:
for _, s := range v.AsStringSlice() {
if stringNeedsTruncation(limit, s) {
return true
}
}
case attribute.SLICE:
return slices.ContainsFunc(v.AsSlice(), func(e attribute.Value) bool { return needsTruncation(limit, e) })
case attribute.MAP:
return slices.ContainsFunc(
v.AsMap(),
func(kv attribute.KeyValue) bool { return needsTruncation(limit, kv.Value) },
)
}
return false
}
// truncate returns a truncated version of s such that it contains less than
// the limit number of characters. Truncation is applied by returning the limit
// number of valid characters contained in s.
//
// If limit is negative, it returns the original string.
//
// UTF-8 is supported. When truncating, all invalid characters are dropped
// before applying truncation.
//
// If s already contains less than the limit number of bytes, it is returned
// unchanged. No invalid characters are removed.
func truncate(limit int, s string) string {
// This prioritize performance in the following order based on the most
// common expected use-cases.
//
// - Short values less than the default limit (128).
// - Strings with valid encodings that exceed the limit.
// - No limit.
// - Strings with invalid encodings that exceed the limit.
if limit < 0 || len(s) <= limit {
return s
}
// Optimistically, assume all valid UTF-8.
var b strings.Builder
count := 0
for i, c := range s {
if c != utf8.RuneError {
count++
if count > limit {
return s[:i]
}
continue
}
_, size := utf8.DecodeRuneInString(s[i:])
if size == 1 {
// Invalid encoding.
b.Grow(len(s) - 1)
_, _ = b.WriteString(s[:i])
s = s[i:]
break
}
}
// Fast-path, no invalid input.
if b.Cap() == 0 {
return s
}
// Truncate while validating UTF-8.
for i := 0; i < len(s) && count < limit; {
c := s[i]
if c < utf8.RuneSelf {
// Optimization for single byte runes (common case).
_ = b.WriteByte(c)
i++
count++
continue
}
_, size := utf8.DecodeRuneInString(s[i:])
if size == 1 {
// We checked for all 1-byte runes above, this is a RuneError.
i++
continue
}
_, _ = b.WriteString(s[i : i+size])
i += size
count++
}
return b.String()
}
+31
View File
@@ -0,0 +1,31 @@
// Copyright The OpenTelemetry Authors
// SPDX-License-Identifier: Apache-2.0
// DO NOT MODIFY. Generated by gotmpl.
// source: internal/shared/counter/counter.go.tmpl
// Package counter provides a simple counter for generating unique IDs.
//
// This package is used to generate unique IDs while allowing testing packages
// to reset the counter.
package counter
import "sync/atomic"
// exporterN is a global 0-based count of the number of exporters created.
var exporterN atomic.Int64
// NextExporterID returns the next unique ID for an exporter.
func NextExporterID() int64 {
const inc = 1
return exporterN.Add(inc) - inc
}
// SetExporterID sets the exporter ID counter to v and returns the previous
// value.
//
// This function is useful for testing purposes, allowing you to reset the
// counter. It should not be used in production code.
func SetExporterID(v int64) int64 {
return exporterN.Swap(v)
}
@@ -0,0 +1,127 @@
// Copyright The OpenTelemetry Authors
// SPDX-License-Identifier: Apache-2.0
package observ
import (
"context"
"errors"
"fmt"
"go.opentelemetry.io/otel"
"go.opentelemetry.io/otel/attribute"
"go.opentelemetry.io/otel/metric"
"go.opentelemetry.io/otel/sdk"
"go.opentelemetry.io/otel/sdk/log/internal/x"
semconv "go.opentelemetry.io/otel/semconv/v1.43.0"
"go.opentelemetry.io/otel/semconv/v1.43.0/otelconv"
)
const (
// SchemaURL is the schema URL of the instrumentation.
SchemaURL = semconv.SchemaURL
)
// ErrQueueFull is the attribute value for the "queue_full" error type.
var ErrQueueFull = otelconv.SDKProcessorLogProcessed{}.AttrErrorType("queue_full")
// BLPComponentName returns the component name attribute for a
// BatchLogProcessor with the given ID.
func BLPComponentName(id int64) attribute.KeyValue {
t := otelconv.ComponentTypeBatchingLogProcessor
name := fmt.Sprintf("%s/%d", t, id)
return semconv.OTelComponentName(name)
}
// BLP is the instrumentation for an OTel SDK BatchLogProcessor.
type BLP struct {
reg metric.Registration
processed metric.Int64Counter
processedOpts []metric.AddOption
processedQueueFullOpts []metric.AddOption
}
// NewBLP creates a new BatchLogProcessor instrumentation.
// Returns nil if observability is not enabled.
func NewBLP(id int64, qLen func() int64, qMax int64) (*BLP, error) {
if !x.Observability.Enabled() {
return nil, nil
}
if qLen == nil {
return nil, errors.New("BLP qLen must not be nil")
}
meter := otel.GetMeterProvider().Meter(
ScopeName,
metric.WithInstrumentationVersion(sdk.Version()),
metric.WithSchemaURL(SchemaURL),
)
qCap, err := otelconv.NewSDKProcessorLogQueueCapacity(meter)
if err != nil {
return nil, fmt.Errorf("failed to create BLP queue capacity metric: %w", err)
}
qCapInst := qCap.Inst()
qSize, err := otelconv.NewSDKProcessorLogQueueSize(meter)
if err != nil {
return nil, fmt.Errorf("failed to create BLP queue size metric: %w", err)
}
qSizeInst := qSize.Inst()
cmpntT := semconv.OTelComponentTypeBatchingLogProcessor
cmpnt := BLPComponentName(id)
set := attribute.NewSet(cmpnt, cmpntT)
// Register callback for async metrics
obsOpts := []metric.ObserveOption{metric.WithAttributeSet(set)}
reg, err := meter.RegisterCallback(
func(_ context.Context, o metric.Observer) error {
o.ObserveInt64(qSizeInst, qLen(), obsOpts...)
o.ObserveInt64(qCapInst, qMax, obsOpts...)
return nil
},
qSizeInst,
qCapInst,
)
if err != nil {
return nil, fmt.Errorf("failed to register BLP queue size/capacity callback: %w", err)
}
processed, err := otelconv.NewSDKProcessorLogProcessed(meter)
if err != nil {
_ = reg.Unregister()
return nil, fmt.Errorf("failed to create BLP processed logs metric: %w", err)
}
processedOpts := []metric.AddOption{metric.WithAttributeSet(set)}
setWithError := attribute.NewSet(cmpnt, cmpntT, ErrQueueFull)
processedQueueFullOpts := []metric.AddOption{metric.WithAttributeSet(setWithError)}
return &BLP{
reg: reg,
processed: processed.Inst(),
processedOpts: processedOpts,
processedQueueFullOpts: processedQueueFullOpts,
}, nil
}
func (b *BLP) Shutdown() error {
if b == nil || b.reg == nil {
return nil
}
return b.reg.Unregister()
}
func (b *BLP) Processed(ctx context.Context, n int64) {
if b.processed.Enabled(ctx) {
b.processed.Add(ctx, n, b.processedOpts...)
}
}
func (b *BLP) ProcessedQueueFull(ctx context.Context, n int64) {
if b.processed.Enabled(ctx) {
b.processed.Add(ctx, n, b.processedQueueFullOpts...)
}
}
+1 -1
View File
@@ -3,4 +3,4 @@
// Package observ provides observability instrumentation for the OTel log SDK
// package.
package observ // import "go.opentelemetry.io/otel/sdk/log/internal/observ"
package observ
@@ -1,7 +1,7 @@
// Copyright The OpenTelemetry Authors
// SPDX-License-Identifier: Apache-2.0
package observ // import "go.opentelemetry.io/otel/sdk/log/internal/observ"
package observ
import (
"context"
@@ -14,8 +14,8 @@ import (
"go.opentelemetry.io/otel/metric"
"go.opentelemetry.io/otel/sdk"
"go.opentelemetry.io/otel/sdk/log/internal/x"
semconv "go.opentelemetry.io/otel/semconv/v1.41.0"
"go.opentelemetry.io/otel/semconv/v1.41.0/otelconv"
semconv "go.opentelemetry.io/otel/semconv/v1.43.0"
"go.opentelemetry.io/otel/semconv/v1.43.0/otelconv"
)
const (
+3
View File
@@ -19,6 +19,9 @@ To opt-in, set the environment variable `OTEL_GO_X_OBSERVABILITY` to `true`.
When enabled, the SDK will create the following metrics using the global `MeterProvider`:
- `otel.sdk.log.created`
- `otel.sdk.processor.log.queue.capacity`
- `otel.sdk.processor.log.queue.size`
- `otel.sdk.processor.log.processed`
Please see the [Semantic conventions for OpenTelemetry SDK metrics] documentation for more details on these metrics.
+1 -1
View File
@@ -2,7 +2,7 @@
// SPDX-License-Identifier: Apache-2.0
// Package x documents experimental features for [go.opentelemetry.io/otel/sdk/log].
package x // import "go.opentelemetry.io/otel/sdk/log/internal/x"
package x
import "strings"
+4 -4
View File
@@ -1,11 +1,11 @@
// Code generated by gotmpl. DO NOT MODIFY.
// source: internal/shared/x/x.go.tmpl
// Copyright The OpenTelemetry Authors
// SPDX-License-Identifier: Apache-2.0
// DO NOT MODIFY. Generated by gotmpl.
// source: internal/shared/x/x.go.tmpl
// Package x documents experimental features for [go.opentelemetry.io/otel/sdk/log].
package x // import "go.opentelemetry.io/otel/sdk/log/internal/x"
package x
import (
"os"
+97 -44
View File
@@ -1,7 +1,7 @@
// Copyright The OpenTelemetry Authors
// SPDX-License-Identifier: Apache-2.0
package log // import "go.opentelemetry.io/otel/sdk/log"
package log
import (
"context"
@@ -11,19 +11,19 @@ import (
"time"
"go.opentelemetry.io/otel"
"go.opentelemetry.io/otel/attribute"
"go.opentelemetry.io/otel/log"
"go.opentelemetry.io/otel/log/embedded"
"go.opentelemetry.io/otel/sdk/instrumentation"
semconv "go.opentelemetry.io/otel/semconv/v1.41.0"
semconv "go.opentelemetry.io/otel/semconv/v1.43.0"
"go.opentelemetry.io/otel/trace"
)
var now = time.Now
const (
exceptionTypeKey = string(semconv.ExceptionTypeKey)
exceptionMessageKey = string(semconv.ExceptionMessageKey)
exceptionStacktraceKey = string(semconv.ExceptionStacktraceKey)
exceptionTypeKey = semconv.ExceptionTypeKey
exceptionMessageKey = semconv.ExceptionMessageKey
)
// Compile-time check logger implements log.Logger.
@@ -55,17 +55,50 @@ func newLogger(p *LoggerProvider, scope instrumentation.Scope) *logger {
}
func (l *logger) Emit(ctx context.Context, r log.Record) {
newRecord := l.newRecord(ctx, r)
for _, p := range l.provider.processors {
if err := p.OnEmit(ctx, &newRecord); err != nil {
otel.Handle(err)
processors := l.provider.processors
if len(processors) == 0 {
if l.provider.stopped.Load() {
return
}
// Emit remains observable without processors, but no lifecycle
// admission or SDK record construction is needed.
l.recordCreated(ctx)
return
}
for _, err := range l.emit(ctx, r, processors) {
otel.Handle(err)
}
}
func (l *logger) emit(ctx context.Context, r log.Record, processors []Processor) []error {
if !l.provider.beginProcessorOperation() {
return nil
}
defer l.provider.endProcessorOperation()
l.recordCreated(ctx)
newRecord := l.newRecord(ctx, r)
var errs []error
for _, processor := range processors {
if err := processor.OnEmit(ctx, &newRecord); err != nil {
errs = append(errs, err)
}
}
return errs
}
func (l *logger) recordCreated(ctx context.Context) {
if l.recCntIncr != nil {
l.recCntIncr(ctx)
}
}
// Enabled returns true if at least one Processor held by the LoggerProvider
// that created the logger will process for the provided context and param.
//
// Enabled returns false after the LoggerProvider that created l starts shutdown.
//
// If it is not possible to definitively determine the record will be
// processed, true will be returned by default. A value of false will only be
// returned if it can be positively verified that no Processor will process.
@@ -76,7 +109,13 @@ func (l *logger) Enabled(ctx context.Context, param log.EnabledParameters) bool
EventName: param.EventName,
}
for _, processor := range l.provider.processors {
processors := l.provider.processors
if len(processors) == 0 || !l.provider.beginProcessorOperation() {
return false
}
defer l.provider.endProcessorOperation()
for _, processor := range processors {
if processor.Enabled(ctx, p) {
// At least one Processor will process the Record.
return true
@@ -106,10 +145,6 @@ func (l *logger) newRecord(ctx context.Context, r log.Record) Record {
attributeCountLimit: l.provider.attributeCountLimit,
allowDupKeys: l.provider.allowDupKeys,
}
if l.recCntIncr != nil {
l.recCntIncr(ctx)
}
// This ensures we deduplicate key-value collections in the log body
newRecord.SetBody(r.Body())
@@ -118,47 +153,65 @@ func (l *logger) newRecord(ctx context.Context, r log.Record) Record {
newRecord.observedTimestamp = now()
}
hasExceptionAttr := false
r.WalkAttributes(func(kv log.KeyValue) bool {
// User-provided exception attributes MUST take precedence. Track message
// and type independently so a supplied value suppresses only its own
// derivation.
var hasExceptionMessage, hasExceptionType bool
r.WalkAttributes(func(kv attribute.KeyValue) bool {
switch kv.Key {
case exceptionTypeKey, exceptionMessageKey, exceptionStacktraceKey:
hasExceptionAttr = true
case exceptionMessageKey:
hasExceptionMessage = true
case exceptionTypeKey:
hasExceptionType = true
}
newRecord.AddAttributes(kv)
return true
})
if err := r.Err(); err != nil && !hasExceptionAttr {
addExceptionAttrs(&newRecord, err)
// Avoid inspecting the error for attributes when the caller has
// already supplied the attributes.
if err := r.Err(); err != nil && (!hasExceptionMessage || !hasExceptionType) {
// Derive missing exception attributes by default, as required by the
// Logs SDK specification. Attribute limits may constrain generation,
// so stop once there is no capacity for another attribute.
var attrs [2]attribute.KeyValue
n := 0
// Derived attributes are buffered until flush, so the current attribute
// count stays unchanged while missing values are prepared.
hasLimit := newRecord.hasAttributeCountLimit()
var remaining int
if hasLimit {
remaining = newRecord.attributeCountLimit - newRecord.AttributesLen()
}
if !hasExceptionMessage {
if msg := err.Error(); msg != "" {
if hasLimit && remaining <= n {
goto flush
}
attrs[n] = exceptionMessageKey.String(msg)
n++
}
}
if !hasExceptionType {
if errType := errorType(err); errType != "" {
if hasLimit && remaining <= n {
goto flush
}
attrs[n] = exceptionTypeKey.String(errType)
n++
}
}
flush:
if n > 0 {
newRecord.addAttrs(attrs[:n])
}
}
return newRecord
}
func addExceptionAttrs(r *Record, err error) {
var attrs [2]log.KeyValue
n := 0
if msg := err.Error(); msg != "" {
if r.attributeCountLimit > 0 && r.attributeCountLimit-r.AttributesLen() < n+1 {
goto flush
}
attrs[n] = log.String(exceptionMessageKey, msg)
n++
}
if errType := errorType(err); errType != "" {
if r.attributeCountLimit > 0 && r.attributeCountLimit-r.AttributesLen() < n+1 {
goto flush
}
attrs[n] = log.String(exceptionTypeKey, errType)
n++
}
flush:
if n > 0 {
r.addAttrs(attrs[:n])
}
}
func errorType(err error) string {
if et, ok := err.(interface{ ErrorType() string }); ok {
if s := et.ErrorType(); s != "" {
+23 -6
View File
@@ -1,7 +1,7 @@
// Copyright The OpenTelemetry Authors
// SPDX-License-Identifier: Apache-2.0
package log // import "go.opentelemetry.io/otel/sdk/log"
package log
import (
"context"
@@ -12,13 +12,23 @@ import (
// Processor handles the processing of log records.
//
// Any of the Processor's methods may be called concurrently with itself
// or with other methods. It is the responsibility of the Processor to manage
// this concurrency.
// Enabled, OnEmit, and ForceFlush may be called concurrently with themselves
// or each other. It is the responsibility of the Processor to manage this
// concurrency.
//
// A Processor must be registered only once and with a single
// [LoggerProvider]. Registering the same Processor with multiple providers or
// multiple times with the same provider is not supported.
//
// A [LoggerProvider] stops admitting new operations that invoke Enabled,
// OnEmit, or ForceFlush when shutdown starts. Callers that use a Processor
// directly are responsible for coordinating those calls with Shutdown.
type Processor interface {
// Enabled reports whether the Processor will process for the given context
// and param.
//
// Enabled is called synchronously and should not block.
//
// The param contains a subset of the information that will be available
// in the Record passed to OnEmit, as defined by EnabledParameters.
// A field being unset in param does not imply the corresponding field
@@ -42,12 +52,12 @@ type Processor interface {
// The SDK's Logger.Enabled returns false if all the registered processors
// return false. Otherwise, it returns true.
//
// Implementations of this method need to be safe for a user to call
// concurrently.
Enabled(ctx context.Context, param EnabledParameters) bool
// OnEmit is called when a Record is emitted.
//
// OnEmit is called synchronously and should not block.
//
// OnEmit will be called independent of Enabled. Implementations need to
// validate the arguments themselves before processing.
//
@@ -73,6 +83,13 @@ type Processor interface {
// resources held by the Processor (and any underlying Exporter) should be
// done in this call.
//
// A LoggerProvider calls Shutdown at most once. Before calling it, the
// LoggerProvider waits for all Enabled, OnEmit, and ForceFlush calls it
// admitted to complete. If the LoggerProvider's Shutdown context is canceled
// while waiting, Shutdown is not called.
//
// Shutdown must include the effects of ForceFlush.
//
// The deadline or cancellation of the passed context must be honored. An
// appropriate error should be returned in these situations.
//
+105 -17
View File
@@ -1,7 +1,7 @@
// Copyright The OpenTelemetry Authors
// SPDX-License-Identifier: Apache-2.0
package log // import "go.opentelemetry.io/otel/sdk/log"
package log
import (
"context"
@@ -15,6 +15,7 @@ import (
"go.opentelemetry.io/otel/log/embedded"
"go.opentelemetry.io/otel/log/noop"
"go.opentelemetry.io/otel/sdk/instrumentation"
"go.opentelemetry.io/otel/sdk/log/internal/attrnorm"
"go.opentelemetry.io/otel/sdk/resource"
)
@@ -65,7 +66,11 @@ func newProviderConfig(opts []LoggerProviderOption) providerConfig {
}
// LoggerProvider handles the creation and coordination of Loggers. All Loggers
// created by a LoggerProvider will be associated with the same Resource.
// it creates are associated with the same Resource.
//
// After [LoggerProvider.Shutdown] starts, calls to [log.Logger.Enabled] on
// Loggers created by the LoggerProvider return false, and calls to
// [log.Logger.Emit] perform no operation.
type LoggerProvider struct {
embedded.LoggerProvider
@@ -78,7 +83,10 @@ type LoggerProvider struct {
loggersMu sync.Mutex
loggers map[instrumentation.Scope]*logger
stopped atomic.Bool
stopped atomic.Bool
processorOperationsMu sync.Mutex
processorOperationsActive int
processorOperationsDone chan struct{}
noCmp [0]func() //nolint: unused // This is indeed used.
}
@@ -105,7 +113,7 @@ func NewLoggerProvider(opts ...LoggerProviderOption) *LoggerProvider {
// Logger returns a new [log.Logger] with the provided name and configuration.
//
// If p is shut down, a [noop.Logger] instance is returned.
// Calls made after [LoggerProvider.Shutdown] starts return a [noop.Logger].
//
// This method can be called concurrently.
func (p *LoggerProvider) Logger(name string, opts ...log.LoggerOption) log.Logger {
@@ -118,11 +126,15 @@ func (p *LoggerProvider) Logger(name string, opts ...log.LoggerOption) log.Logge
}
cfg := log.NewLoggerConfig(opts...)
attrs := cfg.InstrumentationAttributes()
if !p.allowDupKeys {
attrs, _ = attrnorm.Set(attrs)
}
scope := instrumentation.Scope{
Name: name,
Version: cfg.InstrumentationVersion(),
SchemaURL: cfg.SchemaURL(),
Attributes: cfg.InstrumentationAttributes(),
Attributes: attrs,
}
p.loggersMu.Lock()
@@ -143,14 +155,40 @@ func (p *LoggerProvider) Logger(name string, opts ...log.LoggerOption) log.Logge
return l
}
// Shutdown shuts down the provider and all processors.
// Shutdown shuts down the provider and all processors in the order they were
// registered.
//
// The first call stops admitting new operations that invoke processor Enabled,
// OnEmit, or ForceFlush methods. It waits for operations already admitted to
// complete before synchronously invoking each processor's Shutdown method. If
// ctx is canceled before the admitted operations complete, Shutdown returns
// ctx.Err() without invoking processor Shutdown.
//
// Concurrent or subsequent Shutdown calls return nil without invoking
// processor Shutdown.
//
// Shutdown must not be called directly or indirectly from a Processor method.
//
// This method can be called concurrently.
func (p *LoggerProvider) Shutdown(ctx context.Context) error {
stopped := p.stopped.Swap(true)
if stopped {
p.processorOperationsMu.Lock()
if p.stopped.Load() {
p.processorOperationsMu.Unlock()
return nil
}
p.processorOperationsDone = make(chan struct{})
p.stopped.Store(true)
if p.processorOperationsActive == 0 {
close(p.processorOperationsDone)
}
p.processorOperationsMu.Unlock()
// All count updates happen while processorOperationsMu is held. Therefore,
// either Shutdown observes zero operations and closes the channel above, or
// the final operation observes stopped and closes it synchronously.
if err := p.waitForProcessorOperations(ctx); err != nil {
return err
}
var err error
for _, p := range p.processors {
@@ -161,11 +199,14 @@ func (p *LoggerProvider) Shutdown(ctx context.Context) error {
// ForceFlush flushes all processors.
//
// Once Shutdown starts, ForceFlush performs no operation and returns nil.
//
// This method can be called concurrently.
func (p *LoggerProvider) ForceFlush(ctx context.Context) error {
if p.stopped.Load() {
if len(p.processors) == 0 || !p.beginProcessorOperation() {
return nil
}
defer p.endProcessorOperation()
var err error
for _, p := range p.processors {
@@ -174,6 +215,47 @@ func (p *LoggerProvider) ForceFlush(ctx context.Context) error {
return err
}
func (p *LoggerProvider) beginProcessorOperation() bool {
p.processorOperationsMu.Lock()
defer p.processorOperationsMu.Unlock()
if p.stopped.Load() {
return false
}
p.processorOperationsActive++
return true
}
func (p *LoggerProvider) endProcessorOperation() {
p.processorOperationsMu.Lock()
defer p.processorOperationsMu.Unlock()
p.processorOperationsActive--
if p.processorOperationsActive == 0 && p.stopped.Load() {
close(p.processorOperationsDone)
}
}
func (p *LoggerProvider) waitForProcessorOperations(ctx context.Context) error {
if err := ctx.Err(); err != nil {
return err
}
return waitForProcessorOperationsCompletion(ctx, p.processorOperationsDone)
}
func waitForProcessorOperationsCompletion(ctx context.Context, done <-chan struct{}) error {
select {
case <-done:
case <-ctx.Done():
// Prefer a completed drain when it races with cancellation.
select {
case <-done:
default:
return ctx.Err()
}
}
return nil
}
// LoggerProviderOption applies a configuration option value to a LoggerProvider.
type LoggerProviderOption interface {
apply(providerConfig) providerConfig
@@ -257,17 +339,23 @@ func WithAttributeValueLengthLimit(limit int) LoggerProviderOption {
})
}
// WithAllowKeyDuplication sets whether deduplication is skipped for log attributes or other key-value collections.
// WithAllowKeyDuplication sets whether deduplication is skipped for log record
// and instrumentation scope key-value collections.
//
// By default, the key-value collections within a log record are deduplicated to comply with the OpenTelemetry Specification.
// Deduplication means that if multiple key–value pairs with the same key are present, only a single pair
// is retained and others are discarded.
// By default, the key-value collections within a log record and
// instrumentation scope are deduplicated to comply with the OpenTelemetry
// Specification.
// Deduplication means that if multiple key-value pairs with the same key are
// present, only a single pair is retained and others are discarded. Resource
// attributes are always deduplicated by go.opentelemetry.io/otel/sdk/resource.
//
// Disabling deduplication with this option can improve performance e.g. of adding attributes to the log record.
// Disabling deduplication with this option can improve performance e.g. of
// adding attributes to the log record.
//
// Note that if you disable deduplication, you are responsible for ensuring that duplicate
// key-value pairs within in a single collection are not emitted,
// or that the telemetry receiver can handle such duplicates.
// Receivers may handle duplicate keys unpredictably. If you disable
// deduplication, you are responsible for ensuring that duplicate keys within a
// single collection are not emitted, or that the telemetry receiver can handle
// such duplicates.
func WithAllowKeyDuplication() LoggerProviderOption {
return loggerProviderOptionFunc(func(cfg providerConfig) providerConfig {
cfg.allowDupKeys = newSetting(true)
+50 -319
View File
@@ -1,18 +1,18 @@
// Copyright The OpenTelemetry Authors
// SPDX-License-Identifier: Apache-2.0
package log // import "go.opentelemetry.io/otel/sdk/log"
package log
import (
"slices"
"strings"
"sync"
"time"
"unicode/utf8"
"go.opentelemetry.io/otel/attribute"
"go.opentelemetry.io/otel/internal/global"
"go.opentelemetry.io/otel/log"
"go.opentelemetry.io/otel/sdk/instrumentation"
"go.opentelemetry.io/otel/sdk/log/internal/attrnorm"
"go.opentelemetry.io/otel/sdk/resource"
"go.opentelemetry.io/otel/trace"
)
@@ -37,14 +37,14 @@ var logKeyValuePairDropped = sync.OnceFunc(func() {
// uniquePool is a pool of unique attributes used for attributes de-duplication.
var uniquePool = sync.Pool{
New: func() any { return new([]log.KeyValue) },
New: func() any { return new([]attribute.KeyValue) },
}
func getUnique() *[]log.KeyValue {
return uniquePool.Get().(*[]log.KeyValue)
func getUnique() *[]attribute.KeyValue {
return uniquePool.Get().(*[]attribute.KeyValue)
}
func putUnique(v *[]log.KeyValue) {
func putUnique(v *[]attribute.KeyValue) {
if cap(*v) <= maxUniqueSize {
clear(*v)
*v = (*v)[:0]
@@ -54,32 +54,18 @@ func putUnique(v *[]log.KeyValue) {
// indexPool is a pool of index maps used for attributes de-duplication.
var indexPool = sync.Pool{
New: func() any { return make(map[string]int) },
New: func() any { return make(map[attribute.Key]int) },
}
func getIndex() map[string]int {
return indexPool.Get().(map[string]int)
func getIndex() map[attribute.Key]int {
return indexPool.Get().(map[attribute.Key]int)
}
func putIndex(index map[string]int) {
func putIndex(index map[attribute.Key]int) {
clear(index)
indexPool.Put(index)
}
// seenPool is a pool of seen keys used for maps de-duplication.
var seenPool = sync.Pool{
New: func() any { return make(map[string]struct{}) },
}
func getSeen() map[string]struct{} {
return seenPool.Get().(map[string]struct{})
}
func putSeen(seen map[string]struct{}) {
clear(seen)
seenPool.Put(seen)
}
// Record is a log record emitted by the Logger.
// A log record with non-empty event name is interpreted as an event record.
//
@@ -95,7 +81,7 @@ type Record struct {
observedTimestamp time.Time
severity log.Severity
severityText string
body log.Value
body attribute.Value
// The fields below are for optimizing the implementation of Attributes and
// AddAttributes. This design is borrowed from the slog Record type:
@@ -104,7 +90,7 @@ type Record struct {
// Allocation optimization: an inline array sized to hold
// the majority of log calls (based on examination of open-source
// code). It holds the start of the list of attributes.
front [attributesInlineCount]log.KeyValue
front [attributesInlineCount]attribute.KeyValue
// The number of attributes in front.
nFront int
@@ -113,7 +99,7 @@ type Record struct {
// Invariants:
// - len(back) > 0 if nFront == len(front)
// - Unused array elements are zero-ed. Used to detect mistakes.
back []log.KeyValue
back []attribute.KeyValue
// dropped is the count of attributes that have been dropped when limits
// were reached.
@@ -200,22 +186,22 @@ func (r *Record) SetSeverityText(text string) {
}
// Body returns the body of the log record.
func (r *Record) Body() log.Value {
func (r *Record) Body() attribute.Value {
return r.body
}
// SetBody sets the body of the log record.
func (r *Record) SetBody(v log.Value) {
func (r *Record) SetBody(v attribute.Value) {
if !r.allowDupKeys {
r.body = r.dedupeBodyCollections(v)
r.body, _ = attrnorm.Value(v)
} else {
r.body = v
}
}
// WalkAttributes walks all attributes the log record holds by calling f for
// each on each [log.KeyValue] in the [Record]. Iteration stops if f returns false.
func (r *Record) WalkAttributes(f func(log.KeyValue) bool) {
// each on each [attribute.KeyValue] in the [Record]. Iteration stops if f returns false.
func (r *Record) WalkAttributes(f func(attribute.KeyValue) bool) {
for i := 0; i < r.nFront; i++ {
if !f(r.front[i]) {
return
@@ -230,7 +216,7 @@ func (r *Record) WalkAttributes(f func(log.KeyValue) bool) {
// AddAttributes adds attributes to the log record.
// Attributes in attrs will overwrite any attribute already added to r with the same key.
func (r *Record) AddAttributes(attrs ...log.KeyValue) {
func (r *Record) AddAttributes(attrs ...attribute.KeyValue) {
n := r.AttributesLen()
if n == 0 {
// Avoid the more complex duplicate map lookups below.
@@ -242,7 +228,7 @@ func (r *Record) AddAttributes(attrs ...log.KeyValue) {
}
}
attrs, drop := head(attrs, r.attributeCountLimit)
attrs, drop := r.head(attrs)
r.addDropped(drop)
r.addAttrs(attrs)
@@ -295,12 +281,12 @@ func (r *Record) AddAttributes(attrs ...log.KeyValue) {
}
if dropped > 0 {
attrs = make([]log.KeyValue, len(*unique))
attrs = make([]attribute.KeyValue, len(*unique))
copy(attrs, *unique)
}
}
if r.attributeCountLimit > 0 && n+len(attrs) > r.attributeCountLimit {
if r.hasAttributeCountLimit() && n+len(attrs) > r.attributeCountLimit {
// Truncate the now unique attributes to comply with limit.
//
// Do not use head(attrs, r.attributeCountLimit - n) here. If
@@ -320,7 +306,7 @@ func (r *Record) AddAttributes(attrs ...log.KeyValue) {
//
// The returned index is taken from the indexPool. It is the callers
// responsibility to return the index to that pool (putIndex) when done.
func (r *Record) attrIndex() map[string]int {
func (r *Record) attrIndex() map[attribute.Key]int {
index := getIndex()
for i := 0; i < r.nFront; i++ {
key := r.front[i].Key
@@ -336,7 +322,7 @@ func (r *Record) attrIndex() map[string]int {
// addAttrs adds attrs to the Record r. This does not validate any limits or
// duplication of attributes, these tasks are left to the caller to handle
// prior to calling.
func (r *Record) addAttrs(attrs []log.KeyValue) {
func (r *Record) addAttrs(attrs []attribute.KeyValue) {
var i int
for i = 0; i < len(attrs) && r.nFront < len(r.front); i++ {
a := attrs[i]
@@ -354,7 +340,7 @@ func (r *Record) addAttrs(attrs []log.KeyValue) {
}
// SetAttributes sets (and overrides) attributes to the log record.
func (r *Record) SetAttributes(attrs ...log.KeyValue) {
func (r *Record) SetAttributes(attrs ...attribute.KeyValue) {
var drop int
r.dropped = 0
if !r.allowDupKeys {
@@ -364,7 +350,7 @@ func (r *Record) SetAttributes(attrs ...log.KeyValue) {
}
}
attrs, drop = head(attrs, r.attributeCountLimit)
attrs, drop = r.head(attrs)
r.addDropped(drop)
r.nFront = 0
@@ -381,17 +367,23 @@ func (r *Record) SetAttributes(attrs ...log.KeyValue) {
}
}
// head returns the first n values of kvs along with the number of elements
// dropped. If n is less than or equal to zero, kvs is returned with 0.
func head(kvs []log.KeyValue, n int) (out []log.KeyValue, dropped int) {
if n > 0 && len(kvs) > n {
return kvs[:n], len(kvs) - n
// head returns the attributes r can retain along with the number dropped.
func (r *Record) head(kvs []attribute.KeyValue) (out []attribute.KeyValue, dropped int) {
if r.attributeCountLimit < 0 || len(kvs) <= r.attributeCountLimit {
return kvs, 0
}
return kvs, 0
if r.attributeCountLimit == 0 {
return nil, len(kvs)
}
return kvs[:r.attributeCountLimit], len(kvs) - r.attributeCountLimit
}
func (r *Record) hasAttributeCountLimit() bool {
return r.attributeCountLimit >= 0
}
// dedup deduplicates kvs front-to-back with the last value saved.
func dedup(kvs []log.KeyValue) (unique []log.KeyValue, dropped int) {
func dedup(kvs []attribute.KeyValue) (unique []attribute.KeyValue, dropped int) {
if len(kvs) <= 1 {
return kvs, 0 // No deduplication needed.
}
@@ -415,7 +407,7 @@ func dedup(kvs []log.KeyValue) (unique []log.KeyValue, dropped int) {
return kvs, 0
}
unique = make([]log.KeyValue, len(*u))
unique = make([]attribute.KeyValue, len(*u))
copy(unique, *u)
return unique, dropped
}
@@ -482,275 +474,14 @@ func (r *Record) Clone() Record {
return res
}
func (r *Record) applyAttrLimitsAndDedup(attr log.KeyValue) log.KeyValue {
attr.Value = r.applyValueLimitsAndDedup(attr.Value)
func (r *Record) applyAttrLimitsAndDedup(attr attribute.KeyValue) attribute.KeyValue {
if !r.allowDupKeys {
var changed bool
attr, changed = attrnorm.KeyValue(attr)
if changed {
logKeyValuePairDropped()
}
}
attr.Value = attrnorm.TruncateValue(r.attributeValueLengthLimit, attr.Value)
return attr
}
func (r *Record) applyValueLimitsAndDedup(val log.Value) log.Value {
switch val.Kind() {
case log.KindString:
s := val.AsString()
if r.attributeValueLengthLimit >= 0 && len(s) > r.attributeValueLengthLimit {
val = log.StringValue(truncate(r.attributeValueLengthLimit, s))
}
case log.KindSlice:
sl := val.AsSlice()
// First check if any limits need to be applied.
if slices.ContainsFunc(sl, r.needsValueLimitsOrDedup) {
// Create a new slice to avoid modifying the original.
newSl := make([]log.Value, len(sl))
for i, item := range sl {
newSl[i] = r.applyValueLimitsAndDedup(item)
}
val = log.SliceValue(newSl...)
}
case log.KindBytes:
bs := val.AsBytes()
if r.attributeValueLengthLimit >= 0 && len(bs) > r.attributeValueLengthLimit {
val = log.BytesValue(bs[:r.attributeValueLengthLimit])
}
case log.KindMap:
kvs := val.AsMap()
var newKvs []log.KeyValue
var dropped int
if !r.allowDupKeys {
// Deduplicate then truncate.
// Do not do at the same time to avoid wasted truncation operations.
newKvs, dropped = dedup(kvs)
if dropped > 0 {
logKeyValuePairDropped()
}
} else {
newKvs = kvs
}
// Check if any attribute limits need to be applied.
needsChange := false
if dropped > 0 {
needsChange = true // Already changed by dedup.
} else {
for _, kv := range newKvs {
if r.needsValueLimitsOrDedup(kv.Value) {
needsChange = true
break
}
}
}
if needsChange {
// Only create new slice if changes are needed.
if dropped == 0 {
// Make a copy to avoid modifying the original.
newKvs = make([]log.KeyValue, len(kvs))
copy(newKvs, kvs)
}
for i := range newKvs {
newKvs[i] = r.applyAttrLimitsAndDedup(newKvs[i])
}
val = log.MapValue(newKvs...)
}
}
return val
}
// needsValueLimitsOrDedup checks if a value would be modified by applyValueLimitsAndDedup.
func (r *Record) needsValueLimitsOrDedup(val log.Value) bool {
switch val.Kind() {
case log.KindString:
return r.attributeValueLengthLimit >= 0 && len(val.AsString()) > r.attributeValueLengthLimit
case log.KindSlice:
if slices.ContainsFunc(val.AsSlice(), r.needsValueLimitsOrDedup) {
return true
}
case log.KindBytes:
bs := val.AsBytes()
if r.attributeValueLengthLimit >= 0 && len(bs) > r.attributeValueLengthLimit {
return true
}
case log.KindMap:
kvs := val.AsMap()
if !r.allowDupKeys && len(kvs) > 1 {
// Check for duplicates.
hasDuplicates := func() bool {
seen := getSeen()
defer putSeen(seen)
for _, kv := range kvs {
if _, ok := seen[kv.Key]; ok {
return true
}
seen[kv.Key] = struct{}{}
}
return false
}()
if hasDuplicates {
return true
}
}
for _, kv := range kvs {
if r.needsValueLimitsOrDedup(kv.Value) {
return true
}
}
}
return false
}
func (r *Record) dedupeBodyCollections(val log.Value) log.Value {
switch val.Kind() {
case log.KindSlice:
sl := val.AsSlice()
// Check if any nested values need deduplication.
if slices.ContainsFunc(sl, r.needsBodyDedup) {
// Create a new slice to avoid modifying the original.
newSl := make([]log.Value, len(sl))
for i, item := range sl {
newSl[i] = r.dedupeBodyCollections(item)
}
val = log.SliceValue(newSl...)
}
case log.KindMap:
kvs := val.AsMap()
newKvs, dropped := dedup(kvs)
// Check if any nested values need deduplication.
needsValueChange := false
for _, kv := range newKvs {
if r.needsBodyDedup(kv.Value) {
needsValueChange = true
break
}
}
if dropped > 0 || needsValueChange {
// Only create new value if changes are needed.
if dropped == 0 {
// Make a copy to avoid modifying the original.
newKvs = make([]log.KeyValue, len(kvs))
copy(newKvs, kvs)
}
for i := range newKvs {
newKvs[i].Value = r.dedupeBodyCollections(newKvs[i].Value)
}
val = log.MapValue(newKvs...)
}
}
return val
}
// needsBodyDedup checks if a value would be modified by dedupeBodyCollections.
func (r *Record) needsBodyDedup(val log.Value) bool {
switch val.Kind() {
case log.KindSlice:
if slices.ContainsFunc(val.AsSlice(), r.needsBodyDedup) {
return true
}
case log.KindMap:
kvs := val.AsMap()
if len(kvs) > 1 {
// Check for duplicates.
hasDuplicates := func() bool {
seen := getSeen()
defer putSeen(seen)
for _, kv := range kvs {
if _, ok := seen[kv.Key]; ok {
return true
}
seen[kv.Key] = struct{}{}
}
return false
}()
if hasDuplicates {
return true
}
}
for _, kv := range kvs {
if r.needsBodyDedup(kv.Value) {
return true
}
}
}
return false
}
// truncate returns a truncated version of s such that it contains less than
// the limit number of characters. Truncation is applied by returning the limit
// number of valid characters contained in s.
//
// If limit is negative, it returns the original string.
//
// UTF-8 is supported. When truncating, all invalid characters are dropped
// before applying truncation.
//
// If s already contains less than the limit number of bytes, it is returned
// unchanged. No invalid characters are removed.
func truncate(limit int, s string) string {
// This prioritize performance in the following order based on the most
// common expected use-cases.
//
// - Short values less than the default limit (128).
// - Strings with valid encodings that exceed the limit.
// - No limit.
// - Strings with invalid encodings that exceed the limit.
if limit < 0 || len(s) <= limit {
return s
}
// Optimistically, assume all valid UTF-8.
var b strings.Builder
count := 0
for i, c := range s {
if c != utf8.RuneError {
count++
if count > limit {
return s[:i]
}
continue
}
_, size := utf8.DecodeRuneInString(s[i:])
if size == 1 {
// Invalid encoding.
b.Grow(len(s) - 1)
_, _ = b.WriteString(s[:i])
s = s[i:]
break
}
}
// Fast-path, no invalid input.
if b.Cap() == 0 {
return s
}
// Truncate while validating UTF-8.
for i := 0; i < len(s) && count < limit; {
c := s[i]
if c < utf8.RuneSelf {
// Optimization for single byte runes (common case).
_ = b.WriteByte(c)
i++
count++
continue
}
_, size := utf8.DecodeRuneInString(s[i:])
if size == 1 {
// We checked for all 1-byte runes above, this is a RuneError.
i++
continue
}
_, _ = b.WriteString(s[i : i+size])
i += size
count++
}
return b.String()
}
+1 -1
View File
@@ -5,7 +5,7 @@
// Use of this source code is governed by a BSD-style
// license that can be found in the LICENSE file.
package log // import "go.opentelemetry.io/otel/sdk/log"
package log
// A ring is an element of a circular list, or ring. Rings do not have a
// beginning or end; a pointer to any ring element serves as reference to the
+1 -1
View File
@@ -1,7 +1,7 @@
// Copyright The OpenTelemetry Authors
// SPDX-License-Identifier: Apache-2.0
package log // import "go.opentelemetry.io/otel/sdk/log"
package log
import (
"fmt"
+3 -3
View File
@@ -1,7 +1,7 @@
// Copyright The OpenTelemetry Authors
// SPDX-License-Identifier: Apache-2.0
package log // import "go.opentelemetry.io/otel/sdk/log"
package log
import (
"context"
@@ -80,13 +80,13 @@ func (s *SimpleProcessor) OnEmit(ctx context.Context, r *Record) (err error) {
return s.exporter.Export(ctx, *records)
}
// Shutdown shuts down the exporter.
// Shutdown flushes the exporter before shutting it down.
func (s *SimpleProcessor) Shutdown(ctx context.Context) error {
if s.exporter == nil {
return nil
}
return s.exporter.Shutdown(ctx)
return shutdownExporter(ctx, s.exporter)
}
// ForceFlush flushes the exporter.
File diff suppressed because it is too large Load Diff
+5 -4
View File
@@ -1131,12 +1131,11 @@ go.opentelemetry.io/otel/semconv/v1.37.0/rpcconv
go.opentelemetry.io/otel/semconv/v1.40.0
go.opentelemetry.io/otel/semconv/v1.40.0/otelconv
go.opentelemetry.io/otel/semconv/v1.41.0
go.opentelemetry.io/otel/semconv/v1.41.0/otelconv
go.opentelemetry.io/otel/semconv/v1.43.0
go.opentelemetry.io/otel/semconv/v1.43.0/httpconv
go.opentelemetry.io/otel/semconv/v1.43.0/otelconv
go.opentelemetry.io/otel/semconv/v1.43.0/rpcconv
# go.opentelemetry.io/otel/exporters/otlp/otlplog/otlploggrpc v0.20.0
# go.opentelemetry.io/otel/exporters/otlp/otlplog/otlploggrpc v0.21.0
## explicit; go 1.25.0
go.opentelemetry.io/otel/exporters/otlp/otlplog/otlploggrpc
go.opentelemetry.io/otel/exporters/otlp/otlplog/otlploggrpc/internal
@@ -1173,7 +1172,7 @@ go.opentelemetry.io/otel/exporters/prometheus/internal
go.opentelemetry.io/otel/exporters/prometheus/internal/counter
go.opentelemetry.io/otel/exporters/prometheus/internal/observ
go.opentelemetry.io/otel/exporters/prometheus/internal/x
# go.opentelemetry.io/otel/log v0.20.0
# go.opentelemetry.io/otel/log v0.21.0
## explicit; go 1.25.0
go.opentelemetry.io/otel/log
go.opentelemetry.io/otel/log/embedded
@@ -1196,9 +1195,11 @@ go.opentelemetry.io/otel/sdk/trace
go.opentelemetry.io/otel/sdk/trace/internal/env
go.opentelemetry.io/otel/sdk/trace/internal/observ
go.opentelemetry.io/otel/sdk/trace/tracetest
# go.opentelemetry.io/otel/sdk/log v0.20.0
# go.opentelemetry.io/otel/sdk/log v0.21.0
## explicit; go 1.25.0
go.opentelemetry.io/otel/sdk/log
go.opentelemetry.io/otel/sdk/log/internal/attrnorm
go.opentelemetry.io/otel/sdk/log/internal/counter
go.opentelemetry.io/otel/sdk/log/internal/observ
go.opentelemetry.io/otel/sdk/log/internal/x
# go.opentelemetry.io/otel/sdk/metric v1.45.0