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

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
4 changes: 2 additions & 2 deletions apps/daemon/internal/agent/claudesdk/options.go
Original file line number Diff line number Diff line change
Expand Up @@ -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 {
Expand Down
2 changes: 1 addition & 1 deletion apps/daemon/internal/agent/mcode/execution.go
Original file line number Diff line number Diff line change
Expand Up @@ -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) {
Expand Down
1 change: 0 additions & 1 deletion apps/daemon/internal/agent/mcode/execution_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -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)
Expand Down
24 changes: 8 additions & 16 deletions apps/daemon/internal/agenthost/environment_linux_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -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"} {
Expand Down
2 changes: 1 addition & 1 deletion apps/daemon/internal/cli/agent_host_linux.go
Original file line number Diff line number Diff line change
Expand Up @@ -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")

Expand Down
11 changes: 8 additions & 3 deletions apps/daemon/internal/dispatch/assignment.go
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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()
Expand Down Expand Up @@ -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 {
Expand Down Expand Up @@ -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))
}()
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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 {
Expand Down
2 changes: 1 addition & 1 deletion apps/daemon/internal/dispatch/executor_handoff_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -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)
Expand Down
9 changes: 8 additions & 1 deletion apps/daemon/internal/dispatch/executor_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -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
})
Expand Down Expand Up @@ -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)
Expand Down
64 changes: 38 additions & 26 deletions apps/daemon/internal/dispatch/native_file_results_test.go
Original file line number Diff line number Diff line change
@@ -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"
)
Expand All @@ -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)
}
2 changes: 1 addition & 1 deletion apps/daemon/internal/dispatch/router.go
Original file line number Diff line number Diff line change
Expand Up @@ -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()

Expand Down
23 changes: 10 additions & 13 deletions apps/daemon/internal/dispatch/runtime_preparation.go
Original file line number Diff line number Diff line change
Expand Up @@ -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 {
Expand Down Expand Up @@ -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)
Expand Down
Loading
Loading