99 "net/url"
1010 "os"
1111 "path/filepath"
12+ "strconv"
1213 "strings"
1314 "time"
1415
@@ -104,27 +105,33 @@ func SyncCmd(ctx context.Context, o *flags.SyncOpts, s *store.Layout, rso *flags
104105 return nil
105106 }
106107
108+ // caches stores opened via hauler.dev/store, keyed by abs path, so docs sharing a target reuse one Layout
109+ targetStores := map [string ]* store.Layout {}
110+
107111 // Everything below runs with a real store (s != nil; the dry-run branch
108112 // above already returned). Force one durable index checkpoint at the end
109113 // of the run since the per-artifact path only fsyncs on
110114 // indexCheckpointInterval -- deferred so it still runs on error paths,
111115 // where a partially-populated index is worth persisting. This does NOT
112116 // run on Ctrl-C (no signal handler is installed), which is fine: process
113- // death doesn't lose page cache, so the index still reaches disk.
117+ // death doesn't lose page cache, so the index still reaches disk. Covers
118+ // any target stores from hauler.dev/store too, not just the primary one.
114119 defer func () {
115120 if err := s .OCI .SaveIndex (); err != nil {
116121 l .Warnf ("failed to save index at end of sync: %v" , err )
117122 }
118123 l .Debugf ("%s" , formatIOStats (s .OCI .Stats ().Snapshot (), s .OCI .BlobConcurrency ()))
119- }()
120124
121- tempOverride := rso .TempOverride
122-
123- if tempOverride == "" {
124- tempOverride = os .Getenv (consts .HaulerTempDir )
125- }
125+ for _ , ts := range targetStores {
126+ if err := ts .OCI .SaveIndex (); err != nil {
127+ l .Warnf ("failed to save index for target store [%s]: %v" , ts .Root , err )
128+ }
129+ l .Debugf ("%s" , formatIOStats (ts .OCI .Stats ().Snapshot (), ts .OCI .BlobConcurrency ()))
130+ }
131+ }()
126132
127- tempDir , err := os .MkdirTemp (tempOverride , consts .DefaultHaulerTempDirName )
133+ // rso.TempOverride is already resolved (flag or HAULER_TEMP_DIR) by Store().
134+ tempDir , err := os .MkdirTemp (rso .TempOverride , consts .DefaultHaulerTempDirName )
128135 if err != nil {
129136 return err
130137 }
@@ -166,7 +173,7 @@ func SyncCmd(ctx context.Context, o *flags.SyncOpts, s *store.Layout, rso *flags
166173 return err
167174 }
168175 defer fi .Close ()
169- err = processContent (ctx , fi , o , s , rso , ro )
176+ err = processContent (ctx , fi , o , s , rso , ro , targetStores )
170177 if err != nil {
171178 return err
172179 }
@@ -216,7 +223,7 @@ func SyncCmd(ctx context.Context, o *flags.SyncOpts, s *store.Layout, rso *flags
216223 }
217224 defer fi .Close ()
218225
219- err = processContent (ctx , fi , o , s , rso , ro )
226+ err = processContent (ctx , fi , o , s , rso , ro , targetStores )
220227 if err != nil {
221228 return err
222229 }
@@ -304,7 +311,7 @@ func derefInsecure(p *bool) bool {
304311 return p != nil && * p
305312}
306313
307- func processContent (ctx context.Context , fi * os.File , o * flags.SyncOpts , s * store.Layout , rso * flags.StoreRootOpts , ro * flags.CliRootOpts ) error {
314+ 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 {
308315 l := log .FromContext (ctx )
309316
310317 reader := yaml .NewYAMLReader (bufio .NewReader (fi ))
@@ -329,7 +336,6 @@ func processContent(ctx context.Context, fi *os.File, o *flags.SyncOpts, s *stor
329336 }
330337
331338 gvk := obj .GroupVersionKind ()
332- l .Infof ("syncing content [%s] with [kind=%s] to store [%s]" , gvk .GroupVersion (), gvk .Kind , o .StoreDir )
333339
334340 switch gvk .Kind {
335341
@@ -340,8 +346,18 @@ func processContent(ctx context.Context, fi *os.File, o *flags.SyncOpts, s *stor
340346 if err := yaml .Unmarshal (doc , & cfg ); err != nil {
341347 return err
342348 }
343- jobs := resolveFileJobs (o , cfg .GetAnnotations (), cfg .Spec .Files )
344- if err := runFileJobs (ctx , s , jobs , o .Concurrency , rso , ro , newSyncProgress (o , ro )); err != nil {
349+ a := cfg .GetAnnotations ()
350+ docStore , err := resolveTargetStore (ctx , a , s , rso , ro , targetStores )
351+ if err != nil {
352+ return err
353+ }
354+ docRso , err := resolveDocRetries (a , rso )
355+ if err != nil {
356+ return err
357+ }
358+ l .Infof ("syncing content [%s] with [kind=%s] to store [%s]" , gvk .GroupVersion (), gvk .Kind , docStore .Root )
359+ jobs := resolveFileJobs (o , a , cfg .Spec .Files )
360+ if err := runFileJobs (ctx , docStore , jobs , o .Concurrency , docRso , ro , newSyncProgress (o , ro )); err != nil {
345361 return err
346362 }
347363
@@ -358,11 +374,20 @@ func processContent(ctx context.Context, fi *os.File, o *flags.SyncOpts, s *stor
358374 }
359375
360376 a := cfg .GetAnnotations ()
377+ docStore , err := resolveTargetStore (ctx , a , s , rso , ro , targetStores )
378+ if err != nil {
379+ return err
380+ }
381+ docRso , err := resolveDocRetries (a , rso )
382+ if err != nil {
383+ return err
384+ }
385+ l .Infof ("syncing content [%s] with [kind=%s] to store [%s]" , gvk .GroupVersion (), gvk .Kind , docStore .Root )
361386 jobs , err := resolveImageJobs (o , a , cfg .Spec .Images )
362387 if err != nil {
363388 return err
364389 }
365- if err := runImageJobs (ctx , s , jobs , o .Concurrency , rso , ro , newSyncProgress (o , ro )); err != nil {
390+ if err := runImageJobs (ctx , docStore , jobs , o .Concurrency , docRso , ro , newSyncProgress (o , ro )); err != nil {
366391 return err
367392 }
368393
@@ -377,11 +402,21 @@ func processContent(ctx context.Context, fi *os.File, o *flags.SyncOpts, s *stor
377402 if err := yaml .Unmarshal (doc , & cfg ); err != nil {
378403 return err
379404 }
380- jobs , err := resolveChartJobs (o , cfg .GetAnnotations (), filepath .Dir (fi .Name ()), cfg .Spec .Charts )
405+ a := cfg .GetAnnotations ()
406+ docStore , err := resolveTargetStore (ctx , a , s , rso , ro , targetStores )
407+ if err != nil {
408+ return err
409+ }
410+ docRso , err := resolveDocRetries (a , rso )
411+ if err != nil {
412+ return err
413+ }
414+ l .Infof ("syncing content [%s] with [kind=%s] to store [%s]" , gvk .GroupVersion (), gvk .Kind , docStore .Root )
415+ jobs , err := resolveChartJobs (o , a , filepath .Dir (fi .Name ()), cfg .Spec .Charts )
381416 if err != nil {
382417 return err
383418 }
384- if err := runChartJobs (ctx , s , jobs , o .Concurrency , rso , ro , newSyncProgress (o , ro )); err != nil {
419+ if err := runChartJobs (ctx , docStore , jobs , o .Concurrency , docRso , ro , newSyncProgress (o , ro )); err != nil {
385420 return err
386421 }
387422
@@ -396,6 +431,64 @@ func processContent(ctx context.Context, fi *os.File, o *flags.SyncOpts, s *stor
396431 return nil
397432}
398433
434+ // resolveTargetStore picks a doc's store based on its hauler.dev/store annotation,
435+ // falling back to def. Opens (or reuses, via targetStores) the target store otherwise.
436+ 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 ) {
437+ target := a [consts .AnnotationTargetStore ]
438+ if target == "" {
439+ return def , nil
440+ }
441+
442+ abs , err := flags .ResolveStoreDir (ctx , ro , target )
443+ if err != nil {
444+ return nil , fmt .Errorf ("failed to resolve target store [%s]: %w" , target , err )
445+ }
446+
447+ if abs == def .Root {
448+ return def , nil
449+ }
450+
451+ if ts , ok := targetStores [abs ]; ok {
452+ return ts , nil
453+ }
454+
455+ // only overriding StoreDir, everything else still comes from rso
456+ altOpts := * rso
457+ altOpts .StoreDir = abs
458+ ts , err := altOpts .Store (ctx , ro )
459+ if err != nil {
460+ return nil , fmt .Errorf ("failed to open target store [%s]: %w" , target , err )
461+ }
462+
463+ targetStores [abs ] = ts
464+ return ts , nil
465+ }
466+
467+ // resolveDocRetries returns a copy of rso with Retries overridden by a doc's
468+ // hauler.dev/retries annotation, or rso unchanged if it's not set. Copy, not
469+ // mutation, so it can't leak into a sibling doc.
470+ func resolveDocRetries (a map [string ]string , rso * flags.StoreRootOpts ) (* flags.StoreRootOpts , error ) {
471+ v , ok := a [consts .AnnotationRetries ]
472+ if ! ok || v == "" {
473+ return rso , nil
474+ }
475+
476+ n , err := strconv .Atoi (v )
477+ if err != nil {
478+ return nil , fmt .Errorf ("invalid %s value %q: %w" , consts .AnnotationRetries , v , err )
479+ }
480+ if n < 0 {
481+ return nil , fmt .Errorf ("%s must be >= 0, got %d" , consts .AnnotationRetries , n )
482+ }
483+ if n == 0 {
484+ n = consts .DefaultRetries
485+ }
486+
487+ docRso := * rso
488+ docRso .Retries = n
489+ return & docRso , nil
490+ }
491+
399492// resolveChartCreds reads credentials for a Chart entry from the env vars
400493// named by UsernameEnv and PasswordEnv. Both fields must be set or both must
401494// be empty; a mix is a configuration error. If both are set, the env vars
0 commit comments