Skip to content
Closed
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
14 changes: 10 additions & 4 deletions pkg/epp/handlers/response.go
Original file line number Diff line number Diff line change
Expand Up @@ -86,11 +86,17 @@ func (s *StreamingServer) HandleResponseBody(ctx context.Context, reqCtx *Reques
}
}
if endOfStream {
metrics.RecordNormalizedTimePerOutputToken(ctx, reqCtx.IncomingModelName, reqCtx.TargetModelName, fairnessID, priority, reqCtx.RequestReceivedTimestamp, reqCtx.ResponseCompleteTimestamp, reqCtx.Usage.CompletionTokens)
metrics.RecordRequestLatencies(ctx, reqCtx.IncomingModelName, reqCtx.TargetModelName, fairnessID, priority, reqCtx.RequestReceivedTimestamp, reqCtx.ResponseCompleteTimestamp)
if !reqCtx.RequestReceivedTimestamp.IsZero() && !reqCtx.ResponseCompleteTimestamp.IsZero() {
metrics.RecordNormalizedTimePerOutputToken(ctx, reqCtx.IncomingModelName, reqCtx.TargetModelName, fairnessID, priority, reqCtx.RequestReceivedTimestamp, reqCtx.ResponseCompleteTimestamp, reqCtx.Usage.CompletionTokens)
metrics.RecordRequestLatencies(ctx, reqCtx.IncomingModelName, reqCtx.TargetModelName, fairnessID, priority, reqCtx.RequestReceivedTimestamp, reqCtx.ResponseCompleteTimestamp)
}
metrics.RecordResponseSizes(reqCtx.IncomingModelName, reqCtx.TargetModelName, fairnessID, priority, reqCtx.ResponseSize)
metrics.RecordRequestTTFT(ctx, reqCtx.IncomingModelName, reqCtx.TargetModelName, fairnessID, priority, reqCtx.modelServerStreaming, reqCtx.RequestReceivedTimestamp, reqCtx.FirstTokenTimestamp)
metrics.RecordRequestTPOT(ctx, reqCtx.IncomingModelName, reqCtx.TargetModelName, fairnessID, priority, reqCtx.RequestReceivedTimestamp, reqCtx.FirstTokenTimestamp, reqCtx.ResponseCompleteTimestamp, reqCtx.Usage.CompletionTokens)
if !reqCtx.RequestReceivedTimestamp.IsZero() && !reqCtx.FirstTokenTimestamp.IsZero() {
metrics.RecordRequestTTFT(ctx, reqCtx.IncomingModelName, reqCtx.TargetModelName, fairnessID, priority, reqCtx.modelServerStreaming, reqCtx.RequestReceivedTimestamp, reqCtx.FirstTokenTimestamp)
}
if !reqCtx.RequestReceivedTimestamp.IsZero() && !reqCtx.FirstTokenTimestamp.IsZero() && !reqCtx.ResponseCompleteTimestamp.IsZero() {
metrics.RecordRequestTPOT(ctx, reqCtx.IncomingModelName, reqCtx.TargetModelName, fairnessID, priority, reqCtx.RequestReceivedTimestamp, reqCtx.FirstTokenTimestamp, reqCtx.ResponseCompleteTimestamp, reqCtx.Usage.CompletionTokens)
}
}
return s.director.HandleResponseBody(ctx, reqCtx, endOfStream)
}
Expand Down
57 changes: 57 additions & 0 deletions pkg/epp/handlers/response_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -309,6 +309,63 @@ func TestHandleResponseBodyWithoutSchedulingRequest(t *testing.T) {
require.Equal(t, uint64(1), histogram.GetSampleCount())
}

// TestHandleResponseBodyZeroTimestamps verifies that end-of-stream handling with zero
// timestamps does not emit latency metrics or produce error log entries.
func TestHandleResponseBodyZeroTimestamps(t *testing.T) {
ctx := logutil.NewTestLoggerIntoContext(context.Background())
eppmetrics.Register()
eppmetrics.Reset()
t.Cleanup(eppmetrics.Reset)

server := &StreamingServer{
parserRegistry: NewParserRegistry([]fwkrh.Parser{openai.NewOpenAIParser()}, logr.Discard()),
}
server.director = &mockDirector{}
// RequestReceivedTimestamp, FirstTokenTimestamp, and ResponseCompleteTimestamp are all zero.
reqCtx := &RequestContext{
IncomingModelName: "incoming-model",
TargetModelName: "target-model",
Response: &Response{
Headers: map[string]string{},
},
SchedulingRequest: &fwksched.InferenceRequest{FairnessID: metadata.DefaultFairnessID},
}

require.NotPanics(t, func() {
server.HandleResponseBody(ctx, reqCtx, []byte("data"), true)
})

// Latency metrics must not be recorded when timestamps are absent.
families, err := ctrlmetrics.Registry.Gather()
require.NoError(t, err)
latencyMetrics := []string{
"llm_d_epp_request_ntpot_seconds",
"llm_d_epp_request_duration_seconds",
"llm_d_epp_request_ttft_seconds",
"llm_d_epp_request_streaming_tpot_seconds",
}
for _, family := range families {
for _, name := range latencyMetrics {
if family.GetName() == name {
for _, metric := range family.GetMetric() {
if metricHasLabels(metric, map[string]string{"model_name": "incoming-model"}) {
t.Errorf("latency metric %q must not be recorded with zero timestamps, got sample count %d",
name, metric.GetHistogram().GetSampleCount())
}
}
}
}
}

// ResponseSize should still be recorded.
sizeHistogram := findHistogramMetric(t, "llm_d_epp_response_size_bytes", map[string]string{
"model_name": "incoming-model",
"target_model_name": "target-model",
"fairness_id": metadata.DefaultFairnessID,
})
require.Equal(t, uint64(1), sizeHistogram.GetSampleCount())
}

func findHistogramMetric(t *testing.T, name string, labels map[string]string) *dto.Histogram {
t.Helper()

Expand Down
Loading