From f9d0fd457d9e6ce86ff4f2318d36e4169625bfd0 Mon Sep 17 00:00:00 2001 From: SaladDay <1203511142@qq.com> Date: Wed, 7 Oct 2026 14:35:42 +0000 Subject: [PATCH] Delete the mcode opt-in and dead helpers --- .../agent/codex/executor_native_test.go | 4 +- .../daemon/internal/agent/installroot/lock.go | 69 ------------------- .../internal/agent/installroot/lock_test.go | 39 ----------- .../internal/agent/installroot/probe.go | 1 + .../internal/agent/mcode/declaration.go | 26 ++++--- .../internal/agent/mcode/declaration_test.go | 21 +++--- .../agent/mcode/discovery_workspace.go | 4 -- apps/daemon/internal/agent/mcode/execution.go | 7 -- .../internal/agent/mcode/installation.go | 2 +- apps/daemon/internal/cli/connect.go | 2 +- .../cli/connect_environment_binding.go | 11 --- apps/daemon/internal/daemonize/logfile.go | 25 +------ .../daemon/internal/daemonize/logfile_test.go | 15 ---- apps/daemon/internal/daemonize/pidfile.go | 11 --- .../daemon/internal/daemonize/pidfile_test.go | 8 ++- .../internal/dispatch/local_directory_test.go | 4 +- apps/daemon/internal/dispatch/preparation.go | 39 +++++------ .../internal/dispatch/preparation_start.go | 4 +- apps/daemon/internal/dispatch/router.go | 2 +- .../dispatch/runtime_preparation_test.go | 2 +- .../internal/dispatch/workspace_export.go | 2 +- .../dispatch/workspace_export_test.go | 4 +- .../workspace_preparation_status_test.go | 2 +- .../internal/dispatch/workspace_write_test.go | 2 +- .../daemon/internal/localworkspace/binding.go | 10 --- .../internal/localworkspace/binding_test.go | 3 +- .../capability_preparation_test.go | 7 +- apps/daemon/internal/paths/paths.go | 9 --- apps/daemon/internal/paths/paths_test.go | 7 +- contracts/agents-api/v1/model_execution.go | 4 -- internal/obs/log/api.go | 17 +---- internal/obs/log/api_test.go | 21 ------ internal/obs/log/background.go | 24 +------ internal/obs/log/carrier.go | 5 -- internal/obs/log/context.go | 10 --- internal/obs/log/discard.go | 13 ---- internal/obs/log/http_test.go | 16 +---- services/core/cmd/oac/init.go | 2 +- services/core/deploy/mcode/Dockerfile | 2 - services/core/deploy/mcode/README.md | 4 +- services/core/internal/api/installation.go | 47 ------------- .../core/internal/api/installation_test.go | 20 +----- .../internal/engine/configuration_test.go | 2 +- .../core/internal/runtimegateway/session.go | 15 ++-- .../internal/runtimegateway/session_test.go | 5 ++ services/core/internal/sandbox/deployment.go | 5 -- .../core/internal/sandbox/e2b/provider.go | 1 - .../internal/sandbox/microsandbox/provider.go | 1 - .../internal/sandbox/node/capacity_test.go | 4 +- .../internal/sandbox/node/docker_live_test.go | 2 +- .../core/internal/sandbox/node/identity.go | 8 --- .../core/internal/sandbox/node/node_test.go | 5 +- .../internal/sandbox/node/recovery_test.go | 10 +-- .../integration/mcode_public_native_test.go | 2 +- 54 files changed, 101 insertions(+), 486 deletions(-) delete mode 100644 apps/daemon/internal/agent/installroot/lock.go delete mode 100644 apps/daemon/internal/agent/installroot/lock_test.go delete mode 100644 internal/obs/log/discard.go diff --git a/apps/daemon/internal/agent/codex/executor_native_test.go b/apps/daemon/internal/agent/codex/executor_native_test.go index 247e6b083..f14f1340e 100644 --- a/apps/daemon/internal/agent/codex/executor_native_test.go +++ b/apps/daemon/internal/agent/codex/executor_native_test.go @@ -3,6 +3,7 @@ package codex import ( "context" "encoding/json" + "log/slog" "os" "os/exec" "strings" @@ -11,7 +12,6 @@ import ( "github.com/MiniMax-AI/OpenAgentCore/apps/daemon/internal/agent" "github.com/MiniMax-AI/OpenAgentCore/internal/agentdaemon/proto" - obslog "github.com/MiniMax-AI/OpenAgentCore/internal/obs/log" ) // This opt-in test uses the real pinned harness and an explicitly configured @@ -47,7 +47,7 @@ func TestExecutorNativeReuse(t *testing.T) { } cfg := defaultSessionConfig() cfg.codexBinary = binary - cfg.logger = obslog.Discard() + cfg.logger = slog.New(slog.DiscardHandler) req := proto.PromptRequestPayload{ AgentKind: "codex", AgentStateKey: "executor-native", DisableExecutionEnvironment: true, DisableSubagents: true, ObserveMessages: true, diff --git a/apps/daemon/internal/agent/installroot/lock.go b/apps/daemon/internal/agent/installroot/lock.go deleted file mode 100644 index 19d9295e4..000000000 --- a/apps/daemon/internal/agent/installroot/lock.go +++ /dev/null @@ -1,69 +0,0 @@ -// Package installroot coordinates adapter installations within one daemon process. -package installroot - -import ( - "context" - "os" - "sync" -) - -type installRootLock struct { - info os.FileInfo - token chan struct{} - users int -} - -var installRoots = struct { - sync.Mutex - entries map[*installRootLock]struct{} -}{entries: make(map[*installRootLock]struct{})} - -// Lock creates the install root if needed and holds its filesystem identity -// until the returned function is called once. -// Waiting callers may cancel without interrupting the current installation. -func Lock(ctx context.Context, root string) (func(), error) { - if err := ctx.Err(); err != nil { - return nil, err - } - if err := os.MkdirAll(root, 0o755); err != nil { - return nil, err - } - info, err := os.Stat(root) - if err != nil { - return nil, err - } - installRoots.Lock() - var entry *installRootLock - for candidate := range installRoots.entries { - if os.SameFile(candidate.info, info) { - entry = candidate - break - } - } - if entry == nil { - entry = &installRootLock{info: info, token: make(chan struct{}, 1)} - entry.token <- struct{}{} - installRoots.entries[entry] = struct{}{} - } - entry.users++ - installRoots.Unlock() - - release := func() { - installRoots.Lock() - entry.users-- - if entry.users == 0 { - delete(installRoots.entries, entry) - } - installRoots.Unlock() - } - select { - case <-ctx.Done(): - release() - return nil, ctx.Err() - case <-entry.token: - return func() { - entry.token <- struct{}{} - release() - }, nil - } -} diff --git a/apps/daemon/internal/agent/installroot/lock_test.go b/apps/daemon/internal/agent/installroot/lock_test.go deleted file mode 100644 index 4359af299..000000000 --- a/apps/daemon/internal/agent/installroot/lock_test.go +++ /dev/null @@ -1,39 +0,0 @@ -package installroot - -import ( - "context" - "os" - "path/filepath" - "strings" - "testing" - "time" -) - -func TestLockCaseAliases(t *testing.T) { - parent := t.TempDir() - alias := strings.ToUpper(parent) - realInfo, err := os.Stat(parent) - if err != nil { - t.Fatal(err) - } - aliasInfo, err := os.Stat(alias) - if err != nil || !os.SameFile(realInfo, aliasInfo) { - t.Skip("filesystem does not expose this case alias") - } - root := filepath.Join(parent, "new-root") - unlock, err := Lock(context.Background(), root) - if err != nil { - t.Fatal(err) - } - defer unlock() - ctx, cancel := context.WithTimeout(context.Background(), 30*time.Millisecond) - defer cancel() - release, err := Lock(ctx, filepath.Join(alias, "NEW-ROOT")) - if err == nil { - release() - t.Fatal("case alias acquired an independent lock") - } - if err != context.DeadlineExceeded { - t.Fatalf("lock error = %v", err) - } -} diff --git a/apps/daemon/internal/agent/installroot/probe.go b/apps/daemon/internal/agent/installroot/probe.go index 21f50ed47..d2932e3f8 100644 --- a/apps/daemon/internal/agent/installroot/probe.go +++ b/apps/daemon/internal/agent/installroot/probe.go @@ -1,3 +1,4 @@ +// Package installroot probes native adapter installations. package installroot import ( diff --git a/apps/daemon/internal/agent/mcode/declaration.go b/apps/daemon/internal/agent/mcode/declaration.go index bd44257ea..fa22848f1 100644 --- a/apps/daemon/internal/agent/mcode/declaration.go +++ b/apps/daemon/internal/agent/mcode/declaration.go @@ -56,20 +56,18 @@ func discoverWithCheck(parent context.Context, options agent.DiscoveryOptions, r return runtime } result.Available, result.Version = true, version - if SupportsExecution(version) { - result.Capabilities.Steering = proto.CapabilitySupported - result.Capabilities.DurableTurns = proto.CapabilitySupported - result.Capabilities.DurableInputReceipts = proto.CapabilitySupported - result.Capabilities.ExecutionControls = proto.CapabilitySupported - result.Capabilities.ProgrammaticToolCallingDisable = proto.CapabilitySupported - result.Capabilities.ToolObservations = proto.CapabilitySupported - result.Capabilities.SubagentControl = proto.CapabilitySupported - // Native preparation verifies the applied admission/tool profile before input. - result.Capabilities.SubagentObservations = proto.CapabilitySupported - result.Capabilities.EnvironmentNone = proto.CapabilitySupported - result.Capabilities.MCPHTTPTools = proto.CapabilitySupported - result.Capabilities.MCPHTTPBearerAuth = proto.CapabilitySupported - } + result.Capabilities.Steering = proto.CapabilitySupported + result.Capabilities.DurableTurns = proto.CapabilitySupported + result.Capabilities.DurableInputReceipts = proto.CapabilitySupported + result.Capabilities.ExecutionControls = proto.CapabilitySupported + result.Capabilities.ProgrammaticToolCallingDisable = proto.CapabilitySupported + result.Capabilities.ToolObservations = proto.CapabilitySupported + result.Capabilities.SubagentControl = proto.CapabilitySupported + // Native preparation verifies the applied admission/tool profile before input. + result.Capabilities.SubagentObservations = proto.CapabilitySupported + result.Capabilities.EnvironmentNone = proto.CapabilitySupported + result.Capabilities.MCPHTTPTools = proto.CapabilitySupported + result.Capabilities.MCPHTTPBearerAuth = proto.CapabilitySupported runtime.Info = result workspace := discoverWorkspace(parent, options, runtime) if runtime.Info.Available { diff --git a/apps/daemon/internal/agent/mcode/declaration_test.go b/apps/daemon/internal/agent/mcode/declaration_test.go index b60d8b348..386c4df2d 100644 --- a/apps/daemon/internal/agent/mcode/declaration_test.go +++ b/apps/daemon/internal/agent/mcode/declaration_test.go @@ -2,6 +2,7 @@ package mcode import ( "context" + "errors" "io" "reflect" "testing" @@ -10,20 +11,20 @@ import ( "github.com/MiniMax-AI/OpenAgentCore/internal/agentdaemon/proto" ) -func TestMCodeExecutionOptInIsVersionBound(t *testing.T) { +// An available runtime is execution-capable; a rejected native version is unavailable. +func TestMCodeExecutionFollowsAvailability(t *testing.T) { for _, tc := range []struct { - enabled, version string - qualified bool - }{{"", "0.4.12", false}, {"1", "0.3.11", false}, {"1", "0.4.12", true}} { - t.Run(tc.enabled+"/"+tc.version, func(t *testing.T) { - t.Setenv("OAC_RUNTIME_MCODE_AGENTS_API", tc.enabled) + version string + check error + }{{SupportedVersion, nil}, {"0.3.11", errors.New("mcode: unsupported version 0.3.11")}} { + t.Run(tc.version, func(t *testing.T) { rc := agent.DiscoveryOptions{Stdout: io.Discard, Stderr: io.Discard} - runtime := discoverWithCheck(t.Context(), rc, Declaration.Info, func(context.Context, string) (string, error) { return tc.version, nil }) - if runtime.Executor == nil || runtime.Info.Capabilities.WorkspaceReadPreparation.IsSupported() { + runtime := discoverWithCheck(t.Context(), rc, Declaration.Info, func(context.Context, string) (string, error) { return tc.version, tc.check }) + info, available := runtime.Info, tc.check == nil + if (runtime.Executor != nil) != available || info.Capabilities.WorkspaceReadPreparation.IsSupported() { t.Fatalf("factories: %+v", runtime) } - info := runtime.Info - if !info.Available || info.Capabilities.EnvironmentNone.IsSupported() != tc.qualified || info.Capabilities.DurableInputReceipts.IsSupported() != tc.qualified || info.Capabilities.SubagentObservations.IsSupported() != tc.qualified { + if info.Available != available || info.Capabilities.EnvironmentNone.IsSupported() != available || info.Capabilities.DurableInputReceipts.IsSupported() != available || info.Capabilities.SubagentObservations.IsSupported() != available { t.Fatalf("capabilities=%+v", info.Capabilities) } if info.Capabilities.NativeSessionRecovery.IsSupported() || info.Capabilities.LocalEnvironment.IsSupported() || info.Capabilities.FunctionTools.IsSupported() { diff --git a/apps/daemon/internal/agent/mcode/discovery_workspace.go b/apps/daemon/internal/agent/mcode/discovery_workspace.go index 11d7a6d55..da4f768d8 100644 --- a/apps/daemon/internal/agent/mcode/discovery_workspace.go +++ b/apps/daemon/internal/agent/mcode/discovery_workspace.go @@ -27,10 +27,6 @@ func discoverWorkspace(parent context.Context, options agent.DiscoveryOptions, r if binding == nil { return nil } - if !runtime.Info.Available || !SupportsExecution(runtime.Info.Version) { - fail(fmt.Errorf("local execution requires the qualified native version")) - return nil - } root, err := paths.Root() if err != nil { fail(err) diff --git a/apps/daemon/internal/agent/mcode/execution.go b/apps/daemon/internal/agent/mcode/execution.go index bcb9e3902..e35d927a2 100644 --- a/apps/daemon/internal/agent/mcode/execution.go +++ b/apps/daemon/internal/agent/mcode/execution.go @@ -7,13 +7,6 @@ import ( "github.com/MiniMax-AI/OpenAgentCore/internal/agentdaemon/proto" ) -// SupportsExecution reports whether the daemon advertises MiniMax Code execution: -// the operator sets OAC_RUNTIME_MCODE_AGENTS_API=1 and the native version is the -// qualified one. Otherwise discovery reports only availability. -func SupportsExecution(version string) bool { - return os.Getenv("OAC_RUNTIME_MCODE_AGENTS_API") == "1" && version == SupportedVersion -} - func validateExecutionRequest(req proto.PromptRequestPayload) error { if !req.DisableExecutionEnvironment || req.AgentStateKey == "" || req.LocalEnvironment != nil || req.RequireExistingNativeSession || len(req.FunctionTools) != 0 || (req.MCPHTTPServers != nil && len(*req.MCPHTTPServers) != 0) { return fmt.Errorf("mcode: unsupported execution configuration") diff --git a/apps/daemon/internal/agent/mcode/installation.go b/apps/daemon/internal/agent/mcode/installation.go index fe75d9258..5f6937fc1 100644 --- a/apps/daemon/internal/agent/mcode/installation.go +++ b/apps/daemon/internal/agent/mcode/installation.go @@ -12,7 +12,7 @@ import ( func Installation() agent.Installation { return agent.Installation{AgentKind: "mcode", Version: SupportedVersion, Supported: func() bool { return runtime.GOOS == "linux" || runtime.GOOS == "darwin" }, Environment: func(dir, node string) map[string]string { - return map[string]string{"OAC_RUNTIME_MCODE_BIN": filepath.Join(dir, "native", "cli.js"), "OAC_RUNTIME_MCODE_NODE": node, "OAC_RUNTIME_MCODE_WORKSPACE_BRIDGE": filepath.Join(dir, "bridge.mjs"), "OAC_RUNTIME_MCODE_AGENTS_API": "1"} + return map[string]string{"OAC_RUNTIME_MCODE_BIN": filepath.Join(dir, "native", "cli.js"), "OAC_RUNTIME_MCODE_NODE": node, "OAC_RUNTIME_MCODE_WORKSPACE_BRIDGE": filepath.Join(dir, "bridge.mjs")} }, Check: func(ctx context.Context, dir, node string, env []string) error { got, err := installroot.Probe(ctx, node, []string{filepath.Join(dir, "native", "cli.js"), "--version"}, env, dir) diff --git a/apps/daemon/internal/cli/connect.go b/apps/daemon/internal/cli/connect.go index a0d6abed1..dc87bb4a9 100644 --- a/apps/daemon/internal/cli/connect.go +++ b/apps/daemon/internal/cli/connect.go @@ -151,7 +151,7 @@ func spawnBackground(ctx context.Context, rc *runContext, profile string, argv [ } else if !errors.Is(err, os.ErrNotExist) && !errors.Is(err, daemonize.ErrStaleOrCorrupt) { return fmt.Errorf("connect: check pidfile: %w", err) } - // Stale pidfile → remove so WritePIDFile starts clean. + // Stale pidfile → remove so Spawn starts clean. _ = daemonize.RemovePIDFile(pidPath) if err := daemonize.EnsureLogFile(logPath); err != nil { diff --git a/apps/daemon/internal/cli/connect_environment_binding.go b/apps/daemon/internal/cli/connect_environment_binding.go index 9fa651798..b661092e0 100644 --- a/apps/daemon/internal/cli/connect_environment_binding.go +++ b/apps/daemon/internal/cli/connect_environment_binding.go @@ -95,17 +95,6 @@ func bindEnvironmentRuntime(remote string, bound environmentEnrollment, credenti return nil } -// Neither selecting private state nor selecting its parent grants workspace access. -func environmentPathsOverlap(first, second string) bool { - for _, pair := range [][2]string{{first, second}, {second, first}} { - relative, err := filepath.Rel(pair[0], pair[1]) - if err == nil && (relative == "." || filepath.IsLocal(relative)) { - return true - } - } - return false -} - func saveEnvironmentBinding(root string, want environmentBinding) error { // Store the binding with the other daemon state. dir := filepath.Join(root, "daemon") diff --git a/apps/daemon/internal/daemonize/logfile.go b/apps/daemon/internal/daemonize/logfile.go index 80124e1fc..a74f0dcfd 100644 --- a/apps/daemon/internal/daemonize/logfile.go +++ b/apps/daemon/internal/daemonize/logfile.go @@ -1,14 +1,14 @@ package daemonize import ( - "bufio" "errors" "fmt" - "github.com/MiniMax-AI/OpenAgentCore/internal/runtimefs" "io" "os" "path/filepath" "time" + + "github.com/MiniMax-AI/OpenAgentCore/internal/runtimefs" ) // TailOptions configures Tail. @@ -149,27 +149,6 @@ func EnsureLogFile(path string) error { return f.Close() } -// MustWriteLine appends one line to path with 0o600 mode, adding a -// trailing newline if missing. Returns errors despite the name — -// kept short because it's used in startup hot paths. -func MustWriteLine(path string, line string) error { - f, err := openPrivateLog(path) - if err != nil { - return err - } - defer f.Close() - bw := bufio.NewWriter(f) - if _, err := bw.WriteString(line); err != nil { - return err - } - if len(line) == 0 || line[len(line)-1] != '\n' { - if _, err := bw.WriteString("\n"); err != nil { - return err - } - } - return bw.Flush() -} - func openPrivateLog(path string) (*os.File, error) { root, err := os.OpenRoot(filepath.Dir(path)) if err != nil { diff --git a/apps/daemon/internal/daemonize/logfile_test.go b/apps/daemon/internal/daemonize/logfile_test.go index a13bdca48..5d1caeb6e 100644 --- a/apps/daemon/internal/daemonize/logfile_test.go +++ b/apps/daemon/internal/daemonize/logfile_test.go @@ -182,21 +182,6 @@ func TestEnsureLogFileCreatesMissingParentDir(t *testing.T) { } } -func TestMustWriteLineAppendsTrailingNewline(t *testing.T) { - dir := privateTempDir(t) - path := filepath.Join(dir, "ml.log") - if err := MustWriteLine(path, "no-newline"); err != nil { - t.Fatalf("MustWriteLine: %v", err) - } - if err := MustWriteLine(path, "has-newline\n"); err != nil { - t.Fatalf("MustWriteLine: %v", err) - } - body, _ := os.ReadFile(path) - if string(body) != "no-newline\nhas-newline\n" { - t.Errorf("body = %q", string(body)) - } -} - // safeBuf wraps bytes.Buffer with a mutex so concurrent reads/writes // during Tail's poll loop don't race. type safeBuf struct { diff --git a/apps/daemon/internal/daemonize/pidfile.go b/apps/daemon/internal/daemonize/pidfile.go index 2224c8542..b6ad008c1 100644 --- a/apps/daemon/internal/daemonize/pidfile.go +++ b/apps/daemon/internal/daemonize/pidfile.go @@ -19,15 +19,6 @@ type processIdentity struct { StopEvent string `json:"stop_event,omitempty"` } -func WritePIDFile(path string, pid int) error { - identity, err := identifyProcess(pid) - if err != nil { - return err - } - identity.StopEvent = os.Getenv(stopEventEnv) - return writeIdentity(path, identity) -} - func writeIdentity(path string, identity processIdentity) error { if path == "" { return errors.New("daemonize: process record path required") @@ -71,8 +62,6 @@ func ReadPIDFile(path string) (int, error) { return identity.PID, err } -func IsAlive(pid int) error { _, err := identifyProcess(pid); return err } - // StopPIDFile never removes ownership before the exact process has exited. // A cleanup timeout leaves the record available for observation and a later stop. func StopPIDFile(path string, timeout time.Duration) error { diff --git a/apps/daemon/internal/daemonize/pidfile_test.go b/apps/daemon/internal/daemonize/pidfile_test.go index 4aec5a769..5f83362d0 100644 --- a/apps/daemon/internal/daemonize/pidfile_test.go +++ b/apps/daemon/internal/daemonize/pidfile_test.go @@ -11,7 +11,11 @@ import ( func TestProcessRecordIdentity(t *testing.T) { path := filepath.Join(privateTempDir(t), "connect.pid") - if err := WritePIDFile(path, os.Getpid()); err != nil { + record, err := identifyProcess(os.Getpid()) + if err == nil { + err = writeIdentity(path, record) + } + if err != nil { t.Fatal(err) } pid, err := ReadPIDFile(path) @@ -34,7 +38,7 @@ func TestProcessRecordIdentity(t *testing.T) { if err = StopPIDFile(path, time.Second); !errors.Is(err, ErrStaleOrCorrupt) { t.Fatal("stale identity was accepted", err) } - if err = IsAlive(os.Getpid()); err != nil { + if _, err = identifyProcess(os.Getpid()); err != nil { t.Fatal("unrelated process was affected", err) } if _, err = os.Stat(path); err != nil { diff --git a/apps/daemon/internal/dispatch/local_directory_test.go b/apps/daemon/internal/dispatch/local_directory_test.go index b7329c2f5..8319fbecc 100644 --- a/apps/daemon/internal/dispatch/local_directory_test.go +++ b/apps/daemon/internal/dispatch/local_directory_test.go @@ -23,7 +23,7 @@ func TestLocalDirectoryPreparationNeedsNoHarnessAndRejectsOtherOwners(t *testing workspace := t.TempDir() environment, session := uuid.NewString(), uuid.NewString() - binding, err := localworkspace.New(environment, session, workspace) + binding, err := localworkspace.NewWithCapabilityDirectory(environment, session, workspace, t.TempDir()) if err != nil { t.Fatal(err) } @@ -103,7 +103,7 @@ func TestLocalDirectoryKeepsNotDirectorySeparateFromFailures(t *testing.T) { t.Fatal(err) } environment, session := uuid.NewString(), uuid.NewString() - binding, err := localworkspace.New(environment, session, workspace) + binding, err := localworkspace.NewWithCapabilityDirectory(environment, session, workspace, t.TempDir()) if err != nil { t.Fatal(err) } diff --git a/apps/daemon/internal/dispatch/preparation.go b/apps/daemon/internal/dispatch/preparation.go index 3316c282f..fdfe21631 100644 --- a/apps/daemon/internal/dispatch/preparation.go +++ b/apps/daemon/internal/dispatch/preparation.go @@ -18,23 +18,22 @@ const preparationRecords = 64 // All mutable fields are protected by Router.mu. owns includes resources whose // cancellation is underway; a slow close cannot bypass the capacity bound. type preparationState struct { - capabilities proto.AgentKindCapabilities - executor *executorState - requestID string - trace string - fingerprint [32]byte - startFingerprint [32]byte - status proto.PreparationStatusPayload - deadline time.Time - timer *time.Timer - ctx context.Context - cancel context.CancelFunc - environmentID string - busy bool - owns bool - closeErr error - workspaceReadOnly bool - handoff *preparedHandoff + capabilities proto.AgentKindCapabilities + executor *executorState + requestID string + trace string + fingerprint [32]byte + startFingerprint [32]byte + status proto.PreparationStatusPayload + deadline time.Time + timer *time.Timer + ctx context.Context + cancel context.CancelFunc + environmentID string + busy bool + owns bool + closeErr error + handoff *preparedHandoff } func (r *Router) handleExecutionPrepare(ctx context.Context, env proto.Envelope) error { @@ -102,7 +101,7 @@ func (r *Router) handleExecutionPrepare(ctx context.Context, env proto.Envelope) return r.rejectPreparation(env, "preparation_capacity") } owner, cancel := context.WithCancel(context.WithoutCancel(ctx)) - p := &preparationState{capabilities: caps, requestID: env.ID, trace: env.Trace, fingerprint: fingerprint, ctx: owner, cancel: cancel, environmentID: req.EnvironmentID(), workspaceReadOnly: true, busy: true, owns: true, deadline: time.Now().Add(r.preparationTimeout)} + p := &preparationState{capabilities: caps, requestID: env.ID, trace: env.Trace, fingerprint: fingerprint, ctx: owner, cancel: cancel, environmentID: req.EnvironmentID(), busy: true, owns: true, deadline: time.Now().Add(r.preparationTimeout)} p.status = proto.PreparationStatusPayload{Handle: uuid.NewString(), Revision: 1, State: "preparing", ExpiresAt: p.deadline.UnixMilli()} r.preparations[p.status.Handle], r.preparationRequests[p.requestID] = p, p p.timer = time.AfterFunc(r.preparationTimeout, func() { r.releasePreparation(p, "expired", "", true) }) @@ -210,7 +209,7 @@ func (r *Router) prunePreparationsLocked() { func (r *Router) publishPreparation(p *preparationState, status proto.PreparationStatusPayload) { r.mu.Lock() - if p.workspaceReadOnly && status.Revision != p.status.Revision { + if p.executor == nil && status.Revision != p.status.Revision { r.mu.Unlock() return } @@ -223,7 +222,7 @@ func (r *Router) publishPreparation(p *preparationState, status proto.Preparatio go func() { defer r.shutdownWG.Done() // A failed terminal notification must not restart incomplete cleanup. - if !r.sendPreparation(p.requestID, p.trace, status) && (!p.workspaceReadOnly || status.State == "preparing" || status.State == "ready") { + if !r.sendPreparation(p.requestID, p.trace, status) && (p.executor != nil || status.State == "preparing" || status.State == "ready") { r.releasePreparation(p, "failed", "status_delivery_failed", false) } }() diff --git a/apps/daemon/internal/dispatch/preparation_start.go b/apps/daemon/internal/dispatch/preparation_start.go index 8b0818b4c..887f42f11 100644 --- a/apps/daemon/internal/dispatch/preparation_start.go +++ b/apps/daemon/internal/dispatch/preparation_start.go @@ -33,12 +33,12 @@ func (r *Router) handleExecutionStart(_ context.Context, env proto.Envelope) err r.mu.Unlock() return r.rejectPreparation(env, "unknown_preparation") } - if p.workspaceReadOnly { + if p.executor == nil { r.mu.Unlock() return r.rejectPreparation(env, "read_only_preparation") } owner := p.executor - if owner == nil || owner.id != input.ExecutorID { + if owner.id != input.ExecutorID { r.mu.Unlock() return r.rejectPreparation(env, "unknown_executor") } diff --git a/apps/daemon/internal/dispatch/router.go b/apps/daemon/internal/dispatch/router.go index 5df1f4b32..272f65f6b 100644 --- a/apps/daemon/internal/dispatch/router.go +++ b/apps/daemon/internal/dispatch/router.go @@ -187,7 +187,7 @@ func adoptEnvelopeTrace(ctx context.Context, env proto.Envelope) context.Context return obslog.WithTrace(ctx, carrier) } } - ctx, _ = obslog.StartBackgroundTrace(ctx, "daemon.envelope") + ctx, _ = obslog.StartBackgroundTrace(ctx) return ctx } diff --git a/apps/daemon/internal/dispatch/runtime_preparation_test.go b/apps/daemon/internal/dispatch/runtime_preparation_test.go index d8a3399f9..7a8a31c70 100644 --- a/apps/daemon/internal/dispatch/runtime_preparation_test.go +++ b/apps/daemon/internal/dispatch/runtime_preparation_test.go @@ -33,7 +33,7 @@ func (s *capabilitiesTestSender) Send(ctx context.Context, env proto.Envelope) e func capabilitiesTestRouter(t *testing.T) (*Router, *capabilitiesTestSender, string, string) { t.Helper() environment, session := uuid.NewString(), uuid.NewString() - binding, err := localworkspace.New(environment, session, t.TempDir()) + binding, err := localworkspace.NewWithCapabilityDirectory(environment, session, t.TempDir(), t.TempDir()) if err != nil { t.Fatal(err) } diff --git a/apps/daemon/internal/dispatch/workspace_export.go b/apps/daemon/internal/dispatch/workspace_export.go index 86fd2cb43..490cbd463 100644 --- a/apps/daemon/internal/dispatch/workspace_export.go +++ b/apps/daemon/internal/dispatch/workspace_export.go @@ -48,7 +48,7 @@ func (r *Router) handleWorkspaceExport(ctx context.Context, env proto.Envelope) } _, code := r.workspaceResourceLocked(proto.WorkspaceReadPayload{Handle: request.Handle, EnvironmentID: request.EnvironmentID}) p := r.preparations[request.Handle] - if u != nil || r.workspaceWrite != nil || !r.localWorkspace.CanExport() || code != "" || p == nil || !p.workspaceReadOnly { + if u != nil || r.workspaceWrite != nil || !r.localWorkspace.CanExport() || code != "" || p == nil || p.executor != nil { r.mu.Unlock() return r.sendWorkspaceExport(ctx, env.ID, proto.WorkspaceExportResultPayload{Outcome: "rejected", ErrorCode: "resource_unavailable"}) } diff --git a/apps/daemon/internal/dispatch/workspace_export_test.go b/apps/daemon/internal/dispatch/workspace_export_test.go index 90dba1e5e..dfbec11e8 100644 --- a/apps/daemon/internal/dispatch/workspace_export_test.go +++ b/apps/daemon/internal/dispatch/workspace_export_test.go @@ -47,7 +47,7 @@ func exporterRouter(t *testing.T, program string) (*Router, exportSender, proto. f.Close() } environment, session := uuid.NewString(), uuid.NewString() - binding, err := localworkspace.New(environment, session, workspace) + binding, err := localworkspace.NewWithCapabilityDirectory(environment, session, workspace, t.TempDir()) if err != nil { t.Fatal(err) } @@ -57,7 +57,7 @@ func exporterRouter(t *testing.T, program string) (*Router, exportSender, proto. t.Fatal(err) } handle := uuid.NewString() - r.preparations[handle] = &preparationState{workspaceReadOnly: true, environmentID: environment, owns: true, ctx: context.Background(), deadline: time.Now().Add(time.Hour), status: proto.PreparationStatusPayload{State: "ready"}} + r.preparations[handle] = &preparationState{environmentID: environment, owns: true, ctx: context.Background(), deadline: time.Now().Add(time.Hour), status: proto.PreparationStatusPayload{State: "ready"}} t.Cleanup(func() { r.mu.Lock() delete(r.preparations, handle) diff --git a/apps/daemon/internal/dispatch/workspace_preparation_status_test.go b/apps/daemon/internal/dispatch/workspace_preparation_status_test.go index 3128a407b..d748b5bbe 100644 --- a/apps/daemon/internal/dispatch/workspace_preparation_status_test.go +++ b/apps/daemon/internal/dispatch/workspace_preparation_status_test.go @@ -22,7 +22,7 @@ func TestReadPreparationRetryCannotPublishStaleStatus(t *testing.T) { timer := time.NewTimer(time.Hour) defer timer.Stop() r := &Router{sender: sender, shutdownCh: make(chan struct{})} - p := &preparationState{workspaceReadOnly: true, owns: true, ctx: ctx, cancel: cancel, timer: timer, + p := &preparationState{owns: true, ctx: ctx, cancel: cancel, timer: timer, status: proto.PreparationStatusPayload{Handle: "reader", Revision: 2, State: "ready"}} // A prepare retry captures this snapshot before the release settles. snapshot := p.status diff --git a/apps/daemon/internal/dispatch/workspace_write_test.go b/apps/daemon/internal/dispatch/workspace_write_test.go index e6e710fe9..7ecab2792 100644 --- a/apps/daemon/internal/dispatch/workspace_write_test.go +++ b/apps/daemon/internal/dispatch/workspace_write_test.go @@ -20,7 +20,7 @@ func localWriterRouter(t *testing.T) (*dispatch.Router, *recSender, proto.Worksp t.Helper() workspace := t.TempDir() environment, session := uuid.NewString(), uuid.NewString() - binding, err := localworkspace.New(environment, session, workspace) + binding, err := localworkspace.NewWithCapabilityDirectory(environment, session, workspace, t.TempDir()) if err != nil { t.Fatal(err) } diff --git a/apps/daemon/internal/localworkspace/binding.go b/apps/daemon/internal/localworkspace/binding.go index 5ee7fbf31..0ead8710a 100644 --- a/apps/daemon/internal/localworkspace/binding.go +++ b/apps/daemon/internal/localworkspace/binding.go @@ -3,11 +3,9 @@ package localworkspace import ( "errors" "os" - "path/filepath" "strings" "sync" - "github.com/MiniMax-AI/OpenAgentCore/apps/daemon/internal/paths" "github.com/MiniMax-AI/OpenAgentCore/internal/agentcapabilities" "github.com/MiniMax-AI/OpenAgentCore/internal/agentdaemon/proto" @@ -26,14 +24,6 @@ type Binding struct { capabilityRoot string } -func New(environment, session, workspace string) (*Binding, error) { - root, err := paths.Root() - if err != nil { - return nil, err - } - return newNativeBinding(environment, session, workspace, filepath.Join(root, "capabilities")) -} - // NewWithCapabilityDirectory freezes paths selected by the Runtime operator. func NewWithCapabilityDirectory(environment, session, workspace, directory string) (*Binding, error) { return newNativeBinding(environment, session, workspace, directory) diff --git a/apps/daemon/internal/localworkspace/binding_test.go b/apps/daemon/internal/localworkspace/binding_test.go index 6d5bc80a5..b23ea8c0f 100644 --- a/apps/daemon/internal/localworkspace/binding_test.go +++ b/apps/daemon/internal/localworkspace/binding_test.go @@ -20,12 +20,11 @@ func testBinding(t *testing.T) (*Binding, proto.PromptRequestPayload) { t.Setenv("OAC_RUNTIME_HOME", private) root := t.TempDir() environment, session := uuid.NewString(), uuid.NewString() - b, err := New(environment, session, root) + b, err := NewWithCapabilityDirectory(environment, session, root, t.TempDir()) if err != nil { t.Fatal(err) } b.networkAccess = "disabled" - b.capabilityRoot = t.TempDir() return b, proto.PromptRequestPayload{LocalEnvironment: &proto.LocalEnvironment{ID: environment, NetworkAccess: "disabled", WorkspaceDirectory: "/workspace", CapabilitySources: &agentcapabilities.Input{}}, AgentStateKey: "agents-api-" + session} } diff --git a/apps/daemon/internal/localworkspace/capability_preparation_test.go b/apps/daemon/internal/localworkspace/capability_preparation_test.go index 0d226d55e..dc061b085 100644 --- a/apps/daemon/internal/localworkspace/capability_preparation_test.go +++ b/apps/daemon/internal/localworkspace/capability_preparation_test.go @@ -40,12 +40,11 @@ func TestPreparationFreezesLocalContentsAcrossReconnect(t *testing.T) { t.Fatalf("first preparation: %v", err) } writeSourceSkill(t, source, "second") - reconnect, err := New(b.environment, b.capabilityIdentity().SessionID, b.workspace) + reconnect, err := NewWithCapabilityDirectory(b.environment, b.capabilityIdentity().SessionID, b.workspace, b.capabilityRoot) if err != nil { t.Fatal(err) } reconnect.networkAccess = b.networkAccess - reconnect.capabilityRoot = b.capabilityRoot again, err := reconnect.Prepare(t.Context(), configured) if err != nil || len(again.LocalEnvironment.Skills) != 1 { t.Fatalf("reconnection: %v", err) @@ -206,11 +205,11 @@ func TestPreparationFreezesToolOnlyEnvironmentAcrossReconnect(t *testing.T) { if err := os.WriteFile(source, []byte(`{"LOCAL_ONLY":"changed"}`), 0600); err != nil { t.Fatal(err) } - reconnect, err := New(b.environment, b.capabilityIdentity().SessionID, b.workspace) + reconnect, err := NewWithCapabilityDirectory(b.environment, b.capabilityIdentity().SessionID, b.workspace, b.capabilityRoot) if err != nil { t.Fatal(err) } - reconnect.networkAccess, reconnect.capabilityRoot = b.networkAccess, b.capabilityRoot + reconnect.networkAccess = b.networkAccess if _, err := reconnect.Prepare(t.Context(), configured); err != nil { t.Fatal(err) } diff --git a/apps/daemon/internal/paths/paths.go b/apps/daemon/internal/paths/paths.go index 6c835ae3d..accdd85c6 100644 --- a/apps/daemon/internal/paths/paths.go +++ b/apps/daemon/internal/paths/paths.go @@ -88,12 +88,3 @@ func LogFile(profile string) (string, error) { } return filepath.Join(dir, "connect.log"), nil } - -// SessionsFile returns the absolute path to sessions.json. -func SessionsFile(profile string) (string, error) { - dir, err := ProfileDir(profile) - if err != nil { - return "", err - } - return filepath.Join(dir, "sessions.json"), nil -} diff --git a/apps/daemon/internal/paths/paths_test.go b/apps/daemon/internal/paths/paths_test.go index 7f4d8fa9e..cd6fc605f 100644 --- a/apps/daemon/internal/paths/paths_test.go +++ b/apps/daemon/internal/paths/paths_test.go @@ -75,10 +75,9 @@ func TestProfileDirAndFiles(t *testing.T) { } cases := map[string]func(string) (string, error){ - "auth.json": paths.AuthFile, - "connect.pid": paths.PIDFile, - "connect.log": paths.LogFile, - "sessions.json": paths.SessionsFile, + "auth.json": paths.AuthFile, + "connect.pid": paths.PIDFile, + "connect.log": paths.LogFile, } for filename, fn := range cases { got, err := fn("test") diff --git a/contracts/agents-api/v1/model_execution.go b/contracts/agents-api/v1/model_execution.go index d3c5f22da..903255c44 100644 --- a/contracts/agents-api/v1/model_execution.go +++ b/contracts/agents-api/v1/model_execution.go @@ -83,7 +83,3 @@ func (p *ModelProviderInput) ValidateHarnessWithRegistry(harness string, registr } return nil } - -func ValidateModelProtocol(protocol, harness string) error { - return builtin.Registry().ValidateProtocol(harness, protocol) -} diff --git a/internal/obs/log/api.go b/internal/obs/log/api.go index 0b12a74d3..6c0102cda 100644 --- a/internal/obs/log/api.go +++ b/internal/obs/log/api.go @@ -1,6 +1,3 @@ -// Direct slog.{Info,Warn,Error,Debug,Default} use outside this package -// is blocked by .golangci.yml forbidigo so the ctx-first signatures -// here are how trace_id auto-injection stays enforceable. package log import ( @@ -8,24 +5,12 @@ import ( "log/slog" ) -// Info logs via slog.Default with ctx attached so ContextHandler can +// Warn logs via slog.Default with ctx attached so ContextHandler can // inject trace_id/span_id from ctx. -func Info(ctx context.Context, msg string, args ...any) { - slog.Default().InfoContext(ctx, msg, args...) -} - func Warn(ctx context.Context, msg string, args ...any) { slog.Default().WarnContext(ctx, msg, args...) } -func Error(ctx context.Context, msg string, args ...any) { - slog.Default().ErrorContext(ctx, msg, args...) -} - -func Debug(ctx context.Context, msg string, args ...any) { - slog.Default().DebugContext(ctx, msg, args...) -} - // Bg returns slog.Default for ctx-less startup/init/shutdown sites. // Using Bg() in any handler-path code is a bug — it bypasses trace // attribution silently. diff --git a/internal/obs/log/api_test.go b/internal/obs/log/api_test.go index 3be844741..5c7c91e3d 100644 --- a/internal/obs/log/api_test.go +++ b/internal/obs/log/api_test.go @@ -3,32 +3,11 @@ package log import ( "bytes" "context" - "encoding/json" "log/slog" "strings" "testing" ) -func TestInfoEmitsTraceID(t *testing.T) { - carrier, _ := ParseTraceparent("00-0af7651916cd43dd8448eb211c80319c-b7ad6b7169203331-01") - var buf bytes.Buffer - prev := slog.Default() - slog.SetDefault(slog.New(NewContextHandler(slog.NewJSONHandler(&buf, nil)))) - t.Cleanup(func() { slog.SetDefault(prev) }) - - Info(WithTrace(context.Background(), carrier), "hello", "k", "v") - var got map[string]any - if err := json.Unmarshal(bytes.TrimSpace(buf.Bytes()), &got); err != nil { - t.Fatalf("unmarshal: %v\nraw=%s", err, buf.String()) - } - if got[AttrTraceID] != "0af7651916cd43dd8448eb211c80319c" { - t.Fatalf("trace_id: %v", got[AttrTraceID]) - } - if got["k"] != "v" { - t.Fatalf("user attr lost: %v", got["k"]) - } -} - // TestBgOmitsTraceID: Bg().Info must not emit trace_id — that's how // "log line came from outside a request" stays distinguishable. func TestBgOmitsTraceID(t *testing.T) { diff --git a/internal/obs/log/background.go b/internal/obs/log/background.go index f999bb9a3..9e0406a95 100644 --- a/internal/obs/log/background.go +++ b/internal/obs/log/background.go @@ -7,31 +7,11 @@ import ( // StartBackgroundTrace returns a child ctx carrying a fresh Carrier // so any log under it picks up a stable trace_id. Use at the top of // every non-HTTP logical-request entrypoint (sweeper tick, WS envelope -// handler, CLI command). For a sub-step that should keep the trace, -// use ChildSpan instead. -func StartBackgroundTrace(parent context.Context, op string) (context.Context, Carrier) { +// handler, CLI command). +func StartBackgroundTrace(parent context.Context) (context.Context, Carrier) { if parent == nil { parent = context.Background() } c := NewCarrier() - ctx := WithTrace(parent, c) - if op != "" { - _ = op - } - return ctx, c -} - -// ChildSpan returns a ctx with the parent's trace_id and a fresh span_id. -// If the parent has no carrier, behaves like StartBackgroundTrace so -// callers don't have to branch on presence. -func ChildSpan(parent context.Context) (context.Context, Carrier) { - if parent == nil { - parent = context.Background() - } - if existing, ok := TraceFromContext(parent); ok { - child := existing.ChildSpan() - return WithTrace(parent, child), child - } - c := NewCarrier() return WithTrace(parent, c), c } diff --git a/internal/obs/log/carrier.go b/internal/obs/log/carrier.go index dc6702632..1b4185163 100644 --- a/internal/obs/log/carrier.go +++ b/internal/obs/log/carrier.go @@ -30,11 +30,6 @@ func NewCarrier() Carrier { } } -// ChildSpan returns a Carrier sharing this Trace but with a fresh Span. -func (c Carrier) ChildSpan() Carrier { - return Carrier{Trace: c.Trace, Span: NewSpanID(), Sampled: c.Sampled} -} - // String formats the Carrier as a W3C traceparent. Returns "" for a // zero-trace Carrier so callers can use it as a presence check. func (c Carrier) String() string { diff --git a/internal/obs/log/context.go b/internal/obs/log/context.go index 9b2c4fcdd..0fbdfd242 100644 --- a/internal/obs/log/context.go +++ b/internal/obs/log/context.go @@ -66,10 +66,6 @@ type ctxLogger struct { ctx context.Context } -func (l ctxLogger) Debug(msg string, args ...any) { - slog.Default().DebugContext(l.ctx, msg, args...) -} - func (l ctxLogger) Info(msg string, args ...any) { slog.Default().InfoContext(l.ctx, msg, args...) } @@ -81,9 +77,3 @@ func (l ctxLogger) Warn(msg string, args ...any) { func (l ctxLogger) Error(msg string, args ...any) { slog.Default().ErrorContext(l.ctx, msg, args...) } - -// With returns a slog.Logger with the supplied attrs bound. Ctx-derived -// trace attrs still apply on subsequent InfoContext calls. -func (l ctxLogger) With(args ...any) *slog.Logger { - return slog.Default().With(args...) -} diff --git a/internal/obs/log/discard.go b/internal/obs/log/discard.go deleted file mode 100644 index 97fef6a5b..000000000 --- a/internal/obs/log/discard.go +++ /dev/null @@ -1,13 +0,0 @@ -package log - -import ( - "io" - "log/slog" -) - -// Discard returns a *slog.Logger that drops every record. Wrapped in -// ContextHandler so tests behave identically to production. -func Discard() *slog.Logger { - inner := slog.NewTextHandler(io.Discard, &slog.HandlerOptions{Level: slog.LevelError + 1}) - return slog.New(NewContextHandler(inner)) -} diff --git a/internal/obs/log/http_test.go b/internal/obs/log/http_test.go index c8930f03b..ca4c998b1 100644 --- a/internal/obs/log/http_test.go +++ b/internal/obs/log/http_test.go @@ -84,7 +84,7 @@ func TestHTTPMiddlewareIgnoresMalformedHeader(t *testing.T) { } func TestStartBackgroundTraceMintsFresh(t *testing.T) { - ctx, c := StartBackgroundTrace(context.Background(), "test.op") + ctx, c := StartBackgroundTrace(context.Background()) if c.Trace.IsZero() { t.Fatalf("StartBackgroundTrace returned zero carrier") } @@ -94,20 +94,6 @@ func TestStartBackgroundTraceMintsFresh(t *testing.T) { } } -// TestChildSpanKeepsTraceRotatesSpan: child shares trace_id, has a -// different span_id. -func TestChildSpanKeepsTraceRotatesSpan(t *testing.T) { - parentCtx, parent := StartBackgroundTrace(context.Background(), "") - _, child := ChildSpan(parentCtx) - if child.Trace != parent.Trace { - t.Fatalf("child should keep trace; parent=%s child=%s", - parent.Trace.String(), child.Trace.String()) - } - if child.Span == parent.Span { - t.Fatalf("child should rotate span; both=%s", child.Span.String()) - } -} - func TestCtxLoggerEmits(t *testing.T) { carrier, _ := ParseTraceparent("00-0af7651916cd43dd8448eb211c80319c-b7ad6b7169203331-01") var buf bytes.Buffer diff --git a/services/core/cmd/oac/init.go b/services/core/cmd/oac/init.go index 432d9da60..9128bc68b 100644 --- a/services/core/cmd/oac/init.go +++ b/services/core/cmd/oac/init.go @@ -138,7 +138,7 @@ type installReceipt struct { } func initialize(root string, release releaseIdentity, fetch func() (map[string][]byte, error)) (err error) { - ctx, _ := log.StartBackgroundTrace(context.Background(), "installation.init") + ctx, _ := log.StartBackgroundTrace(context.Background()) logger := log.With("component", "oac-init", "revision", release.revision) started := time.Now() step, stepStarted := "prepare_directories", started diff --git a/services/core/deploy/mcode/Dockerfile b/services/core/deploy/mcode/Dockerfile index df7c80c9a..774c1be79 100644 --- a/services/core/deploy/mcode/Dockerfile +++ b/services/core/deploy/mcode/Dockerfile @@ -11,8 +11,6 @@ ENV HOME=/home/runtime OAC_RUNTIME_HOME=/home/runtime/.oac \ OAC_RUNTIME_MCODE_NODE=/usr/local/bin/node \ OAC_RUNTIME_MCODE_BIN=/opt/mcode-harness/native/cli.js \ OAC_RUNTIME_MCODE_WORKSPACE_BRIDGE=/opt/mcode-harness/bridge.mjs \ - OAC_RUNTIME_MCODE_AGENTS_API=1 \ - OAC_RUNTIME_MCODE_WORKSPACE=managed \ OAC_RUNTIME_WORKSPACE=/environment/workspace \ OAC_RUNTIME_INITIALIZATION_DIRECTORY=/environment/initialization \ OAC_RUNTIME_PACKAGE_DIRECTORY=/environment/packages diff --git a/services/core/deploy/mcode/README.md b/services/core/deploy/mcode/README.md index eed55d773..58cbbcea0 100644 --- a/services/core/deploy/mcode/README.md +++ b/services/core/deploy/mcode/README.md @@ -6,7 +6,7 @@ Native tools run with the daemon user's permissions; the outer sandbox provides ## Native pin and readiness -The adapter accepts only `@minimax-ai/code` `0.4.12` ([`version.go`](../../../../apps/daemon/internal/agent/mcode/version.go)), built from the upstream source revision in [`source.json`](../../../../packages/mcode-harness/source.json) and run with Node.js 22. The daemon advertises MiniMax execution only when `OAC_RUNTIME_MCODE_AGENTS_API=1` and the version matches ([`execution.go`](../../../../apps/daemon/internal/agent/mcode/execution.go)); the native installer and the Runtime image set it. MiniMax Code runs on Linux and macOS. +The adapter accepts only `@minimax-ai/code` `0.4.12` ([`version.go`](../../../../apps/daemon/internal/agent/mcode/version.go)), built from the upstream source revision in [`source.json`](../../../../packages/mcode-harness/source.json) and run with Node.js 22. MiniMax Code runs on Linux and macOS. Workspace execution also requires the companion's readiness report: private protocol 2, the pinned native version and the pinned source revision ([`workspace_readiness.go`](../../../../apps/daemon/internal/agent/mcode/workspace_readiness.go)). An older companion is rejected even when the upstream version matches. @@ -47,7 +47,7 @@ In the workspace profile, the frozen installation's Skills are linked into the S | Base | Digest-pinned `node:22.23.1-bookworm-slim` with `ca-certificates`, `bash`, `git`, `python3`, `python3-pip` and `ripgrep` | | Programs | `/usr/local/bin/oac-daemon` and the companion at `/opt/mcode-harness` (native CLI at `native/cli.js`, bridge at `bridge.mjs`) | | User | UID/GID 1000 with `HOME=/home/runtime` | -| Environment | `OAC_RUNTIME_HOME=/home/runtime/.oac`, `OAC_RUNTIME_MCODE_NODE`, `OAC_RUNTIME_MCODE_BIN`, `OAC_RUNTIME_MCODE_WORKSPACE_BRIDGE`, `OAC_RUNTIME_MCODE_AGENTS_API=1`, `OAC_RUNTIME_WORKSPACE=/environment/workspace`, `OAC_RUNTIME_INITIALIZATION_DIRECTORY=/environment/initialization`, `OAC_RUNTIME_PACKAGE_DIRECTORY=/environment/packages` | +| Environment | `OAC_RUNTIME_HOME=/home/runtime/.oac`, `OAC_RUNTIME_MCODE_NODE`, `OAC_RUNTIME_MCODE_BIN`, `OAC_RUNTIME_MCODE_WORKSPACE_BRIDGE`, `OAC_RUNTIME_WORKSPACE=/environment/workspace`, `OAC_RUNTIME_INITIALIZATION_DIRECTORY=/environment/initialization`, `OAC_RUNTIME_PACKAGE_DIRECTORY=/environment/packages` | | Entry point | `oac-daemon connect --profile default`, working directory `/environment/workspace` | The build runs the companion's `check.mjs` and the native `--version`. The combined Runtime image uses this image as its base. Sandboxes run it with the [Docker sandbox settings](../../../../docs/sandbox-provider.md#docker-adapter). diff --git a/services/core/internal/api/installation.go b/services/core/internal/api/installation.go index 081a3d7ee..229da8835 100644 --- a/services/core/internal/api/installation.go +++ b/services/core/internal/api/installation.go @@ -1,15 +1,8 @@ package api import ( - "bytes" "context" - "encoding/json" - "errors" - "io" "net/http" - "path/filepath" - "regexp" - "slices" "time" "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/deployment" @@ -62,46 +55,6 @@ type InstallationSetting struct { Restarts []string `json:"restarts"` } -var installationSettingKey = regexp.MustCompile(`^[a-z][a-z0-9_]*(\.[a-z][a-z0-9_]*)*$`) - -const maxInstallationSettings = 64 << 10 - -// ParseInstallationConfiguration validates the installer's settings snapshot. -// A sensitive setting that carries a value is rejected, so the snapshot cannot -// leak a secret through this read. -func ParseInstallationConfiguration(raw []byte) (*InstallationConfiguration, error) { - invalid := errors.New("installation configuration is invalid") - if len(raw) > maxInstallationSettings { - return nil, invalid - } - var value InstallationConfiguration - decoder := json.NewDecoder(bytes.NewReader(raw)) - decoder.DisallowUnknownFields() - if decoder.Decode(&value) != nil || decoder.Decode(new(any)) != io.EOF { - return nil, invalid - } - if value.Path != "" && !filepath.IsAbs(value.Path) || len(value.ApplyCommand) > 4096 || value.Settings == nil { - return nil, invalid - } - if value.Path != "" && (value.ApplyCommand == "" || value.AppliedAt == nil || value.AppliedAt.IsZero()) { - return nil, invalid - } - seen := make(map[string]bool, len(value.Settings)) - for _, setting := range value.Settings { - if !installationSettingKey.MatchString(setting.Key) || seen[setting.Key] || setting.Restarts == nil || - setting.Sensitive != (setting.Configured != nil) || (setting.Sensitive && (setting.Value != nil || setting.Default != nil)) { - return nil, invalid - } - for _, service := range setting.Restarts { - if !slices.Contains([]string{"core", "web", "database"}, service) { - return nil, invalid - } - } - seen[setting.Key] = true - } - return &value, nil -} - // InstallationBindings counts what is bound to the current public URL. type InstallationBindings interface { AddressBindings(context.Context) (deployment.AddressBindings, error) diff --git a/services/core/internal/api/installation_test.go b/services/core/internal/api/installation_test.go index 94cd2bf7c..721703bf8 100644 --- a/services/core/internal/api/installation_test.go +++ b/services/core/internal/api/installation_test.go @@ -19,12 +19,7 @@ func TestInstallationReadNeedsOnlyTheCoreKey(t *testing.T) { fakes.projectsReader.resolveAPIKey = projectKeys(t, callerBinding()).ResolveAPIKey deps.CoreKeys = coreKeys(t, "administrator") public, id := "https://core.example", "5b7c0f3e-0000-4000-8000-000000000001" - settings, err := ParseInstallationConfiguration([]byte(`{"path":"/home/alice/.oac/core/config.json","apply_command":"/home/alice/.oac/core/oac apply", - "applied_at":"2026-09-25T09:30:00Z","settings":[{"key":"ports.core","value":8091,"default":8091,"changeable":true,"sensitive":false,"restarts":["core"]}, - {"key":"core.runtime_history.headers","value":null,"default":null,"configured":true,"changeable":true,"sensitive":true,"restarts":["core"]}]}`)) - if err != nil { - t.Fatal(err) - } + settings := &InstallationConfiguration{Path: "/home/alice/.oac/core/config.json", Settings: []InstallationSetting{{Key: "ports.core", Value: 8091, Default: 8091, Changeable: true, Restarts: []string{"core"}}}} fakes.installationBindings.addressBindings = func(context.Context) (deployment.AddressBindings, error) { return deployment.AddressBindings{Nodes: 2, NodesOnOtherAddress: 1}, nil } @@ -50,19 +45,6 @@ func TestInstallationReadNeedsOnlyTheCoreKey(t *testing.T) { } } -func TestInstallationSnapshotCannotCarryASensitiveValue(t *testing.T) { - for _, setting := range []string{ - `{"key":"core.runtime_history.headers","value":{"authorization":"secret"},"default":null,"configured":true,"changeable":true,"sensitive":true,"restarts":["core"]}`, - `{"key":"core.runtime_history.headers","value":null,"default":null,"changeable":true,"sensitive":true,"restarts":["core"]}`, - `{"key":"ports.core","value":8091,"default":8091,"changeable":true,"sensitive":false,"restarts":["core"],"unknown":true}`, - } { - raw := `{"path":"/c/config.json","apply_command":"/c/oac apply","applied_at":"2026-09-25T09:30:00Z","settings":[` + setting + `]}` - if _, err := ParseInstallationConfiguration([]byte(raw)); err == nil || strings.Contains(err.Error(), "secret") { - t.Fatal("accepted", setting, err) - } - } -} - func TestDeploymentAddressIsNotInput(t *testing.T) { deps, fakes := sandboxFakes(t) initializations := 0 diff --git a/services/core/internal/engine/configuration_test.go b/services/core/internal/engine/configuration_test.go index 61df5799d..74de79937 100644 --- a/services/core/internal/engine/configuration_test.go +++ b/services/core/internal/engine/configuration_test.go @@ -17,7 +17,7 @@ func TestProviderDeclarationsAgreeWithAdmission(t *testing.T) { for _, protocol := range []string{"responses", "anthropic", "unknown"} { provider, supported := declared.Provider(protocol) input := v1.ModelProviderInput{Protocol: protocol, BaseURL: "https://example.test", APIKey: "private-fixture", ContextWindow: 100, MaxOutputTokens: 20} - if (input.ValidateHarness(kind) == nil) != supported || (v1.ValidateModelProtocol(protocol, kind) == nil) != supported { + if (input.ValidateHarness(kind) == nil) != supported { t.Fatalf("protocol %q disagrees with validation", protocol) } if !supported { diff --git a/services/core/internal/runtimegateway/session.go b/services/core/internal/runtimegateway/session.go index 6e7f0f707..4643526a9 100644 --- a/services/core/internal/runtimegateway/session.go +++ b/services/core/internal/runtimegateway/session.go @@ -140,16 +140,11 @@ type Session struct { closed chan struct{} } -// NewSession wires a freshly-upgraded WS connection into a Session. -// The session does NOT start its goroutines automatically — Start runs -// once the handler is ready so the session can't race with response writes. -func NewSession(conn WSConn, deviceID, workspaceID, daemonVersion string, reg *Registry, log SessionLogger) *Session { - return NewSessionWithOwner(conn, deviceID, workspaceID, daemonVersion, reg, log, nil) -} - -// NewSessionWithOwner wires a session with an optional DB-backed owner -// lease. Multi-pod deployments pass the lease returned by -// ClaimAgentDaemonDeviceOwner so heartbeats can fence stale connections. +// NewSessionWithOwner wires a freshly-upgraded WS connection into a Session +// with an optional DB-backed owner lease. Multi-pod deployments pass the lease +// returned by ClaimAgentDaemonDeviceOwner so heartbeats can fence stale +// connections. The session does NOT start its goroutines automatically — Start +// runs once the handler is ready so the session can't race with response writes. func NewSessionWithOwner(conn WSConn, deviceID, workspaceID, daemonVersion string, reg *Registry, log SessionLogger, owner *ownerLease) *Session { if log == nil { log = func(string, ...any) {} diff --git a/services/core/internal/runtimegateway/session_test.go b/services/core/internal/runtimegateway/session_test.go index 71eaca263..e32150e4b 100644 --- a/services/core/internal/runtimegateway/session_test.go +++ b/services/core/internal/runtimegateway/session_test.go @@ -14,6 +14,11 @@ import ( "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/runtimedevice" ) +// NewSession builds an unowned Session over a test connection. +func NewSession(conn WSConn, deviceID, workspaceID, daemonVersion string, reg *Registry, log SessionLogger) *Session { + return NewSessionWithOwner(conn, deviceID, workspaceID, daemonVersion, reg, log, nil) +} + // fakeConn is the WSConn implementation used by session + registry // tests. Concurrency-safe. type fakeConn struct { diff --git a/services/core/internal/sandbox/deployment.go b/services/core/internal/sandbox/deployment.go index 95b146cc0..712d9e27b 100644 --- a/services/core/internal/sandbox/deployment.go +++ b/services/core/internal/sandbox/deployment.go @@ -67,11 +67,6 @@ type RuntimeRelease struct { FirmwareSHA256 string `json:"firmware_sha256"` } -func lowerHex(v string, bytes int) bool { - x, err := hex.DecodeString(v) - return err == nil && len(x) == bytes && hex.EncodeToString(x) == v -} - func (r RuntimeRelease) Validate() error { values := reflect.ValueOf(r) for i, rule := range runtimeContract { diff --git a/services/core/internal/sandbox/e2b/provider.go b/services/core/internal/sandbox/e2b/provider.go index c368d8866..3328eb0db 100644 --- a/services/core/internal/sandbox/e2b/provider.go +++ b/services/core/internal/sandbox/e2b/provider.go @@ -141,7 +141,6 @@ func (c Config) Validate() error { } return nil } -func New(c Config) (*Provider, error) { return NewWithCaller(c, &ProcessCaller{}) } func NewWithCaller(c Config, caller Caller) (*Provider, error) { if c.Validate() != nil || caller == nil { return nil, sandbox.ErrInvalid diff --git a/services/core/internal/sandbox/microsandbox/provider.go b/services/core/internal/sandbox/microsandbox/provider.go index d6df19ec7..6ba4eeafc 100644 --- a/services/core/internal/sandbox/microsandbox/provider.go +++ b/services/core/internal/sandbox/microsandbox/provider.go @@ -18,7 +18,6 @@ type Provider struct { var _ sandbox.SandboxProvider = (*Provider)(nil) var _ sandbox.CheckpointProvider = (*Provider)(nil) -func New(c Config) (*Provider, error) { return NewWithCaller(c, &ProcessCaller{}) } func NewWithCaller(c Config, caller Caller) (*Provider, error) { if c.Validate() != nil || caller == nil { return nil, sandbox.ErrInvalid diff --git a/services/core/internal/sandbox/node/capacity_test.go b/services/core/internal/sandbox/node/capacity_test.go index 1de42f08a..c408cb82a 100644 --- a/services/core/internal/sandbox/node/capacity_test.go +++ b/services/core/internal/sandbox/node/capacity_test.go @@ -34,7 +34,7 @@ func TestCapacityRefreshUsesAuthenticatedCoreIdentity(t *testing.T) { if err != nil || refreshed.Identity.MaxActive != 2 || refreshed.Identity.MaxRetained != 8 || refreshed.Credential != stored.Credential { t.Fatal("approved capacity not refreshed", err) } - persisted, err := LoadIdentity(dir) + persisted, err := readIdentity(dir) if err != nil || persisted != refreshed { t.Fatal("approved capacity not persisted", err) } @@ -42,7 +42,7 @@ func TestCapacityRefreshUsesAuthenticatedCoreIdentity(t *testing.T) { if _, err = RefreshIdentity(t.Context(), dir); err == nil { t.Fatal("foreign identity accepted") } - after, _ := LoadIdentity(dir) + after, _ := readIdentity(dir) if after != persisted { t.Fatal("rejected response changed retained identity") } diff --git a/services/core/internal/sandbox/node/docker_live_test.go b/services/core/internal/sandbox/node/docker_live_test.go index 5862b308c..13ab16b51 100644 --- a/services/core/internal/sandbox/node/docker_live_test.go +++ b/services/core/internal/sandbox/node/docker_live_test.go @@ -140,7 +140,7 @@ func TestDockerNodeTransportLifecycle(t *testing.T) { } running = false wait(t, func() bool { return !hub.Online(id.NodeID) }) - persisted, err := LoadIdentity(dir) + persisted, err := readIdentity(dir) if err != nil || persisted.Credential != credential || persisted.Identity.NodeID != id.NodeID { t.Fatal("node restart changed identity", err) } diff --git a/services/core/internal/sandbox/node/identity.go b/services/core/internal/sandbox/node/identity.go index ec24276a7..3b556033e 100644 --- a/services/core/internal/sandbox/node/identity.go +++ b/services/core/internal/sandbox/node/identity.go @@ -140,14 +140,6 @@ func initIdentity(dir, coreURL string, identity Identity) (StoredIdentity, error } return stored, nil } -func LoadIdentity(dir string) (StoredIdentity, error) { - release, err := lockDirectory(dir) - if err != nil { - return StoredIdentity{}, err - } - defer release() - return readIdentity(dir) -} func endpoint(raw, path string) (string, error) { u, err := url.Parse(raw) diff --git a/services/core/internal/sandbox/node/node_test.go b/services/core/internal/sandbox/node/node_test.go index 977d00c89..aab23036f 100644 --- a/services/core/internal/sandbox/node/node_test.go +++ b/services/core/internal/sandbox/node/node_test.go @@ -120,7 +120,8 @@ func TestLostCreateResponseDoesNotReplayAndReconnectSerializesCleanup(t *testing } }() wait(t, func() bool { return hub.Online(id.NodeID) }) - if _, err := LoadIdentity(dir); err == nil { + if release, err := lockDirectory(dir); err == nil { + release() t.Fatal("running node did not retain lifetime identity lock") } r := reference() @@ -218,7 +219,7 @@ func TestAgentRejectsDuplicateSequenceAndRetainsEpoch(t *testing.T) { if reads != 1 { t.Fatalf("replayed operation %d", reads) } - persisted, e := LoadIdentity(dir) + persisted, e := readIdentity(dir) if e != nil || persisted.OwnerEpoch != 9 { t.Fatalf("epoch not retained: %+v %v", persisted.Identity, e) } diff --git a/services/core/internal/sandbox/node/recovery_test.go b/services/core/internal/sandbox/node/recovery_test.go index 23ba863e3..d150b6422 100644 --- a/services/core/internal/sandbox/node/recovery_test.go +++ b/services/core/internal/sandbox/node/recovery_test.go @@ -68,7 +68,7 @@ func TestCoreRestartFencesOldConnectionAndNodeRestartKeepsIdentity(t *testing.T) t.Fatal(err) } wait(t, func() bool { return !second.Online(id.NodeID) }) - persisted, err := LoadIdentity(dir) + persisted, err := readIdentity(dir) if err != nil || persisted.OwnerEpoch != 5 || persisted.Credential != credential { t.Fatal("restart identity changed") } @@ -264,7 +264,7 @@ func TestCorruptOrMismatchedIdentityNeverRotates(t *testing.T) { if _, err = InitIdentity(dir, stored.CoreURL, mismatch); err == nil { t.Fatal("adopted wrong backend") } - original, err := LoadIdentity(dir) + original, err := readIdentity(dir) if err != nil || original.Credential != stored.Credential { t.Fatal("credential rotated") } @@ -274,10 +274,4 @@ func TestCorruptOrMismatchedIdentityNeverRotates(t *testing.T) { if _, err = InitIdentity(dir, stored.CoreURL, id); err == nil { t.Fatal("corrupt identity replaced") } - if err = os.Remove(filepath.Join(dir, "identity.json")); err != nil { - t.Fatal(err) - } - if _, err = LoadIdentity(dir); err == nil { - t.Fatal("missing identity recreated by load") - } } diff --git a/services/core/tests/integration/mcode_public_native_test.go b/services/core/tests/integration/mcode_public_native_test.go index 85ed9fca6..f592fc60e 100644 --- a/services/core/tests/integration/mcode_public_native_test.go +++ b/services/core/tests/integration/mcode_public_native_test.go @@ -132,7 +132,7 @@ func startNativeEngineDaemon(t *testing.T, h *dispatchHarness, home, binary, eng } old, _ := h.registry.LookupDevice(h.device.ID) cmd := exec.Command(binary, "connect", "--profile", "execution") - cmd.Env = append(os.Environ(), "OAC_RUNTIME_HOME="+home, "OAC_RUNTIME_MCODE_AGENTS_API=1") + cmd.Env = append(os.Environ(), "OAC_RUNTIME_HOME="+home) cmd.Stdout, cmd.Stderr = log, log if err = cmd.Start(); err != nil { log.Close()