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 }
@@ -164,7 +171,7 @@ func SyncCmd(ctx context.Context, o *flags.SyncOpts, s *store.Layout, rso *flags
164171 return err
165172 }
166173 defer fi .Close ()
167- err = processContent (ctx , fi , o , s , rso , ro )
174+ err = processContent (ctx , fi , o , s , rso , ro , targetStores )
168175 if err != nil {
169176 return err
170177 }
@@ -214,7 +221,7 @@ func SyncCmd(ctx context.Context, o *flags.SyncOpts, s *store.Layout, rso *flags
214221 }
215222 defer fi .Close ()
216223
217- err = processContent (ctx , fi , o , s , rso , ro )
224+ err = processContent (ctx , fi , o , s , rso , ro , targetStores )
218225 if err != nil {
219226 return err
220227 }
@@ -278,7 +285,7 @@ func SyncCmd(ctx context.Context, o *flags.SyncOpts, s *store.Layout, rso *flags
278285 return nil
279286}
280287
281- func processContent (ctx context.Context , fi * os.File , o * flags.SyncOpts , s * store.Layout , rso * flags.StoreRootOpts , ro * flags.CliRootOpts ) error {
288+ 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 {
282289 l := log .FromContext (ctx )
283290
284291 reader := yaml .NewYAMLReader (bufio .NewReader (fi ))
@@ -303,7 +310,6 @@ func processContent(ctx context.Context, fi *os.File, o *flags.SyncOpts, s *stor
303310 }
304311
305312 gvk := obj .GroupVersionKind ()
306- l .Infof ("syncing content [%s] with [kind=%s] to store [%s]" , gvk .GroupVersion (), gvk .Kind , o .StoreDir )
307313
308314 switch gvk .Kind {
309315
@@ -314,8 +320,18 @@ func processContent(ctx context.Context, fi *os.File, o *flags.SyncOpts, s *stor
314320 if err := yaml .Unmarshal (doc , & cfg ); err != nil {
315321 return err
316322 }
323+ a := cfg .GetAnnotations ()
324+ docStore , err := resolveTargetStore (ctx , a , s , rso , ro , targetStores )
325+ if err != nil {
326+ return err
327+ }
328+ docRso , err := resolveDocRetries (a , rso )
329+ if err != nil {
330+ return err
331+ }
332+ l .Infof ("syncing content [%s] with [kind=%s] to store [%s]" , gvk .GroupVersion (), gvk .Kind , docStore .Root )
317333 jobs := resolveFileJobs (cfg .Spec .Files )
318- if err := runFileJobs (ctx , s , jobs , o .Concurrency , rso , ro , newSyncProgress (o , ro )); err != nil {
334+ if err := runFileJobs (ctx , docStore , jobs , o .Concurrency , docRso , ro , newSyncProgress (o , ro )); err != nil {
319335 return err
320336 }
321337
@@ -332,11 +348,20 @@ func processContent(ctx context.Context, fi *os.File, o *flags.SyncOpts, s *stor
332348 }
333349
334350 a := cfg .GetAnnotations ()
351+ docStore , err := resolveTargetStore (ctx , a , s , rso , ro , targetStores )
352+ if err != nil {
353+ return err
354+ }
355+ docRso , err := resolveDocRetries (a , rso )
356+ if err != nil {
357+ return err
358+ }
359+ l .Infof ("syncing content [%s] with [kind=%s] to store [%s]" , gvk .GroupVersion (), gvk .Kind , docStore .Root )
335360 jobs , err := resolveImageJobs (o , a , cfg .Spec .Images )
336361 if err != nil {
337362 return err
338363 }
339- if err := runImageJobs (ctx , s , jobs , o .Concurrency , rso , ro , newSyncProgress (o , ro )); err != nil {
364+ if err := runImageJobs (ctx , docStore , jobs , o .Concurrency , docRso , ro , newSyncProgress (o , ro )); err != nil {
340365 return err
341366 }
342367
@@ -351,11 +376,21 @@ func processContent(ctx context.Context, fi *os.File, o *flags.SyncOpts, s *stor
351376 if err := yaml .Unmarshal (doc , & cfg ); err != nil {
352377 return err
353378 }
354- jobs , err := resolveChartJobs (o , cfg .GetAnnotations (), filepath .Dir (fi .Name ()), cfg .Spec .Charts )
379+ a := cfg .GetAnnotations ()
380+ docStore , err := resolveTargetStore (ctx , a , s , rso , ro , targetStores )
381+ if err != nil {
382+ return err
383+ }
384+ docRso , err := resolveDocRetries (a , rso )
385+ if err != nil {
386+ return err
387+ }
388+ l .Infof ("syncing content [%s] with [kind=%s] to store [%s]" , gvk .GroupVersion (), gvk .Kind , docStore .Root )
389+ jobs , err := resolveChartJobs (o , a , filepath .Dir (fi .Name ()), cfg .Spec .Charts )
355390 if err != nil {
356391 return err
357392 }
358- if err := runChartJobs (ctx , s , jobs , o .Concurrency , rso , ro , newSyncProgress (o , ro )); err != nil {
393+ if err := runChartJobs (ctx , docStore , jobs , o .Concurrency , docRso , ro , newSyncProgress (o , ro )); err != nil {
359394 return err
360395 }
361396
@@ -370,6 +405,64 @@ func processContent(ctx context.Context, fi *os.File, o *flags.SyncOpts, s *stor
370405 return nil
371406}
372407
408+ // resolveTargetStore picks a doc's store based on its hauler.dev/store annotation,
409+ // falling back to def. Opens (or reuses, via targetStores) the target store otherwise.
410+ 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 ) {
411+ target := a [consts .AnnotationTargetStore ]
412+ if target == "" {
413+ return def , nil
414+ }
415+
416+ abs , err := flags .ResolveStoreDir (ctx , ro , target )
417+ if err != nil {
418+ return nil , fmt .Errorf ("failed to resolve target store [%s]: %w" , target , err )
419+ }
420+
421+ if abs == def .Root {
422+ return def , nil
423+ }
424+
425+ if ts , ok := targetStores [abs ]; ok {
426+ return ts , nil
427+ }
428+
429+ // only overriding StoreDir, everything else still comes from rso
430+ altOpts := * rso
431+ altOpts .StoreDir = abs
432+ ts , err := altOpts .Store (ctx , ro )
433+ if err != nil {
434+ return nil , fmt .Errorf ("failed to open target store [%s]: %w" , target , err )
435+ }
436+
437+ targetStores [abs ] = ts
438+ return ts , nil
439+ }
440+
441+ // resolveDocRetries returns a copy of rso with Retries overridden by a doc's
442+ // hauler.dev/retries annotation, or rso unchanged if it's not set. Copy, not
443+ // mutation, so it can't leak into a sibling doc.
444+ func resolveDocRetries (a map [string ]string , rso * flags.StoreRootOpts ) (* flags.StoreRootOpts , error ) {
445+ v , ok := a [consts .AnnotationRetries ]
446+ if ! ok || v == "" {
447+ return rso , nil
448+ }
449+
450+ n , err := strconv .Atoi (v )
451+ if err != nil {
452+ return nil , fmt .Errorf ("invalid %s value %q: %w" , consts .AnnotationRetries , v , err )
453+ }
454+ if n < 0 {
455+ return nil , fmt .Errorf ("%s must be >= 0, got %d" , consts .AnnotationRetries , n )
456+ }
457+ if n == 0 {
458+ n = consts .DefaultRetries
459+ }
460+
461+ docRso := * rso
462+ docRso .Retries = n
463+ return & docRso , nil
464+ }
465+
373466// resolveChartCreds reads credentials for a Chart entry from the env vars
374467// named by UsernameEnv and PasswordEnv. Both fields must be set or both must
375468// be empty; a mix is a configuration error. If both are set, the env vars
0 commit comments