(ctx context.Context, text, model string)
| 142 | } |
| 143 | |
| 144 | func (a *Adapter) invokeStream(ctx context.Context, text, model string) (types.BackendResponse, error) { |
| 145 | t0 := time.Now() |
| 146 | payload := map[string]any{ |
| 147 | "model": model, |
| 148 | "messages": []map[string]string{ |
| 149 | {"role": "user", "content": text}, |
| 150 | }, |
| 151 | "stream": true, |
| 152 | } |
| 153 | raw, err := json.Marshal(payload) |
| 154 | if err != nil { |
| 155 | return types.BackendResponse{}, err |
| 156 | } |
| 157 | chatURL := a.chatCompletionsURL() |
| 158 | httpReq, err := http.NewRequestWithContext(ctx, http.MethodPost, chatURL, bytes.NewReader(raw)) |
| 159 | if err != nil { |
| 160 | return types.BackendResponse{}, err |
| 161 | } |
| 162 | httpReq.Header.Set("Content-Type", "application/json") |
| 163 | httpReq.Header.Set("Accept", "text/event-stream") |
| 164 | a.applyHeaders(httpReq) |
| 165 | |
| 166 | resp, err := a.httpClient.Do(httpReq) |
| 167 | if err != nil { |
| 168 | if errors.Is(err, context.DeadlineExceeded) { |
| 169 | return types.BackendResponse{}, upstream.New(upstream.KindTimeout, err.Error()) |
| 170 | } |
| 171 | return types.BackendResponse{}, err |
| 172 | } |
| 173 | defer resp.Body.Close() |
| 174 | if resp.StatusCode < 200 || resp.StatusCode >= 300 { |
| 175 | if resp.StatusCode >= 500 { |
| 176 | return types.BackendResponse{}, upstream.New(upstream.KindUpstream5xx, fmt.Sprintf("openai-compatible stream failed status=%d", resp.StatusCode)) |
| 177 | } |
| 178 | return types.BackendResponse{}, upstream.New(upstream.KindUpstream4xx, fmt.Sprintf("openai-compatible stream failed status=%d", resp.StatusCode)) |
| 179 | } |
| 180 | |
| 181 | ct := strings.ToLower(resp.Header.Get("Content-Type")) |
| 182 | if !strings.Contains(ct, "event-stream") && !strings.Contains(ct, "text/plain") { |
| 183 | // Upstream ignored stream flag; parse as JSON completion. |
| 184 | var parsed struct { |
| 185 | Choices []struct { |
| 186 | Message struct { |
| 187 | Content string `json:"content"` |
| 188 | } `json:"message"` |
| 189 | } `json:"choices"` |
| 190 | } |
| 191 | if err := json.NewDecoder(resp.Body).Decode(&parsed); err != nil { |
| 192 | return types.BackendResponse{}, upstream.New(upstream.KindBackendError, err.Error()) |
| 193 | } |
| 194 | out := "" |
| 195 | if len(parsed.Choices) > 0 { |
| 196 | out = parsed.Choices[0].Message.Content |
| 197 | } |
| 198 | t1 := time.Now() |
| 199 | lat := t1.Sub(t0).Milliseconds() |
| 200 | return types.BackendResponse{ |
| 201 | Output: map[string]any{ |
no test coverage detected