Skip to content
Merged
Show file tree
Hide file tree
Changes from 2 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
2 changes: 1 addition & 1 deletion cmd/hauler/cli/cli.go
Original file line number Diff line number Diff line change
Expand Up @@ -18,7 +18,7 @@ func New(ctx context.Context, ro *flags.CliRootOpts) *cobra.Command {
cmd := &cobra.Command{
Use: "hauler",
Short: "Airgap Swiss Army Knife",
Example: " View the Docs: https://docs.hauler.dev\n Environment Variables: " + consts.HaulerDir + " | " + consts.HaulerTempDir + " | " + consts.HaulerStoreDir + " | " + consts.HaulerIgnoreErrors + " | " + consts.HaulerLogLevel + " | " + consts.HaulerAuditLevel + " | " + consts.HaulerConcurrency + " | " + consts.HaulerBlobConcurrency + "\n Warnings: Hauler commands and flags marked with (EXPERIMENTAL) are not yet stable and may change in the future.",
Example: " View the Docs: https://docs.hauler.dev\n Environment Variables: " + consts.HaulerDir + " | " + consts.HaulerTempDir + " | " + consts.HaulerStoreDir + " | " + consts.HaulerIgnoreErrors + " | " + consts.HaulerRetries + " | " + consts.HaulerLogLevel + " | " + consts.HaulerAuditLevel + " | " + consts.HaulerConcurrency + " | " + consts.HaulerBlobConcurrency + "\n Warnings: Hauler commands and flags marked with (EXPERIMENTAL) are not yet stable and may change in the future.",
PersistentPreRunE: func(cmd *cobra.Command, args []string) error {
// check for log level env variable or flag
if ro.LogLevel == "" {
Expand Down
7 changes: 2 additions & 5 deletions cmd/hauler/cli/store/add.go
Original file line number Diff line number Diff line change
Expand Up @@ -1002,11 +1002,8 @@ func runChartJobs(ctx context.Context, s *store.Layout, jobs []chartJob, concurr

l := log.FromContext(ctx)

tempOverride := rso.TempOverride
if tempOverride == "" {
tempOverride = os.Getenv(consts.HaulerTempDir)
}
tempRoot, err := os.MkdirTemp(tempOverride, consts.DefaultHaulerTempDirName)
// rso.TempOverride is already resolved (flag or HAULER_TEMP_DIR) by Store().
tempRoot, err := os.MkdirTemp(rso.TempOverride, consts.DefaultHaulerTempDirName)
if err != nil {
return fmt.Errorf("failed to create temp dir: %w", err)
}
Expand Down
5 changes: 1 addition & 4 deletions cmd/hauler/cli/store/load.go
Original file line number Diff line number Diff line change
Expand Up @@ -31,12 +31,9 @@ var legacyChunkRe = regexp.MustCompile(`_\d+\.`)
func LoadCmd(ctx context.Context, o *flags.LoadOpts, rso *flags.StoreRootOpts, ro *flags.CliRootOpts) error {
l := log.FromContext(ctx)

// rso.TempOverride is already resolved (flag or HAULER_TEMP_DIR) by Store().
tempOverride := rso.TempOverride

if tempOverride == "" {
tempOverride = os.Getenv(consts.HaulerTempDir)
}

tempDir, err := os.MkdirTemp(tempOverride, consts.DefaultHaulerTempDirName)
if err != nil {
return err
Expand Down
125 changes: 109 additions & 16 deletions cmd/hauler/cli/store/sync.go
Original file line number Diff line number Diff line change
Expand Up @@ -9,6 +9,7 @@ import (
"net/url"
"os"
"path/filepath"
"strconv"
"strings"
"time"

Expand Down Expand Up @@ -104,27 +105,33 @@ func SyncCmd(ctx context.Context, o *flags.SyncOpts, s *store.Layout, rso *flags
return nil
}

// caches stores opened via hauler.dev/store, keyed by abs path, so docs sharing a target reuse one Layout
targetStores := map[string]*store.Layout{}

// Everything below runs with a real store (s != nil; the dry-run branch
// above already returned). Force one durable index checkpoint at the end
// of the run since the per-artifact path only fsyncs on
// indexCheckpointInterval -- deferred so it still runs on error paths,
// where a partially-populated index is worth persisting. This does NOT
// run on Ctrl-C (no signal handler is installed), which is fine: process
// death doesn't lose page cache, so the index still reaches disk.
// death doesn't lose page cache, so the index still reaches disk. Covers
// any target stores from hauler.dev/store too, not just the primary one.
defer func() {
if err := s.OCI.SaveIndex(); err != nil {
l.Warnf("failed to save index at end of sync: %v", err)
}
l.Debugf("%s", formatIOStats(s.OCI.Stats().Snapshot(), s.OCI.BlobConcurrency()))
}()

tempOverride := rso.TempOverride

if tempOverride == "" {
tempOverride = os.Getenv(consts.HaulerTempDir)
}
for _, ts := range targetStores {
if err := ts.OCI.SaveIndex(); err != nil {
l.Warnf("failed to save index for target store [%s]: %v", ts.Root, err)
}
l.Debugf("%s", formatIOStats(ts.OCI.Stats().Snapshot(), ts.OCI.BlobConcurrency()))
}
}()

tempDir, err := os.MkdirTemp(tempOverride, consts.DefaultHaulerTempDirName)
// rso.TempOverride is already resolved (flag or HAULER_TEMP_DIR) by Store().
tempDir, err := os.MkdirTemp(rso.TempOverride, consts.DefaultHaulerTempDirName)
if err != nil {
return err
}
Expand Down Expand Up @@ -164,7 +171,7 @@ func SyncCmd(ctx context.Context, o *flags.SyncOpts, s *store.Layout, rso *flags
return err
}
defer fi.Close()
err = processContent(ctx, fi, o, s, rso, ro)
err = processContent(ctx, fi, o, s, rso, ro, targetStores)
if err != nil {
return err
}
Expand Down Expand Up @@ -214,7 +221,7 @@ func SyncCmd(ctx context.Context, o *flags.SyncOpts, s *store.Layout, rso *flags
}
defer fi.Close()

err = processContent(ctx, fi, o, s, rso, ro)
err = processContent(ctx, fi, o, s, rso, ro, targetStores)
if err != nil {
return err
}
Expand Down Expand Up @@ -278,7 +285,7 @@ func SyncCmd(ctx context.Context, o *flags.SyncOpts, s *store.Layout, rso *flags
return nil
}

func processContent(ctx context.Context, fi *os.File, o *flags.SyncOpts, s *store.Layout, rso *flags.StoreRootOpts, ro *flags.CliRootOpts) error {
func processContent(ctx context.Context, fi *os.File, o *flags.SyncOpts, s *store.Layout, rso *flags.StoreRootOpts, ro *flags.CliRootOpts, targetStores map[string]*store.Layout) error {
l := log.FromContext(ctx)

reader := yaml.NewYAMLReader(bufio.NewReader(fi))
Expand All @@ -303,7 +310,6 @@ func processContent(ctx context.Context, fi *os.File, o *flags.SyncOpts, s *stor
}

gvk := obj.GroupVersionKind()
l.Infof("syncing content [%s] with [kind=%s] to store [%s]", gvk.GroupVersion(), gvk.Kind, o.StoreDir)

switch gvk.Kind {

Expand All @@ -314,8 +320,18 @@ func processContent(ctx context.Context, fi *os.File, o *flags.SyncOpts, s *stor
if err := yaml.Unmarshal(doc, &cfg); err != nil {
return err
}
a := cfg.GetAnnotations()
docStore, err := resolveTargetStore(ctx, a, s, rso, ro, targetStores)
if err != nil {
return err
}
docRso, err := resolveDocRetries(a, rso)
if err != nil {
return err
}
l.Infof("syncing content [%s] with [kind=%s] to store [%s]", gvk.GroupVersion(), gvk.Kind, docStore.Root)
jobs := resolveFileJobs(cfg.Spec.Files)
if err := runFileJobs(ctx, s, jobs, o.Concurrency, rso, ro, newSyncProgress(o, ro)); err != nil {
if err := runFileJobs(ctx, docStore, jobs, o.Concurrency, docRso, ro, newSyncProgress(o, ro)); err != nil {
return err
}

Expand All @@ -332,11 +348,20 @@ func processContent(ctx context.Context, fi *os.File, o *flags.SyncOpts, s *stor
}

a := cfg.GetAnnotations()
docStore, err := resolveTargetStore(ctx, a, s, rso, ro, targetStores)
if err != nil {
return err
}
docRso, err := resolveDocRetries(a, rso)
if err != nil {
return err
}
l.Infof("syncing content [%s] with [kind=%s] to store [%s]", gvk.GroupVersion(), gvk.Kind, docStore.Root)
jobs, err := resolveImageJobs(o, a, cfg.Spec.Images)
if err != nil {
return err
}
if err := runImageJobs(ctx, s, jobs, o.Concurrency, rso, ro, newSyncProgress(o, ro)); err != nil {
if err := runImageJobs(ctx, docStore, jobs, o.Concurrency, docRso, ro, newSyncProgress(o, ro)); err != nil {
return err
}

Expand All @@ -351,11 +376,21 @@ func processContent(ctx context.Context, fi *os.File, o *flags.SyncOpts, s *stor
if err := yaml.Unmarshal(doc, &cfg); err != nil {
return err
}
jobs, err := resolveChartJobs(o, cfg.GetAnnotations(), filepath.Dir(fi.Name()), cfg.Spec.Charts)
a := cfg.GetAnnotations()
docStore, err := resolveTargetStore(ctx, a, s, rso, ro, targetStores)
if err != nil {
return err
}
docRso, err := resolveDocRetries(a, rso)
if err != nil {
return err
}
l.Infof("syncing content [%s] with [kind=%s] to store [%s]", gvk.GroupVersion(), gvk.Kind, docStore.Root)
jobs, err := resolveChartJobs(o, a, filepath.Dir(fi.Name()), cfg.Spec.Charts)
if err != nil {
return err
}
if err := runChartJobs(ctx, s, jobs, o.Concurrency, rso, ro, newSyncProgress(o, ro)); err != nil {
if err := runChartJobs(ctx, docStore, jobs, o.Concurrency, docRso, ro, newSyncProgress(o, ro)); err != nil {
return err
}

Expand All @@ -370,6 +405,64 @@ func processContent(ctx context.Context, fi *os.File, o *flags.SyncOpts, s *stor
return nil
}

// resolveTargetStore picks a doc's store based on its hauler.dev/store annotation,
// falling back to def. Opens (or reuses, via targetStores) the target store otherwise.
func resolveTargetStore(ctx context.Context, a map[string]string, def *store.Layout, rso *flags.StoreRootOpts, ro *flags.CliRootOpts, targetStores map[string]*store.Layout) (*store.Layout, error) {
target := a[consts.AnnotationTargetStore]
if target == "" {
return def, nil
}

abs, err := flags.ResolveStoreDir(ctx, ro, target)
if err != nil {
return nil, fmt.Errorf("failed to resolve target store [%s]: %w", target, err)
}

if abs == def.Root {
return def, nil
}

if ts, ok := targetStores[abs]; ok {
return ts, nil
}

// only overriding StoreDir, everything else still comes from rso
altOpts := *rso
altOpts.StoreDir = abs
ts, err := altOpts.Store(ctx, ro)
if err != nil {
return nil, fmt.Errorf("failed to open target store [%s]: %w", target, err)
}

targetStores[abs] = ts
return ts, nil
}

// resolveDocRetries returns a copy of rso with Retries overridden by a doc's
// hauler.dev/retries annotation, or rso unchanged if it's not set. Copy, not
// mutation, so it can't leak into a sibling doc.
func resolveDocRetries(a map[string]string, rso *flags.StoreRootOpts) (*flags.StoreRootOpts, error) {
v, ok := a[consts.AnnotationRetries]
if !ok || v == "" {
return rso, nil
}

n, err := strconv.Atoi(v)
if err != nil {
return nil, fmt.Errorf("invalid %s value %q: %w", consts.AnnotationRetries, v, err)
}
if n < 0 {
return nil, fmt.Errorf("%s must be >= 0, got %d", consts.AnnotationRetries, n)
}
if n == 0 {
n = consts.DefaultRetries
}

docRso := *rso
docRso.Retries = n
return &docRso, nil
}

// resolveChartCreds reads credentials for a Chart entry from the env vars
// named by UsernameEnv and PasswordEnv. Both fields must be set or both must
// be empty; a mix is a configuration error. If both are set, the env vars
Expand Down
Loading
Loading