diff --git a/apps/daemon/internal/agent/claudesdk/options.go b/apps/daemon/internal/agent/claudesdk/options.go index 875197992..a74d87d7f 100644 --- a/apps/daemon/internal/agent/claudesdk/options.go +++ b/apps/daemon/internal/agent/claudesdk/options.go @@ -57,8 +57,8 @@ func prepareConfiguration(config Config, req agent.PrepareRequest) (startRequest start.Cwd = workspaceCwd(config.Workspace) return start, env, nil } - if req.LocalEnvironment != nil || req.RequireExistingNativeSession { - return startRequest{}, nil, fmt.Errorf("claudesdk: local execution and history recovery require a dedicated workspace") + if req.LocalEnvironment != nil { + return startRequest{}, nil, fmt.Errorf("claudesdk: local execution requires a dedicated workspace") } root, err := paths.Root() if err != nil { diff --git a/apps/daemon/internal/agent/mcode/execution.go b/apps/daemon/internal/agent/mcode/execution.go index c1705278c..e50440256 100644 --- a/apps/daemon/internal/agent/mcode/execution.go +++ b/apps/daemon/internal/agent/mcode/execution.go @@ -8,7 +8,7 @@ import ( ) func validateExecutionRequest(req proto.PromptRequestPayload) error { - if !req.DisableExecutionEnvironment || req.LocalEnvironment != nil || req.RequireExistingNativeSession || req.ExecutionControls == nil { + if !req.DisableExecutionEnvironment || req.LocalEnvironment != nil || req.ExecutionControls == nil { return fmt.Errorf("mcode: unsupported execution configuration") } if !req.DisableSubagents && (req.MaxConcurrentSubagents == nil || *req.MaxConcurrentSubagents < 1) { diff --git a/apps/daemon/internal/agent/mcode/execution_test.go b/apps/daemon/internal/agent/mcode/execution_test.go index e48c5cbbd..76fa40998 100644 --- a/apps/daemon/internal/agent/mcode/execution_test.go +++ b/apps/daemon/internal/agent/mcode/execution_test.go @@ -42,7 +42,6 @@ func TestExecutionRejectsUnqualifiedAuthority(t *testing.T) { for _, change := range []func(*proto.PromptRequestPayload){ func(r *proto.PromptRequestPayload) { r.DisableExecutionEnvironment = false }, func(r *proto.PromptRequestPayload) { r.DisableSubagents = false }, - func(r *proto.PromptRequestPayload) { r.RequireExistingNativeSession = true }, func(r *proto.PromptRequestPayload) { r.ExecutionControls = nil }, } { r := testRequest(t) diff --git a/apps/daemon/internal/agenthost/environment_linux_test.go b/apps/daemon/internal/agenthost/environment_linux_test.go index e0b85a4e1..fb12f1689 100644 --- a/apps/daemon/internal/agenthost/environment_linux_test.go +++ b/apps/daemon/internal/agenthost/environment_linux_test.go @@ -169,30 +169,22 @@ func TestEnvironmentOwnerServesTheSandbox(t *testing.T) { second.release(t, b, id, status.Handle) // An uncertain File mutation quarantines the owner: no later write or - // runtime_prepare on any Router opens an attachment to send one again. + // runtime_prepare, on this Router or another, sends one again. The + // Router keeps nothing of it and shuts down. h.probe(t, b, sandboxfs.OpLink) if r := second.write(t, b, "notes/uncertain.txt", []byte("uncertain")); r.Outcome != "unknown" { t.Fatalf("the interrupted write is %+v, want unknown", r) } - // The Router reports the uncertain write as it shuts down and still - // drains the owner. - second.shutdown() - if !h.drained(t, b) { - t.Fatal("the owner kept its attachment after the Router shut down") + if r := second.write(t, b, "notes/uncertain.txt", []byte("uncertain")); r.Outcome != "unknown" { + t.Fatalf("the write after an uncertain one is %+v, want unknown", r) + } + if err := second.shutdown(); err != nil || !h.drained(t, b) { + t.Fatalf("the Router's shutdown is %v, want a drained owner", err) } third := &daemon{host: h} third.route(t, reg) third.assign(t, b) - if r := third.write(t, b, "notes/uncertain.txt", []byte("uncertain")); r.Outcome != "unknown" || !h.drained(t, b) { - t.Fatalf("the write after an uncertain one is %+v, want unknown without an attachment", r) - } - // An unknown outcome fences the Router's transfers, so the next one runs - // on another. - third.shutdown() - fourth := &daemon{host: h} - fourth.route(t, reg) - fourth.assign(t, b) - if r := fourth.runtimePrepare(t, b, proto.RuntimePreparePayload{Action: "file", File: &proto.RuntimeInitialFile{Path: "/workspace/notes/file.txt"}}, []byte("file")); r.Outcome != "unknown" || !h.drained(t, b) { + if r := third.runtimePrepare(t, b, proto.RuntimePreparePayload{Action: "file", File: &proto.RuntimeInitialFile{Path: "/workspace/notes/file.txt"}}, []byte("file")); r.Outcome != "unknown" || !h.drained(t, b) { t.Fatalf("the runtime_prepare after an uncertain write is %+v, want unknown without an attachment", r) } for _, name := range []string{"uncertain.txt", "file.txt"} { diff --git a/apps/daemon/internal/cli/agent_host_linux.go b/apps/daemon/internal/cli/agent_host_linux.go index 3abc06368..05d6bfac3 100644 --- a/apps/daemon/internal/cli/agent_host_linux.go +++ b/apps/daemon/internal/cli/agent_host_linux.go @@ -90,7 +90,7 @@ func serveAgentHost(parent context.Context, rc *runContext, args []string, decla origin, err := url.Parse(*coreURL) if err != nil || (origin.Scheme != "http" && origin.Scheme != "https") || origin.Host == "" || origin.User != nil || origin.Path != "" || origin.RawQuery != "" || origin.ForceQuery || origin.Fragment != "" || origin.String() != *coreURL { - return fmt.Errorf("agent-host: --core-url %q is not an http or https origin", *coreURL) + return errors.New("agent-host: --core-url is not an http or https origin") } wsOrigin := "ws" + strings.TrimPrefix(*coreURL, "http") diff --git a/apps/daemon/internal/dispatch/assignment.go b/apps/daemon/internal/dispatch/assignment.go index 3cecade55..a2e9d7639 100644 --- a/apps/daemon/internal/dispatch/assignment.go +++ b/apps/daemon/internal/dispatch/assignment.go @@ -11,7 +11,8 @@ import ( ) // assignmentState is the Router's record of one Session's assignment. Router.mu -// protects it. A released assignment stays recorded, so its frames stay fenced. +// protects it. A released assignment stays recorded, so its frames stay fenced; +// a confirmed release drops its owner, grant and resource. type assignmentState struct { ref proto.AssignmentRef environmentID string @@ -60,7 +61,7 @@ func (r *Router) trackWorkLocked(ref proto.AssignmentRef) func() { // admittedEnvironment admits ref for the Session and returns the owner that // its assignment resolved, or nil, or the assignment rejection code. A -// Session's assignment, and so its owner, never changes on a Router. +// superseding bind replaces both once the earlier epoch's work has settled. func (r *Router) admittedEnvironment(ref proto.AssignmentRef, sessionID string) (Environment, string) { r.mu.Lock() defer r.mu.Unlock() @@ -96,7 +97,7 @@ func (r *Router) handleAssignmentBind(ctx context.Context, env proto.Envelope) e r.mu.Lock() if r.suspensions[input.EnvironmentID] != nil { r.mu.Unlock() - return ErrRouterQuiesced + return r.rejectQuiesced(ctx, env) } a := r.assignments[ref.SessionID] switch { @@ -259,6 +260,10 @@ func (r *Router) handleAssignmentRelease(ctx context.Context, env proto.Envelope if err != nil { r.log.Warn("assignment release cleanup unconfirmed", "session_id", ref.SessionID, "err", err) code = proto.CleanupUnconfirmed + } else { + r.mu.Lock() + a.environment, a.grant, a.resource = nil, nil, sandboxbootstrap.Resource{} + r.mu.Unlock() } _ = r.reply(cleanupCtx, env, proto.TypeAssignmentStatus, assignmentStatus(state, code)) }() diff --git a/apps/daemon/internal/dispatch/executor_cancel_receipt_test.go b/apps/daemon/internal/dispatch/executor_cancel_receipt_test.go index fbf320f1d..b1e582314 100644 --- a/apps/daemon/internal/dispatch/executor_cancel_receipt_test.go +++ b/apps/daemon/internal/dispatch/executor_cancel_receipt_test.go @@ -125,7 +125,7 @@ func TestExecutorCancellationReachesNativeBeforeDurableReceiptJoin(t *testing.T) sender := &receiptCancelSender{recSender: &recSender{}, entered: make(chan struct{}), release: make(chan struct{})} owner := &receiptCancelExecutor{turn: make(chan *receiptCancelTurn, 2), cancelFails: mode == "cancel_failure" || mode == "close_failure_retry", closeFailsFirst: mode == "close_failure_retry", closeEntered: make(chan struct{}), closeRelease: make(chan struct{})} reg := agent.NewRegistry() - reg.RegisterKind(proto.SupportedAgentKind{Kind: "reusable", Available: true, Capabilities: prototest.Capabilities(proto.AgentKindCapabilities{EnvironmentNone: proto.CapabilitySupported})}, prototest.ModelConfiguration()) + reg.RegisterKind(proto.SupportedAgentKind{Kind: "reusable", Available: true, Capabilities: prototest.Capabilities(proto.AgentKindCapabilities{EnvironmentNone: proto.CapabilitySupported, NativeSessionRecovery: proto.CapabilitySupported})}, prototest.ModelConfiguration()) reg.RegisterExecutor("reusable", func(context.Context, agent.PrepareRequest) (agent.Executor, error) { return owner, nil }) r, err := dispatch.New(dispatch.Config{Registry: reg, Sender: sender, IdleTimeout: time.Hour}) if err != nil { diff --git a/apps/daemon/internal/dispatch/executor_handoff_test.go b/apps/daemon/internal/dispatch/executor_handoff_test.go index f1e4034a6..cd75fc0e8 100644 --- a/apps/daemon/internal/dispatch/executor_handoff_test.go +++ b/apps/daemon/internal/dispatch/executor_handoff_test.go @@ -91,7 +91,7 @@ func TestPreparedDonePublishesAfterExecutorHandoff(t *testing.T) { } owner := &terminalHandoffExecutor{turns: make(chan *terminalHandoffTurn, 3)} registry := agent.NewRegistry() - registry.RegisterKind(proto.SupportedAgentKind{Kind: "handoff", Available: true, Capabilities: prototest.Capabilities(proto.AgentKindCapabilities{EnvironmentNone: proto.CapabilitySupported})}, prototest.ModelConfiguration()) + registry.RegisterKind(proto.SupportedAgentKind{Kind: "handoff", Available: true, Capabilities: prototest.Capabilities(proto.AgentKindCapabilities{EnvironmentNone: proto.CapabilitySupported, NativeSessionRecovery: proto.CapabilitySupported})}, prototest.ModelConfiguration()) var creates atomic.Int32 registry.RegisterExecutor("handoff", func(context.Context, agent.PrepareRequest) (agent.Executor, error) { creates.Add(1) diff --git a/apps/daemon/internal/dispatch/executor_test.go b/apps/daemon/internal/dispatch/executor_test.go index 16af4da0e..f06e97663 100644 --- a/apps/daemon/internal/dispatch/executor_test.go +++ b/apps/daemon/internal/dispatch/executor_test.go @@ -86,7 +86,7 @@ func executorRouter(t *testing.T, owner *reusableExecutor, idle time.Duration) ( t.Helper() calls := &atomic.Int32{} reg := agent.NewRegistry() - registerExecutorKind(reg, proto.SupportedAgentKind{Kind: "reusable", Available: true, Capabilities: prototest.Capabilities(proto.AgentKindCapabilities{EnvironmentNone: proto.CapabilitySupported})}, func(context.Context, agent.PrepareRequest) (agent.Executor, error) { + registerExecutorKind(reg, proto.SupportedAgentKind{Kind: "reusable", Available: true, Capabilities: prototest.Capabilities(proto.AgentKindCapabilities{EnvironmentNone: proto.CapabilitySupported, NativeSessionRecovery: proto.CapabilitySupported})}, func(context.Context, agent.PrepareRequest) (agent.Executor, error) { calls.Add(1) return owner, nil }) @@ -314,6 +314,13 @@ func TestExecutorRejectsMissingEnvironmentBeforeFactory(t *testing.T) { } } +func TestExecutorRejectsUnsupportedRecoveryBeforeFactory(t *testing.T) { + caps := prototest.Capabilities(proto.AgentKindCapabilities{EnvironmentNone: proto.CapabilitySupported}) + if !rejectsBeforeFactory(t, caps, proto.PromptRequestPayload{AgentKind: "unrecoverable", RequireExistingNativeSession: true}) { + t.Fatal("a kind without native session recovery prepared one") + } +} + func TestExecutorRejectsOutputFromAnotherTurn(t *testing.T) { e := &reusableExecutor{starts: make(chan *reusableTurn, 1)} r, s, _ := executorRouter(t, e, time.Minute) diff --git a/apps/daemon/internal/dispatch/native_file_results_test.go b/apps/daemon/internal/dispatch/native_file_results_test.go index 4f5bc685a..e8b3353dc 100644 --- a/apps/daemon/internal/dispatch/native_file_results_test.go +++ b/apps/daemon/internal/dispatch/native_file_results_test.go @@ -1,11 +1,12 @@ package dispatch import ( - "errors" + "context" + "crypto/sha256" + "encoding/hex" "io/fs" "testing" - "github.com/MiniMax-AI/OpenAgentCore/apps/daemon/internal/agent" "github.com/MiniMax-AI/OpenAgentCore/internal/agentdaemon/proto" "github.com/google/uuid" ) @@ -27,29 +28,40 @@ func TestNativeDirectoryFailureMapping(t *testing.T) { } } } -func TestLocalUploadUnknownRetainsOwner(t *testing.T) { - sender := exportSender{make(chan proto.Envelope, 8)} - r, err := New(Config{Registry: agent.NewRegistry(), Sender: sender}) - if err != nil { - t.Fatal(err) - } - got := workspaceWriteResult(WorkspaceWriteResult{}, errors.New("unconfirmed mutation"), 3) - if got.Outcome != "unknown" { - t.Fatal(got) - } - request := proto.WorkspaceWritePayload{Step: "begin", EnvironmentID: uuid.NewString(), SessionID: uuid.NewString(), Path: "file", SizeBytes: 0, SHA256: "e3b0c44298fc1c149afbf4c8996fb92427ae41e4649b934ca495991b7852b855"} - r.workspaceWrites[request.SessionID] = &workspaceUpload{envelope: proto.Envelope{ID: uuid.NewString()}, finished: true, uncertain: true} - env, _ := proto.NewEnvelope(proto.TypeWorkspaceWrite, uuid.NewString(), request) - env.Assignment.SessionID = request.SessionID - if err = r.Handle(t.Context(), env); err != nil { - t.Fatal(err) - } - reply := <-sender.replies - var result proto.WorkspaceWriteResultPayload - if err = reply.DecodePayload(&result); err != nil || result.Outcome != "rejected" || result.ErrorCode != "write_capacity" { - t.Fatal(result, err) - } - if err = r.Shutdown(t.Context()); err == nil { - t.Fatal("shutdown declared uncertain mutation settled") + +// quarantinedEnvironment is an owner that cannot observe a write's outcome. +type quarantinedEnvironment struct{ stubEnvironment } + +func (quarantinedEnvironment) WriteWorkspaceFile(context.Context, string, []byte) (WorkspaceWriteResult, error) { + return WorkspaceWriteResult{}, ErrWorkspaceWriteUncertain +} + +// An unknown write leaves its uncertainty to the Environment owner: the +// Session's next write reaches the quarantined owner, and Shutdown settles. +func TestUnknownWriteLeavesUncertaintyToTheOwner(t *testing.T) { + r, sender, environment, session := capabilitiesTestRouter(t) + r.assignments[session].environment = quarantinedEnvironment{} + digest := sha256.Sum256([]byte("abc")) + begin := proto.WorkspaceWritePayload{Step: "begin", EnvironmentID: environment, SessionID: session, Path: "file", SizeBytes: 3, SHA256: hex.EncodeToString(digest[:])} + for range 2 { + id := uuid.NewString() + var result proto.WorkspaceWriteResultPayload + for _, step := range []proto.WorkspaceWritePayload{begin, {Step: "chunk", Data: []byte("abc")}, {Step: "commit"}} { + env, err := proto.NewEnvelope(proto.TypeWorkspaceWrite, id, step) + if err != nil { + t.Fatal(err) + } + env.Assignment = capabilityRef + if err := r.Handle(t.Context(), env); err != nil { + t.Fatal(err) + } + if err := (<-sender.frames).DecodePayload(&result); err != nil { + t.Fatal(err) + } + } + if result.Outcome != "unknown" || result.ErrorCode != "write_unconfirmed" { + t.Fatalf("write = %+v, want unknown", result) + } } + shutdownCapabilitiesRouter(t, r) } diff --git a/apps/daemon/internal/dispatch/router.go b/apps/daemon/internal/dispatch/router.go index 309d552eb..31a92ee55 100644 --- a/apps/daemon/internal/dispatch/router.go +++ b/apps/daemon/internal/dispatch/router.go @@ -176,7 +176,7 @@ func (r *Router) Handle(ctx context.Context, env proto.Envelope) error { } if a := r.assignments[env.Assignment.SessionID]; a != nil && r.suspensions[a.environmentID] != nil && env.Type != proto.TypeAssignmentRelease { r.mu.Unlock() - return ErrRouterQuiesced + return r.rejectQuiesced(ctx, env) } r.mu.Unlock() diff --git a/apps/daemon/internal/dispatch/runtime_preparation.go b/apps/daemon/internal/dispatch/runtime_preparation.go index f8ca616dd..581ca8c0e 100644 --- a/apps/daemon/internal/dispatch/runtime_preparation.go +++ b/apps/daemon/internal/dispatch/runtime_preparation.go @@ -18,15 +18,14 @@ const runtimePreparationTimeout = 120 * time.Second // installation data belongs to the bound Environment and is never removed by // transfer cleanup. type runtimePreparationTransfer struct { - id uuid.UUID // the envelope's ID - envelope proto.Envelope - request proto.RuntimePreparePayload - data []byte - ready chan struct{} - cancel context.CancelFunc - finished bool - apply bool - uncertain bool + id uuid.UUID // the envelope's ID + envelope proto.Envelope + request proto.RuntimePreparePayload + data []byte + ready chan struct{} + cancel context.CancelFunc + finished bool + apply bool } func (r *Router) handleRuntimePrepare(ctx context.Context, env proto.Envelope) error { @@ -204,10 +203,8 @@ func (r *Router) runRuntimePreparationTransfer(ctx context.Context, u *runtimePr data = nil r.mu.Lock() r.transferBytes -= u.request.SizeBytes - u.uncertain = result.Outcome == "unknown" - if !u.uncertain { - delete(r.runtimePreparations, u.request.SessionID) - } + // An unknown outcome stays with the Environment owner, which quarantines it. + delete(r.runtimePreparations, u.request.SessionID) r.mu.Unlock() // The result has a separate send budget, independent of an installation timeout. _ = r.sendRuntimePrepareResult(context.WithoutCancel(ctx), u.envelope, result) diff --git a/apps/daemon/internal/dispatch/runtime_preparation_test.go b/apps/daemon/internal/dispatch/runtime_preparation_test.go index f0a027017..840f68cf1 100644 --- a/apps/daemon/internal/dispatch/runtime_preparation_test.go +++ b/apps/daemon/internal/dispatch/runtime_preparation_test.go @@ -321,7 +321,7 @@ func TestRuntimePreparationInitializationReceipts(t *testing.T) { } } -func TestRuntimePreparationResultCategoriesAndUnknownOwnership(t *testing.T) { +func TestRuntimePreparationResultCategories(t *testing.T) { for _, tc := range []struct { err error outcome, code string @@ -338,122 +338,86 @@ func TestRuntimePreparationResultCategoriesAndUnknownOwnership(t *testing.T) { t.Fatalf("unsafe result: %+v", got) } } - r, sender, environment, session := capabilitiesTestRouter(t) - ctx, cancel := context.WithCancel(context.Background()) - id := uuid.NewString() - request := capabilityBegin(environment, session, []byte("abc")) - owner := &runtimePreparationTransfer{id: uuid.MustParse(id), envelope: capabilityEnvelope(t, id, request), request: request, data: []byte("abc"), ready: make(chan struct{}), cancel: cancel, finished: true, apply: true} - close(owner.ready) - r.runtimePreparations[session], r.transferBytes = owner, request.SizeBytes - r.shutdownWG.Add(1) - go r.runRuntimePreparationTransfer(ctx, owner, func(context.Context, uuid.UUID, proto.RuntimePreparePayload, []byte) error { - return context.DeadlineExceeded - }, func() {}) - capabilitiesReceipt(t, sender, id, "unknown") - r.mu.Lock() - owned := r.runtimePreparations[session] == owner && owner.uncertain && owner.data == nil - r.mu.Unlock() - if !owned { - t.Fatal("unknown mutation released its ownership") - } - wait, stop := context.WithTimeout(context.Background(), time.Second) - defer stop() - if err := r.Shutdown(wait); err == nil { - t.Fatal("shutdown claimed uncertain mutation settled") - } } func TestSessionTransferBlocksOnlyItsSession(t *testing.T) { - for _, mode := range []string{"held", "uncertain"} { - t.Run(mode, func(t *testing.T) { - r, sender, environment, session := capabilitiesTestRouter(t) - other := proto.AssignmentRef{SessionID: uuid.NewString(), AssignmentID: "other", Epoch: 1} - bindAssignment(r, other, environment) - next := func(id string) proto.Envelope { - t.Helper() - select { - case env := <-sender.frames: - if env.ID != id { - t.Fatalf("frame %s %s, want %s", env.Type, env.ID, id) - } - return env - case <-time.After(3 * time.Second): - t.Fatal("missing frame", id) - } - return proto.Envelope{} - } - if mode == "held" { - id := uuid.NewString() - if err := r.Handle(t.Context(), capabilityEnvelope(t, id, capabilityBegin(environment, session, []byte("abc")))); err != nil { - t.Fatal(err) - } - capabilitiesReceipt(t, sender, id, "ready") - } else { - r.mu.Lock() - r.workspaceWrites[session] = &workspaceUpload{finished: true, uncertain: true} - r.mu.Unlock() - } - start := func(ref proto.AssignmentRef) string { - env, err := proto.NewEnvelope(proto.TypeExecutionStart, uuid.NewString(), proto.ExecutionStartPayload{Handle: "handle", ExecutorID: "executor", RunID: "run", Input: proto.TextInput("input")}) - if err != nil { - t.Fatal(err) - } - env.Assignment = ref - _ = r.Handle(t.Context(), env) - var status proto.PreparationStatusPayload - if err := next(env.ID).DecodePayload(&status); err != nil { - t.Fatal(err) - } - return status.ErrorCode - } - if got := start(capabilityRef); got != "resource_unavailable" { - t.Fatalf("the transferring Session started a Run: %s", got) - } - if got := start(other); got != "unknown_preparation" { - t.Fatalf("another Session's transfer blocked execution_start: %s", got) - } - // The other Session's transfers proceed while the connection's - // memory bound allows their bodies. - id := uuid.NewString() - request := capabilityBegin(environment, other.SessionID, []byte("abc")) - env := capabilityEnvelope(t, id, request) - env.Assignment = other - if err := r.Handle(t.Context(), env); err != nil { - t.Fatal(err) - } - capabilitiesReceipt(t, sender, id, "ready") - r.mu.Lock() - r.finishRuntimePreparationTransferLocked(r.runtimePreparations[other.SessionID], false) - r.mu.Unlock() - capabilitiesReceipt(t, sender, id, "rejected") - digest := sha256.Sum256([]byte("abc")) - for _, buffered := range []int{transferMemory, 0} { - write, err := proto.NewEnvelope(proto.TypeWorkspaceWrite, uuid.NewString(), proto.WorkspaceWritePayload{Step: "begin", EnvironmentID: environment, SessionID: other.SessionID, Path: "proof", SizeBytes: 3, SHA256: hex.EncodeToString(digest[:])}) - if err != nil { - t.Fatal(err) - } - write.Assignment = other - r.mu.Lock() - r.transferBytes += buffered - r.mu.Unlock() - if err := r.Handle(t.Context(), write); err != nil { - t.Fatal(err) - } - r.mu.Lock() - r.transferBytes -= buffered - r.mu.Unlock() - var result proto.WorkspaceWriteResultPayload - if err := next(write.ID).DecodePayload(&result); err != nil { - t.Fatal(err) - } - if want := map[int]string{0: "ready", transferMemory: "write_capacity"}[buffered]; result.Outcome != want && result.ErrorCode != want { - t.Fatalf("write with %d bytes buffered = %+v", buffered, result) - } + r, sender, environment, session := capabilitiesTestRouter(t) + other := proto.AssignmentRef{SessionID: uuid.NewString(), AssignmentID: "other", Epoch: 1} + bindAssignment(r, other, environment) + next := func(id string) proto.Envelope { + t.Helper() + select { + case env := <-sender.frames: + if env.ID != id { + t.Fatalf("frame %s %s, want %s", env.Type, env.ID, id) } - r.mu.Lock() - delete(r.workspaceWrites, session) - r.mu.Unlock() - shutdownCapabilitiesRouter(t, r) - }) + return env + case <-time.After(3 * time.Second): + t.Fatal("missing frame", id) + } + return proto.Envelope{} + } + held := uuid.NewString() + if err := r.Handle(t.Context(), capabilityEnvelope(t, held, capabilityBegin(environment, session, []byte("abc")))); err != nil { + t.Fatal(err) + } + capabilitiesReceipt(t, sender, held, "ready") + start := func(ref proto.AssignmentRef) string { + env, err := proto.NewEnvelope(proto.TypeExecutionStart, uuid.NewString(), proto.ExecutionStartPayload{Handle: "handle", ExecutorID: "executor", RunID: "run", Input: proto.TextInput("input")}) + if err != nil { + t.Fatal(err) + } + env.Assignment = ref + _ = r.Handle(t.Context(), env) + var status proto.PreparationStatusPayload + if err := next(env.ID).DecodePayload(&status); err != nil { + t.Fatal(err) + } + return status.ErrorCode } + if got := start(capabilityRef); got != "resource_unavailable" { + t.Fatalf("the transferring Session started a Run: %s", got) + } + if got := start(other); got != "unknown_preparation" { + t.Fatalf("another Session's transfer blocked execution_start: %s", got) + } + // The other Session's transfers proceed while the connection's + // memory bound allows their bodies. + id := uuid.NewString() + request := capabilityBegin(environment, other.SessionID, []byte("abc")) + env := capabilityEnvelope(t, id, request) + env.Assignment = other + if err := r.Handle(t.Context(), env); err != nil { + t.Fatal(err) + } + capabilitiesReceipt(t, sender, id, "ready") + r.mu.Lock() + r.finishRuntimePreparationTransferLocked(r.runtimePreparations[other.SessionID], false) + r.mu.Unlock() + capabilitiesReceipt(t, sender, id, "rejected") + digest := sha256.Sum256([]byte("abc")) + for _, buffered := range []int{transferMemory, 0} { + write, err := proto.NewEnvelope(proto.TypeWorkspaceWrite, uuid.NewString(), proto.WorkspaceWritePayload{Step: "begin", EnvironmentID: environment, SessionID: other.SessionID, Path: "proof", SizeBytes: 3, SHA256: hex.EncodeToString(digest[:])}) + if err != nil { + t.Fatal(err) + } + write.Assignment = other + r.mu.Lock() + r.transferBytes += buffered + r.mu.Unlock() + if err := r.Handle(t.Context(), write); err != nil { + t.Fatal(err) + } + r.mu.Lock() + r.transferBytes -= buffered + r.mu.Unlock() + var result proto.WorkspaceWriteResultPayload + if err := next(write.ID).DecodePayload(&result); err != nil { + t.Fatal(err) + } + if want := map[int]string{0: "ready", transferMemory: "write_capacity"}[buffered]; result.Outcome != want && result.ErrorCode != want { + t.Fatalf("write with %d bytes buffered = %+v", buffered, result) + } + } + shutdownCapabilitiesRouter(t, r) } diff --git a/apps/daemon/internal/dispatch/shutdown.go b/apps/daemon/internal/dispatch/shutdown.go index d67d291f9..7b4377b6c 100644 --- a/apps/daemon/internal/dispatch/shutdown.go +++ b/apps/daemon/internal/dispatch/shutdown.go @@ -72,16 +72,6 @@ func (r *Router) runShutdownAttempt(attempt *shutdownAttempt, victims []sessionC attempt.err = errors.Join(attempt.err, err) } r.mu.Lock() - for session, u := range r.runtimePreparations { - if u.uncertain { - attempt.err = errors.Join(attempt.err, fmt.Errorf("dispatch: capability preparation of Session %s remains uncertain", session)) - } - } - for session, u := range r.workspaceWrites { - if u.uncertain { - attempt.err = errors.Join(attempt.err, fmt.Errorf("dispatch: local workspace write of Session %s remains uncertain", session)) - } - } for _, p := range r.preparations { if p.owns { cause := p.closeErr @@ -104,20 +94,30 @@ func (r *Router) runShutdownAttempt(attempt *shutdownAttempt, victims []sessionC // closeEnvironments drains the owners of the unreleased assignments in // environmentID, or in every Environment when it is empty, once their work -// has settled. A drained owner serves its Session again on the next bind. +// has settled. It waits for a released assignment's cleanup, which closes the +// owner itself. A drained owner serves its Session again on the next bind. func (r *Router) closeEnvironments(ctx context.Context, environmentID string) error { r.mu.Lock() - var owners []*assignmentState + var released []*assignmentState + owners := map[*assignmentState]Environment{} for _, a := range r.assignments { - if !a.released && a.environment != nil && (environmentID == "" || a.environmentID == environmentID) { - owners = append(owners, a) + switch { + case environmentID != "" && a.environmentID != environmentID: + case a.released: + released = append(released, a) + case a.environment != nil: + owners[a] = a.environment } } r.mu.Unlock() + for _, a := range released { + a.cleanup.Lock() + a.cleanup.Unlock() + } var errs []error - for _, a := range owners { + for a, environment := range owners { a.cleanup.Lock() - if err := a.environment.Close(ctx); err != nil { + if err := environment.Close(ctx); err != nil { errs = append(errs, fmt.Errorf("dispatch: Environment of Session %s: %w", a.ref.SessionID, err)) } a.cleanup.Unlock() diff --git a/apps/daemon/internal/dispatch/suspend.go b/apps/daemon/internal/dispatch/suspend.go index ed9fed0e0..fd14f1327 100644 --- a/apps/daemon/internal/dispatch/suspend.go +++ b/apps/daemon/internal/dispatch/suspend.go @@ -219,6 +219,12 @@ func (r *Router) drainEnvironment(ctx context.Context, environmentID string) err return nil } +// rejectQuiesced answers a frame of a quiesced Environment, which admits only +// its Sessions' releases. +func (r *Router) rejectQuiesced(ctx context.Context, env proto.Envelope) error { + return errors.Join(ErrRouterQuiesced, r.reply(ctx, env, proto.TypeProtocolError, proto.ProtocolErrorPayload{Type: env.Type, ErrorCode: "resource_unavailable"})) +} + // suspensionCode is the error code of a refused quiesce or resume. func suspensionCode(err error) string { var rejected AssignmentError diff --git a/apps/daemon/internal/dispatch/suspend_test.go b/apps/daemon/internal/dispatch/suspend_test.go index 482c6fe93..f15b11fc0 100644 --- a/apps/daemon/internal/dispatch/suspend_test.go +++ b/apps/daemon/internal/dispatch/suspend_test.go @@ -3,6 +3,7 @@ package dispatch import ( "context" "errors" + "sync" "sync/atomic" "testing" "time" @@ -84,8 +85,8 @@ func TestQuiesceRejectsEveryUnsettledResource(t *testing.T) { } func TestQuiesceDrainsPendingReceiptAndFencesConcurrentAdmission(t *testing.T) { - entered, release := make(chan struct{}), make(chan struct{}) - r := suspensionRouter(t, suspendSender(func(context.Context, proto.Envelope) error { close(entered); <-release; return nil })) + entered, release, once := make(chan struct{}), make(chan struct{}), sync.Once{} + r := suspensionRouter(t, suspendSender(func(context.Context, proto.Envelope) error { once.Do(func() { close(entered); <-release }); return nil })) // Rejection receipts run independently of the preparation resource map. _ = r.Handle(context.Background(), proto.Envelope{Type: proto.TypeExecutionPrepare, ID: "invalid"}) <-entered @@ -203,6 +204,10 @@ func TestQuiescingOneEnvironmentLeavesAnotherRunning(t *testing.T) { if err := prepare(suspendRef); !errors.Is(err, ErrRouterQuiesced) { t.Fatalf("the quiesced Environment admitted %v", err) } + var rejected proto.ProtocolErrorPayload + if frame := <-frames; frame.Type != proto.TypeProtocolError || frame.DecodePayload(&rejected) != nil || rejected.ErrorCode != "resource_unavailable" { + t.Fatalf("quiesced rejection = %+v %+v", frame, rejected) + } if err := prepare(other); errors.Is(err, ErrRouterQuiesced) { t.Fatal("the other Environment stopped admitting work") } diff --git a/apps/daemon/internal/dispatch/workspace_write.go b/apps/daemon/internal/dispatch/workspace_write.go index 3b8ecefb6..e4705ddf0 100644 --- a/apps/daemon/internal/dispatch/workspace_write.go +++ b/apps/daemon/internal/dispatch/workspace_write.go @@ -13,13 +13,12 @@ import ( // Router.mu protects a Session's bounded transfer to its Environment owner. type workspaceUpload struct { - envelope proto.Envelope - request proto.WorkspaceWritePayload - data []byte - ready chan struct{} - finished bool - apply bool - uncertain bool + envelope proto.Envelope + request proto.WorkspaceWritePayload + data []byte + ready chan struct{} + finished bool + apply bool } func (r *Router) handleWorkspaceWrite(ctx context.Context, env proto.Envelope) error { @@ -134,10 +133,8 @@ func (r *Router) runWorkspaceUpload(ctx context.Context, u *workspaceUpload, env } r.mu.Lock() r.transferBytes -= u.request.SizeBytes - u.uncertain = result.Outcome == "unknown" - if !u.uncertain { - delete(r.workspaceWrites, u.request.SessionID) - } + // An unknown outcome stays with the Environment owner, which quarantines it. + delete(r.workspaceWrites, u.request.SessionID) r.mu.Unlock() _ = r.sendWorkspaceWrite(ctx, u.envelope, result) } diff --git a/contracts/agents-api/index.md b/contracts/agents-api/index.md index d7738cfdf..82f3bc2da 100644 --- a/contracts/agents-api/index.md +++ b/contracts/agents-api/index.md @@ -109,7 +109,7 @@ Each item is Core's deliberate or native behavior where the official service beh **Configuration and tools** - Explicit reasoning effort or summary, service tiers other than `auto`, enabled `web_search` and enabled programmatic tool calling are saved but rejected at Session admission. -- Harness support differs as each [declaration](./harness-onboarding.md#declare-support) states. Codex has no structured output or `tool_search`. Claude Code takes no whitespace-only text, `medium` verbosity only and function-result images only inline and in successful results; it rejects structured output with Subagents, MCP, installed capabilities or `tool_search`, and `tool_search` with MCP or installed capabilities. MiniMax Code has no public functions, no service-origin MCP, no image input, no whitespace-only text, no required MCP and `medium` verbosity only, and takes `allowed_tools` null only. Each Harness reserves an MCP `server_label`: `codex_apps` for Codex, `functions` for Claude Code and `oac_workspace` for MiniMax Code. Claude Code also requires labels to match `^[a-zA-Z0-9_-]+$` and `allowed_tools` names to match `^[a-zA-Z0-9_.-]+$`. +- Harness support differs as each [declaration](./harness-onboarding.md#declare-support) states. Codex has no structured output or `tool_search`. Claude Code takes no whitespace-only text, `medium` verbosity only and function-result images only inline and in successful results; it rejects Subagents with MCP, including MCP servers that Environment Plugins install; structured output with Subagents, MCP, installed capabilities or `tool_search`; and `tool_search` with MCP or installed capabilities. MiniMax Code has no public functions, no service-origin MCP, no image input, no whitespace-only text, no required MCP and `medium` verbosity only, and takes `allowed_tools` null only. Each Harness reserves an MCP `server_label`: `codex_apps` for Codex, `functions` for Claude Code and `oac_workspace` for MiniMax Code. Claude Code also requires labels to match `^[a-zA-Z0-9_-]+$` and `allowed_tools` names to match `^[a-zA-Z0-9_.-]+$`. - Model-derived reasoning defaults are not resolved. - MCP tools support the `http` transport only; `stdio` is rejected, and so is an inline `authorization` on a Session MCP transport ([HTTP MCP](./execution-tools.md#http-mcp)). diff --git a/contracts/agents-api/zh/index.md b/contracts/agents-api/zh/index.md index 2fd794cfd..ec74e95a9 100644 --- a/contracts/agents-api/zh/index.md +++ b/contracts/agents-api/zh/index.md @@ -1,7 +1,7 @@ --- title: "Agents API 覆盖台账" source: contracts/agents-api/index.md -source_hash: 3e6abc5b98cd47c332ff6f5c12fc8676a35428da02a6caae2648dc8bab6bc486 +source_hash: b02e730eee7501c7204a5aa24d9c8d8a9b613273a441c717c8b4dc4fb36f1f8b --- Core 旨在以下方固定版本为准支持完整的 OpenAI Agents API([public API rule](https://github.com/MiniMax-AI/OpenAgentCore/blob/main/AGENTS.md#public-api))。本台账记录 Core 对各项资源实现了哪些内容、哪些契约保存其详细信息,并列出相对于 OpenAI 服务的所有已知差异和所有未解决缺口。[API namespaces and credentials](../../../docs/zh/api/index.md) 说明谁调用哪些 API;[Agents API guide](../../../docs/zh/api/public-agent-api.md) 介绍使用方法。 @@ -112,7 +112,7 @@ Core 自身字段位于 `x_agents_core` 中([Core extensions](../../../docs/zh **配置和工具** - 显式指定推理强度或摘要、使用 `auto` 之外的服务层级、启用 `web_search` 或启用程序化工具调用,这些设置都会被保存,但在 Session 准入时会被拒绝。 -- 各 Harness 的支持差异以其[声明](harness-onboarding.md#declare-support)为准。Codex 不支持结构化输出或 `tool_search`。Claude Code 不接受仅含空白的文本,只支持 `medium` 详细程度,函数结果图像只能内联且只能出现在成功结果中;它拒绝结构化输出与子智能体、MCP、已安装能力或 `tool_search` 同时使用,也拒绝 `tool_search` 与 MCP 或已安装能力同时使用。MiniMax Code 不提供公共 functions,没有服务源 MCP,不支持图像输入、仅含空白的文本和必需 MCP,只支持 `medium` 详细程度,且只接受值为 null 的 `allowed_tools`。每个 Harness 都保留一个 MCP `server_label`:Codex 保留 `codex_apps`,Claude Code 保留 `functions`,MiniMax Code 保留 `oac_workspace`。Claude Code 还要求标签匹配 `^[a-zA-Z0-9_-]+$`,`allowed_tools` 中的名称匹配 `^[a-zA-Z0-9_.-]+$`。 +- 各 Harness 的支持差异以其[声明](harness-onboarding.md#declare-support)为准。Codex 不支持结构化输出或 `tool_search`。Claude Code 不接受仅含空白的文本,只支持 `medium` 详细程度,函数结果图像只能内联且只能出现在成功结果中;它拒绝子智能体与 MCP(包括 Environment Plugins 安装的 MCP server)同时使用,拒绝结构化输出与子智能体、MCP、已安装能力或 `tool_search` 同时使用,也拒绝 `tool_search` 与 MCP 或已安装能力同时使用。MiniMax Code 不提供公共 functions,没有服务源 MCP,不支持图像输入、仅含空白的文本和必需 MCP,只支持 `medium` 详细程度,且只接受值为 null 的 `allowed_tools`。每个 Harness 都保留一个 MCP `server_label`:Codex 保留 `codex_apps`,Claude Code 保留 `functions`,MiniMax Code 保留 `oac_workspace`。Claude Code 还要求标签匹配 `^[a-zA-Z0-9_-]+$`,`allowed_tools` 中的名称匹配 `^[a-zA-Z0-9_.-]+$`。 - 由模型推导出的推理默认值不会被解析确定。 - MCP 工具仅支持 `http` 传输,`stdio` 会被拒绝,Session MCP 传输中的内联 `authorization` 也会被拒绝([HTTP MCP](execution-tools.md#http-mcp))。 diff --git a/docs/runtime-protocol.md b/docs/runtime-protocol.md index 6ed8f348e..2ea7cb6ad 100644 --- a/docs/runtime-protocol.md +++ b/docs/runtime-protocol.md @@ -37,7 +37,7 @@ A heartbeat only narrows the Harness's static declaration. Core treats a heartbe | Capability | Core requires it when | | --- | --- | | `environment_none` | The Environment type is `none` | -| `local_environment` | The Environment type is `openai_hosted` or `self_hosted`, or an idle Files directory read needs a read-only preparation | +| `local_environment` | The Environment type is `openai_hosted` or `self_hosted` | | `native_session_recovery` | A Session with a started Turn has no recorded native Session ID | | `text_verbosity` | The Agent requests a verbosity other than `medium` | | `structured_output` | The Agent requests `json_schema` output | @@ -114,13 +114,13 @@ Before a Session's first operation on a connection, including Environment initia A Session's first bind resolves its Environment owner, which holds the Environment's resources and performs every effect on them; it does not change afterwards. The owner checks each `execution_prepare` configuration against the Environment, including the read-only profile, and fills the installed capabilities before an Executor starts. It applies `runtime_prepare`, lists directories for `workspace_read`, writes files for `workspace_write` and exports outputs for `workspace_export`. The Runtime's dispatcher keeps admission, transfer framing and fencing, and never substitutes another implementation. A self-hosted Runtime's owner is its bound local workspace, which outlives each assignment; the Runtime rejects the bind of any other Session with `assignment_conflict`. An agent host's owner works in the Session's sandbox through File and Process on a [Link](./sandbox-link-protocol.md) attachment of its own, opened under the bind's grant on first use. It lasts from the Session's first bind until its home is removed and outlives the Session's Executors and connections. Quiescing its Environment, releasing the assignment or superseding its epoch closes its attachment; a bind under a later assignment or epoch takes the owner over once that attachment is closed, and a bind it cannot take over, including one of another Environment, fails with `assignment_conflict`. It runs each setup step as the Process operation whose ID is the `runtime_prepare` envelope ID. A `workspace_write` or `runtime_prepare` that cannot reach the sandbox before any effect ends `rejected` with `resource_unavailable`. It fails a Plugin whose MCP server declares literal `http_headers`, or is a stdio server with `env_vars`, before staging any of it. A file mutation or setup step whose outcome it cannot observe quarantines the owner until the home is removed: it sends no further mutation, and every later `workspace_write` and `runtime_prepare` ends `unknown`. A Session without an owner supports none of these operations, and the Runtime rejects each with its typed code: `unsupported_read_preparation` for a read-only preparation, `invalid_configuration` for an Executor configuration with a `local_environment`, `runtime_preparation_unsupported` for `runtime_prepare`, `write_unsupported` for `workspace_write`, and `read_unsupported` for `workspace_read` and `workspace_export`. -Core records a release and advances the epoch before it sends anything, which withdraws the assignment's attach grant; it then has the relay revoke the Environment's Link resource at its current generation, so the attachments opened under the grant close before the Runtime receives the release. Deleting a Session releases its assignment with `remove_home: true`; releasing its Environment sends `false`. A deletion never revokes a shared Runtime credential. `assignment_release` fences the assignment at once. The Runtime then stops the Session's work: a transfer still receiving its body, or committed but not yet applied, ends with `assignment_stale`; it releases read-only preparations and waits until every workspace read, write, export and Runtime preparation has sent its result. It closes the Session's Executors, then releases what its Environment owner holds and, when asked, removes the native home; only then does it reply `released` or `home_removed`. Unfinished cleanup replies `failed` with `cleanup_unconfirmed`, and a retry at the same epoch repeats it. A Runtime that declares `home_removal` unsupported answers `remove_home: true` with `unsupported_operation`, and Core asks it only to release. Core records the release as applied from a matching `released` or `home_removed`, or at once when no Runtime is left to act on it: a release to a Runtime without authority is settled when recorded, and revoking a Runtime settles its releases. Core resends every unacknowledged release to a Runtime when it connects; a release that fails backs off, and the release due longest goes first, so failing releases cannot delay the rest. `environment_quiesce` applies to the named Environment only: it fails with `resource_busy` while one of the Environment's Sessions has work in progress; otherwise the Runtime closes the Environment's Executors, has its owners release what they hold and replies `environment_quiesced`. Until the matching `environment_resume`, which carries the assignment that quiesced it, the Runtime admits for that Environment only the releases of its Sessions; Sessions of other Environments keep running. +Core records a release and advances the epoch before it sends anything, which withdraws the assignment's attach grant; it then has the relay revoke the Environment's Link resource at its current generation, so the attachments opened under the grant close before the Runtime receives the release. Deleting a Session releases its assignment with `remove_home: true`; releasing its Environment sends `false`. A deletion never revokes a shared Runtime credential. `assignment_release` fences the assignment at once. The Runtime then stops the Session's work: a transfer still receiving its body, or committed but not yet applied, ends with `assignment_stale`; it releases read-only preparations and waits until every workspace read, write, export and Runtime preparation has sent its result. It closes the Session's Executors, then releases what its Environment owner holds and, when asked, removes the native home; only then does it reply `released` or `home_removed`. Unfinished cleanup replies `failed` with `cleanup_unconfirmed`, and a retry at the same epoch repeats it. A Runtime that declares `home_removal` unsupported answers `remove_home: true` with `unsupported_operation`, and Core asks it only to release. Core records the release as applied from a matching `released` or `home_removed`, or at once when no Runtime is left to act on it: a release to a Runtime without authority is settled when recorded, and revoking a Runtime settles its releases. Core resends every unacknowledged release to a Runtime when it connects; a release that fails backs off, and the release due longest goes first, so failing releases cannot delay the rest. `environment_quiesce` applies to the named Environment only: it fails with `resource_busy` while one of the Environment's Sessions has work in progress; otherwise the Runtime closes the Environment's Executors, has its owners release what they hold and replies `environment_quiesced`. Until the matching `environment_resume`, which carries the assignment that quiesced it, the Runtime admits for that Environment only the releases of its Sessions and answers any other of its frames, including a bind, with `protocol_error` `resource_unavailable`; Sessions of other Environments keep running. The Runtime answers a Core frame it cannot route with `protocol_error`, which echoes the request's ID and carries its type and an error code. ## Preparation and execution order -Environment initialization uses `runtime_prepare` on every connection, managed or user-owned; the [Environment contract](../contracts/agents-api/environments.md#runtime-capability-preparation) owns what is prepared and when. For a file or archive, send `begin`, wait for `ready`, send ordered chunks and await each matching `received` offset, then send `commit` and await `completed`. Initialization and finalization have typed headers without file data. Validate the expected outcome, offset, size and finite error code with the shared validator. A Session has at most one transfer of each kind at a time, and a transfer, or one whose outcome is unknown, blocks only its own Session's `execution_start` and transfers and its Environment's quiesce. The bodies staged by all of a connection's Sessions together are bounded by the sum of the `workspace_write` and `runtime_prepare` bounds; a `begin` beyond it is rejected with `write_capacity` or `runtime_preparation_capacity`. A chunk receipt confirms staged bytes, not installation; a completed commit confirms that operation, not that a later Turn ran. +Environment initialization uses `runtime_prepare` on every connection, managed or user-owned; the [Environment contract](../contracts/agents-api/environments.md#runtime-capability-preparation) owns what is prepared and when. For a file or archive, send `begin`, wait for `ready`, send ordered chunks and await each matching `received` offset, then send `commit` and await `completed`. Initialization and finalization have typed headers without file data. Validate the expected outcome, offset, size and finite error code with the shared validator. A Session has at most one transfer of each kind at a time, and a transfer blocks only its own Session's `execution_start` and transfers and its Environment's quiesce until it sends its result, whatever the outcome. The bodies staged by all of a connection's Sessions together are bounded by the sum of the `workspace_write` and `runtime_prepare` bounds; a `begin` beyond it is rejected with `write_capacity` or `runtime_preparation_capacity`. A chunk receipt confirms staged bytes, not installation; a completed commit confirms that operation, not that a later Turn ran. An execution Turn runs in five steps: @@ -202,7 +202,7 @@ Request payloads are bounded at 8 KiB and correlation IDs at 128 bytes before ad Core runs an idle directory read on the Worker's Session scheduling reservation and targets the exact Run during active execution. It keeps the reservation through the bounded read and release, returns data only after a confirmed close (an incomplete read or uncertain cleanup returns unavailable without data), releases the reservation before delivering the result, and revokes the scoped read credential on completion or failure. The Runtime keeps uncertain cleanup ownership and capacity. The [Environment Files contract](../contracts/agents-api/environment-files.md) owns public authorization, paths and pagination. -`workspace_write` transfers a complete bounded body in acknowledged 64 KiB frames before the native writer runs, verifies the declared digest and runs no model. The private transfer bound is 50 MiB, separate from the public 5 MiB decoded inline bound that the API checks before any Runtime work. The Runtime excludes the Session's execution while it receives or applies a write; a malformed, incomplete or expired transfer never reaches the installer. An exact commit or rejection receipt releases the mutation owner. A missing or ambiguous receipt keeps the uncertainty: observer cancellation and local process exit cannot prove that nothing changed. Before public admission Core durably reserves the write under the Session lock and blocks successor mutations across restarts until exact settlement; the request is never replayed. On every platform the owner lists directories, creates files and exports outputs in the daemon, with no external helper or staging directory. +`workspace_write` transfers a complete bounded body in acknowledged 64 KiB frames before the native writer runs, verifies the declared digest and runs no model. The private transfer bound is 50 MiB, separate from the public 5 MiB decoded inline bound that the API checks before any Runtime work. The Runtime excludes the Session's execution while it receives or applies a write; a malformed, incomplete or expired transfer never reaches the installer. An exact commit or rejection receipt settles the write. A missing or ambiguous receipt leaves the uncertainty with the [Environment owner](#session-assignments): observer cancellation and local process exit cannot prove that nothing changed. Before public admission Core durably reserves the write under the Session lock and blocks successor mutations across restarts until exact settlement; the request is never replayed. On every platform the owner lists directories, creates files and exports outputs in the daemon, with no external helper or staging directory. ## MCP connection authority diff --git a/docs/zh/runtime-protocol.md b/docs/zh/runtime-protocol.md index 4a030e550..9bb451de0 100644 --- a/docs/zh/runtime-protocol.md +++ b/docs/zh/runtime-protocol.md @@ -1,7 +1,7 @@ --- title: "Core–Runtime 协议" source: docs/runtime-protocol.md -source_hash: 103462f4f02379eeb587294e3124e366ede0743cc174d335f9b5355e37d84c62 +source_hash: b823b0a4d564b9191687005326939a7db10803ef99667d088336427b906d06d5 --- 此协议在 Runtime daemon 获取机器凭据后连接 Core 与 daemon,定义 daemon 连接上消息的含义和顺序。wire 类型、限制和验证器仅在 [`internal/agentdaemon/proto`](https://github.com/MiniMax-AI/OpenAgentCore/tree/main/internal/agentdaemon/proto) 中定义一次;Core 的 [gateway](https://github.com/MiniMax-AI/OpenAgentCore/tree/main/services/core/internal/runtimegateway) 与参考 Runtime 的 [dispatcher](https://github.com/MiniMax-AI/OpenAgentCore/tree/main/apps/daemon/internal/dispatch) 都使用它们,因此无需同步第二套 payload schema。签发凭据和打开连接的 HTTP 路由见[机器连接 API](../../contracts/agents-api/zh/machine-api.md)。 @@ -39,7 +39,7 @@ heartbeat 只能收窄 Harness 的静态声明。某个 kind 没有内置声明 | 能力 | Core 何时要求 | | --- | --- | | `environment_none` | Environment 类型为 `none` | -| `local_environment` | Environment 类型为 `openai_hosted` 或 `self_hosted`,或空闲 Files 目录读取需要只读 preparation | +| `local_environment` | Environment 类型为 `openai_hosted` 或 `self_hosted` | | `native_session_recovery` | Session 已启动过 Turn,但未记录原生 Session ID | | `text_verbosity` | Agent 请求 `medium` 以外的 verbosity | | `structured_output` | Agent 请求 `json_schema` 输出 | @@ -116,13 +116,13 @@ Usage frame 和最终 usage snapshot 都携带当前执行的累计测量,替 Session 的第一次绑定确定其 Environment owner,此后不再改变;owner 持有 Environment 的资源,并执行对这些资源的每个作用。owner 根据 Environment 检查每个 `execution_prepare` 配置(包括只读 profile),并在 Executor 启动前填入已安装的能力。它应用 `runtime_prepare`,为 `workspace_read` 列举目录,为 `workspace_write` 写入文件,为 `workspace_export` 导出输出。Runtime 的 dispatcher 保留准入、传输分帧和 fencing,从不替换为其他实现。self-hosted Runtime 的 owner 是其绑定的本地工作区,该工作区比每个分配存续得更久;Runtime 以 `assignment_conflict` 拒绝任何其他 Session 的绑定。agent host 的 owner 通过自己的一个 [Link](./sandbox-link-protocol.md) attachment,用 File 和 Process 在 Session 的沙箱中工作;该 attachment 在首次使用时凭绑定的 grant 打开。owner 从 Session 的第一次绑定存续到其 home 被删除,比 Session 的 Executor 和连接存续得更久。其 Environment 被 quiesce、分配被释放或其 epoch 被取代时关闭其 attachment;之后的分配或 epoch 下的绑定在该 attachment 关闭后接管 owner,它无法接管的绑定(包括其他 Environment 的绑定)以 `assignment_conflict` 失败。它把每个 setup 步骤作为 Process 操作运行,操作 ID 即 `runtime_prepare` 的 envelope ID。在产生任何作用前无法连到沙箱的 `workspace_write` 或 `runtime_prepare` 以 `rejected` 和 `resource_unavailable` 结束。若 Plugin 的 MCP server 声明了字面量 `http_headers`,或是带 `env_vars` 的 stdio server,owner 会在暂存其任何内容之前使其失败。无法观察到结果的文件变更或 setup 步骤会隔离 owner,直到 home 被删除:它不再发送任何变更,之后每个 `workspace_write` 和 `runtime_prepare` 都以 `unknown` 结束。没有 owner 的 Session 不支持上述任何操作,Runtime 以各自的类型化错误码拒绝:只读 preparation 为 `unsupported_read_preparation`;带 `local_environment` 的 Executor 配置为 `invalid_configuration`;`runtime_prepare` 为 `runtime_preparation_unsupported`;`workspace_write` 为 `write_unsupported`;`workspace_read` 和 `workspace_export` 为 `read_unsupported`。 -Core 先记录释放并推进 epoch,再发送任何消息;记录即撤回该分配的 attach grant。随后 Core 让 relay 吊销 Environment 的 Link resource 的当前 generation,使凭该 grant 打开的 attachment 在 Runtime 收到释放之前关闭。删除 Session 以 `remove_home: true` 释放其分配;释放其 Environment 发送 `false`。删除从不吊销共享的 Runtime 凭据。`assignment_release` 立即约束该分配。随后 Runtime 停止 Session 的工作:仍在接收内容、或已提交但尚未应用的传输以 `assignment_stale` 结束;它释放只读准备,并等待每个 workspace 读取、写入、导出和 Runtime 准备发送结果。它关闭 Session 的 Executor,随后释放其 Environment owner 持有的资源,并在要求时删除原生 home;此后才回复 `released` 或 `home_removed`。未完成的清理回复 `failed` 和 `cleanup_unconfirmed`,同一 epoch 的重试会重复清理。声明 `home_removal` 不支持的 Runtime 以 `unsupported_operation` 回答 `remove_home: true`,Core 只要求它释放。Core 根据匹配的 `released` 或 `home_removed` 记录释放已应用;没有 Runtime 能处理该释放时立即记录:发给无授权 Runtime 的释放在记录时即结清,吊销 Runtime 会结清它的释放。Core 在 Runtime 连接时重发所有未确认的释放;失败的释放退避重试,等待最久的释放先发送,因此失败的释放不会拖延其他释放。`environment_quiesce` 只作用于指定的 Environment:该 Environment 的某个 Session 仍有进行中的工作时,它以 `resource_busy` 失败;否则 Runtime 关闭该 Environment 的 Executor,让其 owner 释放所持有的资源,并回复 `environment_quiesced`。在匹配的 `environment_resume`(携带使其 quiesce 的分配)到达之前,Runtime 对该 Environment 只准入其 Session 的释放;其他 Environment 的 Session 继续运行。 +Core 先记录释放并推进 epoch,再发送任何消息;记录即撤回该分配的 attach grant。随后 Core 让 relay 吊销 Environment 的 Link resource 的当前 generation,使凭该 grant 打开的 attachment 在 Runtime 收到释放之前关闭。删除 Session 以 `remove_home: true` 释放其分配;释放其 Environment 发送 `false`。删除从不吊销共享的 Runtime 凭据。`assignment_release` 立即约束该分配。随后 Runtime 停止 Session 的工作:仍在接收内容、或已提交但尚未应用的传输以 `assignment_stale` 结束;它释放只读准备,并等待每个 workspace 读取、写入、导出和 Runtime 准备发送结果。它关闭 Session 的 Executor,随后释放其 Environment owner 持有的资源,并在要求时删除原生 home;此后才回复 `released` 或 `home_removed`。未完成的清理回复 `failed` 和 `cleanup_unconfirmed`,同一 epoch 的重试会重复清理。声明 `home_removal` 不支持的 Runtime 以 `unsupported_operation` 回答 `remove_home: true`,Core 只要求它释放。Core 根据匹配的 `released` 或 `home_removed` 记录释放已应用;没有 Runtime 能处理该释放时立即记录:发给无授权 Runtime 的释放在记录时即结清,吊销 Runtime 会结清它的释放。Core 在 Runtime 连接时重发所有未确认的释放;失败的释放退避重试,等待最久的释放先发送,因此失败的释放不会拖延其他释放。`environment_quiesce` 只作用于指定的 Environment:该 Environment 的某个 Session 仍有进行中的工作时,它以 `resource_busy` 失败;否则 Runtime 关闭该 Environment 的 Executor,让其 owner 释放所持有的资源,并回复 `environment_quiesced`。在匹配的 `environment_resume`(携带使其 quiesce 的分配)到达之前,Runtime 对该 Environment 只准入其 Session 的释放,并以 `protocol_error` `resource_unavailable` 回答它的其他任何 frame(包括绑定);其他 Environment 的 Session 继续运行。 Runtime 对无法路由的 Core frame 回复 `protocol_error`,回显请求 ID,并携带其类型和错误码。 ## 准备与执行顺序 {#preparation-and-execution-order} -无论托管还是用户自有环境,每条连接都通过 `runtime_prepare` 初始化 Environment;[Environment 契约](../../contracts/agents-api/zh/environments.md#runtime-capability-preparation)负责准备内容和时机。传输文件或 archive 时,发送 `begin`,等待 `ready`,发送有序 chunk 并等待每个匹配的 `received` offset,再发送 `commit` 并等待 `completed`。初始化和终结阶段使用不含文件数据的类型化 header。使用共享 validator 验证预期结果、offset、size 和有限错误 code。每个 Session 同一时间每种 transfer 至多一个;一个 transfer,或结果未知的 transfer,只阻塞其所属 Session 的 `execution_start` 和 transfer,以及其 Environment 的 quiesce。一条连接上所有 Session 同时暂存的 body 总量以 `workspace_write` 与 `runtime_prepare` 限制之和为上限;超出上限的 `begin` 以 `write_capacity` 或 `runtime_preparation_capacity` 被拒绝。chunk 回执确认暂存字节,不确认安装;完成的 commit 确认该操作,不证明后续 Turn 已运行。 +无论托管还是用户自有环境,每条连接都通过 `runtime_prepare` 初始化 Environment;[Environment 契约](../../contracts/agents-api/zh/environments.md#runtime-capability-preparation)负责准备内容和时机。传输文件或 archive 时,发送 `begin`,等待 `ready`,发送有序 chunk 并等待每个匹配的 `received` offset,再发送 `commit` 并等待 `completed`。初始化和终结阶段使用不含文件数据的类型化 header。使用共享 validator 验证预期结果、offset、size 和有限错误 code。每个 Session 同一时间每种 transfer 至多一个;无论结果如何,一个 transfer 在发送结果之前只阻塞其所属 Session 的 `execution_start` 和 transfer,以及其 Environment 的 quiesce。一条连接上所有 Session 同时暂存的 body 总量以 `workspace_write` 与 `runtime_prepare` 限制之和为上限;超出上限的 `begin` 以 `write_capacity` 或 `runtime_preparation_capacity` 被拒绝。chunk 回执确认暂存字节,不确认安装;完成的 commit 确认该操作,不证明后续 Turn 已运行。 执行 Turn 分为五步: @@ -204,7 +204,7 @@ Core 在 Turn outcome 中将接受的值保存为 `engine_error_code` 和 `engin Core 在 Worker 的 Session 调度预约上运行空闲目录读取,活动执行时针对精确 Run。它在有限时 read 与 release 期间保留预约,仅在确认 close 后返回数据(不完整读取或不确定清理返回 unavailable,不包含数据),在交付结果前释放预约,并在完成或失败后撤销限定作用域的读取凭据。Runtime 保留不确定清理的所有权和容量。[Environment Files 契约](../../contracts/agents-api/zh/environment-files.md)负责公开授权、路径和分页。 -`workspace_write` 在原生 writer 运行前,通过已确认的 64 KiB frame 传输完整且有界的 body,验证声明的 digest,不运行模型。私有 transfer 限制为 50 MiB,与公开 API 在任何 Runtime 工作前检查的 5 MiB decoded inline 限制独立。Runtime 在接收或应用写入时排除该 Session 的执行;格式错误、不完整或到期的 transfer 不会到达 installer。精确的 commit 或拒绝回执释放 mutation owner。缺失或有歧义的回执保留不确定性:observer 取消和本地进程退出不能证明没有改变任何内容。公开准入前,Core 在 Session lock 下持久预约写入,跨重启阻止后继 mutation,直到精确结算;请求不重放。在各平台上,owner 都在 daemon 内列举目录、创建文件和导出输出,不使用外部 helper 或 staging directory。 +`workspace_write` 在原生 writer 运行前,通过已确认的 64 KiB frame 传输完整且有界的 body,验证声明的 digest,不运行模型。私有 transfer 限制为 50 MiB,与公开 API 在任何 Runtime 工作前检查的 5 MiB decoded inline 限制独立。Runtime 在接收或应用写入时排除该 Session 的执行;格式错误、不完整或到期的 transfer 不会到达 installer。精确的 commit 或拒绝回执结算该写入。缺失或有歧义的回执把不确定性留给 [Environment owner](#session-assignments):observer 取消和本地进程退出不能证明没有改变任何内容。公开准入前,Core 在 Session lock 下持久预约写入,跨重启阻止后继 mutation,直到精确结算;请求不重放。在各平台上,owner 都在 daemon 内列举目录、创建文件和导出输出,不使用外部 helper 或 staging directory。 ## MCP 连接权限 {#mcp-connection-authority} diff --git a/internal/agentdaemon/proto/selection.go b/internal/agentdaemon/proto/selection.go index 9e995075c..47d30a820 100644 --- a/internal/agentdaemon/proto/selection.go +++ b/internal/agentdaemon/proto/selection.go @@ -34,6 +34,8 @@ type Selection struct { // selected. Environment string InstalledCapabilities bool + // NativeSessionRecovery requires the Session's existing native history. + NativeSessionRecovery bool MultiAgent bool Functions bool DeferredFunctions bool @@ -100,6 +102,8 @@ func ValidateSelection(d Declaration, s Selection) error { return reject("environment", "The harness does not support environment none.") case s.Environment == "local" && !c.LocalEnvironment.IsSupported(): return reject("environment", "The harness does not support this execution environment.") + case s.NativeSessionRecovery && !c.NativeSessionRecovery.IsSupported(): + return reject("", "The harness does not support native session recovery.") case s.MultiAgent && !c.SubagentObservations.IsSupported(): return reject("agent.multi_agent", "The harness does not support multi-agent execution.") case s.MultiAgent && (s.Functions || selectedMCP): @@ -219,7 +223,7 @@ func blankTextRune(r rune) bool { // without the Turn's messages and the Environment's installed MCP servers, // which its owner resolves only when it prepares an Executor. func (r PromptRequestPayload) Selection() Selection { - s := Selection{MultiAgent: r.ObserveSubagentIdentities, Functions: len(r.FunctionTools) > 0, ToolSearch: r.ToolSearch} + s := Selection{NativeSessionRecovery: r.RequireExistingNativeSession, MultiAgent: r.ObserveSubagentIdentities, Functions: len(r.FunctionTools) > 0, ToolSearch: r.ToolSearch} for _, tool := range r.FunctionTools { s.DeferredFunctions = s.DeferredFunctions || tool.DeferLoading } diff --git a/internal/harnessconfig/builtin/selection_test.go b/internal/harnessconfig/builtin/selection_test.go index 13c213a8f..391a6ddf5 100644 --- a/internal/harnessconfig/builtin/selection_test.go +++ b/internal/harnessconfig/builtin/selection_test.go @@ -47,6 +47,7 @@ func TestSelectionsAgainstEachDeclaration(t *testing.T) { {"local environment", proto.Selection{Environment: "local"}, nil, true}, {"installed capabilities", proto.Selection{Environment: "local", InstalledCapabilities: true}, nil, true}, {"multi-agent", proto.Selection{Environment: "none", MultiAgent: true}, nil, true}, + {"native session recovery", proto.Selection{Environment: "none", NativeSessionRecovery: true}, map[string]string{"mcode": ""}, true}, {"function tools", proto.Selection{Environment: "none", Functions: true}, map[string]string{"mcode": tool}, true}, {"tool search", search, map[string]string{"codex": tool, "mcode": tool}, true}, {"json_schema output", proto.Selection{Environment: "none", OutputSchema: schema}, map[string]string{"codex": format, "mcode": format}, true}, diff --git a/services/core/internal/agents/configuration.go b/services/core/internal/agents/configuration.go index 73211692f..b6c20f3ff 100644 --- a/services/core/internal/agents/configuration.go +++ b/services/core/internal/agents/configuration.go @@ -96,7 +96,7 @@ func validateModelExecution(configuration json.RawMessage, provider *v1.ModelPro if !ok || json.Unmarshal(configuration, &agent) != nil { return ErrInvalidInput } - err := proto.ValidateSelection(declared.Declaration, v1.HarnessSelection(v1.Agent{MultiAgent: agent.MultiAgent, Text: agent.Text, Tools: agent.Tools}, nil)) + err := proto.ValidateSelection(declared.Declaration, HarnessSelection(v1.Agent{MultiAgent: agent.MultiAgent, Text: agent.Text, Tools: agent.Tools}, nil)) var invalid *proto.SelectionError if errors.As(err, &invalid) { // A saved Agent's fields are top-level request fields. diff --git a/contracts/agents-api/v1/harness_selection.go b/services/core/internal/agents/harness_selection.go similarity index 91% rename from contracts/agents-api/v1/harness_selection.go rename to services/core/internal/agents/harness_selection.go index 22658e3a2..21f4fead5 100644 --- a/contracts/agents-api/v1/harness_selection.go +++ b/services/core/internal/agents/harness_selection.go @@ -1,8 +1,9 @@ -package v1 +package agents import ( "encoding/json" + v1 "github.com/MiniMax-AI/OpenAgentCore/contracts/agents-api/v1" "github.com/MiniMax-AI/OpenAgentCore/internal/agentdaemon/proto" ) @@ -10,7 +11,7 @@ import ( // proto.ValidateSelection checks against a Harness declaration. Callers // validate the tool declarations themselves. An MCP server with a // credential_id is authenticated. -func HarnessSelection(agent Agent, environment *Environment) proto.Selection { +func HarnessSelection(agent v1.Agent, environment *v1.Environment) proto.Selection { selection := proto.Selection{MultiAgent: agent.MultiAgent.Enabled, TextVerbosity: agent.Text.Verbosity} if agent.Text.Format.Type == "json_schema" { selection.OutputSchema = agent.Text.Format.Schema diff --git a/services/core/internal/api/mcp_configuration.go b/services/core/internal/api/mcp_configuration.go index 8764c4f8d..bc89abad9 100644 --- a/services/core/internal/api/mcp_configuration.go +++ b/services/core/internal/api/mcp_configuration.go @@ -3,10 +3,10 @@ package api import ( "encoding/json" "errors" - "net/url" "strings" v1 "github.com/MiniMax-AI/OpenAgentCore/contracts/agents-api/v1" + "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/execution" ) const mcpHTTPOnly = "MCP currently supports HTTP transport only." @@ -43,8 +43,7 @@ func resolveMCPTool(raw json.RawMessage, saved bool) (json.RawMessage, error) { if len(input.RequestMetadata) != 0 { return nil, errors.New("Nonempty MCP request_metadata is not supported yet.") } - u, err := url.Parse(transport.ServerURL) - if err != nil || u.Hostname() == "" || (u.Scheme != "http" && u.Scheme != "https") || u.User != nil || strings.Contains(transport.ServerURL, "#") || u.RawQuery != "" || u.ForceQuery { + if !execution.ValidMCPServerURL(transport.ServerURL) { return nil, errors.New("MCP server_url must be an absolute HTTP(S) URL without credentials, query or fragment.") } if transport.Authorization != nil { diff --git a/services/core/internal/execution/mcp.go b/services/core/internal/execution/mcp.go index 40abec568..ef636eafb 100644 --- a/services/core/internal/execution/mcp.go +++ b/services/core/internal/execution/mcp.go @@ -77,8 +77,7 @@ func executionTools(raw []json.RawMessage) (executionToolSet, error) { if decoder.Decode(&tool) != nil || strings.TrimSpace(tool.ServerLabel) == "" || names[tool.ServerLabel] || (tool.ConnectionOrigin != "service" && tool.ConnectionOrigin != "environment") || len(tool.RequestMetadata) != 0 || tool.Transport.Type != "http" || tool.Transport.Headers != nil { return executionToolSet{}, errors.New("unsupported execution MCP configuration") } - u, err := url.Parse(tool.Transport.ServerURL) - if err != nil || u.Hostname() == "" || (u.Scheme != "http" && u.Scheme != "https") || u.User != nil || strings.Contains(tool.Transport.ServerURL, "#") || u.RawQuery != "" || u.ForceQuery { + if !ValidMCPServerURL(tool.Transport.ServerURL) { return executionToolSet{}, errors.New("unsupported execution MCP URL") } if tool.AllowedTools != nil { @@ -95,3 +94,9 @@ func executionTools(raw []json.RawMessage) (executionToolSet, error) { resolved, err := functionTools(functions) return executionToolSet{Functions: resolved, MCP: servers, Search: search, DisableProgrammatic: disableProgrammatic}, err } + +// ValidMCPServerURL is Core's one rule for an MCP server_url. +func ValidMCPServerURL(raw string) bool { + u, err := url.Parse(raw) + return err == nil && u.Hostname() != "" && (u.Scheme == "http" || u.Scheme == "https") && u.User == nil && !strings.Contains(raw, "#") && u.RawQuery == "" && !u.ForceQuery +} diff --git a/services/core/internal/execution/recovery_test.go b/services/core/internal/execution/recovery_test.go index 443af39c4..17f826cf8 100644 --- a/services/core/internal/execution/recovery_test.go +++ b/services/core/internal/execution/recovery_test.go @@ -1,6 +1,7 @@ package execution import ( + "errors" "testing" "github.com/MiniMax-AI/OpenAgentCore/internal/agentdaemon/proto" @@ -15,8 +16,9 @@ func TestExistingSessionRecoveryRequiresVerifiedCapability(t *testing.T) { wantRecovery := started && nativeID == "" req, err := (&Dispatcher{SessionsReader: frozenProvider{engine: engine}}).executionRequest(t.Context(), sessions.Session{ID: "session", Engine: engine}, Snapshot{}, proto.Declaration{Capabilities: proto.AgentKindCapabilities{NativeSessionRecovery: proto.CapabilityFromBool(capable)}}, sessions.ExecutionBinding{HasStartedTurn: started, NativeSessionID: nativeID}) if wantRecovery && !capable { - if err == nil { - t.Fatal("unverified recovery admitted", engine) + var rejected *proto.SelectionError + if !errors.As(err, &rejected) { + t.Fatal("unverified recovery admitted", engine, err) } continue } diff --git a/services/core/internal/execution/request.go b/services/core/internal/execution/request.go index 5387705b8..d9baec6d0 100644 --- a/services/core/internal/execution/request.go +++ b/services/core/internal/execution/request.go @@ -2,7 +2,6 @@ package execution import ( "context" - "errors" "github.com/MiniMax-AI/OpenAgentCore/internal/agentdaemon/proto" "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/sessions" @@ -11,8 +10,8 @@ import ( func (d *Dispatcher) executionRequest(ctx context.Context, session sessions.Session, snapshot Snapshot, declaration proto.Declaration, bound sessions.ExecutionBinding) (proto.PromptRequestPayload, error) { recoverNativeSession := bound.HasStartedTurn && bound.NativeSessionID == "" - if recoverNativeSession && !declaration.Capabilities.NativeSessionRecovery.IsSupported() { - return proto.PromptRequestPayload{}, errors.New("native session recovery is unavailable") + if err := proto.ValidateSelection(declaration, proto.Selection{NativeSessionRecovery: recoverNativeSession}); err != nil { + return proto.PromptRequestPayload{}, err } tools, err := executionTools(snapshot.Agent.Tools) if err != nil { diff --git a/services/core/internal/execution/support.go b/services/core/internal/execution/support.go index f158b9faa..d048f97c0 100644 --- a/services/core/internal/execution/support.go +++ b/services/core/internal/execution/support.go @@ -5,9 +5,9 @@ import ( "errors" "strings" - v1 "github.com/MiniMax-AI/OpenAgentCore/contracts/agents-api/v1" "github.com/MiniMax-AI/OpenAgentCore/internal/agentdaemon/proto" "github.com/MiniMax-AI/OpenAgentCore/internal/harnessconfig/builtin" + "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/agents" "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/runtimegateway" "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/sessions" ) @@ -59,7 +59,7 @@ func ValidateSessionConfiguration(engine string, configuration json.RawMessage) // harnessSelection projects the frozen Session configuration. A server with a // frozen credential binding is authenticated. func harnessSelection(snapshot Snapshot) (proto.Selection, error) { - selection := v1.HarnessSelection(snapshot.Agent, snapshot.Environment) + selection := agents.HarnessSelection(snapshot.Agent, snapshot.Environment) selected, err := selectedMCPCredentials(snapshot) if err != nil { return proto.Selection{}, err