Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
6 changes: 4 additions & 2 deletions AGENTS.md
Original file line number Diff line number Diff line change
Expand Up @@ -12,10 +12,11 @@
stdio (os.Stdin/os.Stdout)
json.Decoder / json.Encoder ← newline-delimited JSON-RPC 2.0
newline-delimited JSON-RPC 2.0 ← one message per line (MCP stdio framing)
Server.runWithIO() ← dispatch loop
├── unparseable line → -32700 / -32600, loop keeps serving
├── "initialize" → handshake
├── "tools/list" → registered tools metadata
├── "tools/call" → dispatch to Tool.Handler
Expand All @@ -41,7 +42,8 @@ Server.runWithIO() ← dispatch loop
- **Zero dependencies.** Only Go stdlib — `encoding/json`, `bufio`, `context`, `fmt`, `io`, `os`. Never add a third-party import to go.mod.
- **Interfaces over structs.** Handler signatures use `context.Context` + maps for extensibility. Future: typed generics.
- **Tests use pipes.** Integration tests simulate stdio with `io.Pipe()`. E2E tests spawn a real subprocess via `os/exec`.
- **Error codes follow JSON-RPC 2.0.** `-32601` = method not found, `-32602` = invalid params, `-32000` = application error.
- **Error codes follow JSON-RPC 2.0.** `-32700` = parse error (broken JSON), `-32600` = invalid request (well-formed JSON that is not a Request object, id `null`), `-32601` = method not found, `-32602` = invalid params, `-32000` = application error.
- **Bad input never kills the loop.** A malformed line is answered in-band and the server keeps serving; `RunWithIO` returns only on EOF or a read/write failure.
- **Go naming.** Exported types are PascalCase. Unexported internals are camelCase. Test functions are `TestXxx`.
- **Protocol version pinned.** `2024-11-05` hardcoded — update manually when MCP spec revs.

Expand Down
74 changes: 70 additions & 4 deletions gomcp/server.go
Original file line number Diff line number Diff line change
Expand Up @@ -9,8 +9,11 @@
package gomcp

import (
"bufio"
"bytes"
"context"
"encoding/json"
"errors"
"fmt"
"io"
"os"
Expand Down Expand Up @@ -80,17 +83,41 @@ func (s *Server) Run() error {

// RunWithIO starts the MCP server with custom I/O readers and writers,
// useful for testing with pipes or buffers.
//
// The input is newline-delimited JSON-RPC 2.0 (one message per line, the
// MCP stdio transport framing). A line that fails to parse never kills the
// dispatch loop: it is answered with a JSON-RPC error (-32700 for broken
// JSON, -32600 for well-formed JSON that is not a Request object) and the
// server keeps serving. RunWithIO returns only on clean EOF (nil), a read
// failure on r, or a write failure on w.
func (s *Server) RunWithIO(r io.Reader, w io.Writer) error {
decoder := json.NewDecoder(r)
encoder := json.NewEncoder(w)
br := bufio.NewReader(r)
// lineBuf is reused across messages. Handlers run synchronously before
// the next read, so nothing retains a reference to it across iterations.
var lineBuf []byte

for {
var req JSONRPCRequest
if err := decoder.Decode(&req); err != nil {
line, err := readMessage(br, lineBuf)
if err != nil {
if err == io.EOF {
return nil
}
return fmt.Errorf("decode error: %w", err)
return fmt.Errorf("read error: %w", err)
}
lineBuf = line[:0] // reclaim capacity for the next iteration

msg := bytes.TrimSpace(line)
if len(msg) == 0 {
continue // blank separator line
}

var req JSONRPCRequest
if derr := json.Unmarshal(msg, &req); derr != nil {
if werr := writeDecodeError(encoder, derr); werr != nil {
return fmt.Errorf("write error response: %w", werr)
}
continue
}

// Notifications have no ID — we silently consume them
Expand Down Expand Up @@ -138,6 +165,45 @@ func (s *Server) RunWithIO(r io.Reader, w io.Writer) error {
}
}

// readMessage reads one newline-terminated message from br into buf,
// returning the bytes without the trailing newline. A final message not
// terminated by EOF is still returned; a subsequent call then reports
// io.EOF. Lines longer than the reader's buffer are accumulated, so there
// is no message-size limit (matching the previous json.Decoder behavior).
func readMessage(br *bufio.Reader, buf []byte) ([]byte, error) {
buf = buf[:0]
for {
chunk, err := br.ReadSlice('\n')
buf = append(buf, chunk...)
switch err {
case nil:
return buf, nil
case bufio.ErrBufferFull:
continue // line longer than the buffer: keep accumulating
case io.EOF:
if len(buf) > 0 {
return buf, nil
}
return nil, io.EOF
default:
return nil, err
}
}
}

// writeDecodeError answers an unparseable inbound line per JSON-RPC 2.0:
// -32700 (Parse error) for broken JSON, -32600 (Invalid Request) for
// well-formed JSON that cannot be a Request object. The id is null — the
// request was never understood, so there is nothing to correlate against.
func writeDecodeError(encoder *json.Encoder, derr error) error {
code, message := -32700, "Parse error"
var typeErr *json.UnmarshalTypeError
if errors.As(derr, &typeErr) {
code, message = -32600, "Invalid Request"
}
return encoder.Encode(NewJSONRPCError(nil, code, message))
}

// handleInitialize responds to the MCP initialize handshake.
func (s *Server) handleInitialize(req JSONRPCRequest, encoder *json.Encoder) error {
s.mu.Lock()
Expand Down
164 changes: 152 additions & 12 deletions gomcp/server_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -524,32 +524,172 @@ func TestUnknownMethod(t *testing.T) {

// ----- Edge Cases -----

func TestDecodeError(t *testing.T) {
// ----- Malformed Input Resilience -----
//
// The dispatch loop must never die on a bad line: per JSON-RPC 2.0, broken
// JSON gets -32700 and well-formed JSON that is not a Request object gets
// -32600, and the server keeps serving subsequent messages.

// readResponses decodes n JSON-RPC messages from out.
func readResponses(t *testing.T, out *io.PipeReader, n int) []map[string]any {
t.Helper()
var msgs []map[string]any
for i := 0; i < n; i++ {
var m map[string]any
if err := json.NewDecoder(out).Decode(&m); err != nil {
t.Fatalf("failed to decode response %d/%d: %v", i+1, n, err)
}
msgs = append(msgs, m)
}
return msgs
}

func TestParseErrorKeepsServing(t *testing.T) {
inReader, inWriter := io.Pipe()
outReader, outWriter := io.Pipe()

srv := NewServer("test-server", "1.0.0")
done := make(chan error, 1)
go func() { done <- srv.RunWithIO(inReader, outWriter) }()

errCh := make(chan error, 1)
go func() {
errCh <- srv.RunWithIO(inReader, outWriter)
inWriter.Write([]byte(`this is not json` + "\n"))
inWriter.Write([]byte(`{"jsonrpc":"2.0","id":2,"method":"ping"}` + "\n"))
inWriter.Close()
}()

// Send malformed JSON
msgs := readResponses(t, outReader, 2)
errObj := msgs[0]["error"].(map[string]any)
if errObj["code"].(float64) != -32700 {
t.Errorf("expected code -32700, got %v", errObj["code"])
}
if msgs[0]["id"] != nil {
t.Errorf("expected id null for parse error, got %v", msgs[0]["id"])
}
if msgs[1]["id"].(float64) != 2 {
t.Errorf("server did not stay alive after parse error; got %v", msgs[1])
}

if err := <-done; err != nil {
t.Fatalf("RunWithIO: %v", err)
}
}

func TestTypeMismatchKeepsServing(t *testing.T) {
// Regression: {"jsonrpc":1,...} used to kill the dispatch loop with a
// fatal decode error, leaving the request unanswered and every later
// message stranded behind a blocked writer.
inReader, inWriter := io.Pipe()
outReader, outWriter := io.Pipe()

srv := NewServer("test-server", "1.0.0")
done := make(chan error, 1)
go func() { done <- srv.RunWithIO(inReader, outWriter) }()

go func() {
inWriter.Write([]byte(`this is not json` + "\n"))
inWriter.Write([]byte(`{"jsonrpc":1,"id":2,"method":"ping"}` + "\n"))
inWriter.Write([]byte(`{"jsonrpc":"2.0","id":3,"method":"ping"}` + "\n"))
inWriter.Close()
}()

// Drain the output pipe to avoid blocking
go io.Copy(io.Discard, outReader)
msgs := readResponses(t, outReader, 2)
errObj := msgs[0]["error"].(map[string]any)
if errObj["code"].(float64) != -32600 {
t.Errorf("expected code -32600, got %v", errObj["code"])
}
if msgs[0]["id"] != nil {
t.Errorf("expected id null, got %v", msgs[0]["id"])
}
if msgs[1]["id"].(float64) != 3 {
t.Errorf("server did not answer the valid follow-up request; got %v", msgs[1])
}

err := <-errCh
if err == nil {
t.Fatal("expected decode error, got nil")
if err := <-done; err != nil {
t.Fatalf("RunWithIO: %v", err)
}
if !strings.Contains(err.Error(), "decode error") {
t.Errorf("expected 'decode error' in message, got: %v", err)
}

func TestNonObjectRequestKeepsServing(t *testing.T) {
inReader, inWriter := io.Pipe()
outReader, outWriter := io.Pipe()

srv := NewServer("test-server", "1.0.0")
done := make(chan error, 1)
go func() { done <- srv.RunWithIO(inReader, outWriter) }()

go func() {
inWriter.Write([]byte(`[1,2,3]` + "\n"))
inWriter.Write([]byte(`{"jsonrpc":"2.0","id":7,"method":"ping"}` + "\n"))
inWriter.Close()
}()

msgs := readResponses(t, outReader, 2)
errObj := msgs[0]["error"].(map[string]any)
if errObj["code"].(float64) != -32600 {
t.Errorf("expected code -32600, got %v", errObj["code"])
}
if msgs[1]["id"].(float64) != 7 {
t.Errorf("server did not stay alive after non-object request; got %v", msgs[1])
}

if err := <-done; err != nil {
t.Fatalf("RunWithIO: %v", err)
}
}

func TestMultipleBadLinesThenGood(t *testing.T) {
inReader, inWriter := io.Pipe()
outReader, outWriter := io.Pipe()

srv := NewServer("test-server", "1.0.0")
done := make(chan error, 1)
go func() { done <- srv.RunWithIO(inReader, outWriter) }()

go func() {
inWriter.Write([]byte("{oops\n"))
inWriter.Write([]byte("\n")) // blank separator: skipped silently
inWriter.Write([]byte(`"just a string"` + "\n"))
inWriter.Write([]byte(`{"jsonrpc":"2.0","id":9,"method":"ping"}` + "\r\n"))
inWriter.Close()
}()

// "{oops" → -32700; blank line → nothing; string → -32600; CRLF ping → pong.
msgs := readResponses(t, outReader, 3)
if code := msgs[0]["error"].(map[string]any)["code"].(float64); code != -32700 {
t.Errorf("expected -32700 for broken JSON, got %v", code)
}
if code := msgs[1]["error"].(map[string]any)["code"].(float64); code != -32600 {
t.Errorf("expected -32600 for non-request JSON, got %v", code)
}
if msgs[2]["id"].(float64) != 9 {
t.Errorf("CRLF-terminated request was not served; got %v", msgs[2])
}

if err := <-done; err != nil {
t.Fatalf("RunWithIO: %v", err)
}
}

func TestFinalLineWithoutNewline(t *testing.T) {
inReader, inWriter := io.Pipe()
outReader, outWriter := io.Pipe()

srv := NewServer("test-server", "1.0.0")
done := make(chan error, 1)
go func() { done <- srv.RunWithIO(inReader, outWriter) }()

go func() {
inWriter.Write([]byte(`{"jsonrpc":"2.0","id":1,"method":"ping"}`)) // no trailing \n
inWriter.Close()
}()

msgs := readResponses(t, outReader, 1)
if msgs[0]["id"].(float64) != 1 {
t.Errorf("final unterminated message was dropped; got %v", msgs[0])
}

if err := <-done; err != nil {
t.Fatalf("RunWithIO should return nil on clean EOF, got: %v", err)
}
}

Expand Down
Loading