diff options
Diffstat (limited to 'internal/mcp')
| -rw-r--r-- | internal/mcp/handlers_prompt_test.go | 5 | ||||
| -rw-r--r-- | internal/mcp/server.go | 16 | ||||
| -rw-r--r-- | internal/mcp/server_test.go | 21 |
3 files changed, 36 insertions, 6 deletions
diff --git a/internal/mcp/handlers_prompt_test.go b/internal/mcp/handlers_prompt_test.go index f3a9d87..d3f34e2 100644 --- a/internal/mcp/handlers_prompt_test.go +++ b/internal/mcp/handlers_prompt_test.go @@ -3,6 +3,7 @@ package mcp import ( "bytes" + "context" "encoding/json" "fmt" "io" @@ -882,14 +883,14 @@ func TestServer_Run_InvalidJSON(t *testing.T) { // Run in background done := make(chan error, 1) go func() { - done <- server.Run() + done <- server.Run(context.Background()) }() // Wait for processing to complete select { case <-done: case <-time.After(2 * time.Second): - t.Fatal("server.Run() did not return in time") + t.Fatal("server.Run(context.Background()) did not return in time") } // Should have written error response diff --git a/internal/mcp/server.go b/internal/mcp/server.go index f8042ac..e3c8723 100644 --- a/internal/mcp/server.go +++ b/internal/mcp/server.go @@ -3,6 +3,7 @@ package mcp import ( "bufio" + "context" "encoding/json" "errors" "fmt" @@ -67,9 +68,20 @@ func NewServer(r io.Reader, w io.Writer, logger *log.Logger, store promptstore.P } // Run starts the server main loop, reading and dispatching requests. -// Returns on EOF or fatal error, after waiting for all in-flight handlers. -func (s *Server) Run() error { +// Returns on EOF, on a cancelled ctx, or on a fatal error, after waiting for +// all in-flight handlers. +// +// ctx ties the serve loop to the process lifecycle: once it is cancelled (e.g. +// SIGINT/SIGTERM at main) the loop stops accepting further requests after the +// current blocking read returns and drains outstanding handlers before exiting. +func (s *Server) Run(ctx context.Context) error { for { + // Stop promptly when the caller cancels (shutdown signal); the loop + // otherwise blocks in readMessage until the next request or EOF. + if ctx != nil && ctx.Err() != nil { + s.inflight.Wait() + return nil + } body, err := s.readMessage() if errors.Is(err, io.EOF) { s.inflight.Wait() // drain handlers before signalling callers diff --git a/internal/mcp/server_test.go b/internal/mcp/server_test.go index 00a7823..8374a24 100644 --- a/internal/mcp/server_test.go +++ b/internal/mcp/server_test.go @@ -3,6 +3,7 @@ package mcp import ( "bytes" + "context" "encoding/json" "fmt" "io" @@ -397,12 +398,28 @@ func TestServer_Run(t *testing.T) { logger := log.New(io.Discard, "", 0) server := NewServer(inBuf, outBuf, logger, store, nil) - err := server.Run() + err := server.Run(context.Background()) if err != nil { t.Errorf("Run() error = %v, want nil on EOF", err) } }) + t.Run("exits on cancelled context", func(t *testing.T) { + store := &mockPromptStore{prompts: make(map[string]*promptstore.Prompt)} + // A pipe with no data would otherwise block in readMessage; an + // already-cancelled context makes Run return before reading. + pr, _ := io.Pipe() + outBuf := &bytes.Buffer{} + logger := log.New(io.Discard, "", 0) + server := NewServer(pr, outBuf, logger, store, nil) + + ctx, cancel := context.WithCancel(context.Background()) + cancel() + if err := server.Run(ctx); err != nil { + t.Errorf("Run() error = %v, want nil on cancelled ctx", err) + } + }) + t.Run("processes initialize request", func(t *testing.T) { store := &mockPromptStore{prompts: make(map[string]*promptstore.Prompt)} inBuf := &bytes.Buffer{} @@ -434,7 +451,7 @@ func TestServer_Run(t *testing.T) { // so Run() will complete naturally once it has written the response. done := make(chan error, 1) go func() { - done <- server.Run() + done <- server.Run(context.Background()) }() // Wait for Run() to return (signalled by EOF on the input buffer). |
