| 1 | package plugin |
| 2 | |
| 3 | import ( |
| 4 | "encoding/json" |
| 5 | "net/http" |
| 6 | "net/http/httptest" |
| 7 | "testing" |
| 8 | "time" |
| 9 | ) |
| 10 | |
| 11 | func TestHTTPTransportCloseCancelsActiveRequest(t *testing.T) { |
| 12 | requestStarted := make(chan struct{}) |
| 13 | requestCancelled := make(chan struct{}) |
| 14 | server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { |
| 15 | if r.Method == http.MethodGet { |
| 16 | w.WriteHeader(http.StatusMethodNotAllowed) |
| 17 | return |
| 18 | } |
| 19 | var req struct { |
| 20 | ID *int `json:"id"` |
| 21 | Method string `json:"method"` |
| 22 | } |
| 23 | if err := json.NewDecoder(r.Body).Decode(&req); err != nil { |
| 24 | http.Error(w, "bad body", http.StatusBadRequest) |
| 25 | return |
| 26 | } |
| 27 | switch req.Method { |
| 28 | case "initialize": |
| 29 | w.Header().Set("Mcp-Session-Id", "close-active") |
| 30 | writeHTTPRPCResult(w, req.ID, map[string]any{ |
| 31 | "protocolVersion": testLegacyProtocolVersion, |
| 32 | "serverInfo": map[string]any{"name": "close-active", "version": "1"}, |
| 33 | }) |
| 34 | case "notifications/initialized": |
| 35 | w.WriteHeader(http.StatusAccepted) |
| 36 | case "tools/list": |
| 37 | writeHTTPRPCResult(w, req.ID, map[string]any{"tools": []any{}}) |
| 38 | case "tools/call": |
| 39 | close(requestStarted) |
| 40 | <-r.Context().Done() |
| 41 | close(requestCancelled) |
| 42 | default: |
| 43 | http.Error(w, "unknown method", http.StatusBadRequest) |
| 44 | } |
| 45 | })) |
| 46 | defer server.Close() |
| 47 | |
| 48 | transport, err := newHTTPTransport(Spec{Name: "close-active", Type: "http", URL: server.URL}) |
| 49 | if err != nil { |
| 50 | t.Fatal(err) |
| 51 | } |
| 52 | if _, err := transport.call(t.Context(), "tools/list", map[string]any{}); err != nil { |
| 53 | t.Fatalf("tools/list: %v", err) |
| 54 | } |
| 55 | callDone := make(chan error, 1) |
| 56 | go func() { |
| 57 | _, err := transport.call(t.Context(), "tools/call", map[string]any{"name": "work", "arguments": map[string]any{}}) |
| 58 | callDone <- err |
| 59 | }() |
| 60 | select { |
| 61 | case <-requestStarted: |
| 62 | case <-time.After(time.Second): |
| 63 | t.Fatal("tools/call did not reach the server") |
| 64 | } |
| 65 | |
| 66 | closed := make(chan struct{}) |
| 67 | go func() { |
| 68 | transport.close() |
| 69 | close(closed) |
| 70 | }() |
| 71 | select { |
| 72 | case <-closed: |
| 73 | case <-time.After(4 * time.Second): |
| 74 | t.Fatal("transport close did not finish within its cleanup budget") |
| 75 | } |
| 76 | select { |
| 77 | case <-requestCancelled: |
| 78 | case <-time.After(time.Second): |
| 79 | t.Fatal("transport close did not cancel the active HTTP request") |
| 80 | } |
| 81 | select { |
| 82 | case err := <-callDone: |
| 83 | if err == nil { |
| 84 | t.Fatal("cancelled tools/call unexpectedly succeeded") |
| 85 | } |
| 86 | case <-time.After(time.Second): |
| 87 | t.Fatal("cancelled tools/call did not return") |
| 88 | } |
| 89 | } |
| 90 |