nolint:dupl
(w http.ResponseWriter, r *http.Request, client *daggerClient)
| 271 | |
| 272 | //nolint:dupl |
| 273 | func (ps *PubSub) MetricsSubscribeHandler(w http.ResponseWriter, r *http.Request, client *daggerClient) error { |
| 274 | return ps.sseHandler(w, r, client, func(ctx context.Context, db *clientdb.DB, lastID string) (*sse.Event, bool, error) { |
| 275 | var since int64 |
| 276 | if lastID != "" { |
| 277 | _, err := fmt.Sscanf(lastID, "%d", &since) |
| 278 | if err != nil { |
| 279 | return nil, false, fmt.Errorf("invalid last ID: %w", err) |
| 280 | } |
| 281 | } |
| 282 | metrics, err := db.SelectMetricsSince(ctx, clientdb.SelectMetricsSinceParams{ |
| 283 | ID: since, |
| 284 | Limit: otlpBatchSize, |
| 285 | }) |
| 286 | if err != nil { |
| 287 | return nil, false, fmt.Errorf("select metrics: %w", err) |
| 288 | } |
| 289 | |
| 290 | if len(metrics) == 0 { |
| 291 | return nil, false, nil |
| 292 | } |
| 293 | since = metrics[len(metrics)-1].ID |
| 294 | // Marshal the metrics to OTLP. |
| 295 | payload, err := protojson.Marshal(&colmetricspb.ExportMetricsServiceRequest{ |
| 296 | ResourceMetrics: clientdb.MetricsToPB(metrics), |
| 297 | }) |
| 298 | if err != nil { |
| 299 | return nil, false, fmt.Errorf("marshal metrics: %w", err) |
| 300 | } |
| 301 | return &sse.Event{ |
| 302 | Name: "metrics", |
| 303 | ID: fmt.Sprintf("%d", since), |
| 304 | Data: payload, |
| 305 | }, true, nil |
| 306 | }) |
| 307 | } |
| 308 | |
| 309 | type clientSpans struct { |
| 310 | *PubSub |
nothing calls this directly
no test coverage detected