(ctx context.Context, state *AgentState, executor *StreamingToolExecutor, out chan<- StreamEvent)
| 844 | block := map[string]any{"type": "thinking", "thinking": text, "signature": stringFromAny(event["signature"])} |
| 845 | turn.ThinkingContent = append(turn.ThinkingContent, block) |
| 846 | ev := NewStreamEvent("thinking", text, map[string]any{"signature": stringFromAny(event["signature"])}) |
| 847 | return StreamAction{Event: &ev, ConsecutiveAPIErrors: consecutiveAPIErrors} |
| 848 | case "redacted_thinking": |
| 849 | turn.ThinkingContent = append(turn.ThinkingContent, map[string]any{"type": "redacted_thinking", "data": stringFromAny(event["data"])}) |
| 850 | ev := NewStreamEvent("thinking", "[redacted thinking]", nil) |
| 851 | return StreamAction{Event: &ev, ConsecutiveAPIErrors: consecutiveAPIErrors} |
| 852 | case "thinking_delta": |
| 853 | text := stringFromAny(event["text"]) |
| 854 | if len(turn.ThinkingContent) > 0 && turn.ThinkingContent[len(turn.ThinkingContent)-1]["type"] == "thinking" { |
| 855 | last := turn.ThinkingContent[len(turn.ThinkingContent)-1] |
| 856 | last["thinking"] = stringFromAny(last["thinking"]) + text |
| 857 | } |
| 858 | ev := NewStreamEvent("thinking", text, nil) |
| 859 | return StreamAction{Event: &ev, ConsecutiveAPIErrors: consecutiveAPIErrors} |
| 860 | case "signature_delta": |
| 861 | if len(turn.ThinkingContent) > 0 && turn.ThinkingContent[len(turn.ThinkingContent)-1]["type"] == "thinking" { |
| 862 | last := turn.ThinkingContent[len(turn.ThinkingContent)-1] |
| 863 | last["signature"] = stringFromAny(last["signature"]) + stringFromAny(event["signature"]) |
| 864 | } |
| 865 | case "tool_use": |
| 866 | input := map[string]any{} |
| 867 | if raw, ok := event["input"].(map[string]any); ok { |
| 868 | input = raw |
| 869 | } |
| 870 | tc := coretools.ToolCall{ID: stringFromAny(event["id"]), Name: stringFromAny(event["name"]), Input: input} |
| 871 | turn.ToolCalls = append(turn.ToolCalls, tc) |
| 872 | executor.AddTool(tc) |
| 873 | ev := NewStreamEvent("tool_call", tc.Name, map[string]any{"id": tc.ID, "input": tc.Input}) |
| 874 | return StreamAction{Event: &ev, ConsecutiveAPIErrors: consecutiveAPIErrors} |
| 875 | case "usage": |
| 876 | turn.InputTokens = intFromAny(event["input_tokens"]) |
| 877 | turn.OutputTokens = intFromAny(event["output_tokens"]) |
| 878 | case "stop_reason": |
| 879 | reason := stringFromAny(event["stop_reason"]) |
| 880 | if reason == "" { |
| 881 | reason = stringFromAny(event["reason"]) |
| 882 | } |
| 883 | if IsOutputTruncatedStopReason(reason) { |
| 884 | turn.OutputTruncated = true |
| 885 | } |
| 886 | case "model_fallback": |
| 887 | metadata := map[string]any{ |
| 888 | "primary_model": event["primary_model"], |
| 889 | "fallback_model": event["fallback_model"], |
| 890 | "reason": event["reason"], |
| 891 | } |
| 892 | ev := NewStreamEvent("model_fallback", stringFromAny(event["fallback_model"]), metadata) |
| 893 | return StreamAction{Event: &ev, ConsecutiveAPIErrors: consecutiveAPIErrors} |
| 894 | case "error": |
| 895 | msg := stringFromAny(event["message"]) |
| 896 | metadata := map[string]any{} |
| 897 | for _, key := range []string{"error_type", "provider", "status_code", "status", "raw_error", "request_url", "operation"} { |
| 898 | if value, ok := event[key]; ok { |
| 899 | metadata[key] = value |
| 900 | } |
| 901 | } |
| 902 | if IsFatalAPIError(msg) { |
| 903 | e.rememberAPIError(state, msg) |
no test coverage detected