Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
19 commits
Select commit Hold shift + click to select a range
da158c4
fix(workspace): lock workspace during ExecOneShot to prevent races wi…
skevetter Aug 3, 2026
3fbf972
fix(workspace): wrap long t.Fatalf line in exec_test.go to satisfy ll…
skevetter Aug 3, 2026
485b37d
feat(mcp): bound concurrent workspace_exec/create/start operations wi…
skevetter Aug 3, 2026
5edf9cc
fix(up): --ide-launch=skip also skips IDE server install when --ide i…
skevetter Aug 3, 2026
1bfa505
fix(e2e): assert resolved IDE name, not a container path, for skip-la…
skevetter Aug 3, 2026
f9804d4
test(e2e): add MCP stdio JSON-RPC e2e coverage for workspace_list/wor…
skevetter Aug 3, 2026
373c55c
fix(provider): fsync parent directory after atomic rename for crash d…
skevetter Aug 3, 2026
2902ca7
fix(e2e): add missing !windows build tag to skip_launch_no_install.go
skevetter Aug 3, 2026
f7a0a02
fix: address CodeRabbit findings on semaphore gating and lock test co…
skevetter Aug 3, 2026
d98626b
chore: add plan doc and gitignore SDD scratch workspace
skevetter Aug 3, 2026
76af939
fix: extract semaphore-error check to reduce cyclop complexity
skevetter Aug 3, 2026
dd4f609
refactor: trim comments to non-obvious WHY, drop plan doc from PR
skevetter Aug 3, 2026
9f05c58
fix: correlate MCP e2e responses by ID and harden semaphore acquire
skevetter Aug 3, 2026
9e26371
fix: honor ctx cancellation in MCPClient.CallTool
skevetter Aug 3, 2026
9862da4
fix(lint): extract jsonRPCError to satisfy revive nested-structs
skevetter Aug 3, 2026
3d4d518
fix: trim ResolveDockerCommand comment to one line
skevetter Aug 3, 2026
6b2efdb
fix: replace unicode arrow with ASCII in ResolveDockerCommand comment
skevetter Aug 3, 2026
700a317
Merge branch 'main' into feat/mcp-exec-safety-and-coverage
skevetter Aug 3, 2026
4fe604f
style: update comments
skevetter Aug 4, 2026
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
33 changes: 33 additions & 0 deletions cmd/mcp/semaphore.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,33 @@
package mcp

import (
"context"
"fmt"
)

type opSemaphore struct {
slots chan struct{}
}

func newOpSemaphore(max int) *opSemaphore {
if max <= 0 {
max = 1
}
return &opSemaphore{slots: make(chan struct{}, max)}
}

func (s *opSemaphore) acquire(ctx context.Context) (func(), error) {
if err := ctx.Err(); err != nil {
return nil, fmt.Errorf("waiting for a free operation slot: %w", err)
}
select {
case s.slots <- struct{}{}:
if err := ctx.Err(); err != nil {
<-s.slots
return nil, fmt.Errorf("waiting for a free operation slot: %w", err)
}
return func() { <-s.slots }, nil
case <-ctx.Done():
return nil, fmt.Errorf("waiting for a free operation slot: %w", ctx.Err())
}
Comment on lines +23 to +32

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

🩺 Stability & Availability | 🟠 Major | ⚡ Quick win

Reject an already-canceled context before acquiring a slot.

When ctx is already canceled and a slot is free, both cases in Line 27 are ready. Go can select the channel send. The handler can then start a canceled workspace operation and consume a permit.

  • cmd/mcp/semaphore.go#L27-L32: Check ctx.Err() before the select. If cancellation is observed after the channel send, remove the token and return the context error.
  • cmd/mcp/semaphore_test.go#L82-L96: Add a test that cancels a context before calling acquire while semaphore capacity is available.
📍 Affects 2 files
  • cmd/mcp/semaphore.go#L27-L32 (this comment)
  • cmd/mcp/semaphore_test.go#L82-L96
🤖 Prompt for AI Agents
Verify each finding against current code. Fix only still-valid issues, skip the
rest with a brief reason, keep changes minimal, and validate.

In `@cmd/mcp/semaphore.go` around lines 27 - 32, Update acquire in
cmd/mcp/semaphore.go#L27-L32 to check ctx.Err() before selecting, and recheck
after acquiring a slot; if cancellation is detected, release the token and
return the context error. Add a test in cmd/mcp/semaphore_test.go#L82-L96
covering a pre-canceled context with available semaphore capacity, verifying no
permit remains consumed.

}
96 changes: 96 additions & 0 deletions cmd/mcp/semaphore_test.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,96 @@
package mcp

import (
"context"
"errors"
"sync/atomic"
"testing"
"time"
)

func TestOpSemaphore_LimitsConcurrency(t *testing.T) {
sem := newOpSemaphore(2)
var inFlight, maxInFlight atomic.Int32

track := func() {
cur := inFlight.Add(1)
for {
m := maxInFlight.Load()
if cur <= m || maxInFlight.CompareAndSwap(m, cur) {
break
}
}
time.Sleep(20 * time.Millisecond)
inFlight.Add(-1)
}

done := make(chan struct{}, 5)
for range 5 {
go func() {
release, err := sem.acquire(context.Background())
if err != nil {
t.Errorf("acquire failed: %v", err)
done <- struct{}{}
return
}
track()
release()
done <- struct{}{}
}()
}
for range 5 {
<-done
}

if got := maxInFlight.Load(); got > 2 {
t.Fatalf("max concurrent = %d, want <= 2", got)
}
}

func TestOpSemaphore_ReleaseAllowsNextAcquire(t *testing.T) {
sem := newOpSemaphore(1)
release1, err := sem.acquire(context.Background())
if err != nil {
t.Fatalf("first acquire failed: %v", err)
}

acquired := make(chan struct{})
go func() {
release2, err := sem.acquire(context.Background())
if err != nil {
t.Errorf("second acquire failed: %v", err)
return
}
close(acquired)
release2()
}()

select {
case <-acquired:
t.Fatal("second acquire succeeded while first slot was held")
case <-time.After(50 * time.Millisecond):
}

release1()
select {
case <-acquired:
case <-time.After(time.Second):
t.Fatal("second acquire never succeeded after release")
}
}

func TestOpSemaphore_AcquireRespectsContextCancel(t *testing.T) {
sem := newOpSemaphore(1)
release, err := sem.acquire(context.Background())
if err != nil {
t.Fatalf("first acquire failed: %v", err)
}
defer release()

ctx, cancel := context.WithTimeout(context.Background(), 20*time.Millisecond)
defer cancel()
_, err = sem.acquire(ctx)
if !errors.Is(err, context.DeadlineExceeded) {
t.Fatalf("acquire error = %v, want context deadline exceeded", err)
}
}
23 changes: 17 additions & 6 deletions cmd/mcp/serve.go
Original file line number Diff line number Diff line change
Expand Up @@ -21,6 +21,9 @@ type ServeCmd struct {
ExecTimeoutDefault time.Duration
ExecTimeoutMax time.Duration
ExecOutputCap int
MaxConcurrentOps int

opSem *opSemaphore
}

// NewServeCmd builds the `serve` subcommand.
Expand Down Expand Up @@ -54,14 +57,21 @@ func NewServeCmd(globalFlags *flags.GlobalFlags) *cobra.Command {
100*1024,
"Per-stream byte cap for workspace_exec output; excess is replaced with a truncation marker",
),
cliflags.Int(
&cmd.MaxConcurrentOps,
names.MaxConcurrentOps,
8,
"Maximum number of concurrent workspace_exec/workspace_create/workspace_start "+
"operations; excess calls wait for a free slot",
),
)
return cobraCmd
}

// Run wires up the MCP server and serves over stdio until ctx is cancelled.
func (cmd *ServeCmd) Run(ctx context.Context) error {
log.Debugf("starting MCP server (timeout default=%s max=%s cap=%dB)",
cmd.ExecTimeoutDefault, cmd.ExecTimeoutMax, cmd.ExecOutputCap)
log.Debugf("starting MCP server (timeout default=%s max=%s cap=%dB maxops=%d)",
cmd.ExecTimeoutDefault, cmd.ExecTimeoutMax, cmd.ExecOutputCap, cmd.MaxConcurrentOps)

// Reserve real stdout for the JSON-RPC frame; redirect os.Stdout to stderr
// so any stray write elsewhere in the process can't corrupt the transport.
Expand All @@ -79,13 +89,14 @@ func (cmd *ServeCmd) Run(ctx context.Context) error {
Version: version.GetVersion(),
}, nil)

cmd.registerTools(server)
cmd.opSem = newOpSemaphore(cmd.MaxConcurrentOps)
cmd.registerTools(server, cmd.opSem)

return server.Run(ctx, transport)
}

func (cmd *ServeCmd) registerTools(s *sdkmcp.Server) {
registerWorkspaceTools(s, cmd.GlobalFlags)
registerExecTool(s, cmd)
func (cmd *ServeCmd) registerTools(s *sdkmcp.Server, sem *opSemaphore) {
registerWorkspaceTools(s, cmd.GlobalFlags, sem)
registerExecTool(s, cmd, sem)
registerProviderTools(s, cmd.GlobalFlags)
}
82 changes: 81 additions & 1 deletion cmd/mcp/serve_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -2,6 +2,7 @@ package mcp

import (
"context"
"strings"
"testing"
"time"

Expand All @@ -19,7 +20,7 @@ func TestServer_ListsAllTools(t *testing.T) {
server := sdkmcp.NewServer(&sdkmcp.Implementation{Name: "devsy-test", Version: "test"}, nil)
g := &flags.GlobalFlags{}
serveCmd := &ServeCmd{GlobalFlags: g, ExecOutputCap: 1024}
serveCmd.registerTools(server)
serveCmd.registerTools(server, newOpSemaphore(8))

clientTransport, serverTransport := sdkmcp.NewInMemoryTransports()

Expand Down Expand Up @@ -57,3 +58,82 @@ func TestServer_ListsAllTools(t *testing.T) {
t.Errorf("expected %d tools, got %d: %+v", len(wantNames), len(tools.Tools), have)
}
}

func TestServer_WorkspaceExecRespectsSemaphore(t *testing.T) {
home := t.TempDir()
t.Setenv("DEVSY_HOME", home)

ctx, cancel := context.WithTimeout(context.Background(), 10*time.Second)
defer cancel()

server := sdkmcp.NewServer(&sdkmcp.Implementation{Name: "devsy-test", Version: "test"}, nil)
g := &flags.GlobalFlags{}
serveCmd := &ServeCmd{GlobalFlags: g, ExecOutputCap: 1024, MaxConcurrentOps: 1}
sem := newOpSemaphore(serveCmd.MaxConcurrentOps)
serveCmd.registerTools(server, sem)

clientTransport, serverTransport := sdkmcp.NewInMemoryTransports()

serverErr := make(chan error, 1)
go func() {
serverErr <- server.Run(ctx, serverTransport)
}()

client := sdkmcp.NewClient(&sdkmcp.Implementation{Name: "test-client", Version: "0"}, nil)
session, err := client.Connect(ctx, clientTransport, nil)
if err != nil {
t.Fatalf("connect: %v", err)
}
t.Cleanup(func() { _ = session.Close() })

execArgs := map[string]any{
"name": "some-workspace",
"command": []string{"echo", "hi"},
}

release, err := sem.acquire(context.Background())
if err != nil {
t.Fatalf("acquire: %v", err)
}

blockedCtx, blockedCancel := context.WithTimeout(ctx, 200*time.Millisecond)
defer blockedCancel()
_, callErr := session.CallTool(blockedCtx, &sdkmcp.CallToolParams{
Name: "workspace_exec",
Arguments: execArgs,
})
if callErr == nil {
t.Fatal("expected workspace_exec call to fail while the only semaphore slot is held")
}

release()

res, callErr := session.CallTool(ctx, &sdkmcp.CallToolParams{
Name: "workspace_exec",
Arguments: execArgs,
})
if callErr != nil {
t.Fatalf(
"expected workspace_exec call to reach the handler after release, got transport error: %v",
callErr,
)
}
assertNotSemaphoreError(t, res)
}

// A non-semaphore error (e.g. workspace not found) is expected here.
func assertNotSemaphoreError(t *testing.T, res *sdkmcp.CallToolResult) {
t.Helper()
if !res.IsError {
return
}
var msg string
for _, c := range res.Content {
if tc, ok := c.(*sdkmcp.TextContent); ok {
msg = tc.Text
}
}
if strings.Contains(msg, "waiting for a free operation slot") {
t.Fatalf("workspace_exec failed due to the semaphore even after release: %s", msg)
}
}
Comment thread
coderabbitai[bot] marked this conversation as resolved.
9 changes: 8 additions & 1 deletion cmd/mcp/tools_exec.go
Original file line number Diff line number Diff line change
Expand Up @@ -42,7 +42,7 @@ type execOutput struct {
Error *ErrorPayload `json:"error,omitempty"`
}

func registerExecTool(s *sdkmcp.Server, cmd *ServeCmd) {
func registerExecTool(s *sdkmcp.Server, cmd *ServeCmd, sem *opSemaphore) {
sdkmcp.AddTool(s, &sdkmcp.Tool{
Name: "workspace_exec",
Description: "Run a one-shot command in a running workspace container. The " +
Expand All @@ -59,6 +59,13 @@ func registerExecTool(s *sdkmcp.Server, cmd *ServeCmd) {
if len(in.Command) == 0 {
return errorResult(fmt.Errorf("command is required")), execOutput{}, nil
}

release, err := sem.acquire(ctx)
if err != nil {
return errorResult(err), execOutput{}, nil
}
defer release()

stdout := NewBoundedBuffer(cmd.ExecOutputCap)
stderr := NewBoundedBuffer(cmd.ExecOutputCap)

Expand Down
Loading
Loading