Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
23 commits
Select commit Hold shift + click to select a range
fbff1bc
Share the PTY mode table with the process broker
SaladDay Oct 1, 2026
44c9c50
Add the process shim and its broker
SaladDay Oct 1, 2026
0e1bc8a
Make the process broker's waits cancellable and its settlement complete
SaladDay Oct 1, 2026
1b71fad
Create a view's devpts instance before the view
SaladDay Oct 1, 2026
093fcf6
Harden the process broker against an untrusted view
SaladDay Oct 1, 2026
0c1a082
Take the process broker's view paths from the agent view layout
SaladDay Oct 1, 2026
c45b20c
Retry refused requests, never replay an uncertain Cancel, and end a l…
SaladDay Oct 1, 2026
f39d77d
Define the process relay's IPC with the shim and the broker
SaladDay Oct 1, 2026
233e71d
Move the view's descriptors into an unprivileged process relay
SaladDay Oct 1, 2026
6596239
Create a view's devpts instance in the launcher again
SaladDay Oct 1, 2026
e218b18
Reserve the process relay's name in the agent view's shim directory
SaladDay Oct 1, 2026
292ff90
Start the process relay with the view and bound the view's teardown
SaladDay Oct 1, 2026
67d26ef
Publish each relay invocation's Open in ID order
SaladDay Oct 1, 2026
faba667
Write terminal output only once the terminal is raw
SaladDay Oct 1, 2026
8b11d59
Forward stdin only once the program's Started event arrives
SaladDay Oct 1, 2026
8d56e7b
Resume stdin after an uncertain write from the accepted offset
SaladDay Oct 1, 2026
cd3fd0d
Close pipe stdin when an uncertain write finds the leader exited
SaladDay Oct 1, 2026
0bbd0f6
Cancel after lost stdin only when its failure reply was sent
SaladDay Oct 1, 2026
644e419
Check the exit and reserve the stdin failure reply in one step
SaladDay Oct 1, 2026
35df6f1
Hold the ID-order test between an invocation's ID and its Open
SaladDay Oct 1, 2026
92e2046
Set up the raw-mode and Starting tests with barriers, not timers
SaladDay Oct 1, 2026
bfe0799
Count the relay's output pumps before the control goroutine runs
SaladDay Oct 1, 2026
249bb78
Note that terminal output after the Exit is processed again
SaladDay Oct 1, 2026
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
17 changes: 17 additions & 0 deletions apps/daemon/cmd/oac-process-shim/main.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,17 @@
// Command oac-process-shim runs a declared sandbox executable from a Session
// view through the process broker, and is the Session's process relay. See
// apps/daemon/internal/processshim.
package main

import (
"os"

"github.com/MiniMax-AI/OpenAgentCore/apps/daemon/internal/processshim"
)

func main() {
if processshim.Relaying() {
os.Exit(processshim.Relay())
}
os.Exit(processshim.Run(processshim.SocketPath))
}
7 changes: 5 additions & 2 deletions apps/daemon/internal/agent/harness.go
Original file line number Diff line number Diff line change
Expand Up @@ -113,8 +113,11 @@ const (
ViewShimName = "bin"
// ViewHomeName is the Session home under ViewPrivateRoot.
ViewHomeName = "home"
// ViewRunName is the process broker's socket directory under ViewPrivateRoot.
// ViewRunName is the process relay's socket directory under ViewPrivateRoot.
ViewRunName = "run"
// ViewRelayName is the process relay's name in the shim directory, which
// no shim takes.
ViewRelayName = "oac-process-shim"
// ViewProcRoot and ViewDevRoot are the view's own /proc and minimal /dev.
ViewProcRoot = "/proc"
ViewDevRoot = "/dev"
Expand Down Expand Up @@ -335,7 +338,7 @@ func (v View) Validate() error {
}
}
for i, n := range v.Shims {
if !isPathComponent(n) || slices.Contains(v.Shims[:i], n) {
if !isPathComponent(n) || n == ViewRelayName || slices.Contains(v.Shims[:i], n) {
return invalidView("shim %q", n)
}
}
Expand Down
1 change: 1 addition & 0 deletions apps/daemon/internal/agent/view_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -32,6 +32,7 @@ func TestViewValidate(t *testing.T) {
"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") },
"shim named as the relay": func(v *agent.View) { v.Shims = append(v.Shims, agent.ViewRelayName) },
"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") },
Expand Down
5 changes: 1 addition & 4 deletions apps/daemon/internal/agenthost/broker.go
Original file line number Diff line number Diff line change
Expand Up @@ -11,7 +11,7 @@ import (
// over the Process service. One broker serves a Session from its first launch
// until teardown.
type processBroker interface {
// Start begins serving the Session's run directory.
// Start begins serving the Session.
Start(brokerConfig) error
// Close cancels and releases the remote operations that remain and stops
// serving.
Expand All @@ -20,9 +20,6 @@ type processBroker interface {

// brokerConfig is what a Session's broker serves.
type brokerConfig struct {
// RunDir is the host directory the view presents read-only at
// agent.ViewPrivateRoot/agent.ViewRunName.
RunDir string
// UID and GID are the Session's.
UID, GID uint32
// Names maps each shim name to the program it runs in the sandbox, found
Expand Down
8 changes: 4 additions & 4 deletions apps/daemon/internal/agenthost/doc.go
Original file line number Diff line number Diff line change
Expand Up @@ -17,10 +17,10 @@
// the gateway listening in the view's network namespace.
//
// Each view presents the closure directories read-only and executable, the
// Session home read-write and noexec, the broker's run directory read-only,
// the agent host's /etc/passwd, group, hosts, resolv.conf and nsswitch.conf,
// the agent host's CA directory at its host path, then the adapter's overlays
// and masks and the process shim. Everything else is the world.
// Session home read-write and noexec, the agent host's /etc/passwd, group,
// hosts, resolv.conf and nsswitch.conf, the agent host's CA directory at its
// host path, then the adapter's overlays and masks and the process shim with
// its relay. Everything else is the world.
//
// The agent host owns the Session's Link attachment: it opens each stream
// with the Session's binding, renews the lease and fails the Session when the
Expand Down
8 changes: 3 additions & 5 deletions apps/daemon/internal/agenthost/launch_linux.go
Original file line number Diff line number Diff line change
Expand Up @@ -243,7 +243,7 @@ func (s *session) startBroker(grace time.Duration) error {
paths[p] = p
}
b := s.deps.broker()
err := b.Start(brokerConfig{RunDir: s.dir.entry(runEntry), UID: s.uid, GID: s.uid, Names: names, Paths: paths,
err := b.Start(brokerConfig{UID: s.uid, GID: s.uid, Names: names, Paths: paths,
Pass: slices.Clone(view.ForwardEnv), Sandbox: s.in.Environment.Sandbox, Tool: s.in.Environment.Tool,
Dial: s.openProcess, CancelGrace: grace})
if err != nil {
Expand All @@ -255,7 +255,7 @@ func (s *session) startBroker(grace time.Duration) error {
return nil
}

// spec builds the view: the closure, home and run directories, the agent
// spec builds the view: the closure and home directories, the agent
// host's /etc files and CA directory, the adapter's overlays and masks, the
// shim and the gateway in the view's network namespace.
func (s *session) spec(viewCtx context.Context, world *worldfs.World, opts clirunner.StartOptions, ends *stdio) sessionview.Spec {
Expand All @@ -264,9 +264,7 @@ func (s *session) spec(viewCtx context.Context, world *worldfs.World, opts cliru
for _, m := range view.Closure {
private = append(private, sessionview.PrivateDir{Name: m.Name, HostDir: m.HostDir, Exec: true})
}
private = append(private,
sessionview.PrivateDir{Name: agent.ViewHomeName, HostDir: s.dir.entry(homeEntry), Writable: true},
sessionview.PrivateDir{Name: agent.ViewRunName, HostDir: s.dir.entry(runEntry)})
private = append(private, sessionview.PrivateDir{Name: agent.ViewHomeName, HostDir: s.dir.entry(homeEntry), Writable: true})
var overlays []sessionview.Overlay
for _, name := range etcFiles {
overlays = append(overlays, sessionview.Overlay{Path: "/etc/" + name, Source: s.dir.entry(etcEntry, name)})
Expand Down
7 changes: 3 additions & 4 deletions apps/daemon/internal/agenthost/sessiondir_linux.go
Original file line number Diff line number Diff line change
Expand Up @@ -47,7 +47,6 @@ func freeUID(id uint32) {
// The entries of a Session directory.
const (
homeEntry = agent.ViewHomeName // the Session home, owned by the Session uid
runEntry = agent.ViewRunName // the process broker's run directory
etcEntry = "etc" // the /etc files
maskEntry = "mask" // an empty file and an empty directory that masks present
stagingEntry = "staging" // sessionview's staging parent
Expand All @@ -62,8 +61,8 @@ func (d sessionDir) entry(name ...string) string {
return filepath.Join(append([]string{string(d)}, name...)...)
}

// createSessionDir creates the Session directory with its home, run, etc,
// mask and staging entries.
// createSessionDir creates the Session directory with its home, etc, mask
// and staging entries.
func createSessionDir(stateDir string, id sandboxwire.ID, uid uint32) (sessionDir, error) {
parent := sessionsDir(stateDir)
if err := os.MkdirAll(parent, 0o700); err != nil {
Expand All @@ -86,7 +85,7 @@ func (d sessionDir) populate(uid uint32) error {
dirs := []struct {
name string
mode fs.FileMode
}{{homeEntry, 0o700}, {runEntry, 0o755}, {etcEntry, 0o755}, {maskEntry, 0o755}, {filepath.Join(maskEntry, "dir"), 0o555}, {stagingEntry, 0o700}}
}{{homeEntry, 0o700}, {etcEntry, 0o755}, {maskEntry, 0o755}, {filepath.Join(maskEntry, "dir"), 0o555}, {stagingEntry, 0o700}}
for _, e := range dirs {
if err := os.Mkdir(d.entry(e.name), e.mode); err != nil {
return err
Expand Down
169 changes: 169 additions & 0 deletions apps/daemon/internal/processbroker/broker_linux.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,169 @@
//go:build linux

package processbroker

import (
"bufio"
"context"
"errors"
"fmt"
"log/slog"
"net"
"sync"

"github.com/MiniMax-AI/OpenAgentCore/apps/daemon/internal/processshim"
"github.com/MiniMax-AI/OpenAgentCore/internal/sandboxwire"
)

// Broker serves one Session's shim invocations, which its relay hands it.
type Broker struct {
cfg Config
log *slog.Logger
conn *net.UnixConn
link *link

ctx context.Context
cancel context.CancelFunc
wg sync.WaitGroup // the reader and the invocations
sendMu sync.Mutex
done chan struct{}

mu sync.Mutex
invs map[uint64]*invocation
lastID uint64
closing bool
err error

closeOnce sync.Once
}

// Start serves the relay at cfg.Relay until Close or the relay's loss.
func Start(cfg Config) (*Broker, error) {
if err := cfg.validate(); err != nil {
return nil, err
}
log := cfg.Logger
if log == nil {
log = slog.Default()
}
c, err := net.FileConn(cfg.Relay)
if err != nil {
return nil, fmt.Errorf("processbroker: relay connection: %w", err)
}
uc, ok := c.(*net.UnixConn)
if !ok {
c.Close()
return nil, fmt.Errorf("%w: the relay connection is not a Unix socket", ErrInvalidConfig)
}
ctx, cancel := context.WithCancel(context.Background())
b := &Broker{
cfg: cfg, log: log, conn: uc, link: newLink(cfg.Dial, log),
ctx: ctx, cancel: cancel, done: make(chan struct{}), invs: map[uint64]*invocation{},
}
b.wg.Add(1)
go b.read()
return b, nil
}

// Close ends every invocation and the relay connection, and waits for the
// invocations. It never waits on the relay: the relay answers a waiting shim
// with 255. Remote operations are left to the Session.
func (b *Broker) Close() error {
b.closeOnce.Do(func() {
b.mu.Lock()
b.closing = true
b.mu.Unlock()
b.cancel()
b.conn.CloseWrite()
b.conn.Close()
b.wg.Wait()
b.link.close()
})
return nil
}

// Done closes when the broker stops serving: after Close, or when the relay
// is lost.
func (b *Broker) Done() <-chan struct{} { return b.done }

// Err returns an error wrapping ErrRelayLost once the relay was lost before
// Close, and nil otherwise.
func (b *Broker) Err() error {
b.mu.Lock()
defer b.mu.Unlock()
return b.err
}

// read dispatches the relay's messages until the connection ends or breaks
// the IPC. It reads without a control buffer, so the kernel closes any
// descriptor the relay attaches, and it never blocks on an invocation.
func (b *Broker) read() {
defer b.wg.Done()
r := bufio.NewReaderSize(b.conn, 64<<10)
var err error
for err == nil {
var f sandboxwire.Frame
if f, err = sandboxwire.ReadFrame(r, processshim.MaxFrameBytes); err != nil {
break
}
var m processshim.RelayMessage
if m, err = processshim.DecodeRelay(f); err == nil {
err = b.dispatch(m)
}
}
b.mu.Lock()
lost := !b.closing
if lost {
b.err = fmt.Errorf("%w: %w", ErrRelayLost, err)
}
b.mu.Unlock()
if lost {
b.log.Error("process relay lost", "error", err)
}
b.cancel()
b.conn.Close()
close(b.done)
}

func (b *Broker) dispatch(m processshim.RelayMessage) error {
b.mu.Lock()
defer b.mu.Unlock()
id := m.Invocation()
if open, ok := m.(processshim.Open); ok {
switch {
case id <= b.lastID:
return fmt.Errorf("%w: open of invocation %d after %d", processshim.ErrProtocol, id, b.lastID)
case len(b.invs) >= processshim.MaxInvocations:
return fmt.Errorf("%w: more than %d invocations", processshim.ErrProtocol, processshim.MaxInvocations)
}
b.lastID = id
inv := b.newInvocation(open)
b.invs[id] = inv
b.wg.Add(1)
go inv.serve()
return nil
}
inv := b.invs[id]
switch {
case inv != nil:
return inv.receive(m)
case id > b.lastID:
return fmt.Errorf("%w: message for unopened invocation %d", processshim.ErrProtocol, id)
}
return nil // an invocation that has ended
}

// unregister stops dispatching to the invocation.
func (b *Broker) unregister(id uint64) {
b.mu.Lock()
defer b.mu.Unlock()
delete(b.invs, id)
}

var errEnded = errors.New("invocation ended")

func (b *Broker) send(m processshim.BrokerMessage) error {
b.sendMu.Lock()
defer b.sendMu.Unlock()
return sandboxwire.WriteFrame(b.conn, processshim.Frame(m))
}
Loading
Loading