streamingRequest makes a request to the given streamable server with the given url, sessionID, and method. If provided, the in messages are encoded in the request body. A single message is encoded as a JSON object. Multiple messages are batched as a JSON array. Any received messages are sent to th
(ctx context.Context, serverURL, sessionID string, out chan<- jsonrpc.Message)
| 1456 | // returned, sessionID and status code may still be set if the error occurs |
| 1457 | // after the response headers have been received. |
| 1458 | func (s streamableRequest) do(ctx context.Context, serverURL, sessionID string, out chan<- jsonrpc.Message) (string, int, []byte, error) { |
| 1459 | defer close(out) |
| 1460 | |
| 1461 | var body []byte |
| 1462 | if len(s.messages) == 1 { |
| 1463 | data, err := jsonrpc2.EncodeMessage(s.messages[0]) |
| 1464 | if err != nil { |
| 1465 | return "", 0, nil, fmt.Errorf("encoding message: %w", err) |
| 1466 | } |
| 1467 | body = data |
| 1468 | } else { |
| 1469 | var rawMsgs []json.RawMessage |
| 1470 | for _, msg := range s.messages { |
| 1471 | data, err := jsonrpc2.EncodeMessage(msg) |
| 1472 | if err != nil { |
| 1473 | return "", 0, nil, fmt.Errorf("encoding message: %w", err) |
| 1474 | } |
| 1475 | rawMsgs = append(rawMsgs, data) |
| 1476 | } |
| 1477 | data, err := json.Marshal(rawMsgs) |
| 1478 | if err != nil { |
| 1479 | return "", 0, nil, fmt.Errorf("marshaling batch: %w", err) |
| 1480 | } |
| 1481 | body = data |
| 1482 | } |
| 1483 | |
| 1484 | req, err := http.NewRequestWithContext(ctx, s.method, serverURL, bytes.NewReader(body)) |
| 1485 | if err != nil { |
| 1486 | return "", 0, nil, fmt.Errorf("creating request: %w", err) |
| 1487 | } |
| 1488 | if sessionID != "" { |
| 1489 | req.Header.Set(sessionIDHeader, sessionID) |
| 1490 | } |
| 1491 | if s.method == http.MethodPost { |
| 1492 | req.Header.Set("Content-Type", "application/json") |
| 1493 | } |
| 1494 | req.Header.Set("Accept", "application/json, text/event-stream") |
| 1495 | maps.Copy(req.Header, s.headers) |
| 1496 | |
| 1497 | if req.Header.Get("Content-Type") == "" { |
| 1498 | req.Header.Del("Content-Type") |
| 1499 | } |
| 1500 | |
| 1501 | resp, err := http.DefaultClient.Do(req) |
| 1502 | if err != nil { |
| 1503 | return "", 0, nil, fmt.Errorf("request failed: %v", err) |
| 1504 | } |
| 1505 | defer resp.Body.Close() |
| 1506 | |
| 1507 | newSessionID := resp.Header.Get(sessionIDHeader) |
| 1508 | |
| 1509 | contentType := baseMediaType(resp.Header.Get("Content-Type")) |
| 1510 | var respBody []byte |
| 1511 | if contentType == "text/event-stream" { |
| 1512 | r := readerInto{resp.Body, new(bytes.Buffer)} |
| 1513 | for evt, err := range scanEvents(r) { |
| 1514 | if err != nil { |
| 1515 | return newSessionID, resp.StatusCode, nil, fmt.Errorf("reading events: %v", err) |
no test coverage detected