| 1 | package responses |
| 2 | |
| 3 | import ( |
| 4 | "cmp" |
| 5 | "context" |
| 6 | "encoding/json" |
| 7 | |
| 8 | "reasonix/internal/provider" |
| 9 | ) |
| 10 | |
| 11 | // finishFunctionCall completes call from a finished function_call item, sending |
| 12 | // it once however many sources report it finished. |
| 13 | func finishFunctionCall(ctx context.Context, out chan<- provider.Chunk, call *streamedCall, item *sseItem) bool { |
| 14 | call.id = cmp.Or(item.CallID, call.id) |
| 15 | call.name = cmp.Or(item.Name, call.name) |
| 16 | call.arguments = cmp.Or(item.Arguments, call.arguments) |
| 17 | return completeFunctionCall(ctx, out, call) |
| 18 | } |
| 19 | |
| 20 | // completeFunctionCall sends call once. One with no name yet, such as an |
| 21 | // arguments.done event that named no item, stays open for a later event to name. |
| 22 | func completeFunctionCall(ctx context.Context, out chan<- provider.Chunk, call *streamedCall) bool { |
| 23 | if call.completed || call.name == "" { |
| 24 | return true |
| 25 | } |
| 26 | call.completed = true |
| 27 | return sendChunk(ctx, out, provider.Chunk{Type: provider.ChunkToolCall, ToolCall: &provider.ToolCall{ID: call.id, Name: call.name, Arguments: call.arguments}}) |
| 28 | } |
| 29 | |
| 30 | // outputItems decodes a terminal response's output one item at a time, so an |
| 31 | // item a relay sent in a shape the protocol does not define is dropped alone |
| 32 | // instead of taking the whole terminal event with it. |
| 33 | func outputItems(response *sseResponse) []sseItem { |
| 34 | if response == nil { |
| 35 | return nil |
| 36 | } |
| 37 | items := make([]sseItem, 0, len(response.Output)) |
| 38 | for _, raw := range response.Output { |
| 39 | var item sseItem |
| 40 | if json.Unmarshal(raw, &item) == nil { |
| 41 | items = append(items, item) |
| 42 | } |
| 43 | } |
| 44 | return items |
| 45 | } |
| 46 | |
| 47 | // unclosedOutputCalls lists the function calls a terminal response's output |
| 48 | // reports but the stream never closed: that list is the response's own account |
| 49 | // of what it issued. A call already closed or listed under the same call id is |
| 50 | // left out, and so is one the response marks unfinished or lists without a call id. |
| 51 | func unclosedOutputCalls(response *sseResponse, calls map[string]*streamedCall) []*sseItem { |
| 52 | closed := make(map[string]bool, len(calls)) |
| 53 | for _, call := range calls { |
| 54 | if call.completed && call.id != "" { |
| 55 | closed[call.id] = true |
| 56 | } |
| 57 | } |
| 58 | var out []*sseItem |
| 59 | items := outputItems(response) |
| 60 | for i := range items { |
| 61 | item := &items[i] |
| 62 | if item.Type != "function_call" || item.CallID == "" || |
| 63 | (item.Status != "" && item.Status != "completed") || closed[item.CallID] { |
| 64 | continue |
| 65 | } |
| 66 | closed[item.CallID] = true |
| 67 | out = append(out, item) |
| 68 | } |
| 69 | return out |
| 70 | } |
| 71 |