diff --git a/apps/daemon/internal/agenthost/agenthost_linux_test.go b/apps/daemon/internal/agenthost/agenthost_linux_test.go index a1c92a984..05c3167e9 100644 --- a/apps/daemon/internal/agenthost/agenthost_linux_test.go +++ b/apps/daemon/internal/agenthost/agenthost_linux_test.go @@ -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() }) diff --git a/apps/daemon/internal/agenthost/doc.go b/apps/daemon/internal/agenthost/doc.go index 90599726e..b409761cf 100644 --- a/apps/daemon/internal/agenthost/doc.go +++ b/apps/daemon/internal/agenthost/doc.go @@ -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 diff --git a/apps/daemon/internal/agenthost/environment_linux_test.go b/apps/daemon/internal/agenthost/environment_linux_test.go index de271aa69..83fbbee45 100644 --- a/apps/daemon/internal/agenthost/environment_linux_test.go +++ b/apps/daemon/internal/agenthost/environment_linux_test.go @@ -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. diff --git a/apps/daemon/internal/agenthost/host_linux_test.go b/apps/daemon/internal/agenthost/host_linux_test.go index c7ef7ddf6..4df5d2f1a 100644 --- a/apps/daemon/internal/agenthost/host_linux_test.go +++ b/apps/daemon/internal/agenthost/host_linux_test.go @@ -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" @@ -193,7 +195,7 @@ 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) @@ -201,7 +203,7 @@ func TestRemoveHome(t *testing.T) { 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 { @@ -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) } } diff --git a/apps/daemon/internal/agenthost/host_other.go b/apps/daemon/internal/agenthost/host_other.go index 2a59377f8..63650e538 100644 --- a/apps/daemon/internal/agenthost/host_other.go +++ b/apps/daemon/internal/agenthost/host_other.go @@ -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. @@ -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) } diff --git a/apps/daemon/internal/agenthost/sessiondir_linux.go b/apps/daemon/internal/agenthost/sessiondir_linux.go index e827f8f4a..65d069c17 100644 --- a/apps/daemon/internal/agenthost/sessiondir_linux.go +++ b/apps/daemon/internal/agenthost/sessiondir_linux.go @@ -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) } diff --git a/apps/daemon/internal/agenthost/view_linux_test.go b/apps/daemon/internal/agenthost/view_linux_test.go index a8cb478d7..d3dfc3945 100644 --- a/apps/daemon/internal/agenthost/view_linux_test.go +++ b/apps/daemon/internal/agenthost/view_linux_test.go @@ -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 { diff --git a/apps/daemon/internal/cli/agent_host_linux.go b/apps/daemon/internal/cli/agent_host_linux.go new file mode 100644 index 000000000..3abc06368 --- /dev/null +++ b/apps/daemon/internal/cli/agent_host_linux.go @@ -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) + }) +} diff --git a/apps/daemon/internal/cli/agent_host_linux_test.go b/apps/daemon/internal/cli/agent_host_linux_test.go new file mode 100644 index 000000000..32fd13362 --- /dev/null +++ b/apps/daemon/internal/cli/agent_host_linux_test.go @@ -0,0 +1,150 @@ +//go:build linux + +package cli + +import ( + "context" + "encoding/json" + "encoding/pem" + "errors" + "io" + "net/http" + "net/http/httptest" + "os" + "path/filepath" + "slices" + "strings" + "testing" + "time" + + "github.com/google/uuid" + "github.com/gorilla/websocket" + "golang.org/x/sys/unix" + + "github.com/MiniMax-AI/OpenAgentCore/apps/daemon/internal/agent" + "github.com/MiniMax-AI/OpenAgentCore/apps/daemon/internal/transport" + "github.com/MiniMax-AI/OpenAgentCore/internal/agentdaemon/proto" + "github.com/MiniMax-AI/OpenAgentCore/internal/agentdaemon/proto/prototest" +) + +// TestAgentHostReportsItsDeclarations runs the agent host as its container +// does, against a Core peer, with an image that installs no Harness: once +// without declarations and once with two, of which one has a view. The first +// heartbeat declares exactly the kinds with a view, and home removal. +func TestAgentHostReportsItsDeclarations(t *testing.T) { + if os.Getenv("OAC_TEST_AGENTHOST") != "1" { + t.Skip("set OAC_TEST_AGENTHOST=1 and run the test as root in a throwaway container with the agent-host container's flags") + } + issuer := httptest.NewTLSServer(nil) + issuer.Close() + for path, content := range map[string][]byte{ + agentHostManifest: []byte(`{"node": "/usr/local/bin/node", "harnesses": {}}`), + filepath.Join(agentHostCADir, "ca.crt"): pem.EncodeToMemory(&pem.Block{Type: "CERTIFICATE", Bytes: issuer.Certificate().Raw}), + } { + if err := os.MkdirAll(filepath.Dir(path), 0o755); err != nil { + t.Fatal(err) + } + if err := os.WriteFile(path, content, 0o644); err != nil { + t.Fatal(err) + } + } + // The agent host presents the credential as written; it is not decoded. + runtimeID, credential := uuid.NewString(), "c2VjcmV0K/8=" + identity := filepath.Join(t.TempDir(), "identity.json") + if err := os.WriteFile(identity, []byte(`{"runtime_id": "`+runtimeID+`", "credential": "`+credential+`"}`), 0o600); err != nil { + t.Fatal(err) + } + + heartbeats := make(chan proto.HeartbeatPayload, 1) + core := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + if r.Header.Get("Authorization") != "Bearer "+credential { + http.Error(w, "credential", http.StatusUnauthorized) + return + } + switch r.URL.Path { + case "/api/v1/agent-daemon/bootstrap": + _ = json.NewEncoder(w).Encode(transport.BootstrapResponse{DeviceID: runtimeID, HeartbeatSeconds: 60}) + case "/api/v1/agent-daemon/ws": + if r.URL.Query().Get("device_id") != runtimeID { + http.Error(w, "device", http.StatusUnauthorized) + return + } + peer, err := (&websocket.Upgrader{}).Upgrade(w, r, nil) + if err != nil { + return + } + defer peer.Close() + for { + var env proto.Envelope + if err := peer.ReadJSON(&env); err != nil { + return + } + var heartbeat proto.HeartbeatPayload + if env.Type == proto.TypeHeartbeat && env.DecodePayload(&heartbeat) == nil { + select { + case heartbeats <- heartbeat: + default: + } + } + } + default: + http.NotFound(w, r) + } + })) + defer core.Close() + + declare := func(kind string, view *agent.View) agent.Declaration { + info := proto.SupportedAgentKind{Kind: kind, Available: true, Capabilities: prototest.Capabilities(proto.AgentKindCapabilities{ + LocalEnvironment: proto.CapabilitySupported, EnvironmentNone: proto.CapabilitySupported})} + return agent.Declaration{Info: info, Configuration: prototest.ModelConfiguration(), + Discover: func(context.Context, agent.DiscoveryOptions, proto.SupportedAgentKind) *agent.Runtime { + return &agent.Runtime{Info: info, View: view} + }} + } + view := &agent.View{Proxy: agent.ViewProxyEnv, Executor: func(context.Context, proto.PromptRequestPayload, agent.ViewSession) (agent.Executor, error) { + return nil, errors.New("no Executor") + }} + // The image may install no Harness: the agent host still connects and + // declares no kind. + for _, c := range []struct { + declarations []agent.Declaration + kinds []string + }{{nil, nil}, {[]agent.Declaration{declare("viewed", view), declare("unviewed", nil)}, []string{"viewed"}}} { + ctx, cancel := context.WithCancel(t.Context()) + served := make(chan error, 1) + go func() { + rc := &runContext{stdin: strings.NewReader(""), stdout: io.Discard, stderr: os.Stderr} + served <- serveAgentHost(ctx, rc, []string{"--identity-file", identity, "--core-url", core.URL}, c.declarations) + }() + select { + case heartbeat := <-heartbeats: + var kinds []string + for _, kind := range heartbeat.SupportedAgentKinds { + kinds = append(kinds, kind.Kind) + } + if !slices.Equal(kinds, c.kinds) || heartbeat.HomeRemoval != proto.CapabilitySupported { + t.Errorf("heartbeat declares %q with home removal %q, want %q with home removal", kinds, heartbeat.HomeRemoval, c.kinds) + } + case err := <-served: + t.Fatalf("agent host stopped before its first heartbeat: %v", err) + case <-time.After(30 * time.Second): + t.Fatal("no heartbeat") + } + cancel() + if err := <-served; err != nil { + t.Fatalf("agent host stopped with %v, want nil after its signal", err) + } + if err := unix.Unmount(agentHostCgroup, 0); err != nil { + t.Fatal(err) + } + } +} + +// TestAgentHostExitsWhileCoreStaysUnreachable checks the bound after which +// the agent host exits for its supervisor to restart it. +func TestAgentHostExitsWhileCoreStaysUnreachable(t *testing.T) { + dial := func(context.Context) (*transport.Conn, error) { return nil, errors.New("connection refused") } + if err := serveConnections(t.Context(), "ws://127.0.0.1:1", dial, 10*time.Millisecond, nil); !errors.Is(err, context.DeadlineExceeded) { + t.Fatalf("serveConnections = %v, want the unreachable bound", err) + } +} diff --git a/apps/daemon/internal/cli/agent_host_other.go b/apps/daemon/internal/cli/agent_host_other.go new file mode 100644 index 000000000..0b75b8eb4 --- /dev/null +++ b/apps/daemon/internal/cli/agent_host_other.go @@ -0,0 +1,14 @@ +//go:build !linux + +package cli + +import ( + "fmt" + + "github.com/MiniMax-AI/OpenAgentCore/apps/daemon/internal/agenthost" +) + +// runAgentHost reports that the agent host needs Linux. +func runAgentHost(*runContext, []string) error { + return fmt.Errorf("%w: agent-host", agenthost.ErrUnsupported) +} diff --git a/apps/daemon/internal/cli/connect.go b/apps/daemon/internal/cli/connect.go index 157566dde..70bb9ef1b 100644 --- a/apps/daemon/internal/cli/connect.go +++ b/apps/daemon/internal/cli/connect.go @@ -263,12 +263,25 @@ func mainLoopRemote(parent context.Context, rc *runContext, profile string, prof if control != nil { return runSuspendLoop(rootCtx, dial, registry, local, boot, agentCLIs, control) } + return serveConnections(rootCtx, wsURL, dial, 0, func(conn *transport.Conn) error { + return pumpConn(rootCtx, conn, dispatch.Config{Registry: registry, Environments: localEnvironments(local)}, boot) + }) +} + +// serveConnections dials Core with dial and serves each connection until ctx +// ends or Core rejects the credential. With a positive unreachable bound it +// fails once Core has stayed unreachable that long. +func serveConnections(ctx context.Context, wsURL string, dial transport.DialFn, unreachable time.Duration, serve func(*transport.Conn) error) error { for { - if err := rootCtx.Err(); err != nil { + if err := ctx.Err(); err != nil { return nil } - conn, err := transport.Reconnect(rootCtx, dial, transport.DefaultBackoff, func(attempt int, lastDelay time.Duration, lastErr error) { + dialCtx, stop := ctx, context.CancelFunc(func() {}) + if unreachable > 0 { + dialCtx, stop = context.WithTimeout(ctx, unreachable) + } + conn, err := transport.Reconnect(dialCtx, dial, transport.DefaultBackoff, func(attempt int, lastDelay time.Duration, lastErr error) { switch { case attempt == 1: obslog.Bg().Info("connecting", "ws_url", wsURL) @@ -281,21 +294,24 @@ func mainLoopRemote(parent context.Context, rc *runContext, profile string, prof obslog.Bg().Warn("dial retry", "attempt", attempt, "delay", lastDelay) } }) + stop() if err != nil { - if errors.Is(err, context.Canceled) { + if ctx.Err() != nil { return nil } if errors.Is(err, transport.ErrPermanent) { return fmt.Errorf("connect: permanent error (reissue the daemon credential): %w", err) } + if errors.Is(err, context.DeadlineExceeded) { + return fmt.Errorf("connect: Core unreachable for %s: %w", unreachable, err) + } return fmt.Errorf("connect: dial: %w", err) } obslog.Bg().Info("ws connected", "device_id", conn.DeviceID()) - // pumpConn returns on conn close (peer hangup, transport - // error, root ctx cancel). Loop back into Reconnect unless - // root ctx is cancelled. - pumpErr := pumpConn(rootCtx, conn, registry, local, boot, agentCLIs) + // serve returns on conn close (peer hangup, transport error, ctx + // cancel). Loop back into Reconnect unless ctx is cancelled. + pumpErr := serve(conn) if pumpErr != nil { obslog.Bg().Warn("ws session ended", "err", pumpErr) } else { @@ -305,7 +321,7 @@ func mainLoopRemote(parent context.Context, rc *runContext, profile string, prof // Server-initiated clean close (e.g. shutdown) → exit; // otherwise loop back and reconnect. - if rootCtx.Err() != nil { + if ctx.Err() != nil { return nil } // Permanent error (e.g. runtime deleted) → exit instead of @@ -315,7 +331,7 @@ func mainLoopRemote(parent context.Context, rc *runContext, profile string, prof } // Small breather before redialing so a flapping server doesn't // get a tight loop of upgrade requests. - _ = transport.Sleep(rootCtx, 1*time.Second) + _ = transport.Sleep(ctx, 1*time.Second) } } @@ -328,17 +344,14 @@ func localEnvironments(local *localworkspace.Binding) func(proto.AssignmentRef, return local.Resolve } -// pumpConn runs the per-connection workload: a dispatch.Router fed by -// conn.Recv(), heartbeats every boot.HeartbeatInterval(), and a -// confirmed router.Shutdown before returning ownership to the reconnect loop. -// Failed cleanup keeps this exact Router alive, including after a shutdown signal. -func pumpConn(parentCtx context.Context, conn *transport.Conn, registry *agent.Registry, local *localworkspace.Binding, boot *transport.BootstrapResponse, agentCLIs agentCLIDiscovery) error { - router, err := dispatch.New(dispatch.Config{ - Registry: registry, - Sender: conn, - Log: obslog.Bg(), - Environments: localEnvironments(local), - }) +// pumpConn runs the per-connection workload: a dispatch.Router of cfg's +// Harness kinds and Environment owners fed by conn.Recv(), heartbeats every +// boot.HeartbeatInterval(), and a confirmed router.Shutdown before returning +// ownership to the reconnect loop. Failed cleanup keeps this exact Router +// alive, including after a shutdown signal. +func pumpConn(parentCtx context.Context, conn *transport.Conn, cfg dispatch.Config, boot *transport.BootstrapResponse) error { + cfg.Sender, cfg.Log = conn, obslog.Bg() + router, err := dispatch.New(cfg) if err != nil { return fmt.Errorf("router init: %w", err) } @@ -352,8 +365,8 @@ func pumpConn(parentCtx context.Context, conn *transport.Conn, registry *agent.R Timestamp: time.Now().Unix(), ActiveRequests: router.ActiveRuns(), DaemonVersion: Version, - SupportedAgentKinds: registry.SupportedAgentKinds(), - HomeRemoval: proto.CapabilityUnsupported, + SupportedAgentKinds: cfg.Registry.SupportedAgentKinds(), + HomeRemoval: proto.CapabilityFromBool(cfg.RemoveHome != nil), } }, obslog.Bg().With("component", "heartbeat")) diff --git a/apps/daemon/internal/cli/connect_cleanup_test.go b/apps/daemon/internal/cli/connect_cleanup_test.go index d707804f6..bd1c1d4e7 100644 --- a/apps/daemon/internal/cli/connect_cleanup_test.go +++ b/apps/daemon/internal/cli/connect_cleanup_test.go @@ -9,6 +9,7 @@ import ( "strings" "github.com/MiniMax-AI/OpenAgentCore/apps/daemon/internal/agent" + "github.com/MiniMax-AI/OpenAgentCore/apps/daemon/internal/dispatch" "github.com/MiniMax-AI/OpenAgentCore/apps/daemon/internal/transport" "github.com/MiniMax-AI/OpenAgentCore/internal/agentdaemon/proto" "github.com/MiniMax-AI/OpenAgentCore/internal/agentdaemon/proto/prototest" @@ -133,7 +134,7 @@ func testDisconnectedPumpCleanup(t *testing.T, suspend bool) { return } defer conn.Close() - finished <- pumpConn(ctx, conn, registry, nil, boot, agentCLIDiscovery{}) + finished <- pumpConn(ctx, conn, dispatch.Config{Registry: registry}, boot) }() var peer *websocket.Conn select { diff --git a/apps/daemon/internal/cli/root.go b/apps/daemon/internal/cli/root.go index cd77a7d2b..2b3323e33 100644 --- a/apps/daemon/internal/cli/root.go +++ b/apps/daemon/internal/cli/root.go @@ -41,6 +41,7 @@ var commands = []command{ {name: "runtime-mcp-exec", summary: "Execute an installed MCP server", run: runRuntimeMCP}, {name: "placement", summary: "Enroll or retire an explicitly managed local execution placement", run: runPlacement}, {name: "connect", summary: "Open the reverse WebSocket and start serving prompts", run: runConnect}, + {name: "agent-host", summary: "Serve Sessions from the agent-host container", run: runAgentHost}, {name: "status", summary: "Print the credential profile and daemon state", run: runStatus}, {name: "stop", summary: "Stop a background `connect -b` daemon", run: runStop}, {name: "logs", summary: "Tail the background daemon's log file", run: runLogs}, diff --git a/apps/daemon/internal/cli/root_test.go b/apps/daemon/internal/cli/root_test.go index c27176d35..b235f3d2f 100644 --- a/apps/daemon/internal/cli/root_test.go +++ b/apps/daemon/internal/cli/root_test.go @@ -66,6 +66,7 @@ func TestSubcommandsAreRegistered(t *testing.T) { "runtime-mcp-exec": false, "placement": false, "connect": false, + "agent-host": false, "status": false, "stop": false, "logs": false, diff --git a/apps/daemon/internal/dispatch/assignment.go b/apps/daemon/internal/dispatch/assignment.go index edebffd51..3cecade55 100644 --- a/apps/daemon/internal/dispatch/assignment.go +++ b/apps/daemon/internal/dispatch/assignment.go @@ -22,12 +22,16 @@ type assignmentState struct { // environment is the owner resolved from the bind, or nil. environment Environment released bool - // work counts the Session's admitted reads, writes, exports and Runtime - // preparations until each has sent its terminal result. A release waits - // for it, and the released assignment admits no more. - work sync.WaitGroup - // cleanup serializes release cleanups, so a retried release never closes - // the owner while an earlier Close runs. + // superseding is the bind that superseded the assignment at ref, held + // released until the earlier epoch's work and owner have settled. + superseding *proto.AssignmentBindPayload + // work counts the Session's admitted reads, writes, exports, Runtime + // preparations, Run terminals and cancellation receipts until each has + // sent its terminal result. A release, a superseding bind and a quiesce + // wait for it, and a released assignment admits no more. + work dispatchWork + // cleanup serializes release and supersede cleanups, so a retry never + // closes the owner while an earlier Close runs. cleanup sync.Mutex } @@ -83,34 +87,122 @@ func (r *Router) handleAssignmentBind(ctx context.Context, env proto.Envelope) e var input proto.AssignmentBindPayload ref, code := env.Assignment, "" if env.ID == "" || env.DecodeRequest(&input) != nil || input.Validate() != nil || !ref.Valid() { - code = "invalid_request" - } else { - var resource sandboxbootstrap.Resource - if input.Resource != nil { - resource = *input.Resource - } - r.mu.Lock() - a := r.assignments[ref.SessionID] - switch { - case a == nil: - var environment Environment - if r.environments != nil { - if environment = r.environments(ref, input); environment == nil { - code = proto.AssignmentConflict - break - } + return r.reply(ctx, env, proto.TypeAssignmentStatus, assignmentStatus(proto.AssignmentBound, "invalid_request")) + } + var resource sandboxbootstrap.Resource + if input.Resource != nil { + resource = *input.Resource + } + r.mu.Lock() + if r.suspensions[input.EnvironmentID] != nil { + r.mu.Unlock() + return ErrRouterQuiesced + } + a := r.assignments[ref.SessionID] + switch { + case a == nil: + var environment Environment + if r.environments != nil { + if environment = r.environments(ref, input); environment == nil { + code = proto.AssignmentConflict + break } - r.assignments[ref.SessionID] = &assignmentState{ref: ref, environmentID: input.EnvironmentID, resource: resource, grant: input.AttachGrant, environment: environment} - case a.ref.AssignmentID == ref.AssignmentID && (ref.Epoch < a.ref.Epoch || ref.Epoch == a.ref.Epoch && a.released): - code = proto.AssignmentStale - case a.ref != ref || a.environmentID != input.EnvironmentID || a.resource != resource || !bytes.Equal(a.grant, input.AttachGrant): + } + r.assignments[ref.SessionID] = &assignmentState{ref: ref, environmentID: input.EnvironmentID, resource: resource, grant: input.AttachGrant, environment: environment} + case a.ref.AssignmentID != ref.AssignmentID: + code = proto.AssignmentConflict + case ref.Epoch > a.ref.Epoch && (!a.released || a.superseding != nil): + // The bind supersedes the earlier epoch, or a pending supersede. A + // Session's Environment never changes, so another one is refused + // before anything is fenced. + if input.EnvironmentID != a.environmentID { code = proto.AssignmentConflict + break } + preparations := r.fenceSessionWorkLocked(ref.SessionID) + a.ref, a.released, a.superseding = ref, true, &input + r.shutdownWG.Add(1) r.mu.Unlock() + go r.supersede(ctx, env, a, preparations) + return nil + case ref.Epoch == a.ref.Epoch && a.superseding != nil: + if !sameBind(*a.superseding, input) { + code = proto.AssignmentConflict + break + } + // A retry repeats the pending supersede. + r.shutdownWG.Add(1) + r.mu.Unlock() + go r.supersede(ctx, env, a, nil) + return nil + case ref.Epoch < a.ref.Epoch || ref.Epoch == a.ref.Epoch && a.released: + code = proto.AssignmentStale + case a.ref != ref || a.environmentID != input.EnvironmentID || a.resource != resource || !bytes.Equal(a.grant, input.AttachGrant): + code = proto.AssignmentConflict } + r.mu.Unlock() return r.reply(ctx, env, proto.TypeAssignmentStatus, assignmentStatus(proto.AssignmentBound, code)) } +// supersede settles the earlier epoch's work, closes the Session's Executor +// and owner, and then binds the pending supersede at env's assignment. It +// replies assignment_stale when a release or a later bind took over meanwhile. +// A failed cleanup keeps the supersede pending, so a retry repeats it. +func (r *Router) supersede(ctx context.Context, env proto.Envelope, a *assignmentState, preparations []*preparationState) { + defer r.shutdownWG.Done() + ref := env.Assignment + cleanupCtx, stop := r.shutdownContext(context.WithoutCancel(ctx)) + defer stop() + for _, p := range preparations { + r.releasePreparation(p, "failed", proto.AssignmentStale, true) + } + a.work.Wait() + a.cleanup.Lock() + defer a.cleanup.Unlock() + r.mu.Lock() + pending, environment := a.ref == ref && a.superseding != nil, a.environment + r.mu.Unlock() + var err error + if pending { + err = r.closeSessionExecutor(ref.SessionID) + } + if pending && err == nil && environment != nil { + err = environment.Close(cleanupCtx) + } + code := "" + r.mu.Lock() + switch { + case a.ref != ref || a.released && a.superseding == nil: + code = proto.AssignmentStale + case a.superseding == nil: + // An earlier attempt bound it. + case err != nil: + r.log.Warn("assignment supersede cleanup unconfirmed", "session_id", ref.SessionID, "err", err) + code = proto.CleanupUnconfirmed + default: + input := *a.superseding + var owner Environment + if r.environments != nil { + if owner = r.environments(ref, input); owner == nil { + code = proto.AssignmentConflict + break + } + } + var resource sandboxbootstrap.Resource + if input.Resource != nil { + resource = *input.Resource + } + a.resource, a.grant, a.environment = resource, input.AttachGrant, owner + a.released, a.superseding = false, nil + } + r.mu.Unlock() + _ = r.reply(cleanupCtx, env, proto.TypeAssignmentStatus, assignmentStatus(proto.AssignmentBound, code)) +} + +func sameBind(a, b proto.AssignmentBindPayload) bool { + return a.EnvironmentID == b.EnvironmentID && (a.Resource == nil) == (b.Resource == nil) && (a.Resource == nil || *a.Resource == *b.Resource) && bytes.Equal(a.AttachGrant, b.AttachGrant) +} + // handleAssignmentRelease fences the assignment, then settles the Session's // work and Executor, closes its Environment owner and removes its home before // it replies. A retry at the same epoch repeats the cleanup. @@ -135,7 +227,7 @@ func (r *Router) handleAssignmentRelease(ctx context.Context, env proto.Envelope case ref.Epoch < a.ref.Epoch: code = proto.AssignmentStale default: - a.ref, a.released = ref, true + a.ref, a.released, a.superseding = ref, true, nil } if code != "" { r.mu.Unlock() @@ -179,11 +271,11 @@ func (r *Router) handleAssignmentRelease(ctx context.Context, env proto.Envelope // the release releases, which also cancels their exports. Router.mu must be // held. func (r *Router) fenceSessionWorkLocked(sessionID string) []*preparationState { - if u := r.workspaceWrite; u != nil && u.envelope.Assignment.SessionID == sessionID && !u.finished { + if u := r.workspaceWrites[sessionID]; u != nil && !u.finished { u.finished = true close(u.ready) } - if u := r.runtimePreparation; u != nil && u.envelope.Assignment.SessionID == sessionID && !u.finished { + if u := r.runtimePreparations[sessionID]; u != nil && !u.finished { r.finishRuntimePreparationTransferLocked(u, false) } var preparations []*preparationState diff --git a/apps/daemon/internal/dispatch/assignment_test.go b/apps/daemon/internal/dispatch/assignment_test.go index 930348dae..2aed11991 100644 --- a/apps/daemon/internal/dispatch/assignment_test.go +++ b/apps/daemon/internal/dispatch/assignment_test.go @@ -262,3 +262,61 @@ func TestReleaseRetryAndShutdownCloseOwnersOnce(t *testing.T) { t.Fatalf("released owner closes = %d (overlapped %t), unreleased owner closes = %d", released.closes.Load(), released.overlapped.Load(), unreleased.closes.Load()) } } + +func TestSupersedingBindFencesTheEarlierEpoch(t *testing.T) { + h := newHarness(t) + defer h.router.Shutdown(context.Background()) + startRun(t, h.router, h.sender, "fake_alpha", "s") + sess := <-h.gotSess + bind := func(id string, epoch uint64, payload proto.AssignmentBindPayload) proto.AssignmentStatusPayload { + env := scoped(t, "s", proto.TypeAssignmentBind, id, payload) + env.Assignment.Epoch = epoch + if err := h.router.Handle(t.Context(), env); err != nil { + t.Fatal(err) + } + return waitAssignmentStatus(t, h.sender, id) + } + if got := bind("supersede", 2, proto.AssignmentBindPayload{}); got.State != proto.AssignmentBound || got.ErrorCode != "" { + t.Fatalf("superseding bind = %+v", got) + } + // The earlier epoch's Run ended before the bind replied. + terminal, bound := -1, -1 + for i, frame := range h.sender.snapshot() { + switch { + case frame.ID == "s" && (frame.Type == proto.TypeDone || frame.Type == proto.TypeError): + terminal = i + case frame.ID == "supersede": + bound = i + } + } + if terminal < 0 || terminal > bound || sess.cancels() != 1 { + t.Fatalf("Run terminal at %d, bind reply at %d, cancels = %d", terminal, bound, sess.cancels()) + } + for id, test := range map[string]struct { + epoch uint64 + payload proto.AssignmentBindPayload + code string + }{ + "lower": {1, proto.AssignmentBindPayload{}, proto.AssignmentStale}, + "changed": {2, proto.AssignmentBindPayload{EnvironmentID: uuid.NewString()}, proto.AssignmentConflict}, + "other environment": {3, proto.AssignmentBindPayload{EnvironmentID: uuid.NewString()}, proto.AssignmentConflict}, + "repeated": {2, proto.AssignmentBindPayload{}, ""}, + } { + if got := bind(id, test.epoch, test.payload); got.ErrorCode != test.code { + t.Fatalf("%s bind = %+v", id, got) + } + } + stale := scoped(t, "s", proto.TypeExecutionPrepare, "stale", noEnvironmentPreparation("s", proto.PromptRequestPayload{AgentKind: "fake_alpha"})) + if err := h.router.Handle(t.Context(), stale); err == nil { + t.Fatal("the superseded epoch admitted a preparation") + } + if got := waitPreparationStatus(t, h.sender, "stale", "rejected", ""); got.ErrorCode != proto.AssignmentStale { + t.Fatalf("superseded preparation = %+v", got) + } + current := stale + current.ID, current.Assignment.Epoch = "current", 2 + if err := h.router.Handle(t.Context(), current); err != nil { + t.Fatal(err) + } + waitPreparationStatus(t, h.sender, "current", "ready", "") +} diff --git a/apps/daemon/internal/dispatch/cancellation.go b/apps/daemon/internal/dispatch/cancellation.go index c54e9eb45..adcc02015 100644 --- a/apps/daemon/internal/dispatch/cancellation.go +++ b/apps/daemon/internal/dispatch/cancellation.go @@ -30,8 +30,12 @@ func (r *Router) handlePromptCancel(ctx context.Context, env proto.Envelope) err handoff := state.preparedHandoff release, attempt := r.claimPreparedReleaseLocked(state, true, "", true) if request.DeliveryID != "" { + done := r.trackWorkLocked(env.Assignment) r.shutdownWG.Add(1) - go r.sendPreparedCancellation(state, handoff, release, attempt, env, request.DeliveryID) + go func() { + defer done() + r.sendPreparedCancellation(state, handoff, release, attempt, env, request.DeliveryID) + }() } r.mu.Unlock() return nil diff --git a/apps/daemon/internal/dispatch/executor.go b/apps/daemon/internal/dispatch/executor.go index 72b7c7e56..bc75c1f69 100644 --- a/apps/daemon/internal/dispatch/executor.go +++ b/apps/daemon/internal/dispatch/executor.go @@ -83,7 +83,7 @@ func (r *Router) handleExecutorPrepare(ctx context.Context, env proto.Envelope, requestFingerprint := sha256.Sum256(encoded) req.Assignment = env.Assignment r.mu.Lock() - if r.closed || r.suspension != nil { + if r.closed { r.mu.Unlock() return ErrRouterClosed } @@ -208,7 +208,7 @@ func (r *Router) prepareExecutor(p *preparationState, req proto.PromptRequestPay owner.native, owner.preparing = native, false close(owner.prepared) r.log.Info("executor native_prepare", "executor_id", owner.id, "session_id", owner.sessionID, "duration_ms", time.Since(started).Milliseconds(), "success", err == nil && native != nil) - ready := native != nil && err == nil && !owner.invalid && !r.closed && r.suspension == nil && p.status.State == "preparing" && p.ctx.Err() == nil + ready := native != nil && err == nil && !owner.invalid && !r.closed && p.status.State == "preparing" && p.ctx.Err() == nil p.busy = false if ready { p.status.State, p.status.Revision = "ready", p.status.Revision+1 diff --git a/apps/daemon/internal/dispatch/native_file_results_test.go b/apps/daemon/internal/dispatch/native_file_results_test.go index a3d15e8bd..4f5bc685a 100644 --- a/apps/daemon/internal/dispatch/native_file_results_test.go +++ b/apps/daemon/internal/dispatch/native_file_results_test.go @@ -37,9 +37,10 @@ func TestLocalUploadUnknownRetainsOwner(t *testing.T) { if got.Outcome != "unknown" { t.Fatal(got) } - r.workspaceWrite = &workspaceUpload{envelope: proto.Envelope{ID: uuid.NewString()}, finished: true, uncertain: true} 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) } diff --git a/apps/daemon/internal/dispatch/preparation.go b/apps/daemon/internal/dispatch/preparation.go index 5f31485a5..76ca2ab2b 100644 --- a/apps/daemon/internal/dispatch/preparation.go +++ b/apps/daemon/internal/dispatch/preparation.go @@ -83,7 +83,7 @@ func (r *Router) handleExecutionPrepare(ctx context.Context, env proto.Envelope) r.mu.Unlock() return r.rejectPreparation(env, code) } - if u := r.runtimePreparation; u != nil && u.envelope.Assignment.SessionID == input.SessionID { + if r.runtimePreparations[input.SessionID] != nil { r.mu.Unlock() return r.rejectPreparation(env, "resource_unavailable") } @@ -168,7 +168,7 @@ func (r *Router) handleExecutionRelease(_ context.Context, env proto.Envelope) e func (r *Router) releasePreparation(p *preparationState, state, code string, publish bool) { r.mu.Lock() - if r.closed || r.suspension != nil { + if r.closed { r.mu.Unlock() return } @@ -221,7 +221,7 @@ func (r *Router) publishPreparation(p *preparationState, status proto.Preparatio r.mu.Unlock() return } - if r.closed || r.suspension != nil { + if r.closed { r.mu.Unlock() return } diff --git a/apps/daemon/internal/dispatch/preparation_start.go b/apps/daemon/internal/dispatch/preparation_start.go index 557234eca..c1e268d98 100644 --- a/apps/daemon/internal/dispatch/preparation_start.go +++ b/apps/daemon/internal/dispatch/preparation_start.go @@ -20,11 +20,11 @@ func (r *Router) handleExecutionStart(_ context.Context, env proto.Envelope) err encoded, _ := json.Marshal(input) fingerprint := sha256.Sum256(encoded) r.mu.Lock() - if r.closed || r.suspension != nil { + if r.closed { r.mu.Unlock() return ErrRouterClosed } - if r.workspaceWrite != nil || r.workspaceExport != nil || r.runtimePreparation != nil { + if r.environmentTransferLocked(env.Assignment.SessionID) { r.mu.Unlock() return r.rejectPreparation(env, "resource_unavailable") } diff --git a/apps/daemon/internal/dispatch/prepared_handoff.go b/apps/daemon/internal/dispatch/prepared_handoff.go index 772cce4a5..e5fdf3856 100644 --- a/apps/daemon/internal/dispatch/prepared_handoff.go +++ b/apps/daemon/internal/dispatch/prepared_handoff.go @@ -236,6 +236,8 @@ func (r *Router) runPreparedRelease(state *sessionState, handoff *preparedHandof r.mu.Unlock() r.mu.Lock() outputErr, terminal, closed := handoff.outputErr, handoff.terminal, r.closed + // The retired Run's terminal is the Session's work until it is sent. + tail := r.trackWorkLocked(state.assignment) r.mu.Unlock() // Done can become visible before Sender.Send returns. Commit the settled // owner and retire this Run before publication so an immediate successor @@ -281,6 +283,7 @@ func (r *Router) runPreparedRelease(state *sessionState, handoff *preparedHandof p.closeErr = terminalErr close(release.settled) r.mu.Unlock() + tail() } func (r *Router) forwardPreparedOutput(state *sessionState) { diff --git a/apps/daemon/internal/dispatch/router.go b/apps/daemon/internal/dispatch/router.go index 138ffe3be..1fe1c9ef2 100644 --- a/apps/daemon/internal/dispatch/router.go +++ b/apps/daemon/internal/dispatch/router.go @@ -33,9 +33,8 @@ type Router struct { log *slog.Logger admission sync.RWMutex - suspension *proto.EnvironmentSuspendPayload - suspendedBy proto.AssignmentRef // the assignment that quiesced mu sync.Mutex + suspensions map[string]*suspension // Environment ID → its quiescence assignments map[string]*assignmentState // SessionID → assignment sessions map[string]*sessionState // RunID → state applied map[string]appliedFunctionResult // RunID and call ID → applied result @@ -48,9 +47,12 @@ type Router struct { preparations map[string]*preparationState preparationRequests map[string]*preparationState preparationTimeout time.Duration - runtimePreparation *runtimePreparationTransfer - workspaceWrite *workspaceUpload - workspaceExport *workspaceExport + // Each Session has at most one transfer of each kind; transferBytes + // counts the bodies that admitted writes and Runtime preparations buffer. + runtimePreparations map[string]*runtimePreparationTransfer // SessionID → + workspaceWrites map[string]*workspaceUpload // SessionID → + workspaceExports map[string]*workspaceExport // SessionID → + transferBytes int workspaceReads map[string]string // read ID → SessionID environments func(proto.AssignmentRef, proto.AssignmentBindPayload) Environment removeHome func(sessionID string) error @@ -91,9 +93,10 @@ type Config struct { IdleTimeout time.Duration PreparationTimeout time.Duration // Environments resolves the Environment owner of a Session's first bind on - // this Router, under the Router's lock, without I/O. A nil owner rejects - // the bind: the Runtime does not serve that Session. Nil Environments - // leaves every Session without an owner. + // this Router, and of a bind that supersedes its assignment, under the + // Router's lock, without I/O. A nil owner rejects the bind: the Runtime + // does not serve that Session. Nil Environments leaves every Session + // without an owner. Environments func(proto.AssignmentRef, proto.AssignmentBindPayload) Environment // RemoveHome removes the Session's native home once its Executors have // closed. Nil declares that assignment_release does not accept RemoveHome. @@ -136,19 +139,34 @@ func New(cfg Config) (*Router, error) { preparations: make(map[string]*preparationState), preparationRequests: make(map[string]*preparationState), preparationTimeout: preparationTimeout, + suspensions: make(map[string]*suspension), + runtimePreparations: make(map[string]*runtimePreparationTransfer), + workspaceWrites: make(map[string]*workspaceUpload), + workspaceExports: make(map[string]*workspaceExport), environments: cfg.Environments, removeHome: cfg.RemoveHome, }, nil } +// transferMemory bounds the bodies that a connection's admitted workspace +// writes and Runtime preparations buffer together. +const transferMemory = proto.WorkspaceWriteMaxBytes + proto.RuntimePrepareMaxBytes + // Handle dispatches one inbound Envelope. Errors are returned for // programmer-visible problems (bad shape, registry miss); transient -// session-level failures are logged and swallowed. A quiesced Router -// still handles assignment_release. +// session-level failures are logged and swallowed. A quiesced Environment's +// Sessions still have their assignment_release handled. // // Adopts env.Trace into ctx so every downstream log under it inherits // the same trace_id, making a single grep cover both sides. func (r *Router) Handle(ctx context.Context, env proto.Envelope) error { + ctx = adoptEnvelopeTrace(ctx, env) + switch env.Type { + case proto.TypeEnvironmentQuiesce: + return r.handleQuiesce(ctx, env) + case proto.TypeEnvironmentResume: + return r.handleResume(ctx, env) + } r.admission.RLock() defer r.admission.RUnlock() r.mu.Lock() @@ -156,14 +174,12 @@ func (r *Router) Handle(ctx context.Context, env proto.Envelope) error { r.mu.Unlock() return ErrRouterClosed } - if r.suspension != nil && env.Type != proto.TypeAssignmentRelease { + if a := r.assignments[env.Assignment.SessionID]; a != nil && r.suspensions[a.environmentID] != nil && env.Type != proto.TypeAssignmentRelease { r.mu.Unlock() return ErrRouterQuiesced } r.mu.Unlock() - ctx = adoptEnvelopeTrace(ctx, env) - switch env.Type { case proto.TypeAssignmentBind: return r.handleAssignmentBind(ctx, env) diff --git a/apps/daemon/internal/dispatch/runtime_preparation.go b/apps/daemon/internal/dispatch/runtime_preparation.go index 8db2ec1f4..f8ca616dd 100644 --- a/apps/daemon/internal/dispatch/runtime_preparation.go +++ b/apps/daemon/internal/dispatch/runtime_preparation.go @@ -14,8 +14,9 @@ import ( const runtimePreparationTimeout = 120 * time.Second -// Router.mu protects one connection-local transfer. Partial installation data -// belongs to the bound Environment and is never removed by transfer cleanup. +// Router.mu protects a Session's Runtime preparation transfer. Partial +// 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 @@ -36,9 +37,10 @@ func (r *Router) handleRuntimePrepare(ctx context.Context, env proto.Envelope) e var request proto.RuntimePreparePayload if len(env.Payload) > proto.RuntimePrepareMaxFrameBytes || env.DecodeRequest(&request) != nil || !proto.ValidRuntimePrepareRequest(request) { r.mu.Lock() - pending := r.runtimePreparation != nil && r.runtimePreparation.envelope.ID == env.ID - if pending && !r.runtimePreparation.finished { - r.finishRuntimePreparationTransferLocked(r.runtimePreparation, false) + u := r.runtimePreparations[env.Assignment.SessionID] + pending := u != nil && u.envelope.ID == env.ID + if pending && !u.finished { + r.finishRuntimePreparationTransferLocked(u, false) } r.mu.Unlock() if pending { @@ -48,13 +50,14 @@ func (r *Router) handleRuntimePrepare(ctx context.Context, env proto.Envelope) e return r.sendRuntimePrepareResult(ctx, env, rejectedRuntimePreparation("invalid_request")) } r.mu.Lock() - if r.closed || r.suspension != nil { + if r.closed { r.mu.Unlock() return ErrRouterClosed } + u := r.runtimePreparations[env.Assignment.SessionID] if request.Step == "begin" { - if r.runtimePreparation != nil { - duplicate := r.runtimePreparation.envelope.ID == env.ID + if u != nil || r.transferBytes > transferMemory-request.SizeBytes { + duplicate := u != nil && u.envelope.ID == env.ID r.mu.Unlock() if duplicate { return errors.New("dispatch: Runtime preparation already admitted") @@ -79,7 +82,8 @@ func (r *Router) handleRuntimePrepare(ctx context.Context, env proto.Envelope) e id: id, envelope: env, request: request, data: make([]byte, 0, request.SizeBytes), ready: make(chan struct{}), cancel: cancel, } - r.runtimePreparation = u + r.runtimePreparations[request.SessionID] = u + r.transferBytes += request.SizeBytes done := r.trackWorkLocked(env.Assignment) r.shutdownWG.Add(1) r.mu.Unlock() @@ -90,7 +94,6 @@ func (r *Router) handleRuntimePrepare(ctx context.Context, env proto.Envelope) e } return nil } - u := r.runtimePreparation if u == nil || u.envelope.ID != env.ID || u.envelope.Assignment != env.Assignment { r.mu.Unlock() return r.sendRuntimePrepareResult(ctx, env, rejectedRuntimePreparation("resource_unavailable")) @@ -141,9 +144,7 @@ func (r *Router) runtimePreparationResourcesBusyLocked(sessionID string) bool { // environmentTransferLocked reports whether the Session has a workspace write, // a workspace export or a Runtime preparation. Router.mu must be held. func (r *Router) environmentTransferLocked(sessionID string) bool { - return r.workspaceWrite != nil && r.workspaceWrite.envelope.Assignment.SessionID == sessionID || - r.workspaceExport != nil && r.workspaceExport.request.Assignment.SessionID == sessionID || - r.runtimePreparation != nil && r.runtimePreparation.envelope.Assignment.SessionID == sessionID + return r.workspaceWrites[sessionID] != nil || r.workspaceExports[sessionID] != nil || r.runtimePreparations[sessionID] != nil } // sessionWorkLocked reports whether the Session has a Run, a workspace read or @@ -202,9 +203,10 @@ func (r *Router) runRuntimePreparationTransfer(ctx context.Context, u *runtimePr // Release the potentially large body before waiting on transport delivery. data = nil r.mu.Lock() + r.transferBytes -= u.request.SizeBytes u.uncertain = result.Outcome == "unknown" - if !u.uncertain && r.runtimePreparation == u { - r.runtimePreparation = nil + if !u.uncertain { + delete(r.runtimePreparations, u.request.SessionID) } r.mu.Unlock() // The result has a separate send budget, independent of an installation timeout. diff --git a/apps/daemon/internal/dispatch/runtime_preparation_test.go b/apps/daemon/internal/dispatch/runtime_preparation_test.go index da997c58f..f0a027017 100644 --- a/apps/daemon/internal/dispatch/runtime_preparation_test.go +++ b/apps/daemon/internal/dispatch/runtime_preparation_test.go @@ -103,7 +103,7 @@ func TestRuntimePreparationTransferValidatesCompleteBodyBeforeMutation(t *testin } capabilitiesReceipt(t, sender, id, "ready") r.mu.Lock() - owner := r.runtimePreparation + owner := r.runtimePreparations[session] r.mu.Unlock() duplicate := uuid.NewString() if err := r.Handle(t.Context(), capabilityEnvelope(t, duplicate, request)); err != nil { @@ -133,7 +133,7 @@ func TestRuntimePreparationTransferValidatesCompleteBodyBeforeMutation(t *testin capabilitiesReceipt(t, sender, id, "rejected") r.mu.Lock() defer r.mu.Unlock() - if owner.apply || owner.data != nil || r.runtimePreparation != nil { + if owner.apply || owner.data != nil || len(r.runtimePreparations) != 0 { t.Fatal("invalid body retained or admitted a mutation") } }) @@ -164,7 +164,7 @@ func TestRuntimePreparationBeginRequiresExactBindingAndBounds(t *testing.T) { t.Fatal(err) } capabilitiesReceipt(t, sender, id, "rejected") - if r.runtimePreparation != nil { + if len(r.runtimePreparations) != 0 { t.Fatal("invalid scope allocated a transfer") } } @@ -177,9 +177,9 @@ func TestRuntimePreparationPreparationExcludesOwnedResources(t *testing.T) { own := proto.Envelope{Assignment: capabilityRef} switch mode { case "write": - r.workspaceWrite = &workspaceUpload{envelope: own} + r.workspaceWrites[session] = &workspaceUpload{envelope: own} case "export": - r.workspaceExport = &workspaceExport{request: own} + r.workspaceExports[session] = &workspaceExport{request: own} case "read": r.workspaceReads = map[string]string{"read": session} case "run": @@ -199,17 +199,17 @@ func TestRuntimePreparationPreparationExcludesOwnedResources(t *testing.T) { if mode == "other session" { capabilitiesReceipt(t, sender, id, "ready") r.mu.Lock() - r.finishRuntimePreparationTransferLocked(r.runtimePreparation, false) + r.finishRuntimePreparationTransferLocked(r.runtimePreparations[session], false) r.mu.Unlock() capabilitiesReceipt(t, sender, id, "rejected") } else if got := capabilitiesReceipt(t, sender, id, "rejected"); got.ErrorCode != "resource_unavailable" { t.Fatal(got) } - if r.runtimePreparation != nil { + if len(r.runtimePreparations) != 0 { t.Fatal("busy Runtime admitted capability preparation") } - r.workspaceWrite = nil - r.workspaceExport = nil + clear(r.workspaceWrites) + clear(r.workspaceExports) r.workspaceReads = nil clear(r.sessions) clear(r.executors) @@ -248,10 +248,10 @@ func TestRuntimePreparationUploadBlocksWorkspaceWriteAndSuspension(t *testing.T) t.Fatal("missing write rejection") } r.mu.Lock() - owner := r.runtimePreparation + owner := r.runtimePreparations[session] r.mu.Unlock() shutdownCapabilitiesRouter(t, r) - if owner.data != nil || r.runtimePreparation != nil { + if owner.data != nil || len(r.runtimePreparations) != 0 { t.Fatal("disconnect retained uncommitted body") } } @@ -263,7 +263,7 @@ func TestRuntimePreparationCancellationKeepsOwnershipUntilApplyStops(t *testing. request := proto.RuntimePreparePayload{Step: "begin", Action: "finalize", EnvironmentID: environment, SessionID: session, Sources: &agentcapabilities.Input{}} owner := &runtimePreparationTransfer{id: uuid.MustParse(id), envelope: capabilityEnvelope(t, id, request), request: request, ready: make(chan struct{}), cancel: cancel, finished: true, apply: true} close(owner.ready) - r.runtimePreparation = owner + r.runtimePreparations[session], r.transferBytes = owner, request.SizeBytes r.shutdownWG.Add(1) started, interrupted, release := make(chan struct{}), make(chan struct{}), make(chan struct{}) retained := filepath.Join(t.TempDir(), "installed.json") @@ -286,7 +286,7 @@ func TestRuntimePreparationCancellationKeepsOwnershipUntilApplyStops(t *testing. } <-interrupted r.mu.Lock() - owned := r.runtimePreparation == owner + owned := r.runtimePreparations[session] == owner r.mu.Unlock() if !owned { t.Fatal("cancel released unsettled capability ownership") @@ -297,7 +297,7 @@ func TestRuntimePreparationCancellationKeepsOwnershipUntilApplyStops(t *testing. if _, err := os.Stat(retained); err != nil { t.Fatal("shutdown deleted installation result") } - if r.runtimePreparation != nil { + if len(r.runtimePreparations) != 0 { t.Fatal("confirmed completion retained capacity") } } @@ -344,14 +344,14 @@ func TestRuntimePreparationResultCategoriesAndUnknownOwnership(t *testing.T) { 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.runtimePreparation = owner + 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.runtimePreparation == owner && owner.uncertain && owner.data == nil + owned := r.runtimePreparations[session] == owner && owner.uncertain && owner.data == nil r.mu.Unlock() if !owned { t.Fatal("unknown mutation released its ownership") @@ -362,3 +362,98 @@ func TestRuntimePreparationResultCategoriesAndUnknownOwnership(t *testing.T) { 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.mu.Lock() + delete(r.workspaceWrites, session) + r.mu.Unlock() + shutdownCapabilitiesRouter(t, r) + }) + } +} diff --git a/apps/daemon/internal/dispatch/shutdown.go b/apps/daemon/internal/dispatch/shutdown.go index 1c0e96836..d67d291f9 100644 --- a/apps/daemon/internal/dispatch/shutdown.go +++ b/apps/daemon/internal/dispatch/shutdown.go @@ -31,8 +31,8 @@ func (r *Router) Shutdown(ctx context.Context) error { victims := r.sessionCancellationsLocked() if !r.closed { r.closed = true - if r.runtimePreparation != nil { - r.runtimePreparation.cancel() + for _, u := range r.runtimePreparations { + u.cancel() } // Prepared release claims exist before this signal can interrupt output. close(r.shutdownCh) @@ -72,11 +72,15 @@ func (r *Router) runShutdownAttempt(attempt *shutdownAttempt, victims []sessionC attempt.err = errors.Join(attempt.err, err) } r.mu.Lock() - if r.runtimePreparation != nil && r.runtimePreparation.uncertain { - attempt.err = errors.Join(attempt.err, errors.New("dispatch: capability preparation remains uncertain")) + 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)) + } } - if r.workspaceWrite != nil && r.workspaceWrite.uncertain { - attempt.err = errors.Join(attempt.err, errors.New("dispatch: local workspace write remains uncertain")) + 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 { diff --git a/apps/daemon/internal/dispatch/suspend.go b/apps/daemon/internal/dispatch/suspend.go index 6ff583c86..ed9fed0e0 100644 --- a/apps/daemon/internal/dispatch/suspend.go +++ b/apps/daemon/internal/dispatch/suspend.go @@ -4,6 +4,7 @@ import ( "context" "errors" "strings" + "time" "github.com/MiniMax-AI/OpenAgentCore/internal/agentdaemon/proto" ) @@ -17,11 +18,116 @@ type AssignmentError string func (e AssignmentError) Error() string { return "dispatch: " + string(e) } -// Quiesce serializes against admission, then drains every admitted output and -// receipt and the Environment's owners before acknowledging suspension. Busy rejection leaves admission open; -// a drain timeout keeps it closed until the caller shuts the connection down. -// ref must admit work in the suspended Environment. +// quiesceTimeout bounds the drain of an environment_quiesce that Handle +// admitted. +const quiesceTimeout = 5 * time.Second + +// suspension is one quiesced Environment: the request and the assignment +// that quiesced it, which its resume must carry. +type suspension struct { + request proto.EnvironmentSuspendPayload + by proto.AssignmentRef +} + +// Quiesce quiesces the request's Environment for a caller that then +// suspends the whole Runtime: after fencing the Environment it waits for every +// output and receipt admitted on the connection, then closes the +// Environment's Executors and owners before it returns. Busy rejection leaves +// admission open; a failed drain keeps the Environment quiesced until a +// matching resume or shutdown. ref must admit work in the Environment. func (r *Router) Quiesce(ctx context.Context, ref proto.AssignmentRef, request proto.EnvironmentSuspendPayload) error { + if err := r.fenceEnvironment(ref, request); err != nil { + return err + } + if err := r.shutdownWG.waitContext(ctx); err != nil { + return err + } + return r.drainEnvironment(ctx, request.EnvironmentID) +} + +// Resume reopens the quiesced Environment and replaces the sender, after the +// caller authenticated a new connection and Core confirmed the exact +// suspension on it under the assignment that quiesced. +func (r *Router) Resume(ref proto.AssignmentRef, request proto.EnvironmentSuspendPayload, sender Sender) error { + if sender == nil { + return errors.New("dispatch: no sender") + } + r.mu.Lock() + defer r.mu.Unlock() + if err := r.resumeLocked(ref, request); err != nil { + return err + } + r.sender = sender + return nil +} + +// handleQuiesce quiesces an Environment for Core on this connection. It +// replies environment_quiesced once the drain settles, without blocking +// Handle; a failed drain replies resource_busy. +func (r *Router) handleQuiesce(ctx context.Context, env proto.Envelope) error { + var request proto.EnvironmentSuspendPayload + if env.ID == "" || len(env.ID) > 128 || env.DecodeRequest(&request) != nil || request.Rollback { + return r.replySuspension(ctx, env, proto.TypeEnvironmentQuiesced, request, "invalid_request") + } + if err := r.fenceEnvironment(env.Assignment, request); err != nil { + return r.replySuspension(ctx, env, proto.TypeEnvironmentQuiesced, request, suspensionCode(err)) + } + r.shutdownWG.Add(1) + go func() { + defer r.shutdownWG.Done() + drain, stop := r.shutdownContext(context.WithoutCancel(ctx)) + defer stop() + drain, cancel := context.WithTimeout(drain, quiesceTimeout) + defer cancel() + code := "" + if err := r.drainEnvironment(drain, request.EnvironmentID); err != nil { + r.log.Warn("environment quiesce unconfirmed", "environment_id", request.EnvironmentID, "err", err) + code = "resource_busy" + } + _ = r.replySuspension(context.WithoutCancel(ctx), env, proto.TypeEnvironmentQuiesced, request, code) + }() + return nil +} + +// handleResume reopens the quiesced Environment that the resume names. A +// rollback of an Environment this connection has not quiesced is accepted +// and changes nothing. +func (r *Router) handleResume(ctx context.Context, env proto.Envelope) error { + var request proto.EnvironmentSuspendPayload + if env.ID == "" || len(env.ID) > 128 || env.DecodeRequest(&request) != nil { + return r.replySuspension(ctx, env, proto.TypeEnvironmentResumed, request, "invalid_request") + } + r.mu.Lock() + err := r.resumeLocked(env.Assignment, request) + if errors.Is(err, errNotSuspended) && request.Rollback { + err = nil + } + r.mu.Unlock() + return r.replySuspension(ctx, env, proto.TypeEnvironmentResumed, request, suspensionCode(err)) +} + +var errNotSuspended = errors.New("dispatch: environment not suspended") + +// resumeLocked reopens the Environment that request quiesced under ref. +// Router.mu must be held. +func (r *Router) resumeLocked(ref proto.AssignmentRef, request proto.EnvironmentSuspendPayload) error { + s := r.suspensions[request.EnvironmentID] + switch { + case r.closed: + return ErrRouterClosed + case s == nil || !s.request.SameSuspension(request): + return errNotSuspended + case ref != s.by: + return AssignmentError(proto.AssignmentConflict) + } + delete(r.suspensions, request.EnvironmentID) + return nil +} + +// fenceEnvironment quiesces the request's Environment once none of its +// Sessions has unsettled work: from then on Handle admits only their +// releases and the matching resume. +func (r *Router) fenceEnvironment(ref proto.AssignmentRef, request proto.EnvironmentSuspendPayload) error { if strings.TrimSpace(request.EnvironmentID) == "" || strings.TrimSpace(request.SuspendID) == "" || len(request.SuspendID) > 128 { return errors.New("dispatch: invalid suspension identity") } @@ -30,80 +136,103 @@ func (r *Router) Quiesce(ctx context.Context, ref proto.AssignmentRef, request p } defer r.admission.Unlock() r.mu.Lock() - if r.closed { - r.mu.Unlock() + defer r.mu.Unlock() + switch { + case r.closed: return ErrRouterClosed - } - if r.suspension != nil { - r.mu.Unlock() + case r.suspensions[request.EnvironmentID] != nil: return ErrRouterQuiesced } if code := r.admitLocked(ref, ref.SessionID, request.EnvironmentID); code != "" { - r.mu.Unlock() return AssignmentError(code) } - if r.runtimePreparation != nil || len(r.sessions) != 0 || len(r.workspaceReads) != 0 || r.workspaceWrite != nil || r.workspaceExport != nil { - r.mu.Unlock() - return ErrRouterBusy + in := func(sessionID string) bool { + a := r.assignments[sessionID] + return a != nil && a.environmentID == request.EnvironmentID + } + for session := range r.assignments { + if in(session) && r.environmentTransferLocked(session) { + return ErrRouterBusy + } + } + for _, state := range r.sessions { + if state.environmentID == request.EnvironmentID { + return ErrRouterBusy + } + } + for _, session := range r.workspaceReads { + if in(session) { + return ErrRouterBusy + } } for _, p := range r.preparations { - if p.owns || p.busy { - r.mu.Unlock() + if (p.owns || p.busy) && p.environmentID == request.EnvironmentID { return ErrRouterBusy } } for _, owner := range r.executors { - if owner.invalid || owner.preparing || owner.run != nil || owner.admission != nil || owner.environmentID != request.EnvironmentID { - r.mu.Unlock() + if owner.environmentID == request.EnvironmentID && (owner.invalid || owner.preparing || owner.run != nil || owner.admission != nil) { return ErrRouterBusy } } - r.suspension, r.suspendedBy = &request, ref - for _, p := range r.preparations { - if p.timer != nil { - p.timer.Stop() + r.suspensions[request.EnvironmentID] = &suspension{request: request, by: ref} + return nil +} + +// drainEnvironment waits until the quiesced Environment's Sessions have sent +// every result they owe, then closes their Executors and releases their +// Environment owners. +func (r *Router) drainEnvironment(ctx context.Context, environmentID string) error { + r.mu.Lock() + var work []*dispatchWork + for _, a := range r.assignments { + if a.environmentID == environmentID { + work = append(work, &a.work) } } + var executors []*executorState for _, owner := range r.executors { - owner.idleLease++ - if owner.timer != nil { - owner.timer.Stop() + if owner.environmentID == environmentID { + owner.invalid, owner.closeReason = true, "quiesced" + executors = append(executors, owner) } } r.mu.Unlock() - err := r.shutdownWG.waitContext(ctx) - if err == nil { - err = r.closeEnvironments(ctx, request.EnvironmentID) + for _, w := range work { + if err := w.waitContext(ctx); err != nil { + return err + } } - r.mu.Lock() - if r.closed { - err = ErrRouterClosed + for _, owner := range executors { + if err := r.closeExecutor(owner); err != nil { + return err + } + } + if err := r.closeEnvironments(ctx, environmentID); err != nil { + return err } - r.mu.Unlock() - return err -} - -// Resume opens admission only after the caller authenticated a new connection -// and Core confirmed the exact suspension identity on that connection under the -// assignment that quiesced. -func (r *Router) Resume(ref proto.AssignmentRef, request proto.EnvironmentSuspendPayload, sender Sender) error { - r.admission.Lock() - defer r.admission.Unlock() r.mu.Lock() defer r.mu.Unlock() if r.closed { return ErrRouterClosed } - if sender == nil || r.suspension == nil || !r.suspension.SameSuspension(request) { - return errors.New("dispatch: suspension identity mismatch") - } - if ref != r.suspendedBy { - return AssignmentError(proto.AssignmentConflict) - } - r.sender = sender - r.suspension, r.suspendedBy = nil, proto.AssignmentRef{} - for _, owner := range r.executors { - r.scheduleExecutorIdleLocked(owner) - } return nil } + +// suspensionCode is the error code of a refused quiesce or resume. +func suspensionCode(err error) string { + var rejected AssignmentError + switch { + case err == nil: + return "" + case errors.As(err, &rejected): + return string(rejected) + case errors.Is(err, errNotSuspended): + return "not_suspended" + } + return "resource_busy" +} + +func (r *Router) replySuspension(ctx context.Context, env proto.Envelope, typ string, request proto.EnvironmentSuspendPayload, code string) error { + return r.reply(ctx, env, typ, proto.EnvironmentSuspendResultPayload{EnvironmentID: request.EnvironmentID, SuspendID: request.SuspendID, Accepted: code == "", ErrorCode: code}) +} diff --git a/apps/daemon/internal/dispatch/suspend_test.go b/apps/daemon/internal/dispatch/suspend_test.go index efe1d5714..482c6fe93 100644 --- a/apps/daemon/internal/dispatch/suspend_test.go +++ b/apps/daemon/internal/dispatch/suspend_test.go @@ -52,13 +52,16 @@ func suspensionRouter(t *testing.T, sender Sender) *Router { } func TestQuiesceRejectsEveryUnsettledResource(t *testing.T) { + session := suspendRef.SessionID cases := map[string]func(*Router){ - "active": func(r *Router) { r.sessions["run"] = &sessionState{} }, - "preparing": func(r *Router) { r.preparations["p"] = &preparationState{owns: true} }, - "receipt": func(r *Router) { r.preparations["p"] = &preparationState{busy: true} }, - "read": func(r *Router) { r.workspaceReads = map[string]string{"read": suspendRef.SessionID} }, - "write": func(r *Router) { r.workspaceWrite = &workspaceUpload{} }, - "export": func(r *Router) { r.workspaceExport = &workspaceExport{} }, + "active": func(r *Router) { r.sessions["run"] = &sessionState{environmentID: "env"} }, + "preparing": func(r *Router) { r.preparations["p"] = &preparationState{owns: true, environmentID: "env"} }, + "receipt": func(r *Router) { r.preparations["p"] = &preparationState{busy: true, environmentID: "env"} }, + "read": func(r *Router) { r.workspaceReads = map[string]string{"read": session} }, + "write": func(r *Router) { r.workspaceWrites[session] = &workspaceUpload{} }, + "export": func(r *Router) { r.workspaceExports[session] = &workspaceExport{} }, + "runtime preparation": func(r *Router) { r.runtimePreparations[session] = &runtimePreparationTransfer{} }, + "executor": func(r *Router) { r.executors[session] = &executorState{environmentID: "env", preparing: true} }, } for name, setup := range cases { t.Run(name, func(t *testing.T) { @@ -72,8 +75,10 @@ func TestQuiesceRejectsEveryUnsettledResource(t *testing.T) { r.sessions = map[string]*sessionState{} r.preparations = map[string]*preparationState{} r.workspaceReads = nil - r.workspaceWrite = nil - r.workspaceExport = nil + clear(r.workspaceWrites) + clear(r.workspaceExports) + clear(r.runtimePreparations) + clear(r.executors) }) } } @@ -90,7 +95,7 @@ func TestQuiesceDrainsPendingReceiptAndFencesConcurrentAdmission(t *testing.T) { deadline := time.After(time.Second) for { r.mu.Lock() - parked := r.suspension != nil + parked := r.suspensions["env"] != nil r.mu.Unlock() if parked { break @@ -151,33 +156,83 @@ func TestResumeRequiresExactSuspensionAndAssignment(t *testing.T) { } } -func TestQuiescePreservesIdleExecutorAgainstExpiredTimer(t *testing.T) { - sender := suspendSender(func(context.Context, proto.Envelope) error { return nil }) - r := suspensionRouter(t, sender) - native := &suspendedExecutor{} - owner := &executorState{id: "executor", sessionID: "session", environmentID: "env", native: native, cancel: func() {}} +func TestQuiescingOneEnvironmentLeavesAnotherRunning(t *testing.T) { + frames := make(chan proto.Envelope, 16) + r := suspensionRouter(t, suspendSender(func(_ context.Context, env proto.Envelope) error { frames <- env; return nil })) + other := proto.AssignmentRef{SessionID: "other", AssignmentID: "other", Epoch: 1} + bindAssignment(r, other, "other") + quiesced, running := &suspendedExecutor{}, &suspendedExecutor{} r.mu.Lock() - r.executors[owner.sessionID] = owner - r.scheduleExecutorIdleLocked(owner) - oldLease := owner.idleLease + r.executors[suspendRef.SessionID] = &executorState{id: "quiesced", sessionID: suspendRef.SessionID, environmentID: "env", native: quiesced, cancel: func() {}} + r.executors[other.SessionID] = &executorState{id: "running", sessionID: other.SessionID, environmentID: "other", native: running, cancel: func() {}} + // The other Environment's Session has a Run in progress. + r.sessions["run"] = &sessionState{assignment: other, environmentID: "other"} r.mu.Unlock() + suspend := func(typ, id string, ref proto.AssignmentRef, request proto.EnvironmentSuspendPayload) proto.EnvironmentSuspendResultPayload { + t.Helper() + env, err := proto.NewEnvelope(typ, id, request) + if err != nil { + t.Fatal(err) + } + env.Assignment = ref + if err := r.Handle(t.Context(), env); err != nil { + t.Fatal(err) + } + for { + select { + case frame := <-frames: + var result proto.EnvironmentSuspendResultPayload + if frame.ID == id && frame.DecodePayload(&result) == nil { + return result + } + case <-time.After(3 * time.Second): + t.Fatalf("%s %s has no result", typ, id) + } + } + } + prepare := func(ref proto.AssignmentRef) error { + return r.Handle(t.Context(), proto.Envelope{Type: proto.TypeExecutionPrepare, ID: "prepare", Assignment: ref}) + } request := proto.EnvironmentSuspendPayload{EnvironmentID: "env", SuspendID: "attempt"} - if err := r.Quiesce(context.Background(), suspendRef, request); err != nil { - t.Fatal(err) + if got := suspend(proto.TypeEnvironmentQuiesce, "quiesce", suspendRef, request); !got.Accepted { + t.Fatalf("quiesce = %+v", got) } - r.expireIdleExecutor(owner, oldLease) - if native.closed.Load() != 0 { - t.Fatal("pre-snapshot timer closed retained owner") + if quiesced.closed.Load() != 1 || running.closed.Load() != 0 { + t.Fatalf("closed Executors: quiesced %d, running %d", quiesced.closed.Load(), running.closed.Load()) } - if err := r.Resume(suspendRef, request, sender); err != nil { - t.Fatal(err) + if err := prepare(suspendRef); !errors.Is(err, ErrRouterQuiesced) { + t.Fatalf("the quiesced Environment admitted %v", err) + } + if err := prepare(other); errors.Is(err, ErrRouterQuiesced) { + t.Fatal("the other Environment stopped admitting work") + } + for id, test := range map[string]struct { + ref proto.AssignmentRef + request proto.EnvironmentSuspendPayload + code string + }{ + "foreign": {other, request, proto.AssignmentConflict}, + "unpaused": {other, proto.EnvironmentSuspendPayload{EnvironmentID: "other", SuspendID: "attempt"}, "not_suspended"}, + "rollback": {other, proto.EnvironmentSuspendPayload{EnvironmentID: "other", SuspendID: "attempt", Rollback: true}, ""}, + "obsolete": {suspendRef, proto.EnvironmentSuspendPayload{EnvironmentID: "env", SuspendID: "obsolete"}, "not_suspended"}, + } { + if got := suspend(proto.TypeEnvironmentResume, id, test.ref, test.request); got.ErrorCode != test.code || got.Accepted != (test.code == "") { + t.Fatalf("%s resume = %+v", id, got) + } + } + if got := suspend(proto.TypeEnvironmentResume, "resume", suspendRef, request); !got.Accepted { + t.Fatalf("resume = %+v", got) + } + if err := prepare(suspendRef); errors.Is(err, ErrRouterQuiesced) { + t.Fatal("resume did not reopen the Environment") } r.mu.Lock() - newLease := owner.idleLease + _, run := r.sessions["run"] + // The synthetic Run owns no goroutine; remove it before cleanup. + clear(r.sessions) r.mu.Unlock() - r.expireIdleExecutor(owner, newLease) - if native.closed.Load() != 1 { - t.Fatal("normal idle expiration was not restored") + if !run || running.closed.Load() != 0 { + t.Fatal("quiescing one Environment ended another's work") } } diff --git a/apps/daemon/internal/dispatch/workspace_export.go b/apps/daemon/internal/dispatch/workspace_export.go index eeb2d912c..e7271f30e 100644 --- a/apps/daemon/internal/dispatch/workspace_export.go +++ b/apps/daemon/internal/dispatch/workspace_export.go @@ -26,7 +26,7 @@ func (r *Router) handleWorkspaceExport(ctx context.Context, env proto.Envelope) r.mu.Unlock() return ErrRouterClosed } - u := r.workspaceExport + u := r.workspaceExports[env.Assignment.SessionID] if request.Step != "begin" { if u == nil || u.request.ID != env.ID || u.request.Assignment != env.Assignment { r.mu.Unlock() @@ -49,7 +49,7 @@ func (r *Router) handleWorkspaceExport(ctx context.Context, env proto.Envelope) } environment, code := r.workspaceResourceLocked(env.Assignment, proto.WorkspaceReadPayload{Handle: request.Handle, EnvironmentID: request.EnvironmentID}) p := r.preparations[request.Handle] - if code == "" && (u != nil || r.workspaceWrite != nil && r.workspaceWrite.envelope.Assignment.SessionID == env.Assignment.SessionID || p.executor != nil) { + if code == "" && (u != nil || r.workspaceWrites[env.Assignment.SessionID] != nil || p.executor != nil) { code = "resource_unavailable" } if code != "" { @@ -60,7 +60,7 @@ func (r *Router) handleWorkspaceExport(ctx context.Context, env proto.Envelope) owner, cancel := context.WithTimeout(owner, 180*time.Second) u = &workspaceExport{request: env, requests: make(chan proto.WorkspaceExportPayload, 1), cancel: func() { cancel(); stop() }} u.requests <- request - r.workspaceExport = u + r.workspaceExports[env.Assignment.SessionID] = u done := r.trackWorkLocked(env.Assignment) r.shutdownWG.Add(1) r.mu.Unlock() @@ -90,8 +90,8 @@ func (r *Router) runWorkspaceExport(ctx context.Context, u *workspaceExport, env _ = reader.Close() <-exported r.mu.Lock() - if r.workspaceExport == u { - r.workspaceExport = nil + if r.workspaceExports[u.request.Assignment.SessionID] == u { + delete(r.workspaceExports, u.request.Assignment.SessionID) } r.mu.Unlock() select { @@ -129,8 +129,8 @@ func (r *Router) runWorkspaceExport(ctx context.Context, u *workspaceExport, env // The next owner may start immediately after receiving completion. <-exported r.mu.Lock() - if r.workspaceExport == u { - r.workspaceExport = nil + if r.workspaceExports[u.request.Assignment.SessionID] == u { + delete(r.workspaceExports, u.request.Assignment.SessionID) } r.mu.Unlock() } diff --git a/apps/daemon/internal/dispatch/workspace_export_test.go b/apps/daemon/internal/dispatch/workspace_export_test.go index 19e9f2843..88ca1b55b 100644 --- a/apps/daemon/internal/dispatch/workspace_export_test.go +++ b/apps/daemon/internal/dispatch/workspace_export_test.go @@ -195,7 +195,7 @@ func TestWorkspaceExportCancelUnblocksWriterAndReleasesCapacity(t *testing.T) { deadline := time.Now().Add(3 * time.Second) for time.Now().Before(deadline) { r.mu.Lock() - active := r.workspaceExport != nil + active := len(r.workspaceExports) != 0 r.mu.Unlock() if !active { return @@ -220,7 +220,9 @@ func TestCanceledExportAnswersTheRequestCoreAwaits(t *testing.T) { // Core asks for the next chunk as the preparation that owns the export is released. r.mu.Lock() release() - r.workspaceExport.requests <- proto.WorkspaceExportPayload{Step: "next", Offset: int64(len(first.Data))} + for _, u := range r.workspaceExports { + u.requests <- proto.WorkspaceExportPayload{Step: "next", Offset: int64(len(first.Data))} + } r.mu.Unlock() if got := readExport(t, s); got.Outcome != "failed" && got.Outcome != "chunk" { t.Fatal("the canceled export answered with", got) diff --git a/apps/daemon/internal/dispatch/workspace_write.go b/apps/daemon/internal/dispatch/workspace_write.go index 3bf06b4e2..3b8ecefb6 100644 --- a/apps/daemon/internal/dispatch/workspace_write.go +++ b/apps/daemon/internal/dispatch/workspace_write.go @@ -11,7 +11,7 @@ import ( "github.com/google/uuid" ) -// Router.mu protects this single bounded transfer to an Environment owner. +// Router.mu protects a Session's bounded transfer to its Environment owner. type workspaceUpload struct { envelope proto.Envelope request proto.WorkspaceWritePayload @@ -32,7 +32,8 @@ func (r *Router) handleWorkspaceWrite(ctx context.Context, env proto.Envelope) e // A malformed frame on an already admitted operation cannot claim that // its earlier commit did not execute. r.mu.Lock() - pending := r.workspaceWrite != nil && r.workspaceWrite.envelope.ID == env.ID + u := r.workspaceWrites[env.Assignment.SessionID] + pending := u != nil && u.envelope.ID == env.ID r.mu.Unlock() if pending { return errors.New("dispatch: malformed pending write frame") @@ -44,14 +45,14 @@ func (r *Router) handleWorkspaceWrite(ctx context.Context, env proto.Envelope) e r.mu.Unlock() return ErrRouterClosed } + u := r.workspaceWrites[env.Assignment.SessionID] if request.Step == "begin" { - if r.workspaceExport != nil && r.workspaceExport.request.Assignment.SessionID == request.SessionID || - r.runtimePreparation != nil && r.runtimePreparation.envelope.Assignment.SessionID == request.SessionID { + if r.workspaceExports[request.SessionID] != nil || r.runtimePreparations[request.SessionID] != nil { r.mu.Unlock() return r.sendWorkspaceWrite(ctx, env, rejectedWorkspaceWrite("resource_unavailable")) } - if r.workspaceWrite != nil { - duplicate := r.workspaceWrite.envelope.ID == env.ID + if u != nil || r.transferBytes > transferMemory-request.SizeBytes { + duplicate := u != nil && u.envelope.ID == env.ID r.mu.Unlock() if duplicate { return errors.New("dispatch: workspace write already admitted") @@ -73,14 +74,14 @@ func (r *Router) handleWorkspaceWrite(ctx context.Context, env proto.Envelope) e return r.sendWorkspaceWrite(ctx, env, rejectedWorkspaceWrite("resource_unavailable")) } u := &workspaceUpload{envelope: env, request: request, data: make([]byte, 0, request.SizeBytes), ready: make(chan struct{})} - r.workspaceWrite = u + r.workspaceWrites[request.SessionID] = u + r.transferBytes += request.SizeBytes done := r.trackWorkLocked(env.Assignment) r.shutdownWG.Add(1) r.mu.Unlock() go r.runWorkspaceUpload(context.WithoutCancel(ctx), u, environment, done) return r.sendWorkspaceWrite(ctx, env, proto.WorkspaceWriteResultPayload{Outcome: "ready"}) } - u := r.workspaceWrite if u == nil || u.envelope.ID != env.ID || u.envelope.Assignment != env.Assignment { r.mu.Unlock() return r.sendWorkspaceWrite(ctx, env, rejectedWorkspaceWrite("resource_unavailable")) @@ -132,9 +133,10 @@ func (r *Router) runWorkspaceUpload(ctx context.Context, u *workspaceUpload, env result = workspaceWriteResult(write, err, u.request.SizeBytes) } r.mu.Lock() + r.transferBytes -= u.request.SizeBytes u.uncertain = result.Outcome == "unknown" if !u.uncertain { - r.workspaceWrite = nil + delete(r.workspaceWrites, u.request.SessionID) } r.mu.Unlock() _ = r.sendWorkspaceWrite(ctx, u.envelope, result) diff --git a/docs/configuration.md b/docs/configuration.md index ae09a0fba..ded27d690 100644 --- a/docs/configuration.md +++ b/docs/configuration.md @@ -146,7 +146,14 @@ The agent host runs each Session's Harness outside the sandbox, in a view of its Docker's default seccomp profile stays: with `CAP_SYS_ADMIN` it allows `clone3`, `mount` and `unshare`. The container gets no Docker socket and publishes no port. -The agent host needs a cgroup v2 directory delegated to it. It starts each view in a cgroup of its own there, ends the view with `cgroup.kill` and removes the cgroup. When it starts it ends and removes every cgroup in the directory, because each counts as a view's, so nothing else may use it. The directory must be writable and must not contain the agent host's own process. Docker mounts the container's cgroup read-only. With the flags above, cgroup v2 mounted again inside the container (`mount -t cgroup2 cgroup2