From 4813189cbe57c11e924fa196beebc52e05f2c733 Mon Sep 17 00:00:00 2001 From: Sagnik Ghosh Date: Thu, 23 Jul 2026 16:51:09 +0530 Subject: [PATCH] feat(ingester): forward metric rollups to the sink webhook The sink previously received only fingerprinted occurrences, so the cloud dashboard had error counts but no request denominator. Include the 1-minute usage rollup points in the same sink payload (new optional `metrics` field, version unchanged) and forward from /v1/metrics too, which previously never reached the sink at all. Generated-By: PostHog Code Task-Id: 6e0925be-ecd5-42ee-9fde-895c4678bde4 --- packages/otlp-ingester/src/server.ts | 24 +++++++++++++++++++----- 1 file changed, 19 insertions(+), 5 deletions(-) diff --git a/packages/otlp-ingester/src/server.ts b/packages/otlp-ingester/src/server.ts index 7d92335..e14947a 100644 --- a/packages/otlp-ingester/src/server.ts +++ b/packages/otlp-ingester/src/server.ts @@ -21,6 +21,7 @@ import { import { decodeMetricsRequest, decodeTraceRequest } from "./otlp-proto.js"; import type { IngestContext, + RuntimeMetricPoint, RuntimeOccurrence, RuntimeOccurrenceInput, } from "./types.js"; @@ -157,9 +158,17 @@ export function createIngesterApp(config: IngesterConfig): IngesterApp { })); } - /** Best-effort forward of fingerprinted occurrences for issue grouping. */ - function forwardToSink(ctx: IngestContext, occurrences: RuntimeOccurrence[]) { - if (!config.sinkUrl || occurrences.length === 0) return; + /** + * Best-effort forward for the cloud dashboard: fingerprinted occurrences + * feed issue grouping, metric points feed the request/error-rate rollups. + */ + function forwardToSink( + ctx: IngestContext, + occurrences: RuntimeOccurrence[], + metricPoints: RuntimeMetricPoint[] = [], + ) { + if (!config.sinkUrl) return; + if (occurrences.length === 0 && metricPoints.length === 0) return; void fetch(config.sinkUrl, { method: "POST", headers: { @@ -176,6 +185,10 @@ export function createIngesterApp(config: IngesterConfig): IngesterApp { ...o, occurredAt: o.occurredAt.toISOString(), })), + metrics: metricPoints.map((p) => ({ + ...p, + bucketAt: p.bucketAt.toISOString(), + })), }), signal: AbortSignal.timeout(10_000), }).catch((err) => { @@ -224,7 +237,7 @@ export function createIngesterApp(config: IngesterConfig): IngesterApp { storageError(res, err); return; } - forwardToSink(ctx, fingerprinted); + forwardToSink(ctx, fingerprinted, metricPoints); otlpSuccess(req, res); }); @@ -249,6 +262,7 @@ export function createIngesterApp(config: IngesterConfig): IngesterApp { storageError(res, err); return; } + forwardToSink(ctx, [], metricPoints); otlpSuccess(req, res); }); @@ -284,7 +298,7 @@ export function createIngesterApp(config: IngesterConfig): IngesterApp { storageError(res, err); return; } - forwardToSink(ctx, fingerprinted); + forwardToSink(ctx, fingerprinted, metricPoints); res.status(202).json({ accepted: fingerprinted.length }); });