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
37 changes: 17 additions & 20 deletions packages/orchestrator/pkg/sandbox/build_upload.go
Original file line number Diff line number Diff line change
Expand Up @@ -27,6 +27,7 @@ type Upload struct {
root storage.CompressConfig
objectMetadata storage.ObjectMetadata
future *utils.ErrorOnce
useV4 bool
}

func NewUpload(
Expand All @@ -39,11 +40,11 @@ func NewUpload(
useCase string,
objectMetadata storage.ObjectMetadata,
) (*Upload, error) {
mem, err := resolveCompressConfig(ctx, cfg, ff, storage.MemfileName, snap.MemfileDiffHeader.Metadata.BlockSize, useCase)
mem, memV4, err := resolveCompressConfig(ctx, cfg, ff, storage.MemfileName, snap.MemfileDiffHeader.Metadata.BlockSize, useCase)
if err != nil {
return nil, fmt.Errorf("resolve memfile compress config: %w", err)
}
root, err := resolveCompressConfig(ctx, cfg, ff, storage.RootfsName, snap.RootfsDiffHeader.Metadata.BlockSize, useCase)
root, rootV4, err := resolveCompressConfig(ctx, cfg, ff, storage.RootfsName, snap.RootfsDiffHeader.Metadata.BlockSize, useCase)
if err != nil {
return nil, fmt.Errorf("resolve rootfs compress config: %w", err)
}
Comment thread
levb marked this conversation as resolved.
Expand All @@ -57,6 +58,7 @@ func NewUpload(
mem: mem,
root: root,
objectMetadata: objectMetadata,
useV4: memV4 || rootV4,
}

if uploads != nil {
Expand All @@ -71,7 +73,7 @@ func NewUpload(
}

func (u *Upload) Run(ctx context.Context) error {
if !u.mem.IsCompressionEnabled() && !u.root.IsCompressionEnabled() {
if !u.mem.IsCompressionEnabled() && !u.root.IsCompressionEnabled() && !u.useV4 {
return u.runV3(ctx)
Comment thread
levb marked this conversation as resolved.
}

Expand Down Expand Up @@ -111,20 +113,14 @@ func (u *Upload) publish(ctx context.Context, t build.DiffType, h *headers.Heade
}

// resolveCompressConfig returns the effective compression config for a given
// file type and use case. Feature flags override the base config when active.
// Returns zero-value CompressConfig when compression is disabled.
//
// fileType and useCase are added to the LD evaluation context so that
// LaunchDarkly targeting rules can differentiate (e.g. compress memfile
// but not rootfs, or compress builds but not pauses). blockSize is the
// in-VM read granularity for this fileType (from the diff header) and
// constrains the legal frame sizes — see validateCompressConfig.
//
// The resolved config is validated; an invalid env or LD-derived config
// surfaces as an error so the upload fails fast rather than streaming with
// a misconfigured frame size.
func resolveCompressConfig(ctx context.Context, base storage.CompressConfig, ff *featureflags.Client, fileType string, blockSize uint64, useCase string) (storage.CompressConfig, error) {
// file type and use case, plus whether the V4 header layout should be used for
// an uncompressed upload. Feature flags override the base config when active.
// Returns zero-value CompressConfig when compression is disabled. fileType,
// useCase are added to the LD evaluation context; blockSize constrains legal
// frame sizes — see validateCompressConfig.
func resolveCompressConfig(ctx context.Context, base storage.CompressConfig, ff *featureflags.Client, fileType string, blockSize uint64, useCase string) (storage.CompressConfig, bool, error) {
resolved := base
var useV4 bool

if ff != nil {
var extra []ldcontext.Context
Expand All @@ -136,8 +132,9 @@ func resolveCompressConfig(ctx context.Context, base storage.CompressConfig, ff
}
ctx = featureflags.AddToContext(ctx, extra...)

v := ff.JSONFlag(ctx, featureflags.CompressConfigFlag).AsValueMap()
useV4 = ff.BoolFlag(ctx, featureflags.V4HeaderForUncompressedFlag)

v := ff.JSONFlag(ctx, featureflags.CompressConfigFlag).AsValueMap()
if v.Get("compressBuilds").BoolValue() {
ct := v.Get("compressionType").StringValue()
ldCfg := storage.CompressConfig{
Expand All @@ -156,14 +153,14 @@ func resolveCompressConfig(ctx context.Context, base storage.CompressConfig, ff
}

if !resolved.IsCompressionEnabled() {
return storage.CompressConfig{}, nil
return storage.CompressConfig{}, useV4, nil
}

if err := validateCompressConfig(resolved, blockSize); err != nil {
return storage.CompressConfig{}, err
return storage.CompressConfig{}, false, err
}

return resolved, nil
return resolved, useV4, nil
}

// validateCompressConfig checks that the resolved config is internally
Expand Down
57 changes: 57 additions & 0 deletions packages/orchestrator/pkg/sandbox/build_upload_test.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,57 @@
//go:build linux

package sandbox

import (
"testing"

"github.com/launchdarkly/go-server-sdk/v7/testhelpers/ldtestdata"
"github.com/stretchr/testify/require"

"github.com/e2b-dev/infra/packages/shared/pkg/featureflags"
"github.com/e2b-dev/infra/packages/shared/pkg/storage"
)

func newV4HeaderFF(t *testing.T, on bool) *featureflags.Client {
t.Helper()

td := ldtestdata.DataSource()
td.Update(td.Flag(featureflags.V4HeaderForUncompressedFlag.Key()).VariationForAll(on))

ff, err := featureflags.NewClientWithDatasource(td)
require.NoError(t, err)

t.Cleanup(func() {
_ = ff.Close(t.Context())
})

return ff
}

func resolveV4(t *testing.T, ff *featureflags.Client) bool {
t.Helper()
_, useV4, err := resolveCompressConfig(t.Context(), storage.CompressConfig{}, ff, storage.MemfileName, 4096, storage.UseCaseBuild)
require.NoError(t, err)

return useV4
}

func TestResolveCompressConfig_V4_NilClient(t *testing.T) {
t.Parallel()

require.False(t, resolveV4(t, nil))
}

func TestResolveCompressConfig_V4_FlagOff(t *testing.T) {
t.Parallel()

ff := newV4HeaderFF(t, false)
require.False(t, resolveV4(t, ff))
}

func TestResolveCompressConfig_V4_FlagOn(t *testing.T) {
t.Parallel()

ff := newV4HeaderFF(t, true)
require.True(t, resolveV4(t, ff))
}
19 changes: 13 additions & 6 deletions packages/orchestrator/pkg/sandbox/build_upload_v4.go
Original file line number Diff line number Diff line change
Expand Up @@ -5,6 +5,7 @@ package sandbox
import (
"context"
"fmt"
"os"

"github.com/google/uuid"
"golang.org/x/sync/errgroup"
Expand Down Expand Up @@ -62,17 +63,23 @@ func (u *Upload) uploadFramed(
var selfBuild headers.BuildData

if srcPath != "" {
ft, checksum, err := storage.UploadFramed(ctx, u.store, u.paths.DataFile(string(fileType), cfg.CompressionType()), seekableTypeFor(fileType), srcPath, storage.WithCompressConfig(cfg), storage.WithMetadata(u.objectMetadata))
ft, checksum, err := storage.UploadFramed(ctx, u.store, u.paths.DataFile(string(fileType), cfg.CompressionType()), seekableTypeFor(fileType), srcPath, storage.WithCompressConfig(cfg), storage.WithMetadata(u.objectMetadata), storage.WithChecksumSHA256())
if err != nil {
return fmt.Errorf("%s upload: %w", fileType, err)
}

// FrameTable count, not os.Stat: sparse memfile diffs stream less than
// they appear on disk.
selfBuild = headers.BuildData{Size: ft.UncompressedSize(), Checksum: checksum}
if ft.IsCompressed() {
selfBuild.FrameData = ft
// Compressed: frame-table byte count, since sparse memfile diffs stream
// fewer bytes than they occupy on disk. Uncompressed has no table.
size := ft.UncompressedSize()
if !ft.IsCompressed() {
Comment thread
levb marked this conversation as resolved.
Comment thread
levb marked this conversation as resolved.
info, statErr := os.Stat(srcPath)
if statErr != nil {
return fmt.Errorf("%s stat: %w", fileType, statErr)
}
size = info.Size()
}

selfBuild = headers.BuildData{Size: size, Checksum: checksum, FrameData: ft}
}

h := srcHeader.CloneForUpload(headers.MetadataVersionV4)
Expand Down
8 changes: 7 additions & 1 deletion packages/shared/pkg/featureflags/flags.go
Original file line number Diff line number Diff line change
Expand Up @@ -126,6 +126,11 @@ var (
FreePageReportingFlag = NewBoolFlag("free-page-reporting", false)

NetworkTransformRulesFlag = NewBoolFlag("network-transform-rules", env.IsDevelopment())

// V4HeaderForUncompressedFlag forces the V4 header layout on uncompressed
// uploads. Independent of compress-config: it changes the header format,
// not whether data is compressed.
V4HeaderForUncompressedFlag = NewBoolFlag("v4-header-for-uncompressed", false)
)

type IntFlag struct {
Expand Down Expand Up @@ -383,7 +388,8 @@ func GetTrackedTemplatesSet(ctx context.Context, ff *Client) map[string]struct{}

// CompressConfigFlag controls compression during template builds.
// When compressBuilds is true, builds upload exclusively compressed data
// (no uncompressed fallback). When false, exclusively uncompressed with V3 headers.
// (no uncompressed fallback). When false, exclusively uncompressed with V3
// headers (unless V4HeaderForUncompressedFlag is set).
var CompressConfigFlag = NewJSONFlag("compress-config", ldvalue.FromJSONMarshal(map[string]any{
"compressBuilds": false,
"compressionType": "",
Expand Down
4 changes: 1 addition & 3 deletions packages/shared/pkg/storage/compress_upload.go
Original file line number Diff line number Diff line change
Expand Up @@ -167,11 +167,9 @@ func compressStream(ctx context.Context, in io.Reader, cfg CompressConfig, uploa
return nil, [32]byte{}, fmt.Errorf("complete upload: %w", err)
}

var checksum [32]byte
copy(checksum[:], hasher.Sum(nil))
ft := NewFrameTable(cfg.CompressionType(), frameSizes)

return ft, checksum, nil
return ft, sum256(hasher), nil
}

func readLoop(ctx context.Context, in io.Reader, cfg CompressConfig, hasher io.Writer, q chan<- *part) error {
Expand Down
34 changes: 32 additions & 2 deletions packages/shared/pkg/storage/gcp_multipart.go
Original file line number Diff line number Diff line change
Expand Up @@ -7,6 +7,7 @@ import (
"encoding/base64"
"encoding/xml"
"fmt"
"hash"
"io"
"math"
"math/rand"
Expand Down Expand Up @@ -382,7 +383,7 @@ func (m *MultipartUploader) completeUpload(ctx context.Context, uploadID string,
return nil
}

func (m *MultipartUploader) UploadFileInParallel(ctx context.Context, filePath string, maxConcurrency int) (int64, error) {
func (m *MultipartUploader) UploadFileInParallel(ctx context.Context, filePath string, maxConcurrency int, hasher hash.Hash) (int64, error) {
// Open file
file, err := os.Open(filePath)
if err != nil {
Expand All @@ -409,9 +410,38 @@ func (m *MultipartUploader) UploadFileInParallel(ctx context.Context, filePath s
return 0, fmt.Errorf("failed to initiate upload: %w", err)
}

// Hash on a sibling goroutine while parts upload — the read overlaps the
// upload, adding no wall-clock latency. Own file handle (separate from the
// part uploaders' ReadAt); opened here so a failed upload can close it.
var hashFile *os.File
if hasher != nil {
hashFile, err = os.Open(filePath)
if err != nil {
return 0, fmt.Errorf("failed to open file for checksum: %w", err)
}
defer hashFile.Close()
}

var eg errgroup.Group
if hashFile != nil {
eg.Go(func() error {
if _, err := io.Copy(hasher, hashFile); err != nil {
return fmt.Errorf("failed to checksum file: %w", err)
}

return nil
})
}

parts, err := m.uploadParts(ctx, maxConcurrency, numParts, fileSize, file, uploadID)
if hashFile != nil && err != nil {
hashFile.Close() // cancel the now-pointless io.Copy
}
if hashErr := eg.Wait(); err == nil {
err = hashErr
}
if err != nil {
return 0, fmt.Errorf("failed to upload parts: %w", err)
return 0, fmt.Errorf("failed to upload file: %w", err)
}

if err := m.completeUpload(ctx, uploadID, parts); err != nil {
Expand Down
Loading
Loading