Something went wrong. Try again.
Monorepo for Tangled tangled.org
Something went wrong. Try again.
5.9 kB · 162 lines
Go
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163package observability
import ( "context" "fmt" "net/http"
"github.com/go-chi/chi/v5" "go.opentelemetry.io/otel" "go.opentelemetry.io/otel/attribute" "go.opentelemetry.io/otel/exporters/otlp/otlptrace/otlptracehttp" "go.opentelemetry.io/otel/propagation" "go.opentelemetry.io/otel/sdk/resource" sdktrace "go.opentelemetry.io/otel/sdk/trace" "go.opentelemetry.io/otel/trace" "tangled.org/core/spindle/config")
const ( UserDIDKey = "tangled.user.did" OwnerDIDKey = "tangled.owner.did" RepoDIDKey = "tangled.repo.did" TargetRepoDIDKey = "tangled.repo.target.did" PipelineIDKey = "tangled.pipeline.id" WorkflowIDKey = "tangled.workflow.id" WorkflowEngineKey = "tangled.workflow.engine" LeaseIDKey = "tangled.lease.id" ExecutorNodeIDKey = "tangled.executor.node_id" JobIDKey = "tangled.job.id" RequestIDKey = "tangled.request.id" CollectionKey = "atproto.collection" RKeyKey = "atproto.rkey" StepNameKey = "tangled.step.name" StepIndexKey = "tangled.step.index")
func InitTracing(ctx context.Context, cfg config.Tracing) (func(context.Context) error, error) { otel.SetTextMapPropagator(propagation.TraceContext{}) if cfg.Endpoint == "" { return func(context.Context) error { return nil }, nil }
opts := []otlptracehttp.Option{otlptracehttp.WithEndpoint(cfg.Endpoint)} if cfg.Insecure { opts = append(opts, otlptracehttp.WithInsecure()) } exporter, err := otlptracehttp.New(ctx, opts...) if err != nil { return nil, fmt.Errorf("creating OTLP trace exporter: %w", err) }
res, err := resource.Merge( resource.Default(), resource.NewSchemaless(attribute.String("service.name", cfg.ServiceName)), ) if err != nil { return nil, fmt.Errorf("creating tracing resource: %w", err) }
provider := sdktrace.NewTracerProvider( sdktrace.WithBatcher(exporter), sdktrace.WithResource(res), sdktrace.WithSampler(sdktrace.ParentBased(sdktrace.TraceIDRatioBased(cfg.SampleRatio))), ) otel.SetTracerProvider(provider)
return func(shutdownCtx context.Context) error { if err := provider.Shutdown(shutdownCtx); err != nil { return fmt.Errorf("shutting down tracer provider: %w", err) } return nil }, nil}
func Tracer() trace.Tracer { return otel.Tracer("tangled.org/core/spindle")}
func OTelRouteMiddleware(next http.Handler) http.Handler { return http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { next.ServeHTTP(w, r)
rctx := chi.RouteContext(r.Context()) if rctx == nil { return } pattern := rctx.RoutePattern() if pattern == "" { return } span := trace.SpanFromContext(r.Context()) span.SetName("HTTP " + pattern) span.SetAttributes(attribute.String("http.route", pattern)) })}
func InjectToTraceparentAndTracestate(ctx context.Context) (string, string) { carrier := propagation.MapCarrier{} propagation.TraceContext{}.Inject(ctx, carrier) return carrier.Get("traceparent"), carrier.Get("tracestate")}
func ExtractFromTraceparentAndTracestate(ctx context.Context, traceparent, tracestate string) context.Context { if traceparent == "" { return ctx } carrier := propagation.MapCarrier{ "traceparent": traceparent, "tracestate": tracestate, } return propagation.TraceContext{}.Extract(ctx, carrier)}
const ( ReqWorkflowsKey = "tangled.resource.requested.workflows" ReqVCPUsKey = "tangled.resource.requested.vcpus" ReqMemoryMiBKey = "tangled.resource.requested.memory_mib" ReqDiskMiBKey = "tangled.resource.requested.disk_mib" ReqCacheBytesKey = "tangled.resource.requested.cache_bytes"
ActCPUUsecKey = "tangled.resource.actual.cpu_usec" ActMemoryCurrentBytesKey = "tangled.resource.actual.memory_current_bytes" ActMemoryPeakBytesKey = "tangled.resource.actual.memory_peak_bytes" ActSwapCurrentBytesKey = "tangled.resource.actual.swap_current_bytes" ActSwapPeakBytesKey = "tangled.resource.actual.swap_peak_bytes" ActPIDsCurrentKey = "tangled.resource.actual.pids_current" ActIOReadBytesKey = "tangled.resource.actual.io_read_bytes" ActIOWriteBytesKey = "tangled.resource.actual.io_write_bytes" ActIOReadOpsKey = "tangled.resource.actual.io_read_ops" ActIOWriteOpsKey = "tangled.resource.actual.io_write_ops" ActVolumeAllocatedBytesKey = "tangled.resource.actual.volume_allocated_bytes" ActCgroupAvailableKey = "tangled.resource.actual.cgroup_available" ActVolumeAvailableKey = "tangled.resource.actual.volume_available")
func RequestedResourceAttrs(workflows, vcpus, memoryMiB, diskMiB, cacheBytes int64) []attribute.KeyValue { return []attribute.KeyValue{ attribute.Int64(ReqWorkflowsKey, workflows), attribute.Int64(ReqVCPUsKey, vcpus), attribute.Int64(ReqMemoryMiBKey, memoryMiB), attribute.Int64(ReqDiskMiBKey, diskMiB), attribute.Int64(ReqCacheBytesKey, cacheBytes), }}
func ActualResourceAttrs(cpuUsec, memoryCurrent, memoryPeak, swapCurrent, swapPeak, pidsCurrent, ioReadBytes, ioWriteBytes, ioReadOps, ioWriteOps, volumeAllocated uint64, cgroupAvailable, volumeAvailable bool) []attribute.KeyValue { return []attribute.KeyValue{ attribute.Int64(ActCPUUsecKey, int64(cpuUsec)), attribute.Int64(ActMemoryCurrentBytesKey, int64(memoryCurrent)), attribute.Int64(ActMemoryPeakBytesKey, int64(memoryPeak)), attribute.Int64(ActSwapCurrentBytesKey, int64(swapCurrent)), attribute.Int64(ActSwapPeakBytesKey, int64(swapPeak)), attribute.Int64(ActPIDsCurrentKey, int64(pidsCurrent)), attribute.Int64(ActIOReadBytesKey, int64(ioReadBytes)), attribute.Int64(ActIOWriteBytesKey, int64(ioWriteBytes)), attribute.Int64(ActIOReadOpsKey, int64(ioReadOps)), attribute.Int64(ActIOWriteOpsKey, int64(ioWriteOps)), attribute.Int64(ActVolumeAllocatedBytesKey, int64(volumeAllocated)), attribute.Bool(ActCgroupAvailableKey, cgroupAvailable), attribute.Bool(ActVolumeAvailableKey, volumeAvailable), }}