Files
ombi-mcp/internal/transport/sse.go
gronod 49bc0775de
Build and publish / Test and build (windows) (push) Failing after 11s
Build and publish / Test and build (darwin) (push) Successful in 1m32s
Build and publish / Test and build (linux) (push) Canceled after 0s
Build and publish / Build and publish Docker image (push) Canceled after 0s
Add HTTP/SSE transport alongside stdio (Phase 09)
Network MCP clients (e.g. browser-based frontends) can't spawn stdio
subprocesses, so the server now optionally serves the 2024-11-05 MCP
HTTP/SSE transport via MCP_TRANSPORT=sse and MCP_PORT (default 8080),
backed by the SDK's SSEHandler behind permissive CORS. stdio remains
the default and is behaviourally unchanged.
2026-09-18 21:13:09 +01:00

79 lines
2.4 KiB
Go

// Package transport wires the MCP server onto network transports.
// The stdio transport is driven directly from the entrypoint; this
// package owns the HTTP/SSE transport defined by the 2024-11-05 MCP
// specification.
package transport
import (
"context"
"errors"
"fmt"
"net"
"net/http"
"time"
"github.com/modelcontextprotocol/go-sdk/mcp"
)
// NewSSEHandler builds an http.Handler serving MCP over HTTP/SSE:
//
// GET /sse opens a session and streams server->client messages
// as SSE events; the first event is an 'endpoint' event
// carrying the session's POST URL (?sessionid=...).
// POST /sse or /messages?sessionid=... accepts client->server JSON-RPC
// messages for an established session.
//
// A permissive CORS middleware wraps the mux because network MCP servers
// are frequently queried by browser-based clients.
func NewSSEHandler(srv *mcp.Server) http.Handler {
mcpHandler := mcp.NewSSEHandler(func(*http.Request) *mcp.Server { return srv }, nil)
mux := http.NewServeMux()
mux.Handle("/sse", mcpHandler)
mux.Handle("/messages", mcpHandler)
return cors(mux)
}
func cors(next http.Handler) http.Handler {
return http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
h := w.Header()
h.Set("Access-Control-Allow-Origin", "*")
h.Set("Access-Control-Allow-Methods", "GET, POST, OPTIONS")
h.Set("Access-Control-Allow-Headers", "Content-Type, Accept, Authorization")
if r.Method == http.MethodOptions {
w.WriteHeader(http.StatusNoContent)
return
}
next.ServeHTTP(w, r)
})
}
// RunSSE serves the MCP SSE transport on the given port until ctx is
// cancelled, then shuts down gracefully.
func RunSSE(ctx context.Context, srv *mcp.Server, port int) error {
httpSrv := &http.Server{
Addr: fmt.Sprintf(":%d", port),
Handler: NewSSEHandler(srv),
ReadHeaderTimeout: 10 * time.Second,
}
lc := net.ListenConfig{}
ln, err := lc.Listen(ctx, "tcp", httpSrv.Addr)
if err != nil {
return fmt.Errorf("listen %s: %w", httpSrv.Addr, err)
}
errCh := make(chan error, 1)
go func() {
errCh <- httpSrv.Serve(ln)
}()
select {
case <-ctx.Done():
shutdownCtx, cancel := context.WithTimeout(context.Background(), 10*time.Second)
defer cancel()
return httpSrv.Shutdown(shutdownCtx)
case err := <-errCh:
if errors.Is(err, http.ErrServerClosed) {
return nil
}
return err
}
}