From b0181e77842abfd178bf9c065f67eb48802931fc Mon Sep 17 00:00:00 2001 From: SaladDay <1203511142@qq.com> Date: Wed, 7 Oct 2026 16:47:57 +0000 Subject: [PATCH 1/4] Delete the adapters' test-only paths Every Codex Turn is a Session built by Executor.StartTurn, so the Session-without-Executor branches ran only in tests: Cancel's cancelNativeWork path and the executor guards in emitTerminalFailure, beginOperation/endOperation and onServerRequest. Delete them with the executor marker field, the bare-Session cancellation tests that encoded their semantics, and the Session part of the RPC close test. Prepared only lived until newExecutor took it on the same call, so fold it into Executor: newExecutor prepares the base Session and plan directly, preparation failure releases them in place, and the owner watcher that only ran before the immediate transfer goes with TestPreparedCloseWaitsForOwnerCleanup. Inline startNative into StartTurn and delete session_run.go. The subagent cancellation, blocked steering and readiness cancellation tests now run on the Turn path. Delete Claude's test-only prepare and startRequest.Input, which only prepare set; executor_prepare never carries input. Its tests call prepareConfiguration, the factory's option path, and the Test*Factory* tests are named after the Turn behavior they check. --- .../claudesdk/execution_controls_test.go | 26 +-- .../agent/claudesdk/executor_fixture_test.go | 10 +- .../agent/claudesdk/functions_test.go | 2 +- .../internal/agent/claudesdk/local_test.go | 16 +- .../agent/claudesdk/mcp_bearer_test.go | 8 +- .../agent/claudesdk/mcp_environment_test.go | 2 +- .../internal/agent/claudesdk/mcp_test.go | 4 +- .../internal/agent/claudesdk/options.go | 13 -- .../internal/agent/claudesdk/options_test.go | 3 +- .../agent/claudesdk/preparation_test.go | 2 +- .../agent/claudesdk/restrictions_test.go | 2 +- .../internal/agent/claudesdk/session_test.go | 6 +- .../agent/claudesdk/subagents_test.go | 8 +- .../claudesdk/workspace_live_linux_test.go | 2 +- .../agent/claudesdk/workspace_test.go | 8 +- .../agent/codex/environment_retired_test.go | 11 +- apps/daemon/internal/agent/codex/executor.go | 71 ++++--- .../agent/codex/executor_native_test.go | 4 +- .../internal/agent/codex/executor_test.go | 12 +- .../internal/agent/codex/executor_turn.go | 27 +-- .../internal/agent/codex/preparation.go | 40 ++-- .../agent/codex/preparation_close_test.go | 46 +---- .../agent/codex/preparation_helpers_test.go | 4 +- .../agent/codex/preparation_router_test.go | 10 +- apps/daemon/internal/agent/codex/prepared.go | 43 ---- .../internal/agent/codex/recovery_test.go | 8 +- .../internal/agent/codex/rpc_close_test.go | 11 -- apps/daemon/internal/agent/codex/session.go | 20 +- .../internal/agent/codex/session_cancel.go | 59 ------ .../agent/codex/session_cancel_test.go | 154 --------------- .../agent/codex/session_cancel_write_test.go | 185 ------------------ .../internal/agent/codex/session_log_test.go | 5 +- .../internal/agent/codex/session_run.go | 52 ----- .../codex/session_steering_lifecycle_test.go | 16 +- .../agent/codex/session_usage_live_test.go | 35 +--- .../agent/codex/subagent_cancel_test.go | 6 +- .../agent/codex/subagent_observations_test.go | 11 +- .../agent/codex/terminal_cleanup_test.go | 6 +- 38 files changed, 184 insertions(+), 764 deletions(-) delete mode 100644 apps/daemon/internal/agent/codex/prepared.go delete mode 100644 apps/daemon/internal/agent/codex/session_cancel.go delete mode 100644 apps/daemon/internal/agent/codex/session_cancel_test.go delete mode 100644 apps/daemon/internal/agent/codex/session_cancel_write_test.go delete mode 100644 apps/daemon/internal/agent/codex/session_run.go diff --git a/apps/daemon/internal/agent/claudesdk/execution_controls_test.go b/apps/daemon/internal/agent/claudesdk/execution_controls_test.go index 0fea3251d..609d6995c 100644 --- a/apps/daemon/internal/agent/claudesdk/execution_controls_test.go +++ b/apps/daemon/internal/agent/claudesdk/execution_controls_test.go @@ -17,20 +17,20 @@ func TestExecutionControlsPreserveNativeDefaultsAndInstructions(t *testing.T) { root := t.TempDir() t.Setenv("OAC_RUNTIME_HOME", root) config := Config{Entrypoint: filepath.Join(root, "worker"), StateDir: filepath.Join(root, "state")} - request := proto.PromptRequestPayload{ModelProvider: fixtureProvider(), RunID: "run", Input: proto.TextInput("Original input."), AgentSessionID: "native-session", Model: "native-model", SystemPrompt: "Keep these exact instructions.\nDo not replace them."} - ordinary, _, err := prepare(config, request) + request := proto.PromptRequestPayload{ModelProvider: fixtureProvider(), AgentSessionID: "native-session", Model: "native-model", SystemPrompt: "Keep these exact instructions.\nDo not replace them."} + ordinary, _, err := prepareConfiguration(config, request) if err != nil { t.Fatal(err) } request.ExecutionControls = &proto.ExecutionControls{WebSearch: "disabled", TextVerbosity: "medium"} before, _ := json.Marshal(request) - controlled, _, err := prepare(config, request) + controlled, _, err := prepareConfiguration(config, request) if err != nil { t.Fatal(err) } after, _ := json.Marshal(request) if !reflect.DeepEqual(ordinary, controlled) || string(before) != string(after) { - t.Fatal("default controls changed native input, instructions, continuation or caller options") + t.Fatal("default controls changed instructions, continuation or caller options") } } @@ -81,21 +81,21 @@ func TestStructuredOutputConfigurationReachesNativeUnchanged(t *testing.T) { t.Setenv("OAC_RUNTIME_HOME", root) config := Config{Entrypoint: filepath.Join(root, "worker"), StateDir: filepath.Join(root, "state")} schema := json.RawMessage(`{"type":"object","properties":{"n":{"const":9007199254740992}}}`) - request := proto.PromptRequestPayload{ModelProvider: fixtureProvider(), RunID: "run", Input: proto.TextInput("Original input."), ObserveMessages: true, DisableSubagents: true, Model: "model", SystemPrompt: "Original instructions.", ExecutionControls: &proto.ExecutionControls{WebSearch: "disabled", TextVerbosity: "medium", OutputFormat: &proto.OutputFormat{Type: "json_schema", Schema: schema}}} - start, _, err := prepare(config, request) + request := proto.PromptRequestPayload{ModelProvider: fixtureProvider(), ObserveMessages: true, DisableSubagents: true, Model: "model", SystemPrompt: "Original instructions.", ExecutionControls: &proto.ExecutionControls{WebSearch: "disabled", TextVerbosity: "medium", OutputFormat: &proto.OutputFormat{Type: "json_schema", Schema: schema}}} + start, _, err := prepareConfiguration(config, request) if err != nil { t.Fatal(err) } - if start.OutputFormat == nil || string(start.OutputFormat.Schema) != string(schema) || *start.Input[0].Content[0].Text != *request.Input[0].Content[0].Text || start.SystemPrompt != "Original instructions." { + if start.OutputFormat == nil || string(start.OutputFormat.Schema) != string(schema) || start.SystemPrompt != "Original instructions." { t.Fatal("native configuration changed") } request.ExecutionControls.OutputFormat.Schema = json.RawMessage(`{"type":"object","const":9007199254740993}`) - if _, _, err := prepare(config, request); err == nil { + if _, _, err := prepareConfiguration(config, request); err == nil { t.Fatal("lossy schema accepted") } request.ExecutionControls.OutputFormat.Schema = schema request.DisableSubagents = false - if _, _, err := prepare(config, request); err == nil { + if _, _, err := prepareConfiguration(config, request); err == nil { t.Fatal("unqualified subagent combination accepted") } } @@ -104,21 +104,21 @@ func TestToolDiscoveryPreservesFrozenFunctionsAndRejectsOtherProfiles(t *testing root := t.TempDir() t.Setenv("OAC_RUNTIME_HOME", root) config := Config{Entrypoint: filepath.Join(root, "worker"), StateDir: filepath.Join(root, "state")} - request := proto.PromptRequestPayload{ModelProvider: fixtureProvider(), RunID: "run", Input: proto.TextInput("Original input."), DisableSubagents: true, ToolSearch: true, Model: "model", FunctionTools: []proto.FunctionTool{ + request := proto.PromptRequestPayload{ModelProvider: fixtureProvider(), DisableSubagents: true, ToolSearch: true, Model: "model", FunctionTools: []proto.FunctionTool{ {Name: "lookup", Description: "Lookup", Parameters: json.RawMessage(`{"type":"object","properties":{"ticket":{"const":"original"}}}`), DeferLoading: true}, {Name: "clock", Description: "Clock", Parameters: json.RawMessage(`{"type":"object"}`)}, }} - start, _, err := prepare(config, request) + start, _, err := prepareConfiguration(config, request) if err != nil || !start.ToolSearch || !reflect.DeepEqual(start.Functions, request.FunctionTools) { t.Fatal("function discovery changed native definitions", err) } request.DisableSubagents = false - if _, _, err := prepare(config, request); err == nil { + if _, _, err := prepareConfiguration(config, request); err == nil { t.Fatal("unqualified combination admitted") } request.DisableSubagents = true request.ToolSearch = false - if _, _, err := prepare(config, request); err == nil { + if _, _, err := prepareConfiguration(config, request); err == nil { t.Fatal("deferred definitions became eager") } } diff --git a/apps/daemon/internal/agent/claudesdk/executor_fixture_test.go b/apps/daemon/internal/agent/claudesdk/executor_fixture_test.go index 15d683ccf..e5bc554ae 100644 --- a/apps/daemon/internal/agent/claudesdk/executor_fixture_test.go +++ b/apps/daemon/internal/agent/claudesdk/executor_fixture_test.go @@ -30,17 +30,17 @@ func startSingleTurn(ctx context.Context, config Config, req proto.PromptRequest return turn, err } -func helperTurn(scanner *bufio.Scanner, request *startRequest) (func(bridgeEvent), func()) { +func helperTurn(scanner *bufio.Scanner) (func(bridgeEvent), func()) { _ = json.NewEncoder(os.Stdout).Encode(bridgeEvent{Type: "executor_ready", Protocol: 3}) if !scanner.Scan() { os.Exit(4) } var start struct { - Type string `json:"type"` - TurnID string `json:"turn_id"` - Input json.RawMessage `json:"input"` + Type string `json:"type"` + TurnID string `json:"turn_id"` + Input proto.MessageInput `json:"input"` } - if json.Unmarshal(scanner.Bytes(), &start) != nil || start.Type != "turn_start" || start.TurnID == "" || json.Unmarshal(start.Input, &request.Input) != nil { + if json.Unmarshal(scanner.Bytes(), &start) != nil || start.Type != "turn_start" || start.TurnID == "" { os.Exit(5) } return helperTurnOutput(scanner, start.TurnID) diff --git a/apps/daemon/internal/agent/claudesdk/functions_test.go b/apps/daemon/internal/agent/claudesdk/functions_test.go index d0ff4d060..c813e715f 100644 --- a/apps/daemon/internal/agent/claudesdk/functions_test.go +++ b/apps/daemon/internal/agent/claudesdk/functions_test.go @@ -16,7 +16,7 @@ import ( "github.com/MiniMax-AI/OpenAgentCore/internal/agentdaemon/proto" ) -func TestFunctionFactoryNativeReceipts(t *testing.T) { +func TestFunctionTurnNativeReceipts(t *testing.T) { for _, mode := range []string{"functions-success", "functions-wrong-receipt", "functions-no-receipt", "functions-cancel"} { t.Run(mode, func(t *testing.T) { root := t.TempDir() diff --git a/apps/daemon/internal/agent/claudesdk/local_test.go b/apps/daemon/internal/agent/claudesdk/local_test.go index a05ddac28..8ad52ac55 100644 --- a/apps/daemon/internal/agent/claudesdk/local_test.go +++ b/apps/daemon/internal/agent/claudesdk/local_test.go @@ -19,17 +19,17 @@ func TestLocalWorkspaceBindingNetworkAndRequiredHistory(t *testing.T) { req := workspaceRequest() req.LocalEnvironment = &proto.LocalEnvironment{ID: "environment", NetworkAccess: "enabled", WorkspaceRoot: config.Workspace.Directory} req.RequireExistingNativeSession = true - start, _, err := prepare(config, req) + start, _, err := prepareConfiguration(config, req) if err != nil || !start.RequireHistory || start.Workspace.NetworkAccess != "enabled" { t.Fatal(start, err) } req.LocalEnvironment.NetworkAccess = "disabled" - if _, _, err := prepare(config, req); err == nil { + if _, _, err := prepareConfiguration(config, req); err == nil { t.Fatal("accepted different Runtime network policy") } req.LocalEnvironment.NetworkAccess = "enabled" config.Workspace.PublicDirectory = config.Workspace.HomeDir - if _, _, err := prepare(config, req); err == nil { + if _, _, err := prepareConfiguration(config, req); err == nil { t.Fatal("accepted a different public workspace") } alias := filepath.Join(filepath.Dir(config.Workspace.Directory), "alias") @@ -37,7 +37,7 @@ func TestLocalWorkspaceBindingNetworkAndRequiredHistory(t *testing.T) { t.Fatal(err) } config.Workspace.PublicDirectory = alias - if _, _, err := prepare(config, req); err != nil { + if _, _, err := prepareConfiguration(config, req); err != nil { t.Fatal("same workspace alias rejected", err) } @@ -49,12 +49,12 @@ func TestRestrictedWorkspacePolicyUsesExactBoundAuthority(t *testing.T) { config.Workspace.AllowedDomains = []string{"Example.com", "api.example.com", "example.com"} req := workspaceRequest() req.LocalEnvironment = &proto.LocalEnvironment{ID: "environment", NetworkAccess: "restricted", AllowedDomains: []string{"api.example.com", "EXAMPLE.COM"}, WorkspaceRoot: config.Workspace.Directory} - if _, _, err := prepare(config, req); err == nil { + if _, _, err := prepareConfiguration(config, req); err == nil { t.Fatal("Runtime must not promise inner network isolation") } for _, domains := range [][]string{{"example.com"}, {"example.org"}, nil} { req.LocalEnvironment.AllowedDomains = domains - if _, _, err := prepare(config, req); err == nil { + if _, _, err := prepareConfiguration(config, req); err == nil { t.Fatal("different policy entered bound Runtime", domains) } } @@ -65,7 +65,7 @@ func TestWorkspaceProviderCredentialsReplaceAmbientSelection(t *testing.T) { original := slices.Clone(config.Env) req := workspaceRequest() req.ModelProvider = &modelprovider.Provider{Protocol: modelprovider.Anthropic, BaseURL: "https://provider.example/anthropic", APIKey: "selected-secret"} - start, env, err := prepare(config, req) + start, env, err := prepareConfiguration(config, req) if err != nil { t.Fatal(err) } @@ -78,7 +78,7 @@ func TestWorkspaceProviderCredentialsReplaceAmbientSelection(t *testing.T) { } for _, baseURL := range []string{"http://provider.example", "https://user:pass@provider.example"} { req.ModelProvider = &modelprovider.Provider{Protocol: modelprovider.Anthropic, BaseURL: baseURL, APIKey: "secret"} - if _, _, err := prepare(config, req); err == nil || strings.Contains(err.Error(), "secret") { + if _, _, err := prepareConfiguration(config, req); err == nil || strings.Contains(err.Error(), "secret") { t.Fatal("unsafe provider accepted or disclosed") } } diff --git a/apps/daemon/internal/agent/claudesdk/mcp_bearer_test.go b/apps/daemon/internal/agent/claudesdk/mcp_bearer_test.go index 60b9a8304..19684da9f 100644 --- a/apps/daemon/internal/agent/claudesdk/mcp_bearer_test.go +++ b/apps/daemon/internal/agent/claudesdk/mcp_bearer_test.go @@ -22,10 +22,10 @@ func TestMCPBearerUsesFreshOwnedEnvironmentReferences(t *testing.T) { {ConnectionOrigin: "service", ServerLabel: "second", ServerURL: "https://second.example/mcp", BearerToken: &tokens[1]}, {ConnectionOrigin: "service", ServerLabel: "anonymous", ServerURL: "http://anonymous.example/mcp"}, } - req := proto.PromptRequestPayload{ModelProvider: fixtureProvider(), RunID: "run", Input: proto.TextInput("hello"), DisableExecutionEnvironment: true, MCPHTTPServers: &servers, Model: "fixture"} + req := proto.PromptRequestPayload{ModelProvider: fixtureProvider(), DisableExecutionEnvironment: true, MCPHTTPServers: &servers, Model: "fixture"} seen := map[string]bool{} for range 2 { - start, env, err := prepare(config, req) + start, env, err := prepareConfiguration(config, req) if err != nil { t.Fatal(err) } @@ -66,8 +66,8 @@ func TestMCPBearerRejectsInvalidCredentialBeforeStateCreation(t *testing.T) { t.Setenv("OAC_RUNTIME_HOME", root) config := Config{Entrypoint: filepath.Join(root, "main.js"), StateDir: filepath.Join(root, "state")} servers := []proto.MCPHTTPServer{{ConnectionOrigin: "service", ServerLabel: "fixture", ServerURL: "https://example.invalid/mcp", BearerToken: &token}} - req := proto.PromptRequestPayload{ModelProvider: fixtureProvider(), RunID: "run", Input: proto.TextInput("hello"), DisableExecutionEnvironment: true, MCPHTTPServers: &servers, Model: "fixture"} - if _, _, err := prepare(config, req); err == nil || err.Error() != "claudesdk: unsupported HTTPS MCP bearer credential" { + req := proto.PromptRequestPayload{ModelProvider: fixtureProvider(), DisableExecutionEnvironment: true, MCPHTTPServers: &servers, Model: "fixture"} + if _, _, err := prepareConfiguration(config, req); err == nil || err.Error() != "claudesdk: unsupported HTTPS MCP bearer credential" { t.Fatal("invalid bearer accepted or unsafe error returned") } entries, err := os.ReadDir(root) diff --git a/apps/daemon/internal/agent/claudesdk/mcp_environment_test.go b/apps/daemon/internal/agent/claudesdk/mcp_environment_test.go index ab84b3b9e..3b9f0fea9 100644 --- a/apps/daemon/internal/agent/claudesdk/mcp_environment_test.go +++ b/apps/daemon/internal/agent/claudesdk/mcp_environment_test.go @@ -20,7 +20,7 @@ func TestEnvironmentMCPUsesInstalledLauncherAndSelectedCredential(t *testing.T) {InstallationRoot: "/private/runtime/capabilities", WorkspaceRoot: "/private/runtime/workspace", PackageRoot: "plugins/local", Server: agentplugin.MCPServer{Name: "local", Type: "stdio", Command: "untrusted-package-command", Args: []string{"package-argument"}, EnvVars: []string{"MCP_TOKEN"}}}, {InstallationRoot: "/private/runtime/capabilities", WorkspaceRoot: "/private/runtime/workspace", PackageRoot: "plugins/remote", Server: agentplugin.MCPServer{Name: "remote", Type: "http", URL: "https://example.invalid/mcp", BearerTokenEnvVar: "MCP_TOKEN"}, BearerToken: &token}, }} - start, env, err := prepare(config, req) + start, env, err := prepareConfiguration(config, req) if err != nil { t.Fatal(err) } diff --git a/apps/daemon/internal/agent/claudesdk/mcp_test.go b/apps/daemon/internal/agent/claudesdk/mcp_test.go index 1bbc8af5f..fd393182a 100644 --- a/apps/daemon/internal/agent/claudesdk/mcp_test.go +++ b/apps/daemon/internal/agent/claudesdk/mcp_test.go @@ -16,7 +16,7 @@ func TestHTTPMCPDeclaration(t *testing.T) { t.Setenv("OAC_RUNTIME_HOME", root) config := Config{Entrypoint: filepath.Join(root, "main.js"), StateDir: filepath.Join(root, "state")} servers := []proto.MCPHTTPServer{{ConnectionOrigin: "service", ServerLabel: "fixture", ServerURL: "https://example.invalid/mcp"}} - req := proto.PromptRequestPayload{ModelProvider: fixtureProvider(), RunID: "run", Input: proto.TextInput("hello"), DisableExecutionEnvironment: true, MCPHTTPServers: &servers, Model: "fixture"} + req := proto.PromptRequestPayload{ModelProvider: fixtureProvider(), DisableExecutionEnvironment: true, MCPHTTPServers: &servers, Model: "fixture"} tools := []string{"echo"} switch mode { case "selected": @@ -46,7 +46,7 @@ func TestHTTPMCPDeclaration(t *testing.T) { case "environment": req.DisableExecutionEnvironment = false } - start, _, err := prepare(config, req) + start, _, err := prepareConfiguration(config, req) valid := mode == "unrestricted" || mode == "selected" || mode == "empty" || mode == "nil-slice" || mode == "auth" || mode == "required" if (err == nil) != valid { t.Fatalf("unexpected admission: %v", err) diff --git a/apps/daemon/internal/agent/claudesdk/options.go b/apps/daemon/internal/agent/claudesdk/options.go index d7e61beaa..bcf23e384 100644 --- a/apps/daemon/internal/agent/claudesdk/options.go +++ b/apps/daemon/internal/agent/claudesdk/options.go @@ -29,7 +29,6 @@ type startRequest struct { Subagents *subagentOptions `json:"subagents,omitempty"` OutputFormat *proto.OutputFormat `json:"output_format,omitempty"` Type string `json:"type"` - Input proto.MessageInput `json:"input,omitempty"` Model string `json:"model"` SystemPrompt string `json:"system_prompt"` Cwd string `json:"cwd"` @@ -41,18 +40,6 @@ type startRequest struct { RequireHistory bool `json:"require_history,omitempty"` } -func prepare(config Config, req proto.PromptRequestPayload) (startRequest, []string, error) { - if req.RunID == "" || req.Input.Validate() != nil { - return startRequest{}, nil, fmt.Errorf("claudesdk: run id and prompt are required") - } - start, env, err := prepareConfiguration(config, req) - if err != nil { - return startRequest{}, nil, err - } - start.Input = req.Input - return start, env, nil -} - func prepareConfiguration(config Config, req proto.PromptRequestPayload) (startRequest, []string, error) { start, provider, err := prepareOptions(req, req.MCPHTTPServers != nil || (req.LocalEnvironment != nil && len(req.LocalEnvironment.MCP) != 0)) if err != nil { diff --git a/apps/daemon/internal/agent/claudesdk/options_test.go b/apps/daemon/internal/agent/claudesdk/options_test.go index 9fd82449e..5535b6191 100644 --- a/apps/daemon/internal/agent/claudesdk/options_test.go +++ b/apps/daemon/internal/agent/claudesdk/options_test.go @@ -29,8 +29,7 @@ func TestModelAndProviderAreRequired(t *testing.T) { root := t.TempDir() t.Setenv("OAC_RUNTIME_HOME", root) config := Config{Entrypoint: filepath.Join(root, "main.js"), StateDir: filepath.Join(root, "state")} - tc.req.RunID, tc.req.Input = "run", proto.TextInput("hello") - start, _, err := prepare(config, tc.req) + start, _, err := prepareConfiguration(config, tc.req) if !errors.Is(err, tc.err) || (err == nil) != (tc.err == nil) || start.SystemPrompt != tc.want { t.Fatalf("system prompt %q, error %v", start.SystemPrompt, err) } diff --git a/apps/daemon/internal/agent/claudesdk/preparation_test.go b/apps/daemon/internal/agent/claudesdk/preparation_test.go index e11848dcd..67fb78215 100644 --- a/apps/daemon/internal/agent/claudesdk/preparation_test.go +++ b/apps/daemon/internal/agent/claudesdk/preparation_test.go @@ -60,7 +60,7 @@ func TestPreparationWaitsForReceiptAndRetainsConfiguration(t *testing.T) { if err := json.Unmarshal(raw, &frozen); err != nil { t.Fatal(err) } - if frozen.Model != "fixture" || frozen.Resume != "native-session" || frozen.Workspace == nil || len(frozen.Input) != 0 { + if frozen.Model != "fixture" || frozen.Resume != "native-session" || frozen.Workspace == nil { t.Fatal("configuration-only request was not retained") } config.Env[0] = "HTTPS_PROXY=http://changed.example" diff --git a/apps/daemon/internal/agent/claudesdk/restrictions_test.go b/apps/daemon/internal/agent/claudesdk/restrictions_test.go index 6fc389905..18055762a 100644 --- a/apps/daemon/internal/agent/claudesdk/restrictions_test.go +++ b/apps/daemon/internal/agent/claudesdk/restrictions_test.go @@ -12,7 +12,7 @@ import ( "github.com/MiniMax-AI/OpenAgentCore/internal/agentdaemon/proto" ) -func TestTextFactoryAcceptsRestrictiveCapabilities(t *testing.T) { +func TestTextTurnAcceptsRestrictiveCapabilities(t *testing.T) { for _, test := range []struct { name string environment, subagents bool diff --git a/apps/daemon/internal/agent/claudesdk/session_test.go b/apps/daemon/internal/agent/claudesdk/session_test.go index 811e8041f..0a85c3e6e 100644 --- a/apps/daemon/internal/agent/claudesdk/session_test.go +++ b/apps/daemon/internal/agent/claudesdk/session_test.go @@ -16,7 +16,7 @@ import ( "github.com/MiniMax-AI/OpenAgentCore/internal/agentdaemon/proto" ) -func TestTextFactoryCompletionAndFailures(t *testing.T) { +func TestTextTurnCompletionAndFailures(t *testing.T) { for _, mode := range []string{"success", "wrong-resume", "missing", "malformed", "process-failed", "after-result", "bridge-error"} { t.Run(mode, func(t *testing.T) { root := t.TempDir() @@ -73,7 +73,7 @@ func TestTextFactoryCompletionAndFailures(t *testing.T) { } } -func TestTextFactoryRejectsUnsupportedInput(t *testing.T) { +func TestUnsupportedRequestRejectedBeforeLaunch(t *testing.T) { for _, kind := range []string{"execution-controls", "tool", "outside"} { t.Run(kind, func(t *testing.T) { root := t.TempDir() @@ -129,7 +129,7 @@ func runSDKHelper() { if json.Unmarshal(scanner.Bytes(), &request) != nil || request.Type != "executor_prepare" || request.Model != "fake-model" || request.SystemPrompt != "instructions" { os.Exit(3) } - encode, finish := helperTurn(scanner, &request) + encode, finish := helperTurn(scanner) defer finish() mode := os.Getenv("SDK_HELPER_MODE") if strings.HasPrefix(mode, "cancellation-") { diff --git a/apps/daemon/internal/agent/claudesdk/subagents_test.go b/apps/daemon/internal/agent/claudesdk/subagents_test.go index 6cb3e7467..115fa05f2 100644 --- a/apps/daemon/internal/agent/claudesdk/subagents_test.go +++ b/apps/daemon/internal/agent/claudesdk/subagents_test.go @@ -9,18 +9,18 @@ import ( func TestSubagentConfigurationUsesExplicitRequestAndFrozenLimit(t *testing.T) { config := workspaceFixture(t) req := workspaceRequest() - start, _, err := prepare(config, req) + start, _, err := prepareConfiguration(config, req) if err != nil || start.Subagents != nil { t.Fatal("ordinary execution changed", err) } req.DisableSubagents, req.ObserveSubagentIdentities = false, true - start, _, err = prepare(config, req) + start, _, err = prepareConfiguration(config, req) if err != nil || start.Subagents == nil || start.Subagents.MaxConcurrent != 6 { t.Fatal("missing default native admission limit", err) } limit := 2 req.MaxConcurrentSubagents = &limit - start, _, err = prepare(config, req) + start, _, err = prepareConfiguration(config, req) limit = 4 if err != nil || start.Subagents.MaxConcurrent != 2 { t.Fatal("subagent configuration was not frozen", err) @@ -38,7 +38,7 @@ func TestSubagentConfigurationRejectsUnqualifiedAuthority(t *testing.T) { req := workspaceRequest() req.DisableSubagents, req.ObserveSubagentIdentities = false, true change(&req) - if _, _, err := prepare(config, req); err == nil { + if _, _, err := prepareConfiguration(config, req); err == nil { t.Fatal("unqualified subagent combination accepted") } } diff --git a/apps/daemon/internal/agent/claudesdk/workspace_live_linux_test.go b/apps/daemon/internal/agent/claudesdk/workspace_live_linux_test.go index 7a93ef048..f1e2f6fa7 100644 --- a/apps/daemon/internal/agent/claudesdk/workspace_live_linux_test.go +++ b/apps/daemon/internal/agent/claudesdk/workspace_live_linux_test.go @@ -19,7 +19,7 @@ import ( // Run only inside a separately qualified outer placement, with its pinned native // dependencies. This fixture does not create isolation or public admission. -func TestLiveClaudeWorkspaceFactory(t *testing.T) { +func TestLiveClaudeWorkspaceTurns(t *testing.T) { configFile := os.Getenv("OAC_TEST_CLAUDE_WORKSPACE_LIVE_CONFIG") if configFile == "" { t.Skip("requires explicit qualified placement and real provider configuration") diff --git a/apps/daemon/internal/agent/claudesdk/workspace_test.go b/apps/daemon/internal/agent/claudesdk/workspace_test.go index 28d978fca..b343b09a3 100644 --- a/apps/daemon/internal/agent/claudesdk/workspace_test.go +++ b/apps/daemon/internal/agent/claudesdk/workspace_test.go @@ -44,7 +44,7 @@ func TestWorkspaceTrustedBindingAndEnvironment(t *testing.T) { t.Setenv("OAC_TEST_PARENT_SECRET", "parent-only") t.Setenv("ANTHROPIC_API_KEY", "unselected-provider") config.Env = append(config.Env, "ANTHROPIC_BASE_URL=https://unselected.example") - start, env, err := prepare(config, workspaceRequest()) + start, env, err := prepareConfiguration(config, workspaceRequest()) if err != nil { t.Fatal(err) } @@ -115,7 +115,7 @@ func TestWorkspaceRejectsConflictsBeforeSideEffects(t *testing.T) { config.Env = append(config.Env, "NO_PROXY=bad\x00value") } - if _, _, err := prepare(config, req); err == nil { + if _, _, err := prepareConfiguration(config, req); err == nil { t.Fatal("invalid binding or request accepted") } entries, err := os.ReadDir(config.StateDir) @@ -130,7 +130,7 @@ func TestWorkspaceRetainsDeclaredFunctions(t *testing.T) { config := workspaceFixture(t) req := workspaceRequest() req.FunctionTools = []proto.FunctionTool{{Name: "lookup", Parameters: json.RawMessage(`{"type":"object"}`)}} - start, _, err := prepare(config, req) + start, _, err := prepareConfiguration(config, req) if err != nil { t.Fatal(err) } @@ -147,7 +147,7 @@ func TestPublicMCPUsesWorkspaceProjectionWithoutCredentialCopy(t *testing.T) { token := "vault-selected-canary" tools := []string{"prove"} req.MCPHTTPServers = &[]proto.MCPHTTPServer{{ConnectionOrigin: "environment", ServerLabel: "remote", ServerURL: "https://example.test/mcp", AllowedTools: &tools, Required: true, BearerToken: &token}} - start, env, err := prepare(config, req) + start, env, err := prepareConfiguration(config, req) if err != nil { t.Fatal(err) } diff --git a/apps/daemon/internal/agent/codex/environment_retired_test.go b/apps/daemon/internal/agent/codex/environment_retired_test.go index f3e070843..8ff7e6b4f 100644 --- a/apps/daemon/internal/agent/codex/environment_retired_test.go +++ b/apps/daemon/internal/agent/codex/environment_retired_test.go @@ -1,6 +1,7 @@ package codex import ( + "context" "os" "path/filepath" "strings" @@ -12,7 +13,7 @@ import ( func TestReadOnlyPreparationRejectedBeforeNativeSetup(t *testing.T) { req, cfg, root := preparationFixture(t) req.WorkspaceReadOnly = true - prepared, err := newPreparation(t.Context(), req, cfg) + prepared, err := newExecutor(t.Context(), req, cfg) if err == nil || prepared != nil { t.Fatal("read-only request admitted", err) } @@ -34,8 +35,8 @@ func TestRetiredNativeTransportEnvironmentRejectedBeforeState(t *testing.T) { req.LocalEnvironment = &proto.LocalEnvironment{ID: "local"} } t.Setenv(key, "retired-private-value") - p, err := newPreparation(t.Context(), req, cfg) - if p != nil || err == nil || !strings.Contains(err.Error(), "retired executor transport") || strings.Contains(err.Error(), "retired-private-value") { + e, err := newExecutor(t.Context(), req, cfg) + if e != nil || err == nil || !strings.Contains(err.Error(), "retired executor transport") || strings.Contains(err.Error(), "retired-private-value") { t.Fatal("transport override admitted or disclosed", err) } if len(preparationFrames(t, root)) != 0 { @@ -52,11 +53,11 @@ func TestRetiredNativeTransportEnvironmentRejectedBeforeState(t *testing.T) { func TestNativeNoneSelectorRemainsSupported(t *testing.T) { req, cfg, root := preparationFixture(t) t.Setenv("CODEX_EXEC_SERVER_URL", "none") - p, err := newPreparation(t.Context(), req, cfg) + e, err := newExecutor(t.Context(), req, cfg) if err != nil { t.Fatal(err) } - defer p.Close() + defer e.Close(context.Background()) assertPreparationOnly(t, root) statuses := 0 for _, frame := range preparationFrames(t, root) { diff --git a/apps/daemon/internal/agent/codex/executor.go b/apps/daemon/internal/agent/codex/executor.go index d8997e525..3faacfcf3 100644 --- a/apps/daemon/internal/agent/codex/executor.go +++ b/apps/daemon/internal/agent/codex/executor.go @@ -3,6 +3,7 @@ package codex import ( "context" "errors" + "fmt" "strings" "sync" "time" @@ -14,32 +15,20 @@ import ( // Executor owns the prepared process, fixed plan and native thread. Each StartTurn // creates independent receipt, observation, cancellation and output ownership. type Executor struct { - mu sync.Mutex - prepared *Prepared - active *Session - closed bool - closeMu sync.Mutex + mu sync.Mutex + base *Session + plan SessionPlan + resumeID string + requireExistingNativeSession bool + active *Session + closed bool + closeMu sync.Mutex } func PrepareExecutor(ctx context.Context, req proto.PromptRequestPayload) (agent.Executor, error) { return newExecutor(ctx, req, defaultSessionConfig()) } -func newExecutor(ctx context.Context, req proto.PromptRequestPayload, cfg sessionConfig) (*Executor, error) { - p, err := newPreparation(ctx, req, cfg) - if err != nil { - if p != nil { - return &Executor{prepared: p}, err - } - return nil, err - } - p.mu.Lock() - p.started = true - close(p.transferred) - p.mu.Unlock() - return &Executor{prepared: p}, nil -} - func (e *Executor) StartTurn(ctx context.Context, runID string, input proto.MessageInput, out chan<- proto.Envelope) (agent.Turn, error) { if out == nil || strings.TrimSpace(runID) == "" || input.Validate() != nil { return nil, errors.New("codex: start requires a run identity, input and output") @@ -48,7 +37,7 @@ func (e *Executor) StartTurn(ctx context.Context, runID string, input proto.Mess return nil, err } e.mu.Lock() - base := e.prepared.session + base := e.base if e.closed || base.cancelCtx.Err() != nil || !base.rpc.Alive() { e.mu.Unlock() return nil, errors.New("codex: executor unavailable") @@ -68,10 +57,10 @@ func (e *Executor) StartTurn(ctx context.Context, runID string, input proto.Mess } functions := &functionCalls{definitions: base.functions.definitions, names: base.functions.names, pending: map[string]*pendingFunction{}} turnCtx, cancel := context.WithCancel(base.cancelCtx) - s := &Session{executor: e, nativeHome: base.nativeHome, + s := &Session{nativeHome: base.nativeHome, functions: functions, observeMessages: base.observeMessages, observeSubagentIdentities: base.observeSubagentIdentities, cfg: base.cfg, rpc: base.rpc, cancelCtx: turnCtx, cancelFn: cancel, - waitDone: make(chan struct{}), outputDone: make(chan struct{}), cleanup: func() {}, + waitDone: make(chan struct{}), outputDone: make(chan struct{}), bufs: NewItemBuffers(), resolvedModel: base.resolvedModel, runID: runID, out: out} if previous != nil { s.threadID = previous.currentThreadID() @@ -93,13 +82,39 @@ func (e *Executor) StartTurn(ctx context.Context, runID string, input proto.Mess } s.registerHandlers() e.mu.Unlock() - req := proto.PromptRequestPayload{RunID: runID, Input: input, AgentSessionID: e.prepared.resumeID, RequireExistingNativeSession: e.prepared.requireExistingNativeSession} // Ownership precedes any native submission. Even an uncertain start returns the // exact Turn so its caller can await settlement without replaying the input. - err := s.startNative(ctx, e.prepared.plan, req) + var err error + if s.currentThreadID() == "" { + err = s.resolveThread(proto.PromptRequestPayload{AgentSessionID: e.resumeID, RequireExistingNativeSession: e.requireExistingNativeSession}, e.plan) + } + var native []UserInput + if err == nil { + native, err = nativeInput(input) + } + model := strings.TrimSpace(s.resolvedModel) + if err == nil && model == "" { + err = errors.New("codex: collaboration mode requires a resolved model") + } + if err == nil { + var developerInstructions *string + if instructions := e.plan.SystemPrompt; instructions != "" { + developerInstructions = &instructions + } + params := TurnStartParams{ThreadID: s.currentThreadID(), Input: native, CollaborationMode: &CollaborationMode{ + Mode: CollaborationModeDefault, + Settings: CollaborationModeSettings{ReasoningEffort: e.plan.ModelReasoningEffort, Model: model, DeveloperInstructions: developerInstructions}, + }} + startCtx, stop := context.WithTimeout(ctx, 10*time.Second) + _, err = s.rpc.requestWithResult(startCtx, "turn/start", params, s.bindTurnResult) + stop() + if err != nil { + s.cfg.logger.Warn("codex: turn/start ack failed", "run_id", runID, "err", err) + err = fmt.Errorf("codex: turn/start: %w", err) + } + } if err != nil { s.emitTerminal("codex: native start failed", true) - s.finishAfterTerminal() } go s.settleExecutorTurn(err) return s, err @@ -113,7 +128,7 @@ func (e *Executor) Close(ctx context.Context) error { go func() { e.closeMu.Lock() defer e.closeMu.Unlock() - base := e.prepared.session + base := e.base e.mu.Lock() running := e.active e.mu.Unlock() @@ -158,7 +173,7 @@ func (e *Executor) Close(ctx context.Context) error { err = base.rpc.awaitReaders(ctx) } if err == nil { - e.prepared.plan.Cleanup() + e.plan.Cleanup() } done <- err }() diff --git a/apps/daemon/internal/agent/codex/executor_native_test.go b/apps/daemon/internal/agent/codex/executor_native_test.go index 762ef7389..405b2b066 100644 --- a/apps/daemon/internal/agent/codex/executor_native_test.go +++ b/apps/daemon/internal/agent/codex/executor_native_test.go @@ -71,7 +71,7 @@ func TestExecutorNativeReuse(t *testing.T) { } }) t.Logf("native_version=0.153.4 prepare_ms=%d", time.Since(began).Milliseconds()) - pid := e.prepared.session.rpc.process.Cmd.Process.Pid + pid := e.base.rpc.process.Cmd.Process.Pid var thread string run := func(id, prompt string, old agent.Turn, interrupt bool) (agent.Turn, string) { t.Helper() @@ -159,7 +159,7 @@ func TestExecutorNativeReuse(t *testing.T) { if err := e.Close(ctx); err != nil { t.Fatal("native final cleanup failed") } - if e.prepared.session.rpc.Alive() { + if e.base.rpc.Alive() { t.Fatal("native child remained alive after owner close") } } diff --git a/apps/daemon/internal/agent/codex/executor_test.go b/apps/daemon/internal/agent/codex/executor_test.go index 9c88a0579..7404a823b 100644 --- a/apps/daemon/internal/agent/codex/executor_test.go +++ b/apps/daemon/internal/agent/codex/executor_test.go @@ -107,7 +107,7 @@ func TestExecutorNormalTurnsKeepProcessAndThread(t *testing.T) { if counts["initialize"] != 1 || counts["thread/start"] != 1 || counts["turn/start"] != 2 || counts["turn/interrupt"] != 0 || len(pids) != 1 { t.Fatal(counts, pids) } - if !e.prepared.session.rpc.Alive() { + if !e.base.rpc.Alive() { t.Fatal("normal completion closed executor") } } @@ -124,7 +124,7 @@ func TestExecutorFreezesPreparedConfiguration(t *testing.T) { t.Fatal(err) } assertPreparationOnly(t, root) - cwd := e.prepared.plan.Cwd + cwd := e.plan.Cwd // Caller-owned data cannot revise the prepared native configuration. req.Model = "different-model" req.AgentSessionID = "different-thread" @@ -183,7 +183,7 @@ func TestExecutorUnavailableOwnerStartsNoTurn(t *testing.T) { case "owner cancelled": cancelOwner() case "rpc exited": - err = e.prepared.session.rpc.Close() + err = e.base.rpc.Close() } if err != nil { t.Fatal(err) @@ -292,7 +292,7 @@ func TestExecutorCloseRetainsPlanUntilReaped(t *testing.T) { t.Cleanup(func() { process.Cancel(); reap.Do(rpc.waitChild) }) _, cancel := context.WithCancel(t.Context()) var cleaned atomic.Bool - e := &Executor{prepared: &Prepared{session: &Session{rpc: rpc, cancelFn: cancel}, plan: SessionPlan{Cleanup: func() { cleaned.Store(true) }}}} + e := &Executor{base: &Session{rpc: rpc, cancelFn: cancel}, plan: SessionPlan{Cleanup: func() { cleaned.Store(true) }}} ctx, stop := context.WithTimeout(t.Context(), 20*time.Millisecond) defer stop() if err = e.Close(ctx); !errors.Is(err, context.DeadlineExceeded) { @@ -328,7 +328,7 @@ func TestExecutorCloseSeparatesTurnFailureFromResourceCleanup(t *testing.T) { if err := e.Close(ctx); err != nil { t.Fatal("failed Turn prevented confirmed resource cleanup", err) } - if e.prepared.session.rpc.Alive() || len(preparedCatalogs(t, root)) != 0 { + if e.base.rpc.Alive() || len(preparedCatalogs(t, root)) != 0 { t.Fatal("Close did not release native process and plan") } if _, err := turn.AwaitSettlement(ctx); err == nil { @@ -360,7 +360,7 @@ func TestExecutorCloseTerminatesAfterMissingCancellationTerminal(t *testing.T) { if err := e.Close(retry); err != nil { t.Fatal("Close never retired the native process", err) } - if e.prepared.session.rpc.Alive() || len(preparedCatalogs(t, root)) != 0 { + if e.base.rpc.Alive() || len(preparedCatalogs(t, root)) != 0 { t.Fatal("unresponsive native process or plan still owned after Close") } select { diff --git a/apps/daemon/internal/agent/codex/executor_turn.go b/apps/daemon/internal/agent/codex/executor_turn.go index 60746293f..60e5f6319 100644 --- a/apps/daemon/internal/agent/codex/executor_turn.go +++ b/apps/daemon/internal/agent/codex/executor_turn.go @@ -88,9 +88,6 @@ func (s *Session) settleExecutorTurn(startErr error) { } func (s *Session) beginOperation() bool { - if s.executor == nil { - return true - } s.operationMu.Lock() defer s.operationMu.Unlock() if s.operationsClosed { @@ -99,13 +96,9 @@ func (s *Session) beginOperation() bool { s.operations.Add(1) return true } -func (s *Session) endOperation() { - if s.executor != nil { - s.operations.Done() - } -} +func (s *Session) endOperation() { s.operations.Done() } -func (s *Session) cancelExecutorTurn(ctx context.Context) error { +func (s *Session) Cancel(ctx context.Context) error { s.operationMu.Lock() if s.operationsClosed { s.operationMu.Unlock() @@ -151,15 +144,13 @@ var _ agent.Turn = (*Session)(nil) // Server callbacks retain their originating Turn, whose ID must match exactly. func (s *Session) onServerRequest(method string, handler ServerRequestHandler) { s.rpc.OnServerRequest(method, func(raw json.RawMessage, id any) (any, error) { - if s.executor != nil { - var scope struct { - ThreadID string `json:"threadId"` - TurnID string `json:"turnId"` - } - if json.Unmarshal(raw, &scope) != nil || !s.isRootThread(scope.ThreadID) || s.terminal.Load() || - scope.TurnID == "" || !s.isRootTurn(scope.ThreadID, scope.TurnID) { - return nil, errors.New("codex: request does not belong to the active turn") - } + var scope struct { + ThreadID string `json:"threadId"` + TurnID string `json:"turnId"` + } + if json.Unmarshal(raw, &scope) != nil || !s.isRootThread(scope.ThreadID) || s.terminal.Load() || + scope.TurnID == "" || !s.isRootTurn(scope.ThreadID, scope.TurnID) { + return nil, errors.New("codex: request does not belong to the active turn") } if !s.beginOperation() { return nil, errors.New("codex: turn has settled") diff --git a/apps/daemon/internal/agent/codex/preparation.go b/apps/daemon/internal/agent/codex/preparation.go index d47acc718..3a2d5651a 100644 --- a/apps/daemon/internal/agent/codex/preparation.go +++ b/apps/daemon/internal/agent/codex/preparation.go @@ -11,7 +11,7 @@ import ( obslog "github.com/MiniMax-AI/OpenAgentCore/internal/obs/log" ) -func newPreparation(parent context.Context, req proto.PromptRequestPayload, cfg sessionConfig) (*Prepared, error) { +func newExecutor(parent context.Context, req proto.PromptRequestPayload, cfg sessionConfig) (*Executor, error) { if req.ExecutionControls != nil && req.ExecutionControls.OutputFormat != nil { return nil, errors.New("codex: structured output is not qualified") } @@ -74,59 +74,55 @@ func newPreparation(parent context.Context, req proto.PromptRequestPayload, cfg rpc: rpc, cancelCtx: cancelCtx, cancelFn: cancelFn, - waitDone: make(chan struct{}), - cleanup: sync.OnceFunc(plan.Cleanup), - bufs: NewItemBuffers(), resolvedModel: plan.Model, } - plan.Cleanup = s.cleanup - p := &Prepared{ - session: s, plan: plan, - resumeID: req.AgentSessionID, requireExistingNativeSession: req.RequireExistingNativeSession, - transferred: make(chan struct{}), - } + plan.Cleanup = sync.OnceFunc(plan.Cleanup) + e := &Executor{base: s, plan: plan, resumeID: req.AgentSessionID, requireExistingNativeSession: req.RequireExistingNativeSession} initParams := InitializeParams{ ClientInfo: InitializeClientInfo{Name: "oac-daemon", Version: "0.0.0"}, Capabilities: &InitializeCapabilities{ExperimentalAPI: true}, } if _, err := rpc.Start(cancelCtx, initParams); err != nil { - return p.preparationFailed(fmt.Errorf("codex: rpc start: %w", err)) + return e.preparationFailed(fmt.Errorf("codex: rpc start: %w", err)) } if req.ExecutionControls != nil && req.ExecutionControls.DisableProgrammaticToolCalling { if err := verifyProgrammaticToolsDisabled(cancelCtx, rpc); err != nil { - return p.preparationFailed(err) + return e.preparationFailed(err) } } if req.DisableExecutionEnvironment { if err := verifyNoExecutionEnvironment(cancelCtx, rpc); err != nil { - return p.preparationFailed(err) + return e.preparationFailed(err) } } if s.observeSubagentIdentities { if err := verifySubagentObservationProfile(cancelCtx, rpc, plan.Cwd); err != nil { - return p.preparationFailed(err) + return e.preparationFailed(err) } } if plan.mcpServers != nil { if err := verifyMCPConfig(cancelCtx, rpc, plan); err != nil { - return p.preparationFailed(err) + return e.preparationFailed(err) } } if len(skillRoots) > 0 { if err := setSkillExtraRoots(cancelCtx, rpc, skillRoots); err != nil { - return p.preparationFailed(fmt.Errorf("codex: register skill root: %w", err)) + return e.preparationFailed(fmt.Errorf("codex: register skill root: %w", err)) } } - go p.watchOwner() - return p, nil + return e, nil } -func (p *Prepared) preparationFailed(cause error) (*Prepared, error) { - if err := p.Close(); err != nil { - return p, errors.Join(cause, err) - } +// preparationFailed releases the unused process and plan. An unconfirmed release +// keeps them in the returned Executor for a later Close. +func (e *Executor) preparationFailed(cause error) (*Executor, error) { + e.base.cancelFn() + if err := e.base.rpc.Close(); err != nil { + return e, errors.Join(cause, err) + } + e.plan.Cleanup() return nil, cause } diff --git a/apps/daemon/internal/agent/codex/preparation_close_test.go b/apps/daemon/internal/agent/codex/preparation_close_test.go index 70cc5d769..f4263dbc2 100644 --- a/apps/daemon/internal/agent/codex/preparation_close_test.go +++ b/apps/daemon/internal/agent/codex/preparation_close_test.go @@ -2,58 +2,20 @@ package codex import ( "context" - "sync" "testing" "time" ) -func TestPreparedCloseWaitsForOwnerCleanup(t *testing.T) { - req, cfg, root := preparationFixture(t) - owner, cancel := context.WithCancel(t.Context()) - defer cancel() - p, err := newPreparation(owner, req, cfg) - if err != nil { - t.Fatal(err) - } - defer p.Close() - entered, release := make(chan struct{}), make(chan struct{}) - var releaseOnce sync.Once - unblock := func() { releaseOnce.Do(func() { close(release) }) } - defer unblock() - cleanup := p.plan.Cleanup - p.plan.Cleanup = sync.OnceFunc(func() { close(entered); <-release; cleanup() }) - cancel() - select { - case <-entered: - case <-time.After(3 * time.Second): - t.Fatal("owner watcher did not enter cleanup") - } - done := make(chan struct{}) - go func() { _ = p.Close(); close(done) }() - select { - case <-done: - t.Fatal("Close returned while owner cleanup was still running") - case <-time.After(50 * time.Millisecond): - } - unblock() - select { - case <-done: - case <-time.After(3 * time.Second): - t.Fatal("Close did not finish after cleanup") - } - waitPreparedRelease(t, p, root) -} - -func TestPreparedSessionCancellationDuringReadiness(t *testing.T) { +func TestPreparationCancellationDuringReadiness(t *testing.T) { req, cfg, root := preparationFixture(t) t.Setenv("OAC_TEST_PREPARATION_BLOCK", "1") owner, cancel := context.WithCancel(t.Context()) defer cancel() result := make(chan error, 1) go func() { - p, err := newPreparation(owner, req, cfg) - if p != nil { - _ = p.Close() + e, err := newExecutor(owner, req, cfg) + if e != nil { + _ = e.Close(context.Background()) } result <- err }() diff --git a/apps/daemon/internal/agent/codex/preparation_helpers_test.go b/apps/daemon/internal/agent/codex/preparation_helpers_test.go index 984f25b7b..24952d1d8 100644 --- a/apps/daemon/internal/agent/codex/preparation_helpers_test.go +++ b/apps/daemon/internal/agent/codex/preparation_helpers_test.go @@ -112,10 +112,10 @@ func preparedCatalogs(t *testing.T, root string) []string { return files } -func waitPreparedRelease(t *testing.T, p *Prepared, root string) { +func waitExecutorRelease(t *testing.T, e *Executor, root string) { t.Helper() select { - case <-p.session.rpc.Done(): + case <-e.base.rpc.Done(): case <-time.After(4 * time.Second): t.Fatal("prepared child was not released") } diff --git a/apps/daemon/internal/agent/codex/preparation_router_test.go b/apps/daemon/internal/agent/codex/preparation_router_test.go index e965460a2..b743ec312 100644 --- a/apps/daemon/internal/agent/codex/preparation_router_test.go +++ b/apps/daemon/internal/agent/codex/preparation_router_test.go @@ -58,13 +58,13 @@ func TestPreparationRouterRetainsActualNativeChild(t *testing.T) { req.LocalEnvironment = &proto.LocalEnvironment{ID: environment, WorkspaceDirectory: "/workspace", NetworkAccess: "enabled", CapabilitySources: &agentcapabilities.Input{}} registry := agent.NewRegistry() registry.RegisterKind(proto.SupportedAgentKind{Kind: "codex", Available: true, Capabilities: prototest.Capabilities(proto.AgentKindCapabilities{LocalEnvironment: proto.CapabilitySupported, FunctionTools: proto.CapabilitySupported})}, harnessconfiguration.Configuration()) - prepared := make(chan *Prepared, 1) + prepared := make(chan *Executor, 1) registry.RegisterExecutor("codex", func(ctx context.Context, req proto.PromptRequestPayload) (agent.Executor, error) { e, err := newExecutor(ctx, req, cfg) if err != nil { return nil, err } - prepared <- e.prepared + prepared <- e return e, nil }) sender := make(preparationWireSender, 64) @@ -117,7 +117,7 @@ func TestPreparationRouterRetainsActualNativeChild(t *testing.T) { ready := await("ready") p := <-prepared assertPreparationOnly(t, root) - pid := p.session.rpc.process.Cmd.Process.Pid + pid := p.base.rpc.process.Cmd.Process.Pid if r.ActiveRuns() != 0 { t.Fatal("preparation became a Run") } @@ -140,7 +140,7 @@ func TestPreparationRouterRetainsActualNativeChild(t *testing.T) { } send(proto.TypeExecutionStart, "prepare-request", input) send(proto.TypeExecutionRelease, "prepare-request", proto.ExecutionReleasePayload{Handle: ready.Handle}) - if !p.session.rpc.Alive() || r.ActiveRuns() != 1 { + if !p.base.rpc.Alive() || r.ActiveRuns() != 1 { t.Fatal("release cancelled transferred native session") } send(proto.TypePromptCancel, "actual-run", proto.PromptCancelPayload{}) @@ -150,7 +150,7 @@ func TestPreparationRouterRetainsActualNativeChild(t *testing.T) { if err := r.Shutdown(ctx); err != nil { t.Fatal(err) } - waitPreparedRelease(t, p, root) + waitExecutorRelease(t, p, root) if !start { assertPreparationOnly(t, root) } diff --git a/apps/daemon/internal/agent/codex/prepared.go b/apps/daemon/internal/agent/codex/prepared.go deleted file mode 100644 index d069e9311..000000000 --- a/apps/daemon/internal/agent/codex/prepared.go +++ /dev/null @@ -1,43 +0,0 @@ -package codex - -import "sync" - -// Prepared owns a connected native resource until an Executor takes it. -// It observes owner cancellation and RPC exit, not continuous executor readiness. -type Prepared struct { - mu sync.Mutex - session *Session - plan SessionPlan - resumeID string - requireExistingNativeSession bool - started bool - transferred chan struct{} -} - -// Close waits for unused teardown and plan cleanup, including another caller's -// ongoing Close. After an Executor takes the resource it is inert; use -// Executor.Close to release it. -func (p *Prepared) Close() error { - p.mu.Lock() - if p.started { - p.mu.Unlock() - return nil - } - p.mu.Unlock() - p.session.cancelFn() - err := p.session.rpc.Close() - if err == nil { - p.plan.Cleanup() - } - return err -} - -func (p *Prepared) watchOwner() { - select { - case <-p.session.cancelCtx.Done(): - case <-p.session.rpc.Done(): - case <-p.transferred: - return - } - _ = p.Close() -} diff --git a/apps/daemon/internal/agent/codex/recovery_test.go b/apps/daemon/internal/agent/codex/recovery_test.go index 4c9b158db..bd84ec2bd 100644 --- a/apps/daemon/internal/agent/codex/recovery_test.go +++ b/apps/daemon/internal/agent/codex/recovery_test.go @@ -165,8 +165,8 @@ func TestPreparedRecoveryCannotStartWithoutExistingHistory(t *testing.T) { if err != nil { t.Fatal(err) } - if e.prepared.plan.Cwd != cwd { - t.Fatalf("cwd = %q, want %q", e.prepared.plan.Cwd, cwd) + if e.plan.Cwd != cwd { + t.Fatalf("cwd = %q, want %q", e.plan.Cwd, cwd) } assertPreparationOnly(t, root) out := make(chan proto.Envelope, 16) @@ -202,8 +202,8 @@ func TestRecoveryRequiresWritableAgentState(t *testing.T) { case "read-only": req.WorkspaceReadOnly = true } - if p, err := newPreparation(t.Context(), req, cfg); err == nil { - p.Close() + if e, err := newExecutor(t.Context(), req, cfg); err == nil { + _ = e.Close(t.Context()) t.Fatal("invalid recovery admitted") } if len(preparationFrames(t, root)) != 0 { diff --git a/apps/daemon/internal/agent/codex/rpc_close_test.go b/apps/daemon/internal/agent/codex/rpc_close_test.go index 036b1c898..02ccbe802 100644 --- a/apps/daemon/internal/agent/codex/rpc_close_test.go +++ b/apps/daemon/internal/agent/codex/rpc_close_test.go @@ -78,18 +78,7 @@ func TestJSONRPCClientCloseCanRetryUnreapedChild(t *testing.T) { } } closeConcurrently(true) - ownerCtx, cancelOwner := context.WithCancel(t.Context()) - s := &Session{rpc: client, cancelCtx: ownerCtx, cancelFn: cancelOwner, - cfg: defaultSessionConfig(), bufs: NewItemBuffers()} - ctx, cancelWait := context.WithTimeout(t.Context(), 20*time.Millisecond) - defer cancelWait() - if err := s.Cancel(ctx); !errors.Is(err, context.DeadlineExceeded) { - t.Fatalf("Session cancellation before reap: %v", err) - } reap.Do(client.waitChild) - if err := s.Cancel(t.Context()); err != nil { - t.Fatalf("Session cancellation retry after reap: %v", err) - } closeConcurrently(false) if _, exited := process.ExitCode(); client.process != process || !exited { t.Fatal("Close did not retain and reap its original child") diff --git a/apps/daemon/internal/agent/codex/session.go b/apps/daemon/internal/agent/codex/session.go index 1a5db1094..4c05f9cc9 100644 --- a/apps/daemon/internal/agent/codex/session.go +++ b/apps/daemon/internal/agent/codex/session.go @@ -39,18 +39,17 @@ func defaultSessionConfig() sessionConfig { } } -// Session implements agent.Session. State lifecycle: +// Session is one Executor Turn. State lifecycle: // -// 1. Preparation initializes RPC and verifies the selected environment. -// 2. Start transfers that RPC and wires notification/server-request handlers. -// 3. thread/start or thread/resume runs (resume falls back to start). +// 1. Executor preparation initializes RPC and verifies the selected environment. +// 2. StartTurn wires notification/server-request handlers on that RPC. +// 3. The first Turn runs thread/start or thread/resume. // 4. turn/start delivers the user prompt; subsequent stream notifications // fan out to proto.Envelope via session_items.go. -// 5. turn/completed emits TypeDone + closes out. Cancel can short-cut -// this by killing the child early. +// 5. turn/completed emits TypeDone and closes out. Cancel interrupts the +// native turn; settlement decides whether the Executor stays reusable. type Session struct { retiredTurns map[string]bool - executor *Executor outputDone chan struct{} nativeSettled atomic.Bool settlement agent.TurnSettlement @@ -82,7 +81,6 @@ type Session struct { outMu sync.RWMutex outClosed bool waitDone chan struct{} - cleanup func() threadIDMu sync.Mutex threadID string @@ -220,7 +218,6 @@ func (s *Session) onTurnCompleted(raw json.RawMessage) { body = appendOnNewline(body, errText) } s.emitTerminalFailure(body, true, classifyTurnError(p.Turn.Error)) - s.finishAfterTerminal() return } var completedAt *int64 @@ -286,7 +283,6 @@ func (s *Session) onTurnFailed(raw json.RawMessage) { "turn_status", p.Turn.Status, "last_err_text_present", s.peekLastErrText() != "") s.emitTerminal("codex: turn failed", true) - s.finishAfterTerminal() } func (s *Session) onErrorNotif(raw json.RawMessage) { @@ -369,9 +365,7 @@ func (s *Session) emitTerminal(message string, asError bool) { } func (s *Session) emitTerminalFailure(message string, asError bool, failure proto.ErrorPayload) { - if s.executor != nil { - defer s.finishAfterTerminal() - } + defer s.finishAfterTerminal() if !s.terminal.CompareAndSwap(false, true) { return } diff --git a/apps/daemon/internal/agent/codex/session_cancel.go b/apps/daemon/internal/agent/codex/session_cancel.go deleted file mode 100644 index 091eae119..000000000 --- a/apps/daemon/internal/agent/codex/session_cancel.go +++ /dev/null @@ -1,59 +0,0 @@ -package codex - -import ( - "context" - "errors" - "time" -) - -func (s *Session) Cancel(ctx context.Context) error { - if s.executor != nil { - return s.cancelExecutorTurn(ctx) - } - s.cancelled.Store(true) - s.cancelOnce.Do(func() { - s.cancelReady = make(chan struct{}) - go s.cancelNativeWork() - }) - // A caller deadline does not terminate the owner that is still collecting - // native child cancellation facts. Another call can await the same cleanup. - select { - case <-s.cancelReady: - default: - select { - case <-s.cancelReady: - case <-ctx.Done(): - return ctx.Err() - } - } - // Close already owns shutdown once and permits another wait for process exit. - // Keep actual observation failures, but never cache a caller's wait timeout. - closed := make(chan error, 1) - go func() { closed <- s.rpc.Close() }() - select { - case err := <-closed: - return errors.Join(s.cancelErr, err) - case <-ctx.Done(): - return ctx.Err() - } -} - -func (s *Session) cancelNativeWork() { - defer close(s.cancelReady) - s.stopFunctionCalls() - turnID, active := s.stopSteering() - // Best effort: a known Turn must use its native identity. An explicit - // empty ID invokes native startup cancellation before turn/started. - if threadID := s.currentThreadID(); threadID != "" && active { - ctx, cancel := context.WithTimeout(context.Background(), 2*time.Second) - defer cancel() - _, _ = s.rpc.request(ctx, "turn/interrupt", TurnInterruptParams{ThreadID: threadID, TurnID: turnID}, func(frame any) error { - return s.rpc.writeFrameContext(ctx, frame) - }) - } - if s.subagents != nil { - s.cancelErr = s.cancelSubagentWork(context.Background()) - } - s.cancelFn() - _ = s.rpc.Close() -} diff --git a/apps/daemon/internal/agent/codex/session_cancel_test.go b/apps/daemon/internal/agent/codex/session_cancel_test.go deleted file mode 100644 index 32e794b6a..000000000 --- a/apps/daemon/internal/agent/codex/session_cancel_test.go +++ /dev/null @@ -1,154 +0,0 @@ -package codex - -import ( - "context" - "encoding/json" - "errors" - "io" - "sync" - "testing" - "time" - - "github.com/MiniMax-AI/OpenAgentCore/internal/agentdaemon/proto" -) - -type interruptRequest struct { - ID string `json:"id"` - Method string `json:"method"` - Params struct { - ThreadID string `json:"threadId"` - TurnID *string `json:"turnId"` - } `json:"params"` -} - -func TestCancelUsesNativeTurnIdentity(t *testing.T) { - for _, name := range []string{"running", "foreign notification", "startup", "completed", "before thread", "rejected", "response timeout"} { - t.Run(name, func(t *testing.T) { - s, client, server := cancellationTestSession(t) - if name != "before thread" { - s.setThreadID("native-thread") - } - if name != "startup" && name != "before thread" { - s.onTurnStarted(json.RawMessage(`{"threadId":"native-thread","turn":{"id":"native-turn"}}`)) - } - if name == "foreign notification" { - s.onTurnStarted(json.RawMessage(`{"threadId":"foreign-thread","turn":{"id":"foreign-turn"}}`)) - } - if name == "completed" { - s.onTurnCompleted(json.RawMessage(`{"threadId":"native-thread","turn":{"id":"native-turn","status":"completed"}}`)) - } - requests := collectCancellationRequests(t, server, name) - var calls sync.WaitGroup - for range 3 { - calls.Add(1) - go func() { - defer calls.Done() - if err := s.Cancel(context.Background()); err != nil { - t.Errorf("best-effort cancel: %v", err) - } - }() - } - finished := make(chan struct{}) - go func() { calls.Wait(); close(finished) }() - select { - case <-finished: - case <-time.After(4 * time.Second): - t.Fatal("cancellation did not finish after the response deadline") - } - if client.Alive() || s.cancelCtx.Err() == nil { - t.Fatal("cancellation did not close the native client and context") - } - got := <-requests - if name == "completed" || name == "before thread" { - if len(got) != 0 { - t.Fatalf("unexpected native interrupt: %+v", got) - } - } else { - want := "native-turn" - if name == "startup" { - want = "" - } - if len(got) != 1 || got[0].Method != "turn/interrupt" || got[0].Params.ThreadID != "native-thread" || got[0].Params.TurnID == nil || *got[0].Params.TurnID != want { - t.Fatalf("incorrect or repeated native cancellation: %+v", got) - } - } - if name != "before thread" && s.CancellationOutcome().Metadata[proto.DoneMetaAgentSessionID] != "native-thread" { - t.Fatal("cancellation lost native continuation identity") - } - }) - } -} - -func TestCancelRacingTurnStartedKeepsValidNativeTarget(t *testing.T) { - for range 32 { - s, _, server := cancellationTestSession(t) - s.setThreadID("native-thread") - requests := collectCancellationRequests(t, server, "running") - start := make(chan struct{}) - observed := make(chan struct{}) - go func() { - <-start - s.onTurnStarted(json.RawMessage(`{"threadId":"native-thread","turn":{"id":"native-turn"}}`)) - close(observed) - }() - close(start) - if err := s.Cancel(context.Background()); err != nil { - t.Fatal(err) - } - <-observed - got := <-requests - if len(got) != 1 || got[0].Params.ThreadID != "native-thread" || got[0].Params.TurnID == nil { - t.Fatalf("missing cancellation identity: %+v", got) - } - if id := *got[0].Params.TurnID; id != "" && id != "native-turn" { - t.Fatalf("foreign native Turn ID: %q", id) - } - if id, active := s.stopSteering(); active || id != *got[0].Params.TurnID { - t.Fatalf("late notification changed a stopped target: %q, %v", id, active) - } - } -} - -func cancellationTestSession(t *testing.T) (*Session, *TestClient, ServerSide) { - t.Helper() - client, server, cleanup := NewTestClient() - t.Cleanup(cleanup) - ctx, cancel := context.WithCancel(context.Background()) - t.Cleanup(cancel) - s := &Session{rpc: client.JSONRPCClient, cancelCtx: ctx, cancelFn: cancel, - cfg: defaultSessionConfig(), out: make(chan proto.Envelope, 8), bufs: NewItemBuffers()} - return s, client, server -} - -func collectCancellationRequests(t *testing.T, server ServerSide, mode string) <-chan []interruptRequest { - t.Helper() - done := make(chan []interruptRequest, 1) - go func() { - var requests []interruptRequest - defer func() { done <- requests }() - decoder := json.NewDecoder(server.FromClient) - for { - var request interruptRequest - if err := decoder.Decode(&request); err != nil { - if !errors.Is(err, io.EOF) { - t.Errorf("read cancellation: %v", err) - } - return - } - requests = append(requests, request) - if mode == "response timeout" { - continue - } - reply := map[string]any{"id": request.ID, "result": map[string]any{}} - if mode == "rejected" { - delete(reply, "result") - reply["error"] = map[string]any{"code": -32600, "message": "no active turn to interrupt"} - } - if err := json.NewEncoder(server.ToClient).Encode(reply); err != nil { - t.Errorf("reply to cancellation: %v", err) - return - } - } - }() - return done -} diff --git a/apps/daemon/internal/agent/codex/session_cancel_write_test.go b/apps/daemon/internal/agent/codex/session_cancel_write_test.go deleted file mode 100644 index bb2abe700..000000000 --- a/apps/daemon/internal/agent/codex/session_cancel_write_test.go +++ /dev/null @@ -1,185 +0,0 @@ -package codex - -import ( - "context" - "encoding/json" - "io" - "os" - "strings" - "testing" - "time" -) - -func TestCancelReleasesBlockedControlWrite(t *testing.T) { - for _, concurrent := range []bool{false, true} { - name := "interrupt" - if concurrent { - name = "another write" - } - t.Run(name, func(t *testing.T) { - s, client, server := cancellationTestSession(t) - s.setThreadID("native-thread") - s.onTurnStarted(json.RawMessage(`{"threadId":"native-thread","turn":{"id":"native-turn"}}`)) - _, peer, peerServer := cancellationTestSession(t) - cancelDone := make(chan struct{}) - writeDone := make(chan struct{}) - var writeErr error - t.Cleanup(func() { - _ = client.Close() - select { - case <-cancelDone: - case <-time.After(4 * time.Second): - t.Error("cancellation did not stop after test cleanup") - } - if concurrent { - select { - case <-writeDone: - case <-time.After(4 * time.Second): - t.Error("concurrent writer did not stop after test cleanup") - } - } - }) - if concurrent { - go func() { - writeErr = client.Notify("blocked", nil) - close(writeDone) - }() - readControlPrefix(t, server.FromClient) - } - go func() { - defer close(cancelDone) - if err := s.Cancel(context.Background()); err != nil { - t.Errorf("best-effort cancellation: %v", err) - } - }() - if !concurrent { - readControlPrefix(t, server.FromClient) - } - // The pipe remains undrained after one byte, so the write cannot complete. - select { - case <-cancelDone: - case <-time.After(4 * time.Second): - t.Fatal("blocked control write exceeded the interrupt budget without cleanup") - } - if client.Alive() || s.cancelCtx.Err() == nil { - t.Fatal("cancellation did not close its client and context") - } - client.pendingMu.Lock() - pending := len(client.pending) - client.pendingMu.Unlock() - if pending != 0 { - t.Fatal("blocked cancellation left pending requests") - } - if concurrent { - select { - case <-writeDone: - if writeErr == nil { - t.Fatal("partial concurrent write reported success") - } - case <-time.After(time.Second): - t.Fatal("blocked concurrent writer was not released") - } - } - replied := make(chan error, 1) - go func() { - _, err := SendCannedResponse(t.Context(), peerServer, "alive") - replied <- err - }() - ctx, cancel := context.WithTimeout(t.Context(), time.Second) - defer cancel() - result, err := peer.Request(ctx, "echo", nil) - if err != nil || string(result) != `"alive"` || !peer.Alive() { - t.Fatal("cancellation affected an independent client", err) - } - if err := <-replied; err != nil { - t.Fatal(err) - } - }) - } -} - -func readControlPrefix(t *testing.T, reader io.Reader) { - t.Helper() - read := make(chan error, 1) - go func() { - var prefix [1]byte - _, err := io.ReadFull(reader, prefix[:]) - read <- err - }() - select { - case err := <-read: - if err != nil { - t.Fatal("control write did not start", err) - } - case <-time.After(4 * time.Second): - t.Fatal("control write did not enter the pipe") - } -} - -func TestCancelReleasesBlockedNativeProcess(t *testing.T) { - client := NewJSONRPCClient(JSONRPCConfig{ - Binary: os.Args[0], ExtraArgs: []string{"-test.run=TestJSONRPCClientFakeCodexProcess", "--"}, - Env: append(os.Environ(), "CODEX_RPC_FAKE_PROCESS=1", "CODEX_RPC_FAKE_BLOCK_WRITE=1", "GORACE=atexit_sleep_ms=0"), - LogTag: "cancel-blocked-process", RequestTimeout: 2 * time.Second, - }) - t.Cleanup(func() { _ = client.Close() }) - blocked := make(chan struct{}) - client.OnNotification("test/write_blocked", func(json.RawMessage) { close(blocked) }) - if _, err := client.Start(t.Context(), InitializeParams{ClientInfo: InitializeClientInfo{Name: "test", Version: "0"}}); err != nil { - t.Fatal(err) - } - written := make(chan error, 1) - go func() { written <- client.Notify("blocked", strings.Repeat("x", 1<<20)) }() - t.Cleanup(func() { - _ = client.Close() - select { - case <-written: - case <-time.After(4 * time.Second): - t.Error("native pipe writer did not stop during cleanup") - } - }) - select { - case <-blocked: - case <-time.After(4 * time.Second): - t.Fatal("native child did not stop reading the oversized frame") - } - ctx, cancel := context.WithCancel(context.Background()) - defer cancel() - s := &Session{rpc: client, cancelCtx: ctx, cancelFn: cancel, - cfg: defaultSessionConfig(), bufs: NewItemBuffers()} - s.setThreadID("native-thread") - s.onTurnStarted(json.RawMessage(`{"threadId":"native-thread","turn":{"id":"native-turn"}}`)) - cancelDone := make(chan struct{}) - go func() { _ = s.Cancel(context.Background()); close(cancelDone) }() - t.Cleanup(func() { - _ = client.Close() - select { - case <-cancelDone: - case <-time.After(4 * time.Second): - t.Error("native cancellation did not stop during cleanup") - } - }) - select { - case <-cancelDone: - case <-time.After(4 * time.Second): - t.Fatal("blocked native stdin prevented cancellation cleanup") - } - select { - case <-client.Done(): - case <-time.After(time.Second): - t.Fatal("cancelled native child was not reaped") - } - if _, exited := client.process.ExitCode(); client.Alive() || !exited || ctx.Err() == nil { - t.Fatal("native process or cancellation context remained active") - } - // The cleanup consumes the writer's result after proving it was released. - select { - case err := <-written: - written <- err - if err == nil { - t.Fatal("undrained native frame reported success") - } - case <-time.After(time.Second): - t.Fatal("cancelled native writer remained blocked") - } -} diff --git a/apps/daemon/internal/agent/codex/session_log_test.go b/apps/daemon/internal/agent/codex/session_log_test.go index d0ec7977c..ae51381c7 100644 --- a/apps/daemon/internal/agent/codex/session_log_test.go +++ b/apps/daemon/internal/agent/codex/session_log_test.go @@ -247,7 +247,10 @@ func drainEnvelopes(out <-chan proto.Envelope) []proto.Envelope { var got []proto.Envelope for { select { - case env := <-out: + case env, open := <-out: + if !open { + return got + } got = append(got, env) default: return got diff --git a/apps/daemon/internal/agent/codex/session_run.go b/apps/daemon/internal/agent/codex/session_run.go deleted file mode 100644 index 3394acf90..000000000 --- a/apps/daemon/internal/agent/codex/session_run.go +++ /dev/null @@ -1,52 +0,0 @@ -package codex - -import ( - "context" - "fmt" - "strings" - "time" - - "github.com/MiniMax-AI/OpenAgentCore/internal/agentdaemon/proto" -) - -func (s *Session) startNative(ctx context.Context, plan SessionPlan, req proto.PromptRequestPayload) error { - if s.currentThreadID() == "" { - if err := s.resolveThread(req, plan); err != nil { - return err - } - } - - input, err := nativeInput(req.Input) - if err != nil { - return err - } - turnParams := TurnStartParams{ - ThreadID: s.currentThreadID(), - Input: input, - } - model := strings.TrimSpace(s.resolvedModel) - if model == "" { - return fmt.Errorf("codex: collaboration mode requires a resolved model") - } - var developerInstructions *string - if plan.SystemPrompt != "" { - developerInstructions = &plan.SystemPrompt - } - turnParams.CollaborationMode = &CollaborationMode{ - Mode: CollaborationModeDefault, - Settings: CollaborationModeSettings{ - ReasoningEffort: plan.ModelReasoningEffort, - Model: model, - DeveloperInstructions: developerInstructions, - }, - } - turnCtx, turnCancel := context.WithTimeout(ctx, 10*time.Second) - _, ackErr := s.rpc.requestWithResult(turnCtx, "turn/start", turnParams, s.bindTurnResult) - turnCancel() - if ackErr != nil { - s.cfg.logger.Warn("codex: turn/start ack failed", "run_id", s.runID, "err", ackErr) - return fmt.Errorf("codex: turn/start: %w", ackErr) - } - - return nil -} diff --git a/apps/daemon/internal/agent/codex/session_steering_lifecycle_test.go b/apps/daemon/internal/agent/codex/session_steering_lifecycle_test.go index 6dc663ea0..de8ea9aaa 100644 --- a/apps/daemon/internal/agent/codex/session_steering_lifecycle_test.go +++ b/apps/daemon/internal/agent/codex/session_steering_lifecycle_test.go @@ -83,13 +83,8 @@ func TestBlockedSteeringWriteEndsRunWithTerminalFrames(t *testing.T) { ctx, cancel := context.WithTimeout(context.Background(), 3*time.Second) defer cancel() out := make(chan proto.Envelope, 8) - turnCtx, cancelTurn := context.WithCancel(ctx) - s := &Session{ - executor: &Executor{}, runID: "run", rpc: client.JSONRPCClient, cancelCtx: turnCtx, cancelFn: cancelTurn, out: out, resolvedModel: "synthetic", - cfg: sessionConfig{logger: obslog.Bg()}, waitDone: make(chan struct{}), outputDone: make(chan struct{}), cleanup: func() {}, - bufs: NewItemBuffers(), - } - s.registerHandlers() + e := &Executor{base: &Session{rpc: client.JSONRPCClient, cancelCtx: ctx, functions: &functionCalls{}, resolvedModel: "synthetic", + cfg: sessionConfig{logger: obslog.Bg()}}, plan: SessionPlan{Model: "synthetic"}} ready := make(chan error, 1) go func() { decoder := json.NewDecoder(server.FromClient) @@ -120,11 +115,14 @@ func TestBlockedSteeringWriteEndsRunWithTerminalFrames(t *testing.T) { } ready <- nil }() - // Executor.StartTurn starts and settles each Turn through these two steps. - go s.settleExecutorTurn(s.startNative(ctx, SessionPlan{Model: "synthetic"}, proto.PromptRequestPayload{Input: proto.TextInput("first")})) + turn, err := e.StartTurn(ctx, "run", proto.TextInput("first"), out) + if err != nil { + t.Fatal(err) + } if err := <-ready; err != nil { t.Fatal(err) } + s := turn.(*Session) callCtx, callCancel := context.WithTimeout(ctx, 50*time.Millisecond) defer callCancel() if err := s.Steer(callCtx, proto.PromptSteerPayload{InputID: "blocked", Input: proto.TextInput("extra")}); err == nil { diff --git a/apps/daemon/internal/agent/codex/session_usage_live_test.go b/apps/daemon/internal/agent/codex/session_usage_live_test.go index 58ccc1d44..a58f1c3b4 100644 --- a/apps/daemon/internal/agent/codex/session_usage_live_test.go +++ b/apps/daemon/internal/agent/codex/session_usage_live_test.go @@ -1,7 +1,6 @@ package codex import ( - "context" "encoding/json" "io" "runtime" @@ -53,9 +52,10 @@ func TestUsagePublishedBeforeCompletion(t *testing.T) { } func TestUsageBackpressureKeepsSnapshotReadableAndCompletionOrdered(t *testing.T) { - s, _, server := cancellationTestSession(t) + client, server, cleanup := NewTestClient() + t.Cleanup(cleanup) out := make(chan proto.Envelope, 1) - s.out = out + s := &Session{rpc: client.JSONRPCClient, cancelCtx: t.Context(), cfg: defaultSessionConfig(), out: out, bufs: NewItemBuffers()} s.out <- proto.Envelope{Type: "occupied"} s.registerHandlers() s.setThreadID("thread") @@ -92,35 +92,6 @@ func TestUsageBackpressureKeepsSnapshotReadableAndCompletionOrdered(t *testing.T } } -func TestCancellationReleasesUsageBackpressure(t *testing.T) { - s, _, server := cancellationTestSession(t) - s.out = make(chan proto.Envelope, 1) - s.out <- proto.Envelope{Type: "occupied"} - s.registerHandlers() - s.setThreadID("thread") - s.onTurnStarted(json.RawMessage(`{"threadId":"thread","turn":{"id":"turn"}}`)) - requests := collectCancellationRequests(t, server, "response timeout") - written := make(chan error, 1) - go func() { - _, err := io.WriteString(server.ToClient, activeUsageNotification+"\n") - written <- err - }() - waitForUsageBackpressure(t, s) - assertReadableUsageSnapshot(t, s) - ctx, cancel := context.WithTimeout(t.Context(), 4*time.Second) - defer cancel() - if err := s.Cancel(ctx); err != nil { - t.Fatal("cancellation could not release usage backpressure", err) - } - if s.cancelCtx.Err() == nil || len(<-requests) != 1 { - t.Fatal("native cancellation did not settle") - } - if err := <-written; err != nil { - t.Fatal(err) - } - assertReadableUsageSnapshot(t, s) -} - func waitForUsageBackpressure(t *testing.T, s *Session) { t.Helper() deadline := time.Now().Add(time.Second) diff --git a/apps/daemon/internal/agent/codex/subagent_cancel_test.go b/apps/daemon/internal/agent/codex/subagent_cancel_test.go index f953f65bb..50406eb36 100644 --- a/apps/daemon/internal/agent/codex/subagent_cancel_test.go +++ b/apps/daemon/internal/agent/codex/subagent_cancel_test.go @@ -21,7 +21,7 @@ func TestSubagentCancelDeadlineRetainsOwnerAndRetry(t *testing.T) { f.mu.Lock() f.interruptGate = gate f.mu.Unlock() - s.emitDoneAt("frozen", nil, nil) + s.onTurnCompleted(rootCompleted) for range 4 { select { case e := <-out: @@ -86,12 +86,12 @@ func TestSubagentCancelRetainsActualObservationFailure(t *testing.T) { t.Run(map[bool]string{false: "invalid history", true: "observer stopped"}[interrupted], func(t *testing.T) { s, f, out := observationSession(t, "completed") if interrupted { - s.subagents.cancel() + s.cancelFn() } else { if err := os.WriteFile(filepath.Join(f.home, "sessions", "root.jsonl"), []byte("invalid\n"), 0600); err != nil { t.Fatal(err) } - s.emitDoneAt("frozen", nil, nil) + s.onTurnCompleted(rootCompleted) } select { case <-s.subagents.done: diff --git a/apps/daemon/internal/agent/codex/subagent_observations_test.go b/apps/daemon/internal/agent/codex/subagent_observations_test.go index a96258ca7..ce61264e4 100644 --- a/apps/daemon/internal/agent/codex/subagent_observations_test.go +++ b/apps/daemon/internal/agent/codex/subagent_observations_test.go @@ -75,11 +75,14 @@ func observationSession(t *testing.T, status string) (*Session, *subagentFixture t.Fatal(err) } f.persist(t) - s := &Session{runID: "run", nativeHome: agent.ViewDir{Host: f.home, View: f.home}, rpc: client.JSONRPCClient, out: out, cancelCtx: ctx, cancelFn: cancel, cfg: defaultSessionConfig(), bufs: NewItemBuffers()} + s := &Session{runID: "run", nativeHome: agent.ViewDir{Host: f.home, View: f.home}, rpc: client.JSONRPCClient, out: out, cancelCtx: ctx, cancelFn: cancel, cfg: defaultSessionConfig(), bufs: NewItemBuffers(), + waitDone: make(chan struct{}), outputDone: make(chan struct{})} s.setThreadID("root") s.beginRootTurn("root", "root-turn") s.startSubagentObservations() s.registerHandlers() + // Settle as Executor.StartTurn does after a confirmed native start. + go s.settleExecutorTurn(nil) go func() { decoder := json.NewDecoder(server.FromClient) for { @@ -98,6 +101,8 @@ func observationSession(t *testing.T, status string) (*Session, *subagentFixture result = map[string]any{"thread": f.history(id)} case "thread/turns/list": result = map[string]any{"data": f.history(id).Turns, "nextCursor": nil} + case "thread/backgroundTerminals/list": + result = map[string]any{"data": []any{}} case "thread/list": result = map[string]any{"data": []any{map[string]any{"id": "child", "parentThreadId": "root", "createdAt": 100, "agentNickname": "Child", "source": map[string]any{"subAgent": map[string]any{"thread_spawn": map[string]any{"parent_thread_id": "root"}}}}}, "nextCursor": nil} case "turn/interrupt": @@ -132,6 +137,8 @@ func observationSession(t *testing.T, status string) (*Session, *subagentFixture return s, f, out } +var rootCompleted = json.RawMessage(`{"threadId":"root","turn":{"id":"root-turn","status":"completed"}}`) + func collectObserved(t *testing.T, out <-chan proto.Envelope) []proto.Envelope { t.Helper() var all []proto.Envelope @@ -209,7 +216,7 @@ func TestSubagentRootFirstRetainsChildUntilTerminal(t *testing.T) { func TestSubagentCancellationAfterRootFrozenCollectsNativeTerminal(t *testing.T) { s, f, out := observationSession(t, "inProgress") - s.emitDoneAt("frozen", nil, nil) + s.onTurnCompleted(rootCompleted) for i := 0; i < 4; i++ { select { case <-out: diff --git a/apps/daemon/internal/agent/codex/terminal_cleanup_test.go b/apps/daemon/internal/agent/codex/terminal_cleanup_test.go index c0a208c6b..ef9c70689 100644 --- a/apps/daemon/internal/agent/codex/terminal_cleanup_test.go +++ b/apps/daemon/internal/agent/codex/terminal_cleanup_test.go @@ -39,7 +39,7 @@ func TestExecutorCancellationRequiresNativeTerminalCleanup(t *testing.T) { } for range out { } - if !e.prepared.session.rpc.Alive() { + if !e.base.rpc.Alive() { t.Fatal("native owner lost before cleanup/reuse") } frames := preparationFrames(t, root) @@ -65,7 +65,7 @@ func TestExecutorCancellationRequiresNativeTerminalCleanup(t *testing.T) { t.Fatal("confirmed cleanup prevented reuse", err) } } else { - if err := e.Close(ctx); err == nil || !e.prepared.session.rpc.Alive() { + if err := e.Close(ctx); err == nil || !e.base.rpc.Alive() { t.Fatal("failed cleanup released the owned native process", err) } allow() @@ -107,7 +107,7 @@ func TestExecutorCloseCleansNativeTerminalsAfterNormalCompletion(t *testing.T) { terminated = true } } - if !terminated || e.prepared.session.rpc.Alive() { + if !terminated || e.base.rpc.Alive() { t.Fatal("Executor Close left native terminals or owner alive") } } From 5b8089dfdc45b0bba04630d78b6ca190a138d234 Mon Sep 17 00:00:00 2001 From: SaladDay <1203511142@qq.com> Date: Wed, 7 Oct 2026 16:50:02 +0000 Subject: [PATCH 2/4] Return a nil Codex Executor on rejected preparation PrepareExecutor returned newExecutor's nil *Executor as a non-nil agent.Executor when preparation was rejected before any resource existed, such as a read-only or structured-output request. The daemon then kept it as the native Executor and its cleanup called Close on a nil receiver. Return a nil interface, and test read-only rejection through the factory. --- .../daemon/internal/agent/codex/environment_retired_test.go | 4 ++-- apps/daemon/internal/agent/codex/executor.go | 6 +++++- 2 files changed, 7 insertions(+), 3 deletions(-) diff --git a/apps/daemon/internal/agent/codex/environment_retired_test.go b/apps/daemon/internal/agent/codex/environment_retired_test.go index 8ff7e6b4f..a7d2f1721 100644 --- a/apps/daemon/internal/agent/codex/environment_retired_test.go +++ b/apps/daemon/internal/agent/codex/environment_retired_test.go @@ -11,9 +11,9 @@ import ( ) func TestReadOnlyPreparationRejectedBeforeNativeSetup(t *testing.T) { - req, cfg, root := preparationFixture(t) + req, _, root := preparationFixture(t) req.WorkspaceReadOnly = true - prepared, err := newExecutor(t.Context(), req, cfg) + prepared, err := PrepareExecutor(t.Context(), req) if err == nil || prepared != nil { t.Fatal("read-only request admitted", err) } diff --git a/apps/daemon/internal/agent/codex/executor.go b/apps/daemon/internal/agent/codex/executor.go index 3faacfcf3..0748de54f 100644 --- a/apps/daemon/internal/agent/codex/executor.go +++ b/apps/daemon/internal/agent/codex/executor.go @@ -26,7 +26,11 @@ type Executor struct { } func PrepareExecutor(ctx context.Context, req proto.PromptRequestPayload) (agent.Executor, error) { - return newExecutor(ctx, req, defaultSessionConfig()) + e, err := newExecutor(ctx, req, defaultSessionConfig()) + if e == nil { + return nil, err + } + return e, err } func (e *Executor) StartTurn(ctx context.Context, runID string, input proto.MessageInput, out chan<- proto.Envelope) (agent.Turn, error) { From a0b3c599a4afe2a116e65e55a02b4b31bac3974f Mon Sep 17 00:00:00 2001 From: SaladDay <1203511142@qq.com> Date: Wed, 7 Oct 2026 17:04:51 +0000 Subject: [PATCH 3/4] Delete the unset Codex harness binary override No deployment, image or script sets OAC_RUNTIME_CODEX_HARNESS_BIN, so the local Environment declaration never read a value from it. --- apps/daemon/internal/agent/codex/environment_local.go | 3 +-- 1 file changed, 1 insertion(+), 2 deletions(-) diff --git a/apps/daemon/internal/agent/codex/environment_local.go b/apps/daemon/internal/agent/codex/environment_local.go index 9fa3f3c2e..9f59aada0 100644 --- a/apps/daemon/internal/agent/codex/environment_local.go +++ b/apps/daemon/internal/agent/codex/environment_local.go @@ -1,7 +1,6 @@ package codex import ( - "os" "runtime" "github.com/MiniMax-AI/OpenAgentCore/apps/daemon/internal/localworkspace" @@ -9,7 +8,7 @@ import ( // SupportsLocalEnvironment checks deployment prerequisites, not public admission. func SupportsLocalEnvironment(version string) bool { - if !SupportsNativeSessionRecovery(version) || os.Getenv("OAC_RUNTIME_CODEX_HARNESS_BIN") != "" || (runtime.GOOS != "linux" && runtime.GOOS != "darwin" && runtime.GOOS != "windows") { + if !SupportsNativeSessionRecovery(version) || (runtime.GOOS != "linux" && runtime.GOOS != "darwin" && runtime.GOOS != "windows") { return false } binding, err := localworkspace.Load() From f2f0709374de814d4607f21fe5be7d1a592ff534 Mon Sep 17 00:00:00 2001 From: SaladDay <1203511142@qq.com> Date: Wed, 7 Oct 2026 17:05:37 +0000 Subject: [PATCH 4/4] Wait for the fake Codex server before removing its home The Turn Cancel can send turn/interrupt before settlement closes the client, and settlement does not wait for that reply. The observation fixture's server could then persist history after the test removed its temporary home. Cleanup now waits for the server goroutine. --- .../daemon/internal/agent/codex/subagent_observations_test.go | 4 ++++ 1 file changed, 4 insertions(+) diff --git a/apps/daemon/internal/agent/codex/subagent_observations_test.go b/apps/daemon/internal/agent/codex/subagent_observations_test.go index ce61264e4..1056176a8 100644 --- a/apps/daemon/internal/agent/codex/subagent_observations_test.go +++ b/apps/daemon/internal/agent/codex/subagent_observations_test.go @@ -83,7 +83,9 @@ func observationSession(t *testing.T, status string) (*Session, *subagentFixture s.registerHandlers() // Settle as Executor.StartTurn does after a confirmed native start. go s.settleExecutorTurn(nil) + served := make(chan struct{}) go func() { + defer close(served) decoder := json.NewDecoder(server.FromClient) for { var request JsonRpcRequest @@ -128,6 +130,8 @@ func observationSession(t *testing.T, status string) (*Session, *subagentFixture t.Cleanup(func() { cancel() cleanup() + // A reply may still be persisting after settlement closed the client. + <-served select { case <-s.subagents.done: case <-time.After(time.Second):