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
9 changes: 1 addition & 8 deletions apps/daemon/internal/agenthost/agenthost_linux_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -178,15 +178,8 @@ func newDaemon(t *testing.T, cfg Config, d deps) *daemon {
func (dm *daemon) route(t *testing.T, reg *agent.Registry) {
t.Helper()
dm.frames, dm.opened = map[string]chan proto.Envelope{}, map[string]*session{}
removeHome := func(session string) error {
id, err := canonicalID(session)
if err != nil {
return err
}
return dm.host.RemoveHome(id)
}
var err error
if dm.router, err = dispatch.New(dispatch.Config{Registry: reg, Sender: dm, Environments: dm.host.Environments, RemoveHome: removeHome, Log: slog.New(slog.DiscardHandler)}); err != nil {
if dm.router, err = dispatch.New(dispatch.Config{Registry: reg, Sender: dm, Environments: dm.host.Environments, RemoveHome: dm.host.RemoveHome, Log: slog.New(slog.DiscardHandler)}); err != nil {
t.Fatal(err)
}
t.Cleanup(func() { dm.shutdown() })
Expand Down
9 changes: 4 additions & 5 deletions apps/daemon/internal/agenthost/doc.go
Original file line number Diff line number Diff line change
Expand Up @@ -3,11 +3,10 @@
// Session's Link attachment. It needs Linux; elsewhere Open returns
// ErrUnsupported.
//
// The process that runs the agent host calls Open once at startup, hands
// Host.Registry and Host.Environments to the daemon's dispatch, which drives
// each Turn of each Executor, and calls Close once every Executor has closed. No production
// caller constructs Host.Registry yet; oac-daemon connect still runs
// Harnesses in the sandbox. Open takes two installation locks, which the
// The process that runs the agent host, oac-daemon agent-host, calls Open
// once at startup, hands Host.Registry and Host.Environments to the daemon's
// dispatch, which drives each Turn of each Executor, and calls Close once
// every Executor has closed. Open takes two installation locks, which the
// Host holds until Close: a flock on StateDir/lock for the Session
// directories and one on the Config.ViewCgroups directory for the view
// cgroups. Under the locks, Open checks the requirements that need nothing
Expand Down
46 changes: 46 additions & 0 deletions apps/daemon/internal/agenthost/environment_linux_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -204,6 +204,52 @@ func TestEnvironmentOwnerServesTheSandbox(t *testing.T) {
}
}

// TestSupersedingBindRebindsTheOwner checks that a bind of the Session's
// assignment at a higher epoch with a new grant drains the owner's
// attachment and gives the owner the new binding, under which its next write
// attaches to the sandbox again.
func TestSupersedingBindRebindsTheOwner(t *testing.T) {
if os.Getenv(gateEnv) != "1" {
t.Skipf("set %s=1 and run the test binary as root in a throwaway container; see the view suite", gateEnv)
}
sb := startSandbox(t, os.Getenv(sandboxIOEnv))
for _, name := range []string{"before.txt", "after.txt"} {
if err := os.RemoveAll(path.Join(sandboxWorkspace, name)); err != nil {
t.Fatal(err)
}
}
if err := os.MkdirAll(sandboxWorkspace, 0o777); err != nil {
t.Fatal(err)
}
cfg := Config{StateDir: t.TempDir(), RelayURL: sb.url, RuntimeID: sandboxwire.NewID(), Credential: []byte("runtime-credential"), Harnesses: agent.NewRegistry()}
sb.auth.AddRuntime(cfg.Credential, cfg.RuntimeID)
sb.ready(t, cfg)
dm := &daemon{host: &Host{cfg: cfg, owners: owners{d: deps{dial: relayDial(cfg)}}}}
dm.route(t, agent.NewRegistry())
b := sb.bind(cfg.RuntimeID, time.Minute)
dm.assign(t, b)
if r := dm.write(t, b, "before.txt", []byte("before")); r.Outcome != "completed" || dm.host.drained(t, b) {
t.Fatalf("the write at epoch 1 is %+v", r)
}

next := b
next.AssignmentEpoch, next.AttachGrant = 2, []byte("grant-"+sandboxwire.NewID().String())
sb.grant(next, cfg.RuntimeID, time.Minute)
dm.assign(t, next)
o := dm.host.owner(t, next)
rebound, drained := sameBinding(o.binding, next), o.link == nil
o.release()
if !rebound || !drained {
t.Fatalf("after the superseding bind the owner has the new binding %t and no attachment %t", rebound, drained)
}
if r := dm.write(t, next, "after.txt", []byte("after")); r.Outcome != "completed" {
t.Fatalf("the write at epoch 2 is %+v", r)
}
if body, err := os.ReadFile(path.Join(sandboxWorkspace, "after.txt")); err != nil || string(body) != "after" {
t.Fatalf("the sandbox has %q, %v", body, err)
}
}

// TestUnreachableSandboxRejectsRuntimePreparation checks that a
// runtime_prepare whose owner cannot reach the sandbox ends rejected, which
// leaves the Router free to run another Session's and to shut down.
Expand Down
8 changes: 5 additions & 3 deletions apps/daemon/internal/agenthost/host_linux_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -16,6 +16,8 @@ import (
"syscall"
"testing"

"github.com/google/uuid"

"github.com/MiniMax-AI/OpenAgentCore/apps/daemon/internal/agent"
"github.com/MiniMax-AI/OpenAgentCore/apps/daemon/internal/sessionview"
"github.com/MiniMax-AI/OpenAgentCore/apps/daemon/internal/sessionview/sessionviewtest"
Expand Down Expand Up @@ -193,15 +195,15 @@ func TestRemoveHome(t *testing.T) {
opened <- e
}()
<-claimed
if err := h.RemoveHome(b.SessionID); !errors.Is(err, ErrSessionExists) {
if err := h.RemoveHome(uuid.UUID(b.SessionID).String()); !errors.Is(err, ErrSessionExists) {
t.Fatalf("RemoveHome during an open = %v, want ErrSessionExists", err)
}
close(proceed)
e := <-opened
if e == nil {
t.Fatal("open returned no Executor to close")
}
if err := h.RemoveHome(b.SessionID); !errors.Is(err, ErrSessionExists) {
if err := h.RemoveHome(uuid.UUID(b.SessionID).String()); !errors.Is(err, ErrSessionExists) {
t.Errorf("RemoveHome before the Executor closed = %v, want ErrSessionExists", err)
}
if _, err := os.Stat(f.session.Home.Host); err != nil {
Expand All @@ -221,7 +223,7 @@ func TestRemoveHome(t *testing.T) {
}
// Removal is idempotent: an absent home is removed.
for range 2 {
if err := h.RemoveHome(b.SessionID); err != nil {
if err := h.RemoveHome(uuid.UUID(b.SessionID).String()); err != nil {
t.Fatalf("RemoveHome = %v", err)
}
}
Expand Down
3 changes: 1 addition & 2 deletions apps/daemon/internal/agenthost/host_other.go
Original file line number Diff line number Diff line change
Expand Up @@ -8,7 +8,6 @@ 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/sandboxwire"
)

// Open reports that the agent host needs Linux.
Expand All @@ -27,6 +26,6 @@ func (*Host) Environments(proto.AssignmentRef, proto.AssignmentBindPayload) disp
}

// RemoveHome reports that the agent host needs Linux.
func (*Host) RemoveHome(sandboxwire.ID) error {
func (*Host) RemoveHome(string) error {
return fmt.Errorf("%w: remove home", ErrUnsupported)
}
17 changes: 11 additions & 6 deletions apps/daemon/internal/agenthost/sessiondir_linux.go
Original file line number Diff line number Diff line change
Expand Up @@ -156,12 +156,17 @@ func sweep(stateDir string) error {
return nil
}

// RemoveHome drains and forgets the Session's Environment owner, then removes
// the Session's directory with its home once the Session's processes have
// settled. It returns ErrSessionExists while an Executor of the Session has
// not closed, and an Executor of the Session does not open while RemoveHome
// runs. A Session without a directory has nothing to remove.
func (h *Host) RemoveHome(id sandboxwire.ID) error {
// RemoveHome drains and forgets the Environment owner of the Session with
// the canonical ID session, then removes the Session's directory with its
// home once the Session's processes have settled, as
// dispatch.Config.RemoveHome. It returns ErrSessionExists while an Executor
// of the Session has not closed, and an Executor of the Session does not open
// while RemoveHome runs. A Session without a directory has nothing to remove.
func (h *Host) RemoveHome(session string) error {
id, err := canonicalID(session)
if err != nil {
return fmt.Errorf("%w: remove home: %w", ErrInvalidSession, err)
}
if !claimSession(id) {
return fmt.Errorf("%w: remove home", ErrSessionExists)
}
Expand Down
7 changes: 6 additions & 1 deletion apps/daemon/internal/agenthost/view_linux_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -566,10 +566,15 @@ func (sb *sandbox) stop() {
// runtimeID for lease.
func (sb *sandbox) bind(runtimeID sandboxwire.ID, lease time.Duration) Binding {
b := newBinding(sb.resource)
sb.grant(b, runtimeID, lease)
return b
}

// grant lets runtimeID attach b's Session to the resource for lease.
func (sb *sandbox) grant(b Binding, runtimeID sandboxwire.ID, lease time.Duration) {
sb.auth.AddGrant(b.AttachGrant, sandboxlinktest.Grant{RuntimeID: runtimeID, Resource: b.Resource, SessionID: b.SessionID,
AssignmentID: b.AssignmentID, AssignmentEpoch: b.AssignmentEpoch, Lease: lease,
Services: []sandboxlink.Service{sandboxlink.ServiceFile, sandboxlink.ServiceProcess, sandboxlink.ServiceNetwork}})
return b
}

func (sb *sandbox) dial(t *testing.T, cfg Config) *sandboxlink.AttachLink {
Expand Down
157 changes: 157 additions & 0 deletions apps/daemon/internal/cli/agent_host_linux.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,157 @@
//go:build linux

package cli

import (
"bytes"
"context"
"encoding/json"
"errors"
"fmt"
"log/slog"
"net/url"
"os"
"path/filepath"
"strings"
"time"

"github.com/google/uuid"
"golang.org/x/sys/unix"

"github.com/MiniMax-AI/OpenAgentCore/apps/daemon/internal/agent"
"github.com/MiniMax-AI/OpenAgentCore/apps/daemon/internal/agent/claudesdk"
"github.com/MiniMax-AI/OpenAgentCore/apps/daemon/internal/agent/codex"
"github.com/MiniMax-AI/OpenAgentCore/apps/daemon/internal/agent/mcode"
"github.com/MiniMax-AI/OpenAgentCore/apps/daemon/internal/agenthost"
"github.com/MiniMax-AI/OpenAgentCore/apps/daemon/internal/daemonize"
"github.com/MiniMax-AI/OpenAgentCore/apps/daemon/internal/dispatch"
"github.com/MiniMax-AI/OpenAgentCore/apps/daemon/internal/paths"
"github.com/MiniMax-AI/OpenAgentCore/apps/daemon/internal/transport"
"github.com/MiniMax-AI/OpenAgentCore/internal/agentdaemon/proto"
obslog "github.com/MiniMax-AI/OpenAgentCore/internal/obs/log"
"github.com/MiniMax-AI/OpenAgentCore/internal/sandboxwire"
)

// The agent-host image's layout (deploy/distribution/AgentHost.Dockerfile)
// and the container's own state.
const (
agentHostManifest = "/opt/oac/harnesses.json"
agentHostShim = "/opt/oac/bin/oac-process-shim"
agentHostCADir = "/usr/share/ca-certificates/mozilla"
// agentHostState keeps each Session's home across restarts.
agentHostState = "/var/lib/oac/agent-host"
// agentHostCgroup is where the agent host mounts the container's own
// cgroup v2 hierarchy; its views directory is delegated to the views.
agentHostCgroup = "/run/oac/cgroup"
// agentHostUnreachable is how long Core may stay unreachable before the
// agent host exits, so that its supervisor restarts it.
agentHostUnreachable = 2 * time.Minute
)

// agentHostUIDs is the range the agent host runs its Executors as.
var agentHostUIDs = agenthost.UIDRange{First: 70000, Count: 4096}

func runAgentHost(rc *runContext, args []string) error {
return serveAgentHost(context.Background(), rc, args, harnessDeclarations)
}

// serveAgentHost runs the agent host in its container: the Harnesses that
// the image's manifest installs and that declare a view serve Sessions bound
// to the Runtime of the identity, whose Environments are the agent host's.
func serveAgentHost(parent context.Context, rc *runContext, args []string, declarations []agent.Declaration) error {
flags := newFlagSet("agent-host")
identityFile := flags.String("identity-file", "", "path to the agent host's identity JSON")
coreURL := flags.String("core-url", "", "Core origin, https or a loopback http origin")
if err := flags.Parse(args); err != nil {
return fmt.Errorf("agent-host: parse flags: %w", err)
}
if flags.NArg() != 0 {
return fmt.Errorf("agent-host: unexpected arguments %q", flags.Args())
}
var identity struct {
RuntimeID string `json:"runtime_id"`
Credential string `json:"credential"`
}
raw, err := os.ReadFile(*identityFile)
if err != nil {
return fmt.Errorf("agent-host: identity: %w", err)
}
decoder := json.NewDecoder(bytes.NewReader(raw))
decoder.DisallowUnknownFields()
if err := decoder.Decode(&identity); err != nil {
return fmt.Errorf("agent-host: identity: %w", err)
}
runtimeID, err := uuid.Parse(identity.RuntimeID)
if err != nil || runtimeID.String() != identity.RuntimeID || identity.Credential == "" {
return errors.New("agent-host: identity needs a canonical runtime_id and a credential")
}
// Core derives the same URLs from OAC_PUBLIC_URL; its Link exists only
// on https or a loopback origin, which agenthost.Open checks.
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)
}
wsOrigin := "ws" + strings.TrimPrefix(*coreURL, "http")

obslog.Init(obslog.Config{Format: "text", Level: slog.LevelInfo, Out: rc.stderr})
ctx, cancel := daemonize.NotifyContext(parent)
defer cancel()

env, err := agent.ManifestEnvironment(agentHostManifest, claudesdk.Installation(), codex.Installation(), mcode.Installation())
if err != nil {
return fmt.Errorf("agent-host: %w", err)
}
for name, value := range env {
if err := os.Setenv(name, value); err != nil {
return fmt.Errorf("agent-host: %w", err)
}
}
harnesses := agent.NewRegistry()
for _, declaration := range declarations {
runtime := declaration.Discover(ctx, agent.DiscoveryOptions{Profile: paths.DefaultProfile, Stdout: rc.stdout, Stderr: rc.stderr}, declaration.Info)
if runtime == nil || runtime.View == nil {
continue
}
harnesses.RegisterKind(runtime.Info, declaration.Configuration)
harnesses.RegisterView(runtime.Info.Kind, *runtime.View)
}

// With a private cgroup namespace, cgroup v2 mounted again is the
// container's own cgroup, writable unlike Docker's mount of it.
views := filepath.Join(agentHostCgroup, "views")
if err := os.MkdirAll(agentHostCgroup, 0o755); err != nil {
return fmt.Errorf("agent-host: %w", err)
}
if err := unix.Mount("cgroup2", agentHostCgroup, "cgroup2", 0, ""); err != nil {
return fmt.Errorf("%w: agent-host: mount cgroup v2: %w", agenthost.ErrUnsupported, err)
}
if err := os.Mkdir(views, 0o755); err != nil && !errors.Is(err, os.ErrExist) {
return fmt.Errorf("%w: agent-host: view cgroups: %w", agenthost.ErrUnsupported, err)
}
host, err := agenthost.Open(agenthost.Config{StateDir: agentHostState, UIDs: agentHostUIDs, ViewCgroups: views,
RelayURL: wsOrigin + "/api/v1/sandbox-link", RuntimeID: sandboxwire.ID(runtimeID), Credential: []byte(identity.Credential),
Harnesses: harnesses, Shim: agentHostShim, CADir: agentHostCADir, Log: obslog.Bg()})
if err != nil {
return fmt.Errorf("agent-host: %w", err)
}
defer host.Close()

// Bootstrap on each dial, so a Core that is still starting is retried
// with the connection's backoff. The agent host dials Core's origin, not
// the bootstrap's public ws_url.
wsURL := wsOrigin + "/api/v1/agent-daemon/ws"
var boot *transport.BootstrapResponse
dial := func(ctx context.Context) (*transport.Conn, error) {
b, err := transport.Bootstrap(ctx, *coreURL+"/api/v1", identity.RuntimeID, identity.Credential, Version)
if err != nil {
return nil, err
}
boot = b
return transport.Dial(ctx, transport.DialOptions{WSURL: wsURL, DeviceID: identity.RuntimeID, Credential: identity.Credential, DaemonVersion: proto.Version})
}
cfg := dispatch.Config{Registry: host.Registry(), Environments: host.Environments, RemoveHome: host.RemoveHome}
return serveConnections(ctx, wsURL, dial, agentHostUnreachable, func(conn *transport.Conn) error {
return pumpConn(ctx, conn, cfg, boot)
})
}
Loading
Loading