Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
100 changes: 89 additions & 11 deletions apps/daemon/internal/agenthost/environment_linux.go
Original file line number Diff line number Diff line change
Expand Up @@ -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) {
Expand All @@ -744,39 +750,111 @@ 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")
}
if err != nil {
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)
if jobs[i].err != nil {
cancel()
}
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
}
Expand Down
Loading
Loading