diff --git a/apps/daemon/internal/agent/claudesdk/declaration.go b/apps/daemon/internal/agent/claudesdk/declaration.go index e99650049..1360c2cbb 100644 --- a/apps/daemon/internal/agent/claudesdk/declaration.go +++ b/apps/daemon/internal/agent/claudesdk/declaration.go @@ -62,7 +62,7 @@ func discoverWithCheck(parent context.Context, options agent.DiscoveryOptions, d } out := &agent.Runtime{Info: descriptor, Session: func(context.Context, proto.PromptRequestPayload, chan<- proto.Envelope) (agent.Session, error) { return nil, fmt.Errorf("claude_sdk: configured runtime is unavailable") - }} + }, View: nil} var config Config fail := func(err error) *agent.Runtime { diff --git a/apps/daemon/internal/agent/clirunner/handle.go b/apps/daemon/internal/agent/clirunner/handle.go new file mode 100644 index 000000000..fd952579d --- /dev/null +++ b/apps/daemon/internal/agent/clirunner/handle.go @@ -0,0 +1,106 @@ +package clirunner + +import ( + "context" + "errors" + "fmt" + "io" + "os" + "sync" + "sync/atomic" + "syscall" + "time" +) + +// Handle is a running process that clirunner did not start, such as a Harness in an agent-host Session view. Its end also ends every descendant. +type Handle interface { + // Signal delivers sig to the process and every descendant it still has. Once the process itself has exited it delivers nothing and returns an error that matches os.ErrProcessDone, even while its descendants are still ending. + Signal(syscall.Signal) error + // Wait returns once the process and its descendants have ended. The code is -1 when a signal ended the process. An error means the exit is unknown. + Wait() (int, error) + // Close kills whatever still runs and releases the handle. It never closes the stdio ends in HandleOptions, which the Process owns. It is safe to call more than once and after Wait. + Close() error +} + +// HandleOptions are the stdio ends of a Handle's process and its cancellation grace. The Process owns the stdio ends and closes them in Wait. +type HandleOptions struct { + Parent context.Context + // Stdin is nil when the process has no stdin pipe. + Stdin io.WriteCloser + Stdout io.ReadCloser + Stderr io.ReadCloser + KillTimeout time.Duration +} + +// FromHandle returns a Process that owns h. Cancel, or cancelling Parent, sends TERM and closes h after KillTimeout; cancelling a process that has ended does nothing. Wait closes the stdio ends and h, and returns the context error when the TERM reached the running process and it then exited 0, as exec.CommandContext does. A TERM that finds the process exited, which Signal reports with os.ErrProcessDone, leaves its result as it was. +func FromHandle(h Handle, opts HandleOptions) (*Process, error) { + if h == nil || opts.Stdout == nil || opts.Stderr == nil { + return nil, errors.New("clirunner: handle, stdout and stderr required") + } + if opts.Parent == nil { + opts.Parent = context.Background() + } + if opts.KillTimeout <= 0 { + opts.KillTimeout = 3 * time.Second + } + ctx, cancel := context.WithCancel(opts.Parent) + p := &Process{Stdin: opts.Stdin, Stdout: opts.Stdout, Stderr: opts.Stderr, ctx: ctx, cancel: cancel, done: make(chan struct{}), killAfter: opts.KillTimeout} + var terminateOnce sync.Once + var interrupted atomic.Bool + p.cancelProcess = func() error { + terminateOnce.Do(func() { + err := h.Signal(syscall.SIGTERM) + interrupted.Store(err == nil) + go func() { + if err == nil || errors.Is(err, os.ErrProcessDone) { + timer := time.NewTimer(p.killAfter) + defer timer.Stop() + select { + case <-p.done: + return + case <-timer.C: + } + } + _ = h.Close() + }() + }) + return nil + } + stop := context.AfterFunc(ctx, func() { _ = p.cancelProcess() }) + var waitErr error + go func() { + code, err := h.Wait() + stop() + // A cancel in progress settles whether it interrupted the process, and none starts after the end. + terminateOnce.Do(func() {}) + switch { + case err != nil: + waitErr = fmt.Errorf("clirunner: wait: %w", err) + case code != 0: + waitErr = fmt.Errorf("clirunner: exit code %d", code) + case interrupted.Load(): + // Cancel runs before the context is cancelled. + if waitErr = ctx.Err(); waitErr == nil { + waitErr = context.Canceled + } + } + if err == nil { + p.exitCode, p.exited = code, true + } + close(p.done) + }() + var waitOnce sync.Once + p.waitProcess = func() error { + <-p.done + waitOnce.Do(func() { + closePipe(p.Stdin) + closePipe(p.Stdout) + closePipe(p.Stderr) + if err := h.Close(); err != nil { + waitErr = errors.Join(waitErr, fmt.Errorf("clirunner: close: %w", err)) + } + }) + return waitErr + } + return p, nil +} diff --git a/apps/daemon/internal/agent/clirunner/handle_test.go b/apps/daemon/internal/agent/clirunner/handle_test.go new file mode 100644 index 000000000..5d8900ec0 --- /dev/null +++ b/apps/daemon/internal/agent/clirunner/handle_test.go @@ -0,0 +1,177 @@ +package clirunner + +import ( + "context" + "errors" + "io" + "os" + "slices" + "strings" + "sync" + "syscall" + "testing" + "time" +) + +func TestHandleProcessCancellation(t *testing.T) { + const grace = 100 * time.Millisecond + for name, termExit := range map[string]int{"ignores TERM": -1, "exits 0 on TERM": 0, "exits 3 on TERM": 3} { + t.Run(name, func(t *testing.T) { + h := newFakeHandle(termExit) + p, err := FromHandle(h, HandleOptions{Stdout: emptyReader(), Stderr: emptyReader(), KillTimeout: grace}) + if err != nil { + t.Fatal(err) + } + started := time.Now() + p.Cancel() + if signals := h.received(); !slices.Equal(signals, []syscall.Signal{syscall.SIGTERM}) { + t.Fatalf("signals = %v, want TERM", signals) + } + <-p.Done() + waitErr := p.Wait() + code, ok := p.ExitCode() + switch { + case termExit < 0: + if h.closedAt.Sub(started) < grace || waitErr == nil || ok { + t.Fatalf("closed after %v, Wait = %v, ExitCode ok = %v; want close after the grace and an unknown exit", h.closedAt.Sub(started), waitErr, ok) + } + case termExit == 0: + if !errors.Is(waitErr, context.Canceled) || !ok || code != 0 { + t.Fatalf("Wait = %v, ExitCode = %d, %v; want the context error and exit 0", waitErr, code, ok) + } + default: + if waitErr == nil || errors.Is(waitErr, context.Canceled) || !ok || code != termExit { + t.Fatalf("Wait = %v, ExitCode = %d, %v; want exit %d", waitErr, code, ok, termExit) + } + } + if !h.isClosed() { + t.Fatal("Wait did not close the handle") + } + }) + } +} + +func TestHandleProcessCancelAfterExit(t *testing.T) { + h := newFakeHandle(-1) + stdin, stdout, stderr := &closer{}, &closer{Reader: strings.NewReader("output")}, &closer{Reader: strings.NewReader("")} + p, err := FromHandle(h, HandleOptions{Stdin: stdin, Stdout: stdout, Stderr: stderr}) + if err != nil { + t.Fatal(err) + } + h.exit <- 0 + <-p.Done() + p.Cancel() + if signals := h.received(); len(signals) != 0 || h.isClosed() || stdout.isClosed() { + t.Fatalf("Cancel after exit sent %v and closed the handle %v, stdout %v; want nothing", signals, h.isClosed(), stdout.isClosed()) + } + if out, err := io.ReadAll(p.Stdout); err != nil || string(out) != "output" { + t.Fatalf("read %q, %v; want the whole output", out, err) + } + for range 2 { + if err := p.Wait(); err != nil { + t.Fatalf("Wait = %v, want success", err) + } + } + if !stdin.isClosed() || !stdout.isClosed() || !stderr.isClosed() || !h.isClosed() { + t.Fatal("Wait did not close stdin, stdout, stderr and the handle") + } +} + +// TestHandleProcessCancelAfterLeaderExit checks that a cancel that finds the process exited 0 while its descendants still end leaves the success. +func TestHandleProcessCancelAfterLeaderExit(t *testing.T) { + h := newFakeHandle(-1) + h.leaderExited = true + p, err := FromHandle(h, HandleOptions{Stdout: emptyReader(), Stderr: emptyReader()}) + if err != nil { + t.Fatal(err) + } + p.Cancel() + h.exit <- 0 + if err := p.Wait(); err != nil { + t.Fatalf("Wait = %v, want success", err) + } +} + +type fakeHandle struct { + termExit int + leaderExited bool + exit chan int + closed chan struct{} + closeOnce sync.Once + closedAt time.Time + mu sync.Mutex + signals []syscall.Signal +} + +func newFakeHandle(termExit int) *fakeHandle { + return &fakeHandle{termExit: termExit, exit: make(chan int, 1), closed: make(chan struct{})} +} + +func (h *fakeHandle) Signal(sig syscall.Signal) error { + if h.leaderExited { + return os.ErrProcessDone + } + h.mu.Lock() + h.signals = append(h.signals, sig) + h.mu.Unlock() + if sig == syscall.SIGTERM && h.termExit >= 0 { + h.exit <- h.termExit + } + return nil +} + +func (h *fakeHandle) Wait() (int, error) { + select { + case code := <-h.exit: + return code, nil + case <-h.closed: + return 0, errors.New("closed") + } +} + +func (h *fakeHandle) Close() error { + h.closeOnce.Do(func() { + h.closedAt = time.Now() + close(h.closed) + }) + return nil +} + +func (h *fakeHandle) isClosed() bool { + select { + case <-h.closed: + return true + default: + return false + } +} + +func (h *fakeHandle) received() []syscall.Signal { + h.mu.Lock() + defer h.mu.Unlock() + return slices.Clone(h.signals) +} + +// closer is a stdio end that records Close. +type closer struct { + io.Reader + mu sync.Mutex + closed bool +} + +func (c *closer) Write(b []byte) (int, error) { return len(b), nil } + +func (c *closer) Close() error { + c.mu.Lock() + defer c.mu.Unlock() + c.closed = true + return nil +} + +func (c *closer) isClosed() bool { + c.mu.Lock() + defer c.mu.Unlock() + return c.closed +} + +func emptyReader() io.ReadCloser { return io.NopCloser(strings.NewReader("")) } diff --git a/apps/daemon/internal/agent/clirunner/process.go b/apps/daemon/internal/agent/clirunner/process.go index e693df7f5..d9d40efc3 100644 --- a/apps/daemon/internal/agent/clirunner/process.go +++ b/apps/daemon/internal/agent/clirunner/process.go @@ -38,6 +38,10 @@ type Process struct { waitOnce sync.Once cancelProcess func() error waitProcess func() error + + // A handle-backed process records its exit here before closing done. + exitCode int + exited bool } func Start(opts StartOptions) (*Process, error) { @@ -144,17 +148,39 @@ func (p *Process) Cancel() { } func (p *Process) Wait() error { - if p == nil || p.Cmd == nil { + if p == nil { return nil } if p.waitProcess != nil { return p.waitProcess() } + if p.Cmd == nil { + return nil + } err := p.Cmd.Wait() p.waitOnce.Do(func() { close(p.done) }) return err } +// ExitCode returns the exit code once Done is closed, or -1 when a signal ended the process. ok is false before then and when the exit is unknown. +func (p *Process) ExitCode() (code int, ok bool) { + if p == nil { + return 0, false + } + select { + case <-p.done: + default: + return 0, false + } + if p.Cmd == nil { + return p.exitCode, p.exited + } + if p.Cmd.ProcessState == nil { + return 0, false + } + return p.Cmd.ProcessState.ExitCode(), true +} + func closePipe(p io.Closer) { if p != nil { _ = p.Close() diff --git a/apps/daemon/internal/agent/codex/declaration.go b/apps/daemon/internal/agent/codex/declaration.go index 2c962fe2c..7020b66c0 100644 --- a/apps/daemon/internal/agent/codex/declaration.go +++ b/apps/daemon/internal/agent/codex/declaration.go @@ -45,13 +45,13 @@ var Declaration = agent.Declaration{Info: proto.SupportedAgentKind{ MCPHTTPRequired: proto.CapabilityUnsupported, MCPHTTPBearerAuth: proto.CapabilitySupported, }, -}, Configuration: configuration.Configuration(), Discover: discover} +}, Configuration: configuration.Configuration(), ConnectionOptions: []string{"mcp_servers", "env"}, Discover: discover} func discover(ctx context.Context, options agent.DiscoveryOptions, info proto.SupportedAgentKind) *agent.Runtime { return discoverWithCheck(ctx, options, info, CheckCLIAvailable) } func discoverWithCheck(parent context.Context, options agent.DiscoveryOptions, info proto.SupportedAgentKind, check func(context.Context, string) (string, error)) *agent.Runtime { - runtime := &agent.Runtime{Info: info, Session: Factory, SessionCapabilityContext: true, ExecutorCapabilityContext: true} + runtime := &agent.Runtime{Info: info, Session: Factory, SessionCapabilityContext: true, ExecutorCapabilityContext: true, View: nil} ctx, cancel := context.WithTimeout(parent, 15*time.Second) defer cancel() version, err := check(ctx, "") diff --git a/apps/daemon/internal/agent/codex/executor_native_test.go b/apps/daemon/internal/agent/codex/executor_native_test.go index 227ad311c..a9eba21ad 100644 --- a/apps/daemon/internal/agent/codex/executor_native_test.go +++ b/apps/daemon/internal/agent/codex/executor_native_test.go @@ -69,7 +69,7 @@ func TestExecutorNativeReuse(t *testing.T) { } }) t.Logf("native_version=0.153.4 prepare_ms=%d", time.Since(began).Milliseconds()) - pid := e.prepared.session.rpc.cmd.Process.Pid + pid := e.prepared.session.rpc.process.Cmd.Process.Pid var thread string run := func(id, prompt string, old agent.Turn, interrupt bool) (agent.Turn, string) { t.Helper() @@ -135,7 +135,7 @@ func TestExecutorNativeReuse(t *testing.T) { if thread == "" { thread = s.currentThreadID() } - if thread == "" || s.currentThreadID() != thread || !s.rpc.Alive() || s.rpc.cmd.Process.Pid != pid { + if thread == "" || s.currentThreadID() != thread || !s.rpc.Alive() || s.rpc.process.Cmd.Process.Pid != pid { t.Fatal("native owner/thread changed") } t.Logf("turn=%s first_event_ms=%d settled_ms=%d same_process=true same_thread=true", id, first.Milliseconds(), time.Since(started).Milliseconds()) diff --git a/apps/daemon/internal/agent/codex/executor_test.go b/apps/daemon/internal/agent/codex/executor_test.go index 01245f944..dd8d4f518 100644 --- a/apps/daemon/internal/agent/codex/executor_test.go +++ b/apps/daemon/internal/agent/codex/executor_test.go @@ -170,7 +170,7 @@ func TestExecutorCloseRetainsPlanUntilReaped(t *testing.T) { t.Fatal(err) } rpc := NewJSONRPCClient(JSONRPCConfig{}) - rpc.process, rpc.cmd, rpc.stdin, rpc.alive = process, process.Cmd, process.Stdin, true + rpc.process, rpc.stdin, rpc.alive = process, process.Stdin, true var reap sync.Once t.Cleanup(func() { process.Cancel(); reap.Do(rpc.waitChild) }) diff --git a/apps/daemon/internal/agent/codex/preparation_router_test.go b/apps/daemon/internal/agent/codex/preparation_router_test.go index 660c94504..7f7ccb53e 100644 --- a/apps/daemon/internal/agent/codex/preparation_router_test.go +++ b/apps/daemon/internal/agent/codex/preparation_router_test.go @@ -119,7 +119,7 @@ func TestPreparationRouterRetainsActualNativeChild(t *testing.T) { ready := await("ready") p := <-prepared assertPreparationOnly(t, root) - pid := p.session.rpc.cmd.Process.Pid + pid := p.session.rpc.process.Cmd.Process.Pid if r.ActiveRuns() != 0 { t.Fatal("preparation became a Run") } diff --git a/apps/daemon/internal/agent/codex/prepared_test.go b/apps/daemon/internal/agent/codex/prepared_test.go index e729c9613..c52bb3bbb 100644 --- a/apps/daemon/internal/agent/codex/prepared_test.go +++ b/apps/daemon/internal/agent/codex/prepared_test.go @@ -27,7 +27,7 @@ func TestPreparedSessionTransfersSameResourceOnce(t *testing.T) { if len(preparedCatalogs(t, root)) != 1 { t.Fatal("preparation did not retain its model catalog") } - pid := p.session.rpc.cmd.Process.Pid + pid := p.session.rpc.process.Cmd.Process.Pid cfgPreparedCwd := p.plan.Cwd // Caller-owned data cannot revise the prepared native configuration. req.AgentOptions["model"] = "different-model" @@ -42,7 +42,7 @@ func TestPreparedSessionTransfersSameResourceOnce(t *testing.T) { } session := started.(*Session) defer session.Cancel(context.Background()) - if session.rpc != p.session.rpc || session.rpc.cmd.Process.Pid != pid { + if session.rpc != p.session.rpc || session.rpc.process.Cmd.Process.Pid != pid { t.Fatal("start replaced the prepared native resource") } if err := p.Close(); err != nil || !session.rpc.Alive() { diff --git a/apps/daemon/internal/agent/codex/rpc.go b/apps/daemon/internal/agent/codex/rpc.go index 2102eda9e..dc3bba9cb 100644 --- a/apps/daemon/internal/agent/codex/rpc.go +++ b/apps/daemon/internal/agent/codex/rpc.go @@ -10,7 +10,6 @@ import ( "fmt" "io" "log/slog" - "os/exec" "sync" "time" @@ -82,7 +81,6 @@ type JSONRPCClient struct { cfg JSONRPCConfig process *clirunner.Process - cmd *exec.Cmd stdin io.WriteCloser stdout io.ReadCloser stderr io.ReadCloser @@ -189,7 +187,6 @@ func (c *JSONRPCClient) Start(ctx context.Context, init InitializeParams) (Initi } c.mu.Lock() c.process = process - c.cmd = process.Cmd c.stdin, c.stdout, c.stderr = process.Stdin, process.Stdout, process.Stderr c.alive = true c.mu.Unlock() @@ -465,8 +462,7 @@ func (c *JSONRPCClient) waitChild() { err := c.process.Wait() c.mu.Lock() c.alive = false - if c.cmd.ProcessState != nil { - code := c.cmd.ProcessState.ExitCode() + if code, ok := c.process.ExitCode(); ok { c.exitCode = &code } c.mu.Unlock() diff --git a/apps/daemon/internal/agent/codex/rpc_close_test.go b/apps/daemon/internal/agent/codex/rpc_close_test.go index 80e404a85..5fd3aaefd 100644 --- a/apps/daemon/internal/agent/codex/rpc_close_test.go +++ b/apps/daemon/internal/agent/codex/rpc_close_test.go @@ -19,10 +19,10 @@ func TestJSONRPCClientCloseCanRetryUnreapedChild(t *testing.T) { if err != nil { t.Fatal(err) } - cmd, stdin, stdout := process.Cmd, process.Stdin, process.Stdout + stdin, stdout := process.Stdin, process.Stdout client := NewJSONRPCClient(JSONRPCConfig{}) input := &countedCloseWriter{WriteCloser: stdin} - client.process, client.cmd, client.stdin, client.stdout, client.alive = process, cmd, input, stdout, true + client.process, client.stdin, client.stdout, client.alive = process, input, stdout, true var reap sync.Once t.Cleanup(func() { process.Cancel() @@ -91,7 +91,7 @@ func TestJSONRPCClientCloseCanRetryUnreapedChild(t *testing.T) { t.Fatalf("Session cancellation retry after reap: %v", err) } closeConcurrently(false) - if client.cmd != cmd || cmd.ProcessState == nil { + if _, exited := process.ExitCode(); client.process != process || !exited { t.Fatal("Close did not retain and reap its original child") } if got := input.closes.Load(); got != 1 { diff --git a/apps/daemon/internal/agent/codex/rpc_process_test.go b/apps/daemon/internal/agent/codex/rpc_process_test.go index eb0a7647d..d499a8bd2 100644 --- a/apps/daemon/internal/agent/codex/rpc_process_test.go +++ b/apps/daemon/internal/agent/codex/rpc_process_test.go @@ -26,7 +26,7 @@ func TestJSONRPCClientCloseReapsChildProcess(t *testing.T) { if err != nil { t.Fatalf("start fake process: %v", err) } - if client.cmd == nil || client.cmd.Process == nil { + if client.process == nil || client.process.Cmd.Process == nil { t.Fatal("client did not spawn a child process") } @@ -38,8 +38,8 @@ func TestJSONRPCClientCloseReapsChildProcess(t *testing.T) { case <-time.After(2 * time.Second): t.Fatal("child process was not reaped") } - if client.cmd.ProcessState == nil { - t.Fatalf("process state not exited after close: %#v", client.cmd.ProcessState) + if _, exited := client.process.ExitCode(); !exited { + t.Fatal("process exit not recorded after close") } } diff --git a/apps/daemon/internal/agent/codex/session_cancel_write_test.go b/apps/daemon/internal/agent/codex/session_cancel_write_test.go index 7c0eed406..53cbb1392 100644 --- a/apps/daemon/internal/agent/codex/session_cancel_write_test.go +++ b/apps/daemon/internal/agent/codex/session_cancel_write_test.go @@ -169,7 +169,7 @@ func TestCancelReleasesBlockedNativeProcess(t *testing.T) { case <-time.After(time.Second): t.Fatal("cancelled native child was not reaped") } - if client.Alive() || client.cmd.ProcessState == nil || ctx.Err() == nil { + if _, exited := client.process.ExitCode(); client.Alive() || !exited || ctx.Err() == nil { t.Fatal("native process or cancellation context remained active") } // The cleanup consumes the writer's result after proving it was released. diff --git a/apps/daemon/internal/agent/contract_declarations_test.go b/apps/daemon/internal/agent/contract_declarations_test.go index fc642d3f4..76d404b7d 100644 --- a/apps/daemon/internal/agent/contract_declarations_test.go +++ b/apps/daemon/internal/agent/contract_declarations_test.go @@ -2,6 +2,7 @@ package agent_test import ( "encoding/json" + "fmt" "go/ast" "go/parser" "go/token" @@ -123,6 +124,33 @@ func TestPublicHarnessContractDeclarations(t *testing.T) { t.Errorf("%s needs explicit compile assertions on %v; got %v", name, want, assertions[name]) } } + // Each discovered Runtime decides its agent-host view explicitly, even when it declares none. + declaration, err := parser.ParseFile(token.NewFileSet(), filepath.Join(entry.Configuration, "declaration.go"), nil, 0) + if err != nil { + t.Fatal(err) + } + runtimes := 0 + ast.Inspect(declaration, func(node ast.Node) bool { + literal, ok := node.(*ast.CompositeLit) + if !ok { + return true + } + selector, ok := literal.Type.(*ast.SelectorExpr) + if !ok || selector.Sel.Name != "Runtime" { + return true + } + runtimes++ + for _, element := range literal.Elts { + if field, ok := element.(*ast.KeyValueExpr); ok && fmt.Sprint(field.Key) == "View" { + return true + } + } + t.Error("agent.Runtime literal must set View explicitly") + return true + }) + if runtimes == 0 { + t.Error("declaration.go constructs no agent.Runtime") + } }) } } diff --git a/apps/daemon/internal/agent/harness.go b/apps/daemon/internal/agent/harness.go index 7d0f715dd..2528f5342 100644 --- a/apps/daemon/internal/agent/harness.go +++ b/apps/daemon/internal/agent/harness.go @@ -27,10 +27,20 @@ package agent import ( "context" + "errors" + "fmt" "io" - + "net/netip" + "net/url" + "path" + "path/filepath" + "slices" + "strings" + + "github.com/MiniMax-AI/OpenAgentCore/apps/daemon/internal/agent/clirunner" "github.com/MiniMax-AI/OpenAgentCore/internal/agentdaemon/proto" "github.com/MiniMax-AI/OpenAgentCore/internal/harnessconfig" + "github.com/MiniMax-AI/OpenAgentCore/internal/modelprovider" ) // Declaration is the complete startup contract for a Harness implementation. @@ -40,7 +50,11 @@ import ( type Declaration struct { Info proto.SupportedAgentKind Configuration harnessconfig.Configuration - Discover func(context.Context, DiscoveryOptions, proto.SupportedAgentKind) *Runtime + // ConnectionOptions lists the AgentOptions keys whose values carry MCP + // servers, endpoints, credentials or environment values. An agent-host + // view rejects a request that sets any of them with ErrViewHandoff. + ConnectionOptions []string + Discover func(context.Context, DiscoveryOptions, proto.SupportedAgentKind) *Runtime } // DiscoveryOptions provides process context without naming an implementation. @@ -61,6 +75,9 @@ type Runtime struct { WorkspaceReadPreparation bool SessionCapabilityContext bool ExecutorCapabilityContext bool + // View declares how the Harness runs in an agent-host Session view. + // A nil View means the agent host rejects the kind with ErrUnsupportedOperation. + View *View } // Register installs a discovered Runtime with its declaration's configuration. @@ -75,6 +92,354 @@ func (r *Registry) Register(declaration Declaration, runtime Runtime) { if runtime.Preparation != nil { r.RegisterPreparation(runtime.Info.Kind, runtime.WorkspaceReadPreparation, runtime.Preparation) } + if runtime.View != nil { + r.RegisterView(runtime.Info.Kind, declaration.ConnectionOptions, *runtime.View) + } +} + +// Agent-host Session view. The agent host runs the Harness in a per-Session +// view: the sandbox world at /, the closure, home and shims under +// ViewPrivateRoot, and a loopback-only network whose model, MCP and proxy +// endpoints belong to the Session's credential gateway. The declaration is +// data; the agent host builds each view from it and the Session. + +// The view layout. This is its one definition: sessionview builds views from +// it, and View.Validate keeps declarations out of the trees it reserves. +const ( + // ViewPrivateRoot holds the closure mounts, the shim directory, the home + // and the run directory. + ViewPrivateRoot = "/.oac" + // ViewShimName is the shim directory under ViewPrivateRoot. + ViewShimName = "bin" + // ViewHomeName is the Session home under ViewPrivateRoot. + ViewHomeName = "home" + // ViewRunName is the process broker's socket directory under ViewPrivateRoot. + ViewRunName = "run" + // ViewProcRoot and ViewDevRoot are the view's own /proc and minimal /dev. + ViewProcRoot = "/proc" + ViewDevRoot = "/dev" +) + +// ViewReserved reports whether the view path p is at or beneath a tree the +// view builds itself: ViewPrivateRoot, ViewProcRoot or ViewDevRoot. +func ViewReserved(p string) bool { + for _, root := range [...]string{ViewPrivateRoot, ViewProcRoot, ViewDevRoot} { + if p == root || isWithin(p, root) { + return true + } + } + return false +} + +// viewOwnedEnv are the variables the view or the process broker sets for a +// forwarded process; ForwardEnv never names them. Proxy variables match in +// any case. +var ( + viewOwnedEnv = []string{"HOME", "PATH", "TMPDIR", "LANG", "LD_LIBRARY_PATH"} + viewOwnedProxy = []string{"HTTP_PROXY", "HTTPS_PROXY", "ALL_PROXY", "NO_PROXY"} +) + +// ErrInvalidView marks a View declaration that View.Validate rejects. +var ErrInvalidView = errors.New("agent: invalid view declaration") + +// ErrViewHandoff rejects a view request that carries a model provider other +// than the Session's gateway, MCP outside ViewSession.MCP, an MCP credential or +// a connection option. +var ErrViewHandoff = errors.New("agent: view request carries a connection outside the Session's gateway") + +// View declares how the Harness runs in an agent-host Session view. View +// paths are absolute and clean. The closure, Exec overlays and the shim are +// the only executable mounts, and all are read-only; the world, the home and +// every other mount are noexec. +type View struct { + // Closure lists the host directories presented read-only and executable + // at ViewPrivateRoot/. + Closure []ViewMount + // Overlays present trusted host files or directories read-only at view paths. + Overlays []ViewOverlay + // Masks present view paths empty and read-only. + Masks []ViewMask + // LocalExec lists every view path the Harness process tree executes + // locally. Each lies in the closure or in an Exec overlay. + LocalExec []string + // Shims are names on ViewPrivateRoot/ViewShimName. Each runs that name on + // the Environment's tool PATH in the sandbox. + Shims []string + // ShimPaths are view paths the shim is bound over. Each runs the same + // path in the sandbox. + ShimPaths []string + // ForwardEnv names the Harness variables a forwarded process keeps. It + // never names a variable the view or the broker sets: HOME, PATH, TMPDIR, + // LANG, LD_LIBRARY_PATH or a proxy variable. The Environment's tool + // environment wins over a forwarded variable of the same name. + ForwardEnv []string + Proxy ViewProxy + Executor ViewExecutorFactory +} + +// ViewMount presents HostDir at ViewPrivateRoot/. +type ViewMount struct { + Name string + HostDir string +} + +// Path returns the mount's view path. +func (m ViewMount) Path() string { return ViewPrivateRoot + "/" + m.Name } + +// ViewOverlay presents the trusted host file or directory Source at Path. +// Exec makes it executable, as the ELF interpreter of a dynamic closure +// binary must be. +type ViewOverlay struct { + Path string + Source string + Exec bool +} + +// ViewMask presents Path as an empty directory when Dir is set, otherwise as +// an empty file. +type ViewMask struct { + Path string + Dir bool +} + +// ViewProxy declares how the Harness reaches the network beyond its model and +// MCP endpoints. The zero value is invalid. +type ViewProxy uint8 + +const ( + // ViewProxyNone gives the view no generic proxy. Admission rejects a + // request that enables a feature needing one. + ViewProxyNone ViewProxy = iota + 1 + // ViewProxyEnv gives the view a generic proxy. The adapter has qualified + // that every request its Harness makes locally honours HTTPS_PROXY and + // HTTP_PROXY. + ViewProxyEnv +) + +// ViewExecutorFactory prepares the Session's Executor in its view. The agent +// host has already pointed the request's model provider at the Session's +// gateway, with the placeholder in place of the key, and moved its MCP into +// ViewSession.MCP: the request carries neither MCPHTTPServers nor +// LocalEnvironment.MCP. +type ViewExecutorFactory func(context.Context, proto.PromptRequestPayload, ViewSession) (Executor, error) + +// ViewSession is what the agent host gives a view Executor factory. +type ViewSession struct { + // Home is the per-Session native home, read-write and noexec in the view. + // It persists across the Session's Executors. + Home ViewDir + // Proxy is http://127.0.0.1: for ViewProxyEnv and empty for + // ViewProxyNone. + Proxy string + // MCP is the Session's effective MCP, resolved once from the public + // declarations and the installed Environment MCP. Each HTTP binding's + // ServerURL is its loopback gateway URL, and it carries no BearerToken and + // no HTTPHeaders; the gateway adds them. A stdio binding is as resolved and + // runs in the sandbox through the declared shims. A view Executor takes MCP + // only from here. + MCP []MCPBinding + // Launch replaces clirunner.Start. Each call builds one view and runs + // Binary, which must be a LocalExec path, in it. Dir is a world path, + // OwnProcessGroup is true, and Env is the complete Harness environment. + // Cancel sends TERM to every process in the view and closes the view after + // KillTimeout; a Cancel after the Harness exited leaves its exit as it was. + // When the Harness exits while other processes remain, the view sends them + // TERM unless Cancel already did, and ends once they exit or KillTimeout + // passes from the first TERM. + Launch func(clirunner.StartOptions) (*clirunner.Process, error) +} + +// checkViewHandoff enforces, before the factory runs, that the view request +// reaches the network only through the Session's gateway: the model provider +// is the gateway with the placeholder key, MCP arrives only in session.MCP and +// without credentials, and no connection option is set. +func checkViewHandoff(req proto.PromptRequestPayload, prepared harnessconfig.PreparedConfiguration, connection []string, session ViewSession) error { + if provider := prepared.Provider; provider == nil || provider.APIKey != modelprovider.Placeholder || !isGatewayURL(provider.BaseURL, false) { + return fmt.Errorf("%w: the model provider is not the Session's gateway", ErrViewHandoff) + } + for _, key := range connection { + if _, set := req.AgentOptions[key]; set { + return fmt.Errorf("%w: option %q", ErrViewHandoff, key) + } + } + if req.MCPHTTPServers != nil || (req.LocalEnvironment != nil && len(req.LocalEnvironment.MCP) > 0) { + return fmt.Errorf("%w: MCP outside ViewSession.MCP", ErrViewHandoff) + } + for _, binding := range session.MCP { + if binding.BearerToken != nil || len(binding.HTTPHeaders) > 0 || (binding.Transport == "http" && !isGatewayURL(binding.ServerURL, true)) { + return fmt.Errorf("%w: MCP binding %q is not a credential-free gateway endpoint", ErrViewHandoff, binding.ServerLabel) + } + } + return nil +} + +// isGatewayURL reports whether raw is a plain HTTP URL on a loopback address +// and port, with a path only when withPath is set. +func isGatewayURL(raw string, withPath bool) bool { + endpoint, err := url.Parse(raw) + if err != nil || endpoint.Scheme != "http" || endpoint.User != nil || endpoint.Port() == "" || endpoint.RawQuery != "" || endpoint.Fragment != "" || (!withPath && endpoint.Path != "") { + return false + } + addr, err := netip.ParseAddr(endpoint.Hostname()) + return err == nil && addr.IsLoopback() +} + +// ViewDir is one directory as the adapter writes it on the host and as the +// Harness sees it in the view. +type ViewDir struct { + Host string + View string +} + +// Validate checks the declaration without touching the host. +func (v View) Validate() error { + if v.Executor == nil { + return invalidView("executor is required") + } + if v.Proxy != ViewProxyNone && v.Proxy != ViewProxyEnv { + return invalidView("proxy %d", v.Proxy) + } + names := map[string]bool{ViewShimName: true, ViewHomeName: true, ViewRunName: true} + for _, m := range v.Closure { + if !isPathComponent(m.Name) || names[m.Name] { + return invalidView("closure name %q", m.Name) + } + names[m.Name] = true + if !isHostPath(m.HostDir) { + return invalidView("closure %s host directory %q", m.Name, m.HostDir) + } + } + claimed := slices.Clone(v.ShimPaths) + for _, o := range v.Overlays { + if !isHostPath(o.Source) { + return invalidView("overlay %s source %q", o.Path, o.Source) + } + claimed = append(claimed, o.Path) + } + for _, m := range v.Masks { + claimed = append(claimed, m.Path) + } + for i, p := range claimed { + if !isViewPath(p) || p == "/" || ViewReserved(p) { + return invalidView("view path %q", p) + } + for _, q := range claimed[:i] { + if p == q || isWithin(p, q) || isWithin(q, p) { + return invalidView("view paths %s and %s overlap", q, p) + } + } + } + for i, p := range v.LocalExec { + if !isViewPath(p) || slices.Contains(v.LocalExec[:i], p) || !v.executable(p) { + return invalidView("local exec %q is not a unique path in the closure or an Exec overlay", p) + } + } + for i, n := range v.Shims { + if !isPathComponent(n) || slices.Contains(v.Shims[:i], n) { + return invalidView("shim %q", n) + } + } + for i, n := range v.ForwardEnv { + if n == "" || strings.ContainsAny(n, "=\x00") || slices.Contains(v.ForwardEnv[:i], n) || viewOwnsEnv(n) { + return invalidView("forwarded variable %q", n) + } + } + return nil +} + +func (v View) executable(p string) bool { + for _, m := range v.Closure { + if isWithin(p, m.Path()) { + return true + } + } + for _, o := range v.Overlays { + if o.Exec && (p == o.Path || isWithin(p, o.Path)) { + return true + } + } + return false +} + +func (v View) clone() View { + v.Closure = slices.Clone(v.Closure) + v.Overlays = slices.Clone(v.Overlays) + v.Masks = slices.Clone(v.Masks) + v.LocalExec = slices.Clone(v.LocalExec) + v.Shims = slices.Clone(v.Shims) + v.ShimPaths = slices.Clone(v.ShimPaths) + v.ForwardEnv = slices.Clone(v.ForwardEnv) + return v +} + +func viewOwnsEnv(name string) bool { + return slices.Contains(viewOwnedEnv, name) || slices.ContainsFunc(viewOwnedProxy, func(proxy string) bool { return strings.EqualFold(name, proxy) }) +} + +func invalidView(format string, args ...any) error { + return fmt.Errorf("%w: %s", ErrInvalidView, fmt.Sprintf(format, args...)) +} + +func isViewPath(p string) bool { + return strings.HasPrefix(p, "/") && path.Clean(p) == p && !strings.ContainsRune(p, 0) +} + +func isHostPath(p string) bool { + return filepath.IsAbs(p) && filepath.Clean(p) == p && !strings.ContainsRune(p, 0) +} + +func isPathComponent(name string) bool { + return name != "" && name != "." && name != ".." && !strings.ContainsAny(name, "/\x00") +} + +// isWithin reports whether p is strictly beneath dir. +func isWithin(p, dir string) bool { + return strings.HasPrefix(p, dir+"/") +} + +// RegisterView validates and installs the kind's agent-host view after +// RegisterKind. connection is the declaration's ConnectionOptions. Its +// Executor factory validates the model configuration like RegisterExecutor and +// then checks the request against the view's gateway rule. +func (r *Registry) RegisterView(kind string, connection []string, view View) { + if err := view.Validate(); err != nil { + panic(err) + } + r.mu.Lock() + defer r.mu.Unlock() + configuration, declared := r.configurations[kind] + if !declared { + panic("agent.Registry.RegisterView: registered kind required") + } + view = view.clone() + connection = slices.Clone(connection) + factory := view.Executor + view.Executor = func(ctx context.Context, req proto.PromptRequestPayload, session ViewSession) (Executor, error) { + prepared, err := configuration.Prepare(req.AgentOptions) + if err != nil { + return nil, err + } + if err := checkViewHandoff(req, prepared, connection, session); err != nil { + return nil, err + } + return factory(ctx, req, session) + } + r.views[kind] = view +} + +// ResolveView returns the kind's agent-host view. A kind registered without +// one wraps ErrUnsupportedOperation. +func (r *Registry) ResolveView(kind string) (View, error) { + r.mu.RLock() + defer r.mu.RUnlock() + if _, ok := r.kinds[kind]; !ok { + return View{}, fmt.Errorf("%w: %q", ErrUnsupportedKind, kind) + } + view, ok := r.views[kind] + if !ok { + return View{}, fmt.Errorf("%w: %q declares no agent-host view", ErrUnsupportedOperation, kind) + } + return view.clone(), nil } // Model configuration has one shared contract, authored in @@ -268,6 +633,7 @@ func (r *Registry) RegisterKind(info proto.SupportedAgentKind, configuration har } delete(r.preparers, kind) delete(r.executors, kind) + delete(r.views, kind) info.Capabilities.Preparation = proto.CapabilityUnsupported info.Capabilities.WorkspaceReadPreparation = proto.CapabilityUnsupported r.kinds[kind] = info diff --git a/apps/daemon/internal/agent/mcode/declaration.go b/apps/daemon/internal/agent/mcode/declaration.go index bfab385d3..69209159f 100644 --- a/apps/daemon/internal/agent/mcode/declaration.go +++ b/apps/daemon/internal/agent/mcode/declaration.go @@ -42,13 +42,13 @@ var Declaration = agent.Declaration{Info: proto.SupportedAgentKind{Kind: "mcode" MCPHTTPTools: proto.CapabilityUnsupported, MCPHTTPRequired: proto.CapabilityUnsupported, MCPHTTPBearerAuth: proto.CapabilityUnsupported, -}}, Configuration: configuration.Configuration(), Discover: discover} +}}, Configuration: configuration.Configuration(), ConnectionOptions: []string{"mcp_servers", "env"}, Discover: discover} func discover(ctx context.Context, options agent.DiscoveryOptions, info proto.SupportedAgentKind) *agent.Runtime { return discoverWithCheck(ctx, options, info, CheckCLIAvailable) } func discoverWithCheck(parent context.Context, options agent.DiscoveryOptions, result proto.SupportedAgentKind, check func(context.Context, string) (string, error)) *agent.Runtime { - runtime := &agent.Runtime{Info: result, Session: Factory, SessionCapabilityContext: true, ExecutorCapabilityContext: true} + runtime := &agent.Runtime{Info: result, Session: Factory, SessionCapabilityContext: true, ExecutorCapabilityContext: true, View: nil} ctx, cancel := context.WithTimeout(parent, 15*time.Second) defer cancel() diff --git a/apps/daemon/internal/agent/registry.go b/apps/daemon/internal/agent/registry.go index 4c6b5e815..ec66a03ff 100644 --- a/apps/daemon/internal/agent/registry.go +++ b/apps/daemon/internal/agent/registry.go @@ -29,6 +29,7 @@ type Registry struct { factories map[string]Factory preparers map[string]PreparationFactory executors map[string]ExecutorFactory + views map[string]View kinds map[string]proto.SupportedAgentKind configurations map[string]harnessconfig.Configuration } @@ -38,6 +39,7 @@ func NewRegistry() *Registry { factories: make(map[string]Factory), preparers: make(map[string]PreparationFactory), executors: make(map[string]ExecutorFactory), + views: make(map[string]View), kinds: make(map[string]proto.SupportedAgentKind), configurations: make(map[string]harnessconfig.Configuration), } diff --git a/apps/daemon/internal/agent/view_test.go b/apps/daemon/internal/agent/view_test.go new file mode 100644 index 000000000..453842971 --- /dev/null +++ b/apps/daemon/internal/agent/view_test.go @@ -0,0 +1,138 @@ +package agent_test + +import ( + "context" + "errors" + "slices" + "testing" + + "github.com/MiniMax-AI/OpenAgentCore/apps/daemon/internal/agent" + "github.com/MiniMax-AI/OpenAgentCore/internal/agentdaemon/proto" + "github.com/MiniMax-AI/OpenAgentCore/internal/agentdaemon/proto/prototest" + "github.com/MiniMax-AI/OpenAgentCore/internal/harnessconfig" + "github.com/MiniMax-AI/OpenAgentCore/internal/modelprovider" +) + +func TestViewValidate(t *testing.T) { + if err := validView(t).Validate(); err != nil { + t.Fatalf("valid view: %v", err) + } + for name, change := range map[string]func(*agent.View){ + "nil executor": func(v *agent.View) { v.Executor = nil }, + "unset proxy": func(v *agent.View) { v.Proxy = 0 }, + "reserved closure name": func(v *agent.View) { v.Closure[0].Name = agent.ViewHomeName }, + "relative closure directory": func(v *agent.View) { v.Closure[0].HostDir = "opt/harness" }, + "local exec outside the closure": func(v *agent.View) { v.LocalExec = append(v.LocalExec, "/usr/bin/node") }, + "local exec in read-only overlay": func(v *agent.View) { v.LocalExec = append(v.LocalExec, "/etc/ssl/certs/tool") }, + "mask inside an overlay": func(v *agent.View) { + v.Masks = append(v.Masks, agent.ViewMask{Path: "/etc/ssl/certs/private", Dir: true}) + }, + "shim path equal to a mask": func(v *agent.View) { v.ShimPaths = append(v.ShimPaths, "/etc/harness") }, + "overlay in the private root": func(v *agent.View) { v.Overlays[1].Path = "/.oac/certs" }, + "mask in /proc": func(v *agent.View) { v.Masks[0].Path = "/proc/cpuinfo" }, + "unclean view path": func(v *agent.View) { v.Masks[0].Path = "/etc/../etc/harness" }, + "duplicate shim": func(v *agent.View) { v.Shims = append(v.Shims, "git") }, + "forwarded assignment": func(v *agent.View) { v.ForwardEnv = append(v.ForwardEnv, "A=B") }, + "forwarded broker variable": func(v *agent.View) { v.ForwardEnv = append(v.ForwardEnv, "PATH") }, + "forwarded proxy variable": func(v *agent.View) { v.ForwardEnv = append(v.ForwardEnv, "https_proxy") }, + } { + view := validView(t) + change(&view) + if err := view.Validate(); !errors.Is(err, agent.ErrInvalidView) { + t.Errorf("%s: Validate = %v, want ErrInvalidView", name, err) + } + } +} + +func TestRegistryResolvesOnlyDeclaredViews(t *testing.T) { + reg := agent.NewRegistry() + declared := validView(t) + for kind, view := range map[string]*agent.View{"with_view": &declared, "without_view": nil} { + info := proto.SupportedAgentKind{Kind: kind, Available: true, Capabilities: prototest.Capabilities(proto.AgentKindCapabilities{})} + reg.Register(agent.Declaration{Info: info}, agent.Runtime{Info: info, Session: stubFactory(kind), View: view}) + } + if _, err := reg.ResolveView("without_view"); !errors.Is(err, agent.ErrUnsupportedOperation) { + t.Fatalf("ResolveView without a view = %v, want ErrUnsupportedOperation", err) + } + view, err := reg.ResolveView("with_view") + if err != nil || !slices.Equal(view.LocalExec, declared.LocalExec) || view.Executor == nil { + t.Fatalf("ResolveView = %+v, %v; want the declared view", view, err) + } +} + +func TestViewExecutorReceivesOnlyGatewayConnections(t *testing.T) { + reg := agent.NewRegistry() + declared := validView(t) + declared.Executor = func(context.Context, proto.PromptRequestPayload, agent.ViewSession) (agent.Executor, error) { + return nil, errReached + } + info := proto.SupportedAgentKind{Kind: "viewed", Available: true, Capabilities: prototest.Capabilities(proto.AgentKindCapabilities{})} + declaration := agent.Declaration{ + Info: info, + Configuration: harnessconfig.Configuration{Providers: []harnessconfig.Provider{{Protocol: string(modelprovider.Responses)}}}, + ConnectionOptions: []string{"mcp_servers"}, + } + reg.Register(declaration, agent.Runtime{Info: info, Session: stubFactory("viewed"), View: &declared}) + view, err := reg.ResolveView("viewed") + if err != nil { + t.Fatal(err) + } + options := func(baseURL, key string, extra map[string]any) map[string]any { + o := map[string]any{"model": "m", "model_provider": map[string]any{"protocol": "responses", "base_url": baseURL, "api_key": key}} + for k, v := range extra { + o[k] = v + } + return o + } + gatewayOptions := options("http://127.0.0.1:4101", modelprovider.Placeholder, nil) + token := "secret" + gateway := agent.MCPBinding{ServerLabel: "docs", Transport: "http", ServerURL: "http://127.0.0.1:4100/mcp/docs"} + withBearer, withHeaders, remote := gateway, gateway, gateway + withBearer.BearerToken = &token + withHeaders.HTTPHeaders = map[string]string{"X-Api-Key": token} + remote.ServerURL = "https://mcp.example.com/docs" + for name, c := range map[string]struct { + req proto.PromptRequestPayload + mcp []agent.MCPBinding + }{ + "gateway": {req: proto.PromptRequestPayload{AgentOptions: gatewayOptions}, mcp: []agent.MCPBinding{gateway}}, + "model key": {req: proto.PromptRequestPayload{AgentOptions: options("http://127.0.0.1:4101", token, nil)}}, + "model endpoint": {req: proto.PromptRequestPayload{AgentOptions: options("https://api.example.com", modelprovider.Placeholder, nil)}}, + "connection option": {req: proto.PromptRequestPayload{AgentOptions: options("http://127.0.0.1:4101", modelprovider.Placeholder, map[string]any{"mcp_servers": map[string]any{}})}}, + "request MCP": {req: proto.PromptRequestPayload{AgentOptions: gatewayOptions, MCPHTTPServers: &[]proto.MCPHTTPServer{}}}, + "installed MCP": {req: proto.PromptRequestPayload{AgentOptions: gatewayOptions, LocalEnvironment: &proto.LocalEnvironment{MCP: []proto.EnvironmentMCP{{}}}}}, + "bearer": {req: proto.PromptRequestPayload{AgentOptions: gatewayOptions}, mcp: []agent.MCPBinding{withBearer}}, + "credential headers": {req: proto.PromptRequestPayload{AgentOptions: gatewayOptions}, mcp: []agent.MCPBinding{withHeaders}}, + "direct MCP": {req: proto.PromptRequestPayload{AgentOptions: gatewayOptions}, mcp: []agent.MCPBinding{remote}}, + } { + want := agent.ErrViewHandoff + if name == "gateway" { + want = errReached + } + if _, err := view.Executor(context.Background(), c.req, agent.ViewSession{MCP: c.mcp}); !errors.Is(err, want) { + t.Errorf("%s: Executor = %v, want %v", name, err, want) + } + } +} + +var errReached = errors.New("factory reached") + +func validView(t *testing.T) agent.View { + host := t.TempDir() + return agent.View{ + Closure: []agent.ViewMount{{Name: "harness", HostDir: host}}, + Overlays: []agent.ViewOverlay{ + {Path: "/lib64/ld-linux-x86-64.so.2", Source: host, Exec: true}, + {Path: "/etc/ssl/certs", Source: host}, + }, + Masks: []agent.ViewMask{{Path: "/etc/harness", Dir: true}}, + LocalExec: []string{"/.oac/harness/bin/node", "/lib64/ld-linux-x86-64.so.2"}, + Shims: []string{"bash", "git"}, + ShimPaths: []string{"/bin/sh"}, + ForwardEnv: []string{"GIT_EDITOR"}, + Proxy: agent.ViewProxyEnv, + Executor: func(context.Context, proto.PromptRequestPayload, agent.ViewSession) (agent.Executor, error) { + return nil, errors.New("not started") + }, + } +} diff --git a/apps/daemon/internal/sessionview/build_linux.go b/apps/daemon/internal/sessionview/build_linux.go index c148ee2f3..dab12b6b7 100644 --- a/apps/daemon/internal/sessionview/build_linux.go +++ b/apps/daemon/internal/sessionview/build_linux.go @@ -8,6 +8,8 @@ import ( "strings" "golang.org/x/sys/unix" + + "github.com/MiniMax-AI/OpenAgentCore/apps/daemon/internal/agent" ) // builder mounts the local pieces onto the world. Every target is resolved beneath its parent mount without following symlinks, so the sandbox cannot redirect a mount. @@ -44,7 +46,7 @@ func (b *builder) build(spec *launchSpec) error { if err != nil { return err } - if err := b.bind(src, b.root, privateRoot+"/"+d.Name, bindAttr(d.Writable, d.Exec, false)); err != nil { + if err := b.bind(src, b.root, agent.ViewPrivateRoot+"/"+d.Name, bindAttr(d.Writable, d.Exec, false)); err != nil { return err } } @@ -76,7 +78,7 @@ func (b *builder) build(spec *launchSpec) error { return err } defer unix.Close(proc) - if err := b.attach(proc, b.root, "/proc", true); err != nil { + if err := b.attach(proc, b.root, agent.ViewProcRoot, true); err != nil { return err } return b.dev() @@ -84,7 +86,7 @@ func (b *builder) build(spec *launchSpec) error { // shimDir presents the shim at /.oac/bin/ on a read-only tmpfs. func (b *builder) shimDir(names []string, shim int) error { - dir := privateRoot + "/" + shimName + dir := agent.ViewPrivateRoot + "/" + agent.ViewShimName mnt, err := newFS("tmpfs", [][2]string{{"mode", "0755"}, {"size", "64k"}}, attrNoSuid|attrNoDev|attrNoExec) if err != nil { return err @@ -111,7 +113,7 @@ func (b *builder) dev() error { return err } defer unix.Close(mnt) - if err := b.attach(mnt, b.root, "/dev", true); err != nil { + if err := b.attach(mnt, b.root, agent.ViewDevRoot, true); err != nil { return err } for _, n := range devNodes { @@ -153,7 +155,7 @@ func (b *builder) dev() error { return &Error{Kind: ErrLauncher, Op: "symlink", Path: "/dev/" + l[0], Err: err} } } - return readOnly(mnt, "/dev") + return readOnly(mnt, agent.ViewDevRoot) } // source opens a trusted host path. diff --git a/apps/daemon/internal/sessionview/control_linux.go b/apps/daemon/internal/sessionview/control_linux.go index 694be509e..266dc19c5 100644 --- a/apps/daemon/internal/sessionview/control_linux.go +++ b/apps/daemon/internal/sessionview/control_linux.go @@ -12,6 +12,7 @@ import ( "os" "sync" "syscall" + "time" "golang.org/x/sys/unix" ) @@ -38,25 +39,28 @@ type launchSpec struct { UID uint32 GID uint32 Groups []uint32 + Grace time.Duration } type msgKind uint8 const ( - msgMounted msgKind = iota + 1 // launcher: the world is mounted; carries the /dev/fuse and netns fds - msgProceed // daemon: the world serves and the network is set up - msgStarted // launcher: the process runs - msgFailed // launcher: construction failed - msgExited // launcher: the process ended - msgSignal // daemon: signal the process + msgMounted msgKind = iota + 1 // launcher: the world is mounted; carries the /dev/fuse and netns fds + msgProceed // daemon: the world serves and the network is set up + msgStarted // launcher: the process runs + msgFailed // launcher: construction failed + msgExited // launcher: the process ended + msgSignal // daemon: signal the process + msgSignaled // launcher: whether msgSignal reached the running process ) type message struct { - Kind msgKind - Pid int - Signal syscall.Signal - Exit Exit - Fail failure + Kind msgKind + Pid int + Signal syscall.Signal + Delivered bool + Exit Exit + Fail failure } // failure carries a launcher *Error across the control socket. diff --git a/apps/daemon/internal/sessionview/doc.go b/apps/daemon/internal/sessionview/doc.go index 448b9bdcb..2d6271a2e 100644 --- a/apps/daemon/internal/sessionview/doc.go +++ b/apps/daemon/internal/sessionview/doc.go @@ -4,7 +4,7 @@ // // The exec guard is narrow. Mount flags alone decide which files can be executed: the world and every writable mount are nosuid and noexec, so executing a file from the file system works only from read-only mounts declared executable, such as the Harness directory, the shim and exec-flagged overlays. The seccomp filter denies creating a user namespace and every setns, so the process cannot create or enter another user namespace. Executing from a memfd and code that an allowed interpreter runs are outside this guard. // -// The daemon calls [Init] first thing in main. [Start] re-executes the daemon binary as the launcher, which becomes PID 1 of the view: it builds the view, starts the process, forwards signals, reaps orphans and exits with the process status. Its exit tears the view down. +// The daemon calls [Init] first thing in main. [Start] re-executes the daemon binary as the launcher, which becomes PID 1 of the view: it builds the view, starts the process, delivers signals to every process in the view, reaps orphans and exits with the process status. Once the process has exited, Signal reports [ErrExited] and delivers nothing. When the process exits while others remain, the launcher sends them TERM unless one already went to the view, and waits for them up to [Process].Grace from the first TERM. Its exit kills what remains and tears the view down. // // The package works only on Linux. Elsewhere [Start] returns [ErrUnsupported]. package sessionview diff --git a/apps/daemon/internal/sessionview/errors.go b/apps/daemon/internal/sessionview/errors.go index 6ea36a66c..e7dcc5bab 100644 --- a/apps/daemon/internal/sessionview/errors.go +++ b/apps/daemon/internal/sessionview/errors.go @@ -21,6 +21,8 @@ var ( ErrExec = errors.New("sessionview: process start failed") ErrLauncher = errors.New("sessionview: launcher failed") ErrClosed = errors.New("sessionview: view closed") + // ErrExited is Signal's result once the process has exited. It also matches os.ErrProcessDone. + ErrExited = errors.New("sessionview: process exited") ) // errorKinds fixes the wire code of each kind the launcher reports. diff --git a/apps/daemon/internal/sessionview/launcher_linux.go b/apps/daemon/internal/sessionview/launcher_linux.go index 64cb22df6..5b8fffc0e 100644 --- a/apps/daemon/internal/sessionview/launcher_linux.go +++ b/apps/daemon/internal/sessionview/launcher_linux.go @@ -10,8 +10,8 @@ import ( "runtime" "strings" "sync" - "sync/atomic" "syscall" + "time" "golang.org/x/sys/unix" ) @@ -31,7 +31,11 @@ type launcher struct { ctl *control proceed chan struct{} proceedOnce sync.Once - pidfd atomic.Int64 + + // mu orders signals against the process's exit. running holds from the process's start until it is reaped; termAt is when TERM first went to the view. + mu sync.Mutex + running bool + termAt time.Time } func runLauncher() int { @@ -41,7 +45,6 @@ func runLauncher() int { return 1 } l := &launcher{ctl: ctl, proceed: make(chan struct{})} - l.pidfd.Store(-1) code, err := l.run() if err != nil { _ = ctl.send(message{Kind: msgFailed, Fail: failureOf(err)}) @@ -82,11 +85,9 @@ func (l *launcher) run() (int, error) { if err != nil { return 0, err } - pidfd, err := unix.PidfdOpen(pid, 0) - if err != nil { - return 0, &Error{Kind: ErrExec, Op: "pidfd_open", Err: err} - } - l.pidfd.Store(int64(pidfd)) + l.mu.Lock() + l.running = true + l.mu.Unlock() for _, fd := range []int{stdinFD, stdoutFD, stderrFD} { unix.Close(fd) } @@ -94,7 +95,11 @@ func (l *launcher) run() (int, error) { return 0, &Error{Kind: ErrLauncher, Op: "report start", Err: err} } l.forwardSignals() - return l.reap(pid) + code, err := l.reap(pid) + if err == nil { + l.drain(spec.Grace) + } + return code, err } func readSpec() (*launchSpec, error) { @@ -119,15 +124,23 @@ func (l *launcher) serveControl() { case msgProceed: l.proceedOnce.Do(func() { close(l.proceed) }) case msgSignal: - l.signal(m.Signal) + _ = l.ctl.send(message{Kind: msgSignaled, Delivered: l.signal(m.Signal)}) } } } -func (l *launcher) signal(sig syscall.Signal) { - if fd := l.pidfd.Load(); fd >= 0 { - _ = unix.PidfdSendSignal(int(fd), sig, nil, 0) +// signal delivers sig to every process in the view while the process runs and reports whether it did. As PID 1 of the view, the launcher reaches them all with kill(-1) and is itself spared. Once the process has been reaped, only the drain signals what remains. +func (l *launcher) signal(sig syscall.Signal) bool { + l.mu.Lock() + defer l.mu.Unlock() + if !l.running { + return false } + if sig == unix.SIGTERM && l.termAt.IsZero() { + l.termAt = time.Now() + } + _ = unix.Kill(-1, sig) + return true } func (l *launcher) forwardSignals() { @@ -140,6 +153,34 @@ func (l *launcher) forwardSignals() { }() } +// drain gives the processes left after the process exits TERM and up to grace from that TERM to exit, reaping them. When TERM already went to the view, they keep what remains of the grace and get no second TERM. The launcher's exit then kills whatever remains. +func (l *launcher) drain(grace time.Duration) { + l.mu.Lock() + termAt := l.termAt + l.mu.Unlock() + if !termAt.IsZero() { + grace -= time.Since(termAt) + } + if grace <= 0 || (termAt.IsZero() && unix.Kill(-1, unix.SIGTERM) != nil) { + return + } + reaped := make(chan struct{}) + go func() { + defer close(reaped) + for { + if _, err := unix.Wait4(-1, nil, 0, nil); err != nil && err != unix.EINTR { + return + } + } + }() + timer := time.NewTimer(grace) + defer timer.Stop() + select { + case <-reaped: + case <-timer.C: + } +} + // reap collects every child, since orphans in the view reparent to PID 1, until the process exits. func (l *launcher) reap(pid int) (int, error) { for { @@ -154,6 +195,9 @@ func (l *launcher) reap(pid int) (int, error) { if wpid != pid { continue } + l.mu.Lock() + l.running = false + l.mu.Unlock() var exit Exit code := ws.ExitStatus() if ws.Signaled() { diff --git a/apps/daemon/internal/sessionview/spec.go b/apps/daemon/internal/sessionview/spec.go index f8773f53c..7c28774de 100644 --- a/apps/daemon/internal/sessionview/spec.go +++ b/apps/daemon/internal/sessionview/spec.go @@ -6,6 +6,9 @@ import ( "path/filepath" "strings" "syscall" + "time" + + "github.com/MiniMax-AI/OpenAgentCore/apps/daemon/internal/agent" ) // Spec declares one view and the process it runs. @@ -16,6 +19,8 @@ type Spec struct { Shim Shim Process Process Network Network + // StagingParent is an existing absolute host directory, such as the Session directory, in which the view creates its staging directory and removes it at teardown. + StagingParent string } // World starts serving the view's root file system on dev, a /dev/fuse connection that the launcher has already mounted as mount describes. The server runs in the daemon, outside the view, and never mounts or unmounts anything. sessionview closes dev after Stop returns. @@ -43,7 +48,7 @@ type Mountpoint struct { Dir bool } -// PrivateDir binds a host directory at /.oac/. A writable directory is never executable. +// PrivateDir binds a host directory at agent.ViewPrivateRoot/. A writable directory is never executable. type PrivateDir struct { Name string HostDir string @@ -58,7 +63,7 @@ type Overlay struct { Exec bool } -// Shim presents a static binary at /.oac/bin/ for each name and binds it over each absolute view path. +// Shim presents a static binary at agent.ViewPrivateRoot/agent.ViewShimName/ for each name and binds it over each absolute view path. type Shim struct { Binary string Names []string @@ -76,6 +81,8 @@ type Process struct { Groups []uint32 // A nil Stdin, Stdout or Stderr is a pipe whose other end the View exposes. Stdin, Stdout, Stderr *os.File + // Grace is how long processes left in the view when the process exits have to exit, counted from the first TERM the view sent them, before the view ends. Zero ends the view at once. + Grace time.Duration } // Network configures the view's network namespace, which has only loopback up. @@ -91,21 +98,18 @@ type Exit struct { CoreDumped bool } -const ( - privateRoot = "/.oac" - shimName = "bin" -) - -// reservedTrees are built by the launcher; overlays and shim paths stay out of them. -var reservedTrees = []string{privateRoot, "/proc", "/dev"} - func (s *Spec) validate() error { if s.World == nil { return invalid("world is required") } + if info, err := hostSource(s.StagingParent); err != nil { + return err + } else if !info.IsDir() { + return invalid("staging parent %s is not a directory", s.StagingParent) + } names := map[string]bool{} for _, d := range s.Private { - if !isComponent(d.Name) || d.Name == shimName || names[d.Name] { + if !isComponent(d.Name) || d.Name == agent.ViewShimName || names[d.Name] { return invalid("private directory name %q", d.Name) } names[d.Name] = true @@ -165,6 +169,9 @@ func (p *Process) validate() error { if len(p.Args) == 0 { return invalid("process argv is empty") } + if p.Grace < 0 { + return invalid("process grace %v is negative", p.Grace) + } if p.UID == 0 || p.GID == 0 { return invalid("process must not run as uid or gid 0") } @@ -185,10 +192,8 @@ func viewPath(p string) error { if !isViewAbs(p) || p == "/" { return invalid("view path %q must be absolute, clean and not /", p) } - for _, r := range reservedTrees { - if p == r || within(p, r) { - return invalid("view path %s is inside %s", p, r) - } + if agent.ViewReserved(p) { + return invalid("view path %s is inside a tree the launcher builds", p) } return nil } diff --git a/apps/daemon/internal/sessionview/view_linux.go b/apps/daemon/internal/sessionview/view_linux.go index 1e31c87da..21b302943 100644 --- a/apps/daemon/internal/sessionview/view_linux.go +++ b/apps/daemon/internal/sessionview/view_linux.go @@ -15,6 +15,8 @@ import ( "syscall" "golang.org/x/sys/unix" + + "github.com/MiniMax-AI/OpenAgentCore/apps/daemon/internal/agent" ) // launcherArg0 marks the re-executed daemon binary as a launcher. @@ -39,6 +41,9 @@ type View struct { reapOnce sync.Once reapErr error + signalMu sync.Mutex + signaled chan bool + exited atomic.Bool closing atomic.Bool closeOnce sync.Once done chan struct{} @@ -54,7 +59,7 @@ func Start(ctx context.Context, spec Spec) (*View, error) { if err := Probe(); err != nil { return nil, err } - v := &View{done: make(chan struct{})} + v := &View{done: make(chan struct{}), signaled: make(chan bool, 1)} if err := v.launch(&spec); err != nil { v.abort() return nil, err @@ -85,7 +90,7 @@ func (v *View) launch(spec *Spec) error { } } }() - v.staging, err = os.MkdirTemp("", "oac-view-*") + v.staging, err = os.MkdirTemp(spec.StagingParent, "oac-view-*") if err != nil { return &Error{Kind: ErrLauncher, Op: "staging", Err: err} } @@ -125,7 +130,7 @@ func (v *View) launch(spec *Spec) error { ls := &launchSpec{ Staging: v.staging, Private: spec.Private, Overlays: spec.Overlays, Shim: spec.Shim, Path: spec.Process.Path, Args: spec.Process.Args, Env: spec.Process.Env, Dir: spec.Process.Dir, - UID: spec.Process.UID, GID: spec.Process.GID, Groups: spec.Process.Groups, + UID: spec.Process.UID, GID: spec.Process.GID, Groups: spec.Process.Groups, Grace: spec.Process.Grace, } // A launcher that dies early breaks the pipe; the handshake reports that. go func() { @@ -267,6 +272,9 @@ func (v *View) watch() { switch m.Kind { case msgExited: exit = &m.Exit + v.exited.Store(true) + case msgSignaled: + v.signaled <- m.Delivered case msgFailed: failed = m.Fail.err() } @@ -294,8 +302,14 @@ func (v *View) Wait() (Exit, error) { return v.exit, v.err } -// Signal delivers sig to the process. +// Signal delivers sig to every process in the view while the process runs. Once the process has exited it delivers nothing and returns ErrExited, even while the processes it left still drain. func (v *View) Signal(sig syscall.Signal) error { + v.signalMu.Lock() + defer v.signalMu.Unlock() + exited := &Error{Kind: ErrExited, Op: "signal", Err: os.ErrProcessDone} + if v.exited.Load() { + return exited + } select { case <-v.done: return ErrClosed @@ -304,7 +318,18 @@ func (v *View) Signal(sig syscall.Signal) error { if err := v.ctl.send(message{Kind: msgSignal, Signal: sig}); err != nil { return &Error{Kind: ErrLauncher, Op: "signal", Err: err} } - return nil + select { + case delivered := <-v.signaled: + if !delivered { + return exited + } + return nil + case <-v.done: + if v.exited.Load() { + return exited + } + return ErrClosed + } } // Close kills the view, waits for its teardown and closes the pipes it created. @@ -334,12 +359,12 @@ func (v *View) Stderr() *os.File { return v.pipes[2] } func (s *Spec) mountpoints() []Mountpoint { var m []Mountpoint for _, d := range s.Private { - m = append(m, Mountpoint{Path: privateRoot + "/" + d.Name, Dir: true}) + m = append(m, Mountpoint{Path: agent.ViewPrivateRoot + "/" + d.Name, Dir: true}) } m = append(m, - Mountpoint{Path: privateRoot + "/" + shimName, Dir: true}, - Mountpoint{Path: "/proc", Dir: true}, - Mountpoint{Path: "/dev", Dir: true}, + Mountpoint{Path: agent.ViewPrivateRoot + "/" + agent.ViewShimName, Dir: true}, + Mountpoint{Path: agent.ViewProcRoot, Dir: true}, + Mountpoint{Path: agent.ViewDevRoot, Dir: true}, ) for _, o := range s.Overlays { info, err := os.Stat(o.Source) diff --git a/apps/daemon/internal/sessionview/view_linux_test.go b/apps/daemon/internal/sessionview/view_linux_test.go index d1672e27a..328ee8f00 100644 --- a/apps/daemon/internal/sessionview/view_linux_test.go +++ b/apps/daemon/internal/sessionview/view_linux_test.go @@ -106,6 +106,9 @@ func TestViewSignalAndTeardown(t *testing.T) { if n := processesWith(t, token); n != 1 { t.Fatalf("%d grandchildren before exit, want 1", n) } + if staged, _ := os.ReadDir(f.staging); len(staged) != 1 { + t.Fatalf("staging parent holds %v, want the view's staging directory", staged) + } if err := v.Signal(syscall.SIGTERM); err != nil { t.Fatalf("Signal: %v", err) } @@ -120,11 +123,51 @@ func TestViewSignalAndTeardown(t *testing.T) { default: t.Error("world server still serving") } - if left, _ := filepath.Glob(filepath.Join(os.TempDir(), "oac-view-*")); len(left) != 0 { + if left, _ := os.ReadDir(f.staging); len(left) != 0 { t.Errorf("staging directories left: %v", left) } } +// TestViewDescendantsKeepTheGrace checks that a helper still cleaning up when the Harness exits gets TERM and finishes within the grace. +func TestViewDescendantsKeepTheGrace(t *testing.T) { + requireView(t) + f := newFixture(t) + w := &loopbackWorld{dir: f.world} + spec := f.spec(w, "cleanup") + spec.Process.Grace = 10 * time.Second + v, err := Start(context.Background(), spec) + if err != nil { + t.Fatalf("Start: %v", err) + } + defer v.Close() + if line, err := bufio.NewReader(v.Stdout()).ReadString('\n'); err != nil || line != "ready\n" { + t.Fatalf("helper said %q, %v", line, err) + } + started := time.Now() + if err := v.Signal(syscall.SIGTERM); err != nil { + t.Fatalf("Signal: %v", err) + } + // The Harness's exit shows while its helper still cleans up. + for err := v.Signal(0); !errors.Is(err, os.ErrProcessDone); err = v.Signal(0) { + if err != nil || time.Since(started) > 5*time.Second { + t.Fatalf("Signal after the Harness exited = %v, want os.ErrProcessDone", err) + } + time.Sleep(5 * time.Millisecond) + } + if _, err := os.Stat(filepath.Join(f.world, "data", "cleaned")); err == nil { + t.Error("Signal reported the exit only after the helper finished") + } + if exit, err := v.Wait(); err != nil || exit != (Exit{Code: 7}) { + t.Fatalf("Wait = %+v, %v; want exit code 7", exit, err) + } + if elapsed := time.Since(started); elapsed >= spec.Process.Grace { + t.Errorf("view ended after %v, want once the helper exited", elapsed) + } + if got, err := os.ReadFile(filepath.Join(f.world, "data", "cleaned")); err != nil || string(got) != "done" { + t.Errorf("helper cleanup = %q, %v; want it finished", got, err) + } +} + // TestViewRefusesSymlinkedMountpoint checks that a sandbox symlink on the way to a mountpoint fails the view instead of redirecting the mount. func TestViewRefusesSymlinkedMountpoint(t *testing.T) { requireView(t) @@ -191,7 +234,11 @@ func TestStartRejectsInvalidSpec(t *testing.T) { for name, spec := range map[string]Spec{ "relative overlay": {Overlays: []Overlay{{Path: "etc/resolv.conf", Source: "/etc/hosts"}}}, "writable executable private": {Private: []PrivateDir{{Name: "home", HostDir: t.TempDir(), Writable: true, Exec: true}}}, + "missing staging parent": {StagingParent: filepath.Join(t.TempDir(), "missing")}, } { + if spec.StagingParent == "" { + spec.StagingParent = t.TempDir() + } spec.World = (&loopbackWorld{}).serve spec.Process = Process{Path: "/bin/true", Args: []string{"true"}, Dir: "/", UID: viewID, GID: viewID} if _, err := Start(context.Background(), spec); !errors.Is(err, ErrInvalidSpec) { @@ -211,7 +258,7 @@ func requireView(t *testing.T) { } type fixture struct { - self, world, harness, home, run, overlay string + self, world, harness, home, run, overlay, staging string } // newFixture lays out a world with the mountpoints the real world frontend presents synthetically, plus the local sources. @@ -228,6 +275,7 @@ func newFixture(t *testing.T) *fixture { home: filepath.Join(base, "home"), run: filepath.Join(base, "run"), overlay: filepath.Join(base, "overlay"), + staging: filepath.Join(base, "staging"), } for _, d := range []string{".oac/harness", ".oac/home", ".oac/run", ".oac/bin", "proc", "dev", "bin", "usr/bin", "data", "etc/oac-overlay"} { mkdir(t, filepath.Join(f.world, d)) @@ -245,6 +293,7 @@ func newFixture(t *testing.T) *fixture { } } mkdir(t, f.overlay) + mkdir(t, f.staging) writeFile(t, filepath.Join(f.overlay, "greeting"), "from the overlay") return f } @@ -272,6 +321,7 @@ func (f *fixture) spec(w *loopbackWorld, mode string, env ...string) Spec { GID: viewID, Stderr: os.Stderr, }, + StagingParent: f.staging, } } @@ -393,6 +443,32 @@ func runHelper(mode string) int { case "sleep": time.Sleep(time.Hour) return 0 + case "cleanup": + // The Harness exits on TERM at once while its helper still cleans up. + sigs := make(chan os.Signal, 1) + signal.Notify(sigs, syscall.SIGTERM) + child := exec.Command("/.oac/harness/harness") + child.Env = []string{helperEnv + "=slow-term"} + child.Stdout = os.Stdout + if err := child.Start(); err != nil { + fmt.Fprintln(os.Stderr, err) + return 1 + } + <-sigs + return 7 + case "slow-term": + // A second TERM ends it before the cleanup finishes. + sigs := make(chan os.Signal, 1) + signal.Notify(sigs, syscall.SIGTERM) + fmt.Println("ready") + <-sigs + signal.Reset(syscall.SIGTERM) + time.Sleep(300 * time.Millisecond) + if err := os.WriteFile("/data/cleaned", []byte("done"), 0o644); err != nil { + fmt.Fprintln(os.Stderr, err) + return 1 + } + return 0 case "noop": return 0 } diff --git a/contracts/agents-api/harness-onboarding.md b/contracts/agents-api/harness-onboarding.md index 0365dba09..649499676 100644 --- a/contracts/agents-api/harness-onboarding.md +++ b/contracts/agents-api/harness-onboarding.md @@ -147,7 +147,7 @@ A Harness that supports the Subagent reads implements the [neutral observation c ## Register the adapter -Registration is static and requires a build. Export one `agent.Declaration` from `apps/daemon/internal/agent//declaration.go`, then add it to `harnessDeclarations` in [`cli/agent_discovery.go`](../../apps/daemon/internal/cli/agent_discovery.go). The declaration contains the kind and complete capability descriptor, the shared model `Configuration` and a `Discover` function. Discovery receives the profile and diagnostic writers, owns native configuration and availability checks, and returns the installed `agent.Runtime` with its descriptor and session, preparation and Executor factories. Return nil when the adapter is not configured; return an unavailable descriptor with a session factory when configured prerequisites fail. Keep version gates and factory-selection conditions inside the adapter. +Registration is static and requires a build. Export one `agent.Declaration` from `apps/daemon/internal/agent//declaration.go`, then add it to `harnessDeclarations` in [`cli/agent_discovery.go`](../../apps/daemon/internal/cli/agent_discovery.go). The declaration contains the kind and complete capability descriptor, the shared model `Configuration`, the `ConnectionOptions` an agent-host view rejects ([Endpoints and proxy](#endpoints-and-proxy)) and a `Discover` function. Discovery receives the profile and diagnostic writers, owns native configuration and availability checks, and returns the installed `agent.Runtime` with its descriptor and session, preparation and Executor factories. Return nil when the adapter is not configured; return an unavailable descriptor with a session factory when configured prerequisites fail. Keep version gates and factory-selection conditions inside the adapter. `Runtime.SessionCapabilityContext` and `Runtime.ExecutorCapabilityContext` explicitly request capability-download URL resolution and scoped product-upload context for the corresponding execution factory. Preparation never receives those effects. An adapter that supports product workspace authoring declares `WorkspaceAuthoring` itself; common registration does not grant it. @@ -158,9 +158,12 @@ Registration is static and requires a build. Export one `agent.Declaration` from | 1 | `RegisterKind(proto.SupportedAgentKind, harnessconfig.Configuration, agent.Factory)` | Kind, availability, version, `AgentKindCapabilities`, the model configuration declaration and the direct-call factory. It resets the other registrations, so call it first. | | 2 | `RegisterExecutor(kind, agent.ExecutorFactory)` | The Executor and Turn lifecycle used for execution; derives the `Preparation` capability | | 3 | `RegisterPreparation(kind, workspaceRead, agent.PreparationFactory)` | Optional: separate read-only workspace preparation for qualified workspace operations | +| 4 | `RegisterView(kind, connectionOptions, agent.View)` | Optional: the agent-host view declaration from `Runtime.View`, with the declaration's `ConnectionOptions`. It panics with `ErrInvalidView` when `View.Validate` fails. Its Executor factory validates the model configuration like `RegisterExecutor` and enforces the [gateway rule](#endpoints-and-proxy). | The direct-call `agent.Factory` delegates to the same Executor implementation. +`Runtime.View` declares how the Harness runs in an agent-host Session view, described in [Run in an agent-host view](#run-in-an-agent-host-view). Every adapter sets it explicitly; `View: nil` means the agent host rejects the kind, and `Registry.ResolveView` returns an error wrapping `ErrUnsupportedOperation`. `TestPublicHarnessContractDeclarations` requires the field in each declaration. + Every `proto.AgentKindCapabilities` field must be explicitly `proto.CapabilitySupported` or `proto.CapabilityUnsupported`, even for an unavailable Harness. `proto.CapabilityUnspecified` is invalid: zero values and omitted fields never mean Unsupported. An installation probe may set an individual field with `proto.CapabilityFromBool`; it must not populate unmentioned or future fields. Availability stays separate in `SupportedAgentKind.Available`. Registration validates the complete declaration before changing the registry, and the wire carries an explicit boolean for every field, so omitted and null fields are invalid. A new field requires a decision in every production declaration. Runtime consumers use `IsSupported()` and reject unsupported requests before native operations; an interface assertion verifies implementation, never support. Every declaration must match the behavior verified for that installation; the [Core–Runtime protocol](../../docs/runtime-protocol.md#capability-declarations) owns how declarations travel and are frozen. The admission mapping is explicit. `Steering` controls non-durable `Steerer` input. `DurableInputReceipts` controls `DurableSteerer` input and also requires the Turn settlement contract; neither implies the other, and Core's public text profile requires both. `Permissions` qualifies permission and user-choice responses together and requires both native response paths. Workspace declarations describe the authorized resource owner, including the common Runtime workspace implementation. Runtime registration does not grant Core qualification; the service profile does. @@ -253,6 +256,84 @@ Adapters run native tools unattended with the launching user's permissions: Code Network admission follows [Restricted network](environments.md#restricted-network). +## Run in an agent-host view + +An agent host runs the Harness outside the sandbox, in a per-Session view. The view shows the sandbox's files at `/` over the [File access protocol](../../docs/file-access-protocol.md), and the Harness's own files under `/.oac`. Programs the Harness does not declare local run in the sandbox over the [Process protocol](../../docs/process-protocol.md). The network has loopback only, where the Session's credential gateway listens. The adapter declares what its Harness needs in `Runtime.View`, and the agent host builds each view from that declaration and the Session. Support is the declared field: a kind without a View is rejected with `ErrUnsupportedOperation`. + +### The declaration + +| Field | Declares | +| --- | --- | +| `Closure` | Host directories presented read-only and executable at `/.oac/` (`ViewMount.Path`) | +| `Overlays` | Trusted host files or directories presented read-only at view paths; `Exec` makes one executable | +| `Masks` | View paths presented empty and read-only, as a directory when `Dir` is set | +| `LocalExec` | Every view path the Harness process tree executes locally | +| `Shims` | Names on `/.oac/bin`; each runs that name on the Environment's tool `PATH` in the sandbox | +| `ShimPaths` | View paths the shim is bound over; each runs the same path in the sandbox | +| `ForwardEnv` | Harness variables that a process run in the sandbox keeps | +| `Proxy` | `ViewProxyEnv` or `ViewProxyNone` | +| `Executor` | The `ViewExecutorFactory` that prepares the Session's Executor in its view | + +`View.Validate` checks the declaration without touching the host: + +- view and host paths are absolute and clean; +- closure names are single path components other than `bin`, `home` and `run`, which the agent host uses for the shims, the Session home and the process broker; +- shim paths, overlays and masks do not overlap each other or `/`, and stay out of the trees the view builds itself: `/.oac`, `/proc` and `/dev` (`ViewReserved`); +- each `LocalExec` entry lies in a closure directory or an `Exec` overlay; +- shim names and `ForwardEnv` names are unique, a variable name contains no `=`, and `ForwardEnv` names no variable the view or the broker sets ([Environment](#environment)); +- `Proxy` is one of the two values and `Executor` is non-nil. + +`harness.go` defines the view layout once, and `sessionview` builds views from it. The agent host checks its own overlays, such as `/etc/passwd`, against the declaration when it builds the view. + +### Executables + +Only mount flags grant execution. The closure, `Exec` overlays and the shim are read-only and are the only executable mounts; the sandbox's files and the home are noexec. `Launch` accepts only a `LocalExec` path as `Binary`. A dynamic binary, such as `node`, needs its ELF interpreter as an `Exec` overlay at its `PT_INTERP` path, and every library it loads in the closure, reached through `LD_LIBRARY_PATH`. Nothing loads from the sandbox's files. + +### Shims + +The agent host derives the process broker's table from the declaration: `/.oac/bin/` runs `` in the sandbox, and each `ShimPaths` entry runs the same path there. A Harness that runs tools by name finds them through a `PATH` that lists `/.oac/bin`. The view's `/etc/passwd` gives the Session user the login shell `/bin/bash` and the home `/.oac/home`. + +### Environment + +`Launch` takes the complete Harness environment in `StartOptions.Env`. The agent host's own environment never passes through, so a view adapter does not start from `os.Environ()`. A process run in the sandbox gets the broker's environment: the `ForwardEnv` variables from the Harness, the Environment's fixed sandbox values (`HOME`, `PATH`, `TMPDIR` and `LANG`) and the Environment's tool environment. The broker is the only home of the tool environment, and a view adapter passes none of it to the Harness. + +`ForwardEnv` never names a variable the view or the broker sets: `HOME`, `PATH`, `TMPDIR`, `LANG`, `LD_LIBRARY_PATH`, or `HTTP_PROXY`, `HTTPS_PROXY`, `ALL_PROXY` and `NO_PROXY` in any case. When the Environment's tool environment also sets a forwarded variable, the tool environment's value wins. + +### Endpoints and proxy + +Before it calls the factory, the agent host points the request's `model_provider` at the Session's [credential gateway](model-execution.md#credential-gateway): `base_url` is `http://127.0.0.1:` with no path and `api_key` is `modelprovider.Placeholder`. It resolves the Session's MCP once, from the public declarations and the installed Environment MCP, into `ViewSession.MCP`, and removes both from the request. Each HTTP binding points at its gateway URL and carries no bearer and no headers; the gateway adds the declared credential and headers. A stdio binding is as resolved and runs in the sandbox through the declared shims. A view Executor takes MCP only from `ViewSession.MCP` and never resolves the request. The adapter renders the provider and the bindings as it does for a local Harness and never sees a real credential. + +An adapter whose Harness reads MCP servers, endpoints, credentials or environment values from other `AgentOptions` keys lists those keys in `Declaration.ConnectionOptions`. The Registry checks each view request once, before the factory, and rejects it with `ErrViewHandoff` when its model provider is missing or is not the gateway with the placeholder, when it sets a connection option, when it carries MCP outside `ViewSession.MCP`, or when an HTTP binding is not a credential-free loopback endpoint. + +With `ViewProxyEnv`, `ViewSession.Proxy` is the gateway's proxy URL. The adapter sets `HTTPS_PROXY` and `HTTP_PROXY` to it and `NO_PROXY` to `127.0.0.1,localhost`, each in upper and lower case. Declare `ViewProxyEnv` only after qualifying that every request the Harness makes locally honours these variables. A request that ignores them fails to connect, because the view has no route out. + +With `ViewProxyNone`, `ViewSession.Proxy` is empty and the view has no generic proxy. Admission rejects a request that enables a feature needing one with `ErrUnsupportedOperation`. Web tools that the provider executes keep provider origin. + +### Home + +`ViewSession.Home` is the per-Session native home. The adapter writes at `Home.Host`, and the Harness sees the same directory at `Home.View` (`/.oac/home`), read-write and noexec. It persists across the Session's Executors. Lay out native directories and write configuration under it before calling `Launch`. `Launch` gives the tree to the Session user without following links; after that, read the home without following links. + +### Launch + +`ViewSession.Launch` replaces `clirunner.Start`. Each call builds one view and runs `Binary` in it, and at most one view per Session is live at a time. `Dir` is a path in the sandbox, `OwnProcessGroup` is true and `Env` is the complete environment. The returned `clirunner.Process` follows [Native process ownership](#native-process-ownership): + +- Cancel sends TERM to every process in the view and closes the view after `KillTimeout`. A Cancel that finds the Harness exited leaves its exit as it was, even while the processes it left still end. +- When the Harness exits while other processes remain, the view sends them TERM unless Cancel already did, and ends once they exit or `KillTimeout` passes from the first TERM. +- `Wait` closes the stdio ends, returns the context error when Cancel's TERM reached the running Harness and it then exited 0, and `ExitCode` reports the exit once `Done` closes. + +### Qualify the view + +Run the adapter's Turns, cancellation and continuation in a view, then qualify each declared entry: + +| Entry | Qualification | +| --- | --- | +| `Closure`, `Overlays`, `LocalExec` | Every local exec succeeds from a declared path. Each dynamic binary's interpreter overlay matches its `PT_INTERP`, and `LD_LIBRARY_PATH` resolves every library in the closure. | +| `Masks` | The Harness reads none of the sandbox's files at the masked paths. | +| `Shims`, `ShimPaths` | Each tool the Harness runs by name or path runs in the sandbox, and its output, exit status and signals reach the Harness. | +| `ForwardEnv` | A process run in the sandbox keeps each declared variable and no other Harness variable. | +| `Proxy` | With `ViewProxyEnv`, every local request, such as web fetches, downloads and update checks, goes through the proxy. With `ViewProxyNone`, a request enabling a feature that needs it is rejected. | +| `Home` | Native history and configuration stay under `/.oac/home`, and a later Executor in the same Session continues from them. | + ## Native references | Harness | Adapter | Native transport | Runtime guide |