(w http.ResponseWriter, r *http.Request, client *daggerClient)
| 195 | const otlpBatchSize = 1000 |
| 196 | |
| 197 | func (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 |
| 237 | func (ps *PubSub) LogsSubscribeHandler(w http.ResponseWriter, r *http.Request, client *daggerClient) error { |
nothing calls this directly
no test coverage detected