MCPcopy Create free account
hub / github.com/dagger/dagger / TracesSubscribeHandler

Method TracesSubscribeHandler

engine/server/telemetry.go:197–234  ·  view source on GitHub ↗
(w http.ResponseWriter, r *http.Request, client *daggerClient)

Source from the content-addressed store, hash-verified

195const otlpBatchSize = 1000
196
197func (ps *PubSub) TracesSubscribeHandler(w http.ResponseWriter, r *http.Request, client *daggerClient) error {
198 return ps.sseHandler(w, r, client, func(ctx context.Context, db *clientdb.DB, lastID string) (*sse.Event, bool, error) {
199 var since int64
200 if lastID != "" {
201 _, err := fmt.Sscanf(lastID, "%d", &since)
202 if err != nil {
203 return nil, false, fmt.Errorf("invalid last ID: %w", err)
204 }
205 }
206 spans, err := db.SelectSpansSince(ctx, clientdb.SelectSpansSinceParams{
207 ID: since,
208 Limit: otlpBatchSize,
209 })
210 if err != nil {
211 return nil, false, fmt.Errorf("select spans: %w", err)
212 }
213 if len(spans) == 0 {
214 return nil, false, nil
215 }
216 roSpans := make([]sdktrace.ReadOnlySpan, len(spans))
217 for i, span := range spans {
218 roSpans[i] = span.ReadOnly()
219 since = span.ID
220 }
221 // Marshal the spans to OTLP.
222 payload, err := protojson.Marshal(&coltracepb.ExportTraceServiceRequest{
223 ResourceSpans: telemetry.SpansToPB(roSpans),
224 })
225 if err != nil {
226 return nil, false, fmt.Errorf("marshal spans: %w", err)
227 }
228 return &sse.Event{
229 Name: "spans",
230 ID: fmt.Sprintf("%d", since),
231 Data: payload,
232 }, true, nil
233 })
234}
235
236//nolint:dupl
237func (ps *PubSub) LogsSubscribeHandler(w http.ResponseWriter, r *http.Request, client *daggerClient) error {

Callers

nothing calls this directly

Calls 4

sseHandlerMethod · 0.95
SelectSpansSinceMethod · 0.80
ReadOnlyMethod · 0.80
MarshalMethod · 0.65

Tested by

no test coverage detected