From 0d814803c45610e8bb212b29e8f93ccbafb82f5c Mon Sep 17 00:00:00 2001 From: SaladDay <1203511142@qq.com> Date: Fri, 9 Oct 2026 11:43:42 +0000 Subject: [PATCH 1/2] Read capability files with bounded concurrency --- .../internal/agenthost/tree_linux_test.go | 304 ++++++++++++++++++ apps/daemon/internal/agenthost/world_linux.go | 61 ++-- 2 files changed, 345 insertions(+), 20 deletions(-) create mode 100644 apps/daemon/internal/agenthost/tree_linux_test.go diff --git a/apps/daemon/internal/agenthost/tree_linux_test.go b/apps/daemon/internal/agenthost/tree_linux_test.go new file mode 100644 index 000000000..a77226531 --- /dev/null +++ b/apps/daemon/internal/agenthost/tree_linux_test.go @@ -0,0 +1,304 @@ +//go:build linux + +package agenthost + +import ( + "context" + "errors" + "fmt" + "os" + "path/filepath" + "slices" + "sync/atomic" + "testing" + "time" + + "github.com/MiniMax-AI/OpenAgentCore/apps/sandboxio/fileservicetest" + "github.com/MiniMax-AI/OpenAgentCore/internal/agentbundle" + "github.com/MiniMax-AI/OpenAgentCore/internal/sandboxfs" + "github.com/MiniMax-AI/OpenAgentCore/internal/sandboxwire" +) + +type treeService struct { + sandboxfs.Service + opens, releases, active, peak atomic.Int32 + read func(context.Context) error + opened func(context.Context) + listed func(*sandboxfs.ReadDirResponse) + readDone chan struct{} + releaseError bool +} + +func (s *treeService) Open(ctx context.Context, a sandboxfs.Attachment, q *sandboxfs.OpenRequest) (*sandboxfs.OpenResponse, error) { + r, err := s.Service.Open(ctx, a, q) + if err == nil { + s.opens.Add(1) + if s.opened != nil { + s.opened(ctx) + } + } + return r, err +} +func (s *treeService) Read(ctx context.Context, a sandboxfs.Attachment, q *sandboxfs.ReadRequest) (*sandboxfs.ReadResponse, error) { + n := s.active.Add(1) + defer func() { + s.active.Add(-1) + if s.readDone != nil { + s.readDone <- struct{}{} + } + }() + for old := s.peak.Load(); n > old; old = s.peak.Load() { + if s.peak.CompareAndSwap(old, n) { + break + } + } + if s.read != nil { + if err := s.read(ctx); err != nil { + return nil, err + } + } + return s.Service.Read(ctx, a, q) +} +func (s *treeService) ReadDir(ctx context.Context, a sandboxfs.Attachment, q *sandboxfs.ReadDirRequest) (*sandboxfs.ReadDirResponse, error) { + r, err := s.Service.ReadDir(ctx, a, q) + if err == nil && s.listed != nil { + s.listed(r) + } + return r, err +} +func (s *treeService) Release(ctx context.Context, a sandboxfs.Attachment, q *sandboxfs.ReleaseRequest) (*sandboxfs.ReleaseResponse, error) { + if s.releaseError { + return nil, sandboxfs.NewErrnoFailure(sandboxfs.ErrnoIO, sandboxwire.EffectNone, "release refused") + } + r, err := s.Service.Release(ctx, a, q) + if err == nil { + s.releases.Add(1) + } + return r, err +} +func treeWorld(t *testing.T, s *treeService, files map[string]string) (*world, string) { + t.Helper() + dir := t.TempDir() + for name, body := range files { + path := filepath.Join(dir, name) + if err := os.MkdirAll(filepath.Dir(path), 0700); err != nil { + t.Fatal(err) + } + if err := os.WriteFile(path, []byte(body), 0500); err != nil { + t.Fatal(err) + } + } + server, err := fileservicetest.New(dir) + if err != nil { + t.Fatal(err) + } + t.Cleanup(func() { server.Close() }) + server.Intercept(func(real sandboxfs.Service) sandboxfs.Service { s.Service = real; return s }) + stream, err := server.Dial(t.Context()) + if err != nil { + t.Fatal(err) + } + uncertain := false + w, err := attachWorld(t.Context(), stream, &uncertain) + if err != nil { + t.Fatal(err) + } + t.Cleanup(func() { w.c.Close() }) + return w, dir +} +func treeReleased(t *testing.T, w *world, s *treeService) { + t.Helper() + if s.active.Load() != 0 || s.opens.Load() != s.releases.Load() { + t.Fatalf("unfinished reads: active=%d open=%d release=%d", s.active.Load(), s.opens.Load(), s.releases.Load()) + } + nodes := make([]sandboxfs.NodeRef, 0, len(w.refs)) + for node := range w.refs { + nodes = append(nodes, node) + } + if err := w.forget(t.Context()); err != nil { + t.Fatal(err) + } + for _, node := range nodes { + _, err := w.c.GetAttr(t.Context(), &sandboxfs.GetAttrRequest{Target: sandboxfs.Target{Kind: sandboxfs.TargetNode, Node: node}}) + var f *sandboxfs.Failure + if !errors.As(err, &f) || f.Code != sandboxfs.CodeStaleNode { + t.Fatalf("forgotten node remained live: %v", err) + } + } +} + +func TestReadTreeConcurrent(t *testing.T) { + for _, mode := range []string{"success", "read_error", "parent_cancel", "open_cancel"} { + t.Run(mode, func(t *testing.T) { + s := &treeService{readDone: make(chan struct{}, 8)} + files := map[string]string{} + for i := range 8 { + files[fmt.Sprintf("dir/%02d", i)] = fmt.Sprint(i) + } + w, _ := treeWorld(t, s, files) + ctx, cancel := context.WithCancel(t.Context()) + defer cancel() + entered := make(chan struct{}, 8) + gate := make(chan struct{}) + gateClosed := false + defer func() { + if !gateClosed { + close(gate) + } + }() + if mode == "open_cancel" { + s.opened = func(ctx context.Context) { entered <- struct{}{}; <-ctx.Done() } + } else { + var first atomic.Bool + s.read = func(ctx context.Context) error { + entered <- struct{}{} + select { + case <-ctx.Done(): + return ctx.Err() + case <-gate: + } + if mode == "read_error" && first.CompareAndSwap(false, true) { + return sandboxfs.NewErrnoFailure(sandboxfs.ErrnoIO, sandboxwire.EffectNone, "read refused") + } + return nil + } + } + type result struct { + files []agentbundle.File + err error + } + done := make(chan result, 1) + go func() { f, e := w.readTree(ctx, w.root, true); done <- result{f, e} }() + for range 4 { + select { + case <-entered: + case <-time.After(5 * time.Second): + t.Fatal("four reads did not start") + } + } + if s.opens.Load() != 4 { + t.Fatalf("opened %d files before releasing bound", s.opens.Load()) + } + if mode == "parent_cancel" || mode == "open_cancel" { + // A completed call fences the preceding request writes: this + // case exercises cancellation after dispatch, not a torn frame. + if _, err := w.c.Describe(t.Context(), &sandboxfs.DescribeRequest{}); err != nil { + t.Fatal(err) + } + cancel() + } else { + close(gate) + gateClosed = true + } + var r result + select { + case r = <-done: + case <-time.After(5 * time.Second): + t.Fatal("read did not join") + } + if mode == "success" { + if r.err != nil || len(r.files) != len(files) { + t.Fatalf("tree result: %v %v", r.files, r.err) + } + names := make([]string, 0, len(files)) + for name := range files { + names = append(names, name) + } + slices.Sort(names) + for i, f := range r.files { + if f.Path != names[i] || string(f.Data) != files[f.Path] || !f.Executable { + t.Fatalf("incorrect ordered file: %+v", f) + } + } + } else if r.err == nil || r.files != nil { + t.Fatalf("failure returned tree: %v %v", r.files, r.err) + } + if mode != "open_cancel" && s.peak.Load() != 4 { + t.Fatalf("read concurrency %d, want 4", s.peak.Load()) + } + if mode != "open_cancel" { + n := 8 + if mode == "parent_cancel" { + n = 4 + } + for range n { + select { + case <-s.readDone: + case <-time.After(5 * time.Second): + t.Fatal("File read handler did not finish") + } + } + } + treeReleased(t, w, s) + }) + } +} + +func TestReadTreeValidation(t *testing.T) { + for _, mode := range []string{"writable", "symlink", "size_sum", "entry_count", "grow", "shrink", "readback_tamper", "release_error"} { + t.Run(mode, func(t *testing.T) { + s := &treeService{} + w, dir := treeWorld(t, s, map[string]string{"a": "original", "b": "second"}) + switch mode { + case "writable": + if err := os.Chmod(filepath.Join(dir, "a"), 0600); err != nil { + t.Fatal(err) + } + case "symlink": + if err := os.Symlink("a", filepath.Join(dir, "c")); err != nil { + t.Fatal(err) + } + case "size_sum": + s.listed = func(r *sandboxfs.ReadDirResponse) { + for i := range r.Entries { + r.Entries[i].Entry.Attr.Size = uint64(agentbundle.MaxExpandedBytes) + } + } + case "entry_count": + for i := range agentbundle.MaxFiles { + if err := os.WriteFile(filepath.Join(dir, fmt.Sprint(i)), nil, 0400); err != nil { + t.Fatal(err) + } + } + case "grow", "shrink": + s.listed = func(r *sandboxfs.ReadDirResponse) { + for i := range r.Entries { + if string(r.Entries[i].Name) == "a" { + if mode == "grow" { + r.Entries[i].Entry.Attr.Size-- + } else { + r.Entries[i].Entry.Attr.Size++ + } + } + } + } + case "readback_tamper": + if _, err := w.readTree(t.Context(), w.root, true); err != nil { + t.Fatal(err) + } + if err := os.Chmod(filepath.Join(dir, "a"), 0600); err != nil { + t.Fatal(err) + } + case "release_error": + s.releaseError = true + } + files, err := w.readTree(t.Context(), w.root, true) + if err == nil || files != nil { + t.Fatalf("invalid tree accepted: %v %v", files, err) + } + if mode == "release_error" { + if !w.ended() { + t.Fatal("failed release left attachment reusable") + } + if _, err := w.c.Describe(t.Context(), &sandboxfs.DescribeRequest{}); !errors.Is(err, sandboxfs.ErrTransport) { + t.Fatalf("failed attachment admitted a call: %v", err) + } + return + } + if slices.Contains([]string{"writable", "symlink", "size_sum", "entry_count"}, mode) && s.opens.Load() != 0 { + t.Fatalf("opened files before validating metadata: %d", s.opens.Load()) + } + treeReleased(t, w, s) + }) + } +} diff --git a/apps/daemon/internal/agenthost/world_linux.go b/apps/daemon/internal/agenthost/world_linux.go index 56d48a87d..3ee055a18 100644 --- a/apps/daemon/internal/agenthost/world_linux.go +++ b/apps/daemon/internal/agenthost/world_linux.go @@ -17,6 +17,7 @@ import ( "github.com/MiniMax-AI/OpenAgentCore/internal/sandboxfs" "github.com/MiniMax-AI/OpenAgentCore/internal/sandboxwire" "github.com/google/uuid" + "golang.org/x/sync/errgroup" ) var ( @@ -184,14 +185,14 @@ func (w *world) mkdir(ctx context.Context, dir sandboxfs.NodeRef, name string) ( return w.track(r.Entry), nil } -// open opens the regular file e for reading and returns its handle. +// open returns the regular file's attempted handle even when its outcome is unknown. func (w *world) open(ctx context.Context, e sandboxfs.Entry) (sandboxfs.HandleID, error) { if !isType(e.Attr, sandboxfs.ModeRegular) { return 0, fs.ErrInvalid } h := w.handles.Next() if _, err := w.c.Open(ctx, &sandboxfs.OpenRequest{Handle: h, Node: e.Node, Access: sandboxfs.AccessRead, Flags: sandboxfs.OpenNoFollow}); err != nil { - return 0, err + return h, err } return h, nil } @@ -237,16 +238,28 @@ func (w *world) readEntry(ctx context.Context, e sandboxfs.Entry, limit int64) ( return nil, fs.ErrInvalid } h, err := w.open(ctx, e) - if err != nil { + var failure *sandboxfs.Failure + if errors.As(err, &failure) && failure.Effect == sandboxwire.EffectNone { return nil, err } - defer w.release(ctx, h) b := bytes.NewBuffer([]byte{}) - n, err := w.read(ctx, h, limit, b) - if err == nil && uint64(n) != e.Attr.Size { - err = fs.ErrInvalid + if err == nil { + var n int64 + n, err = w.read(ctx, h, limit, b) + if err == nil && uint64(n) != e.Attr.Size { + err = fs.ErrInvalid + } } - return b.Bytes(), err + // Release also joins an Open whose reply was lost to cancellation. The + // server orders release after acquisition of this handle. + cleanup, cancel := context.WithTimeout(context.WithoutCancel(ctx), closeBound) + defer cancel() + _, releaseErr := w.c.Release(cleanup, &sandboxfs.ReleaseRequest{Handle: h}) + if releaseErr != nil { + // The next owner operation drains this attachment before reusing it. + w.c.Close() + } + return b.Bytes(), errors.Join(err, releaseErr) } // create makes name in dir, exclusively, with mode, writes data to it and @@ -396,6 +409,22 @@ func (w *world) readTree(ctx context.Context, dir sandboxfs.NodeRef, immutable b if err := t.walk(ctx, dir, ""); err != nil { return nil, err } + // Enumeration owns node references; only reads of the retained entries + // run concurrently. Join every read before the owner can forget them. + var reads errgroup.Group + reads.SetLimit(4) + for i, entry := range t.nodes { + reads.Go(func() error { + body, err := w.readEntry(ctx, entry, int64(entry.Attr.Size)) + if err == nil { + t.files[i].Data = body + } + return err + }) + } + if err := reads.Wait(); err != nil { + return nil, err + } return t.files, nil } @@ -403,6 +432,7 @@ type treeReader struct { w *world immutable bool files []agentbundle.File + nodes []sandboxfs.Entry entries int total int } @@ -425,18 +455,9 @@ func (t *treeReader) walk(ctx context.Context, dir sandboxfs.NodeRef, prefix str case !isType(attr, sandboxfs.ModeRegular) || t.immutable && attr.Mode&0o222 != 0 || attr.Size > uint64(agentbundle.MaxExpandedBytes-t.total): return fs.ErrInvalid default: - h, err := t.w.open(ctx, *e.Entry) - if err != nil { - return err - } - var b strings.Builder - n, err := t.w.read(ctx, h, int64(agentbundle.MaxExpandedBytes-t.total), &b) - t.w.release(ctx, h) - if err != nil || uint64(n) != attr.Size { - return fs.ErrInvalid - } - t.total += int(n) - t.files = append(t.files, agentbundle.File{Path: name, Data: []byte(b.String()), Executable: attr.Mode&0o111 != 0}) + t.total += int(attr.Size) + t.nodes = append(t.nodes, *e.Entry) + t.files = append(t.files, agentbundle.File{Path: name, Executable: attr.Mode&0o111 != 0}) } } return nil From f65e3e8744669b774b6051dc602105209898f3ca Mon Sep 17 00:00:00 2001 From: SaladDay <1203511142@qq.com> Date: Fri, 9 Oct 2026 11:48:19 +0000 Subject: [PATCH 2/2] Respect declared File handle limits in capability reads --- .../internal/agenthost/tree_linux_test.go | 48 +++++++++++++++---- apps/daemon/internal/agenthost/world_linux.go | 2 +- 2 files changed, 40 insertions(+), 10 deletions(-) diff --git a/apps/daemon/internal/agenthost/tree_linux_test.go b/apps/daemon/internal/agenthost/tree_linux_test.go index a77226531..7f3e8bdff 100644 --- a/apps/daemon/internal/agenthost/tree_linux_test.go +++ b/apps/daemon/internal/agenthost/tree_linux_test.go @@ -27,10 +27,27 @@ type treeService struct { listed func(*sandboxfs.ReadDirResponse) readDone chan struct{} releaseError bool + maxHandles uint32 + held, refused atomic.Int32 } +func (s *treeService) Describe(ctx context.Context, a sandboxfs.Attachment, q *sandboxfs.DescribeRequest) (*sandboxfs.DescribeResponse, error) { + r, err := s.Service.Describe(ctx, a, q) + if err == nil && s.maxHandles != 0 { + r.Capabilities.MaxOpenHandles = s.maxHandles + } + return r, err +} func (s *treeService) Open(ctx context.Context, a sandboxfs.Attachment, q *sandboxfs.OpenRequest) (*sandboxfs.OpenResponse, error) { + if s.maxHandles != 0 && s.held.Add(1) > int32(s.maxHandles) { + s.held.Add(-1) + s.refused.Add(1) + return nil, sandboxfs.NewFailure(sandboxfs.CodeResourceExhausted, sandboxwire.EffectNone, "declared handle limit reached") + } r, err := s.Service.Open(ctx, a, q) + if err != nil && s.maxHandles != 0 { + s.held.Add(-1) + } if err == nil { s.opens.Add(1) if s.opened != nil { @@ -73,6 +90,9 @@ func (s *treeService) Release(ctx context.Context, a sandboxfs.Attachment, q *sa r, err := s.Service.Release(ctx, a, q) if err == nil { s.releases.Add(1) + if s.maxHandles != 0 { + s.held.Add(-1) + } } return r, err } @@ -128,9 +148,16 @@ func treeReleased(t *testing.T, w *world, s *treeService) { } func TestReadTreeConcurrent(t *testing.T) { - for _, mode := range []string{"success", "read_error", "parent_cancel", "open_cancel"} { - t.Run(mode, func(t *testing.T) { - s := &treeService{readDone: make(chan struct{}, 8)} + for _, tc := range []struct { + mode string + limit int + }{ + {"success", 4}, {"read_error", 4}, {"parent_cancel", 4}, {"open_cancel", 4}, + {"success", 1}, {"success", 2}, {"success", 3}, + } { + mode, limit := tc.mode, tc.limit + t.Run(fmt.Sprintf("%s/handles_%d", mode, limit), func(t *testing.T) { + s := &treeService{readDone: make(chan struct{}, 8), maxHandles: uint32(limit)} files := map[string]string{} for i := range 8 { files[fmt.Sprintf("dir/%02d", i)] = fmt.Sprint(i) @@ -169,14 +196,14 @@ func TestReadTreeConcurrent(t *testing.T) { } done := make(chan result, 1) go func() { f, e := w.readTree(ctx, w.root, true); done <- result{f, e} }() - for range 4 { + for range limit { select { case <-entered: case <-time.After(5 * time.Second): - t.Fatal("four reads did not start") + t.Fatal("declared number of reads did not start") } } - if s.opens.Load() != 4 { + if s.opens.Load() != int32(limit) { t.Fatalf("opened %d files before releasing bound", s.opens.Load()) } if mode == "parent_cancel" || mode == "open_cancel" { @@ -213,13 +240,13 @@ func TestReadTreeConcurrent(t *testing.T) { } else if r.err == nil || r.files != nil { t.Fatalf("failure returned tree: %v %v", r.files, r.err) } - if mode != "open_cancel" && s.peak.Load() != 4 { - t.Fatalf("read concurrency %d, want 4", s.peak.Load()) + if mode != "open_cancel" && s.peak.Load() != int32(limit) { + t.Fatalf("read concurrency %d, want %d", s.peak.Load(), limit) } if mode != "open_cancel" { n := 8 if mode == "parent_cancel" { - n = 4 + n = limit } for range n { select { @@ -229,6 +256,9 @@ func TestReadTreeConcurrent(t *testing.T) { } } } + if s.held.Load() != 0 || s.refused.Load() != 0 { + t.Fatalf("handle bound exceeded or leaked: held=%d refused=%d", s.held.Load(), s.refused.Load()) + } treeReleased(t, w, s) }) } diff --git a/apps/daemon/internal/agenthost/world_linux.go b/apps/daemon/internal/agenthost/world_linux.go index 3ee055a18..ee24f60b1 100644 --- a/apps/daemon/internal/agenthost/world_linux.go +++ b/apps/daemon/internal/agenthost/world_linux.go @@ -412,7 +412,7 @@ func (w *world) readTree(ctx context.Context, dir sandboxfs.NodeRef, immutable b // Enumeration owns node references; only reads of the retained entries // run concurrently. Join every read before the owner can forget them. var reads errgroup.Group - reads.SetLimit(4) + reads.SetLimit(int(min(uint32(4), w.caps.MaxOpenHandles))) for i, entry := range t.nodes { reads.Go(func() error { body, err := w.readEntry(ctx, entry, int64(entry.Attr.Size))