MCPcopy Create free account
hub / github.com/InferCore/InferCore / invokeStream

Method invokeStream

internal/adapters/vllm/vllm.go:144–314  ·  view source on GitHub ↗
(ctx context.Context, text, model string)

Source from the content-addressed store, hash-verified

142}
143
144func (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{

Callers 1

InvokeMethod · 0.95

Calls 5

chatCompletionsURLMethod · 0.95
applyHeadersMethod · 0.95
NewFunction · 0.92
CloseMethod · 0.65
ErrorMethod · 0.45

Tested by

no test coverage detected