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
9 changes: 7 additions & 2 deletions internal/exec/stack_processor_utils.go
Original file line number Diff line number Diff line change
Expand Up @@ -1214,9 +1214,14 @@ func processYAMLConfigFileWithContextInternal(

// Check if the import is a remote URL.
if isRemote {
// Resolve the per-import cache TTL, falling back to the global default.
importTTL := importStruct.TTL
if importTTL == "" && atmosConfig != nil {
importTTL = atmosConfig.Imports.TTL
}
// Download the remote import.
log.Debug("Downloading remote stack import", "uri", imp, "file", relativeFilePath, "nested_imports", nestedImports)
remoteMatches, err := stackimports.ResolveRemoteImportNested(atmosConfig, imp, nestedImports)
log.Debug("Downloading remote stack import", "uri", imp, "file", relativeFilePath, "nested_imports", nestedImports, "ttl", importTTL)
remoteMatches, err := stackimports.ResolveRemoteImportNested(atmosConfig, imp, nestedImports, importTTL)
if err != nil {
if importStruct.SkipIfMissing {
log.Debug("Skipping missing remote import", "uri", imp)
Expand Down
11 changes: 11 additions & 0 deletions pkg/datafetcher/schema/atmos/manifest/1.0.json
Original file line number Diff line number Diff line change
Expand Up @@ -147,6 +147,17 @@
"skip_if_missing": {
"type": "boolean"
},
"nested_imports": {
"type": "string",
"enum": [
"local",
"remote"
]
},
"ttl": {
"type": "string",
"description": "Cache duration for the cloned source repo of a remote (git) import. When set, the clone is reused across Atmos invocations until it expires (e.g. with a warm CI cache). Within a single invocation a source repo is cloned at most once regardless of TTL. Examples: '0s', '5m', '1h', '7d', 'daily'."
},
"context": {
"type": "object",
"additionalProperties": true
Expand Down
11 changes: 11 additions & 0 deletions pkg/datafetcher/schema/config/global/1.0.json
Original file line number Diff line number Diff line change
Expand Up @@ -140,6 +140,17 @@
"skip_if_missing": {
"type": "boolean"
},
"nested_imports": {
"type": "string",
"enum": [
"local",
"remote"
]
},
"ttl": {
"type": "string",
"description": "Cache duration for the cloned source repo of a remote (git) import. When set, the clone is reused across Atmos invocations until it expires (e.g. with a warm CI cache). Within a single invocation a source repo is cloned at most once regardless of TTL. Examples: '0s', '5m', '1h', '7d', 'daily'."
},
"context": {
"type": "object",
"additionalProperties": true
Expand Down
4 changes: 4 additions & 0 deletions pkg/datafetcher/schema/stacks/stack-config/1.0.json
Original file line number Diff line number Diff line change
Expand Up @@ -156,6 +156,10 @@
"remote"
]
},
"ttl": {
"type": "string",
"description": "Cache duration for the cloned source repo of a remote (git) import. When set, the clone is reused across Atmos invocations until it expires (e.g. with a warm CI cache). Within a single invocation a source repo is cloned at most once regardless of TTL. Examples: '0s', '5m', '1h', '7d', 'daily'."
},
"context": {
"type": "object",
"additionalProperties": true
Expand Down
38 changes: 38 additions & 0 deletions pkg/duration/duration.go
Original file line number Diff line number Diff line change
Expand Up @@ -170,3 +170,41 @@ func ParseDuration(s string) (time.Duration, error) {

return time.Duration(seconds) * time.Second, nil
}

// IsZeroTTL reports whether the TTL string represents a zero duration.
//
// ParseDuration rejects "0" (a zero duration is not a valid period for cleanup
// scheduling), but for cache TTLs a zero value is a legitimate, common case meaning
// "always expired / never reuse". Callers detect it with this helper before parsing.
func IsZeroTTL(ttl string) bool {
defer perf.Track(nil, "duration.IsZeroTTL")()

switch strings.TrimSpace(ttl) {
case "0", "0s", "0m", "0h", "0d":
return true
default:
return false
}
}

// IsExpired reports whether a cache entry last refreshed at updatedAt has exceeded ttl.
//
// A zero TTL (see IsZeroTTL) is always expired. A non-zero TTL is parsed with
// ParseDuration; on a parse error the error is returned and callers should fail safe
// by treating the entry as expired to avoid serving stale data. This helper makes no
// decision for an empty ttl (""). The meaning of "no TTL configured" differs between
// subsystems (reuse-forever versus refresh-each-run), so callers handle "" explicitly.
func IsExpired(updatedAt time.Time, ttl string) (bool, error) {
defer perf.Track(nil, "duration.IsExpired")()

if IsZeroTTL(ttl) {
return true, nil
}

ttlDuration, err := ParseDuration(ttl)
if err != nil {
return false, err
}

return time.Since(updatedAt) > ttlDuration, nil
}
63 changes: 63 additions & 0 deletions pkg/duration/duration_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -148,3 +148,66 @@ func TestParseDuration_OverflowProtection(t *testing.T) {
})
}
}

func TestIsZeroTTL(t *testing.T) {
tests := []struct {
name string
input string
expected bool
}{
// Zero values.
{name: "bare zero", input: "0", expected: true},
{name: "zero seconds", input: "0s", expected: true},
{name: "zero minutes", input: "0m", expected: true},
{name: "zero hours", input: "0h", expected: true},
{name: "zero days", input: "0d", expected: true},

// Zero values with surrounding whitespace (regression for trimming).
{name: "whitespace around zero seconds", input: " 0s ", expected: true},
{name: "whitespace around bare zero", input: " 0 ", expected: true},
{name: "tabs and newlines around zero hours", input: "\t0h\n", expected: true},

// Non-zero / non-TTL values.
{name: "one hour", input: "1h", expected: false},
{name: "thirty minutes", input: "30m", expected: false},
{name: "empty string", input: "", expected: false},
{name: "keyword", input: "daily", expected: false},
{name: "whitespace only", input: " ", expected: false},
}

for _, tt := range tests {
t.Run(tt.name, func(t *testing.T) {
assert.Equal(t, tt.expected, IsZeroTTL(tt.input), "unexpected result for input %q", tt.input)
})
}
}

func TestIsExpired(t *testing.T) {
now := time.Now()

tests := []struct {
name string
updatedAt time.Time
ttl string
wantExpired bool
wantErr bool
}{
{name: "zero TTL is always expired", updatedAt: now, ttl: "0s", wantExpired: true},
{name: "zero TTL with whitespace is always expired", updatedAt: now, ttl: " 0s ", wantExpired: true},
{name: "fresh entry within TTL", updatedAt: now.Add(-1 * time.Minute), ttl: "1h", wantExpired: false},
{name: "stale entry beyond TTL", updatedAt: now.Add(-2 * time.Hour), ttl: "1h", wantExpired: true},
{name: "invalid TTL returns error", updatedAt: now, ttl: "nonsense", wantErr: true},
}

for _, tt := range tests {
t.Run(tt.name, func(t *testing.T) {
expired, err := IsExpired(tt.updatedAt, tt.ttl)
if tt.wantErr {
assert.Error(t, err, "expected error for ttl %q", tt.ttl)
return
}
require.NoError(t, err, "unexpected error for ttl %q", tt.ttl)
assert.Equal(t, tt.wantExpired, expired, "unexpected expiry for ttl %q", tt.ttl)
})
}
}
12 changes: 7 additions & 5 deletions pkg/provisioner/source/provision_hook.go
Original file line number Diff line number Diff line change
Expand Up @@ -360,19 +360,21 @@ func checkMetadataChanges(metadata *workdir.WorkdirMetadata, sourceSpec *schema.
}

// isSourceCacheExpired checks if the source cache has expired based on TTL.
// A TTL of "0" or "0s" means always expired (always re-pull).
// A TTL of "0" or "0s" means always expired (always re-pull). The expiry decision
// is delegated to the shared duration.IsExpired helper; this wrapper adds the
// source-provisioning-specific human-readable reason strings.
func isSourceCacheExpired(ttl string, updatedAt time.Time) (bool, string) {
// Handle zero TTL explicitly (always expired).
// Handle zero TTL explicitly (always expired) so we can tailor the reason.
if isZeroTTL(ttl) {
return true, fmt.Sprintf("Source cache expired (TTL: %s, always re-pull)", ttl)
}

ttlDuration, err := duration.ParseDuration(ttl)
expired, err := duration.IsExpired(updatedAt, ttl)
if err != nil {
return true, fmt.Sprintf("Invalid source TTL %q; forcing re-provision to avoid stale cache", ttl)
}

if time.Since(updatedAt) > ttlDuration {
if expired {
return true, fmt.Sprintf("Source cache expired (TTL: %s, last updated: %s)",
ttl, updatedAt.Format(time.RFC3339))
}
Expand All @@ -381,7 +383,7 @@ func isSourceCacheExpired(ttl string, updatedAt time.Time) (bool, string) {

// isZeroTTL checks if the TTL string represents a zero duration.
func isZeroTTL(ttl string) bool {
return ttl == "0" || ttl == "0s" || ttl == "0m" || ttl == "0h" || ttl == "0d"
return duration.IsZeroTTL(ttl)
}

// isLocalSource determines if a source URI refers to a local path.
Expand Down
16 changes: 16 additions & 0 deletions pkg/schema/schema.go
Original file line number Diff line number Diff line change
Expand Up @@ -62,6 +62,7 @@ type AtmosConfiguration struct {
BasePathSource string `yaml:"-" json:"-" mapstructure:"-"` // "runtime" if from env var/CLI/provider, "" if from config file.
Components Components `yaml:"components" json:"components" mapstructure:"components"`
Stacks Stacks `yaml:"stacks" json:"stacks" mapstructure:"stacks"`
Imports ImportsSettings `yaml:"imports,omitempty" json:"imports,omitempty" mapstructure:"imports"`
Workflows Workflows `yaml:"workflows,omitempty" json:"workflows,omitempty" mapstructure:"workflows"`
Logs Logs `yaml:"logs,omitempty" json:"logs,omitempty" mapstructure:"logs"`
Errors ErrorsConfig `yaml:"errors,omitempty" json:"errors,omitempty" mapstructure:"errors"`
Expand Down Expand Up @@ -1262,13 +1263,28 @@ const (
StackImportNestedImportsRemote = "remote"
)

// ImportsSettings holds global defaults applied when processing stack imports.
type ImportsSettings struct {
// TTL is the default cache duration for the cloned source repos of remote (git)
// imports, applied when a per-import `ttl` is not set. See StackImport.TTL for the
// caching semantics and accepted formats.
TTL string `yaml:"ttl,omitempty" json:"ttl,omitempty" mapstructure:"ttl"`
}

type StackImport struct {
Path string `yaml:"path" json:"path" mapstructure:"path"`
Context AtmosSectionMapType `yaml:"context" json:"context" mapstructure:"context"`
SkipTemplatesProcessing bool `yaml:"skip_templates_processing,omitempty" json:"skip_templates_processing,omitempty" mapstructure:"skip_templates_processing"`
IgnoreMissingTemplateValues bool `yaml:"ignore_missing_template_values,omitempty" json:"ignore_missing_template_values,omitempty" mapstructure:"ignore_missing_template_values"`
SkipIfMissing bool `yaml:"skip_if_missing,omitempty" json:"skip_if_missing,omitempty" mapstructure:"skip_if_missing"`
NestedImports string `yaml:"nested_imports,omitempty" json:"nested_imports,omitempty" mapstructure:"nested_imports"`
// TTL is the cache duration for the cloned source repo of a remote (git) import.
// When set, the cloned source is reused across Atmos invocations until it expires,
// so a warm cache (e.g. GitHub Actions cache of the imports cache dir) skips the
// re-clone. Within a single invocation a source repo is always cloned at most once,
// regardless of TTL. When unset, the source is re-cloned once per invocation.
// Examples: "0s" (always re-fetch), "5m", "1h", "7d", or keywords like "daily".
TTL string `yaml:"ttl,omitempty" json:"ttl,omitempty" mapstructure:"ttl"`
}

// Dependencies
Expand Down
72 changes: 42 additions & 30 deletions pkg/stack/imports/remote.go
Original file line number Diff line number Diff line change
@@ -1,7 +1,6 @@
package imports

import (
"context"
"crypto/sha256"
"errors"
"fmt"
Expand All @@ -14,7 +13,6 @@ import (

"github.com/bmatcuk/doublestar/v4"
errUtils "github.com/cloudposse/atmos/errors"
"github.com/cloudposse/atmos/pkg/auth/broker"
"github.com/cloudposse/atmos/pkg/cache"
"github.com/cloudposse/atmos/pkg/downloader"
"github.com/cloudposse/atmos/pkg/perf"
Expand Down Expand Up @@ -46,8 +44,12 @@ type RemoteImporter struct {
cache *cache.FileCache
memCache map[string]string // In-memory cache for session.
matchCache map[string][]RemoteImportMatch
memMu sync.RWMutex
sourceMu sync.Mutex
// sessionFetched tracks source repos already fetched in this process so that a
// repo shared by many subdir imports is cloned at most once per invocation,
// regardless of TTL. Guarded by sourceMu.
sessionFetched map[string]bool
memMu sync.RWMutex
sourceMu sync.Mutex
}

// RemoteImporterOption is a functional option for configuring RemoteImporter.
Expand Down Expand Up @@ -85,11 +87,12 @@ func NewRemoteImporter(atmosConfig *schema.AtmosConfiguration, opts ...RemoteImp
fd := downloader.NewGoGetterDownloader(atmosConfig)

r := &RemoteImporter{
atmosConfig: atmosConfig,
downloader: fd,
cache: fileCache,
memCache: make(map[string]string),
matchCache: make(map[string][]RemoteImportMatch),
atmosConfig: atmosConfig,
downloader: fd,
cache: fileCache,
memCache: make(map[string]string),
matchCache: make(map[string][]RemoteImportMatch),
sessionFetched: make(map[string]bool),
}

// Apply options.
Expand Down Expand Up @@ -169,9 +172,21 @@ func (r *RemoteImporter) Download(uri string) (string, error) {
}

// Resolve fetches a remote import and returns all local stack files it resolves to.
//
// The cached source clone is refreshed on every invocation (no cross-run reuse); use
// ResolveRemoteImportNested with a TTL when cross-run cache reuse is desired.
func (r *RemoteImporter) Resolve(uri string) ([]RemoteImportMatch, error) {
defer perf.Track(nil, "imports.RemoteImporter.Resolve")()

return r.resolve(uri, "")
}

// resolve fetches a remote import and returns all local stack files it resolves to.
// The ttl controls cross-run reuse of the cloned source repo for git subdir imports
// (see ensureSourceDir); an empty ttl refreshes the clone once per invocation.
func (r *RemoteImporter) resolve(uri, ttl string) ([]RemoteImportMatch, error) {
defer perf.Track(nil, "imports.RemoteImporter.resolve")()

if !IsRemote(uri) {
return nil, errUtils.Build(errUtils.ErrInvalidRemoteImport).
WithExplanation("URI is not a remote URL").
Expand Down Expand Up @@ -204,7 +219,7 @@ func (r *RemoteImporter) Resolve(uri string) ([]RemoteImportMatch, error) {
return matches, nil
}

matches, err := r.resolveGitSubdir(uri, sourceURI, subdir)
matches, err := r.resolveGitSubdir(uri, sourceURI, subdir, ttl)
if err != nil {
return nil, err
}
Expand Down Expand Up @@ -247,36 +262,29 @@ func cloneMatches(matches []RemoteImportMatch) []RemoteImportMatch {
return cloned
}

func (r *RemoteImporter) resolveGitSubdir(originalURI, sourceURI, subdir string) ([]RemoteImportMatch, error) {
// Provision ambient credential brokers (e.g. Atmos Pro github/sts) before the eager detect below.
// detectGitSource runs CustomGitDetector.Detect directly — not through the downloader — so unlike
// downloader.Fetch it does not itself trigger the broker. Without provisioning here, the detector
// runs before the broker has exported its GIT_CONFIG_* insteadOf rewrite into the process env,
// misses the rewrite, and injects the ambient GITHUB_TOKEN — shadowing the minted least-privilege
// token and 404ing a cross-repo import. EnsureCredentials is process-once and gated (CI + config).
if r.atmosConfig != nil {
broker.EnsureCredentials(context.Background(), r.atmosConfig)
}

sourceURI = r.detectGitSource(sourceURI)
tempDir, err := os.MkdirTemp(r.cache.BaseDir(), uriToTempName(originalURI)+".dir-")
// resolveGitSubdir resolves a git import that targets a subdirectory (or glob) within
// a repo. It sources files from the shared, deduplicated clone produced by
// ensureSourceDir (cloned once per source repo per invocation, persisted across runs
// when a TTL is set) rather than re-cloning into a throwaway temp dir per import. The
// resolved files are copied into the file cache so the returned paths stay stable and
// the default ("local" nested imports) semantics are unchanged.
//
// Ambient credential brokers (e.g. Atmos Pro github/sts) are provisioned inside
// ensureSourceDir, before its detectGitSource call — see the note there.
func (r *RemoteImporter) resolveGitSubdir(originalURI, sourceURI, subdir, ttl string) ([]RemoteImportMatch, error) {
sourceRoot, err := r.ensureSourceDir(sourceURI, ttl)
if err != nil {
return nil, err
}
defer os.RemoveAll(tempDir)

if err := r.downloader.Fetch(sourceURI, tempDir, downloader.ClientModeDir, defaultDownloadTimeout); err != nil {
return nil, err
}

files, err := resolveStackFiles(tempDir, subdir)
files, err := resolveStackFiles(sourceRoot, subdir)
if err != nil {
return nil, err
}

matches := make([]RemoteImportMatch, 0, len(files))
for _, file := range files {
rel, err := filepath.Rel(tempDir, file)
rel, err := filepath.Rel(sourceRoot, file)
if err != nil {
return nil, err
}
Expand Down Expand Up @@ -458,6 +466,10 @@ func (r *RemoteImporter) ClearCache() error {
r.memCache = make(map[string]string)
r.matchCache = make(map[string][]RemoteImportMatch)

r.sourceMu.Lock()
r.sessionFetched = make(map[string]bool)
r.sourceMu.Unlock()

// Clear persistent cache.
return r.cache.Clear()
}
Expand Down
Loading
Loading