From 033cf83d5914f700d4c4538444d4a8e4e4c974c1 Mon Sep 17 00:00:00 2001 From: Dorin Geman Date: Fri, 12 Sep 2025 16:28:35 +0300 Subject: [PATCH 1/7] feat: add streaming endpoint for inference requests Signed-off-by: Dorin Geman --- pkg/inference/scheduling/scheduler.go | 1 + pkg/metrics/openai_recorder.go | 110 ++++++++++++++++++++++++++ 2 files changed, 111 insertions(+) diff --git a/pkg/inference/scheduling/scheduler.go b/pkg/inference/scheduling/scheduler.go index b6d7ac089..377f16a98 100644 --- a/pkg/inference/scheduling/scheduler.go +++ b/pkg/inference/scheduling/scheduler.go @@ -118,6 +118,7 @@ func (s *Scheduler) routeHandlers() map[string]http.HandlerFunc { m["POST "+inference.InferencePrefix+"/{backend}/_configure"] = s.Configure m["POST "+inference.InferencePrefix+"/_configure"] = s.Configure m["GET "+inference.InferencePrefix+"/requests"] = s.openAIRecorder.GetRecordsHandler() + m["GET "+inference.InferencePrefix+"/requests/stream"] = s.openAIRecorder.StreamRequestsHandler() return m } diff --git a/pkg/metrics/openai_recorder.go b/pkg/metrics/openai_recorder.go index f7ac8c3ce..00ec7ce25 100644 --- a/pkg/metrics/openai_recorder.go +++ b/pkg/metrics/openai_recorder.go @@ -69,6 +69,10 @@ type OpenAIRecorder struct { records map[string]*ModelData // key is model ID modelManager *models.Manager // for resolving model tags to IDs m sync.RWMutex + + // streaming + subscribers map[string]chan *RequestResponsePair + subMutex sync.RWMutex } func NewOpenAIRecorder(log logging.Logger, modelManager *models.Manager) *OpenAIRecorder { @@ -76,6 +80,7 @@ func NewOpenAIRecorder(log logging.Logger, modelManager *models.Manager) *OpenAI log: log, modelManager: modelManager, records: make(map[string]*ModelData), + subscribers: make(map[string]chan *RequestResponsePair), } } @@ -188,6 +193,7 @@ func (r *OpenAIRecorder) RecordResponse(id, model string, rw http.ResponseWriter record.Response = response record.Error = "" // Ensure Error is empty for successful responses } + go r.broadcastToSubscribers(record) return } } @@ -352,6 +358,110 @@ func (r *OpenAIRecorder) getRecordsByModel(model string) []ModelRecordsResponse return nil } +func (r *OpenAIRecorder) broadcastToSubscribers(record *RequestResponsePair) { + r.subMutex.RLock() + defer r.subMutex.RUnlock() + + for _, ch := range r.subscribers { + select { + case ch <- record: + default: + // The channel is full, skip this subscriber. + } + } +} + +func (r *OpenAIRecorder) StreamRequestsHandler() http.HandlerFunc { + return func(w http.ResponseWriter, req *http.Request) { + // Set SSE headers. + w.Header().Set("Content-Type", "text/event-stream") + w.Header().Set("Cache-Control", "no-cache") + w.Header().Set("Connection", "keep-alive") + + // Create subscriber channel. + subscriberID := fmt.Sprintf("sub_%d", time.Now().UnixNano()) + ch := make(chan *RequestResponsePair, 100) + + // Register subscriber. + r.subMutex.Lock() + r.subscribers[subscriberID] = ch + r.subMutex.Unlock() + + // Clean up on disconnect. + defer func() { + r.subMutex.Lock() + delete(r.subscribers, subscriberID) + close(ch) + r.subMutex.Unlock() + }() + + // Optional: Send existing records first. + model := req.URL.Query().Get("model") + if includeExisting := req.URL.Query().Get("include_existing"); includeExisting == "true" { + r.sendExistingRecords(w, model) + } + + flusher, ok := w.(http.Flusher) + if !ok { + http.Error(w, "Streaming not supported", http.StatusInternalServerError) + return + } + + // Send heartbeat to establish connection. + fmt.Fprintf(w, "event: connected\ndata: {\"status\": \"connected\"}\n\n") + flusher.Flush() + + for { + select { + case record, ok := <-ch: + if !ok { + return + } + + // Filter by model if specified. + if model != "" && record.Model != model { + continue + } + + // Send as SSE event. + jsonData, err := json.Marshal(record) + if err != nil { + continue + } + + fmt.Fprintf(w, "event: new_request\ndata: %s\n\n", jsonData) + flusher.Flush() + + case <-req.Context().Done(): + // Client disconnected. + return + } + } + } +} + +func (r *OpenAIRecorder) sendExistingRecords(w http.ResponseWriter, model string) { + var records []ModelRecordsResponse + + if model == "" { + records = r.getAllRecords() + } else { + records = r.getRecordsByModel(model) + } + + if records != nil { + for _, modelResponse := range records { + for _, record := range modelResponse.Records { + jsonData, err := json.Marshal(record) + if err != nil { + continue + } + fmt.Fprintf(w, "event: existing_request\ndata: %s\n\n", jsonData) + } + } + } +} + func (r *OpenAIRecorder) RemoveModel(model string) { modelID := r.modelManager.ResolveModelID(model) From a0ddf2e5f18610f8c5debe5a2f67b19740b57b26 Mon Sep 17 00:00:00 2001 From: Dorin Geman Date: Mon, 15 Sep 2025 16:09:33 +0300 Subject: [PATCH 2/7] refactor: combine the requests endpoints Differentiate regular and streaming based on the Accept Header. Signed-off-by: Dorin Geman --- pkg/inference/scheduling/scheduler.go | 1 - pkg/metrics/openai_recorder.go | 197 ++++++++++++++------------ 2 files changed, 104 insertions(+), 94 deletions(-) diff --git a/pkg/inference/scheduling/scheduler.go b/pkg/inference/scheduling/scheduler.go index 377f16a98..b6d7ac089 100644 --- a/pkg/inference/scheduling/scheduler.go +++ b/pkg/inference/scheduling/scheduler.go @@ -118,7 +118,6 @@ func (s *Scheduler) routeHandlers() map[string]http.HandlerFunc { m["POST "+inference.InferencePrefix+"/{backend}/_configure"] = s.Configure m["POST "+inference.InferencePrefix+"/_configure"] = s.Configure m["GET "+inference.InferencePrefix+"/requests"] = s.openAIRecorder.GetRecordsHandler() - m["GET "+inference.InferencePrefix+"/requests/stream"] = s.openAIRecorder.StreamRequestsHandler() return m } diff --git a/pkg/metrics/openai_recorder.go b/pkg/metrics/openai_recorder.go index 00ec7ce25..251b2409b 100644 --- a/pkg/metrics/openai_recorder.go +++ b/pkg/metrics/openai_recorder.go @@ -280,36 +280,116 @@ func (r *OpenAIRecorder) convertStreamingResponse(streamingBody string) string { func (r *OpenAIRecorder) GetRecordsHandler() http.HandlerFunc { return func(w http.ResponseWriter, req *http.Request) { - w.Header().Set("Content-Type", "application/json") + acceptHeader := req.Header.Get("Accept") - model := req.URL.Query().Get("model") + // Check if client wants Server-Sent Events + if acceptHeader == "text/event-stream" { + r.handleStreamingRequests(w, req) + return + } - if model == "" { - // Retrieve all records for all models. - allRecords := r.getAllRecords() - if allRecords == nil { - // No records found. - http.Error(w, "No records found", http.StatusNotFound) - return - } - if err := json.NewEncoder(w).Encode(allRecords); err != nil { - http.Error(w, fmt.Sprintf("Failed to encode all records: %v", err), - http.StatusInternalServerError) + // Default to JSON response + r.handleJSONRequests(w, req) + } +} + +func (r *OpenAIRecorder) handleJSONRequests(w http.ResponseWriter, req *http.Request) { + w.Header().Set("Content-Type", "application/json") + + model := req.URL.Query().Get("model") + + if model == "" { + // Retrieve all records for all models. + allRecords := r.getAllRecords() + if allRecords == nil { + // No records found. + http.Error(w, "No records found", http.StatusNotFound) + return + } + if err := json.NewEncoder(w).Encode(allRecords); err != nil { + http.Error(w, fmt.Sprintf("Failed to encode all records: %v", err), + http.StatusInternalServerError) + return + } + } else { + // Retrieve records for the specified model. + records := r.getRecordsByModel(model) + if records == nil { + // No records found for the specified model. + http.Error(w, fmt.Sprintf("No records found for model '%s'", model), http.StatusNotFound) + return + } + if err := json.NewEncoder(w).Encode(records); err != nil { + http.Error(w, fmt.Sprintf("Failed to encode records for model '%s': %v", model, err), + http.StatusInternalServerError) + return + } + } +} + +func (r *OpenAIRecorder) handleStreamingRequests(w http.ResponseWriter, req *http.Request) { + // Set SSE headers. + w.Header().Set("Content-Type", "text/event-stream") + w.Header().Set("Cache-Control", "no-cache") + w.Header().Set("Connection", "keep-alive") + + // Create subscriber channel. + subscriberID := fmt.Sprintf("sub_%d", time.Now().UnixNano()) + ch := make(chan *RequestResponsePair, 100) + + // Register subscriber. + r.subMutex.Lock() + r.subscribers[subscriberID] = ch + r.subMutex.Unlock() + + // Clean up on disconnect. + defer func() { + r.subMutex.Lock() + delete(r.subscribers, subscriberID) + close(ch) + r.subMutex.Unlock() + }() + + // Optional: Send existing records first. + model := req.URL.Query().Get("model") + if includeExisting := req.URL.Query().Get("include_existing"); includeExisting == "true" { + r.sendExistingRecords(w, model) + } + + flusher, ok := w.(http.Flusher) + if !ok { + http.Error(w, "Streaming not supported", http.StatusInternalServerError) + return + } + + // Send heartbeat to establish connection. + fmt.Fprintf(w, "event: connected\ndata: {\"status\": \"connected\"}\n\n") + flusher.Flush() + + for { + select { + case record, ok := <-ch: + if !ok { return } - } else { - // Retrieve records for the specified model. - records := r.getRecordsByModel(model) - if records == nil { - // No records found for the specified model. - http.Error(w, fmt.Sprintf("No records found for model '%s'", model), http.StatusNotFound) - return + + // Filter by model if specified. + if model != "" && record.Model != model { + continue } - if err := json.NewEncoder(w).Encode(records); err != nil { - http.Error(w, fmt.Sprintf("Failed to encode records for model '%s': %v", model, err), - http.StatusInternalServerError) - return + + // Send as SSE event. + jsonData, err := json.Marshal(record) + if err != nil { + continue } + + fmt.Fprintf(w, "event: new_request\ndata: %s\n\n", jsonData) + flusher.Flush() + + case <-req.Context().Done(): + // Client disconnected. + return } } } @@ -371,75 +451,6 @@ func (r *OpenAIRecorder) broadcastToSubscribers(record *RequestResponsePair) { } } -func (r *OpenAIRecorder) StreamRequestsHandler() http.HandlerFunc { - return func(w http.ResponseWriter, req *http.Request) { - // Set SSE headers. - w.Header().Set("Content-Type", "text/event-stream") - w.Header().Set("Cache-Control", "no-cache") - w.Header().Set("Connection", "keep-alive") - - // Create subscriber channel. - subscriberID := fmt.Sprintf("sub_%d", time.Now().UnixNano()) - ch := make(chan *RequestResponsePair, 100) - - // Register subscriber. - r.subMutex.Lock() - r.subscribers[subscriberID] = ch - r.subMutex.Unlock() - - // Clean up on disconnect. - defer func() { - r.subMutex.Lock() - delete(r.subscribers, subscriberID) - close(ch) - r.subMutex.Unlock() - }() - - // Optional: Send existing records first. - model := req.URL.Query().Get("model") - if includeExisting := req.URL.Query().Get("include_existing"); includeExisting == "true" { - r.sendExistingRecords(w, model) - } - - flusher, ok := w.(http.Flusher) - if !ok { - http.Error(w, "Streaming not supported", http.StatusInternalServerError) - return - } - - // Send heartbeat to establish connection. - fmt.Fprintf(w, "event: connected\ndata: {\"status\": \"connected\"}\n\n") - flusher.Flush() - - for { - select { - case record, ok := <-ch: - if !ok { - return - } - - // Filter by model if specified. - if model != "" && record.Model != model { - continue - } - - // Send as SSE event. - jsonData, err := json.Marshal(record) - if err != nil { - continue - } - - fmt.Fprintf(w, "event: new_request\ndata: %s\n\n", jsonData) - flusher.Flush() - - case <-req.Context().Done(): - // Client disconnected. - return - } - } - } -} - func (r *OpenAIRecorder) sendExistingRecords(w http.ResponseWriter, model string) { var records []ModelRecordsResponse From d61cffbf0a0c652215879d10f9d383b57f0906af Mon Sep 17 00:00:00 2001 From: Dorin Geman Date: Tue, 16 Sep 2025 15:09:34 +0300 Subject: [PATCH 3/7] refactor: on streaming return the same struct as on regular Signed-off-by: Dorin Geman --- pkg/metrics/openai_recorder.go | 55 ++++++++++++++++++++++++---------- 1 file changed, 40 insertions(+), 15 deletions(-) diff --git a/pkg/metrics/openai_recorder.go b/pkg/metrics/openai_recorder.go index 251b2409b..58243108e 100644 --- a/pkg/metrics/openai_recorder.go +++ b/pkg/metrics/openai_recorder.go @@ -18,6 +18,9 @@ import ( // per model. const maximumRecordsPerModel = 10 +// subscriberChannelBuffer is the buffer size for subscriber channels. +const subscriberChannelBuffer = 100 + type responseRecorder struct { http.ResponseWriter body *bytes.Buffer @@ -71,7 +74,7 @@ type OpenAIRecorder struct { m sync.RWMutex // streaming - subscribers map[string]chan *RequestResponsePair + subscribers map[string]chan []ModelRecordsResponse subMutex sync.RWMutex } @@ -80,7 +83,7 @@ func NewOpenAIRecorder(log logging.Logger, modelManager *models.Manager) *OpenAI log: log, modelManager: modelManager, records: make(map[string]*ModelData), - subscribers: make(map[string]chan *RequestResponsePair), + subscribers: make(map[string]chan []ModelRecordsResponse), } } @@ -193,7 +196,18 @@ func (r *OpenAIRecorder) RecordResponse(id, model string, rw http.ResponseWriter record.Response = response record.Error = "" // Ensure Error is empty for successful responses } - go r.broadcastToSubscribers(record) + // Create ModelRecordsResponse with this single updated record to match + // what the non-streaming endpoint returns - []ModelRecordsResponse. + // See getAllRecords and getRecordsByModel. + modelResponse := []ModelRecordsResponse{{ + Count: 1, + Model: model, + ModelData: ModelData{ + Config: modelData.Config, + Records: []*RequestResponsePair{record}, + }, + }} + go r.broadcastToSubscribers(modelResponse) return } } @@ -335,7 +349,7 @@ func (r *OpenAIRecorder) handleStreamingRequests(w http.ResponseWriter, req *htt // Create subscriber channel. subscriberID := fmt.Sprintf("sub_%d", time.Now().UnixNano()) - ch := make(chan *RequestResponsePair, 100) + ch := make(chan []ModelRecordsResponse, subscriberChannelBuffer) // Register subscriber. r.subMutex.Lock() @@ -368,18 +382,18 @@ func (r *OpenAIRecorder) handleStreamingRequests(w http.ResponseWriter, req *htt for { select { - case record, ok := <-ch: + case modelRecords, ok := <-ch: if !ok { return } // Filter by model if specified. - if model != "" && record.Model != model { + if model != "" && len(modelRecords) > 0 && modelRecords[0].Model != model { continue } // Send as SSE event. - jsonData, err := json.Marshal(record) + jsonData, err := json.Marshal(modelRecords) if err != nil { continue } @@ -438,13 +452,13 @@ func (r *OpenAIRecorder) getRecordsByModel(model string) []ModelRecordsResponse return nil } -func (r *OpenAIRecorder) broadcastToSubscribers(record *RequestResponsePair) { +func (r *OpenAIRecorder) broadcastToSubscribers(modelResponses []ModelRecordsResponse) { r.subMutex.RLock() defer r.subMutex.RUnlock() for _, ch := range r.subscribers { select { - case ch <- record: + case ch <- modelResponses: default: // The channel is full, skip this subscriber. } @@ -461,13 +475,24 @@ func (r *OpenAIRecorder) sendExistingRecords(w http.ResponseWriter, model string } if records != nil { - for _, modelResponse := range records { - for _, record := range modelResponse.Records { - jsonData, err := json.Marshal(record) - if err != nil { - continue + // Send each individual request-response pair as a separate event. + for _, modelRecord := range records { + for _, requestRecord := range modelRecord.Records { + // Create a ModelRecordsResponse with a single record to match + // what the non-streaming endpoint returns - []ModelRecordsResponse. + // See getAllRecords and getRecordsByModel. + singleRecord := []ModelRecordsResponse{{ + Count: 1, + Model: modelRecord.Model, + ModelData: ModelData{ + Config: modelRecord.Config, + Records: []*RequestResponsePair{requestRecord}, + }, + }} + jsonData, err := json.Marshal(singleRecord) + if err == nil { + fmt.Fprintf(w, "event: existing_request\ndata: %s\n\n", jsonData) } - fmt.Fprintf(w, "event: existing_request\ndata: %s\n\n", jsonData) } } } From 98361c12e32d1251467ae3d1293de40933016c0c Mon Sep 17 00:00:00 2001 From: Dorin Geman Date: Tue, 16 Sep 2025 15:29:13 +0300 Subject: [PATCH 4/7] OpenAIRecorder: add clarifying comment about modelRecords size assumption Explains that modelRecords is expected to have size 1 due to how broadcastToSubscribers is called, avoiding the need for a second model config query. Signed-off-by: Dorin Geman --- pkg/metrics/openai_recorder.go | 2 ++ 1 file changed, 2 insertions(+) diff --git a/pkg/metrics/openai_recorder.go b/pkg/metrics/openai_recorder.go index 58243108e..7ca7567bd 100644 --- a/pkg/metrics/openai_recorder.go +++ b/pkg/metrics/openai_recorder.go @@ -388,6 +388,8 @@ func (r *OpenAIRecorder) handleStreamingRequests(w http.ResponseWriter, req *htt } // Filter by model if specified. + // modelRecords is assumed to have size 1 because that's how we call broadcastToSubscribers. + // We do this so we don't need to query a 2nd time for the model config. if model != "" && len(modelRecords) > 0 && modelRecords[0].Model != model { continue } From 78f30f96c43141b869b0083391e69b4d74df03f7 Mon Sep 17 00:00:00 2001 From: Dorin Geman Date: Tue, 16 Sep 2025 15:35:21 +0300 Subject: [PATCH 5/7] OpenAIRecorder: send error event for JSON marshaling failure Signed-off-by: Dorin Geman --- pkg/metrics/openai_recorder.go | 8 +++++++- 1 file changed, 7 insertions(+), 1 deletion(-) diff --git a/pkg/metrics/openai_recorder.go b/pkg/metrics/openai_recorder.go index 7ca7567bd..717599189 100644 --- a/pkg/metrics/openai_recorder.go +++ b/pkg/metrics/openai_recorder.go @@ -397,6 +397,9 @@ func (r *OpenAIRecorder) handleStreamingRequests(w http.ResponseWriter, req *htt // Send as SSE event. jsonData, err := json.Marshal(modelRecords) if err != nil { + errorMsg := fmt.Sprintf(`{"error": "Failed to marshal record: %v"}`, err) + fmt.Fprintf(w, "event: error\ndata: %s\n\n", errorMsg) + flusher.Flush() continue } @@ -492,7 +495,10 @@ func (r *OpenAIRecorder) sendExistingRecords(w http.ResponseWriter, model string }, }} jsonData, err := json.Marshal(singleRecord) - if err == nil { + if err != nil { + errorMsg := fmt.Sprintf(`{"error": "Failed to marshal existing record: %v"}`, err) + fmt.Fprintf(w, "event: error\ndata: %s\n\n", errorMsg) + } else { fmt.Fprintf(w, "event: existing_request\ndata: %s\n\n", jsonData) } } From 2fca1f5caccf521578cd39c2e8f2926dbfcfd290 Mon Sep 17 00:00:00 2001 From: Dorin Geman Date: Tue, 16 Sep 2025 15:40:42 +0300 Subject: [PATCH 6/7] OpenAIRecorder: add error handling and logging for SSE response writes Signed-off-by: Dorin Geman --- pkg/metrics/openai_recorder.go | 22 +++++++++++++++++----- 1 file changed, 17 insertions(+), 5 deletions(-) diff --git a/pkg/metrics/openai_recorder.go b/pkg/metrics/openai_recorder.go index 717599189..37024344b 100644 --- a/pkg/metrics/openai_recorder.go +++ b/pkg/metrics/openai_recorder.go @@ -377,7 +377,9 @@ func (r *OpenAIRecorder) handleStreamingRequests(w http.ResponseWriter, req *htt } // Send heartbeat to establish connection. - fmt.Fprintf(w, "event: connected\ndata: {\"status\": \"connected\"}\n\n") + if _, err := fmt.Fprintf(w, "event: connected\ndata: {\"status\": \"connected\"}\n\n"); err != nil { + r.log.Errorf("Failed to write connected event to response: %v", err) + } flusher.Flush() for { @@ -397,13 +399,18 @@ func (r *OpenAIRecorder) handleStreamingRequests(w http.ResponseWriter, req *htt // Send as SSE event. jsonData, err := json.Marshal(modelRecords) if err != nil { + r.log.Errorf("Failed to marshal record for streaming: %v", err) errorMsg := fmt.Sprintf(`{"error": "Failed to marshal record: %v"}`, err) - fmt.Fprintf(w, "event: error\ndata: %s\n\n", errorMsg) + if _, writeErr := fmt.Fprintf(w, "event: error\ndata: %s\n\n", errorMsg); writeErr != nil { + r.log.Errorf("Failed to write error event to response: %v", writeErr) + } flusher.Flush() continue } - fmt.Fprintf(w, "event: new_request\ndata: %s\n\n", jsonData) + if _, err := fmt.Fprintf(w, "event: new_request\ndata: %s\n\n", jsonData); err != nil { + r.log.Errorf("Failed to write new_request event to response: %v", err) + } flusher.Flush() case <-req.Context().Done(): @@ -496,10 +503,15 @@ func (r *OpenAIRecorder) sendExistingRecords(w http.ResponseWriter, model string }} jsonData, err := json.Marshal(singleRecord) if err != nil { + r.log.Errorf("Failed to marshal existing record for streaming: %v", err) errorMsg := fmt.Sprintf(`{"error": "Failed to marshal existing record: %v"}`, err) - fmt.Fprintf(w, "event: error\ndata: %s\n\n", errorMsg) + if _, writeErr := fmt.Fprintf(w, "event: error\ndata: %s\n\n", errorMsg); writeErr != nil { + r.log.Errorf("Failed to write error event to response: %v", writeErr) + } } else { - fmt.Fprintf(w, "event: existing_request\ndata: %s\n\n", jsonData) + if _, writeErr := fmt.Fprintf(w, "event: existing_request\ndata: %s\n\n", jsonData); writeErr != nil { + r.log.Errorf("Failed to write existing_request event to response: %v", writeErr) + } } } } From 296301f2b58f7d457fe2d2b7c47d47ef8d1f8bc2 Mon Sep 17 00:00:00 2001 From: Dorin Geman Date: Tue, 16 Sep 2025 16:14:28 +0300 Subject: [PATCH 7/7] OpenAIRecorder: return an empty list instead of 404 Signed-off-by: Dorin Geman --- pkg/metrics/openai_recorder.go | 8 ++------ 1 file changed, 2 insertions(+), 6 deletions(-) diff --git a/pkg/metrics/openai_recorder.go b/pkg/metrics/openai_recorder.go index 37024344b..a7669be22 100644 --- a/pkg/metrics/openai_recorder.go +++ b/pkg/metrics/openai_recorder.go @@ -316,9 +316,7 @@ func (r *OpenAIRecorder) handleJSONRequests(w http.ResponseWriter, req *http.Req // Retrieve all records for all models. allRecords := r.getAllRecords() if allRecords == nil { - // No records found. - http.Error(w, "No records found", http.StatusNotFound) - return + allRecords = []ModelRecordsResponse{} } if err := json.NewEncoder(w).Encode(allRecords); err != nil { http.Error(w, fmt.Sprintf("Failed to encode all records: %v", err), @@ -329,9 +327,7 @@ func (r *OpenAIRecorder) handleJSONRequests(w http.ResponseWriter, req *http.Req // Retrieve records for the specified model. records := r.getRecordsByModel(model) if records == nil { - // No records found for the specified model. - http.Error(w, fmt.Sprintf("No records found for model '%s'", model), http.StatusNotFound) - return + records = []ModelRecordsResponse{} } if err := json.NewEncoder(w).Encode(records); err != nil { http.Error(w, fmt.Sprintf("Failed to encode records for model '%s': %v", model, err),