Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
11 changes: 11 additions & 0 deletions pulsar-function-go/pf/context.go
Original file line number Diff line number Diff line change
Expand Up @@ -202,6 +202,17 @@ func (c *FunctionContext) RecordMetric(metricName string, metricValue float64) {
v.(prometheus.Observer).Observe(metricValue)
}

// GetMetricsRegistry returns the Prometheus registry that backs the function's
// metrics endpoint, so a function can register additional collectors (counters,
// gauges, histograms, ...) that are then exposed on the same endpoint alongside
// the SDK's own metrics. Metric names prefixed with "pulsar_function_" are
// reserved for SDK metrics. A collector whose fully-qualified name collides with
// an already-registered metric fails to register: Register returns an error,
// while MustRegister panics.
func (c *FunctionContext) GetMetricsRegistry() prometheus.Registerer {
return reg
}

// An unexported type to be used as the key for types in this package. This
// prevents collisions with keys defined in other packages.
type key struct{}
Expand Down
49 changes: 49 additions & 0 deletions pulsar-function-go/pf/stats_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -216,3 +216,52 @@ func TestInstanceControlMetrics(t *testing.T) {
assert.EqualValuesf(t, value+1, metrics.UserMetrics[label], "user metric %s != %d", label, value+1)
}
}

func TestGetMetricsRegistry_CustomCollector(t *testing.T) {
gi := newGoInstance()
metricsServicer := NewMetricsServicer(gi)
metricsServicer.serve()

// A Counter is something the user_metric Summary cannot express.
customCounter := prometheus.NewCounter(prometheus.CounterOpts{
Name: "pulsar_function_go_test_custom_counter_total",
Help: "A user-registered counter exposed via FunctionContext.GetMetricsRegistry.",
})
err := gi.context.GetMetricsRegistry().Register(customCounter)
assert.NoError(t, err)
defer gi.context.GetMetricsRegistry().Unregister(customCounter)
customCounter.Add(42)

time.Sleep(time.Second * 1)
resp, err := http.Get(fmt.Sprintf("http://localhost:%d/metrics", gi.context.GetMetricsPort()))
assert.Equal(t, nil, err)
assert.NotEqual(t, nil, resp)
assert.Equal(t, 200, resp.StatusCode)
body, err := io.ReadAll(resp.Body)
assert.Equal(t, nil, err)
assert.Containsf(t, string(body), "\npulsar_function_go_test_custom_counter_total 42\n",
"custom collector should be exposed on /metrics")
resp.Body.Close()

gi.close()
metricsServicer.close()
}

func TestGetMetricsRegistry_NameCollision(t *testing.T) {
gi := newGoInstance()

// Collides with the built-in received_total gauge registered on reg in init().
colliding := prometheus.NewCounter(prometheus.CounterOpts{
Name: PulsarFunctionMetricsPrefix + TotalReceived,
Help: "collector whose name collides with a built-in SDK metric",
})

// Register is the documented safe path: it returns an error, it does not panic.
err := gi.context.GetMetricsRegistry().Register(colliding)
assert.Error(t, err)

// MustRegister panics on the same collision.
assert.Panics(t, func() {
gi.context.GetMetricsRegistry().MustRegister(colliding)
})
}