Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
86 changes: 86 additions & 0 deletions apps/daemon/internal/agent/mcode/preparation_log_test.go
Original file line number Diff line number Diff line change
@@ -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")
}
}
})
}
}
21 changes: 20 additions & 1 deletion apps/daemon/internal/agent/mcode/session.go
Original file line number Diff line number Diff line change
Expand Up @@ -3,6 +3,7 @@ package mcode
import (
"context"
"encoding/json"
"errors"
"fmt"
"net/url"
"slices"
Expand All @@ -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 {
Expand Down Expand Up @@ -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 {
Expand Down Expand Up @@ -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
Expand All @@ -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.
Expand All @@ -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
}
Expand Down
4 changes: 4 additions & 0 deletions apps/daemon/internal/agent/mcode/session_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -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":
Expand Down
4 changes: 2 additions & 2 deletions packages/mcode-harness/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -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.

Expand Down Expand Up @@ -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:

Expand Down
35 changes: 30 additions & 5 deletions packages/mcode-harness/native-mcp-lifecycle.test.mjs
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand All @@ -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 = {
Expand All @@ -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));
}
Expand Down Expand Up @@ -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'))});` },
Expand All @@ -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' } }]);
Expand Down Expand Up @@ -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;
Expand Down
8 changes: 7 additions & 1 deletion packages/mcode-harness/patch-native.mjs
Original file line number Diff line number Diff line change
Expand Up @@ -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)))
Expand Down
Loading