Skip to content

Commit 227f6ca

Browse files
authored
refactor(epp): unexport RequestContext protocol and response state (llm-d#2457)
The stream state machine, drop reason, and response-phase bookkeeping fields have no users outside the handlers package. Unexporting them gives the ext_proc protocol state a compile-time boundary, making the struct split proposed in llm-d#1189 unnecessary. Also deletes the request trailer response field and constants, which were never assigned. Signed-off-by: Luke Van Drie <lukevandrie@google.com>
1 parent 229b8ea commit 227f6ca

4 files changed

Lines changed: 104 additions & 111 deletions

File tree

pkg/epp/handlers/response.go

Lines changed: 11 additions & 11 deletions
Original file line numberDiff line numberDiff line change
@@ -50,19 +50,19 @@ func (s *StreamingServer) HandleResponseBody(ctx context.Context, reqCtx *Reques
5050

5151
fairnessID, priority := extractFairnessAndPriority(reqCtx)
5252

53-
reqCtx.ResponseSize += len(responseBytes)
53+
reqCtx.responseSize += len(responseBytes)
5454

55-
if reqCtx.FirstTokenTimestamp.IsZero() && len(responseBytes) > 0 {
56-
reqCtx.FirstTokenTimestamp = time.Now()
55+
if reqCtx.firstTokenTimestamp.IsZero() && len(responseBytes) > 0 {
56+
reqCtx.firstTokenTimestamp = time.Now()
5757
}
5858

5959
if reqCtx.modelServerStreaming && len(responseBytes) > 0 {
6060
now := time.Now()
61-
if !reqCtx.LastChunkReceivedTimestamp.IsZero() {
62-
itl := now.Sub(reqCtx.LastChunkReceivedTimestamp).Seconds()
61+
if !reqCtx.lastChunkReceivedTimestamp.IsZero() {
62+
itl := now.Sub(reqCtx.lastChunkReceivedTimestamp).Seconds()
6363
metrics.RecordInterTokenLatency(ctx, reqCtx.IncomingModelName, reqCtx.TargetModelName, fairnessID, priority, itl)
6464
}
65-
reqCtx.LastChunkReceivedTimestamp = now
65+
reqCtx.lastChunkReceivedTimestamp = now
6666
}
6767

6868
var parsedResp *fwkrh.ParsedResponse
@@ -91,11 +91,11 @@ func (s *StreamingServer) HandleResponseBody(ctx context.Context, reqCtx *Reques
9191
}
9292
}
9393
if endOfStream {
94-
metrics.RecordNormalizedTimePerOutputToken(ctx, reqCtx.IncomingModelName, reqCtx.TargetModelName, fairnessID, priority, reqCtx.RequestReceivedTimestamp, reqCtx.ResponseCompleteTimestamp, reqCtx.Usage.CompletionTokens)
95-
metrics.RecordRequestLatencies(ctx, reqCtx.IncomingModelName, reqCtx.TargetModelName, fairnessID, priority, reqCtx.RequestReceivedTimestamp, reqCtx.ResponseCompleteTimestamp)
96-
metrics.RecordResponseSizes(reqCtx.IncomingModelName, reqCtx.TargetModelName, fairnessID, priority, reqCtx.ResponseSize)
97-
metrics.RecordRequestTTFT(ctx, reqCtx.IncomingModelName, reqCtx.TargetModelName, fairnessID, priority, reqCtx.modelServerStreaming, reqCtx.RequestReceivedTimestamp, reqCtx.FirstTokenTimestamp)
98-
metrics.RecordRequestTPOT(ctx, reqCtx.IncomingModelName, reqCtx.TargetModelName, fairnessID, priority, reqCtx.modelServerStreaming, reqCtx.RequestReceivedTimestamp, reqCtx.FirstTokenTimestamp, reqCtx.ResponseCompleteTimestamp, reqCtx.Usage.CompletionTokens)
94+
metrics.RecordNormalizedTimePerOutputToken(ctx, reqCtx.IncomingModelName, reqCtx.TargetModelName, fairnessID, priority, reqCtx.RequestReceivedTimestamp, reqCtx.responseCompleteTimestamp, reqCtx.Usage.CompletionTokens)
95+
metrics.RecordRequestLatencies(ctx, reqCtx.IncomingModelName, reqCtx.TargetModelName, fairnessID, priority, reqCtx.RequestReceivedTimestamp, reqCtx.responseCompleteTimestamp)
96+
metrics.RecordResponseSizes(reqCtx.IncomingModelName, reqCtx.TargetModelName, fairnessID, priority, reqCtx.responseSize)
97+
metrics.RecordRequestTTFT(ctx, reqCtx.IncomingModelName, reqCtx.TargetModelName, fairnessID, priority, reqCtx.modelServerStreaming, reqCtx.RequestReceivedTimestamp, reqCtx.firstTokenTimestamp)
98+
metrics.RecordRequestTPOT(ctx, reqCtx.IncomingModelName, reqCtx.TargetModelName, fairnessID, priority, reqCtx.modelServerStreaming, reqCtx.RequestReceivedTimestamp, reqCtx.firstTokenTimestamp, reqCtx.responseCompleteTimestamp, reqCtx.Usage.CompletionTokens)
9999
}
100100
return s.director.HandleResponseBody(ctx, reqCtx, endOfStream)
101101
}

pkg/epp/handlers/response_test.go

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -291,7 +291,7 @@ func TestHandleResponseBodyWithoutSchedulingRequest(t *testing.T) {
291291
TargetModelName: "target-model",
292292
Priority: 3,
293293
RequestReceivedTimestamp: timeBaseline,
294-
ResponseCompleteTimestamp: timeBaseline.Add(time.Second),
294+
responseCompleteTimestamp: timeBaseline.Add(time.Second),
295295
Response: &Response{
296296
Headers: map[string]string{},
297297
},
@@ -642,7 +642,7 @@ func TestResponseSizeAccumulation(t *testing.T) {
642642
endOfStream := i == len(tt.chunks)-1
643643
server.HandleResponseBody(ctx, reqCtx, chunk, endOfStream)
644644
}
645-
assert.Equal(t, tt.wantResponseSize, reqCtx.ResponseSize)
645+
assert.Equal(t, tt.wantResponseSize, reqCtx.responseSize)
646646
})
647647
}
648648
}

0 commit comments

Comments
 (0)