@cryptotaxi247 / kubo / commits / 1e3b6c985

feat: add tracing to the commands client

Jorropo committed Mar 28, 2023 at 21:52 UTC 1e3b6c98579cf549f4042e4fa76ec9d155f82bf3
5 files changed +60 -24
cmd/ipfs/main.go
+42 -17
@@ -10,8 +10,15 @@ import (
10 "net/http"
11 "os"
12 "runtime/pprof"
13 + "strings"
14 "time"
15
16 + "github.com/google/uuid"
17 + cmds "github.com/ipfs/go-ipfs-cmds"
18 + "github.com/ipfs/go-ipfs-cmds/cli"
19 + cmdhttp "github.com/ipfs/go-ipfs-cmds/http"
20 + u "github.com/ipfs/go-ipfs-util"
21 + logging "github.com/ipfs/go-log"
22 "github.com/ipfs/kubo/cmd/ipfs/util"
23 oldcmds "github.com/ipfs/kubo/commands"
24 "github.com/ipfs/kubo/core"
@@ -21,22 +28,19 @@ import (
28 "github.com/ipfs/kubo/repo"
29 "github.com/ipfs/kubo/repo/fsrepo"
30 "github.com/ipfs/kubo/tracing"
24 -
25 - cmds "github.com/ipfs/go-ipfs-cmds"
26 - "github.com/ipfs/go-ipfs-cmds/cli"
27 - cmdhttp "github.com/ipfs/go-ipfs-cmds/http"
28 - u "github.com/ipfs/go-ipfs-util"
29 - logging "github.com/ipfs/go-log"
31 ma "github.com/multiformats/go-multiaddr"
32 madns "github.com/multiformats/go-multiaddr-dns"
33 manet "github.com/multiformats/go-multiaddr/net"
33 -
34 - "github.com/google/uuid"
34 + "go.opentelemetry.io/contrib/instrumentation/net/http/otelhttp"
35 "go.opentelemetry.io/otel"
36 + "go.opentelemetry.io/otel/attribute"
37 + "go.opentelemetry.io/otel/codes"
38 + "go.opentelemetry.io/otel/trace"
39 )
40
41 // log is the command logger
42 var log = logging.Logger("cmd/ipfs")
43 +var tracer trace.Tracer
44
45 // declared as a var for testing purposes
46 var dnsResolver = madns.DefaultResolver
@@ -91,7 +95,6 @@ func newUUID(key string) logging.Metadata {
95 func mainRet() (exitCode int) {
96 rand.Seed(time.Now().UnixNano())
97 ctx := logging.ContextWithLoggable(context.Background(), newUUID("session"))
94 - var err error
98
99 tp, err := tracing.NewTracerProvider(ctx)
100 if err != nil {
@@ -103,6 +106,7 @@ func mainRet() (exitCode int) {
106 }
107 }()
108 otel.SetTracerProvider(tp)
109 + tracer = tp.Tracer("Kubo-cli")
110
111 stopFunc, err := profileIfEnabled()
112 if err != nil {
@@ -219,7 +223,7 @@ func apiAddrOption(req *cmds.Request) (ma.Multiaddr, error) {
223 }
224
225 func makeExecutor(req *cmds.Request, env interface{}) (cmds.Executor, error) {
222 - exe := cmds.NewExecutor(req.Root)
226 + exe := tracingWrappedExecutor{cmds.NewExecutor(req.Root)}
227 cctx := env.(*oldcmds.Context)
228
229 // Check if the command is disabled.
@@ -294,23 +298,44 @@ func makeExecutor(req *cmds.Request, env interface{}) (cmds.Executor, error) {
298 opts = append(opts, cmdhttp.ClientWithFallback(exe))
299 }
300
301 + var tpt http.RoundTripper
302 switch network {
303 case "tcp", "tcp4", "tcp6":
304 + tpt = http.DefaultTransport
305 case "unix":
306 path := host
307 host = "unix"
302 - opts = append(opts, cmdhttp.ClientWithHTTPClient(&http.Client{
303 - Transport: &http.Transport{
304 - DialContext: func(_ context.Context, _, _ string) (net.Conn, error) {
305 - return net.Dial("unix", path)
306 - },
308 + tpt = &http.Transport{
309 + DialContext: func(_ context.Context, _, _ string) (net.Conn, error) {
310 + return net.Dial("unix", path)
311 },
308 - }))
312 + }
313 default:
314 return nil, fmt.Errorf("unsupported API address: %s", apiAddr)
315 }
316 + opts = append(opts, cmdhttp.ClientWithHTTPClient(&http.Client{
317 + Transport: otelhttp.NewTransport(tpt,
318 + otelhttp.WithPropagators(tracing.Propagator()),
319 + ),
320 + }))
321 +
322 + return tracingWrappedExecutor{cmdhttp.NewClient(host, opts...)}, nil
323 +}
324
313 - return cmdhttp.NewClient(host, opts...), nil
325 +type tracingWrappedExecutor struct {
326 + exec cmds.Executor
327 +}
328 +
329 +func (twe tracingWrappedExecutor) Execute(req *cmds.Request, re cmds.ResponseEmitter, env cmds.Environment) error {
330 + ctx, span := tracer.Start(req.Context, "cmds."+strings.Join(req.Path, "."), trace.WithAttributes(attribute.StringSlice("Arguments", req.Arguments)))
331 + defer span.End()
332 + req.Context = ctx
333 +
334 + err := twe.exec.Execute(req, re, env)
335 + if err != nil {
336 + span.SetStatus(codes.Error, err.Error())
337 + }
338 + return err
339 }
340
341 func getRepoPath(req *cmds.Request) (string, error) {
core/corehttp/commands.go
+7 -5
@@ -9,15 +9,16 @@ import (
9 "strconv"
10 "strings"
11
12 - version "github.com/ipfs/kubo"
13 - oldcmds "github.com/ipfs/kubo/commands"
14 - "github.com/ipfs/kubo/core"
15 - corecommands "github.com/ipfs/kubo/core/commands"
16 -
12 cmds "github.com/ipfs/go-ipfs-cmds"
13 cmdsHttp "github.com/ipfs/go-ipfs-cmds/http"
14 path "github.com/ipfs/go-path"
15 + version "github.com/ipfs/kubo"
16 + oldcmds "github.com/ipfs/kubo/commands"
17 config "github.com/ipfs/kubo/config"
18 + "github.com/ipfs/kubo/core"
19 + corecommands "github.com/ipfs/kubo/core/commands"
20 + "github.com/ipfs/kubo/tracing"
21 + "go.opentelemetry.io/contrib/instrumentation/net/http/otelhttp"
22 )
23
24 var (
@@ -146,6 +147,7 @@ func commandsOption(cctx oldcmds.Context, command *cmds.Command, allowGet bool)
147 patchCORSVars(cfg, l.Addr())
148
149 cmdHandler := cmdsHttp.NewHandler(&cctx, command, cfg)
150 + cmdHandler = otelhttp.NewHandler(cmdHandler, "corehttp.cmdsHandler", otelhttp.WithPropagators(tracing.Propagator()))
151 mux.Handle(APIPath+"/", cmdHandler)
152 return mux, nil
153 }
core/corehttp/gateway.go
+2
@@ -48,6 +48,8 @@ func GatewayOption(paths ...string) ServeOption {
48 }
49
50 gw := gateway.NewHandler(gwConfig, gwAPI)
51 + // TODO: Add otelhttp.WithPropagators(tracing.Propagator()) option to
52 + // propagate traces through the gateway once we test this feature.
53 gw = otelhttp.NewHandler(gw, "Gateway.Request")
54
55 // By default, our HTTP handler is the gateway handler.
test/sharness/t0230-channel-streaming-http-content-type.sh
+2
@@ -23,6 +23,7 @@ test_ls_cmd() {
23 printf "HTTP/1.1 200 OK\r\n" >expected_output &&
24 printf "Access-Control-Allow-Headers: X-Stream-Output, X-Chunked-Output, X-Content-Length\r\n" >>expected_output &&
25 printf "Access-Control-Expose-Headers: X-Stream-Output, X-Chunked-Output, X-Content-Length\r\n" >>expected_output &&
26 + printf "Connection: close\r\n" >>expected_output &&
27 printf "Content-Type: text/plain\r\n" >>expected_output &&
28 printf "Server: kubo/%s\r\n" $(ipfs version -n) >>expected_output &&
29 printf "Trailer: X-Stream-Error\r\n" >>expected_output &&
@@ -46,6 +47,7 @@ test_ls_cmd() {
47 printf "HTTP/1.1 200 OK\r\n" >expected_output &&
48 printf "Access-Control-Allow-Headers: X-Stream-Output, X-Chunked-Output, X-Content-Length\r\n" >>expected_output &&
49 printf "Access-Control-Expose-Headers: X-Stream-Output, X-Chunked-Output, X-Content-Length\r\n" >>expected_output &&
50 + printf "Connection: close\r\n" >>expected_output &&
51 printf "Content-Type: application/json\r\n" >>expected_output &&
52 printf "Server: kubo/%s\r\n" $(ipfs version -n) >>expected_output &&
53 printf "Trailer: X-Stream-Error\r\n" >>expected_output &&
tracing/tracing.go
+7 -2
@@ -13,6 +13,7 @@ import (
13 "go.opentelemetry.io/otel/exporters/otlp/otlptrace/otlptracegrpc"
14 "go.opentelemetry.io/otel/exporters/otlp/otlptrace/otlptracehttp"
15 "go.opentelemetry.io/otel/exporters/zipkin"
16 + "go.opentelemetry.io/otel/propagation"
17 "go.opentelemetry.io/otel/sdk/resource"
18 "go.opentelemetry.io/otel/sdk/trace"
19 semconv "go.opentelemetry.io/otel/semconv/v1.4.0"
@@ -130,7 +131,7 @@ func NewTracerProvider(ctx context.Context) (shutdownTracerProvider, error) {
131 r, err := resource.Merge(
132 resource.Default(),
133 resource.NewSchemaless(
133 - semconv.ServiceNameKey.String("go-ipfs"),
134 + semconv.ServiceNameKey.String("Kubo"),
135 semconv.ServiceVersionKey.String(version.CurrentVersionNumber),
136 ),
137 )
@@ -144,5 +145,9 @@ func NewTracerProvider(ctx context.Context) (shutdownTracerProvider, error) {
145
146 // Span starts a new span using the standard IPFS tracing conventions.
147 func Span(ctx context.Context, componentName string, spanName string, opts ...traceapi.SpanStartOption) (context.Context, traceapi.Span) {
147 - return otel.Tracer("go-ipfs").Start(ctx, fmt.Sprintf("%s.%s", componentName, spanName), opts...)
148 + return otel.Tracer("Kubo").Start(ctx, fmt.Sprintf("%s.%s", componentName, spanName), opts...)
149 +}
150 +
151 +func Propagator() propagation.TextMapPropagator {
152 + return propagation.NewCompositeTextMapPropagator(propagation.TraceContext{}, propagation.Baggage{})
153 }