From 486da82d812a7a8e96f6440d68a9d3b1dce039ad Mon Sep 17 00:00:00 2001 From: breken-ai <312387581+breken-ai@users.noreply.github.com> Date: Fri, 25 Sep 2026 21:18:09 -0700 Subject: [PATCH 1/4] fix(responses): join streamed tool call argument chunks In a chat completion stream only the first delta of a tool call has its id and name. The argument deltas after it carry just "index" (llama.cpp and OpenAI both stream this way). handleToolCallDelta looked tool calls up by id only, so every argument delta became a new function_call item with a random call_id and no name. With streaming on, a Responses API client got one item with the right name and empty arguments, plus one nameless item per argument chunk. Parse "index" and match on it, falling back to the id for streams that do not send it. Argument deltas also report the output_index of their own item now. Assisted-By: Claude Signed-off-by: breken-ai <312387581+breken-ai@users.noreply.github.com> --- pkg/responses/handler_test.go | 86 +++++++++++++++++++++++++++++++++++ pkg/responses/streaming.go | 42 +++++++++++++---- pkg/responses/transform.go | 3 ++ 3 files changed, 121 insertions(+), 10 deletions(-) diff --git a/pkg/responses/handler_test.go b/pkg/responses/handler_test.go index 255320ee4..7a6f34911 100644 --- a/pkg/responses/handler_test.go +++ b/pkg/responses/handler_test.go @@ -1260,6 +1260,92 @@ func TestHandler_CreateResponse_Streaming_Persistence(t *testing.T) { } } +func TestHandler_CreateResponse_Streaming_ToolCallArgumentChunks(t *testing.T) { + // Chat completion streams send the id and name only in the first chunk of + // each tool call. The argument chunks that follow carry just the index + // (this is what llama.cpp and OpenAI send). + chunk := func(delta string) string { + return "data: {\"id\":\"chatcmpl-1\",\"object\":\"chat.completion.chunk\",\"choices\":[{\"index\":0,\"delta\":" + delta + ",\"finish_reason\":null}]}\n\n" + } + mock := &mockSchedulerHTTP{ + streaming: true, + streamChunks: []string{ + chunk(`{"role":"assistant","content":null}`), + chunk(`{"tool_calls":[{"index":0,"id":"call_a","type":"function","function":{"name":"get_weather"}}]}`), + chunk(`{"tool_calls":[{"index":0,"function":{"arguments":"{\"city\":"}}]}`), + chunk(`{"tool_calls":[{"index":0,"function":{"arguments":"\"Paris\"}"}}]}`), + chunk(`{"tool_calls":[{"index":1,"id":"call_b","type":"function","function":{"name":"get_time"}}]}`), + chunk(`{"tool_calls":[{"index":1,"function":{"arguments":"{\"tz\":\"CET\"}"}}]}`), + "data: {\"id\":\"chatcmpl-1\",\"object\":\"chat.completion.chunk\",\"choices\":[{\"index\":0,\"delta\":{},\"finish_reason\":\"tool_calls\"}]}\n\n", + "data: [DONE]\n\n", + }, + } + + handler := newTestHandler(t, mock) + + reqBody := `{"model": "gpt-4", "input": "Weather and time in Paris?", "stream": true}` + req := httptest.NewRequest(http.MethodPost, "/v1/responses", strings.NewReader(reqBody)) + req.Header.Set("Content-Type", "application/json") + w := httptest.NewRecorder() + + handler.handleCreate(w, req) + + if w.Code != http.StatusOK { + t.Fatalf("status = %d, want %d", w.Code, http.StatusOK) + } + + ids := handler.store.GetResponseIDs() + if len(ids) != 1 { + t.Fatalf("expected one stored response, got %d", len(ids)) + } + persisted, ok := handler.store.Get(ids[0]) + if !ok { + t.Fatal("stored response not found") + } + + type call struct{ callID, name, args string } + var got []call + for _, item := range persisted.Output { + if item.Type == ItemTypeFunctionCall { + got = append(got, call{item.CallID, item.Name, item.Arguments}) + } + } + want := []call{ + {"call_a", "get_weather", `{"city":"Paris"}`}, + {"call_b", "get_time", `{"tz":"CET"}`}, + } + if len(got) != len(want) { + t.Fatalf("function_call items = %+v, want %+v", got, want) + } + for i := range want { + if got[i] != want[i] { + t.Errorf("function_call[%d] = %+v, want %+v", i, got[i], want[i]) + } + } + + // Argument deltas must point at the output item they belong to. + for _, line := range strings.Split(w.Body.String(), "\n") { + data, found := strings.CutPrefix(line, "data: ") + if !found { + continue + } + var ev StreamEvent + if err := json.Unmarshal([]byte(data), &ev); err != nil { + t.Fatalf("bad event %q: %v", data, err) + } + if ev.Type != EventFunctionCallArgsDelta { + continue + } + wantIndex := 0 + if strings.Contains(ev.Delta, "tz") { + wantIndex = 1 + } + if ev.OutputIndex != wantIndex { + t.Errorf("delta %q has output_index %d, want %d", ev.Delta, ev.OutputIndex, wantIndex) + } + } +} + // Benchmark for response creation func BenchmarkHandler_CreateResponse(b *testing.B) { mock := &mockSchedulerHTTP{ diff --git a/pkg/responses/streaming.go b/pkg/responses/streaming.go index e4637fb58..a86e93a30 100644 --- a/pkg/responses/streaming.go +++ b/pkg/responses/streaming.go @@ -22,6 +22,8 @@ type StreamingResponseWriter struct { currentContentIdx int accumulatedContent strings.Builder toolCalls []OutputItem + // toolCallPos maps a streaming tool call index to its position in toolCalls. + toolCallPos map[int]int } // NewStreamingResponseWriter creates a new streaming response writer. @@ -319,16 +321,29 @@ func (s *StreamingResponseWriter) handleContentDelta(content string) { // handleToolCallDelta handles tool call deltas from the chat completion stream. func (s *StreamingResponseWriter) handleToolCallDelta(toolCalls []ChatToolCall) { for _, tc := range toolCalls { - // Find or create the tool call item - var item *OutputItem - for i := range s.toolCalls { - if s.toolCalls[i].CallID == tc.ID { - item = &s.toolCalls[i] - break + // Find or create the tool call item. Argument deltas after the first + // one carry only the index (no ID), so match on the index when present. + pos := -1 + if tc.Index != nil { + if p, ok := s.toolCallPos[*tc.Index]; ok { + pos = p + } + } else if tc.ID != "" { + for i := range s.toolCalls { + if s.toolCalls[i].CallID == tc.ID { + pos = i + break + } } } - if item == nil { + var item *OutputItem + if pos >= 0 { + item = &s.toolCalls[pos] + if item.Name == "" { + item.Name = tc.Function.Name + } + } else { // New tool call callID := tc.ID if callID == "" { @@ -343,14 +358,21 @@ func (s *StreamingResponseWriter) handleToolCallDelta(toolCalls []ChatToolCall) Status: StatusInProgress, } s.toolCalls = append(s.toolCalls, newItem) - item = &s.toolCalls[len(s.toolCalls)-1] + pos = len(s.toolCalls) - 1 + item = &s.toolCalls[pos] + if tc.Index != nil { + if s.toolCallPos == nil { + s.toolCallPos = make(map[int]int) + } + s.toolCallPos[*tc.Index] = pos + } // Send output_item.added for function call s.sendEvent(EventOutputItemAdded, &StreamEvent{ Type: EventOutputItemAdded, SequenceNumber: s.nextSeq(), Item: item, - OutputIndex: len(s.toolCalls) - 1, + OutputIndex: pos, }) } @@ -363,7 +385,7 @@ func (s *StreamingResponseWriter) handleToolCallDelta(toolCalls []ChatToolCall) Type: EventFunctionCallArgsDelta, SequenceNumber: s.nextSeq(), ItemID: item.ID, - OutputIndex: len(s.toolCalls) - 1, + OutputIndex: pos, Delta: tc.Function.Arguments, }) } diff --git a/pkg/responses/transform.go b/pkg/responses/transform.go index cae624ee0..2b52d3f85 100644 --- a/pkg/responses/transform.go +++ b/pkg/responses/transform.go @@ -75,6 +75,9 @@ type ChatFunction struct { // ChatToolCall represents a tool call in chat format. type ChatToolCall struct { + // Index identifies the tool call in streaming deltas. Only the first + // delta of a call carries its ID; later argument deltas carry only Index. + Index *int `json:"index,omitempty"` ID string `json:"id"` Type string `json:"type"` Function ChatFunctionCall `json:"function"` From e7e96e1631a83e4e27677e4d640434021669e197 Mon Sep 17 00:00:00 2001 From: breken-ai <312387581+breken-ai@users.noreply.github.com> Date: Fri, 25 Sep 2026 21:21:08 -0700 Subject: [PATCH 2/4] fix(responses): fall back to the id when the tool call index is new If a call was first seen by id without an index, a later delta that adds the index now extends that call instead of starting a duplicate, and the index is remembered for the index-only deltas that follow. Assisted-By: Claude Signed-off-by: breken-ai <312387581+breken-ai@users.noreply.github.com> --- pkg/responses/handler_test.go | 18 ++++++++++++++++++ pkg/responses/streaming.go | 25 +++++++++++++++++-------- 2 files changed, 35 insertions(+), 8 deletions(-) diff --git a/pkg/responses/handler_test.go b/pkg/responses/handler_test.go index 7a6f34911..62d4c9844 100644 --- a/pkg/responses/handler_test.go +++ b/pkg/responses/handler_test.go @@ -1346,6 +1346,24 @@ func TestHandler_CreateResponse_Streaming_ToolCallArgumentChunks(t *testing.T) { } } +func TestStreamingResponseWriter_ToolCallMatchedByIDThenIndex(t *testing.T) { + // A delta that adds the index to a call first seen by ID must extend + // that call, and later index-only deltas must go to the same call. + w := httptest.NewRecorder() + s := NewStreamingResponseWriter(w, &Response{}, nil) + idx := 0 + s.handleToolCallDelta([]ChatToolCall{{ID: "call_a", Function: ChatFunctionCall{Name: "get_weather"}}}) + s.handleToolCallDelta([]ChatToolCall{{Index: &idx, ID: "call_a", Function: ChatFunctionCall{Arguments: `{"city":`}}}) + s.handleToolCallDelta([]ChatToolCall{{Index: &idx, Function: ChatFunctionCall{Arguments: `"Paris"}`}}}) + + if len(s.toolCalls) != 1 { + t.Fatalf("tool calls = %+v, want one", s.toolCalls) + } + if got := s.toolCalls[0]; got.CallID != "call_a" || got.Name != "get_weather" || got.Arguments != `{"city":"Paris"}` { + t.Errorf("tool call = %+v", got) + } +} + // Benchmark for response creation func BenchmarkHandler_CreateResponse(b *testing.B) { mock := &mockSchedulerHTTP{ diff --git a/pkg/responses/streaming.go b/pkg/responses/streaming.go index a86e93a30..c00f4395a 100644 --- a/pkg/responses/streaming.go +++ b/pkg/responses/streaming.go @@ -322,13 +322,15 @@ func (s *StreamingResponseWriter) handleContentDelta(content string) { func (s *StreamingResponseWriter) handleToolCallDelta(toolCalls []ChatToolCall) { for _, tc := range toolCalls { // Find or create the tool call item. Argument deltas after the first - // one carry only the index (no ID), so match on the index when present. + // one carry only the index (no ID), so match on the index first and + // fall back to the ID. pos := -1 if tc.Index != nil { if p, ok := s.toolCallPos[*tc.Index]; ok { pos = p } - } else if tc.ID != "" { + } + if pos < 0 && tc.ID != "" { for i := range s.toolCalls { if s.toolCalls[i].CallID == tc.ID { pos = i @@ -343,6 +345,7 @@ func (s *StreamingResponseWriter) handleToolCallDelta(toolCalls []ChatToolCall) if item.Name == "" { item.Name = tc.Function.Name } + s.rememberToolCallIndex(tc.Index, pos) } else { // New tool call callID := tc.ID @@ -360,12 +363,7 @@ func (s *StreamingResponseWriter) handleToolCallDelta(toolCalls []ChatToolCall) s.toolCalls = append(s.toolCalls, newItem) pos = len(s.toolCalls) - 1 item = &s.toolCalls[pos] - if tc.Index != nil { - if s.toolCallPos == nil { - s.toolCallPos = make(map[int]int) - } - s.toolCallPos[*tc.Index] = pos - } + s.rememberToolCallIndex(tc.Index, pos) // Send output_item.added for function call s.sendEvent(EventOutputItemAdded, &StreamEvent{ @@ -392,6 +390,17 @@ func (s *StreamingResponseWriter) handleToolCallDelta(toolCalls []ChatToolCall) } } +// rememberToolCallIndex records which item a streaming tool call index refers to. +func (s *StreamingResponseWriter) rememberToolCallIndex(index *int, pos int) { + if index == nil { + return + } + if s.toolCallPos == nil { + s.toolCallPos = make(map[int]int) + } + s.toolCallPos[*index] = pos +} + // finalize completes the streaming response. func (s *StreamingResponseWriter) finalize() { // Finalize any accumulated content From 65d55933cecd3c010cab91d3bd194ed87941e3b1 Mon Sep 17 00:00:00 2001 From: breken-ai <312387581+breken-ai@users.noreply.github.com> Date: Sun, 27 Sep 2026 11:24:50 -0700 Subject: [PATCH 3/4] fix(responses): account for text before tool call indices Assisted-By: Claude Signed-off-by: breken-ai <312387581+breken-ai@users.noreply.github.com> --- pkg/responses/handler_test.go | 9 ++++++--- pkg/responses/streaming.go | 17 +++++++++++++---- 2 files changed, 19 insertions(+), 7 deletions(-) diff --git a/pkg/responses/handler_test.go b/pkg/responses/handler_test.go index 62d4c9844..877ecaeec 100644 --- a/pkg/responses/handler_test.go +++ b/pkg/responses/handler_test.go @@ -1270,7 +1270,7 @@ func TestHandler_CreateResponse_Streaming_ToolCallArgumentChunks(t *testing.T) { mock := &mockSchedulerHTTP{ streaming: true, streamChunks: []string{ - chunk(`{"role":"assistant","content":null}`), + chunk(`{"role":"assistant","content":"Checking the weather."}`), chunk(`{"tool_calls":[{"index":0,"id":"call_a","type":"function","function":{"name":"get_weather"}}]}`), chunk(`{"tool_calls":[{"index":0,"function":{"arguments":"{\"city\":"}}]}`), chunk(`{"tool_calls":[{"index":0,"function":{"arguments":"\"Paris\"}"}}]}`), @@ -1329,6 +1329,9 @@ func TestHandler_CreateResponse_Streaming_ToolCallArgumentChunks(t *testing.T) { if !found { continue } + if data == "[DONE]" { + continue + } var ev StreamEvent if err := json.Unmarshal([]byte(data), &ev); err != nil { t.Fatalf("bad event %q: %v", data, err) @@ -1336,9 +1339,9 @@ func TestHandler_CreateResponse_Streaming_ToolCallArgumentChunks(t *testing.T) { if ev.Type != EventFunctionCallArgsDelta { continue } - wantIndex := 0 + wantIndex := 1 // The assistant text item precedes the tool calls. if strings.Contains(ev.Delta, "tz") { - wantIndex = 1 + wantIndex = 2 } if ev.OutputIndex != wantIndex { t.Errorf("delta %q has output_index %d, want %d", ev.Delta, ev.OutputIndex, wantIndex) diff --git a/pkg/responses/streaming.go b/pkg/responses/streaming.go index c00f4395a..437dcd84c 100644 --- a/pkg/responses/streaming.go +++ b/pkg/responses/streaming.go @@ -370,7 +370,7 @@ func (s *StreamingResponseWriter) handleToolCallDelta(toolCalls []ChatToolCall) Type: EventOutputItemAdded, SequenceNumber: s.nextSeq(), Item: item, - OutputIndex: pos, + OutputIndex: s.toolCallOutputIndex(pos), }) } @@ -383,13 +383,22 @@ func (s *StreamingResponseWriter) handleToolCallDelta(toolCalls []ChatToolCall) Type: EventFunctionCallArgsDelta, SequenceNumber: s.nextSeq(), ItemID: item.ID, - OutputIndex: pos, + OutputIndex: s.toolCallOutputIndex(pos), Delta: tc.Function.Arguments, }) } } } +// toolCallOutputIndex converts a position in toolCalls to the corresponding +// position in response.Output. The assistant message, when present, is first. +func (s *StreamingResponseWriter) toolCallOutputIndex(pos int) int { + if s.currentItemID != "" { + return pos + 1 + } + return pos +} + // rememberToolCallIndex records which item a streaming tool call index refers to. func (s *StreamingResponseWriter) rememberToolCallIndex(index *int, pos int) { if index == nil { @@ -475,7 +484,7 @@ func (s *StreamingResponseWriter) finalize() { Type: EventFunctionCallArgsDone, SequenceNumber: s.nextSeq(), ItemID: tc.ID, - OutputIndex: i, + OutputIndex: s.toolCallOutputIndex(i), Delta: tc.Arguments, }) @@ -484,7 +493,7 @@ func (s *StreamingResponseWriter) finalize() { s.sendEvent(EventOutputItemDone, &StreamEvent{ Type: EventOutputItemDone, SequenceNumber: s.nextSeq(), - OutputIndex: i, + OutputIndex: s.toolCallOutputIndex(i), Item: &tc, }) From 709e000e28f6f8b9ae007d1baf79a92fd1f722e6 Mon Sep 17 00:00:00 2001 From: breken-ai <312387581+breken-ai@users.noreply.github.com> Date: Wed, 7 Oct 2026 17:20:43 -0700 Subject: [PATCH 4/4] fix(responses): keep each item's output index fixed from when it is added A tool call that started before any assistant text was added at output index 0, then switched to index 1 once text arrived, and the stored output listed the message first. Assign each item its output index when it is added, use it for all of its events, and place it at that index in the response output. Tool calls that share a stream index but carry different IDs are now kept as separate calls instead of having their arguments joined. Co-Authored-By: Claude Opus 5.5 Co-authored-by: Eric Curtin Signed-off-by: breken-ai <312387581+breken-ai@users.noreply.github.com> --- pkg/responses/handler_test.go | 73 +++++++++++++++++++++++++++++++++++ pkg/responses/streaming.go | 50 +++++++++++++++++------- 2 files changed, 108 insertions(+), 15 deletions(-) diff --git a/pkg/responses/handler_test.go b/pkg/responses/handler_test.go index 877ecaeec..cf30ae3b8 100644 --- a/pkg/responses/handler_test.go +++ b/pkg/responses/handler_test.go @@ -1349,6 +1349,58 @@ func TestHandler_CreateResponse_Streaming_ToolCallArgumentChunks(t *testing.T) { } } +func TestStreamingResponseWriter_ToolCallBeforeTextKeepsOutputIndex(t *testing.T) { + // A tool call that starts before any assistant text must keep the + // output_index it was added with, and the stored output must list the + // items in that order. + w := httptest.NewRecorder() + resp := &Response{} + s := NewStreamingResponseWriter(w, resp, nil) + idx := 0 + s.handleToolCallDelta([]ChatToolCall{{Index: &idx, ID: "call_a", Function: ChatFunctionCall{Name: "get_weather"}}}) + s.handleToolCallDelta([]ChatToolCall{{Index: &idx, Function: ChatFunctionCall{Arguments: `{"city":`}}}) + s.handleContentDelta("Checking the weather.") + s.handleToolCallDelta([]ChatToolCall{{Index: &idx, Function: ChatFunctionCall{Arguments: `"Paris"}`}}}) + s.finalize() + + indices := map[string]map[int]bool{} + for _, line := range strings.Split(w.Body.String(), "\n") { + data, found := strings.CutPrefix(line, "data: ") + if !found { + continue + } + var ev StreamEvent + if err := json.Unmarshal([]byte(data), &ev); err != nil { + t.Fatalf("bad event %q: %v", data, err) + } + itemID := ev.ItemID + if ev.Item != nil { + itemID = ev.Item.ID + } + if itemID == "" { + continue + } + if indices[itemID] == nil { + indices[itemID] = map[int]bool{} + } + indices[itemID][ev.OutputIndex] = true + } + + if len(resp.Output) != 2 { + t.Fatalf("output = %+v, want a function call and a message", resp.Output) + } + if resp.Output[0].Type != ItemTypeFunctionCall || resp.Output[1].Type != ItemTypeMessage { + t.Errorf("output types = [%s %s], want [%s %s]", + resp.Output[0].Type, resp.Output[1].Type, ItemTypeFunctionCall, ItemTypeMessage) + } + for pos, item := range resp.Output { + got := indices[item.ID] + if len(got) != 1 || !got[pos] { + t.Errorf("%s item %s used output_index %v, want only %d", item.Type, item.ID, got, pos) + } + } +} + func TestStreamingResponseWriter_ToolCallMatchedByIDThenIndex(t *testing.T) { // A delta that adds the index to a call first seen by ID must extend // that call, and later index-only deltas must go to the same call. @@ -1367,6 +1419,27 @@ func TestStreamingResponseWriter_ToolCallMatchedByIDThenIndex(t *testing.T) { } } +func TestStreamingResponseWriter_ToolCallsSharingIndexSplitByID(t *testing.T) { + // Some servers send index 0 for every parallel call and tell them apart + // by ID only. Those must stay separate calls. + w := httptest.NewRecorder() + s := NewStreamingResponseWriter(w, &Response{}, nil) + idx := 0 + s.handleToolCallDelta([]ChatToolCall{{Index: &idx, ID: "call_a", Function: ChatFunctionCall{Name: "get_weather", Arguments: `{"city":"Paris"}`}}}) + s.handleToolCallDelta([]ChatToolCall{{Index: &idx, ID: "call_b", Function: ChatFunctionCall{Name: "get_time", Arguments: `{"tz":`}}}) + s.handleToolCallDelta([]ChatToolCall{{Index: &idx, Function: ChatFunctionCall{Arguments: `"CET"}`}}}) + + if len(s.toolCalls) != 2 { + t.Fatalf("tool calls = %+v, want two", s.toolCalls) + } + if got := s.toolCalls[0]; got.CallID != "call_a" || got.Arguments != `{"city":"Paris"}` { + t.Errorf("tool call 0 = %+v", got) + } + if got := s.toolCalls[1]; got.CallID != "call_b" || got.Name != "get_time" || got.Arguments != `{"tz":"CET"}` { + t.Errorf("tool call 1 = %+v", got) + } +} + // Benchmark for response creation func BenchmarkHandler_CreateResponse(b *testing.B) { mock := &mockSchedulerHTTP{ diff --git a/pkg/responses/streaming.go b/pkg/responses/streaming.go index 437dcd84c..60fe50a60 100644 --- a/pkg/responses/streaming.go +++ b/pkg/responses/streaming.go @@ -24,6 +24,12 @@ type StreamingResponseWriter struct { toolCalls []OutputItem // toolCallPos maps a streaming tool call index to its position in toolCalls. toolCallPos map[int]int + + // Output indices are assigned when an item is added, in arrival order, and + // stay fixed for all of that item's events and its place in the output. + nextOutputIndex int + messageOutputIndex int + toolCallOutputIndices []int } // NewStreamingResponseWriter creates a new streaming response writer. @@ -269,6 +275,7 @@ func (s *StreamingResponseWriter) handleContentDelta(content string) { if s.currentItemID == "" { s.currentItemID = GenerateMessageID() s.currentContentIdx = 0 + s.messageOutputIndex = s.allocOutputIndex() // Send output_item.added item := &OutputItem{ @@ -286,7 +293,7 @@ func (s *StreamingResponseWriter) handleContentDelta(content string) { Type: EventOutputItemAdded, SequenceNumber: s.nextSeq(), Item: item, - OutputIndex: 0, + OutputIndex: s.messageOutputIndex, }) // Send content_part.added @@ -294,7 +301,7 @@ func (s *StreamingResponseWriter) handleContentDelta(content string) { Type: EventContentPartAdded, SequenceNumber: s.nextSeq(), ItemID: s.currentItemID, - OutputIndex: 0, + OutputIndex: s.messageOutputIndex, ContentIndex: 0, Part: &ContentPart{ Type: ContentTypeOutputText, @@ -312,7 +319,7 @@ func (s *StreamingResponseWriter) handleContentDelta(content string) { Type: EventOutputTextDelta, SequenceNumber: s.nextSeq(), ItemID: s.currentItemID, - OutputIndex: 0, + OutputIndex: s.messageOutputIndex, ContentIndex: 0, Delta: content, }) @@ -328,6 +335,11 @@ func (s *StreamingResponseWriter) handleToolCallDelta(toolCalls []ChatToolCall) if tc.Index != nil { if p, ok := s.toolCallPos[*tc.Index]; ok { pos = p + // Some servers reuse one index for several calls and tell them + // apart by ID, so an ID that differs starts a different call. + if tc.ID != "" && s.toolCalls[p].CallID != tc.ID { + pos = -1 + } } } if pos < 0 && tc.ID != "" { @@ -361,6 +373,7 @@ func (s *StreamingResponseWriter) handleToolCallDelta(toolCalls []ChatToolCall) Status: StatusInProgress, } s.toolCalls = append(s.toolCalls, newItem) + s.toolCallOutputIndices = append(s.toolCallOutputIndices, s.allocOutputIndex()) pos = len(s.toolCalls) - 1 item = &s.toolCalls[pos] s.rememberToolCallIndex(tc.Index, pos) @@ -390,13 +403,16 @@ func (s *StreamingResponseWriter) handleToolCallDelta(toolCalls []ChatToolCall) } } -// toolCallOutputIndex converts a position in toolCalls to the corresponding -// position in response.Output. The assistant message, when present, is first. +// allocOutputIndex returns the output index for a newly added item. +func (s *StreamingResponseWriter) allocOutputIndex() int { + i := s.nextOutputIndex + s.nextOutputIndex++ + return i +} + +// toolCallOutputIndex returns the output index assigned to toolCalls[pos]. func (s *StreamingResponseWriter) toolCallOutputIndex(pos int) int { - if s.currentItemID != "" { - return pos + 1 - } - return pos + return s.toolCallOutputIndices[pos] } // rememberToolCallIndex records which item a streaming tool call index refers to. @@ -412,6 +428,9 @@ func (s *StreamingResponseWriter) rememberToolCallIndex(index *int, pos int) { // finalize completes the streaming response. func (s *StreamingResponseWriter) finalize() { + // Items go into the output at the index their events used. + output := make([]OutputItem, s.nextOutputIndex) + // Finalize any accumulated content if s.currentItemID != "" { finalText := s.accumulatedContent.String() @@ -421,7 +440,7 @@ func (s *StreamingResponseWriter) finalize() { Type: EventOutputTextDone, SequenceNumber: s.nextSeq(), ItemID: s.currentItemID, - OutputIndex: 0, + OutputIndex: s.messageOutputIndex, ContentIndex: 0, Part: &ContentPart{ Type: ContentTypeOutputText, @@ -435,7 +454,7 @@ func (s *StreamingResponseWriter) finalize() { Type: EventContentPartDone, SequenceNumber: s.nextSeq(), ItemID: s.currentItemID, - OutputIndex: 0, + OutputIndex: s.messageOutputIndex, ContentIndex: 0, Part: &ContentPart{ Type: ContentTypeOutputText, @@ -448,7 +467,7 @@ func (s *StreamingResponseWriter) finalize() { s.sendEvent(EventOutputItemDone, &StreamEvent{ Type: EventOutputItemDone, SequenceNumber: s.nextSeq(), - OutputIndex: 0, + OutputIndex: s.messageOutputIndex, Item: &OutputItem{ ID: s.currentItemID, Type: ItemTypeMessage, @@ -463,7 +482,7 @@ func (s *StreamingResponseWriter) finalize() { }) // Add to response output - s.response.Output = append(s.response.Output, OutputItem{ + output[s.messageOutputIndex] = OutputItem{ ID: s.currentItemID, Type: ItemTypeMessage, Role: "assistant", @@ -473,7 +492,7 @@ func (s *StreamingResponseWriter) finalize() { Annotations: []Annotation{}, }}, Status: StatusCompleted, - }) + } s.response.OutputText = finalText } @@ -498,8 +517,9 @@ func (s *StreamingResponseWriter) finalize() { }) // Add to response output - s.response.Output = append(s.response.Output, tc) + output[s.toolCallOutputIndex(i)] = tc } + s.response.Output = append(s.response.Output, output...) // Update response status s.response.Status = StatusCompleted