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
10 changes: 8 additions & 2 deletions .github/workflows/release.yml
Original file line number Diff line number Diff line change
Expand Up @@ -74,6 +74,9 @@ jobs:
with:
go-version-file: go.mod
cache: false
- uses: actions/setup-python@v6
with:
python-version: '3.12'
- uses: actions/cache@v6
with:
path: |
Expand All @@ -94,6 +97,9 @@ jobs:
run: |
sudo apt-get update
sudo apt-get install -y build-essential pkg-config libssl-dev pigz
- uses: mlugg/setup-zig@d1434d08867e3ee9daa34448df10607b98908d29
with:
version: 0.14.0
- name: Cache pinned release downloads
uses: actions/cache@v6
with:
Expand Down Expand Up @@ -123,9 +129,9 @@ jobs:
CORE_DISTRIBUTION_OFFLINE: ${{ (github.event_name == 'push' || inputs.offline) && '1' || '0' }}
run: |
inputs="$HOME/.oac/build/release-inputs/inputs.json"
CODEX_CLI_DIR="$(python3 -c 'import json,sys; print(json.load(open(sys.argv[1]))["codex"])' "$inputs")"
CODEX_HARNESS_BUILD_DIR="$(python3 -c 'import json,sys; print(json.load(open(sys.argv[1]))["codex"])' "$inputs")"
MCODE_HARNESS_BUILD_DIR="$(python3 -c 'import json,sys; print(json.load(open(sys.argv[1]))["mcode"])' "$inputs")"
export CODEX_CLI_DIR MCODE_HARNESS_BUILD_DIR
export CODEX_HARNESS_BUILD_DIR MCODE_HARNESS_BUILD_DIR
export CORE_DISTRIBUTION_RELEASE_BASE_URL="https://github.com/$RELEASE_REPOSITORY/releases/download/$RELEASE_TAG"
bash scripts/build-core-distribution.sh
mkdir -p "$HOME/.oac/build/release-upload"
Expand Down
9 changes: 5 additions & 4 deletions apps/daemon/internal/agent/claudesdk/executor.go
Original file line number Diff line number Diff line change
Expand Up @@ -23,11 +23,12 @@ type executor struct {
done chan struct{}
invalid bool
nativeID string
stopMCP func(context.Context, []string) error
}

// startExecutor checks the installed bridge against probe, starts it through
// run and waits until it is ready for Turns.
func startExecutor(ctx context.Context, checked *runtimeCheckCache, probe Config, start startRequest, run func() (*session, error)) (agent.Executor, error) {
func startExecutor(ctx context.Context, checked *runtimeCheckCache, probe Config, start startRequest, stopMCP func(context.Context, []string) error, run func() (*session, error)) (agent.Executor, error) {
info, err := checked.check(ctx, probe)
if err != nil {
return nil, err
Expand All @@ -39,7 +40,7 @@ func startExecutor(ctx context.Context, checked *runtimeCheckCache, probe Config
if err != nil {
return nil, err
}
e := &executor{base: base, start: start, ready: make(chan error, 1), done: make(chan struct{}), nativeID: start.Resume}
e := &executor{base: base, start: start, ready: make(chan error, 1), done: make(chan struct{}), nativeID: start.Resume, stopMCP: stopMCP}
go e.read()
if err = e.write(start); err == nil {
select {
Expand All @@ -60,7 +61,7 @@ func startExecutor(ctx context.Context, checked *runtimeCheckCache, probe Config
}

func validateExecutorFeatures(info RuntimeInfo, start startRequest) error {
if info.Protocol != 3 {
if info.Protocol != 4 {
return errors.New("claudesdk: Executor bridge protocol is unavailable")
}
if start.Workspace != nil && !info.supportsWorkspacePreparation() {
Expand Down Expand Up @@ -114,7 +115,7 @@ func (e *executor) read() {
break
}
if !ready {
if event.Type != "executor_ready" || event.Protocol != 3 || event.TurnID != "" {
if event.Type != "executor_ready" || event.Protocol != 4 || event.TurnID != "" {
if event.Type == "error" {
readyFailure = bridgeFailure(event.Code)
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -74,7 +74,7 @@ func startSingleTurn(ctx context.Context, config testBridge, req proto.PromptReq
}

func helperTurn(scanner *bufio.Scanner) (func(bridgeEvent), func()) {
_ = json.NewEncoder(os.Stdout).Encode(bridgeEvent{Type: "executor_ready", Protocol: 3})
_ = json.NewEncoder(os.Stdout).Encode(bridgeEvent{Type: "executor_ready", Protocol: 4})
if !scanner.Scan() {
os.Exit(4)
}
Expand Down
23 changes: 21 additions & 2 deletions apps/daemon/internal/agent/claudesdk/executor_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -38,7 +38,7 @@ func runPersistentExecutorHelper() {
_ = file.Close()
}
printlnReport := json.NewEncoder(os.Stdout)
_ = printlnReport.Encode(RuntimeInfo{Type: "runtime_ready", Protocol: 3, Node: "fixture", SDK: "fixture", MCP: "fixture", Native: "fixture"})
_ = printlnReport.Encode(RuntimeInfo{Type: "runtime_ready", Protocol: 4, Node: "fixture", SDK: "fixture", MCP: "fixture", Native: "fixture"})
return
}
scanner := bufio.NewScanner(os.Stdin)
Expand All @@ -51,7 +51,7 @@ func runPersistentExecutorHelper() {
}
_ = os.WriteFile(filepath.Join(root, "prepared"), []byte(strconv.Itoa(os.Getpid())), 0600)
encode := func(event bridgeEvent) { _ = json.NewEncoder(os.Stdout).Encode(event) }
encode(bridgeEvent{Type: "executor_ready", Protocol: 3})
encode(bridgeEvent{Type: "executor_ready", Protocol: 4})
if os.Getenv("SDK_EXECUTOR_MODE") == "block" {
time.Sleep(time.Hour)
return
Expand Down Expand Up @@ -117,8 +117,27 @@ func runPersistentExecutorHelper() {
encode(bridgeEvent{Type: "delta", TurnID: active, ItemID: "message", Delta: "steer-written"})
case "turn_cancel":
if active == command.TurnID {
if strings.HasPrefix(os.Getenv("SDK_EXECUTOR_MODE"), "mcp_") {
label := "target"
if os.Getenv("SDK_EXECUTOR_MODE") == "mcp_http" {
label = "http"
}
encode(bridgeEvent{Type: "mcp_stop", TurnID: active, Servers: []string{label}})
continue
}
settle(true)
}
case "mcp_stopped":
var response struct {
Confirmed bool `json:"confirmed"`
}
if command.TurnID != active || json.Unmarshal(scanner.Bytes(), &response) != nil {
os.Exit(6)
}
if !response.Confirmed {
_ = os.Setenv("SDK_EXECUTOR_MODE", "unknown_cancel")
}
settle(true)
}
}
}
Expand Down
29 changes: 29 additions & 0 deletions apps/daemon/internal/agent/claudesdk/executor_turn.go
Original file line number Diff line number Diff line change
Expand Up @@ -60,6 +60,7 @@ func (s *session) runTurn(start startRequest, out chan<- proto.Envelope) {
cancelled := false
reusable := false
reason := "bridge_interrupted"
mcpStopRequested := false
for raw := range s.frames {
var event bridgeEvent
if err := json.Unmarshal(raw, &event); err != nil || event.TurnID != s.runID {
Expand All @@ -85,6 +86,34 @@ func (s *session) runTurn(start startRequest, out chan<- proto.Envelope) {
break
}
switch event.Type {
case "mcp_stop":
if mcpStopRequested || s.owner.stopMCP == nil || !start.validStdioServers(event.Servers) {
failure = fmt.Errorf("claudesdk: invalid MCP stop request")
s.invalidate()
break
}
select {
case <-s.cancelOutput:
default:
failure = fmt.Errorf("claudesdk: unsolicited MCP stop request")
s.invalidate()
}
if failure != nil {
break
}
mcpStopRequested = true
stopErr := s.owner.stopMCP(s.process.Context(), event.Servers)
if stopErr != nil {
s.settlementErr = fmt.Errorf("claudesdk: MCP process-scope settlement is unconfirmed")
}
if err := s.owner.write(struct {
Type string `json:"type"`
TurnID string `json:"turn_id"`
Confirmed bool `json:"confirmed"`
}{"mcp_stopped", s.runID, stopErr == nil}); err != nil {
failure = fmt.Errorf("claudesdk: MCP stop receipt delivery failed")
s.invalidate()
}
case "command_observation":
if err := commands.receive(event, start, s.inputSessionID(), emit); err != nil {
failure = err
Expand Down
109 changes: 109 additions & 0 deletions apps/daemon/internal/agent/claudesdk/mcp_cancellation_test.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,109 @@
//go:build unix

package claudesdk

import (
"context"
"errors"
"slices"
"sync/atomic"
"testing"

"github.com/MiniMax-AI/OpenAgentCore/apps/daemon/internal/agent"
)

func TestExecutorMCPStopWaitsForScopeAndRetainsTurnOwnership(t *testing.T) {
config, req := persistentConfig(t, "mcp_cancel")
resource, err := config.factory()(t.Context(), prepared(t, req))
if err != nil {
t.Fatal(err)
}
defer resource.Close(context.Background())
owner := resource.(*executor)
owner.start.Workspace = &workspaceProfile{MCP: []environmentMCPServer{
{mcpHTTPServer: mcpHTTPServer{ServerLabel: "target"}, Command: agent.ViewAlias(0)},
{mcpHTTPServer: mcpHTTPServer{ServerLabel: "healthy"}, Command: agent.ViewAlias(1)},
}}
entered, release := make(chan []string, 1), make(chan struct{})
var calls atomic.Int32
owner.stopMCP = func(ctx context.Context, labels []string) error {
calls.Add(1)
entered <- slices.Clone(labels)
select {
case <-release:
return nil
case <-ctx.Done():
return ctx.Err()
}
}
turn, out := consumeExecutorTurn(t, owner, "cancelled", "wait")
<-out
waitCtx, stopWait := context.WithCancel(t.Context())
first := make(chan error, 1)
go func() { first <- turn.Cancel(waitCtx) }()
if labels := <-entered; !slices.Equal(labels, []string{"target"}) {
t.Fatal(labels)
}
stopWait()
if err := <-first; !errors.Is(err, context.Canceled) {
t.Fatalf("cancel waiter = %v", err)
}
if _, err := turn.AwaitSettlement(waitCtx); !errors.Is(err, context.Canceled) {
t.Fatalf("scope still open: %v", err)
}
joined := make(chan error, 1)
go func() { joined <- turn.Cancel(t.Context()) }()
close(release)
if err := <-joined; err != nil {
t.Fatal(err)
}
awaitExecutorTurn(t, turn, out, true)
next, nextOut := consumeExecutorTurn(t, owner, "successor", "hello")
if err := turn.Cancel(t.Context()); err != nil {
t.Fatal(err)
}
awaitExecutorTurn(t, next, nextOut, true)
if calls.Load() != 1 {
t.Fatalf("scope stops = %d", calls.Load())
}
}

func TestExecutorMCPStopFailureAndHTTPRemainUnconfirmed(t *testing.T) {
for _, mode := range []string{"mcp_cancel", "mcp_http"} {
t.Run(mode, func(t *testing.T) {
config, req := persistentConfig(t, mode)
resource, err := config.factory()(t.Context(), prepared(t, req))
if err != nil {
t.Fatal(err)
}
defer resource.Close(context.Background())
owner := resource.(*executor)
owner.start.Workspace = &workspaceProfile{MCP: []environmentMCPServer{
{mcpHTTPServer: mcpHTTPServer{ServerLabel: "target"}, Command: agent.ViewAlias(0)},
{mcpHTTPServer: mcpHTTPServer{ServerLabel: "http", ServerURL: "http://gateway/mcp/http"}},
}}
var calls atomic.Int32
owner.stopMCP = func(context.Context, []string) error {
calls.Add(1)
return errors.New("scope not closed")
}
turn, out := consumeExecutorTurn(t, owner, "cancelled", "wait")
<-out
if err := turn.Cancel(t.Context()); err == nil {
t.Fatal("unconfirmed scope cancellation succeeded")
}
for range out {
}
if settled, err := turn.AwaitSettlement(t.Context()); err == nil || settled.Reusable {
t.Fatalf("settlement = %+v, %v", settled, err)
}
want := int32(1)
if mode == "mcp_http" {
want = 0
}
if calls.Load() != want {
t.Fatalf("scope stops = %d, want %d", calls.Load(), want)
}
})
}
}
22 changes: 22 additions & 0 deletions apps/daemon/internal/agent/claudesdk/mcp_environment.go
Original file line number Diff line number Diff line change
Expand Up @@ -42,3 +42,25 @@ func (start startRequest) declaredMCP() []mcpHTTPServer {
}
return servers
}

func (start startRequest) validStdioServers(labels []string) bool {
if start.Workspace == nil || len(labels) == 0 {
return false
}
seen := make(map[string]bool, len(labels))
for _, label := range labels {
if seen[label] {
return false
}
for _, server := range start.Workspace.MCP {
if server.ServerLabel == label && server.Command != "" {
seen[label] = true
break
}
}
if !seen[label] {
return false
}
}
return true
}
Original file line number Diff line number Diff line change
Expand Up @@ -72,7 +72,7 @@ func runPreparationHelper() {
features = append(features, "workspace_structured_output")
}
}
_ = json.NewEncoder(os.Stdout).Encode(RuntimeInfo{Type: "runtime_ready", Protocol: 3, Node: "fixture", SDK: "fixture", MCP: "fixture", Native: "fixture", Features: features})
_ = json.NewEncoder(os.Stdout).Encode(RuntimeInfo{Type: "runtime_ready", Protocol: 4, Node: "fixture", SDK: "fixture", MCP: "fixture", Native: "fixture", Features: features})
return
}
state := os.Getenv("CLAUDE_CONFIG_DIR")
Expand Down Expand Up @@ -106,7 +106,7 @@ func runPreparationHelper() {
time.Sleep(time.Millisecond)
}
}
emit(bridgeEvent{Type: "executor_ready", Protocol: 3})
emit(bridgeEvent{Type: "executor_ready", Protocol: 4})
if !scanner.Scan() {
return
}
Expand Down
2 changes: 1 addition & 1 deletion apps/daemon/internal/agent/claudesdk/readiness.go
Original file line number Diff line number Diff line change
Expand Up @@ -132,7 +132,7 @@ func CheckRuntime(ctx context.Context, config Config) (RuntimeInfo, error) {
return RuntimeInfo{}, fmt.Errorf("claudesdk: runtime check failed")
}
var info RuntimeInfo
if json.Unmarshal(raw, &info) != nil || info.Type != "runtime_ready" || info.Protocol != 3 ||
if json.Unmarshal(raw, &info) != nil || info.Type != "runtime_ready" || info.Protocol != 4 ||
info.Node == "" || info.SDK == "" || info.MCP == "" || info.Native == "" {
return RuntimeInfo{}, fmt.Errorf("claudesdk: invalid runtime readiness report")
}
Expand Down
6 changes: 3 additions & 3 deletions apps/daemon/internal/agent/claudesdk/readiness_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -15,7 +15,7 @@ import (
"github.com/MiniMax-AI/OpenAgentCore/internal/agentdaemon/proto"
)

const readyReport = `{"type":"runtime_ready","protocol":3,"node":"22.22.2","sdk":"0.3.269","mcp":"1.30.0","native":"2.1.269 (Claude Code)"}`
const readyReport = `{"type":"runtime_ready","protocol":4,"node":"22.22.2","sdk":"0.3.269","mcp":"1.30.0","native":"2.1.269 (Claude Code)"}`

func TestRequiredMCPNeedsQualifiedRuntime(t *testing.T) {
root := t.TempDir()
Expand Down Expand Up @@ -107,11 +107,11 @@ func runReadinessHelper() {
case "ready":
_, _ = fmt.Fprintln(os.Stdout, readyReport)
case "ready-http-mcp":
_, _ = fmt.Fprintln(os.Stdout, strings.Replace(readyReport, `"protocol":3`, `"protocol":3,"features":["mcp_http_tools"]`, 1))
_, _ = fmt.Fprintln(os.Stdout, strings.Replace(readyReport, `"protocol":4`, `"protocol":4,"features":["mcp_http_tools"]`, 1))
case "malformed":
_, _ = fmt.Fprintln(os.Stdout, "not-json")
case "wrong-protocol":
_, _ = fmt.Fprintln(os.Stdout, strings.Replace(readyReport, `"protocol":3`, `"protocol":1`, 1))
_, _ = fmt.Fprintln(os.Stdout, strings.Replace(readyReport, `"protocol":4`, `"protocol":3`, 1))
case "missing-version":
_, _ = fmt.Fprintln(os.Stdout, strings.Replace(readyReport, `"mcp":"1.30.0"`, `"mcp":""`, 1))
case "multiple":
Expand Down
2 changes: 1 addition & 1 deletion apps/daemon/internal/agent/claudesdk/runtime_cache.go
Original file line number Diff line number Diff line change
Expand Up @@ -41,7 +41,7 @@ func (c *runtimeCheckCache) check(ctx context.Context, config Config) (RuntimeIn
}
c.mu.Lock()
defer c.mu.Unlock()
if c.stamp == stamp && c.info.Protocol == 3 {
if c.stamp == stamp && c.info.Protocol == 4 {
return c.info, nil
}
info, err := CheckRuntime(ctx, config)
Expand Down
1 change: 1 addition & 0 deletions apps/daemon/internal/agent/claudesdk/session.go
Original file line number Diff line number Diff line change
Expand Up @@ -37,6 +37,7 @@ type bridgeEvent struct {
Confirmed *bool `json:"confirmed"`
Reason string `json:"reason"`
Protocol int `json:"protocol"`
Servers []string `json:"servers"`

Fact json.RawMessage `json:"fact"`
InputID string `json:"input_id"`
Expand Down
2 changes: 1 addition & 1 deletion apps/daemon/internal/agent/claudesdk/view.go
Original file line number Diff line number Diff line change
Expand Up @@ -98,7 +98,7 @@ func newViewExecutorFactory(probe Config, layout viewLayout) agent.ViewExecutorF
if err != nil {
return nil, err
}
return startExecutor(ctx, checked, probe, start, func() (*session, error) {
return startExecutor(ctx, checked, probe, start, view.StopMCP, func() (*session, error) {
return startSession(view.Launch, clirunner.StartOptions{Parent: ctx, Binary: layout.node, Args: []string{layout.bridge}, Dir: start.Cwd, Env: env, NeedStdin: true})
})
}
Expand Down
3 changes: 2 additions & 1 deletion apps/daemon/internal/agent/claudesdk/view_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -113,6 +113,7 @@ func TestViewExecutorLaunchesAClosedGatewayEnvironment(t *testing.T) {
docs := session.MCP
session.MCP = []agent.MCPBinding{{ServerLabel: "local", Transport: "stdio", Stdio: &agent.EnvironmentMCP{
Server: agentplugin.MCPServer{Name: "local", Type: "stdio", Command: agent.ViewAlias(0)}}}}
session.StopMCP = func(context.Context, []string) error { return nil }
installed := req
installed.CapabilityRoot, installed.Skills = agentcapabilities.Directory, []agentcapabilities.InstalledSkill{{InstallationRoot: agentcapabilities.Directory,
Metadata: agentskill.Metadata{Type: "inline", Name: "review", Description: "Review."}, RelativeRoot: "skills/review", PackageRoot: "skills/review"}}
Expand Down Expand Up @@ -210,7 +211,7 @@ func startViewBridge(options clirunner.StartOptions, requests chan<- []byte) (*c
go func() {
line, _ := bufio.NewReader(stdinReader).ReadBytes('\n')
requests <- line
_, _ = fmt.Fprintln(stdoutWriter, `{"type":"executor_ready","protocol":3}`)
_, _ = fmt.Fprintln(stdoutWriter, `{"type":"executor_ready","protocol":4}`)
<-bridge.ended
_ = stdoutWriter.Close()
_ = stdinReader.Close()
Expand Down
Loading
Loading