package mcp import ( "bufio" "bytes" "context" "encoding/json" "fmt" "io" "net/http" "strings" "sync/atomic" ) // HTTPConfig configures a Streamable HTTP client transport. type HTTPConfig struct { Endpoint string Headers map[string]string Client *http.Client } // httpTransport talks MCP over Streamable HTTP: each outbound message is an // HTTP POST whose response is either a single JSON message or an SSE stream. type httpTransport struct { endpoint string headers map[string]string client *http.Client counter int64 sessionID atomic.Value // string notify atomic.Pointer[func(*Message)] } // NewHTTPTransport creates a Streamable HTTP transport to endpoint. func NewHTTPTransport(cfg HTTPConfig) (Transport, error) { if cfg.Endpoint == "" { return nil, fmt.Errorf("http transport: empty endpoint") } c := cfg.Client if c == nil { c = http.DefaultClient } return &httpTransport{endpoint: cfg.Endpoint, headers: cfg.Headers, client: c}, nil } func (t *httpTransport) SetNotificationHandler(h func(*Message)) { if h == nil { t.notify.Store(nil) return } t.notify.Store(&h) } func (t *httpTransport) post(ctx context.Context, msg *Message) (*http.Response, error) { body, err := json.Marshal(msg) if err != nil { return nil, err } req, err := http.NewRequestWithContext(ctx, http.MethodPost, t.endpoint, bytes.NewReader(body)) if err != nil { return nil, err } req.Header.Set("Content-Type", "application/json") req.Header.Set("Accept", "application/json, text/event-stream") for k, v := range t.headers { req.Header.Set(k, v) } if sid, ok := t.sessionID.Load().(string); ok && sid != "" { req.Header.Set("Mcp-Session-Id", sid) } return t.client.Do(req) } func (t *httpTransport) Call(ctx context.Context, method string, params any) (*Message, error) { id := atomic.AddInt64(&t.counter, 1) idRaw := json.RawMessage(fmt.Sprintf("%d", id)) req, err := NewRequest(idRaw, method, params) if err != nil { return nil, err } resp, err := t.post(ctx, req) if err != nil { return nil, err } defer resp.Body.Close() if sid := resp.Header.Get("Mcp-Session-Id"); sid != "" { t.sessionID.Store(sid) } if resp.StatusCode < 200 || resp.StatusCode >= 300 { b, _ := io.ReadAll(io.LimitReader(resp.Body, 4096)) return nil, fmt.Errorf("http %s: status %d: %s", method, resp.StatusCode, strings.TrimSpace(string(b))) } ct := resp.Header.Get("Content-Type") want := normalizeID(idRaw) switch { case strings.HasPrefix(ct, "text/event-stream"): return t.readSSE(resp.Body, want) default: var m Message if err := json.NewDecoder(resp.Body).Decode(&m); err != nil { return nil, fmt.Errorf("http %s: decode: %w", method, err) } return &m, nil } } // readSSE consumes an SSE stream, dispatching peer notifications to the handler // and returning the response whose id matches want. func (t *httpTransport) readSSE(r io.Reader, want string) (*Message, error) { scanner := bufio.NewScanner(r) scanner.Buffer(make([]byte, 0, 64*1024), 16*1024*1024) var data strings.Builder flush := func() (*Message, bool, error) { if data.Len() == 0 { return nil, false, nil } payload := data.String() data.Reset() var m Message if err := json.Unmarshal([]byte(payload), &m); err != nil { return nil, false, nil // ignore non-JSON events } if m.IsResponse() && normalizeID(m.ID) == want { return &m, true, nil } if m.Method != "" { if hp := t.notify.Load(); hp != nil { (*hp)(&m) } } return nil, false, nil } for scanner.Scan() { line := scanner.Text() if line == "" { // event delimiter if m, ok, _ := flush(); ok { return m, nil } continue } if strings.HasPrefix(line, "data:") { data.WriteString(strings.TrimSpace(line[len("data:"):])) } } if err := scanner.Err(); err != nil { return nil, err } if m, ok, _ := flush(); ok { return m, nil } return nil, fmt.Errorf("sse stream ended without response for id %s", want) } func (t *httpTransport) Notify(ctx context.Context, method string, params any) error { note, err := NewNotification(method, params) if err != nil { return err } resp, err := t.post(ctx, note) if err != nil { return err } defer resp.Body.Close() io.Copy(io.Discard, io.LimitReader(resp.Body, 1<<20)) if resp.StatusCode >= 400 { return fmt.Errorf("http notify %s: status %d", method, resp.StatusCode) } return nil } func (t *httpTransport) Close() error { return nil }