diff --git a/apps/daemon/internal/agent/mcode/preparation_log_test.go b/apps/daemon/internal/agent/mcode/preparation_log_test.go new file mode 100644 index 000000000..1e1fd6b2c --- /dev/null +++ b/apps/daemon/internal/agent/mcode/preparation_log_test.go @@ -0,0 +1,86 @@ +package mcode + +import ( + "bytes" + "context" + "encoding/json" + "log/slog" + "strings" + "testing" + "time" +) + +func TestPreparationLogsStagesWithoutNativePayloads(t *testing.T) { + for _, tc := range []struct { + name, scenario, phase, outcome string + resume, success bool + }{ + {"initialize", "prepare-reject-initialize", "initialize", "active", false, false}, + {"new", "prepare-reject-session/new", "session/new", "active", false, false}, + {"load", "prepare-reject-session/load", "session/load", "active", true, false}, + {"model-selection", "unknown-model", "model_selection", "active", false, false}, + {"model-configuration", "prepare-reject-session/set_config_option", "model_configuration", "active", false, false}, + {"success", "complete", "model_configuration", "active", false, true}, + {"cancelled", "hang", "initialize", "cancelled", false, false}, + {"deadline", "hang", "initialize", "deadline_exceeded", false, false}, + } { + t.Run(tc.name, func(t *testing.T) { + var logs bytes.Buffer + previous := slog.Default() + slog.SetDefault(slog.New(slog.NewJSONHandler(&logs, nil))) + defer slog.SetDefault(previous) + req := testRequest(t) + req.SystemPrompt = "PRIVATE_INSTRUCTIONS" + req.ModelProvider.APIKey = "PRIVATE_MODEL_KEY" + if tc.resume { + req.AgentSessionID = "PRIVATE_NATIVE_SESSION" + } + ctx := t.Context() + if tc.outcome == "deadline_exceeded" { + var cancel context.CancelFunc + ctx, cancel = context.WithTimeout(ctx, 200*time.Millisecond) + defer cancel() + } else if tc.outcome == "cancelled" { + var cancel context.CancelFunc + ctx, cancel = context.WithCancel(ctx) + defer cancel() + timer := time.AfterFunc(200*time.Millisecond, cancel) + defer timer.Stop() + } + e, err := prepareExecutor(t, ctx, helperInstall(t, tc.scenario, ""), req) + if (err == nil) != tc.success || (e != nil) != tc.success { + t.Fatalf("preparation result changed: owner=%v err=%v", e != nil, err) + } + if strings.HasPrefix(tc.scenario, "prepare-reject-") && !strings.Contains(err.Error(), "PRIVATE_NATIVE_ERROR private-token") { + t.Fatal("native error was changed") + } + if e != nil { + if err := e.Close(t.Context()); err != nil { + t.Fatal(err) + } + } + var entry map[string]any + if err := json.Unmarshal(bytes.TrimSpace(logs.Bytes()), &entry); err != nil { + t.Fatalf("expected exactly one preparation log: %v", err) + } + if entry["msg"] != "mcode preparation" || entry["phase"] != tc.phase || entry["success"] != tc.success || entry["context_outcome"] != tc.outcome { + t.Fatalf("incorrect preparation outcome: %v", entry) + } + if duration, ok := entry["duration_ms"].(float64); !ok || duration < 0 { + t.Fatalf("invalid preparation duration: %v", entry) + } + for key := range entry { + switch key { + case "time", "level", "msg", "phase", "duration_ms", "success", "context_outcome": + default: + t.Fatalf("unexpected diagnostic field %q", key) + } + } + for _, secret := range []string{"PRIVATE_", "private-token", req.ModelProvider.BaseURL, "m:custom_provider", "fixture"} { + if strings.Contains(logs.String(), secret) { + t.Fatal("private preparation data in log") + } + } + }) + } +} diff --git a/apps/daemon/internal/agent/mcode/session.go b/apps/daemon/internal/agent/mcode/session.go index 6d87ddb21..8a764343c 100644 --- a/apps/daemon/internal/agent/mcode/session.go +++ b/apps/daemon/internal/agent/mcode/session.go @@ -3,6 +3,7 @@ package mcode import ( "context" "encoding/json" + "errors" "fmt" "net/url" "slices" @@ -13,6 +14,7 @@ import ( "github.com/MiniMax-AI/OpenAgentCore/apps/daemon/internal/agent" "github.com/MiniMax-AI/OpenAgentCore/apps/daemon/internal/agent/clirunner" "github.com/MiniMax-AI/OpenAgentCore/internal/agentdaemon/proto" + obslog "github.com/MiniMax-AI/OpenAgentCore/internal/obs/log" ) type Session struct { @@ -73,7 +75,21 @@ func newTurnSession(ctx context.Context, req proto.PromptRequestPayload, opts la return &Session{ctx: ctx, req: req, opts: opts, connection: c, out: out, frames: make(chan rpcFrame, 32), finished: make(chan struct{}), tools: map[string]toolUpdate{}, completedTools: map[string]bool{}, completedMessages: map[string]bool{}} } -func (s *Session) prepareNative() error { +func (s *Session) prepareNative() (err error) { + started := time.Now() + phase := "initialize" + defer func() { + // Native errors and configuration may contain credentials. Log only + // static stages and the owner context, before preparation cleanup. + contextOutcome := "active" + switch { + case errors.Is(s.ctx.Err(), context.DeadlineExceeded): + contextOutcome = "deadline_exceeded" + case errors.Is(s.ctx.Err(), context.Canceled): + contextOutcome = "cancelled" + } + obslog.Ctx(s.ctx).Info("mcode preparation", "phase", phase, "duration_ms", time.Since(started).Milliseconds(), "success", err == nil, "context_outcome", contextOutcome) + }() var initialized struct { ProtocolVersion int `json:"protocolVersion"` Meta struct { @@ -105,6 +121,7 @@ func (s *Session) prepareNative() error { method = "session/load" params["sessionId"] = s.req.AgentSessionID } + phase = method var session sessionResult if err := s.call(method, params, &session, false); err != nil { return err @@ -118,6 +135,7 @@ func (s *Session) prepareNative() error { s.mu.Lock() s.sessionID = session.SessionID s.mu.Unlock() + phase = "model_selection" model, err := advertisedModel(session.ConfigOptions, s.opts.Model) if err != nil && s.req.AgentSessionID != "" && !slices.ContainsFunc(session.ConfigOptions, func(option configOption) bool { return option.ID == "model" }) { // Native load omits the selector when its persisted model was removed; selection still validates against the current catalog. @@ -128,6 +146,7 @@ func (s *Session) prepareNative() error { return err } s.nativeModel = model + phase = "model_configuration" if err := s.call("session/set_config_option", map[string]any{"sessionId": session.SessionID, "configId": "model", "value": model}, nil, false); err != nil { return err } diff --git a/apps/daemon/internal/agent/mcode/session_test.go b/apps/daemon/internal/agent/mcode/session_test.go index a8f02d2e5..f94b9a850 100644 --- a/apps/daemon/internal/agent/mcode/session_test.go +++ b/apps/daemon/internal/agent/mcode/session_test.go @@ -294,6 +294,10 @@ func TestMCodeProcess(t *testing.T) { } } } + if scenario == "prepare-reject-"+frame.Method { + send(rpcFrame{JSONRPC: "2.0", ID: frame.ID, Error: &rpcError{Code: -32603, Message: "PRIVATE_NATIVE_ERROR private-token"}}) + continue + } result := any(map[string]any{}) switch frame.Method { case "initialize": diff --git a/packages/mcode-harness/README.md b/packages/mcode-harness/README.md index 2a1a6f419..076de8bcb 100644 --- a/packages/mcode-harness/README.md +++ b/packages/mcode-harness/README.md @@ -6,7 +6,7 @@ One patch script (`patch-native.mjs`) patches the pinned native CLI source. The ## Workspace tools -The native process, its ACP Session and the workspace tools share the Session's workspace as their working directory; native configuration, Skills and history stay in the private Session data directory. Builtin file tools are disabled, so project files are read and written through the bridge's tools. Native Session creation and loading connect the required workspace bridge and validate its tool inventory before preparation succeeds, using the existing 15-second stdio discovery limit and the same connection pool as Turn execution. An unavailable or incomplete bridge fails preparation; optional MCP servers retain native lazy discovery. The expected inventory is generated from the worker's `--describe` output. Only the adapter registers this bridge; callers cannot supply its command, profile, working directory or environment. Native diff/undo capture is not provided by this path. Common Files and Artifacts use the same bound workspace. +The native process, its ACP Session and the workspace tools share the Session's workspace as their working directory; native configuration, Skills and history stay in the private Session data directory. Builtin file tools are disabled, so project files are read and written through the bridge's tools. Native Session creation and loading connect the required workspace bridge and validate its tool inventory before preparation succeeds, using the existing Session-scoped connection pool and its normal request timeout. The adapter’s ACP request and Core preparation deadlines continue to bound startup. Turn discovery retains its existing shorter timeout. An unavailable or incomplete bridge fails preparation; optional MCP servers retain native lazy discovery. The expected inventory is generated from the worker's `--describe` output. Only the adapter registers this bridge; callers cannot supply its command, profile, working directory or environment. Native diff/undo capture is not provided by this path. Common Files and Artifacts use the same bound workspace. For each tool call, the bridge starts `launch.mjs` with the Session's private profile (`workspace-profile.json`, written by the daemon). The launcher checks that the profile's `workspace` is a canonical absolute path, creates the `scratch` directory, and runs the worker in the workspace with the bridge's environment. An invalid profile rejects the call. @@ -36,7 +36,7 @@ Native tool schemas are retained. Text and image results use standard MCP conten `make check-mcode-harness` runs the package's Node tests and syntax checks. The native prompt regression requires `MCODE_SOURCE` pointing to the pinned, patched native source with upstream dependencies installed; the companion build always runs it. It checks ordinary native guidance, disabled local tools and protected root and child capabilities without a model. Qualify changes with the [Harness acceptance checklist](../../contracts/agents-api/harness-onboarding.md#qualify-the-adapter); synthetic probes and native model runs do not complete public Files/Artifacts or independent Core acceptance. -With the same `MCODE_SOURCE`, `native-mcp-lifecycle.test.mjs` exercises the patched native Session configurations, connection pool and MCP SDK against controlled transports. It checks exact Session/server selection, delayed close events, connection attempts interrupted during initialization, and later lazy reconnection without closing unaffected services. It also checks workspace readiness before preparation completes, rejection of incomplete or failed discovery, and reuse of the prepared connection. +With the same `MCODE_SOURCE`, `native-mcp-lifecycle.test.mjs` exercises the patched native Session configurations, connection pool and MCP SDK against controlled transports. It checks exact Session/server selection, delayed close events, connection attempts interrupted during initialization, and later lazy reconnection without closing unaffected services. It also checks workspace readiness before preparation completes, rejection of incomplete or failed discovery, reuse of the prepared connection, preparation beyond the shorter Turn discovery timeout, configured connection timeouts and shutdown during required discovery. For the packaged Linux regression, provide an operator-owned private profile and artifact directory, then run `native.test.mjs` inside the qualified Docker Runtime: diff --git a/packages/mcode-harness/native-mcp-lifecycle.test.mjs b/packages/mcode-harness/native-mcp-lifecycle.test.mjs index e99e6fa60..092c9f3b9 100644 --- a/packages/mcode-harness/native-mcp-lifecycle.test.mjs +++ b/packages/mcode-harness/native-mcp-lifecycle.test.mjs @@ -4,10 +4,11 @@ import { mkdirSync, mkdtempSync, readFileSync, rmSync, writeFileSync } from 'nod import { createRequire } from 'node:module'; import { join } from 'node:path'; import { homedir } from 'node:os'; +import { setTimeout as delay } from 'node:timers/promises'; import { pathToFileURL } from 'node:url'; const root = process.env.MCODE_SOURCE; -test('native MCP shutdown fences the old scoped transport and reconnects lazily', { skip: !root }, async t => { +test('native MCP shutdown fences the old scoped transport and reconnects lazily', { skip: !root, timeout: 10_000 }, async t => { const require = createRequire(join(root, 'package.json')); const { build } = require('esbuild'); const paths = JSON.parse(readFileSync(join(root, 'tsconfig.standalone.json'))).compilerOptions.paths; @@ -17,6 +18,8 @@ test('native MCP shutdown fences the old scoped transport and reconnects lazily' export const transports = []; let workspaceListed; export const workspaceDiscovery = new Promise(resolve => { workspaceListed = resolve; }); + let workspaceClosing; + export const workspaceClosingDiscovery = new Promise(resolve => { workspaceClosing = resolve; }); export function createTransport(config) { const previous = transports.filter(t => t.name === config.command).length; const transport = { @@ -29,6 +32,11 @@ test('native MCP shutdown fences the old scoped transport and reconnects lazily' if (message.method === 'tools/call' && message.params.name === 'hold') return; if (message.method === 'initialize' && this.name === 'connecting' && previous === 0) return; if (message.method === 'initialize' && this.name === 'workspace-timeout') return; + if (message.method === 'tools/list' && this.name === 'optional-fast') return; + if (message.method === 'initialize' && this.name === 'workspace-closing') { + workspaceClosing(); + return; + } if (message.method === 'tools/list' && this.name === 'workspace-held') { await new Promise(resolve => workspaceListed(resolve)); } @@ -60,7 +68,7 @@ test('native MCP shutdown fences the old scoped transport and reconnects lazily' stdin: { contents: `export { McpConnectionPool } from '@mavis/mcp/runtime/connection-pool'; export { SessionMcpServers } from ${JSON.stringify(join(root, 'packages/local-runtime-v2/src/service/mcp/runtime/session-servers.ts'))}; export { LocalMcpService } from ${JSON.stringify(join(root, 'packages/local-runtime-v2/src/service/mcp/runtime/local-mcp.service.ts'))}; - export { transports, workspaceDiscovery } from 'lifecycle-transport';`, resolveDir: root }, + export { transports, workspaceDiscovery, workspaceClosingDiscovery } from 'lifecycle-transport';`, resolveDir: root }, bundle: true, write: false, platform: 'node', format: 'esm', target: 'node22', tsconfig: join(root, 'tsconfig.standalone.json'), nodePaths: [join(root, 'node_modules')], banner: { js: `import { createRequire as lifecycleCreateRequire } from 'node:module'; const require = lifecycleCreateRequire(${JSON.stringify(join(root, 'package.json'))});` }, @@ -77,10 +85,10 @@ test('native MCP shutdown fences the old scoped transport and reconnects lazily' t.after(() => rmSync(directory, { recursive: true, force: true })); const bundle = join(directory, 'lifecycle.mjs'); writeFileSync(bundle, result.outputFiles[0].text); - const { McpConnectionPool, SessionMcpServers, LocalMcpService, transports, workspaceDiscovery } = await import(pathToFileURL(bundle).href); + const { McpConnectionPool, SessionMcpServers, LocalMcpService, transports, workspaceDiscovery, workspaceClosingDiscovery } = await import(pathToFileURL(bundle).href); const pool = new McpConnectionPool({ getResolvedServer() { throw new Error('Session overrides were lost'); } }); const sessions = new SessionMcpServers(pool); - const service = new LocalMcpService(() => directory, { connectionPool: pool }); + const service = new LocalMcpService(() => directory, { connectionPool: pool, turnStdioDiscoveryTimeoutMs: 20 }); const stdio = name => ({ name, config: { type: 'stdio', command: name, args: [] } }); await sessions.configure('root', [stdio('target'), stdio('untouched'), stdio('oac_workspace'), stdio('connecting'), { name: 'remote', config: { type: 'http', url: 'https://example.test/mcp' } }]); @@ -161,19 +169,36 @@ test('native MCP shutdown fences the old scoped transport and reconnects lazily' let prepared = false; const preparation = service.configureSessionServers('workspace', [workspace('workspace-held')]) .then(() => { prepared = true; }); + const failedPreparation = preparation.catch(() => {}); const releaseList = await workspaceDiscovery; assert.equal(prepared, false, 'preparation must wait for the actual tool inventory'); + await delay(40); + assert.equal(prepared, false); + assert.equal(transports.at(-1).closeRequested, false, 'startup must outlive the Turn discovery budget'); releaseList(); await preparation; + await failedPreparation; const connected = transports.length; const inventory = await service.listNativeToolsForTurn({ sessionId: 'workspace' }); assert.deepEqual(inventory.map(tool => tool.toolName).sort(), ['bash', 'edit', 'glob', 'grep', 'read', 'write'].map(name => 'workspace_' + name)); assert.equal(transports.length, connected, 'Turn discovery must reuse the prepared connection'); for (const command of ['workspace-empty', 'workspace-incomplete', 'workspace-failed', 'workspace-timeout']) { - await assert.rejects(service.configureSessionServers(command, [workspace(command)]), /Required workspace MCP tools/); + await assert.rejects(service.configureSessionServers(command, [workspace(command)]), /Required workspace MCP tools|Failed to connect|timed out/i); assert.deepEqual(await service.listNativeToolsForTurn({ sessionId: command }), [], 'failed preparation must remove its configuration'); } + await service.configureSessionServers('optional-fast', [stdio('optional-fast')]); + assert.deepEqual(await service.listNativeToolsForTurn({ sessionId: 'optional-fast' }), [], + 'optional Turn discovery retains its short timeout'); + const closingPreparation = service.configureSessionServers('closing', [workspace('workspace-closing')]); + const rejectedClose = assert.rejects(closingPreparation, /abort|closed|Connection|connect/i); + await workspaceClosingDiscovery; + // Service shutdown joins its mutation and clears configuration; the outer + // Executor/View owner confirms final process exit after preparation fails. + const close = service.close(); + await Promise.all([close, rejectedClose]); + assert.equal(transports.at(-1).closeRequested, true); + assert.equal(service.sessionServers.get('closing'), undefined); } finally { if (previousPolicy === undefined) delete process.env.OAC_RUNTIME_MCODE_TOOL_POLICY; else process.env.OAC_RUNTIME_MCODE_TOOL_POLICY = previousPolicy; diff --git a/packages/mcode-harness/patch-native.mjs b/packages/mcode-harness/patch-native.mjs index 6db7ca1ea..35bc86b67 100644 --- a/packages/mcode-harness/patch-native.mjs +++ b/packages/mcode-harness/patch-native.mjs @@ -155,7 +155,13 @@ replaceNative(mcpServiceFile, const workspace = this.sessionServers.get(sessionId)?.oac_workspace; if (process.env.OAC_RUNTIME_MCODE_TOOL_POLICY !== 'protected-mcp-v1' || !workspace) return; try { - const tools = await this.listToolsForTurn('oac_workspace', workspace, { sessionId }); + const pool = this.options.connectionPool; + if (!pool) throw new Error('Required workspace MCP connection pool is unavailable.'); + const overrides = await this.resolveConnectionOverrides('oac_workspace', workspace, { sessionId }); + this.assertOpen(); + const tools = await pool.listTools('oac_workspace', overrides, { + signal: this.connectionSignal(workspace), + }); const names = new Set(tools.map((tool) => tool.name)); if (names.size !== requiredWorkspaceTools.length || !requiredWorkspaceTools.every((name) => names.has(name)))