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
2 changes: 2 additions & 0 deletions apps/daemon/internal/agent/codex/preparation_router_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -88,10 +88,12 @@ func TestPreparationRouterRetainsActualNativeChild(t *testing.T) {
if err != nil {
t.Fatal(err)
}
env.Assignment = proto.AssignmentRef{SessionID: session, AssignmentID: "assignment", Epoch: 1}
if err = r.Handle(t.Context(), env); err != nil {
t.Fatal(err)
}
}
send(proto.TypeAssignmentBind, "bind", proto.AssignmentBindPayload{EnvironmentID: environment})
await := func(state string) proto.PreparationStatusPayload {
t.Helper()
timer := time.NewTimer(4 * time.Second)
Expand Down
24 changes: 24 additions & 0 deletions apps/daemon/internal/agenthost/admit_linux_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -251,5 +251,29 @@ func TestViewExecutorReceivesTheGatewayRequest(t *testing.T) {
if _, err := os.Stat(f.session.Home.Host); dials.Load() != 0 || err != nil || len(leftEntries(t, f.cfg)) != 0 {
t.Errorf("%d dials, home %v and transient entries %v after the preparation", dials.Load(), err, leftEntries(t, f.cfg))
}
}

// TestReleaseRemovesTheHome checks that the agent host removes a Session's
// home behind its assignment's release.
func TestReleaseRemovesTheHome(t *testing.T) {
f := newViewFixture(t)
var dials atomic.Int32
d := newDaemon(t, f.cfg, deps{dial: countingDial(&dials), tasks: noTasks})
b := newBinding(newResource())
if _, p := d.prepare(t, b, request("viewed", "/workspace", "https://model.test", "sk-test")); p.State != "failed" {
t.Fatalf("the preparation is %s, want failed with the factory", p.State)
}
if _, err := os.Stat(f.session.Home.Host); err != nil {
t.Fatalf("the home: %v", err)
}
released := ref(b)
released.Epoch++
id := "release"
d.handle(t, released, proto.TypeAssignmentRelease, id, proto.AssignmentReleasePayload{RemoveHome: true})
if status := d.status(t, id); status.State != proto.AssignmentHomeRemoved {
t.Fatalf("the release is %s (%s), want home_removed", status.State, status.ErrorCode)
}
if _, err := os.Stat(filepath.Dir(f.session.Home.Host)); !errors.Is(err, os.ErrNotExist) {
t.Errorf("the Session directory remains: %v", err)
}
}
66 changes: 52 additions & 14 deletions apps/daemon/internal/agenthost/agenthost_linux_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -144,8 +144,9 @@ func leftEntries(t *testing.T, cfg Config) []string {
}

// daemon drives Sessions through a dispatch Router, as the daemon does. It
// binds each request to the Session its state key names and records the
// latest Executor the agent host opened for each Session.
// binds each request to the Session of the assignment the Router admitted it
// under, records the latest Executor the agent host opened for each Session,
// and removes a released Session's home.
type daemon struct {
router *dispatch.Router
// mcp is the installed MCP that the Environment's preparation resolves
Expand Down Expand Up @@ -175,7 +176,16 @@ func newDaemon(t *testing.T, cfg Config, d deps) *daemon {
return e, err
})
var err error
if dm.router, err = dispatch.New(dispatch.Config{Registry: reg, Sender: dm, SessionEnvironments: true, Log: slog.New(slog.DiscardHandler)}); err != nil {
removeHome := func(session string) error {
dm.mu.Lock()
b, ok := dm.bindings[session]
dm.mu.Unlock()
if !ok {
return fmt.Errorf("%w: no binding", ErrInvalidSession)
}
return (&Host{cfg: cfg}).RemoveHome(b.SessionID)
}
if dm.router, err = dispatch.New(dispatch.Config{Registry: reg, Sender: dm, SessionEnvironments: true, RemoveHome: removeHome, Log: slog.New(slog.DiscardHandler)}); err != nil {
t.Fatal(err)
}
t.Cleanup(func() { dm.shutdown() })
Expand All @@ -188,13 +198,18 @@ const stateKeyPrefix = "agents-api-"
func (dm *daemon) bind(req proto.PromptRequestPayload) (Binding, Environment, error) {
dm.mu.Lock()
defer dm.mu.Unlock()
b, ok := dm.bindings[strings.TrimPrefix(req.AgentStateKey, stateKeyPrefix)]
if !ok {
b, ok := dm.bindings[req.Assignment.SessionID]
if !ok || ref(b) != req.Assignment {
return Binding{}, Environment{}, fmt.Errorf("%w: no binding", ErrInvalidSession)
}
return b, Environment{}, nil
}

// ref is the reference of b's assignment.
func ref(b Binding) proto.AssignmentRef {
return proto.AssignmentRef{SessionID: b.SessionID.String(), AssignmentID: b.AssignmentID.String(), Epoch: b.AssignmentEpoch}
}

func (dm *daemon) Send(_ context.Context, e proto.Envelope) error {
dm.frame(e.ID) <- e
return nil
Expand All @@ -211,13 +226,14 @@ func (dm *daemon) frame(id string) chan proto.Envelope {
return ch
}

// handle hands the Router an envelope from Core.
func (dm *daemon) handle(t *testing.T, typ, id string, payload any) {
// handle hands the Router an envelope from Core under the assignment ref.
func (dm *daemon) handle(t *testing.T, ref proto.AssignmentRef, typ, id string, payload any) {
t.Helper()
e, err := proto.NewEnvelope(typ, id, payload)
if err != nil {
t.Fatal(err)
}
e.Assignment = ref
if err := dm.router.Handle(context.Background(), e); err != nil {
t.Fatalf("%s: %v", typ, err)
}
Expand All @@ -235,17 +251,39 @@ func (dm *daemon) next(t *testing.T, id string) proto.Envelope {
}
}

// prepare prepares an Executor of b's Session for req. It returns the
// request ID and the preparation's first status other than preparing.
func (dm *daemon) prepare(t *testing.T, b Binding, req proto.PromptRequestPayload) (string, proto.PreparationStatusPayload) {
// assign binds b's Session to the Router in environment.
func (dm *daemon) assign(t *testing.T, b Binding, environment string) {
t.Helper()
session := b.SessionID.String()
dm.mu.Lock()
dm.bindings[session] = b
dm.bindings[b.SessionID.String()] = b
dm.mu.Unlock()
id := sandboxwire.NewID().String()
dm.handle(t, ref(b), proto.TypeAssignmentBind, id, proto.AssignmentBindPayload{EnvironmentID: environment})
if status := dm.status(t, id); status.State != proto.AssignmentBound {
t.Fatalf("the bind is %s (%s), want bound", status.State, status.ErrorCode)
}
}

// status returns the assignment status sent with id.
func (dm *daemon) status(t *testing.T, id string) proto.AssignmentStatusPayload {
t.Helper()
var status proto.AssignmentStatusPayload
if err := dm.next(t, id).DecodePayload(&status); err != nil {
t.Fatal(err)
}
return status
}

// prepare binds b's Session and prepares an Executor of it for req. It
// returns the request ID and the preparation's first status other than
// preparing.
func (dm *daemon) prepare(t *testing.T, b Binding, req proto.PromptRequestPayload) (string, proto.PreparationStatusPayload) {
t.Helper()
dm.assign(t, b, req.EnvironmentID())
session := b.SessionID.String()
req.AgentStateKey = stateKeyPrefix + session
id := sandboxwire.NewID().String()
dm.handle(t, proto.TypeExecutionPrepare, id, proto.ExecutionPreparePayload{SessionID: session, Configuration: req})
dm.handle(t, ref(b), proto.TypeExecutionPrepare, id, proto.ExecutionPreparePayload{SessionID: session, Configuration: req})
for {
var p proto.PreparationStatusPayload
if err := dm.next(t, id).DecodePayload(&p); err != nil {
Expand All @@ -266,7 +304,7 @@ func (dm *daemon) start(t *testing.T, b Binding, req proto.PromptRequestPayload,
t.Fatalf("the preparation is %s (%s), want ready", p.State, p.ErrorCode)
}
run := "run-" + id
dm.handle(t, proto.TypeExecutionStart, id, proto.ExecutionStartPayload{Handle: p.Handle, ExecutorID: p.ExecutorID, RunID: run, Input: proto.TextInput(text)})
dm.handle(t, ref(b), proto.TypeExecutionStart, id, proto.ExecutionStartPayload{Handle: p.Handle, ExecutorID: p.ExecutorID, RunID: run, Input: proto.TextInput(text)})
for {
if err := dm.next(t, id).DecodePayload(&p); err != nil {
t.Fatal(err)
Expand Down
2 changes: 1 addition & 1 deletion apps/daemon/internal/agenthost/view_linux_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -278,7 +278,7 @@ func TestSessionRunsInAViewOverItsAttachment(t *testing.T) {
t.Errorf("the other Session's command printed %q and exited %d; stderr %s", r.Stdout, r.Code, r.Stderr)
}
// Cancelling the waiting Turn ends it with its view.
d.handle(t, proto.TypePromptCancel, run, proto.PromptCancelPayload{})
d.handle(t, ref(waiting), proto.TypePromptCancel, run, proto.PromptCancelPayload{})
d.done(t, run)
if err := d.shutdown(); err != nil {
t.Fatalf("Shutdown = %v", err)
Expand Down
30 changes: 26 additions & 4 deletions apps/daemon/internal/cli/claude_sdk_live_linux_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -15,6 +15,7 @@ import (
"github.com/MiniMax-AI/OpenAgentCore/apps/daemon/internal/agent"
"github.com/MiniMax-AI/OpenAgentCore/apps/daemon/internal/dispatch"
"github.com/MiniMax-AI/OpenAgentCore/internal/agentdaemon/proto"
"github.com/MiniMax-AI/OpenAgentCore/internal/agentdaemon/proto/prototest"
"github.com/google/uuid"
)

Expand Down Expand Up @@ -68,6 +69,7 @@ func TestLiveRegisteredClaudeSDK(t *testing.T) {
Cancelled bool `json:"cancelled"`
}
nonce := "registered-function-" + uuid.NewString()
ref := proto.AssignmentRef{SessionID: prototest.SessionID, AssignmentID: "registered-acceptance", Epoch: 1}
run := func(index int, prompt, resume string, callFunction, cancelOnText bool) execution {
t.Helper()
reg := agent.NewRegistry()
Expand All @@ -87,21 +89,24 @@ func TestLiveRegisteredClaudeSDK(t *testing.T) {
ctx, cancel := context.WithTimeout(t.Context(), 120*time.Second)
defer cancel()
id := uuid.NewString()
request := proto.PromptRequestPayload{RunID: id, AgentKind: "claude_sdk", Input: proto.TextInput(prompt), AgentStateKey: "registered-acceptance", AgentSessionID: resume, StrictResume: true, ReleaseOnCompletion: true, ObserveMessages: true, ObserveToolObservations: true, DisableExecutionEnvironment: true, DisableSubagents: true, ExecutionControls: &proto.ExecutionControls{WebSearch: "disabled", TextVerbosity: "medium"}, Model: "MiniMax-M3"}
request := proto.PromptRequestPayload{AgentKind: "claude_sdk", AgentStateKey: prototest.StateKey, AgentSessionID: resume, RequireExistingNativeSession: resume != "", StrictResume: true, ReleaseOnCompletion: true, ObserveMessages: true, ObserveToolObservations: true, DisableExecutionEnvironment: true, DisableSubagents: true, ExecutionControls: &proto.ExecutionControls{WebSearch: "disabled", TextVerbosity: "medium"}, Model: "MiniMax-M3"}
if callFunction {
request.FunctionTools = []proto.FunctionTool{{Name: "lookup", Description: "Return a verification value.", Parameters: json.RawMessage(`{"type":"object","properties":{"id":{"type":"string"}},"required":["id"],"additionalProperties":false}`)}}
}
handle := func(kind string, payload any) {
send := func(kind, envID string, payload any) {
t.Helper()
env, err := proto.NewEnvelope(kind, id, payload)
env, err := proto.NewEnvelope(kind, envID, payload)
if err != nil {
t.Fatal(err)
}
env.Assignment = ref
if err := router.Handle(ctx, env); err != nil {
t.Fatal("registered router request failed", err)
}
}
handle(proto.TypePromptRequest, request)
handle := func(kind string, payload any) { t.Helper(); send(kind, id, payload) }
send(proto.TypeAssignmentBind, "bind", proto.AssignmentBindPayload{})
send(proto.TypeExecutionPrepare, "prepare", proto.ExecutionPreparePayload{SessionID: ref.SessionID, Configuration: request})
proof := execution{}
defer func() {
data, _ := json.MarshalIndent(proof, "", " ")
Expand All @@ -117,6 +122,23 @@ func TestLiveRegisteredClaudeSDK(t *testing.T) {
case <-ctx.Done():
t.Fatal("registered execution timed out", ctx.Err())
}
if event.Type == proto.TypeAssignmentStatus {
continue
}
if event.Type == proto.TypePreparationStatus {
var status proto.PreparationStatusPayload
if err := event.DecodePayload(&status); err != nil {
t.Fatal(err)
}
switch status.State {
case "ready":
send(proto.TypeExecutionStart, "prepare", proto.ExecutionStartPayload{Handle: status.Handle, ExecutorID: status.ExecutorID, RunID: id, Input: proto.TextInput(prompt)})
case "started":
default:
t.Fatal("registered preparation failed", status)
}
continue
}
if event.ID != id {
t.Fatal("event identity changed")
}
Expand Down
1 change: 1 addition & 0 deletions apps/daemon/internal/cli/connect.go
Original file line number Diff line number Diff line change
Expand Up @@ -340,6 +340,7 @@ func pumpConn(parentCtx context.Context, conn *transport.Conn, registry *agent.R
ActiveRequests: router.ActiveRuns(),
DaemonVersion: Version,
SupportedAgentKinds: kinds,
HomeRemoval: proto.CapabilityUnsupported,
}
}, obslog.Bg().With("component", "heartbeat"))

Expand Down
28 changes: 28 additions & 0 deletions apps/daemon/internal/cli/connect_cleanup_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -145,11 +145,16 @@ func testDisconnectedPumpCleanup(t *testing.T, suspend bool) {
t.Fatal("initial connection missing")
}
defer peer.Close()
ref := proto.AssignmentRef{SessionID: "cleanup", AssignmentID: "assignment", Epoch: 1}
if got := sendAssignment(t, peer, proto.TypeAssignmentBind, ref, proto.AssignmentBindPayload{}); got.State != proto.AssignmentBound {
t.Fatalf("bind = %+v", got)
}
env, err := proto.NewEnvelope(proto.TypeExecutionPrepare, "prepare", proto.ExecutionPreparePayload{SessionID: "cleanup",
Configuration: prototest.WithModel(proto.PromptRequestPayload{AgentKind: "cleanup", AgentStateKey: "agents-api-cleanup", StrictResume: true, DisableExecutionEnvironment: true})})
if err != nil {
t.Fatal(err)
}
env.Assignment = ref
if err := peer.WriteJSON(env); err != nil {
t.Fatal(err)
}
Expand Down Expand Up @@ -195,3 +200,26 @@ func testDisconnectedPumpCleanup(t *testing.T, suspend bool) {
t.Fatalf("factory=%d close=%d", factories.Load(), owner.closes.Load())
}
}

func sendAssignment(t *testing.T, peer *websocket.Conn, kind string, ref proto.AssignmentRef, payload any) proto.AssignmentStatusPayload {
t.Helper()
env, err := proto.NewEnvelope(kind, kind, payload)
if err != nil {
t.Fatal(err)
}
env.Assignment = ref
if err := peer.WriteJSON(env); err != nil {
t.Fatal(err)
}
_ = peer.SetReadDeadline(time.Now().Add(3 * time.Second))
for {
var reply proto.Envelope
if err := peer.ReadJSON(&reply); err != nil {
t.Fatal(err)
}
var status proto.AssignmentStatusPayload
if reply.Type == proto.TypeAssignmentStatus && reply.ID == kind && reply.Assignment == ref && reply.DecodePayload(&status) == nil {
return status
}
}
}
Loading
Loading