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_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() diff --git a/apps/daemon/internal/agent/codex/environment_retired_test.go b/apps/daemon/internal/agent/codex/environment_retired_test.go index f3e070843..a7d2f1721 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" @@ -10,9 +11,9 @@ import ( ) func TestReadOnlyPreparationRejectedBeforeNativeSetup(t *testing.T) { - req, cfg, root := preparationFixture(t) + req, _, root := preparationFixture(t) req.WorkspaceReadOnly = true - prepared, err := newPreparation(t.Context(), req, cfg) + prepared, err := PrepareExecutor(t.Context(), req) 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..0748de54f 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,30 +15,22 @@ 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 - } + e, err := newExecutor(ctx, req, defaultSessionConfig()) + if e == nil { return nil, err } - p.mu.Lock() - p.started = true - close(p.transferred) - p.mu.Unlock() - return &Executor{prepared: p}, nil + return e, err } func (e *Executor) StartTurn(ctx context.Context, runID string, input proto.MessageInput, out chan<- proto.Envelope) (agent.Turn, error) { @@ -48,7 +41,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 +61,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 +86,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 +132,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 +177,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..1056176a8 100644 --- a/apps/daemon/internal/agent/codex/subagent_observations_test.go +++ b/apps/daemon/internal/agent/codex/subagent_observations_test.go @@ -75,12 +75,17 @@ 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) + served := make(chan struct{}) go func() { + defer close(served) decoder := json.NewDecoder(server.FromClient) for { var request JsonRpcRequest @@ -98,6 +103,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": @@ -123,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): @@ -132,6 +141,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 +220,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") } }