From 4bb2093017423338995eceacdcd5de0655136912 Mon Sep 17 00:00:00 2001 From: SaladDay <1203511142@qq.com> Date: Fri, 9 Oct 2026 19:07:00 +0000 Subject: [PATCH 1/2] Pipeline workspace artifact exports --- .../internal/agenthost/environment_linux.go | 97 ++++- .../internal/agenthost/export_linux_test.go | 376 ++++++++++++++++++ 2 files changed, 462 insertions(+), 11 deletions(-) create mode 100644 apps/daemon/internal/agenthost/export_linux_test.go diff --git a/apps/daemon/internal/agenthost/environment_linux.go b/apps/daemon/internal/agenthost/environment_linux.go index 6b5fa16a8..4b2325ed2 100644 --- a/apps/daemon/internal/agenthost/environment_linux.go +++ b/apps/daemon/internal/agenthost/environment_linux.go @@ -736,6 +736,12 @@ func (x *export) walk(ctx context.Context, dir sandboxfs.NodeRef, name string, d if x.entries += len(entries); x.entries > exportEntries { return errors.New("workspace export exceeds entry bound") } + var files []exportFile + flush := func() error { + err := x.append(ctx, files) + files = nil + return err + } for _, e := range entries { child := name + "/" + string(e.Name) if !proto.ValidWorkspacePath(child) { @@ -744,9 +750,15 @@ func (x *export) walk(ctx context.Context, dir sandboxfs.NodeRef, name string, d switch e.Entry.Attr.Mode & sandboxfs.ModeType { case sandboxfs.ModeSymlink: case sandboxfs.ModeDirectory: - err = x.walk(ctx, e.Entry.Node, child, depth+1) + // Directory handles and lookup references stay with this owner. + if err = flush(); err == nil { + err = x.walk(ctx, e.Entry.Node, child, depth+1) + } case sandboxfs.ModeRegular: - err = x.append(ctx, *e.Entry, child) + files = append(files, exportFile{entry: *e.Entry, name: child}) + if len(files) == int(min(uint32(4), x.w.caps.MaxOpenHandles)) { + err = flush() + } default: err = errors.New("workspace output is not a regular file") } @@ -754,29 +766,92 @@ func (x *export) walk(ctx context.Context, dir sandboxfs.NodeRef, name string, d return err } } + return flush() +} + +type exportFile struct { + entry sandboxfs.Entry + name string +} + +// append overlaps bounded file reads while only this owner writes the archive. +// Each pipe holds at most one File Read response; no complete file is buffered. +func (x *export) append(parent context.Context, files []exportFile) (err error) { + if len(files) == 0 { + return nil + } + ctx, cancel := context.WithCancel(parent) + type pending struct { + reader *io.PipeReader + size chan int64 + err error + } + jobs := make([]pending, len(files)) + var workers sync.WaitGroup + defer func() { + cancel() + for i := range jobs { + _ = jobs[i].reader.Close() + } + workers.Wait() + for i := range jobs { + err = errors.Join(err, jobs[i].err) + } + }() + for i, file := range files { + reader, writer := io.Pipe() + jobs[i].reader, jobs[i].size = reader, make(chan int64, 1) + workers.Add(1) + go func() { + defer workers.Done() + jobs[i].err = x.read(ctx, file.entry, jobs[i].size, writer) + close(jobs[i].size) + _ = writer.CloseWithError(jobs[i].err) + }() + } + buffer := make([]byte, 32<<10) + for i, file := range files { + size, ok := <-jobs[i].size + if !ok { + // The joined worker supplies the error, before any header is written. + return errors.New("workspace output could not be opened") + } + if size > exportBatchBytes-x.bytes { + return errors.New("workspace export exceeds file bound") + } + x.bytes += size + if err = x.archive.WriteHeader(&tar.Header{Name: file.name, Typeflag: tar.TypeReg, Mode: 0o600, Size: size, Format: tar.FormatPAX}); err != nil { + return err + } + if _, err = io.CopyBuffer(x.archive, jobs[i].reader, buffer); err != nil { + return err + } + } return nil } -func (x *export) append(ctx context.Context, e sandboxfs.Entry, name string) error { +func (x *export) read(ctx context.Context, e sandboxfs.Entry, sizeReady chan<- int64, out io.Writer) (err error) { h, err := x.w.open(ctx, e) + var failure *sandboxfs.Failure + if errors.As(err, &failure) && failure.Effect == sandboxwire.EffectNone { + return err + } + // Settle even an Open whose reply was lost; no task outlives the export. + defer func() { err = errors.Join(err, x.w.closeHandle(ctx, h, false)) }() if err != nil { return err } - defer x.w.release(ctx, h) target := sandboxfs.Target{Kind: sandboxfs.TargetHandle, Handle: h} before, err := x.w.c.GetAttr(ctx, &sandboxfs.GetAttrRequest{Target: target}) if err != nil { return err } - size := int64(before.Attr.Size) - if !isType(before.Attr, sandboxfs.ModeRegular) || size > exportFileBytes || size > exportBatchBytes-x.bytes { + if !isType(before.Attr, sandboxfs.ModeRegular) || before.Attr.Size > uint64(exportFileBytes) { return errors.New("workspace export exceeds file bound") } - x.bytes += size - if err := x.archive.WriteHeader(&tar.Header{Name: name, Typeflag: tar.TypeReg, Mode: 0o600, Size: size, Format: tar.FormatPAX}); err != nil { - return err - } - n, err := x.w.read(ctx, h, size, x.archive) + size := int64(before.Attr.Size) + sizeReady <- size + n, err := x.w.read(ctx, h, size, out) if err != nil { return err } diff --git a/apps/daemon/internal/agenthost/export_linux_test.go b/apps/daemon/internal/agenthost/export_linux_test.go new file mode 100644 index 000000000..014275439 --- /dev/null +++ b/apps/daemon/internal/agenthost/export_linux_test.go @@ -0,0 +1,376 @@ +//go:build linux + +package agenthost + +import ( + "archive/tar" + "bytes" + "context" + "io" + "os" + "path/filepath" + "reflect" + "strings" + "sync/atomic" + "testing" + "time" + + "github.com/MiniMax-AI/OpenAgentCore/internal/sandboxfs" + "github.com/MiniMax-AI/OpenAgentCore/internal/sandboxwire" +) + +func exportOwner(w *world) *environment { + return &environment{id: "export", workspace: "/workspace", world: w, lost: new(atomic.Bool), sem: make(chan struct{}, 1)} +} + +func readExport(t *testing.T, data []byte) ([]string, map[string]string) { + t.Helper() + r := tar.NewReader(bytes.NewReader(data)) + var order []string + files := map[string]string{} + for { + h, err := r.Next() + if err == io.EOF { + break + } + if err != nil { + t.Fatal(err) + } + body, err := io.ReadAll(r) + if err != nil { + t.Fatal(err) + } + order = append(order, h.Name) + files[h.Name] = string(body) + } + return order, files +} + +func TestExportOutputsPreservesOrderAndBounds(t *testing.T) { + for _, limit := range []uint32{1, 2, 3, 4} { + t.Run(string(rune('0'+limit)), func(t *testing.T) { + s := &treeService{maxHandles: limit} + files := map[string]string{"workspace/outputs/a": "a", "workspace/outputs/b/one": "1", "workspace/outputs/b/two": "2", "workspace/outputs/c": "", "workspace/outputs/d": strings.Repeat("large", (1<<20)/5+100), "workspace/outputs/e": "e", "workspace/outputs/f": "f", "workspace/outputs/g": "g"} + w, dir := treeWorld(t, s, files) + if err := os.Symlink("/outside", filepath.Join(dir, "workspace/outputs/b/link")); err != nil { + t.Fatal(err) + } + var out bytes.Buffer + if err := exportOwner(w).ExportOutputs(t.Context(), &out); err != nil { + t.Fatal(err) + } + order, got := readExport(t, out.Bytes()) + want := []string{"outputs/a", "outputs/b/one", "outputs/b/two", "outputs/c", "outputs/d", "outputs/e", "outputs/f", "outputs/g"} + if !reflect.DeepEqual(order, want) { + t.Fatalf("order %v", order) + } + for name, body := range files { + if got[strings.TrimPrefix(name, "workspace/")] != body { + t.Fatalf("content %s", name) + } + } + if len(w.refs) != 0 || s.held.Load() != 0 || s.refused.Load() != 0 || s.opens.Load() != s.releases.Load() { + t.Fatalf("leaks/limit: refs=%d held=%d refused=%d opens=%d release=%d", len(w.refs), s.held.Load(), s.refused.Load(), s.opens.Load(), s.releases.Load()) + } + }) + } + for _, empty := range []bool{false, true} { + t.Run(map[bool]string{false: "missing", true: "empty"}[empty], func(t *testing.T) { + w, dir := treeWorld(t, &treeService{}, nil) + p := "workspace" + if empty { + p += "/outputs" + } + if err := os.MkdirAll(filepath.Join(dir, p), 0700); err != nil { + t.Fatal(err) + } + var out bytes.Buffer + if err := exportOwner(w).ExportOutputs(t.Context(), &out); err != nil { + t.Fatal(err) + } + order, _ := readExport(t, out.Bytes()) + if len(order) != 0 { + t.Fatal(order) + } + }) + } +} + +type exportReadService struct { + sandboxfs.Service + afterRead func() + openPossible bool +} + +func (s *exportReadService) Read(ctx context.Context, a sandboxfs.Attachment, q *sandboxfs.ReadRequest) (*sandboxfs.ReadResponse, error) { + r, err := s.Service.Read(ctx, a, q) + if s.afterRead != nil { + s.afterRead() + } + return r, err +} +func (s *exportReadService) 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.openPossible { + return nil, sandboxfs.NewErrnoFailure(sandboxfs.ErrnoIO, sandboxwire.EffectPossible, "uncertain open") + } + return r, err +} + +func TestExportOutputsRejectsChangesAndSettlesOpen(t *testing.T) { + for _, kind := range []string{"mtime", "size", "uncertain-open", "file-limit", "total-limit"} { + t.Run(kind, func(t *testing.T) { + proxy := &exportReadService{openPossible: kind == "uncertain-open"} + s := &treeService{intercept: func(real sandboxfs.Service) sandboxfs.Service { proxy.Service = real; return proxy }} + w, dir := treeWorld(t, s, map[string]string{"workspace/outputs/a": "abc"}) + p := filepath.Join(dir, "workspace/outputs/a") + if err := os.Chmod(p, 0600); err != nil { + t.Fatal(err) + } + switch kind { + case "mtime": + proxy.afterRead = func() { + if err := os.Chtimes(p, time.Unix(1, 0), time.Unix(2, 0)); err != nil { + t.Error(err) + } + } + case "size": + proxy.afterRead = func() { + if err := os.Truncate(p, 1); err != nil { + t.Error(err) + } + } + case "file-limit": + if err := os.Truncate(p, exportFileBytes+1); err != nil { + t.Fatal(err) + } + } + var out bytes.Buffer + var err error + if kind == "total-limit" { + e, eerr := w.lookup(t.Context(), w.root, "workspace") + if eerr != nil { + t.Fatal(eerr) + } + e, eerr = w.lookup(t.Context(), e.Node, "outputs") + if eerr != nil { + t.Fatal(eerr) + } + x := export{w: w, archive: tar.NewWriter(&out), bytes: exportBatchBytes - 2} + err = x.walk(t.Context(), e.Node, "outputs", 0) + } else { + err = exportOwner(w).ExportOutputs(t.Context(), &out) + } + if err == nil { + t.Fatal("invalid export succeeded") + } + if s.opens.Load() != s.releases.Load() { + select { + case <-w.c.Done(): + default: + t.Fatalf("unconfirmed cleanup left admission open: open=%d release=%d", s.opens.Load(), s.releases.Load()) + } + } + }) + } +} + +type exportFailWriter struct{} + +func (exportFailWriter) Write([]byte) (int, error) { return 0, io.ErrClosedPipe } + +func TestExportOutputsJoinsWorkersOnError(t *testing.T) { + for _, kind := range []string{"cancel", "writer", "read", "release"} { + t.Run(kind, func(t *testing.T) { + ctx, cancel := context.WithCancel(t.Context()) + defer cancel() + started := make(chan struct{}, 4) + s := &treeService{releaseError: kind == "release"} + s.read = func(ctx context.Context) error { + started <- struct{}{} + if kind == "cancel" { + <-ctx.Done() + return ctx.Err() + } + if kind == "read" { + return sandboxfs.NewErrnoFailure(sandboxfs.ErrnoIO, sandboxwire.EffectNone, "read failed") + } + return nil + } + w, _ := treeWorld(t, s, map[string]string{"workspace/outputs/a": "a", "workspace/outputs/b": "b", "workspace/outputs/c": "c", "workspace/outputs/d": "d", "workspace/outputs/e": "unadmitted"}) + var out io.Writer = io.Discard + if kind == "writer" { + out = exportFailWriter{} + } + done := make(chan error, 1) + go func() { done <- exportOwner(w).ExportOutputs(ctx, out) }() + if kind == "cancel" { + for range 4 { + select { + case <-started: + case <-time.After(5 * time.Second): + t.Fatal("read not admitted") + } + } + // Fence the client write slot before cancelling: all four requests + // have been sent, rather than cancellation racing a frame write. + if _, err := w.c.GetAttr(t.Context(), &sandboxfs.GetAttrRequest{Target: sandboxfs.Target{Kind: sandboxfs.TargetNode, Node: w.root}}); err != nil { + t.Fatal(err) + } + cancel() + } + select { + case err := <-done: + if err == nil { + t.Fatal("export succeeded") + } + case <-time.After(5 * time.Second): + t.Fatal("export did not join") + } + if s.active.Load() != 0 { + t.Fatal("read still running") + } + if s.opens.Load() > 4 { + t.Fatalf("admitted next batch: %d", s.opens.Load()) + } + if kind == "release" { + select { + case <-w.c.Done(): + default: + t.Fatal("unconfirmed Release left admission open") + } + } else if s.opens.Load() != s.releases.Load() { + select { + case <-w.c.Done(): + default: + t.Fatalf("unconfirmed cleanup left admission open: open=%d release=%d", s.opens.Load(), s.releases.Load()) + } + } + }) + } +} + +func TestExportOutputsKeepsOpenedIdentity(t *testing.T) { + s := &treeService{} + w, dir := treeWorld(t, s, map[string]string{"workspace/outputs/a": "old"}) + p := filepath.Join(dir, "workspace/outputs/a") + s.opened = func(context.Context) { + if err := os.Rename(p, p+".moved"); err != nil { + t.Error(err) + } + if err := os.WriteFile(p, []byte("replacement"), 0600); err != nil { + t.Error(err) + } + } + var out bytes.Buffer + if err := exportOwner(w).ExportOutputs(t.Context(), &out); err != nil { + t.Fatal(err) + } + _, files := readExport(t, out.Bytes()) + if files["outputs/a"] != "old" { + t.Fatal(files) + } +} + +type exportCountWriter struct { + total int64 + largest int +} + +func (w *exportCountWriter) Write(p []byte) (int, error) { + w.total += int64(len(p)) + w.largest = max(w.largest, len(p)) + return len(p), nil +} + +func TestExportOutputsStreamsMaximumFile(t *testing.T) { + s := &treeService{maxHandles: 1} + w, dir := treeWorld(t, s, map[string]string{"workspace/outputs/large": ""}) + p := filepath.Join(dir, "workspace/outputs/large") + if err := os.Chmod(p, 0600); err != nil { + t.Fatal(err) + } + if err := os.Truncate(p, exportFileBytes); err != nil { + t.Fatal(err) + } + out := new(exportCountWriter) + if err := exportOwner(w).ExportOutputs(t.Context(), out); err != nil { + t.Fatal(err) + } + if out.total != exportFileBytes+1536 || out.largest > int(w.caps.MaxReadBytes) { + t.Fatalf("unbounded/incorrect stream: %+v", out) + } + if s.opens.Load() != 1 || s.releases.Load() != 1 { + t.Fatalf("open=%d release=%d", s.opens.Load(), s.releases.Load()) + } +} + +func TestExportOutputsSettlesCancelledAcquisition(t *testing.T) { + ctx, cancel := context.WithCancel(t.Context()) + defer cancel() + entered := make(chan struct{}, 4) + s := &treeService{maxHandles: 4, opened: func(ctx context.Context) { entered <- struct{}{}; <-ctx.Done() }} + w, _ := treeWorld(t, s, map[string]string{"workspace/outputs/a": "a", "workspace/outputs/b": "b", "workspace/outputs/c": "c", "workspace/outputs/d": "d"}) + done := make(chan error, 1) + go func() { done <- exportOwner(w).ExportOutputs(ctx, io.Discard) }() + for range 4 { + select { + case <-entered: + case <-time.After(5 * time.Second): + t.Fatal("open not admitted") + } + } + if _, err := w.c.GetAttr(t.Context(), &sandboxfs.GetAttrRequest{Target: sandboxfs.Target{Kind: sandboxfs.TargetNode, Node: w.root}}); err != nil { + t.Fatal(err) + } + cancel() + select { + case err := <-done: + if err == nil { + t.Fatal("cancelled export succeeded") + } + case <-time.After(5 * time.Second): + t.Fatal("acquisition not settled") + } + if s.opens.Load() != 4 || s.releases.Load() != 4 || s.held.Load() != 0 { + t.Fatalf("open=%d released=%d held=%d", s.opens.Load(), s.releases.Load(), s.held.Load()) + } +} + +type exportWriterFunc func([]byte) (int, error) + +func (f exportWriterFunc) Write(p []byte) (int, error) { return f(p) } + +func TestExportOutputsWriterFailureCancelsBlockedReads(t *testing.T) { + entered := make(chan struct{}, 4) + s := &treeService{maxHandles: 4, read: func(ctx context.Context) error { entered <- struct{}{}; <-ctx.Done(); return ctx.Err() }} + w, _ := treeWorld(t, s, map[string]string{"workspace/outputs/a": "a", "workspace/outputs/b": "b", "workspace/outputs/c": "c", "workspace/outputs/d": "d"}) + writer := exportWriterFunc(func([]byte) (int, error) { + for range 4 { + select { + case <-entered: + case <-time.After(5 * time.Second): + return 0, io.ErrNoProgress + } + } + // Keep the parent live, with every Read frame fully sent and each service + // method blocked on cancellation when the local writer fails. + if _, err := w.c.GetAttr(t.Context(), &sandboxfs.GetAttrRequest{Target: sandboxfs.Target{Kind: sandboxfs.TargetNode, Node: w.root}}); err != nil { + return 0, err + } + return 0, io.ErrClosedPipe + }) + done := make(chan error, 1) + go func() { done <- exportOwner(w).ExportOutputs(t.Context(), writer) }() + select { + case err := <-done: + if err == nil { + t.Fatal("writer failure ignored") + } + case <-time.After(10 * time.Second): + t.Fatal("writer failure left peer reads running") + } + if s.active.Load() != 0 || s.opens.Load() != 4 || s.releases.Load() != 4 || s.held.Load() != 0 { + t.Fatalf("active=%d open=%d release=%d held=%d", s.active.Load(), s.opens.Load(), s.releases.Load(), s.held.Load()) + } +} From 9baa30b42d085e30d800c6135c7b002efb6d24d9 Mon Sep 17 00:00:00 2001 From: SaladDay <1203511142@qq.com> Date: Fri, 9 Oct 2026 19:10:38 +0000 Subject: [PATCH 2/2] Cancel artifact readers when a sibling fails --- .../internal/agenthost/environment_linux.go | 3 + .../internal/agenthost/export_linux_test.go | 59 +++++++++++++++++++ 2 files changed, 62 insertions(+) diff --git a/apps/daemon/internal/agenthost/environment_linux.go b/apps/daemon/internal/agenthost/environment_linux.go index 4b2325ed2..11b2d12dc 100644 --- a/apps/daemon/internal/agenthost/environment_linux.go +++ b/apps/daemon/internal/agenthost/environment_linux.go @@ -805,6 +805,9 @@ func (x *export) append(parent context.Context, files []exportFile) (err error) go func() { defer workers.Done() jobs[i].err = x.read(ctx, file.entry, jobs[i].size, writer) + if jobs[i].err != nil { + cancel() + } close(jobs[i].size) _ = writer.CloseWithError(jobs[i].err) }() diff --git a/apps/daemon/internal/agenthost/export_linux_test.go b/apps/daemon/internal/agenthost/export_linux_test.go index 014275439..7fc98d2a9 100644 --- a/apps/daemon/internal/agenthost/export_linux_test.go +++ b/apps/daemon/internal/agenthost/export_linux_test.go @@ -374,3 +374,62 @@ func TestExportOutputsWriterFailureCancelsBlockedReads(t *testing.T) { t.Fatalf("active=%d open=%d release=%d held=%d", s.active.Load(), s.opens.Load(), s.releases.Load(), s.held.Load()) } } + +type exportLaterErrorService struct { + sandboxfs.Service + failed chan struct{} + blocked chan struct{} +} + +func (s *exportLaterErrorService) Read(ctx context.Context, a sandboxfs.Attachment, q *sandboxfs.ReadRequest) (*sandboxfs.ReadResponse, error) { + r, err := s.Service.Read(ctx, a, q) + if err != nil { + return r, err + } + switch string(r.Data) { + case "a": + close(s.blocked) + <-ctx.Done() + return nil, ctx.Err() + case "b": + select { + case <-s.blocked: + case <-ctx.Done(): + return nil, ctx.Err() + } + close(s.failed) + return nil, sandboxfs.NewErrnoFailure(sandboxfs.ErrnoIO, sandboxwire.EffectNone, "later worker read failed") + } + return r, nil +} +func TestExportOutputsLaterFailureCancelsEarlierRead(t *testing.T) { + proxy := &exportLaterErrorService{failed: make(chan struct{}), blocked: make(chan struct{})} + s := &treeService{intercept: func(real sandboxfs.Service) sandboxfs.Service { proxy.Service = real; return proxy }} + w, _ := treeWorld(t, s, map[string]string{"workspace/outputs/a": "a", "workspace/outputs/b": "b"}) + ctx, cancel := context.WithCancel(t.Context()) + defer cancel() + done := make(chan error, 1) + go func() { done <- exportOwner(w).ExportOutputs(ctx, io.Discard) }() + select { + case <-proxy.failed: + case <-time.After(5 * time.Second): + t.Fatal("later failure not reached") + } + select { + case err := <-done: + if err == nil { + t.Fatal("export unexpectedly succeeded") + } + case <-time.After(500 * time.Millisecond): + cancel() + select { + case <-done: + case <-time.After(5 * time.Second): + t.Fatal("external cancellation did not settle export") + } + t.Fatal("later worker failure failed to cancel earlier blocked Read; export only returned after external cancellation") + } + if s.opens.Load() != 2 || s.releases.Load() != 2 { + t.Fatalf("open=%d released=%d", s.opens.Load(), s.releases.Load()) + } +}