diff --git a/apps/daemon/internal/gateway/view_linux_test.go b/apps/daemon/internal/gateway/view_linux_test.go index 1fc830e8..efa544cd 100644 --- a/apps/daemon/internal/gateway/view_linux_test.go +++ b/apps/daemon/internal/gateway/view_linux_test.go @@ -93,8 +93,9 @@ func TestListenersExistOnlyInTheSession(t *testing.T) { ctx, cancel := context.WithCancel(context.Background()) defer cancel() v, err := sessionview.Start(context.Background(), sessionview.Spec{ - World: (&loopbackWorld{dir: world}).serve, - Private: []sessionview.PrivateDir{{Name: "harness", HostDir: harness, Exec: true}}, + World: (&loopbackWorld{dir: world}).serve, + StagingParent: t.TempDir(), + Private: []sessionview.PrivateDir{{Name: "harness", HostDir: harness, Exec: true}}, Process: sessionview.Process{ Path: "/.oac/harness/harness", Args: []string{"harness"}, Dir: "/", UID: viewID, GID: viewID, Stderr: os.Stderr, Env: []string{harnessEnv + "=1", endpointsEnv + "=" + string(encoded), externalEnv + "=" + net.JoinHostPort(hostAddress(t), port)}, @@ -224,33 +225,37 @@ func copyExecutable(t *testing.T, dst string) { } } -// loopbackWorld serves a directory as the view's world. +// loopbackWorld serves a directory as the view's world. It presents each mountpoint at its declared path. type loopbackWorld struct { dir string served chan struct{} } -func (w *loopbackWorld) serve(dev *os.File, _ sessionview.WorldMount) (sessionview.WorldServer, error) { +func (w *loopbackWorld) serve(_ context.Context, dev *os.File, mount sessionview.WorldMount) (sessionview.WorldServer, sessionview.Presentation, error) { fd, err := unix.Dup(int(dev.Fd())) if err != nil { - return nil, err + return nil, sessionview.Presentation{}, err } root, err := gofs.NewLoopbackRoot(w.dir) if err != nil { unix.Close(fd) - return nil, err + return nil, sessionview.Presentation{}, err } srv, err := fuse.NewServer(gofs.NewNodeFS(root, &gofs.Options{}), fmt.Sprintf("/dev/fd/%d", fd), &fuse.MountOptions{}) if err != nil { unix.Close(fd) - return nil, err + return nil, sessionview.Presentation{}, err } w.served = make(chan struct{}) go func() { srv.Serve() close(w.served) }() - return w, nil + var p sessionview.Presentation + for _, m := range mount.Mountpoints { + p.Targets = append(p.Targets, m.Path) + } + return w, p, nil } func (w *loopbackWorld) Stop() error { diff --git a/apps/daemon/internal/sessionview/build_linux.go b/apps/daemon/internal/sessionview/build_linux.go index dab12b6b..ae5c4131 100644 --- a/apps/daemon/internal/sessionview/build_linux.go +++ b/apps/daemon/internal/sessionview/build_linux.go @@ -12,10 +12,11 @@ import ( "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. +// builder mounts the local pieces onto the world, at the paths the world presents the mountpoints at. Every target is resolved beneath its parent mount without following symlinks, so the sandbox cannot redirect a mount. type builder struct { - root int // the world root, an O_PATH fd - fds []int + root int // the world root, an O_PATH fd + targets map[string]string // each mountpoint's view path to where the world presents it + fds []int } // devNodes are bound from the host's /dev into the view's /dev. @@ -42,11 +43,15 @@ func (b *builder) build(spec *launchSpec) error { } b.root = b.keep(root) for _, d := range spec.Private { + at, err := b.at(agent.ViewPrivateRoot + "/" + d.Name) + if err != nil { + return err + } src, err := b.source(d.HostDir) if err != nil { return err } - if err := b.bind(src, b.root, agent.ViewPrivateRoot+"/"+d.Name, bindAttr(d.Writable, d.Exec, false)); err != nil { + if err := b.bind(src, b.root, at, bindAttr(d.Writable, d.Exec, false)); err != nil { return err } } @@ -60,33 +65,56 @@ func (b *builder) build(spec *launchSpec) error { return err } for _, o := range spec.Overlays { + at, err := b.at(o.Path) + if err != nil { + return err + } src, err := b.source(o.Source) if err != nil { return err } - if err := b.bind(src, b.root, o.Path, bindAttr(false, o.Exec, false)); err != nil { + if err := b.bind(src, b.root, at, bindAttr(false, o.Exec, false)); err != nil { return err } } for _, p := range spec.Shim.Paths { - if err := b.bind(shim, b.root, p, bindAttr(false, true, false)); err != nil { + at, err := b.at(p) + if err != nil { + return err + } + if err := b.bind(shim, b.root, at, bindAttr(false, true, false)); err != nil { return err } } + at, err := b.at(agent.ViewProcRoot) + if err != nil { + return err + } proc, err := newFS("proc", nil, attrNoSuid|attrNoDev|attrNoExec) if err != nil { return err } defer unix.Close(proc) - if err := b.attach(proc, b.root, agent.ViewProcRoot, true); err != nil { + if err := b.attach(proc, b.root, at, true); err != nil { return err } return b.dev() } +// at returns where the world presents the mountpoint at view path p. +func (b *builder) at(p string) (string, error) { + if t, ok := b.targets[p]; ok { + return t, nil + } + return "", &Error{Kind: ErrLauncher, Op: "present", Path: p, Err: errors.New("the world presents no target")} +} + // shimDir presents the shim at /.oac/bin/ on a read-only tmpfs. func (b *builder) shimDir(names []string, shim int) error { - dir := agent.ViewPrivateRoot + "/" + agent.ViewShimName + dir, err := b.at(agent.ViewPrivateRoot + "/" + agent.ViewShimName) + if err != nil { + return err + } mnt, err := newFS("tmpfs", [][2]string{{"mode", "0755"}, {"size", "64k"}}, attrNoSuid|attrNoDev|attrNoExec) if err != nil { return err @@ -108,12 +136,16 @@ func (b *builder) shimDir(names []string, shim int) error { // dev builds a minimal read-only /dev with host device nodes, a new devpts instance and a noexec /dev/shm. func (b *builder) dev() error { + at, err := b.at(agent.ViewDevRoot) + if err != nil { + return err + } mnt, err := newFS("tmpfs", [][2]string{{"mode", "0755"}, {"size", "64k"}}, attrNoSuid|attrNoDev|attrNoExec) if err != nil { return err } defer unix.Close(mnt) - if err := b.attach(mnt, b.root, agent.ViewDevRoot, true); err != nil { + if err := b.attach(mnt, b.root, at, true); err != nil { return err } for _, n := range devNodes { diff --git a/apps/daemon/internal/sessionview/control_linux.go b/apps/daemon/internal/sessionview/control_linux.go index 266dc19c..65fe9065 100644 --- a/apps/daemon/internal/sessionview/control_linux.go +++ b/apps/daemon/internal/sessionview/control_linux.go @@ -46,7 +46,7 @@ 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 + msgProceed // daemon: the world serves and the network is set up; carries the mount targets msgStarted // launcher: the process runs msgFailed // launcher: construction failed msgExited // launcher: the process ended @@ -61,6 +61,7 @@ type message struct { Delivered bool Exit Exit Fail failure + Targets map[string]string // each mountpoint's view path to the path the world presents it at } // 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 2d6271a2..1e296a51 100644 --- a/apps/daemon/internal/sessionview/doc.go +++ b/apps/daemon/internal/sessionview/doc.go @@ -1,6 +1,6 @@ // Package sessionview runs one process inside a per-Session view on the agent host. // -// A view is a private mount, PID and network namespace whose root is the Session's world: a FUSE file system that the daemon serves over a /dev/fuse connection. The launcher adds the local pieces on top of the world: private directories under /.oac, trusted overlays, the command shim, a fresh /proc and a minimal /dev. The process starts with no capabilities, no_new_privs, a seccomp filter and only stdin, stdout and stderr open. Its network namespace has only loopback up. +// A view is a private mount, PID and network namespace whose root is the Session's world: a FUSE file system that the daemon serves over a /dev/fuse connection. The launcher adds the local pieces on top of the world: private directories under /.oac, trusted overlays, the command shim, a fresh /proc and a minimal /dev. The world presents a mountpoint for each piece and reports where, following the sandbox's symlinks, and the launcher mounts at those paths without following any symlink itself. The process starts with no capabilities, no_new_privs, a seccomp filter and only stdin, stdout and stderr open. Its network namespace has only loopback up. // // 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. // diff --git a/apps/daemon/internal/sessionview/launcher_linux.go b/apps/daemon/internal/sessionview/launcher_linux.go index 5b8fffc0..6a8bd752 100644 --- a/apps/daemon/internal/sessionview/launcher_linux.go +++ b/apps/daemon/internal/sessionview/launcher_linux.go @@ -31,6 +31,7 @@ type launcher struct { ctl *control proceed chan struct{} proceedOnce sync.Once + targets map[string]string // set before proceed closes // 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 @@ -69,7 +70,7 @@ func (l *launcher) run() (int, error) { return 0, err } <-l.proceed - b := &builder{root: -1} + b := &builder{root: -1, targets: l.targets} defer b.close() if err := b.build(spec); err != nil { return 0, err @@ -122,7 +123,10 @@ func (l *launcher) serveControl() { } switch m.Kind { case msgProceed: - l.proceedOnce.Do(func() { close(l.proceed) }) + l.proceedOnce.Do(func() { + l.targets = m.Targets + close(l.proceed) + }) case msgSignal: _ = l.ctl.send(message{Kind: msgSignaled, Delivered: l.signal(m.Signal)}) } diff --git a/apps/daemon/internal/sessionview/spec.go b/apps/daemon/internal/sessionview/spec.go index 7c28774d..17204a52 100644 --- a/apps/daemon/internal/sessionview/spec.go +++ b/apps/daemon/internal/sessionview/spec.go @@ -1,6 +1,7 @@ package sessionview import ( + "context" "os" "path" "path/filepath" @@ -23,8 +24,10 @@ type Spec struct { 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. -type World func(dev *os.File, mount WorldMount) (WorldServer, error) +// World starts serving the view's root file system on dev, a /dev/fuse connection that the launcher has already mounted as mount describes, and reports how it presents mount's mountpoints. ctx is Start's: it bounds connecting, attaching and presenting, not the serving. The server runs in the daemon, outside the view, and never mounts or unmounts anything. sessionview closes dev after Stop returns. +// +// A World that fails releases what it acquired, or its error says that it cannot show it did. The owner of the attachment the world serves from then ends that attachment, which releases everything it holds. +type World func(ctx context.Context, dev *os.File, mount WorldMount) (WorldServer, Presentation, error) // WorldServer is a running world. type WorldServer interface { @@ -38,10 +41,28 @@ type WorldMount struct { Options []string // Flags are the mount flags. Flags []string - // Mountpoints are the view paths the launcher mounts over. The world presents each one, and each ancestor as a directory, with no symlinks and a stable identity for the view's lifetime, whether or not the sandbox has the path. A mountpoint that disappears or changes identity detaches what is mounted on it. + // UID and GID are the identity the view's process runs as. + UID, GID uint32 + // Mountpoints are the view paths the launcher mounts over. The world presents each one with its type, whether or not the sandbox has the path, and keeps it and its ancestors stable for the view's lifetime: a mountpoint that disappears or changes identity detaches what is mounted on it. Mountpoints []Mountpoint } +// Presentation is how the world presents the mountpoints. +type Presentation struct { + // Targets holds, for each of WorldMount.Mountpoints in order, the absolute view path without symlinks where the world presents it. The launcher mounts there, still refusing to follow a symlink. + Targets []string + // Links are the sandbox symlinks on the way to a mountpoint. They stay symlinks in the view. + Links []PresentedLink + // Synthesized are the directories the world presents because the sandbox lacks them. + Synthesized []string +} + +// PresentedLink is the symlink at Path, whose target is Target. +type PresentedLink struct { + Path string + Target string +} + // Mountpoint is a view path the world presents as a directory or a regular file. type Mountpoint struct { Path string diff --git a/apps/daemon/internal/sessionview/view_linux.go b/apps/daemon/internal/sessionview/view_linux.go index 21b30294..7303f2a4 100644 --- a/apps/daemon/internal/sessionview/view_linux.go +++ b/apps/daemon/internal/sessionview/view_linux.go @@ -35,6 +35,7 @@ type View struct { cmd *exec.Cmd ctl *control world WorldServer + present Presentation dev *os.File staging string pipes [3]*os.File @@ -65,9 +66,10 @@ func Start(ctx context.Context, spec Spec) (*View, error) { return nil, err } stop := context.AfterFunc(ctx, func() { _ = v.cmd.Process.Kill() }) - err := v.handshake(&spec) + err := v.handshake(ctx, &spec) if !stop() { - err = &Error{Kind: ErrLauncher, Op: "start", Err: ctx.Err()} + // The world's own error stays: it may say that the attachment must be ended. + err = errors.Join(&Error{Kind: ErrLauncher, Op: "start", Err: ctx.Err()}, err) } if err != nil { v.abort() @@ -166,7 +168,7 @@ func (v *View) stdio(p *Process) ([3]*os.File, error) { return child, nil } -func (v *View) handshake(spec *Spec) error { +func (v *View) handshake(ctx context.Context, spec *Spec) error { m, files, err := v.ctl.recv() if err != nil { return v.lost("mount", err) @@ -181,17 +183,22 @@ func (v *View) handshake(spec *Spec) error { v.dev = files[0] netns := files[1] defer netns.Close() - world, err := spec.World(v.dev, WorldMount{Options: fuseOptions, Flags: fuseFlags, Mountpoints: spec.mountpoints()}) + mps := spec.mountpoints() + world, present, err := spec.World(ctx, v.dev, WorldMount{Options: fuseOptions, Flags: fuseFlags, UID: spec.Process.UID, GID: spec.Process.GID, Mountpoints: mps}) if err != nil { return &Error{Kind: ErrWorld, Op: "serve", Err: err} } - v.world = world + v.world, v.present = world, present + targets, err := targetsOf(mps, present) + if err != nil { + return &Error{Kind: ErrWorld, Op: "present", Err: err} + } if spec.Network.Setup != nil { if err := spec.Network.Setup(netns); err != nil { return &Error{Kind: ErrNetwork, Op: "setup", Err: err} } } - if err := v.ctl.send(message{Kind: msgProceed}); err != nil { + if err := v.ctl.send(message{Kind: msgProceed, Targets: targets}); err != nil { return v.lost("proceed", err) } m, files, err = v.ctl.recv() @@ -207,6 +214,22 @@ func (v *View) handshake(spec *Spec) error { return nil } +// targetsOf pairs each mountpoint with the path the world presents it at. +func targetsOf(mps []Mountpoint, p Presentation) (map[string]string, error) { + if len(p.Targets) != len(mps) { + return nil, fmt.Errorf("%d targets for %d mountpoints", len(p.Targets), len(mps)) + } + targets := make(map[string]string, len(mps)) + for i, m := range mps { + t := p.Targets[i] + if !isViewAbs(t) || t == "/" { + return nil, fmt.Errorf("target %q for %s", t, m.Path) + } + targets[m.Path] = t + } + return targets, nil +} + // lost reports a launcher that stopped talking, with its exit status when it has exited. func (v *View) lost(op string, err error) error { if errors.Is(err, io.EOF) { @@ -302,6 +325,9 @@ func (v *View) Wait() (Exit, error) { return v.exit, v.err } +// Presentation reports how the world presented the view's mountpoints. +func (v *View) Presentation() Presentation { return v.present } + // 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() diff --git a/apps/daemon/internal/sessionview/view_linux_test.go b/apps/daemon/internal/sessionview/view_linux_test.go index 328ee8f0..f68a2b97 100644 --- a/apps/daemon/internal/sessionview/view_linux_test.go +++ b/apps/daemon/internal/sessionview/view_linux_test.go @@ -168,7 +168,7 @@ func TestViewDescendantsKeepTheGrace(t *testing.T) { } } -// TestViewRefusesSymlinkedMountpoint checks that a sandbox symlink on the way to a mountpoint fails the view instead of redirecting the mount. +// TestViewRefusesSymlinkedMountpoint checks that the launcher still refuses a symlink on the way to a target, so a world that reports a target it did not resolve cannot redirect a mount. func TestViewRefusesSymlinkedMountpoint(t *testing.T) { requireView(t) f := newFixture(t) @@ -325,33 +325,37 @@ func (f *fixture) spec(w *loopbackWorld, mode string, env ...string) Spec { } } -// loopbackWorld serves a directory as the world, the way the world frontend serves a sandbox. +// loopbackWorld serves a directory as the world, the way the world frontend serves a sandbox. It presents each mountpoint at its declared path. type loopbackWorld struct { dir string served chan struct{} } -func (w *loopbackWorld) serve(dev *os.File, _ WorldMount) (WorldServer, error) { +func (w *loopbackWorld) serve(_ context.Context, dev *os.File, mount WorldMount) (WorldServer, Presentation, error) { fd, err := unix.Dup(int(dev.Fd())) if err != nil { - return nil, err + return nil, Presentation{}, err } root, err := gofs.NewLoopbackRoot(w.dir) if err != nil { unix.Close(fd) - return nil, err + return nil, Presentation{}, err } srv, err := fuse.NewServer(gofs.NewNodeFS(root, &gofs.Options{}), fmt.Sprintf("/dev/fd/%d", fd), &fuse.MountOptions{}) if err != nil { unix.Close(fd) - return nil, err + return nil, Presentation{}, err } w.served = make(chan struct{}) go func() { srv.Serve() close(w.served) }() - return w, nil + var p Presentation + for _, m := range mount.Mountpoints { + p.Targets = append(p.Targets, m.Path) + } + return w, p, nil } func (w *loopbackWorld) Stop() error { diff --git a/apps/daemon/internal/sessionview/view_other.go b/apps/daemon/internal/sessionview/view_other.go index f50c0bef..8770549d 100644 --- a/apps/daemon/internal/sessionview/view_other.go +++ b/apps/daemon/internal/sessionview/view_other.go @@ -20,6 +20,7 @@ func Start(context.Context, Spec) (*View, error) { return nil, ErrUnsupported } // View is a running view. Outside Linux none exists. type View struct{} +func (*View) Presentation() Presentation { return Presentation{} } func (*View) Wait() (Exit, error) { return Exit{}, ErrUnsupported } func (*View) Signal(syscall.Signal) error { return ErrUnsupported } func (*View) Close() error { return nil } diff --git a/apps/daemon/internal/worldfs/conn_linux.go b/apps/daemon/internal/worldfs/conn_linux.go new file mode 100644 index 00000000..4f9e625f --- /dev/null +++ b/apps/daemon/internal/worldfs/conn_linux.go @@ -0,0 +1,139 @@ +//go:build linux + +package worldfs + +import ( + "context" + "errors" + + "github.com/MiniMax-AI/OpenAgentCore/internal/sandboxfs" + "github.com/MiniMax-AI/OpenAgentCore/internal/sandboxwire" +) + +var ( + errDead = errors.New("worldfs: the world is lost") + errInterrupted = errors.New("worldfs: interrupted before the request was sent") +) + +// connect opens a stream and describes the service within ctx. +func (f *frontend) connect(ctx context.Context) (*sandboxfs.Client, *sandboxfs.DescribeResponse, error) { + rw, err := f.dial(ctx) + if err != nil { + return nil, nil, err + } + c := sandboxfs.NewClient(rw) + d, err := c.Describe(ctx, &sandboxfs.DescribeRequest{}) + if err != nil { + c.Close() + return nil, nil, err + } + return c, d, nil +} + +// client returns the stream's client. After the stream failed it redials within ctx and continues only with the same service instance. A failed redial is an [ErrConnect] error, and an interrupt while waiting for the stream or redialing is errInterrupted: either way the request was never sent. +func (f *frontend) client(ctx context.Context, interrupt <-chan struct{}) (*sandboxfs.Client, error) { + select { + case f.connTurn <- struct{}{}: + case <-ctx.Done(): + return nil, &Error{Kind: ErrConnect, Op: "reconnect", Err: ctx.Err()} + case <-interrupt: + return nil, errInterrupted + } + defer func() { <-f.connTurn }() + if f.dead.Load() { + return nil, errDead + } + if c := f.conn; c != nil { + select { + case <-c.Done(): + default: + return c, nil + } + } + dial, cancel := context.WithCancel(ctx) + defer cancel() + if interrupt != nil { + go func() { + select { + case <-interrupt: + cancel() + case <-dial.Done(): + } + }() + } + c, d, err := f.connect(dial) + switch { + case err != nil && interrupted(interrupt): + return nil, errInterrupted + case err != nil: + f.observe(err) + return nil, &Error{Kind: ErrConnect, Op: "reconnect", Err: err} + case d.ServerInstanceID != f.instance: + c.Close() + f.lose(&Error{Kind: ErrInstanceChanged, Op: "reconnect"}, true) + return nil, errDead + } + f.conn = c + f.gen.Add(1) + f.wake() // what the failed stream never sent goes on the new one + return c, nil +} + +// call sends one request. A request the stream failed before sending cannot have taken effect, so it is sent once more on a new stream; nothing else is retried. +func call[Q, R any](f *frontend, ctx context.Context, op func(*sandboxfs.Client, context.Context, Q) (R, error), q Q) (R, error) { + return callUntil(f, ctx, nil, op, q) +} + +// callUntil is call for a request the kernel may interrupt. An interrupt before the request is sent abandons it with errInterrupted; after that the call waits for the request's own outcome (see [sandboxfs.WithInterrupt]). +func callUntil[Q, R any](f *frontend, ctx context.Context, interrupt <-chan struct{}, op func(*sandboxfs.Client, context.Context, Q) (R, error), q Q) (R, error) { + var r R + var err error + sent := ctx + if interrupt != nil { + sent = sandboxfs.WithInterrupt(ctx, interrupt) + } + for range 2 { + var c *sandboxfs.Client + if c, err = f.client(ctx, interrupt); err != nil { + return r, err + } + if interrupted(interrupt) { + return r, errInterrupted + } + if r, err = op(c, sent, q); err == nil { + return r, nil + } + f.observe(err) + if !unsent(err) { + break + } + } + return r, err +} + +func interrupted(interrupt <-chan struct{}) bool { + select { + case <-interrupt: + return true + default: + return false + } +} + +// unsent reports whether err shows that the request never left the frontend: the redial failed, or the stream had failed before the request was written. +func unsent(err error) bool { + var fail *sandboxfs.Failure + return errors.Is(err, ErrConnect) || errors.As(err, &fail) && fail.Effect == sandboxwire.EffectNone && errors.Is(err, sandboxfs.ErrTransport) +} + +// retryable reports whether a request certainly did nothing and may succeed when sent again: it was never sent, or the service refused it for now with ResourceExhausted and EffectNone. +func retryable(err error) bool { + var fail *sandboxfs.Failure + return unsent(err) || errors.As(err, &fail) && fail.Code == sandboxfs.CodeResourceExhausted && fail.Effect == sandboxwire.EffectNone +} + +// noEffect reports whether a request certainly changed nothing. +func noEffect(err error) bool { + var fail *sandboxfs.Failure + return unsent(err) || errors.Is(err, errDead) || errors.Is(err, errInterrupted) || errors.As(err, &fail) && fail.Effect == sandboxwire.EffectNone +} diff --git a/apps/daemon/internal/worldfs/dir_linux.go b/apps/daemon/internal/worldfs/dir_linux.go new file mode 100644 index 00000000..e8e119fe --- /dev/null +++ b/apps/daemon/internal/worldfs/dir_linux.go @@ -0,0 +1,94 @@ +//go:build linux + +package worldfs + +import ( + "github.com/MiniMax-AI/OpenAgentCore/internal/sandboxfs" + "github.com/hanwen/go-fuse/v2/fuse" +) + +func (f *frontend) OpenDir(_ <-chan struct{}, in *fuse.OpenIn, out *fuse.OpenOut) fuse.Status { + n, st := f.node(in.NodeId) + if !st.Ok() { + return st + } + h := &handle{node: n} + if !n.synthetic() { + r, err := call(f, f.ctx, (*sandboxfs.Client).OpenDir, &sandboxfs.OpenDirRequest{Node: n.ref}) + if err != nil { + return status(err) + } + h.server = r.Handle + } + *out = fuse.OpenOut{Fh: f.newHandle(h)} + return fuse.OK +} + +// ReadDir lists a directory without "." and "..". A plain sandbox directory uses the service's cookies as offsets, and a synthetic directory lists its children at offsets 1 to k. A directory with presented children is a [listing]. +func (f *frontend) ReadDir(_ <-chan struct{}, in *fuse.ReadIn, out *fuse.DirEntryList) fuse.Status { + h, st := f.handle(in.Fh) + if !st.Ok() { + return st + } + n := h.node + switch { + case h.server == 0: + for off := in.Offset; off < uint64(len(n.order)); off++ { + c := n.fixed[n.order[off]] + if !out.AddDirEntry(fuse.DirEntry{Mode: c.fileType(), Name: n.order[off], Ino: f.ino(c), Off: off + 1}) { + break + } + } + return fuse.OK + case len(n.order) == 0: + return f.readPlain(h, in, out) + } + return f.readPresented(h, in, out) +} + +func (f *frontend) readPlain(h *handle, in *fuse.ReadIn, out *fuse.DirEntryList) fuse.Status { + r, err := call(f, f.ctx, (*sandboxfs.Client).ReadDir, &sandboxfs.ReadDirRequest{Handle: h.server, Cookie: in.Offset, Limit: min(in.Size, f.caps.MaxReadDirBytes)}) + switch { + case err != nil: + return status(err) + case len(r.Entries) == 0 && !r.End: + return fuse.EIO + } + for _, e := range r.Entries { + if !out.AddDirEntry(fuse.DirEntry{Mode: e.Type, Name: string(e.Name), Ino: e.Ino, Off: e.Cookie}) { + break + } + } + return fuse.OK +} + +func (f *frontend) ino(n *inode) uint64 { + if n.synthetic() { + return f.attrOf(n).Ino + } + return n.attr.Ino +} + +func (f *frontend) ReleaseDir(in *fuse.ReleaseIn) { + if h := f.dropHandle(in.Fh); h != nil && h.server != 0 && !f.closed.Load() { + f.release(cleanup{handle: h.server, dir: true}) + } +} + +// FsyncDir fails with EINVAL, as fsync(2) does on a directory a file system cannot sync, when the service does not declare DirectoryFsync. +func (f *frontend) FsyncDir(_ <-chan struct{}, in *fuse.FsyncIn) fuse.Status { + h, st := f.handle(in.Fh) + switch { + case !st.Ok() || h.server == 0: + return st + case !f.caps.DirectoryFsync: + return fuse.EINVAL + } + _, err := call(f, f.ctx, (*sandboxfs.Client).Fsync, &sandboxfs.FsyncRequest{Handle: h.server, DataOnly: in.FsyncFlags&1 != 0}) + return status(err) +} + +// ReadDirPlus is never negotiated. +func (f *frontend) ReadDirPlus(<-chan struct{}, *fuse.ReadIn, *fuse.DirEntryList) fuse.Status { + return fuse.ENOSYS +} diff --git a/apps/daemon/internal/worldfs/doc.go b/apps/daemon/internal/worldfs/doc.go new file mode 100644 index 00000000..e97316f4 --- /dev/null +++ b/apps/daemon/internal/worldfs/doc.go @@ -0,0 +1,75 @@ +// Package worldfs serves a Session's world, the sandbox file system, as the FUSE file system a sessionview launcher mounts at the view's root. Every kernel request becomes at most one File request (see internal/sandboxfs), apart from the redial and lock recovery described below, so processes in the view read and write sandbox files natively; nothing on the agent host shows through, and nothing is created in the sandbox to support the view. +// +// [World.Serve] implements [sessionview.World]. Within the context sessionview.Start passes, it dials the attachment's File stream, checks Describe, attaches the export and presents the view's mountpoints; it then serves the launcher's /dev/fuse connection with go-fuse's raw API. +// +// # Mapping +// +// A kernel node ID names a NodeRef, and each kernel lookup holds one server lookup reference; FORGET and BATCH_FORGET release the same counts with Forget. A kernel file handle names a HandleID, and a kernel lock owner is the attachment's LockOwner. +// +// LOOKUP Lookup +// FORGET, BATCH_FORGET Forget, batched +// GETATTR GetAttr of the handle when the kernel names one, else of the node +// SETATTR SetAttr; a ctime-only change is GetAttr +// ACCESS Access +// READLINK Readlink +// MKDIR Mkdir +// UNLINK, RMDIR Unlink, Rmdir +// RENAME, RENAME2 Rename: Replace, NoReplace or Exchange; other flags EINVAL +// LINK, SYMLINK Link, Symlink +// CREATE Create; O_EXCL is Exclusive +// OPEN Open +// READ, WRITE Read, Write; a write that stopped after a prefix is a short write +// FLUSH Flush with the closing lock owner +// FSYNC, FSYNCDIR Fsync of the file or directory handle +// RELEASE, RELEASEDIR Release, ReleaseDir +// OPENDIR, READDIR OpenDir, ReadDir +// STATFS StatFS +// GETLK GetLock +// SETLK, SETLKW SetLock, POSIX or flock; SETLKW waits (see Locks) +// MKNOD EOPNOTSUPP +// GETXATTR, LISTXATTR ENOSYS: the kernel stops asking and answers EOPNOTSUPP +// SETXATTR, REMOVEXATTR ENOSYS, as above +// FALLOCATE ENOSYS, as above +// LSEEK ENOSYS: the kernel seeks itself +// COPY_FILE_RANGE ENOSYS: the kernel copies with READ and WRITE +// STATX ENOSYS: the kernel uses GETATTR +// IOCTL ENOTTY +// READDIRPLUS never negotiated +// +// FIFOs and device nodes in the world are opened by the kernel itself and never reach the sandbox; the world is mounted nodev. +// +// # Locks +// +// Lock requests the service does not declare fail with ENOLCK, so the kernel never locks only within the view. A lock request reports what the service did. An interrupt that arrives before the request is sent, while it waits for its turn, for the stream or for a redial, is EINTR and sends nothing. Once the request is sent, the frontend answers an interrupt with CancelRequest and waits for the request's own response: a lock acquired before the cancellation arrived is reported acquired, and a cancelled request is EINTR. When the service answers that a failed request may have changed the lock (EffectPossible), the frontend unlocks the same owner and range before it reports the failure. When it cannot learn what the request did, because the stream failed, or the unlock fails, it fails the handle and the request returns EIO: every later request on the handle fails with EIO too, and its Release drops any lock the service holds on it. +// +// The requests of one owner on one handle take turns, so that a recovery unlock never removes a lock another request reported. A non-blocking request keeps its turn until its recovery is done. A waiting request (SETLKW) gives up its turn while it waits, so unlocks and interrupts get through, and takes it again only to recover. A request unlocks to recover only when no other acquisition of the owner registered after it and no waiting request was out when it registered; otherwise the unlock could remove a lock another request reported, so the handle fails and the request returns EIO. +// +// # Uncached profile +// +// Entry, attribute and negative timeouts are zero, every open returns FOPEN_DIRECT_IO, and the frontend never asks for the writeback cache, KEEP_CACHE or CACHE_DIR. A private mapping of a world file works, and the kernel reads its pages with READ when they fault; a shared mapping fails with ENODEV, since the frontend does not negotiate DIRECT_IO_ALLOW_MMAP. default_permissions stays off: the service decides access. The frontend negotiates only BIG_WRITES, MAX_PAGES (requests up to the service's read and write limit), PARALLEL_DIROPS, ATOMIC_O_TRUNC, POSIX_LOCKS and FLOCK_LOCKS. +// +// # Errors +// +// A sandboxfs errno maps to its Linux errno through one table, StaleNode and StaleHandle map to ESTALE, and every other failure maps to EIO. Nothing a request did not do is reported as done. +// +// # Identity +// +// Attributes owned by the service's uid or gid, from Describe, show the view's uid or gid from [sessionview.WorldMount]; SetAttr maps them back. Other owners pass through unchanged, and the service decides every change. +// +// # Presentation +// +// The launcher mounts private directories, overlays and shims over mountpoints in the world. The frontend presents each one without touching the sandbox: +// +// - The mountpoint itself is a synthetic empty read-only directory or regular file that hides any sandbox entry of that name. +// - A missing ancestor, confirmed by NotFound, is a synthetic directory, 0555 and root-owned, holding only synthetic children. +// - An existing ancestor directory stays the sandbox directory, and Lookup of a presented name returns the presented node. ReadDir lists every sandbox entry at its own position, a presented name as the presented node, and then the presented names the sandbox does not list. +// - An ancestor symlink, such as /bin -> usr/bin, stays a symlink. The frontend resolves it with Walk and Readlink inside the world root, at most 40 hops, and presents the mountpoint at the resolved target: a shim declared at /bin/sh is mounted at /usr/bin/sh. +// +// Every entry on the way to a mountpoint is pinned for the view's lifetime: Lookup returns the same node ID, a pinned symlink is read from its pinned target, and creating, removing or renaming a presented name returns EPERM. When the sandbox replaces or removes a pinned entry, the view keeps the pinned entry and [World.Lost] reports [ErrTopologyChanged]. Lookup finds any replacement or removal by its node. ReadDir finds an entry listed with another type, and a removal when an enumeration from offset 0 reaches the end without the name; it compares no inode numbers, so a replacement of the same type is found on Lookup. A loop, a permission failure, a non-directory ancestor or two mountpoints that meet fail Serve with [ErrMountpoint]. Serve returns the resolved, symlink-free path of each mountpoint and the links and directories it presented. Synthetic nodes have frontend-local node IDs and never reach the service. +// +// # Link loss +// +// When the stream fails, requests in flight fail with EIO. The next request redials, sends Describe and continues with the same node and handle tables only when ServerInstanceID is unchanged; a request is sent again only when the stream never sent it. A Forget, Release or ReleaseDir that certainly did nothing, because the stream never sent it or the service refused it with ResourceExhausted and EffectNone, is queued and sent again after the next redial, or after a backoff of up to 2 seconds; one that may have reached the service is never sent again. A new incarnation or an ended attachment marks the world lost: every later request fails with EIO and [World.Lost] closes. The attachment has ended when the service answers StaleAttachment, or when the redial fails with a Link failure that is not retryable (sandboxlink.Code.Retryable), such as LeaseExpired or StaleGeneration; a redial refused with InstanceChanged is a new incarnation. +// +// Detach cannot overtake a request still running on a failed stream, so a Serve that fails detaches only when the export was attached on the stream still in use and nothing else is in doubt. When Attach may have taken effect without a response, when a request may still run on a failed stream, or when Detach fails, Serve fails with [ErrAttachmentDirty], and the owner of the Link attachment ends it; the service releases everything the attachment holds when its lease ends. [World.Stop] detaches without sending what is queued, since Detach drops every reference and handle the attachment holds. +package worldfs diff --git a/apps/daemon/internal/worldfs/drain_linux.go b/apps/daemon/internal/worldfs/drain_linux.go new file mode 100644 index 00000000..c6551e16 --- /dev/null +++ b/apps/daemon/internal/worldfs/drain_linux.go @@ -0,0 +1,121 @@ +//go:build linux + +package worldfs + +import ( + "context" + "time" + + "github.com/MiniMax-AI/OpenAgentCore/internal/sandboxfs" +) + +// cleanup is a Release or ReleaseDir of a handle the kernel closed. +type cleanup struct { + handle sandboxfs.HandleID + dir bool +} + +// wake runs the drainer without blocking. Whoever queues work calls it after queuing, so the drainer, which takes the queues after each wakeup, always sees the work. +func (f *frontend) wake() { + select { + case f.kick <- struct{}{}: + default: + } +} + +// drain sends what the kernel released: queued Forget references and the Release and ReleaseDir requests that could not be sent. After a pass that leaves something queued it waits, doubling the wait from retryMin to retryMax, unless a redial gives it a new stream first. It runs until Stop. +func (f *frontend) drain() { + defer close(f.drained) + var retry <-chan time.Time + var backoff time.Duration + var failedOn uint64 // the stream generation the last pass that left work queued failed on + for { + select { + case <-f.drainCtx.Done(): + return + case <-retry: + case <-f.kick: + if retry != nil && f.gen.Load() == failedOn { + continue // no new stream: newly queued work waits for the retry with the rest + } + } + if f.flush() { + retry, backoff = nil, 0 + continue + } + // f.gen now counts every redial up to the failure, the pass's own included, so only a stream opened after the failure ends the wait early. + failedOn, backoff = f.gen.Load(), min(max(2*backoff, retryMin), retryMax) + retry = time.After(backoff) + } +} + +// flush sends everything queued and reports whether nothing was left to retry. It stops at the first request it must retry, which goes back to the queue with everything after it. +func (f *frontend) flush() bool { + if !f.sendForgets() { + return false + } + f.mu.Lock() + rs := f.releases + f.releases = nil + f.mu.Unlock() + for i, c := range rs { + if !f.sendRelease(f.drainCtx, c) { + f.mu.Lock() + f.releases = append(rs[i:], f.releases...) + f.mu.Unlock() + return false + } + } + return true +} + +// sendForgets sends the queued references in batches. A batch that certainly did nothing goes back to the queue with the rest; one that may have been applied is not sent again. +func (f *frontend) sendForgets() bool { + f.mu.Lock() + batch := f.forgets + f.forgets = map[sandboxfs.NodeRef]uint64{} + f.mu.Unlock() + entries := make([]sandboxfs.ForgetEntry, 0, len(batch)) + for ref, n := range batch { + entries = append(entries, sandboxfs.ForgetEntry{Node: ref, Count: n}) + } + for len(entries) > 0 { + n := min(len(entries), forgetBatch) + _, err := call(f, f.drainCtx, (*sandboxfs.Client).Forget, &sandboxfs.ForgetRequest{Entries: entries[:n]}) + if err != nil && retryable(err) && !f.dead.Load() { + f.mu.Lock() + for _, e := range entries { + f.forgets[e.Node] += e.Count + } + f.mu.Unlock() + return false + } + entries = entries[n:] + } + return true +} + +// release sends a Release or ReleaseDir for a handle the kernel closed. One that certainly did nothing and may succeed later is queued for the drainer; one that may have reached the service is never sent again. +func (f *frontend) release(c cleanup) { + if f.sendRelease(f.ctx, c) { + return + } + if f.seams.queue != nil { + f.seams.queue() + } + f.mu.Lock() + f.releases = append(f.releases, c) + f.wake() + f.mu.Unlock() +} + +// sendRelease sends c and reports false when it must be sent again. +func (f *frontend) sendRelease(ctx context.Context, c cleanup) bool { + var err error + if c.dir { + _, err = call(f, ctx, (*sandboxfs.Client).ReleaseDir, &sandboxfs.ReleaseDirRequest{Handle: c.handle}) + } else { + _, err = call(f, ctx, (*sandboxfs.Client).Release, &sandboxfs.ReleaseRequest{Handle: c.handle}) + } + return err == nil || !retryable(err) || f.dead.Load() || f.closed.Load() +} diff --git a/apps/daemon/internal/worldfs/errno_linux.go b/apps/daemon/internal/worldfs/errno_linux.go new file mode 100644 index 00000000..bccefd92 --- /dev/null +++ b/apps/daemon/internal/worldfs/errno_linux.go @@ -0,0 +1,69 @@ +//go:build linux + +package worldfs + +import ( + "errors" + "syscall" + + "github.com/MiniMax-AI/OpenAgentCore/internal/sandboxfs" + "github.com/hanwen/go-fuse/v2/fuse" +) + +// errnos maps every sandboxfs errno to its Linux errno. +var errnos = map[sandboxfs.Errno]syscall.Errno{ + sandboxfs.ErrnoPermissionDenied: syscall.EACCES, + sandboxfs.ErrnoOperationNotPermitted: syscall.EPERM, + sandboxfs.ErrnoNotFound: syscall.ENOENT, + sandboxfs.ErrnoExists: syscall.EEXIST, + sandboxfs.ErrnoNotDirectory: syscall.ENOTDIR, + sandboxfs.ErrnoIsDirectory: syscall.EISDIR, + sandboxfs.ErrnoDirectoryNotEmpty: syscall.ENOTEMPTY, + sandboxfs.ErrnoInvalidArgument: syscall.EINVAL, + sandboxfs.ErrnoBadDescriptor: syscall.EBADF, + sandboxfs.ErrnoTooManyOpenFiles: syscall.EMFILE, + sandboxfs.ErrnoNoSpace: syscall.ENOSPC, + sandboxfs.ErrnoQuotaExceeded: syscall.EDQUOT, + sandboxfs.ErrnoReadOnlyFilesystem: syscall.EROFS, + sandboxfs.ErrnoCrossDevice: syscall.EXDEV, + sandboxfs.ErrnoNameTooLong: syscall.ENAMETOOLONG, + sandboxfs.ErrnoSymlinkLoop: syscall.ELOOP, + sandboxfs.ErrnoFileTooLarge: syscall.EFBIG, + sandboxfs.ErrnoOverflow: syscall.EOVERFLOW, + sandboxfs.ErrnoBusy: syscall.EBUSY, + sandboxfs.ErrnoAgain: syscall.EAGAIN, + sandboxfs.ErrnoInterrupted: syscall.EINTR, + sandboxfs.ErrnoIO: syscall.EIO, + sandboxfs.ErrnoNoDevice: syscall.ENODEV, + sandboxfs.ErrnoNoSuchDeviceOrAddress: syscall.ENXIO, + sandboxfs.ErrnoBrokenPipe: syscall.EPIPE, + sandboxfs.ErrnoNotSupported: syscall.EOPNOTSUPP, + sandboxfs.ErrnoNoLocks: syscall.ENOLCK, + sandboxfs.ErrnoDeadlock: syscall.EDEADLK, +} + +// errnoOf returns the Linux errno for a failed request: the mapped errno, ESTALE for a stale node or handle, and EIO for anything else. +func errnoOf(err error) syscall.Errno { + var fail *sandboxfs.Failure + if !errors.As(err, &fail) { + return syscall.EIO + } + switch fail.Code { + case sandboxfs.CodeErrno: + if e, ok := errnos[fail.Errno]; ok { + return e + } + case sandboxfs.CodeStaleNode, sandboxfs.CodeStaleHandle: + return syscall.ESTALE + } + return syscall.EIO +} + +func status(err error) fuse.Status { + if err == nil { + return fuse.OK + } + return fuse.Status(errnoOf(err)) +} + +func errno(e syscall.Errno) fuse.Status { return fuse.Status(e) } diff --git a/apps/daemon/internal/worldfs/errors.go b/apps/daemon/internal/worldfs/errors.go new file mode 100644 index 00000000..855eb961 --- /dev/null +++ b/apps/daemon/internal/worldfs/errors.go @@ -0,0 +1,54 @@ +package worldfs + +import ( + "context" + "errors" + "io" + "strings" +) + +// Dial opens a new File stream to the attachment's service. It returns once ctx ends; ctx bounds the open, not the stream. +type Dial func(ctx context.Context) (io.ReadWriteCloser, error) + +// Error kinds. Every error the package returns, and every reason [World.Err] reports, matches one of them with errors.Is. +var ( + ErrUnsupported = errors.New("worldfs: unsupported platform") + ErrIncompatible = errors.New("worldfs: file service profile not supported") + ErrConnect = errors.New("worldfs: cannot reach the file service") + ErrMountpoint = errors.New("worldfs: mountpoint cannot be presented") + ErrInstanceChanged = errors.New("worldfs: file service instance changed") + ErrAttachmentLost = errors.New("worldfs: attachment ended") + ErrTopologyChanged = errors.New("worldfs: pinned topology changed") + // ErrAttachmentDirty is a failed Serve that cannot show the attachment holds nothing: the owner of the Link attachment must end it. + ErrAttachmentDirty = errors.New("worldfs: the attachment may still hold state") +) + +// Error is a typed worldfs failure. It matches Kind and, when present, Err. +type Error struct { + Kind error + Op string + Path string + Err error +} + +func (e *Error) Error() string { + var b strings.Builder + b.WriteString(e.Kind.Error()) + if e.Op != "" { + b.WriteString(": " + e.Op) + } + if e.Path != "" { + b.WriteString(" " + e.Path) + } + if e.Err != nil { + b.WriteString(": " + e.Err.Error()) + } + return b.String() +} + +func (e *Error) Unwrap() []error { + if e.Err == nil { + return []error{e.Kind} + } + return []error{e.Kind, e.Err} +} diff --git a/apps/daemon/internal/worldfs/files_linux.go b/apps/daemon/internal/worldfs/files_linux.go new file mode 100644 index 00000000..86fda7f3 --- /dev/null +++ b/apps/daemon/internal/worldfs/files_linux.go @@ -0,0 +1,161 @@ +//go:build linux + +package worldfs + +import ( + "syscall" + + "github.com/MiniMax-AI/OpenAgentCore/internal/sandboxfs" + "github.com/hanwen/go-fuse/v2/fuse" +) + +// openFlags maps open(2) flags. ok is false for an invalid access mode. +func openFlags(flags uint32) (acc sandboxfs.AccessMode, of sandboxfs.OpenFlags, ok bool) { + switch flags & syscall.O_ACCMODE { + case syscall.O_RDONLY: + acc = sandboxfs.AccessRead + case syscall.O_WRONLY: + acc = sandboxfs.AccessWrite + case syscall.O_RDWR: + acc = sandboxfs.AccessReadWrite + default: + return 0, 0, false + } + for _, m := range []struct { + bit uint32 + flag sandboxfs.OpenFlags + }{{syscall.O_APPEND, sandboxfs.OpenAppend}, {syscall.O_TRUNC, sandboxfs.OpenTruncate}, {syscall.O_NOFOLLOW, sandboxfs.OpenNoFollow}} { + if flags&m.bit != 0 { + of |= m.flag + } + } + switch { + case flags&syscall.O_SYNC == syscall.O_SYNC: + of |= sandboxfs.OpenSync + case flags&syscall.O_DSYNC != 0: + of |= sandboxfs.OpenDataSync + } + return acc, of, true +} + +func (f *frontend) Open(_ <-chan struct{}, in *fuse.OpenIn, out *fuse.OpenOut) fuse.Status { + n, st := f.node(in.NodeId) + if !st.Ok() { + return st + } + acc, of, ok := openFlags(in.Flags) + if !ok { + return fuse.EINVAL + } + h := &handle{node: n} + if n.synthetic() { + if acc.Writes() || of&sandboxfs.OpenTruncate != 0 { + return fuse.EPERM + } + } else { + r, err := call(f, f.ctx, (*sandboxfs.Client).Open, &sandboxfs.OpenRequest{Node: n.ref, Access: acc, Flags: of}) + if err != nil { + return status(err) + } + h.server = r.Handle + } + *out = fuse.OpenOut{Fh: f.newHandle(h), OpenFlags: fuse.FOPEN_DIRECT_IO} + return fuse.OK +} + +func (f *frontend) Create(_ <-chan struct{}, in *fuse.CreateIn, name string, out *fuse.CreateOut) fuse.Status { + p, st := f.parent(in.NodeId, name) + if !st.Ok() { + return st + } + acc, of, ok := openFlags(in.Flags) + if !ok { + return fuse.EINVAL + } + r, err := call(f, f.ctx, (*sandboxfs.Client).Create, &sandboxfs.CreateRequest{ + Parent: p.ref, Name: []byte(name), Mode: in.Mode & sandboxfs.ModePerm, Access: acc, Flags: of, Exclusive: in.Flags&syscall.O_EXCL != 0, + }) + if err != nil { + return status(err) + } + n := f.adopt(r.Entry, &out.EntryOut) + out.OpenOut = fuse.OpenOut{Fh: f.newHandle(&handle{node: n, server: r.Handle}), OpenFlags: fuse.FOPEN_DIRECT_IO} + return fuse.OK +} + +func (f *frontend) Read(_ <-chan struct{}, in *fuse.ReadIn, _ []byte) (fuse.ReadResult, fuse.Status) { + h, st := f.handle(in.Fh) + if !st.Ok() { + return nil, st + } + if h.server == 0 { + return fuse.ReadResultData(nil), fuse.OK + } + r, err := call(f, f.ctx, (*sandboxfs.Client).Read, &sandboxfs.ReadRequest{Handle: h.server, Offset: in.Offset, Size: min(in.Size, f.caps.MaxReadBytes)}) + if err != nil { + return nil, status(err) + } + return fuse.ReadResultData(r.Data), fuse.OK +} + +// Write reports a short write when the service stopped after a prefix. +func (f *frontend) Write(_ <-chan struct{}, in *fuse.WriteIn, data []byte) (uint32, fuse.Status) { + h, st := f.handle(in.Fh) + if !st.Ok() { + return 0, st + } + if h.server == 0 { + return 0, fuse.EBADF + } + data = data[:min(uint32(len(data)), f.caps.MaxWriteBytes)] + r, err := call(f, f.ctx, (*sandboxfs.Client).Write, &sandboxfs.WriteRequest{Handle: h.server, Offset: in.Offset, Data: data}) + if err != nil { + return 0, status(err) + } + return r.Written, fuse.OK +} + +func (f *frontend) Flush(_ <-chan struct{}, in *fuse.FlushIn) fuse.Status { + h, st := f.handle(in.Fh) + if !st.Ok() || h.server == 0 { + return st + } + _, err := call(f, f.ctx, (*sandboxfs.Client).Flush, &sandboxfs.FlushRequest{Handle: h.server, Owner: sandboxfs.LockOwner(in.LockOwner)}) + return status(err) +} + +func (f *frontend) Fsync(_ <-chan struct{}, in *fuse.FsyncIn) fuse.Status { + h, st := f.handle(in.Fh) + if !st.Ok() || h.server == 0 { + return st + } + _, err := call(f, f.ctx, (*sandboxfs.Client).Fsync, &sandboxfs.FsyncRequest{Handle: h.server, DataOnly: in.FsyncFlags&1 != 0}) + return status(err) +} + +// Release releases a failed handle too, which drops any lock the service holds on it. +func (f *frontend) Release(_ <-chan struct{}, in *fuse.ReleaseIn) { + if h := f.dropHandle(in.Fh); h != nil && h.server != 0 && !f.closed.Load() { + f.release(cleanup{handle: h.server}) + } +} + +// Lseek has no File request. ENOSYS makes the kernel seek itself; SEEK_DATA and SEEK_HOLE then treat the file as one data extent. +func (f *frontend) Lseek(<-chan struct{}, *fuse.LseekIn, *fuse.LseekOut) fuse.Status { + return fuse.ENOSYS +} + +// Fallocate has no File request. ENOSYS makes the kernel stop asking and answer EOPNOTSUPP itself. +func (f *frontend) Fallocate(<-chan struct{}, *fuse.FallocateIn) fuse.Status { + return fuse.ENOSYS +} + +// CopyFileRange has no File request. ENOSYS makes the kernel copy with Read and Write. +func (f *frontend) CopyFileRange(<-chan struct{}, *fuse.CopyFileRangeIn) (uint32, fuse.Status) { + return 0, fuse.ENOSYS +} + +// Ioctl has no File request: no world file takes ioctls. +func (f *frontend) Ioctl(<-chan struct{}, *fuse.IoctlIn, []byte, *fuse.IoctlOut, []byte) fuse.Status { + return errno(syscall.ENOTTY) +} diff --git a/apps/daemon/internal/worldfs/frontend_linux_test.go b/apps/daemon/internal/worldfs/frontend_linux_test.go new file mode 100644 index 00000000..a0642fa6 --- /dev/null +++ b/apps/daemon/internal/worldfs/frontend_linux_test.go @@ -0,0 +1,351 @@ +//go:build linux + +package worldfs + +import ( + "context" + "errors" + "io" + "math" + "net" + "os" + "path/filepath" + "sync" + "sync/atomic" + "syscall" + "testing" + "time" + + "github.com/MiniMax-AI/OpenAgentCore/apps/daemon/internal/sessionview" + "github.com/MiniMax-AI/OpenAgentCore/apps/sandboxio/fileservicetest" + "github.com/MiniMax-AI/OpenAgentCore/internal/sandboxfs" + "github.com/MiniMax-AI/OpenAgentCore/internal/sandboxwire" + "github.com/hanwen/go-fuse/v2/fuse" + "golang.org/x/sys/unix" +) + +// These tests drive the frontend's FUSE methods directly, without a mount, so they need no privileges. + +// counts records what the service ran. It refuses Releases while refuse is positive, and when held is set it closes held once a waiting lock request arrives and ends that request Cancelled with EffectPossible once it is cancelled. +type counts struct { + refuse atomic.Int32 + released atomic.Int32 + held chan struct{} + + mu sync.Mutex + tries []time.Time // when each Release arrived +} + +type counting struct { + sandboxfs.Service + *counts +} + +func (c counting) Release(ctx context.Context, a sandboxfs.Attachment, q *sandboxfs.ReleaseRequest) (*sandboxfs.ReleaseResponse, error) { + c.mu.Lock() + c.tries = append(c.tries, time.Now()) + c.mu.Unlock() + if c.refuse.Add(-1) >= 0 { + return nil, sandboxfs.NewFailure(sandboxfs.CodeResourceExhausted, sandboxwire.EffectNone, "too many requests in flight") + } + r, err := c.Service.Release(ctx, a, q) + if err == nil { + c.released.Add(1) + } + return r, err +} + +func (c counting) SetLock(ctx context.Context, a sandboxfs.Attachment, q *sandboxfs.SetLockRequest) (*sandboxfs.SetLockResponse, error) { + if c.held == nil || !q.Wait { + return c.Service.SetLock(ctx, a, q) + } + close(c.held) + <-ctx.Done() + return nil, sandboxfs.NewFailure(sandboxfs.CodeCancelled, sandboxwire.EffectPossible, "cancelled after it may have locked") +} + +// newServer serves a directory holding the empty file f and returns the file's path. +func newServer(t *testing.T) (*fileservicetest.Server, *counts, string) { + t.Helper() + dir := t.TempDir() + path := filepath.Join(dir, "f") + if err := os.WriteFile(path, nil, 0o644); err != nil { + t.Fatal(err) + } + srv, err := fileservicetest.New(dir) + if err != nil { + t.Fatal(err) + } + t.Cleanup(func() { srv.Close() }) + c := &counts{} + srv.Intercept(func(s sandboxfs.Service) sandboxfs.Service { return counting{s, c} }) + return srv, c, path +} + +// attached returns a frontend attached through dial as Serve leaves it, but neither mounted nor draining. +func attached(t *testing.T, dial Dial) *frontend { + t.Helper() + f := New(fileservicetest.Export, dial).fs + if _, _, err := f.attach(context.Background(), nil); err != nil { + t.Fatal(err) + } + t.Cleanup(f.shutdown) + return f +} + +func draining(t *testing.T, f *frontend) { + go f.drain() + t.Cleanup(func() { + f.stopDrain() + <-f.drained + }) +} + +func open(t *testing.T, f *frontend, name string) uint64 { + t.Helper() + var e fuse.EntryOut + if st := f.Lookup(nil, &fuse.InHeader{NodeId: f.root.id}, name, &e); !st.Ok() { + t.Fatalf("Lookup %s: %v", name, st) + } + var o fuse.OpenOut + if st := f.Open(nil, &fuse.OpenIn{InHeader: fuse.InHeader{NodeId: e.NodeId}, Flags: syscall.O_RDWR}, &o); !st.Ok() { + t.Fatalf("Open %s: %v", name, st) + } + return o.Fh +} + +// flock sends a whole-file flock request, as the kernel does. +func flock(f *frontend, cancel <-chan struct{}, fh, owner uint64, typ uint32, wait bool) fuse.Status { + in := &fuse.LkIn{Fh: fh, Owner: owner, LkFlags: fuse.FUSE_LK_FLOCK, Lk: fuse.FileLock{End: math.MaxInt64, Typ: typ}} + if wait { + return f.SetLkw(cancel, in) + } + return f.SetLk(cancel, in) +} + +func eventually(t *testing.T, what string, ok func() bool) { + t.Helper() + for deadline := time.Now().Add(2 * time.Second); !ok(); time.Sleep(5 * time.Millisecond) { + if time.Now().After(deadline) { + t.Fatalf("%s did not happen within 2s", what) + } + } +} + +// A read lock requested while a failed conversion awaits its recovery unlock survives that unlock. +func TestLockRecoveryKeepsLaterLock(t *testing.T) { + srv, _, path := newServer(t) + f := attached(t, srv.Dial) + fh, other := open(t, f, "f"), open(t, f, "f") + if flock(f, nil, fh, 1, syscall.F_RDLCK, false) != fuse.OK || flock(f, nil, other, 2, syscall.F_RDLCK, false) != fuse.OK { + t.Fatal("read locks failed") + } + later := make(chan fuse.Status, 1) + f.seams.undo = func() { + go func() { later <- flock(f, nil, fh, 1, syscall.F_RDLCK, false) }() + select { + case st := <-later: + later <- st // it ran before the recovery + case <-time.After(200 * time.Millisecond): + } + } + // The other handle's read lock makes the conversion fail, and the failed conversion drops fh's read lock. + if st := flock(f, nil, fh, 1, syscall.F_WRLCK, false); st == fuse.OK { + t.Fatal("conversion succeeded beside another read lock") + } + if st := <-later; st != fuse.OK { + t.Fatalf("later read lock: %v", st) + } + if flock(f, nil, other, 2, syscall.F_UNLCK, false) != fuse.OK { + t.Fatal("unlock failed") + } + native, err := os.Open(path) + if err != nil { + t.Fatal(err) + } + defer native.Close() + if err := unix.Flock(int(native.Fd()), unix.LOCK_EX|unix.LOCK_NB); err != unix.EWOULDBLOCK { + t.Fatalf("native exclusive lock: %v, want EWOULDBLOCK while the later read lock holds", err) + } +} + +// A Release the service refuses for now is sent again until it runs. +func TestRefusedReleaseIsRetried(t *testing.T) { + srv, c, _ := newServer(t) + c.refuse.Store(1) + f := attached(t, srv.Dial) + fh := open(t, f, "f") + draining(t, f) + f.Release(nil, &fuse.ReleaseIn{Fh: fh}) + eventually(t, "the refused Release", func() bool { return c.released.Load() == 1 }) +} + +// A Release the service refuses right after the drainer redialed to send it waits for the backoff before it goes again. +func TestRefusalAfterOwnRedialWaits(t *testing.T) { + srv, c, _ := newServer(t) + var down atomic.Bool + f := attached(t, func(ctx context.Context) (io.ReadWriteCloser, error) { + if down.Load() { + return nil, errors.New("unreachable") + } + return srv.Dial(ctx) + }) + fh := open(t, f, "f") + down.Store(true) + srv.Break() + <-f.conn.Done() + f.Release(nil, &fuse.ReleaseIn{Fh: fh}) // never sent: queued + c.refuse.Store(2) + down.Store(false) + draining(t, f) + eventually(t, "the refused Release", func() bool { return c.released.Load() == 1 }) + c.mu.Lock() + defer c.mu.Unlock() + if gap := c.tries[1].Sub(c.tries[0]); gap < retryMin { + t.Errorf("sent again after %v, want at least %v", gap, retryMin) + } +} + +// A Release queued just after another request redialed, once the drainer has taken that redial's wakeup, is still sent. +func TestReleaseQueuedAfterRedial(t *testing.T) { + srv, c, _ := newServer(t) + var down atomic.Bool + f := attached(t, func(ctx context.Context) (io.ReadWriteCloser, error) { + if down.Load() { + return nil, errors.New("unreachable") + } + return srv.Dial(ctx) + }) + fh := open(t, f, "f") + down.Store(true) + srv.Break() + <-f.conn.Done() + f.seams.queue = func() { + down.Store(false) + if _, err := f.client(context.Background(), nil); err != nil { + t.Errorf("redial: %v", err) + } + <-f.kick + } + f.Release(nil, &fuse.ReleaseIn{Fh: fh}) + draining(t, f) + eventually(t, "the queued Release", func() bool { return c.released.Load() == 1 }) +} + +// stalled returns a frontend whose stream has failed and whose service is unreachable, with fh open, and a channel closed once a redial waits. +func stalled(t *testing.T) (f *frontend, fh uint64, dialing chan struct{}) { + srv, _, _ := newServer(t) + var stall atomic.Bool + var once sync.Once + dialing = make(chan struct{}) + f = attached(t, func(ctx context.Context) (io.ReadWriteCloser, error) { + if stall.Load() { + once.Do(func() { close(dialing) }) + } + return srv.Dial(ctx) + }) + fh = open(t, f, "f") + stall.Store(true) + srv.Stall() + <-f.conn.Done() + return f, fh, dialing +} + +func wantStatus(t *testing.T, got <-chan fuse.Status, want fuse.Status) { + t.Helper() + select { + case st := <-got: + if st != want { + t.Fatalf("status %v, want %v", st, want) + } + case <-time.After(2 * time.Second): + t.Fatalf("no status within 2s, want %v", want) + } +} + +// A waiting lock blocked on redialing an unreachable service returns EINTR once interrupted. +func TestLockInterruptedWhileRedialing(t *testing.T) { + f, fh, dialing := stalled(t) + cancel := make(chan struct{}) + got := make(chan fuse.Status, 1) + go func() { got <- flock(f, cancel, fh, 1, syscall.F_WRLCK, true) }() + <-dialing + close(cancel) + wantStatus(t, got, fuse.EINTR) +} + +// An interrupted waiting lock that was never sent returns at once, although a non-blocking request of its owner holds the turn while it waits for the stream. +func TestUnsentLockInterruptedBehindTurn(t *testing.T) { + f, fh, dialing := stalled(t) + go f.client(f.ctx, nil) // another request redials and holds the stream + <-dialing + h, _ := f.handle(fh) + o := h.order(1) + registered := func(n uint64) func() bool { + return func() bool { + o.mu.Lock() + defer o.mu.Unlock() + return o.acquired == n + } + } + cancel := make(chan struct{}) + got := make(chan fuse.Status, 1) + go func() { got <- flock(f, cancel, fh, 1, syscall.F_WRLCK, true) }() + eventually(t, "the waiting lock's registration", registered(1)) + go flock(f, nil, fh, 1, syscall.F_RDLCK, false) + eventually(t, "the non-blocking lock's turn", registered(2)) + close(cancel) + wantStatus(t, got, fuse.EINTR) +} + +// A waiting lock cancelled after it may have locked, once another acquisition of its owner has completed, cannot be undone safely: it returns EIO. +func TestUnsafeLockRecoveryFails(t *testing.T) { + srv, c, _ := newServer(t) + c.held = make(chan struct{}) + f := attached(t, srv.Dial) + fh := open(t, f, "f") + cancel := make(chan struct{}) + got := make(chan fuse.Status, 1) + go func() { got <- flock(f, cancel, fh, 1, syscall.F_WRLCK, true) }() + <-c.held + if st := flock(f, nil, fh, 1, syscall.F_RDLCK, false); st != fuse.OK { + t.Fatalf("read lock: %v", st) + } + close(cancel) + wantStatus(t, got, fuse.EIO) +} + +// Serve gives up when Start's context ends, even while Describe gets no answer, and Stop then returns at once. +func TestServeEndsWithItsContext(t *testing.T) { + w := New(fileservicetest.Export, func(context.Context) (io.ReadWriteCloser, error) { + c, s := net.Pipe() + go io.Copy(io.Discard, s) // reads requests and never answers + t.Cleanup(func() { s.Close() }) + return c, nil + }) + ctx, cancel := context.WithTimeout(context.Background(), 100*time.Millisecond) + defer cancel() + served := make(chan error, 1) + go func() { + _, _, err := w.Serve(ctx, nil, sessionview.WorldMount{}) + served <- err + }() + select { + case err := <-served: + if !errors.Is(err, ErrConnect) || !errors.Is(err, context.DeadlineExceeded) { + t.Fatalf("Serve = %v, want ErrConnect with the deadline", err) + } + case <-time.After(2 * time.Second): + t.Fatal("Serve still waiting 2s after its context ended") + } + stopped := make(chan error, 1) + go func() { stopped <- w.Stop() }() + select { + case err := <-stopped: + if err != nil { + t.Fatalf("Stop: %v", err) + } + case <-time.After(2 * time.Second): + t.Fatal("Stop blocked after a failed Serve") + } +} diff --git a/apps/daemon/internal/worldfs/listing_linux.go b/apps/daemon/internal/worldfs/listing_linux.go new file mode 100644 index 00000000..d4425f42 --- /dev/null +++ b/apps/daemon/internal/worldfs/listing_linux.go @@ -0,0 +1,124 @@ +//go:build linux + +package worldfs + +import ( + "github.com/MiniMax-AI/OpenAgentCore/internal/sandboxfs" + "github.com/hanwen/go-fuse/v2/fuse" +) + +// listing pages a sandbox directory that has presented children. Every sandbox entry is listed at its own position, a presented name with the presented node's type and inode number, and the presented names the sandbox does not list follow its last entry. An offset names the position after one listed entry. Offsets are assigned when an entry is first listed and kept until ReleaseDir, so a rewind or seekdir never invalidates one, and each READDIR sends at most one ReadDir. +type listing struct { + at []position // at[o-1] is where offset o resumes + offs map[position]uint64 // the offset of each position listed so far + seen map[string]bool // presented names the sandbox has listed + pass *pass // the enumeration since offset 0, while it skips nothing +} + +// position is a place in a listing: after the sandbox entry whose cookie is cookie, or, once the sandbox listing ended, before the presented name order[from]. +type position struct { + cookie uint64 + tail bool + from int +} + +// pass is an enumeration from offset 0. It is complete when it reaches the end of the sandbox listing having resumed only at offsets it listed. +type pass struct { + offs map[uint64]bool // offsets listed in this pass + names map[string]bool // presented names the sandbox listed in this pass +} + +// readPresented lists from in.Offset. A pinned entry the sandbox lists with another type, or that a complete pass does not find, is reported as a topology change. +func (f *frontend) readPresented(h *handle, in *fuse.ReadIn, out *fuse.DirEntryList) fuse.Status { + h.mu.Lock() + defer h.mu.Unlock() + l, n := &h.list, h.node + if l.offs == nil { + l.offs, l.seen = map[position]uint64{}, map[string]bool{} + } + var p position + switch { + case in.Offset == 0: + l.pass = &pass{offs: map[uint64]bool{}, names: map[string]bool{}} + case in.Offset > uint64(len(l.at)): + return fuse.EINVAL + default: + p = l.at[in.Offset-1] + if l.pass != nil && !l.pass.offs[in.Offset] { + l.pass = nil + } + } + if !p.tail { + r, err := call(f, f.ctx, (*sandboxfs.Client).ReadDir, &sandboxfs.ReadDirRequest{Handle: h.server, Cookie: p.cookie, Limit: min(in.Size, f.caps.MaxReadDirBytes)}) + switch { + case err != nil: + return status(err) + case len(r.Entries) == 0 && !r.End: + return fuse.EIO + } + for _, e := range r.Entries { + name := string(e.Name) + ent := fuse.DirEntry{Mode: e.Type, Name: name, Ino: e.Ino} + c := n.fixed[name] + if c != nil { + if !c.synthetic() && e.Type != c.fileType() { + f.lose(&Error{Kind: ErrTopologyChanged, Op: "readdir", Path: c.path}, false) + } + ent.Mode, ent.Ino = c.fileType(), f.ino(c) + } + if !l.add(out, ent, position{cookie: e.Cookie}) { + return fuse.OK + } + if c != nil { + l.seen[name] = true + if l.pass != nil { + l.pass.names[name] = true + } + } + } + if !r.End { + return fuse.OK + } + if l.pass != nil { + for _, name := range n.order { + if c := n.fixed[name]; !c.synthetic() && !l.pass.names[name] { + f.lose(&Error{Kind: ErrTopologyChanged, Op: "readdir", Path: c.path}, false) + } + } + } + p = position{tail: true} + } + listed := l.seen + if l.pass != nil { + listed = l.pass.names + } + for i := p.from; i < len(n.order); i++ { + name := n.order[i] + if listed[name] { + continue + } + c := n.fixed[name] + if !l.add(out, fuse.DirEntry{Mode: c.fileType(), Name: name, Ino: f.ino(c)}, position{tail: true, from: i + 1}) { + break + } + } + return fuse.OK +} + +// add lists e, after which the listing resumes at p. +func (l *listing) add(out *fuse.DirEntryList, e fuse.DirEntry, p position) bool { + off, ok := l.offs[p] + if !ok { + l.at = append(l.at, p) + off = uint64(len(l.at)) + l.offs[p] = off + } + e.Off = off + if !out.AddDirEntry(e) { + return false + } + if l.pass != nil { + l.pass.offs[off] = true + } + return true +} diff --git a/apps/daemon/internal/worldfs/lock_linux.go b/apps/daemon/internal/worldfs/lock_linux.go new file mode 100644 index 00000000..584bd180 --- /dev/null +++ b/apps/daemon/internal/worldfs/lock_linux.go @@ -0,0 +1,211 @@ +//go:build linux + +package worldfs + +import ( + "context" + "errors" + "sync" + "syscall" + + "github.com/MiniMax-AI/OpenAgentCore/internal/sandboxfs" + "github.com/hanwen/go-fuse/v2/fuse" +) + +func lockMode(typ uint32) (sandboxfs.LockMode, bool) { + switch typ { + case syscall.F_RDLCK: + return sandboxfs.LockRead, true + case syscall.F_WRLCK: + return sandboxfs.LockWrite, true + case syscall.F_UNLCK: + return sandboxfs.LockUnlock, true + } + return 0, false +} + +// GetLk tests a POSIX lock. Without the service's POSIX locks, the kernel's lock calls fail with ENOLCK rather than locking only within the view. +func (f *frontend) GetLk(_ <-chan struct{}, in *fuse.LkIn, out *fuse.LkOut) fuse.Status { + h, st := f.handle(in.Fh) + if !st.Ok() { + return st + } + mode, ok := lockMode(in.Lk.Typ) + switch { + case !ok || mode == sandboxfs.LockUnlock: + return fuse.EINVAL + case h.server == 0 || !f.caps.POSIXLocks: + return errno(syscall.ENOLCK) + } + r, err := call(f, f.ctx, (*sandboxfs.Client).GetLock, &sandboxfs.GetLockRequest{ + Handle: h.server, Owner: sandboxfs.LockOwner(in.Owner), Lock: sandboxfs.Lock{Mode: mode, Start: in.Lk.Start, End: in.Lk.End}, + }) + if err != nil { + return status(err) + } + if r.Conflict == nil { + out.Lk = in.Lk + out.Lk.Typ = syscall.F_UNLCK + return fuse.OK + } + typ := uint32(syscall.F_RDLCK) + if r.Conflict.Mode == sandboxfs.LockWrite { + typ = syscall.F_WRLCK + } + out.Lk = fuse.FileLock{Start: r.Conflict.Start, End: r.Conflict.End, Typ: typ} + return fuse.OK +} + +func (f *frontend) SetLk(cancel <-chan struct{}, in *fuse.LkIn) fuse.Status { + return f.setLock(cancel, in, false) +} + +func (f *frontend) SetLkw(cancel <-chan struct{}, in *fuse.LkIn) fuse.Status { + return f.setLock(cancel, in, true) +} + +// lockOrder orders one lock owner's requests on one handle, so that a recovery unlock never removes a lock another request reported. A non-blocking request holds the turn from before it is sent until its recovery is done. A waiting request holds it only to register and, when it must recover, to recover: it waits without it, so interrupts and unlocks get through, and an outcome that needs no recovery returns without it. +type lockOrder struct { + turn chan struct{} + + mu sync.Mutex + acquired uint64 // acquisitions registered + waiting int // waiting requests registered and not yet answered +} + +// place is where a request registered. +type place struct { + acquired uint64 // acquisitions registered up to and including the request + busy bool // a waiting request was out +} + +func (o *lockOrder) take(interrupt <-chan struct{}) bool { + select { + case o.turn <- struct{}{}: + return true + case <-interrupt: + return false + } +} + +func (o *lockOrder) give() { <-o.turn } + +// register records a request. The caller holds the turn. +func (o *lockOrder) register(acquire, wait bool) place { + o.mu.Lock() + defer o.mu.Unlock() + p := place{busy: o.waiting > 0} + if acquire { + o.acquired++ + } + if wait { + o.waiting++ + } + p.acquired = o.acquired + return p +} + +// answered records that a waiting request ended. +func (o *lockOrder) answered() { + o.mu.Lock() + defer o.mu.Unlock() + o.waiting-- +} + +// alone reports whether no other request of the owner may have locked after the request at p was sent: none registered since, and none was waiting when it registered. The caller holds the turn. +func (o *lockOrder) alone(p place) bool { + o.mu.Lock() + defer o.mu.Unlock() + return !p.busy && o.acquired == p.acquired +} + +func (h *handle) order(owner sandboxfs.LockOwner) *lockOrder { + h.locksMu.Lock() + defer h.locksMu.Unlock() + o := h.locks[owner] + if o == nil { + if h.locks == nil { + h.locks = map[sandboxfs.LockOwner]*lockOrder{} + } + o = &lockOrder{turn: make(chan struct{}, 1)} + h.locks[owner] = o + } + return o +} + +// setLock acquires or releases a POSIX or flock lock and reports what the service did, even when the kernel interrupted the request. A request that may have changed the lock but failed is undone with an unlock of the same owner and range, unless another request of the owner may have locked since; then, or when its outcome cannot be learned or the unlock fails, the handle fails and the request returns EIO. +func (f *frontend) setLock(cancel <-chan struct{}, in *fuse.LkIn, wait bool) fuse.Status { + h, st := f.handle(in.Fh) + if !st.Ok() { + return st + } + mode, ok := lockMode(in.Lk.Typ) + if !ok { + return fuse.EINVAL + } + kind, supported := sandboxfs.LockPOSIX, f.caps.POSIXLocks + if in.LkFlags&fuse.FUSE_LK_FLOCK != 0 { + kind, supported = sandboxfs.LockFlock, f.caps.Flock + } + if h.server == 0 || !supported { + return errno(syscall.ENOLCK) + } + q := &sandboxfs.SetLockRequest{ + Handle: h.server, Kind: kind, Owner: sandboxfs.LockOwner(in.Owner), Lock: sandboxfs.Lock{Mode: mode, Start: in.Lk.Start, End: in.Lk.End}, Wait: wait, + } + o := h.order(q.Owner) + if !o.take(cancel) { + return fuse.EINTR + } + p := o.register(mode != sandboxfs.LockUnlock, wait) + held := !wait + if wait { + o.give() + } + _, err := callUntil(f, f.ctx, cancel, (*sandboxfs.Client).SetLock, q) + if wait { + o.answered() + } + recovered := true + if err != nil && !noEffect(err) { + if !held { + o.take(nil) + held = true + } + if f.seams.undo != nil { + f.seams.undo() + } + recovered = f.undoLock(h, q, answered(err) && o.alone(p)) + } + if held { + o.give() + } + var fail *sandboxfs.Failure + switch { + case err == nil: + return fuse.OK + case !recovered: + return fuse.EIO + case errors.Is(err, errInterrupted) || errors.As(err, &fail) && fail.Code == sandboxfs.CodeCancelled || interrupted(cancel) && !answered(err): + return fuse.EINTR + } + return status(err) +} + +// undoLock unlocks the owner's range after a lock request that may have changed it failed, and reports whether it did. The unlock is safe only when the service runs it after that request and no other request of the owner may have locked in between, which safe says; otherwise, or when the unlock fails, the handle fails instead. +func (f *frontend) undoLock(h *handle, q *sandboxfs.SetLockRequest, safe bool) bool { + if safe { + unlock := *q + unlock.Lock.Mode, unlock.Wait = sandboxfs.LockUnlock, false + if _, err := call(f, f.ctx, (*sandboxfs.Client).SetLock, &unlock); err == nil { + return true + } + } + h.failed.Store(true) + return false +} + +// answered reports whether a failure came in the service's response, so a later request on the handle runs after the failed one. After a stream or context failure the request may still be running. +func answered(err error) bool { + return !errors.Is(err, sandboxfs.ErrTransport) && !errors.Is(err, context.Canceled) && !errors.Is(err, context.DeadlineExceeded) +} diff --git a/apps/daemon/internal/worldfs/nodes_linux.go b/apps/daemon/internal/worldfs/nodes_linux.go new file mode 100644 index 00000000..8895d03a --- /dev/null +++ b/apps/daemon/internal/worldfs/nodes_linux.go @@ -0,0 +1,189 @@ +//go:build linux + +package worldfs + +import ( + "sync" + "sync/atomic" + "syscall" + + "github.com/MiniMax-AI/OpenAgentCore/internal/sandboxfs" + "github.com/hanwen/go-fuse/v2/fuse" +) + +// inode is a kernel node ID. Every ID is frontend-local: a server node gets one when first seen, and a synthetic node has no server node. +type inode struct { + id uint64 + ref sandboxfs.NodeRef // zero for a synthetic node + synth uint32 // a synthetic node's file type, else zero + mount bool // a synthetic mountpoint + path string // a presented node's in-root path + attr sandboxfs.Attr // a pinned node's attributes when it was pinned + link []byte // a pinned symlink's target + // fixed holds the presented children: synthetic nodes and the pinned entries on the way to a mountpoint. It never changes once Serve returns. + fixed map[string]*inode + order []string // fixed's names, sorted, in ReadDir order + held bool // root, pinned or synthetic: kept for the view's lifetime + + lookups uint64 // the kernel's lookup count + refs uint64 // server references taken for kernel lookups +} + +func (n *inode) synthetic() bool { return n.synth != 0 } + +func (n *inode) fileType() uint32 { + if n.synthetic() { + return n.synth + } + return n.attr.Mode & sandboxfs.ModeType +} + +// handle is an open file handle the kernel holds. +type handle struct { + node *inode + server sandboxfs.HandleID // zero for a synthetic node + + failed atomic.Bool // the service may hold a lock on the handle that no request reported: every request but Release fails + + locksMu sync.Mutex + locks map[sandboxfs.LockOwner]*lockOrder + + mu sync.Mutex + list listing // a presented directory's offsets +} + +// newInode registers a node for ref, or a synthetic node when ref is zero. The caller holds mu, or Serve has not returned. +func (f *frontend) newInode(ref sandboxfs.NodeRef) *inode { + f.lastID++ + n := &inode{id: f.lastID, ref: ref} + f.nodes[n.id] = n + if ref != (sandboxfs.NodeRef{}) { + f.byRef[ref] = n + } + return n +} + +// node returns the node the kernel names. +func (f *frontend) node(id uint64) (*inode, fuse.Status) { + if f.dead.Load() || f.closed.Load() { + return nil, fuse.EIO + } + f.mu.Lock() + defer f.mu.Unlock() + if n := f.nodes[id]; n != nil { + return n, fuse.OK + } + return nil, fuse.Status(syscall.ESTALE) +} + +// adopt records the lookup reference a response carried and answers the kernel with the entry. +func (f *frontend) adopt(e sandboxfs.Entry, out *fuse.EntryOut) *inode { + f.mu.Lock() + n := f.byRef[e.Node] + if n == nil { + n = f.newInode(e.Node) + } + n.lookups++ + n.refs++ + f.mu.Unlock() + f.entry(n, e.Attr, out) + return n +} + +// lookedUp counts a kernel lookup answered without a server reference. +func (f *frontend) lookedUp(n *inode) { + f.mu.Lock() + n.lookups++ + f.mu.Unlock() +} + +// entry fills an entry reply. Timeouts stay zero, so the kernel revalidates every name and attribute. +func (f *frontend) entry(n *inode, a sandboxfs.Attr, out *fuse.EntryOut) { + *out = fuse.EntryOut{NodeId: n.id} + f.fill(n, a, &out.Attr) +} + +// attrOf returns a synthetic node's attributes: an empty read-only directory or file, root-owned. +func (f *frontend) attrOf(n *inode) sandboxfs.Attr { + a := sandboxfs.Attr{Ino: 1<<63 | n.id, Mode: sandboxfs.ModeDirectory | 0o555, Nlink: 2, Blksize: 4096, Atime: f.born, Mtime: f.born, Ctime: f.born} + if n.synth == sandboxfs.ModeRegular { + a.Mode, a.Nlink = sandboxfs.ModeRegular|0o444, 1 + } + return a +} + +// fill converts attributes. Owners the service acts as appear as the view's identity; synthetic nodes stay root-owned. +func (f *frontend) fill(n *inode, a sandboxfs.Attr, out *fuse.Attr) { + if !n.synthetic() { + if a.UID == f.service.UID { + a.UID = f.view.UID + } + if a.GID == f.service.GID { + a.GID = f.view.GID + } + } + *out = fuse.Attr{ + Ino: a.Ino, Size: a.Size, Blocks: a.Blocks, + Atime: uint64(a.Atime.Sec), Mtime: uint64(a.Mtime.Sec), Ctime: uint64(a.Ctime.Sec), + Atimensec: a.Atime.Nsec, Mtimensec: a.Mtime.Nsec, Ctimensec: a.Ctime.Nsec, + Mode: a.Mode, Nlink: a.Nlink, Rdev: uint32(a.Rdev), Blksize: a.Blksize, + Owner: fuse.Owner{Uid: a.UID, Gid: a.GID}, + } +} + +// Forget releases kernel lookups. The server references they carried are queued for one Forget request. +func (f *frontend) Forget(id, count uint64) { + f.mu.Lock() + defer f.mu.Unlock() + n := f.nodes[id] + if n == nil { + return + } + n.lookups -= min(count, n.lookups) + if r := min(count, n.refs); r > 0 { + n.refs -= r + f.unref(n.ref, r) + } + if n.lookups == 0 && !n.held { + delete(f.nodes, id) + delete(f.byRef, n.ref) + } +} + +// unref queues count server references of ref for a Forget request. The caller holds mu, or Serve has not returned. +func (f *frontend) unref(ref sandboxfs.NodeRef, count uint64) { + f.forgets[ref] += count + f.wake() +} + +func (f *frontend) newHandle(h *handle) uint64 { + f.mu.Lock() + defer f.mu.Unlock() + f.lastFh++ + f.handles[f.lastFh] = h + return f.lastFh +} + +func (f *frontend) handle(fh uint64) (*handle, fuse.Status) { + if f.dead.Load() || f.closed.Load() { + return nil, fuse.EIO + } + f.mu.Lock() + defer f.mu.Unlock() + switch h := f.handles[fh]; { + case h == nil: + return nil, fuse.EBADF + case h.failed.Load(): + return nil, fuse.EIO + default: + return h, fuse.OK + } +} + +func (f *frontend) dropHandle(fh uint64) *handle { + f.mu.Lock() + defer f.mu.Unlock() + h := f.handles[fh] + delete(f.handles, fh) + return h +} diff --git a/apps/daemon/internal/worldfs/ops_linux.go b/apps/daemon/internal/worldfs/ops_linux.go new file mode 100644 index 00000000..57c57c5f --- /dev/null +++ b/apps/daemon/internal/worldfs/ops_linux.go @@ -0,0 +1,355 @@ +//go:build linux + +package worldfs + +import ( + "syscall" + + "github.com/MiniMax-AI/OpenAgentCore/internal/sandboxfs" + "github.com/hanwen/go-fuse/v2/fuse" + "golang.org/x/sys/unix" +) + +func (f *frontend) String() string { return "oac-world" } +func (f *frontend) SetDebug(bool) {} +func (f *frontend) Init(*fuse.Server) {} +func (f *frontend) OnUnmount() {} + +func (f *frontend) Lookup(_ <-chan struct{}, in *fuse.InHeader, name string, out *fuse.EntryOut) fuse.Status { + p, st := f.node(in.NodeId) + if !st.Ok() { + return st + } + if c := p.fixed[name]; c != nil { + if c.synthetic() { + f.lookedUp(c) + f.entry(c, f.attrOf(c), out) + return fuse.OK + } + return f.lookupPinned(p, name, c, out) + } + if p.synthetic() { + return fuse.ENOENT + } + r, err := call(f, f.ctx, (*sandboxfs.Client).Lookup, &sandboxfs.LookupRequest{Parent: p.ref, Name: []byte(name)}) + if err != nil { + return status(err) + } + f.adopt(r.Entry, out) + return fuse.OK +} + +// lookupPinned answers a pinned name with the pinned node. When the sandbox replaced or removed it, the view keeps the pinned node, since a new node ID would detach the mounts beneath it, and the world reports the change. +func (f *frontend) lookupPinned(p *inode, name string, c *inode, out *fuse.EntryOut) fuse.Status { + r, err := call(f, f.ctx, (*sandboxfs.Client).Lookup, &sandboxfs.LookupRequest{Parent: p.ref, Name: []byte(name)}) + switch { + case err == nil && r.Entry.Node == c.ref: + f.adopt(r.Entry, out) + return fuse.OK + case err == nil: + f.mu.Lock() + f.unref(r.Entry.Node, 1) + f.mu.Unlock() + case !isErrno(err, sandboxfs.ErrnoNotFound): + return status(err) + } + f.lose(&Error{Kind: ErrTopologyChanged, Op: "lookup", Path: c.path}, false) + f.lookedUp(c) + f.entry(c, c.attr, out) + return fuse.OK +} + +func (f *frontend) GetAttr(_ <-chan struct{}, in *fuse.GetAttrIn, out *fuse.AttrOut) fuse.Status { + n, st := f.node(in.NodeId) + if !st.Ok() { + return st + } + switch { + case n.synthetic(): + f.fill(n, f.attrOf(n), &out.Attr) + return fuse.OK + case n.link != nil: + f.fill(n, n.attr, &out.Attr) + return fuse.OK + } + t, st := f.target(n, in.Flags()&fuse.FUSE_GETATTR_FH != 0, in.Fh()) + if !st.Ok() { + return st + } + r, err := call(f, f.ctx, (*sandboxfs.Client).GetAttr, &sandboxfs.GetAttrRequest{Target: t}) + if err != nil { + return status(err) + } + f.fill(n, r.Attr, &out.Attr) + return fuse.OK +} + +// target addresses the open handle when the kernel names one, and the node otherwise. A failed handle fails the request. +func (f *frontend) target(n *inode, useFh bool, fh uint64) (sandboxfs.Target, fuse.Status) { + if useFh { + h, st := f.handle(fh) + switch { + case st == fuse.EIO: + return sandboxfs.Target{}, st + case st.Ok() && h.server != 0: + return sandboxfs.Target{Kind: sandboxfs.TargetHandle, Handle: h.server}, fuse.OK + } + } + return sandboxfs.Target{Kind: sandboxfs.TargetNode, Node: n.ref}, fuse.OK +} + +func (f *frontend) SetAttr(_ <-chan struct{}, in *fuse.SetAttrIn, out *fuse.AttrOut) fuse.Status { + n, st := f.node(in.NodeId) + if !st.Ok() { + return st + } + if n.synthetic() || n.link != nil { + return fuse.EPERM + } + t, st := f.target(n, in.Valid&fuse.FATTR_FH != 0, in.Fh) + if !st.Ok() { + return st + } + q := &sandboxfs.SetAttrRequest{Target: t} + if in.Valid&fuse.FATTR_MODE != 0 { + q.Set |= sandboxfs.AttrMode + q.Mode = in.Mode & sandboxfs.ModePerm + } + if in.Valid&fuse.FATTR_UID != 0 { + q.Set |= sandboxfs.AttrUID + q.UID = in.Uid + if q.UID == f.view.UID { + q.UID = f.service.UID + } + } + if in.Valid&fuse.FATTR_GID != 0 { + q.Set |= sandboxfs.AttrGID + q.GID = in.Gid + if q.GID == f.view.GID { + q.GID = f.service.GID + } + } + if in.Valid&fuse.FATTR_SIZE != 0 { + q.Set |= sandboxfs.AttrSize + q.Size = in.Size + } + // The kernel marks a time set to now with both bits. + switch { + case in.Valid&fuse.FATTR_ATIME_NOW != 0: + q.Set |= sandboxfs.AttrAtimeNow + case in.Valid&fuse.FATTR_ATIME != 0: + q.Set |= sandboxfs.AttrAtime + q.Atime = sandboxfs.Timestamp{Sec: int64(in.Atime), Nsec: in.Atimensec} + } + switch { + case in.Valid&fuse.FATTR_MTIME_NOW != 0: + q.Set |= sandboxfs.AttrMtimeNow + case in.Valid&fuse.FATTR_MTIME != 0: + q.Set |= sandboxfs.AttrMtime + q.Mtime = sandboxfs.Timestamp{Sec: int64(in.Mtime), Nsec: in.Mtimensec} + } + switch { + case q.Set&sandboxfs.AttrMode != 0 && !f.caps.SetMode, + q.Set&(sandboxfs.AttrUID|sandboxfs.AttrGID) != 0 && !f.caps.SetOwner, + q.Set&(sandboxfs.AttrAtime|sandboxfs.AttrMtime|sandboxfs.AttrAtimeNow|sandboxfs.AttrMtimeNow) != 0 && !f.caps.SetTimes: + return fuse.EPERM + case q.Set == 0: + // Nothing the File protocol sets, such as a ctime-only update: report the current attributes. + r, err := call(f, f.ctx, (*sandboxfs.Client).GetAttr, &sandboxfs.GetAttrRequest{Target: q.Target}) + if err != nil { + return status(err) + } + f.fill(n, r.Attr, &out.Attr) + return fuse.OK + } + r, err := call(f, f.ctx, (*sandboxfs.Client).SetAttr, q) + if err != nil { + return status(err) + } + f.fill(n, r.Attr, &out.Attr) + return fuse.OK +} + +func (f *frontend) Access(_ <-chan struct{}, in *fuse.AccessIn) fuse.Status { + n, st := f.node(in.NodeId) + if !st.Ok() { + return st + } + mask := sandboxfs.AccessMask(in.Mask & 7) + if n.synthetic() { + if mask&sandboxfs.MayWrite != 0 || mask&sandboxfs.MayExecute != 0 && n.synth != sandboxfs.ModeDirectory { + return fuse.EACCES + } + return fuse.OK + } + _, err := call(f, f.ctx, (*sandboxfs.Client).Access, &sandboxfs.AccessRequest{Node: n.ref, Mask: mask}) + return status(err) +} + +func (f *frontend) Readlink(_ <-chan struct{}, in *fuse.InHeader) ([]byte, fuse.Status) { + n, st := f.node(in.NodeId) + switch { + case !st.Ok(): + return nil, st + case n.link != nil: + return n.link, fuse.OK + case n.synthetic(): + return nil, fuse.EINVAL + } + r, err := call(f, f.ctx, (*sandboxfs.Client).Readlink, &sandboxfs.ReadlinkRequest{Node: n.ref}) + if err != nil { + return nil, status(err) + } + return r.Target, fuse.OK +} + +// parent returns the directory a name is created in or removed from. Presented names and synthetic directories never change. +func (f *frontend) parent(id uint64, name string) (*inode, fuse.Status) { + p, st := f.node(id) + if !st.Ok() { + return nil, st + } + if p.synthetic() || p.fixed[name] != nil { + return nil, fuse.EPERM + } + return p, fuse.OK +} + +func (f *frontend) Mkdir(_ <-chan struct{}, in *fuse.MkdirIn, name string, out *fuse.EntryOut) fuse.Status { + p, st := f.parent(in.NodeId, name) + if !st.Ok() { + return st + } + r, err := call(f, f.ctx, (*sandboxfs.Client).Mkdir, &sandboxfs.MkdirRequest{Parent: p.ref, Name: []byte(name), Mode: in.Mode & sandboxfs.ModePerm}) + if err != nil { + return status(err) + } + f.adopt(r.Entry, out) + return fuse.OK +} + +// Mknod has no File request. The kernel creates regular files with Create. +func (f *frontend) Mknod(<-chan struct{}, *fuse.MknodIn, string, *fuse.EntryOut) fuse.Status { + return errno(syscall.EOPNOTSUPP) +} + +func (f *frontend) Unlink(_ <-chan struct{}, in *fuse.InHeader, name string) fuse.Status { + p, st := f.parent(in.NodeId, name) + if !st.Ok() { + return st + } + _, err := call(f, f.ctx, (*sandboxfs.Client).Unlink, &sandboxfs.UnlinkRequest{Parent: p.ref, Name: []byte(name)}) + return status(err) +} + +func (f *frontend) Rmdir(_ <-chan struct{}, in *fuse.InHeader, name string) fuse.Status { + p, st := f.parent(in.NodeId, name) + if !st.Ok() { + return st + } + _, err := call(f, f.ctx, (*sandboxfs.Client).Rmdir, &sandboxfs.RmdirRequest{Parent: p.ref, Name: []byte(name)}) + return status(err) +} + +func (f *frontend) Rename(_ <-chan struct{}, in *fuse.RenameIn, oldName, newName string) fuse.Status { + p, st := f.parent(in.NodeId, oldName) + if !st.Ok() { + return st + } + np, st := f.parent(in.Newdir, newName) + if !st.Ok() { + return st + } + var mode sandboxfs.RenameMode + switch { + case in.Flags == 0: + mode = sandboxfs.RenameReplace + case in.Flags == unix.RENAME_NOREPLACE && f.caps.RenameNoReplace: + mode = sandboxfs.RenameNoReplace + case in.Flags == unix.RENAME_EXCHANGE && f.caps.RenameExchange: + mode = sandboxfs.RenameExchange + default: + return fuse.EINVAL + } + _, err := call(f, f.ctx, (*sandboxfs.Client).Rename, &sandboxfs.RenameRequest{Parent: p.ref, Name: []byte(oldName), NewParent: np.ref, NewName: []byte(newName), Mode: mode}) + return status(err) +} + +func (f *frontend) Link(_ <-chan struct{}, in *fuse.LinkIn, name string, out *fuse.EntryOut) fuse.Status { + n, st := f.node(in.Oldnodeid) + if !st.Ok() { + return st + } + np, st := f.parent(in.NodeId, name) + if !st.Ok() { + return st + } + if n.synthetic() || !f.caps.HardLinks { + return fuse.EPERM + } + r, err := call(f, f.ctx, (*sandboxfs.Client).Link, &sandboxfs.LinkRequest{Node: n.ref, NewParent: np.ref, NewName: []byte(name)}) + if err != nil { + return status(err) + } + f.adopt(r.Entry, out) + return fuse.OK +} + +func (f *frontend) Symlink(_ <-chan struct{}, in *fuse.InHeader, target, name string, out *fuse.EntryOut) fuse.Status { + p, st := f.parent(in.NodeId, name) + if !st.Ok() { + return st + } + if !f.caps.Symlinks { + return fuse.EPERM + } + r, err := call(f, f.ctx, (*sandboxfs.Client).Symlink, &sandboxfs.SymlinkRequest{Parent: p.ref, Name: []byte(name), Target: []byte(target)}) + if err != nil { + return status(err) + } + f.adopt(r.Entry, out) + return fuse.OK +} + +func (f *frontend) StatFs(_ <-chan struct{}, in *fuse.InHeader, out *fuse.StatfsOut) fuse.Status { + n, st := f.node(in.NodeId) + if !st.Ok() { + return st + } + ref := n.ref + if n.synthetic() { + ref = f.root.ref + } + r, err := call(f, f.ctx, (*sandboxfs.Client).StatFS, &sandboxfs.StatFSRequest{Node: ref}) + if err != nil { + return status(err) + } + *out = fuse.StatfsOut{ + Blocks: r.Blocks, Bfree: r.BlocksFree, Bavail: r.BlocksAvailable, Files: r.Files, Ffree: r.FilesFree, + Bsize: r.BlockSize, NameLen: r.NameMax, Frsize: r.FragmentSize, + } + return fuse.OK +} + +// The File protocol has no extended attributes. ENOSYS makes the kernel stop asking and answer EOPNOTSUPP itself. + +func (f *frontend) GetXAttr(<-chan struct{}, *fuse.InHeader, string, []byte) (uint32, fuse.Status) { + return 0, fuse.ENOSYS +} + +func (f *frontend) ListXAttr(<-chan struct{}, *fuse.InHeader, []byte) (uint32, fuse.Status) { + return 0, fuse.ENOSYS +} + +func (f *frontend) SetXAttr(<-chan struct{}, *fuse.SetXAttrIn, string, []byte) fuse.Status { + return fuse.ENOSYS +} + +func (f *frontend) RemoveXAttr(<-chan struct{}, *fuse.InHeader, string) fuse.Status { + return fuse.ENOSYS +} + +// Statx has no File request. ENOSYS makes the kernel fall back to GetAttr. +func (f *frontend) Statx(<-chan struct{}, *fuse.StatxIn, *fuse.StatxOut) fuse.Status { + return fuse.ENOSYS +} diff --git a/apps/daemon/internal/worldfs/present_linux.go b/apps/daemon/internal/worldfs/present_linux.go new file mode 100644 index 00000000..1d4ba16e --- /dev/null +++ b/apps/daemon/internal/worldfs/present_linux.go @@ -0,0 +1,243 @@ +//go:build linux + +package worldfs + +import ( + "context" + "errors" + "fmt" + "slices" + "strings" + "syscall" + + "github.com/MiniMax-AI/OpenAgentCore/apps/daemon/internal/sessionview" + "github.com/MiniMax-AI/OpenAgentCore/internal/sandboxfs" +) + +// maxHops is how many symlinks resolving one mountpoint may traverse, the kernel's MAXSYMLINKS. +const maxHops = 40 + +// present resolves each mountpoint inside the world within ctx and installs the synthetic nodes the view needs. It runs before serving starts, so it needs no locks. +func (f *frontend) present(ctx context.Context, mps []sessionview.Mountpoint) (sessionview.Presentation, error) { + var p sessionview.Presentation + for _, mp := range mps { + r := resolver{ctx: ctx, f: f, p: &p, mp: mp} + target, err := r.resolve() + r.dropAhead() + if err != nil { + return sessionview.Presentation{}, err + } + p.Targets = append(p.Targets, target) + } + for _, n := range f.nodes { + for name := range n.fixed { + n.order = append(n.order, name) + } + slices.Sort(n.order) + } + return p, nil +} + +type step struct { + n *inode + name string +} + +// resolver walks one mountpoint's path from the world root, like the kernel would in the view. +type resolver struct { + ctx context.Context + f *frontend + p *sessionview.Presentation + mp sessionview.Mountpoint + stack []step + queue []string + hops int + // ahead holds entries one Walk returned for the next names in queue, and the failure that stopped it. + ahead []sandboxfs.Entry + aheadErr error +} + +func (r *resolver) resolve() (string, error) { + comps := strings.Split(strings.TrimPrefix(r.mp.Path, "/"), "/") + final := comps[len(comps)-1] + if !isName(final) { + return "", r.fail(syscall.EINVAL) + } + r.stack = []step{{n: r.f.root}} + r.queue = slices.Clone(comps[:len(comps)-1]) + for len(r.queue) > 0 { + name := r.queue[0] + r.queue = r.queue[1:] + switch name { + case "", ".": + continue + case "..": + if len(r.stack) > 1 { + r.stack = r.stack[:len(r.stack)-1] + } + continue + } + child, err := r.child(name) + if err != nil { + return "", err + } + switch { + case child.mount: + return "", r.fail(fmt.Errorf("%s is the mountpoint %s", r.at(name), child.path)) + case child.link != nil: + if r.hops++; r.hops > maxHops { + return "", r.fail(syscall.ELOOP) + } + r.dropAhead() + if child.link[0] == '/' { + r.stack = r.stack[:1] + } + r.queue = append(strings.Split(string(child.link), "/"), r.queue...) + case child.fileType() != sandboxfs.ModeDirectory: + return "", r.fail(syscall.ENOTDIR) + default: + r.stack = append(r.stack, step{child, name}) + } + } + cur, at := r.top(), r.at(final) + if c := cur.fixed[final]; c != nil { + return "", r.fail(fmt.Errorf("%s is already presented for %s", at, c.path)) + } + typ := sandboxfs.ModeRegular + if r.mp.Dir { + typ = sandboxfs.ModeDirectory + } + r.synthetic(cur, final, typ).mount = true + return at, nil +} + +// child returns cur's entry name: a presented node, a pinned sandbox entry, or a synthetic directory when the sandbox confirms the name is missing. +func (r *resolver) child(name string) (*inode, error) { + cur := r.top() + if c := cur.fixed[name]; c != nil { + r.dropAhead() + return c, nil + } + if cur.synthetic() { + return r.missing(cur, name), nil + } + if len(r.ahead) == 0 && r.aheadErr == nil { + r.walk(name) + } + if len(r.ahead) == 0 { + err := r.aheadErr + r.aheadErr = nil + if !isErrno(err, sandboxfs.ErrnoNotFound) { + return nil, r.serverErr(err) + } + return r.missing(cur, name), nil + } + e := r.ahead[0] + r.ahead = r.ahead[1:] + return r.pin(cur, name, e) +} + +// walk looks up name and the plain names after it in one request. +func (r *resolver) walk(name string) { + names := [][]byte{[]byte(name)} + for _, n := range r.queue { + if !isName(n) || len(names) == int(r.f.caps.MaxWalkComponents) { + break + } + names = append(names, []byte(n)) + } + resp, err := call(r.f, r.ctx, (*sandboxfs.Client).Walk, &sandboxfs.WalkRequest{Parent: r.top().ref, Names: names}) + switch { + case err != nil: + r.aheadErr = err + case resp.Failure != nil: + r.ahead, r.aheadErr = resp.Entries, resp.Failure + default: + r.ahead = resp.Entries + } +} + +// dropAhead releases the entries walked for names the resolution no longer reaches. +func (r *resolver) dropAhead() { + for _, e := range r.ahead { + r.f.unref(e.Node, 1) + } + r.ahead, r.aheadErr = nil, nil +} + +// pin keeps a sandbox entry on the way to a mountpoint for the view's lifetime. A symlink keeps its target too. +func (r *resolver) pin(cur *inode, name string, e sandboxfs.Entry) (*inode, error) { + n := r.f.byRef[e.Node] + if n != nil { + r.f.unref(e.Node, 1) + } else { + n = r.f.newInode(e.Node) + n.attr, n.path = e.Attr, r.at(name) + } + n.held = true + r.fixed(cur, name, n) + if n.fileType() != sandboxfs.ModeSymlink || n.link != nil { + return n, nil + } + resp, err := call(r.f, r.ctx, (*sandboxfs.Client).Readlink, &sandboxfs.ReadlinkRequest{Node: n.ref}) + if err != nil { + return nil, r.serverErr(err) + } + if len(resp.Target) == 0 { + return nil, r.fail(syscall.ENOENT) + } + n.link = resp.Target + r.p.Links = append(r.p.Links, sessionview.PresentedLink{Path: n.path, Target: string(n.link)}) + return n, nil +} + +func (r *resolver) missing(cur *inode, name string) *inode { + n := r.synthetic(cur, name, sandboxfs.ModeDirectory) + r.p.Synthesized = append(r.p.Synthesized, n.path) + return n +} + +func (r *resolver) synthetic(cur *inode, name string, typ uint32) *inode { + n := r.f.newInode(sandboxfs.NodeRef{}) + n.synth, n.held, n.path = typ, true, r.at(name) + r.fixed(cur, name, n) + return n +} + +func (r *resolver) fixed(cur *inode, name string, n *inode) { + if cur.fixed == nil { + cur.fixed = map[string]*inode{} + } + cur.fixed[name] = n +} + +func (r *resolver) top() *inode { return r.stack[len(r.stack)-1].n } + +// at is the in-root path of name in the current directory. +func (r *resolver) at(name string) string { + var b strings.Builder + for _, s := range r.stack[1:] { + b.WriteString("/" + s.name) + } + return b.String() + "/" + name +} + +func (r *resolver) fail(err error) error { + return &Error{Kind: ErrMountpoint, Op: "present", Path: r.mp.Path, Err: err} +} + +// serverErr reports a sandbox errno as a presentation failure and anything else as a connection failure. +func (r *resolver) serverErr(err error) error { + var fail *sandboxfs.Failure + if errors.As(err, &fail) && fail.Code == sandboxfs.CodeErrno { + return r.fail(fmt.Errorf("%w (%w)", errnoOf(err), err)) + } + return &Error{Kind: ErrConnect, Op: "present", Path: r.mp.Path, Err: err} +} + +func isName(s string) bool { return s != "" && s != "." && s != ".." } + +func isErrno(err error, errno sandboxfs.Errno) bool { + var fail *sandboxfs.Failure + return errors.As(err, &fail) && fail.Code == sandboxfs.CodeErrno && fail.Errno == errno +} diff --git a/apps/daemon/internal/worldfs/present_linux_test.go b/apps/daemon/internal/worldfs/present_linux_test.go new file mode 100644 index 00000000..defbac22 --- /dev/null +++ b/apps/daemon/internal/worldfs/present_linux_test.go @@ -0,0 +1,360 @@ +//go:build linux + +package worldfs_test + +import ( + "context" + "errors" + "fmt" + "io" + "io/fs" + "maps" + "os" + "path/filepath" + "slices" + "strings" + "syscall" + "testing" + + "github.com/MiniMax-AI/OpenAgentCore/apps/daemon/internal/sessionview" + "github.com/MiniMax-AI/OpenAgentCore/apps/daemon/internal/worldfs" + "github.com/MiniMax-AI/OpenAgentCore/apps/sandboxio/fileservicetest" + "golang.org/x/sys/unix" +) + +// sandboxTree is a sandbox whose /bin is a symlink to usr/bin, as on merged-/usr distributions. +func sandboxTree(t *testing.T) string { + t.Helper() + backing := t.TempDir() + for _, d := range []string{"usr/bin", "etc", "data"} { + if err := os.MkdirAll(filepath.Join(backing, d), 0o755); err != nil { + t.Fatal(err) + } + } + writeFile(t, filepath.Join(backing, "usr/bin/existing"), "sandbox file") + writeFile(t, filepath.Join(backing, "usr/bin/env"), "sandbox env") + if err := os.Symlink("usr/bin", filepath.Join(backing, "bin")); err != nil { + t.Fatal(err) + } + return backing +} + +func TestPresentation(t *testing.T) { + requireFUSE(t) + backing := sandboxTree(t) + before := snapshot(t, backing) + m, err := serve(t, backing, 0, + sessionview.Mountpoint{Path: "/.oac/harness", Dir: true}, + sessionview.Mountpoint{Path: "/etc/oac-overlay", Dir: true}, + sessionview.Mountpoint{Path: "/bin/sh"}, + sessionview.Mountpoint{Path: "/usr/bin/env"}, + ) + if err != nil { + t.Fatalf("Serve: %v", err) + } + if want := []string{"/.oac/harness", "/etc/oac-overlay", "/usr/bin/sh", "/usr/bin/env"}; !slices.Equal(m.present.Targets, want) { + t.Errorf("Targets = %q, want %q", m.present.Targets, want) + } + if want := []sessionview.PresentedLink{{Path: "/bin", Target: "usr/bin"}}; !slices.Equal(m.present.Links, want) { + t.Errorf("Links = %+v, want %+v", m.present.Links, want) + } + if want := []string{"/.oac"}; !slices.Equal(m.present.Synthesized, want) { + t.Errorf("Synthesized = %q, want %q", m.present.Synthesized, want) + } + + at := func(p string) string { return filepath.Join(m.dir, p) } + wantNode(t, at(".oac"), unix.S_IFDIR|0o555) + wantNode(t, at(".oac/harness"), unix.S_IFDIR|0o555) + wantNode(t, at("etc/oac-overlay"), unix.S_IFDIR|0o555) + wantNode(t, at("usr/bin/sh"), unix.S_IFREG|0o444) + wantNode(t, at("usr/bin/env"), unix.S_IFREG|0o444) + if l, err := os.Readlink(at("bin")); err != nil || l != "usr/bin" { + t.Errorf("readlink bin = %q, %v", l, err) + } + if b, err := os.ReadFile(at("bin/existing")); err != nil || string(b) != "sandbox file" { + t.Errorf("bin/existing = %q, %v", b, err) + } + if b, err := os.ReadFile(at("bin/env")); err != nil || len(b) != 0 { + t.Errorf("presented bin/env = %q, %v; want the empty mountpoint", b, err) + } + if names := list(t, at("usr/bin")); !slices.Equal(names, []string{"env", "existing", "sh"}) { + t.Errorf("usr/bin lists %q", names) + } + if names := list(t, m.dir); !slices.Equal(names, []string{".oac", "bin", "data", "etc", "usr"}) { + t.Errorf("root lists %q", names) + } + + for _, op := range []struct { + name string + err error + }{ + {"remove a mountpoint", os.Remove(at("usr/bin/env"))}, + {"rename a pinned symlink", os.Rename(at("bin"), at("bin2"))}, + {"replace a pinned directory", unix.Rename(at("data"), at("usr"))}, + {"create in a synthetic directory", os.Mkdir(at(".oac/x"), 0o755)}, + {"write a mountpoint", os.WriteFile(at("usr/bin/sh"), nil, 0o644)}, + } { + if !errors.Is(op.err, syscall.EPERM) { + t.Errorf("%s: %v, want EPERM", op.name, op.err) + } + } + writeFile(t, at("bin/new"), "through the link") + if err := os.Remove(filepath.Join(backing, "usr/bin/new")); err != nil { + t.Errorf("file written through the presented link: %v", err) + } + if after := snapshot(t, backing); !maps.Equal(before, after) { + t.Errorf("sandbox changed:\nbefore %v\nafter %v", before, after) + } +} + +// A directory with presented children lists every name once, sends at most one ReadDir per READDIR, and keeps its offsets across a rewind. In hidden every sandbox entry is a mountpoint; paged has one mountpoint, which the sandbox lacks. +func TestPresentedPaging(t *testing.T) { + requireFUSE(t) + backing := t.TempDir() + dirs := []string{"hidden", "paged"} + var mps []sessionview.Mountpoint + want := map[string][]string{} + for _, dir := range dirs { + if err := os.Mkdir(filepath.Join(backing, dir), 0o755); err != nil { + t.Fatal(err) + } + for i := range 120 { + name := fmt.Sprintf("%s-%03d-%s", dir, i, strings.Repeat("x", 60)) + writeFile(t, filepath.Join(backing, dir, name), "") + want[dir] = append(want[dir], name) + if dir == "hidden" { + mps = append(mps, sessionview.Mountpoint{Path: "/hidden/" + name}) + } + } + mps = append(mps, sessionview.Mountpoint{Path: "/" + dir + "/absent"}) + want[dir] = append(want[dir], "absent") + slices.Sort(want[dir]) + } + m, err := serve(t, backing, 0, mps...) + if err != nil { + t.Fatalf("Serve: %v", err) + } + for _, dir := range dirs { + fd, err := unix.Open(filepath.Join(m.dir, dir), unix.O_RDONLY|unix.O_DIRECTORY, 0) + if err != nil { + t.Fatal(err) + } + defer unix.Close(fd) + var pages [][]string + var offs []int64 + for { + // Each getdents is one READDIR. + before := m.readDirs.Load() + names := getdents(t, fd) + if n := m.readDirs.Load() - before; n > 1 { + t.Errorf("%s: a READDIR sent %d ReadDir requests", dir, n) + } + if len(names) == 0 { + break + } + pages = append(pages, names) + off, err := unix.Seek(fd, 0, io.SeekCurrent) + if err != nil { + t.Fatal(err) + } + offs = append(offs, off) + } + if got := slices.Sorted(slices.Values(slices.Concat(pages...))); !slices.Equal(got, want[dir]) { + t.Errorf("%s lists %d names, want each of %d once", dir, len(got), len(want[dir])) + } + if len(pages) < 3 { + t.Fatalf("%s lists in %d pages, want at least 3", dir, len(pages)) + } + if _, err := unix.Seek(fd, 0, io.SeekStart); err != nil { + t.Fatal(err) + } + getdents(t, fd) + if _, err := unix.Seek(fd, offs[1], io.SeekStart); err != nil { + t.Fatal(err) + } + if got := getdents(t, fd); !slices.Equal(got, pages[2]) { + t.Errorf("%s: after a rewind, the offset of the second page lists %d names, want the %d of the third page", dir, len(got), len(pages[2])) + } + } +} + +func getdents(t *testing.T, fd int) []string { + t.Helper() + buf := make([]byte, 64<<10) + n, err := unix.Getdents(fd, buf) + if err != nil { + t.Fatalf("getdents: %v", err) + } + _, _, names := unix.ParseDirent(buf[:n], -1, nil) + return names +} + +func TestPresentationLoop(t *testing.T) { + requireFUSE(t) + backing := t.TempDir() + for _, l := range [][2]string{{"b", "a"}, {"a", "b"}} { + if err := os.Symlink(l[0], filepath.Join(backing, l[1])); err != nil { + t.Fatal(err) + } + } + _, err := serve(t, backing, 0, sessionview.Mountpoint{Path: "/a/x"}) + if !errors.Is(err, worldfs.ErrMountpoint) || !errors.Is(err, syscall.ELOOP) { + t.Fatalf("Serve = %v, want ErrMountpoint and ELOOP", err) + } +} + +func TestTopologyChanged(t *testing.T) { + requireFUSE(t) + backing := sandboxTree(t) + m, err := serve(t, backing, 0, sessionview.Mountpoint{Path: "/bin/sh"}) + if err != nil { + t.Fatalf("Serve: %v", err) + } + if err := os.Remove(filepath.Join(backing, "bin")); err != nil { + t.Fatal(err) + } + if err := os.Symlink("usr/local/bin", filepath.Join(backing, "bin")); err != nil { + t.Fatal(err) + } + if l, err := os.Readlink(filepath.Join(m.dir, "bin")); err != nil || l != "usr/bin" { + t.Errorf("readlink bin = %q, %v; want the pinned target", l, err) + } + waitLost(t, m.world, worldfs.ErrTopologyChanged) + wantNode(t, filepath.Join(m.dir, "usr/bin/sh"), unix.S_IFREG|0o444) +} + +func TestViewEditsWorkspace(t *testing.T) { + requireFUSE(t) + if err := sessionview.Probe(); err != nil { + t.Fatalf("Probe: %v", err) + } + backing := sandboxTree(t) + srv, err := fileservicetest.New(backing) + if err != nil { + t.Fatal(err) + } + defer srv.Close() + harness := t.TempDir() + self, err := os.Executable() + if err != nil { + t.Fatal(err) + } + copyFile(t, self, filepath.Join(harness, "harness")) + + w := worldfs.New(fileservicetest.Export, srv.Dial) + v, err := sessionview.Start(context.Background(), sessionview.Spec{ + World: w.Serve, + StagingParent: t.TempDir(), + Private: []sessionview.PrivateDir{{Name: "harness", HostDir: harness, Exec: true}}, + Shim: sessionview.Shim{Binary: filepath.Join(harness, "harness"), Names: []string{"sh"}, Paths: []string{"/bin/sh"}}, + Process: sessionview.Process{ + Path: "/.oac/harness/harness", + Args: []string{"harness"}, + Env: []string{helperEnv + "=1"}, + Dir: "/data", + UID: 1000, + GID: 1000, + Stderr: os.Stderr, + }, + }) + if err != nil { + t.Fatalf("Start: %v", err) + } + defer v.Close() + out, err := io.ReadAll(v.Stdout()) + if err != nil { + t.Fatal(err) + } + if exit, err := v.Wait(); err != nil || exit != (sessionview.Exit{}) { + t.Fatalf("Wait = %+v, %v; output %s", exit, err, out) + } + if strings.TrimSpace(string(out)) != shimMarker { + t.Errorf("/bin/sh in the view printed %q, want the shim", out) + } + if !slices.Contains(v.Presentation().Targets, "/usr/bin/sh") { + t.Errorf("Targets = %q, want the shim at /usr/bin/sh", v.Presentation().Targets) + } + if b, err := os.ReadFile(filepath.Join(backing, "data/f.txt")); err != nil || string(b) != "edited in the view" { + t.Errorf("data/f.txt = %q, %v", b, err) + } + if _, err := os.Lstat(filepath.Join(backing, "data/f.tmp")); !errors.Is(err, fs.ErrNotExist) { + t.Errorf("data/f.tmp left behind: %v", err) + } + for _, p := range []string{".oac", "proc", "dev", "usr/bin/sh"} { + if _, err := os.Lstat(filepath.Join(backing, p)); !errors.Is(err, fs.ErrNotExist) { + t.Errorf("the view created %s in the sandbox: %v", p, err) + } + } +} + +func wantNode(t *testing.T, p string, mode uint32) { + t.Helper() + var st unix.Stat_t + if err := unix.Lstat(p, &st); err != nil { + t.Errorf("lstat %s: %v", p, err) + return + } + if st.Mode != mode || st.Uid != 0 || st.Gid != 0 { + t.Errorf("%s: mode %o owner %d:%d, want %o 0:0", p, st.Mode, st.Uid, st.Gid, mode) + } +} + +func list(t *testing.T, dir string) []string { + t.Helper() + es, err := os.ReadDir(dir) + if err != nil { + t.Fatal(err) + } + var names []string + for _, e := range es { + names = append(names, e.Name()) + } + return names +} + +// snapshot describes every entry under dir by type, permissions and content. +func snapshot(t *testing.T, dir string) map[string]string { + t.Helper() + m := map[string]string{} + err := filepath.WalkDir(dir, func(p string, d fs.DirEntry, err error) error { + if err != nil { + return err + } + info, err := d.Info() + if err != nil { + return err + } + desc := info.Mode().String() + switch { + case info.Mode()&fs.ModeSymlink != 0: + l, err := os.Readlink(p) + if err != nil { + return err + } + desc += " -> " + l + case info.Mode().IsRegular(): + b, err := os.ReadFile(p) + if err != nil { + return err + } + desc += " " + string(b) + } + m[strings.TrimPrefix(p, dir)] = desc + return nil + }) + if err != nil { + t.Fatal(err) + } + return m +} + +func copyFile(t *testing.T, from, to string) { + t.Helper() + b, err := os.ReadFile(from) + if err != nil { + t.Fatal(err) + } + if err := os.WriteFile(to, b, 0o755); err != nil { + t.Fatal(err) + } +} diff --git a/apps/daemon/internal/worldfs/world_linux.go b/apps/daemon/internal/worldfs/world_linux.go new file mode 100644 index 00000000..41ee6063 --- /dev/null +++ b/apps/daemon/internal/worldfs/world_linux.go @@ -0,0 +1,316 @@ +//go:build linux + +package worldfs + +import ( + "context" + "errors" + "fmt" + "os" + "sync" + "sync/atomic" + "time" + + "github.com/MiniMax-AI/OpenAgentCore/apps/daemon/internal/sessionview" + "github.com/MiniMax-AI/OpenAgentCore/internal/sandboxfs" + "github.com/MiniMax-AI/OpenAgentCore/internal/sandboxlink" + "github.com/MiniMax-AI/OpenAgentCore/internal/sandboxwire" + "github.com/hanwen/go-fuse/v2/fuse" + "golang.org/x/sys/unix" +) + +const ( + // stopWait bounds how long Stop waits for serving to end. + stopWait = 10 * time.Second + // detachWait bounds Detach, with the redial it may need, when Stop or a failed Serve ends the attachment. + detachWait = 5 * time.Second + // retryMin and retryMax bound the backoff before the drainer sends again what it could not send. + retryMin = 20 * time.Millisecond + retryMax = 2 * time.Second + // forgetBatch is the most entries one Forget request carries. + forgetBatch = 4096 +) + +// World serves one Session's world to one view. +type World struct{ fs *frontend } + +// New returns a world that attaches export through the streams dial opens. +func New(export sandboxlink.ExportID, dial Dial) *World { + ctx, cancel := context.WithCancel(context.Background()) + drainCtx, stopDrain := context.WithCancel(ctx) + return &World{fs: &frontend{ + export: export, + dial: dial, + ctx: ctx, + cancel: cancel, + nodes: map[uint64]*inode{}, + byRef: map[sandboxfs.NodeRef]*inode{}, + handles: map[uint64]*handle{}, + forgets: map[sandboxfs.NodeRef]uint64{}, + connTurn: make(chan struct{}, 1), + kick: make(chan struct{}, 1), + drainCtx: drainCtx, + stopDrain: stopDrain, + drained: make(chan struct{}), + lost: make(chan struct{}), + served: make(chan struct{}), + }} +} + +// Serve has the signature of [sessionview.World]. Within ctx it connects, attaches the export and presents the mountpoints; it then starts serving dev. A World serves once. +// +// A Serve that fails after the export may have been attached detaches it only when nothing can still reach the attachment: when Attach may have taken effect without a response, when a request may still run on a failed stream, or when Detach fails, Serve returns [ErrAttachmentDirty] and the owner of the Link attachment must end it. +func (w *World) Serve(ctx context.Context, dev *os.File, mount sessionview.WorldMount) (sessionview.WorldServer, sessionview.Presentation, error) { + p, err := w.fs.serve(ctx, dev, mount) + if err != nil { + return nil, sessionview.Presentation{}, err + } + return w, p, nil +} + +// Stop waits up to 10 seconds for serving to end, which happens once the view's mount namespace is gone. It then detaches, ending every request still waiting on the service after 5 more seconds, so it returns within 15 seconds. It unmounts nothing. After a Serve that failed, or without one, it returns at once. +func (w *World) Stop() error { + w.fs.stopOnce.Do(func() { w.fs.stopErr = w.fs.stop() }) + return w.fs.stopErr +} + +// Lost is closed once the view no longer shows the world faithfully: the service incarnation changed, the attachment ended, or the sandbox changed a pinned entry. The view must then be rebuilt. +func (w *World) Lost() <-chan struct{} { return w.fs.lost } + +// Err reports why Lost closed, or nil while it is open. +func (w *World) Err() error { + select { + case <-w.fs.lost: + return w.fs.lostErr + default: + return nil + } +} + +// frontend is the FUSE file system. go-fuse calls it from its reader goroutines. +type frontend struct { + export sandboxlink.ExportID + dial Dial + ctx context.Context + cancel context.CancelFunc + + connTurn chan struct{} // held to use or replace conn; a channel so that waiting for it can be interrupted + conn *sandboxfs.Client + gen atomic.Uint64 // counts redials + instance sandboxwire.ID + service sandboxfs.Identity // the identity the service acts as + view sandboxfs.Identity // the identity the view's processes run as + caps sandboxfs.Capabilities + + mu sync.Mutex + nodes map[uint64]*inode + byRef map[sandboxfs.NodeRef]*inode + lastID uint64 + handles map[uint64]*handle + lastFh uint64 + root *inode + born sandboxfs.Timestamp + forgets map[sandboxfs.NodeRef]uint64 // references the kernel released, not yet sent + releases []cleanup // handles the kernel released whose Release a failed stream never sent + + kick chan struct{} // wakes the drainer + drainCtx context.Context + stopDrain context.CancelFunc + drained chan struct{} + + started atomic.Bool + serving atomic.Bool // Serve succeeded + attached bool // Attach succeeded + maybe bool // Attach failed after it may have taken effect + dead atomic.Bool // the attachment is gone: no request reaches the service again + closed atomic.Bool // Stop began: kernel requests fail + lostOnce sync.Once + lost chan struct{} + lostErr error + served chan struct{} + stopOnce sync.Once + stopErr error + + // seams let tests run other requests at the points where order matters. They are nil outside tests. + seams struct { + undo func() // after a lock request failed, before its recovery + queue func() // after a cleanup failed, before it is queued + } +} + +func (f *frontend) serve(ctx context.Context, dev *os.File, mount sessionview.WorldMount) (sessionview.Presentation, error) { + if !f.started.CompareAndSwap(false, true) { + return sessionview.Presentation{}, &Error{Kind: ErrConnect, Op: "serve", Err: errors.New("the world already serves a view")} + } + f.view = sandboxfs.Identity{UID: mount.UID, GID: mount.GID} + now := time.Now() + f.born = sandboxfs.Timestamp{Sec: now.Unix(), Nsec: uint32(now.Nanosecond())} + p, opts, err := f.attach(ctx, mount.Mountpoints) + if err == nil { + err = f.start(dev, opts) + } + if err != nil { + return sessionview.Presentation{}, f.abort(ctx, err) + } + f.serving.Store(true) + return p, nil +} + +// attach connects, checks the service's declarations, attaches the export and presents the mountpoints. +func (f *frontend) attach(ctx context.Context, mps []sessionview.Mountpoint) (sessionview.Presentation, *fuse.MountOptions, error) { + c, d, err := f.connect(ctx) + if err != nil { + return sessionview.Presentation{}, nil, &Error{Kind: ErrConnect, Op: "describe", Err: err} + } + f.conn, f.instance, f.service, f.caps = c, d.ServerInstanceID, d.Identity, d.Capabilities + maxIO := min(f.caps.MaxReadBytes, f.caps.MaxWriteBytes) &^ uint32(os.Getpagesize()-1) + switch { + case f.caps.PathProfile != sandboxfs.PathProfileLinuxBytes || f.caps.CacheProfile != sandboxfs.CacheProfileUncached: + return sessionview.Presentation{}, nil, &Error{Kind: ErrIncompatible, Op: "describe", Err: fmt.Errorf("path profile %d, cache profile %d", f.caps.PathProfile, f.caps.CacheProfile)} + case f.caps.ReadOnly: + return sessionview.Presentation{}, nil, &Error{Kind: ErrIncompatible, Op: "describe", Err: errors.New("the export is read-only")} + case maxIO == 0 || f.caps.MaxWalkComponents == 0 || f.caps.MaxReadDirBytes == 0: + return sessionview.Presentation{}, nil, &Error{Kind: ErrIncompatible, Op: "describe", Err: errors.New("read, write, walk or directory limit too small")} + } + a, err := c.Attach(ctx, &sandboxfs.AttachRequest{Export: f.export}) + if err != nil { + var fail *sandboxfs.Failure + f.maybe = !errors.As(err, &fail) || fail.Effect != sandboxwire.EffectNone + return sessionview.Presentation{}, nil, &Error{Kind: ErrConnect, Op: "attach", Path: string(f.export), Err: err} + } + f.attached = true + f.root = f.newInode(a.Root.Node) + f.root.attr, f.root.held = a.Root.Attr, true + p, err := f.present(ctx, mps) + if err != nil { + return sessionview.Presentation{}, nil, err + } + return p, &fuse.MountOptions{ + MaxWrite: int(maxIO), + EnableLocks: true, + // Only what the mapping implements: no read-ahead, page cache or open-less modes, no READDIRPLUS, passthrough or id-mapped mounts. + DisabledCapabilities: fuse.CAP_ASYNC_READ | fuse.CAP_FILE_OPS | fuse.CAP_AUTO_INVAL_DATA | fuse.CAP_READDIRPLUS | + fuse.CAP_NO_OPEN_SUPPORT | fuse.CAP_PASSTHROUGH | fuse.CAP_ALLOW_IDMAP, + ExtraCapabilities: fuse.CAP_ATOMIC_O_TRUNC, + }, nil +} + +// start serves dev on a duplicate descriptor, which go-fuse owns and closes when serving ends. +func (f *frontend) start(dev *os.File, opts *fuse.MountOptions) error { + fd, err := unix.FcntlInt(dev.Fd(), unix.F_DUPFD_CLOEXEC, 3) + if err != nil { + return &Error{Kind: ErrConnect, Op: "dup", Err: err} + } + srv, err := fuse.NewServer(f, fmt.Sprintf("/dev/fd/%d", fd), opts) + if err != nil { + return &Error{Kind: ErrConnect, Op: "init", Err: err} + } + go f.drain() + go func() { + srv.Serve() + close(f.served) + }() + return nil +} + +// abort undoes a failed Serve and returns its error. Detach on the stream the export was attached on releases what the attachment holds only when nothing else can still reach the attachment, so it is sent only then, and any doubt is [ErrAttachmentDirty]. +func (f *frontend) abort(ctx context.Context, err error) error { + close(f.served) + close(f.drained) + defer f.shutdown() + switch { + case f.maybe: + return &Error{Kind: ErrAttachmentDirty, Op: "attach", Err: err} + case !f.attached || f.dead.Load(): + return err + case f.gen.Load() != 0: + // A request on the failed stream may still be running. + return &Error{Kind: ErrAttachmentDirty, Op: "present", Err: err} + } + // Once ctx has ended, a request it abandoned may still run, and Detach under it fails, so the attachment is dirty. + ctx, cancel := context.WithTimeout(ctx, detachWait) + defer cancel() + if _, derr := f.conn.Detach(ctx, &sandboxfs.DetachRequest{}); derr != nil { + return &Error{Kind: ErrAttachmentDirty, Op: "detach", Err: errors.Join(err, derr)} + } + return err +} + +// observe marks the world lost when err shows that the service incarnation or the attachment is gone: a File failure that says so, or a Link failure that is not retryable, such as LeaseExpired or StaleGeneration on a redial. +func (f *frontend) observe(err error) { + var fail *sandboxfs.Failure + if errors.As(err, &fail) { + switch fail.Code { + case sandboxfs.CodeInstanceChanged: + f.lose(&Error{Kind: ErrInstanceChanged, Err: err}, true) + case sandboxfs.CodeStaleAttachment: + f.lose(&Error{Kind: ErrAttachmentLost, Err: err}, true) + } + return + } + code, ok := linkCode(err) + switch { + case !ok || code.Retryable(): + case code == sandboxlink.InstanceChanged: + f.lose(&Error{Kind: ErrInstanceChanged, Err: err}, true) + default: + f.lose(&Error{Kind: ErrAttachmentLost, Err: err}, true) + } +} + +func linkCode(err error) (sandboxlink.Code, bool) { + var e *sandboxlink.Error + if errors.As(err, &e) { + return e.Code, true + } + var c sandboxlink.Code + return c, errors.As(err, &c) +} + +// lose reports why the view must be rebuilt. dead also fails every later request. +func (f *frontend) lose(err *Error, dead bool) { + if dead { + f.dead.Store(true) + } + f.lostOnce.Do(func() { + f.lostErr = err + close(f.lost) + }) +} + +func (f *frontend) stop() error { + if !f.serving.Load() { + f.shutdown() + return nil + } + var errs []error + select { + case <-f.served: + case <-time.After(stopWait): + errs = append(errs, &Error{Kind: ErrConnect, Op: "stop", Err: errors.New("the view's mount still exists")}) + } + f.closed.Store(true) + ctx, cancel := context.WithTimeout(f.ctx, detachWait) + defer cancel() + // At the deadline every request still on f.ctx ends too, such as a redial that holds the stream. + context.AfterFunc(ctx, f.cancel) + // Detach drops every reference and handle the attachment holds, so nothing queued needs sending. + f.stopDrain() + <-f.drained + if !f.dead.Load() { + if _, err := call(f, ctx, (*sandboxfs.Client).Detach, &sandboxfs.DetachRequest{}); err != nil { + errs = append(errs, &Error{Kind: ErrConnect, Op: "detach", Err: err}) + } + } + f.shutdown() + return errors.Join(errs...) +} + +func (f *frontend) shutdown() { + f.cancel() + f.connTurn <- struct{}{} + defer func() { <-f.connTurn }() + if f.conn != nil { + f.conn.Close() + } +} diff --git a/apps/daemon/internal/worldfs/world_linux_test.go b/apps/daemon/internal/worldfs/world_linux_test.go new file mode 100644 index 00000000..56aa9274 --- /dev/null +++ b/apps/daemon/internal/worldfs/world_linux_test.go @@ -0,0 +1,386 @@ +//go:build linux + +package worldfs_test + +import ( + "context" + "encoding/binary" + "errors" + "fmt" + "io" + "os" + "os/exec" + "path/filepath" + "sync/atomic" + "syscall" + "testing" + "time" + + "github.com/MiniMax-AI/OpenAgentCore/apps/daemon/internal/sessionview" + "github.com/MiniMax-AI/OpenAgentCore/apps/daemon/internal/worldfs" + "github.com/MiniMax-AI/OpenAgentCore/apps/sandboxio/fileservicetest" + "github.com/MiniMax-AI/OpenAgentCore/internal/sandboxfs" + "github.com/hanwen/go-fuse/v2/posixtest" + "golang.org/x/sys/unix" +) + +// The world tests mount FUSE and need root with CAP_SYS_ADMIN (and CAP_NET_ADMIN for the view test), /dev/fuse and no AppArmor confinement. Run them in a throwaway container: +// +// CGO_ENABLED=0 go test -c -o /tmp/worldfs.test ./apps/daemon/internal/worldfs +// docker run --rm --cap-add SYS_ADMIN --cap-add NET_ADMIN --device /dev/fuse --security-opt apparmor=unconfined \ +// -e OAC_TEST_WORLDFS=1 -v /tmp/worldfs.test:/t.test:ro debian:bookworm-slim /t.test -test.v +const ( + gateEnv = "OAC_TEST_WORLDFS" + helperEnv = "OAC_WORLDFS_HELPER" + shimMarker = "oac-test-shim" +) + +// The test binary is also the launcher, the Harness and the shim inside the view. +func TestMain(m *testing.M) { + sessionview.Init() + if filepath.Base(os.Args[0]) == "sh" { + fmt.Println(shimMarker) + os.Exit(0) + } + if os.Getenv(helperEnv) != "" { + os.Exit(runHelper()) + } + os.Exit(m.Run()) +} + +// runHelper edits the workspace the way editors do, then runs the shim through the sandbox's /bin symlink. +func runHelper() int { + if err := os.WriteFile("/data/f.tmp", []byte("edited in the view"), 0o644); err != nil { + fmt.Fprintln(os.Stderr, err) + return 1 + } + if err := os.Rename("/data/f.tmp", "/data/f.txt"); err != nil { + fmt.Fprintln(os.Stderr, err) + return 1 + } + out, err := exec.Command("/bin/sh").Output() + if err != nil { + fmt.Fprintln(os.Stderr, err) + return 1 + } + os.Stdout.Write(out) + return 0 +} + +func requireFUSE(t *testing.T) { + t.Helper() + if os.Getenv(gateEnv) != "1" { + t.Skipf("set %s=1 and run the test binary as root in a privileged container; see the comment at the top of this file", gateEnv) + } +} + +type mounted struct { + dir string + world *worldfs.World + srv *fileservicetest.Server + present sessionview.Presentation + readDirs atomic.Int64 // ReadDir requests the world sent + stopped bool // the test unmounted and stopped the world itself +} + +// counted counts the ReadDir requests written to a stream, one frame per write. +type counted struct { + io.ReadWriteCloser + readDirs *atomic.Int64 +} + +func (c counted) Write(b []byte) (int, error) { + if len(b) >= 6 && binary.BigEndian.Uint16(b[4:6]) == uint16(sandboxfs.OpReadDir) { + c.readDirs.Add(1) + } + return c.ReadWriteCloser.Write(b) +} + +// serve mounts a FUSE connection the way the sessionview launcher does and serves the world over backing on it, as view identity id with the given mountpoints. +func serve(t *testing.T, backing string, id uint32, mps ...sessionview.Mountpoint) (*mounted, error) { + t.Helper() + srv, err := fileservicetest.New(backing) + if err != nil { + t.Fatal(err) + } + t.Cleanup(func() { srv.Close() }) + mnt := t.TempDir() + dev, err := os.OpenFile("/dev/fuse", os.O_RDWR, 0) + if err != nil { + t.Fatal(err) + } + opts := fmt.Sprintf("fd=%d,rootmode=40000,user_id=0,group_id=0,allow_other", dev.Fd()) + if err := unix.Mount("oac-world", mnt, "fuse", unix.MS_NOSUID|unix.MS_NODEV, opts); err != nil { + dev.Close() + t.Fatalf("mount: %v", err) + } + m := &mounted{dir: mnt, srv: srv} + m.world = worldfs.New(fileservicetest.Export, func(ctx context.Context) (io.ReadWriteCloser, error) { + rw, err := srv.Dial(ctx) + if err != nil { + return nil, err + } + return counted{rw, &m.readDirs}, nil + }) + ws, p, err := m.world.Serve(context.Background(), dev, sessionview.WorldMount{UID: id, GID: id, Mountpoints: mps}) + m.present = p + t.Cleanup(func() { + if !m.stopped { + if err := unix.Unmount(mnt, unix.MNT_DETACH); err != nil { + t.Errorf("unmount: %v", err) + } + if ws != nil { + if err := ws.Stop(); err != nil { + t.Errorf("Stop: %v", err) + } + } + } + dev.Close() + }) + return m, err +} + +// posixSkips are the posixtest cases outside phase 1, each with the reason. +var posixSkips = map[string]string{ + "Fallocate": "phase 1 excludes allocation: FALLOCATE is ENOSYS and the kernel answers EOPNOTSUPP", + "FallocateKeepSize": "phase 1 excludes allocation: FALLOCATE is ENOSYS and the kernel answers EOPNOTSUPP", + "FcntlFlockSetLk": "the file service declares no POSIX locks, so fcntl locks fail with ENOLCK instead of locking only within the view", + "FcntlFlockLocksFile": "the file service declares no POSIX locks, so fcntl locks fail with ENOLCK instead of locking only within the view", +} + +func TestPOSIX(t *testing.T) { + requireFUSE(t) + for name, fn := range posixtest.All { + t.Run(name, func(t *testing.T) { + if reason, ok := posixSkips[name]; ok { + t.Skip(reason) + } + m, err := serve(t, t.TempDir(), 0) + if err != nil { + t.Fatalf("Serve: %v", err) + } + fn(t, m.dir) + }) + } +} + +func TestOwnerMapping(t *testing.T) { + requireFUSE(t) + const view, other = 1000, 4242 + backing := t.TempDir() + writeFile(t, filepath.Join(backing, "mine"), "") + writeFile(t, filepath.Join(backing, "theirs"), "") + if err := os.Chown(filepath.Join(backing, "theirs"), other, other); err != nil { + t.Fatal(err) + } + m, err := serve(t, backing, view) + if err != nil { + t.Fatalf("Serve: %v", err) + } + wantOwner(t, filepath.Join(m.dir, "mine"), view) + wantOwner(t, filepath.Join(m.dir, "theirs"), other) + + if err := os.Chown(filepath.Join(m.dir, "theirs"), view, view); err != nil { + t.Fatalf("chown to the view identity: %v", err) + } + wantOwner(t, filepath.Join(backing, "theirs"), 0) + wantOwner(t, filepath.Join(m.dir, "theirs"), view) + if err := os.Chown(filepath.Join(m.dir, "mine"), other, other); err != nil { + t.Fatalf("chown to another identity: %v", err) + } + wantOwner(t, filepath.Join(backing, "mine"), other) + + writeFile(t, filepath.Join(m.dir, "new"), "") + wantOwner(t, filepath.Join(m.dir, "new"), view) + wantOwner(t, filepath.Join(backing, "new"), 0) +} + +func TestInstanceChanged(t *testing.T) { + requireFUSE(t) + backing := t.TempDir() + writeFile(t, filepath.Join(backing, "f"), "before") + m, err := serve(t, backing, 0) + if err != nil { + t.Fatalf("Serve: %v", err) + } + f := filepath.Join(m.dir, "f") + if _, err := os.ReadFile(f); err != nil { + t.Fatal(err) + } + // A private mapping reads the file through READ when it faults; a shared one is refused. + mf, err := os.Open(f) + if err != nil { + t.Fatal(err) + } + if b, err := unix.Mmap(int(mf.Fd()), 0, len("before"), unix.PROT_READ, unix.MAP_PRIVATE); err != nil || string(b) != "before" { + t.Errorf("private mapping = %q, %v", b, err) + } else { + unix.Munmap(b) + } + if _, err := unix.Mmap(int(mf.Fd()), 0, len("before"), unix.PROT_READ, unix.MAP_SHARED); !errors.Is(err, syscall.ENODEV) { + t.Errorf("shared mapping: %v, want ENODEV", err) + } + mf.Close() + + // A lost stream to the same incarnation redials. A request that raced the break may fail; the next one must succeed. + m.srv.Break() + if _, err := os.ReadFile(f); err != nil { + if _, err := os.ReadFile(f); err != nil { + t.Fatalf("read after a lost stream: %v", err) + } + } + select { + case <-m.world.Lost(): + t.Fatalf("world lost after a lost stream: %v", m.world.Err()) + default: + } + + if err := m.srv.Restart(); err != nil { + t.Fatal(err) + } + if _, err := os.ReadFile(f); !errors.Is(err, syscall.EIO) { + t.Fatalf("read after restart = %v, want EIO", err) + } + waitLost(t, m.world, worldfs.ErrInstanceChanged) + if _, err := os.Stat(filepath.Join(m.dir, "g")); !errors.Is(err, syscall.EIO) { + t.Errorf("lookup after restart = %v, want EIO", err) + } +} + +// Stop returns within its bound when the service has become unreachable, although Detach must redial. +func TestStopUnreachable(t *testing.T) { + requireFUSE(t) + m, err := serve(t, t.TempDir(), 0) + if err != nil { + t.Fatalf("Serve: %v", err) + } + m.stopped = true + m.srv.Stall() + if err := unix.Unmount(m.dir, unix.MNT_DETACH); err != nil { + t.Fatal(err) + } + stopped := make(chan error, 1) + go func() { stopped <- m.world.Stop() }() + select { + case err := <-stopped: + if !errors.Is(err, worldfs.ErrConnect) { + t.Errorf("Stop = %v, want ErrConnect", err) + } + case <-time.After(20 * time.Second): + t.Fatal("Stop did not return within its 15-second bound") + } +} + +// unanswered never answers Attach; it closes attaching once Attach arrives. +type unanswered struct { + sandboxfs.Service + attaching chan struct{} +} + +func (u unanswered) Attach(ctx context.Context, _ sandboxfs.Attachment, _ *sandboxfs.AttachRequest) (*sandboxfs.AttachResponse, error) { + close(u.attaching) + <-ctx.Done() + return nil, ctx.Err() +} + +// A Start that ends while Attach is unanswered keeps the world's report that the attachment must be ended. +func TestStartEndsDuringAttach(t *testing.T) { + requireFUSE(t) + srv, err := fileservicetest.New(t.TempDir()) + if err != nil { + t.Fatal(err) + } + defer srv.Close() + attaching := make(chan struct{}) + srv.Intercept(func(s sandboxfs.Service) sandboxfs.Service { return unanswered{s, attaching} }) + ctx, cancel := context.WithCancel(context.Background()) + defer cancel() + go func() { + <-attaching + cancel() + }() + _, err = sessionview.Start(ctx, sessionview.Spec{ + World: worldfs.New(fileservicetest.Export, srv.Dial).Serve, + StagingParent: t.TempDir(), + Process: sessionview.Process{Path: "/bin/true", Args: []string{"true"}, Dir: "/", UID: 1000, GID: 1000, Stderr: os.Stderr}, + }) + if !errors.Is(err, worldfs.ErrAttachmentDirty) || !errors.Is(err, context.Canceled) { + t.Fatalf("Start = %v, want ErrAttachmentDirty with the cancellation", err) + } +} + +// An interrupted flock fails with EINTR and leaves no lock on the view's handle. flock(1) waits with a timer whose signal handler does not restart the call, and it locks the descriptor the test keeps open. +func TestLockInterrupted(t *testing.T) { + requireFUSE(t) + flock, err := exec.LookPath("flock") + if err != nil { + t.Skip("needs flock(1) from util-linux") + } + backing := t.TempDir() + writeFile(t, filepath.Join(backing, "f"), "") + m, err := serve(t, backing, 0) + if err != nil { + t.Fatalf("Serve: %v", err) + } + holder, err := os.Open(filepath.Join(backing, "f")) + if err != nil { + t.Fatal(err) + } + defer holder.Close() + if err := unix.Flock(int(holder.Fd()), unix.LOCK_EX); err != nil { + t.Fatal(err) + } + f, err := os.Open(filepath.Join(m.dir, "f")) + if err != nil { + t.Fatal(err) + } + defer f.Close() + + cmd := exec.Command(flock, "--exclusive", "--timeout", "0.5", "3") + cmd.ExtraFiles = []*os.File{f} + var exit *exec.ExitError + if err := cmd.Run(); !errors.As(err, &exit) || exit.ExitCode() != 1 { + t.Fatalf("flock(1) = %v, want exit status 1: the timer interrupted the wait", err) + } + if err := unix.Flock(int(holder.Fd()), unix.LOCK_UN); err != nil { + t.Fatal(err) + } + other, err := os.Open(filepath.Join(backing, "f")) + if err != nil { + t.Fatal(err) + } + defer other.Close() + if err := unix.Flock(int(other.Fd()), unix.LOCK_EX|unix.LOCK_NB); err != nil { + t.Errorf("the view's handle kept a lock after the interrupted flock: %v", err) + } +} + +func waitLost(t *testing.T, w *worldfs.World, kind error) { + t.Helper() + select { + case <-w.Lost(): + case <-time.After(10 * time.Second): + t.Fatal("Lost never closed") + } + if !errors.Is(w.Err(), kind) { + t.Fatalf("Err = %v, want %v", w.Err(), kind) + } +} + +func wantOwner(t *testing.T, p string, id uint32) { + t.Helper() + var st unix.Stat_t + if err := unix.Lstat(p, &st); err != nil { + t.Fatal(err) + } + if st.Uid != id || st.Gid != id { + t.Errorf("%s owned by %d:%d, want %d:%d", p, st.Uid, st.Gid, id, id) + } +} + +func writeFile(t *testing.T, p, data string) { + t.Helper() + if err := os.WriteFile(p, []byte(data), 0o644); err != nil { + t.Fatal(err) + } +} diff --git a/apps/daemon/internal/worldfs/world_other.go b/apps/daemon/internal/worldfs/world_other.go new file mode 100644 index 00000000..4c0b9fb1 --- /dev/null +++ b/apps/daemon/internal/worldfs/world_other.go @@ -0,0 +1,31 @@ +//go:build !linux + +package worldfs + +import ( + "context" + "os" + + "github.com/MiniMax-AI/OpenAgentCore/apps/daemon/internal/sessionview" + "github.com/MiniMax-AI/OpenAgentCore/internal/sandboxlink" +) + +// World serves one Session's world to one view. It is Linux-only. +type World struct{} + +// New returns a world that fails to serve on this platform. +func New(sandboxlink.ExportID, Dial) *World { return &World{} } + +// Serve reports ErrUnsupported. +func (w *World) Serve(context.Context, *os.File, sessionview.WorldMount) (sessionview.WorldServer, sessionview.Presentation, error) { + return nil, sessionview.Presentation{}, &Error{Kind: ErrUnsupported, Op: "serve"} +} + +// Stop does nothing: the world never served. +func (w *World) Stop() error { return nil } + +// Lost is never closed. +func (w *World) Lost() <-chan struct{} { return nil } + +// Err is always nil. +func (w *World) Err() error { return nil } diff --git a/apps/sandboxio/fileservicetest/fileservicetest.go b/apps/sandboxio/fileservicetest/fileservicetest.go new file mode 100644 index 00000000..a8a5d0ce --- /dev/null +++ b/apps/sandboxio/fileservicetest/fileservicetest.go @@ -0,0 +1,124 @@ +//go:build linux + +// Package fileservicetest serves the Linux file service over a directory through in-memory streams, for testing File clients outside apps/sandboxio. +package fileservicetest + +import ( + "context" + "errors" + "io" + "net" + "sync" + + "github.com/MiniMax-AI/OpenAgentCore/apps/sandboxio/internal/fileservice" + "github.com/MiniMax-AI/OpenAgentCore/internal/sandboxfs" + "github.com/MiniMax-AI/OpenAgentCore/internal/sandboxlink" + "github.com/MiniMax-AI/OpenAgentCore/internal/sandboxwire" +) + +// Export is the export the service serves. +const Export = fileservice.Export + +// Server serves one attachment, whose state survives lost streams as it does behind Link. +type Server struct { + dir string + id sandboxwire.ID + + mu sync.Mutex + svc *fileservice.Service + conns []net.Conn + stalled bool + wrap func(sandboxfs.Service) sandboxfs.Service +} + +// New serves the absolute directory dir. The file service sets the process umask to zero. +func New(dir string) (*Server, error) { + svc, err := fileservice.New(dir) + if err != nil { + return nil, err + } + return &Server{dir: dir, id: sandboxwire.NewID(), svc: svc}, nil +} + +// Dial opens a stream to the current service incarnation. While the server is stalled it waits for ctx to end. +func (s *Server) Dial(ctx context.Context) (io.ReadWriteCloser, error) { + s.mu.Lock() + if s.stalled { + s.mu.Unlock() + <-ctx.Done() + return nil, ctx.Err() + } + defer s.mu.Unlock() + if s.svc == nil { + return nil, errors.New("fileservicetest: closed") + } + client, server := net.Pipe() + s.conns = append(s.conns, server) + a := sandboxfs.Attachment{ + ID: s.id, + ServerInstanceID: s.svc.InstanceID(), + Lease: context.Background(), + Exports: []sandboxlink.ExportGrant{{ID: Export}}, + } + var svc sandboxfs.Service = s.svc + if s.wrap != nil { + svc = s.wrap(svc) + } + go sandboxfs.Serve(context.Background(), server, svc, a) + return client, nil +} + +// Intercept wraps the service each later stream serves, so a test can change what it answers. +func (s *Server) Intercept(wrap func(sandboxfs.Service) sandboxfs.Service) { + s.mu.Lock() + defer s.mu.Unlock() + s.wrap = wrap +} + +// Break closes every open stream, as a lost transport does. +func (s *Server) Break() { + s.mu.Lock() + defer s.mu.Unlock() + s.breakLocked() +} + +// Stall closes every open stream and leaves the service unreachable: every later Dial waits for its context to end. +func (s *Server) Stall() { + s.mu.Lock() + defer s.mu.Unlock() + s.breakLocked() + s.stalled = true +} + +func (s *Server) breakLocked() { + for _, c := range s.conns { + c.Close() + } + s.conns = nil +} + +// Restart replaces the service with a new incarnation over the same directory and closes every open stream. +func (s *Server) Restart() error { + s.mu.Lock() + defer s.mu.Unlock() + s.breakLocked() + if s.svc != nil { + s.svc.Close() + } + svc, err := fileservice.New(s.dir) + s.svc = svc + return err +} + +// Close closes every stream and the service. +func (s *Server) Close() error { + s.mu.Lock() + defer s.mu.Unlock() + s.breakLocked() + if s.svc == nil { + return nil + } + err := s.svc.Close() + s.svc = nil + return err +}