mirror of
https://github.com/agent-substrate/substrate.git
synced 2026-10-02 03:24:42 +08:00
Today `atelet` and `ateom` derive the per-actor directories (oci
bundles, runsc state, pid files, checkpoint and restore state,
durable-dir, system-info and volume roots) from the actor UID through
the same package, `internal/ateompath`.
This PR:
1. Add an `ActorDirs` message to RunWorkloadRequest,
RestoreWorkloadRequest, CheckpointWorkloadRequest and
TerminateWorkloadRequest, and have atelet fill it with the directories
it prepared. **The ateoms do not read it yet.**
2. Split `internal/ateompath` into:
- `cmd/atelet/internal/ateletpath`: atelet's pathes, the per-actor
directories (and `ActorDirs`, built from them) plus the directories only
atelet uses.
- `internal/nodepath`: shared pathes, the base dir both mount, the ateom
socket, the OTLP
sockets, and the netns name.
- `internal/ateompath`: what the ateoms still derive from the actor UID,
each function marked with the `ActorDirs` field it duplicates. atelet no
longer imports it.
---------
This is the first part of #1604 to decouple atelet and ateom shared
pathes. Behavior **does not change**: both sides still compute the same
paths, atelet from `ateletpath` and the ateoms from `ateompath`.
Followup PRs will remove ateompath by making `ateom-gvisor` and
`ateom-microvm` read the `ActorDirs` message from RPC.
850 lines
32 KiB
Go
850 lines
32 KiB
Go
// Copyright 2026 Google LLC
|
|
//
|
|
// Licensed under the Apache License, Version 2.0 (the "License");
|
|
// you may not use this file except in compliance with the License.
|
|
// You may obtain a copy of the License at
|
|
//
|
|
// http://www.apache.org/licenses/LICENSE-2.0
|
|
//
|
|
// Unless required by applicable law or agreed to in writing, software
|
|
// distributed under the License is distributed on an "AS IS" BASIS,
|
|
// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
|
|
// See the License for the specific language governing permissions and
|
|
// limitations under the License.
|
|
|
|
// Package imagecache implements the node-local OCI image cache: a
|
|
// content-addressed pool of unpacked image layers shared by every actor on
|
|
// the node, plus the per-bundle overlay spec that tells the ateom runtimes
|
|
// how to compose an actor rootfs from cached layers.
|
|
//
|
|
// The work is split along the existing atelet/ateom privilege boundary:
|
|
//
|
|
// - atelet (plain root, all capabilities dropped) pulls layers and unpacks
|
|
// them into the pool (Store.EnsureImage), and writes a rootfs-overlay.json
|
|
// next to each bundle's config.json (WriteSpec). Whiteout entries are
|
|
// recorded in per-layer metadata rather than materialized, because
|
|
// overlayfs whiteouts are char devices (CAP_MKNOD) with trusted.* xattrs
|
|
// for opaque dirs (CAP_SYS_ADMIN).
|
|
// - ateom (privileged; it already owns every mount on the node) finalizes
|
|
// layers — materializing the recorded whiteout state, once per layer —
|
|
// and mounts the overlay rootfs (SetupBundleRootfs) just before
|
|
// `runsc create` / staging the micro-VM virtio-fs lower.
|
|
//
|
|
// On-disk layout under the cache root (a directory on the BasePath hostPath,
|
|
// so the same absolute paths resolve in atelet and every ateom pod):
|
|
//
|
|
// version layout version marker
|
|
// layers/sha256/<diffid-hex>/
|
|
// fs/ the unpacked layer tree (overlay lowerdir)
|
|
// whiteouts.json whiteout state recorded at unpack time
|
|
// finalized marker written by FinalizeLayer (ateom)
|
|
// manifests/sha256/<digest-hex>.json
|
|
// image config + ordered diffID list
|
|
//
|
|
// Layers land in the pool via unpack-into-tempdir + atomic rename, so a
|
|
// layer directory that exists is always complete; startup recovery only has
|
|
// to sweep orphaned temp dirs.
|
|
package imagecache
|
|
|
|
import (
|
|
"context"
|
|
"encoding/json"
|
|
"errors"
|
|
"fmt"
|
|
"log/slog"
|
|
"net"
|
|
"os"
|
|
"path/filepath"
|
|
"runtime"
|
|
"strconv"
|
|
"strings"
|
|
"sync"
|
|
"time"
|
|
|
|
"github.com/google/go-containerregistry/pkg/authn"
|
|
"github.com/google/go-containerregistry/pkg/name"
|
|
v1 "github.com/google/go-containerregistry/pkg/v1"
|
|
"github.com/google/go-containerregistry/pkg/v1/remote"
|
|
"go.opentelemetry.io/otel/metric"
|
|
"golang.org/x/sync/errgroup"
|
|
"golang.org/x/sync/singleflight"
|
|
|
|
"github.com/agent-substrate/substrate/internal/ateattr"
|
|
)
|
|
|
|
const (
|
|
layoutVersion = "1"
|
|
versionFileName = "version"
|
|
|
|
layerFSDirName = "fs"
|
|
layerWhiteoutsFileName = "whiteouts.json"
|
|
layerFinalizedMarkerName = "finalized"
|
|
// layerSizeFileName holds the layer's byte count, recorded at unpack so
|
|
// sizing the pool never walks a tree. Absent for layers unpacked by
|
|
// older atelets (backfilled lazily, see layerSize). An estimate: the
|
|
// value is the tar-stream length when written at unpack, or the summed
|
|
// file sizes when backfilled, so it may differ from disk usage and
|
|
// across nodes.
|
|
layerSizeFileName = "size"
|
|
|
|
// defaultMinAge is the default eviction minimum age (see WithMinAge).
|
|
defaultMinAge = 2 * time.Minute
|
|
|
|
// layerPullConcurrency bounds concurrent layer download+unpack streams per
|
|
// image pull. Memory use is O(stream buffers) per slot, independent of
|
|
// layer size.
|
|
layerPullConcurrency = 4
|
|
|
|
// defaultPullTimeout is the default per-pull bound (see WithPullTimeout).
|
|
// Generous enough for multi-GiB images on a busy node.
|
|
defaultPullTimeout = 10 * time.Minute
|
|
)
|
|
|
|
// Pull retry backoffs, chosen per registry by retryBackoffFor. Both replace
|
|
// go-containerregistry's default (3 attempts over ~4s, jitter 0.1), which is
|
|
// tuned for a single flaky client, and both apply to every request of a pull,
|
|
// including the /v2/ auth ping. The retryable status codes stay at the
|
|
// library default, which already includes 429. Vars so tests can shrink the
|
|
// waits.
|
|
//
|
|
// The backoff applies per request, so under a sustained throttle a
|
|
// multi-layer pull can accumulate far more wall clock than one request's
|
|
// worst case (~14s shared). The caller's ctx still bounds the total: every
|
|
// request carries it (remote.WithContext in remoteOpts), a canceled ctx
|
|
// fails the next attempt immediately and is never retried, so retrying
|
|
// overshoots a deadline by at most one backoff sleep. Production pulls run
|
|
// under the Run/Restore RPC ctx, which ateapi's resume path caps at its
|
|
// server-wide max RPC deadline — a throttled shared-registry pull surfaces
|
|
// as that RPC's deadline error, not a hang.
|
|
var (
|
|
// dedicatedRegistryBackoff covers registries where the deployment has
|
|
// its own quota (Artifact Registry, ECR, Harbor, self-hosted, …): the
|
|
// server absorbs retry bursts and nobody else competes for the limit,
|
|
// and pulls sit on the Run/Restore critical path, so retries start fast
|
|
// and come often — more attempts in less total wall clock, still quick
|
|
// enough to surface a real outage to the RPC-level retry.
|
|
dedicatedRegistryBackoff = remote.Backoff{
|
|
Duration: 200 * time.Millisecond,
|
|
Factor: 2.0,
|
|
Jitter: 1.0,
|
|
Steps: 4,
|
|
Cap: 2 * time.Second,
|
|
}
|
|
// sharedRegistryBackoff covers communal registries (registry.k8s.io,
|
|
// Docker Hub, …), where pulls are anonymous and rate limits are shared:
|
|
// a whole fleet can be throttled at once, so retries must outlast the
|
|
// throttling wave and full jitter must de-synchronize the herd rather
|
|
// than replay it.
|
|
sharedRegistryBackoff = remote.Backoff{
|
|
Duration: 1 * time.Second,
|
|
Factor: 2.0,
|
|
Jitter: 1.0,
|
|
Steps: 4,
|
|
Cap: 10 * time.Second,
|
|
}
|
|
)
|
|
|
|
// sharedRegistries are the well-known communal registries whose rate limits
|
|
// are shared across all anonymous clients. Everything not listed here is
|
|
// assumed dedicated: fleet-scale pulls realistically hit either the pause
|
|
// image (registry.k8s.io) or the customer's own registry, and enumerating
|
|
// the small stable set of communal hosts beats guessing at every private
|
|
// registry vendor.
|
|
var sharedRegistries = map[string]bool{
|
|
"registry.k8s.io": true,
|
|
// Pulls only ever present index.docker.io here: name.ParseReference
|
|
// normalizes docker.io (and bare refs like "ubuntu") to it before
|
|
// retryBackoffFor runs. The docker.io entry is kept so a lookup by the
|
|
// canonical name classifies the same way.
|
|
"docker.io": true,
|
|
"index.docker.io": true,
|
|
"registry-1.docker.io": true,
|
|
"quay.io": true,
|
|
"ghcr.io": true,
|
|
"public.ecr.aws": true,
|
|
"mcr.microsoft.com": true,
|
|
"registry.gitlab.com": true,
|
|
"cgr.dev": true,
|
|
}
|
|
|
|
// retryBackoffFor picks the pull retry backoff for a registry host.
|
|
func retryBackoffFor(registry string) remote.Backoff {
|
|
if sharedRegistries[registry] {
|
|
return sharedRegistryBackoff
|
|
}
|
|
return dedicatedRegistryBackoff
|
|
}
|
|
|
|
// Store is atelet's handle to the on-disk layer pool. It is safe for
|
|
// concurrent use; concurrent pulls of the same image or layer are collapsed.
|
|
// The store assumes it is the only writer on the node (one atelet per node).
|
|
type Store struct {
|
|
root string
|
|
|
|
// authenticator, when set, is attached to pulls from registries that use
|
|
// GCP credentials (gcr.io / pkg.dev). See remoteOpts.
|
|
authenticator authn.Authenticator
|
|
|
|
localhostRegistryReplacement string
|
|
|
|
// platform overrides the default pull platform (linux/GOARCH), for
|
|
// callers pulling on a different architecture than the images' target.
|
|
platform *v1.Platform
|
|
|
|
// actorsDir is scanned by InUse for bundle overlay specs; empty disables
|
|
// the scan (the root set is then empty).
|
|
actorsDir string
|
|
|
|
// minAge vetoes eviction of any layer or image record younger than this,
|
|
// covering the window between a pull (or cache-hit stat) and the bundle
|
|
// spec write / ateom mount that roots it.
|
|
minAge time.Duration
|
|
|
|
// pullTimeout bounds each pull. Pulls run detached from the contexts of
|
|
// the callers waiting on them (see EnsureImage), so this is the only
|
|
// bound on how long one can run.
|
|
pullTimeout time.Duration
|
|
|
|
// meter, when set, is the meter the store reports on. See WithMeter.
|
|
meter metric.Meter
|
|
|
|
// requests counts EnsureImage lookups by outcome. Nil without a meter,
|
|
// which recordRequest treats as a no-op.
|
|
requests metric.Int64Counter
|
|
|
|
imageSF singleflight.Group
|
|
layerSF singleflight.Group
|
|
|
|
// evictMu serializes EvictUnused passes (concurrent passes would fight
|
|
// over the same candidates for no benefit).
|
|
evictMu sync.Mutex
|
|
|
|
// hitMu closes the hit-vs-evict window: held shared by the hit path
|
|
// (cachedImageHit), exclusive by eviction's record removal
|
|
// (removeStaleRecord), so a hit's last-use touch and eviction's final
|
|
// re-check can never interleave. Uncontended except during a pass.
|
|
hitMu sync.RWMutex
|
|
}
|
|
|
|
// Option configures a Store.
|
|
type Option func(*Store)
|
|
|
|
// WithAuthenticator attaches an authenticator used for gcr.io / pkg.dev
|
|
// registries. A nil authenticator is ignored.
|
|
func WithAuthenticator(a authn.Authenticator) Option {
|
|
return func(s *Store) { s.authenticator = a }
|
|
}
|
|
|
|
// WithLocalhostRegistryReplacement rewrites localhost/loopback registry refs
|
|
// to the given endpoint, mirroring the containerd mirror config used by kind
|
|
// local registries (https://kind.sigs.k8s.io/docs/user/local-registry/).
|
|
func WithLocalhostRegistryReplacement(replacement string) Option {
|
|
return func(s *Store) { s.localhostRegistryReplacement = replacement }
|
|
}
|
|
|
|
// WithPlatform overrides the pull platform (default: linux/GOARCH).
|
|
func WithPlatform(p v1.Platform) Option {
|
|
return func(s *Store) { s.platform = &p }
|
|
}
|
|
|
|
// WithActorsDir points the eviction root-set scan at the node's actors
|
|
// directory (the per-actor state dirs under nodepath.BasePath). Each
|
|
// <actorsDir>/<actorUID>/bundles/<container>/rootfs-overlay.json roots its
|
|
// image and layers against eviction. Empty disables the scan.
|
|
func WithActorsDir(dir string) Option {
|
|
return func(s *Store) { s.actorsDir = dir }
|
|
}
|
|
|
|
// WithMinAge overrides the eviction minimum age (default 2m): layers and
|
|
// image records younger than this are never evicted.
|
|
func WithMinAge(d time.Duration) Option {
|
|
return func(s *Store) { s.minAge = d }
|
|
}
|
|
|
|
// WithPullTimeout overrides the per-pull timeout (default 10m).
|
|
func WithPullTimeout(d time.Duration) Option {
|
|
return func(s *Store) { s.pullTimeout = d }
|
|
}
|
|
|
|
// WithMeter attaches the meter the store reports ate.imagecache.requests on.
|
|
// Without it the store records nothing, so a caller with no metrics pipeline
|
|
// needs no meter provider.
|
|
func WithMeter(m metric.Meter) Option {
|
|
return func(s *Store) { s.meter = m }
|
|
}
|
|
|
|
// Image describes one cached, ready-to-compose image.
|
|
type Image struct {
|
|
// Digest is the manifest digest the caller's ref resolved to (for a
|
|
// multi-arch ref, the index digest as requested, not the per-platform
|
|
// child).
|
|
Digest v1.Hash
|
|
// Config is the OCI image config (entrypoint, env, ...).
|
|
Config v1.Config
|
|
// LayerDirs are the absolute cached layer directories, bottom-most layer
|
|
// first. Each contains the unpacked tree under "fs/".
|
|
LayerDirs []string
|
|
}
|
|
|
|
// imageRecord is the persisted form of a cached image, stored under
|
|
// manifests/<algorithm>/<hex>.json.
|
|
type imageRecord struct {
|
|
Version int `json:"version"`
|
|
Config v1.Config `json:"config"`
|
|
DiffIDs []string `json:"diffIDs"`
|
|
}
|
|
|
|
// New opens (creating if needed) the layer pool rooted at root and runs
|
|
// startup recovery: verifying the layout version and sweeping temp dirs left
|
|
// by unpacks that were in flight when a previous atelet died.
|
|
func New(root string, opts ...Option) (*Store, error) {
|
|
s := &Store{root: root, minAge: defaultMinAge, pullTimeout: defaultPullTimeout}
|
|
for _, o := range opts {
|
|
o(s)
|
|
}
|
|
|
|
if s.meter != nil {
|
|
requests, err := newRequestsCounter(s.meter)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
s.requests = requests
|
|
}
|
|
|
|
for _, d := range []string{s.layersDir(), s.manifestsDir()} {
|
|
if err := os.MkdirAll(d, 0o700); err != nil {
|
|
return nil, fmt.Errorf("while creating image cache dir %q: %w", d, err)
|
|
}
|
|
}
|
|
|
|
versionPath := filepath.Join(root, versionFileName)
|
|
switch b, err := os.ReadFile(versionPath); {
|
|
case err == nil:
|
|
if got := strings.TrimSpace(string(b)); got != layoutVersion {
|
|
// Fail loudly instead of silently mixing layouts; an operator can
|
|
// delete the cache dir to rebuild it (it holds no unique state).
|
|
return nil, fmt.Errorf("image cache at %q has layout version %q, this atelet supports %q", root, got, layoutVersion)
|
|
}
|
|
case errors.Is(err, os.ErrNotExist):
|
|
if err := os.WriteFile(versionPath, []byte(layoutVersion+"\n"), 0o600); err != nil {
|
|
return nil, fmt.Errorf("while writing image cache version marker: %w", err)
|
|
}
|
|
default:
|
|
return nil, fmt.Errorf("while reading image cache version marker: %w", err)
|
|
}
|
|
|
|
if err := s.sweepTempDirs(); err != nil {
|
|
return nil, err
|
|
}
|
|
// Startup-only orphan recovery (see RecoverOrphans for why it must not
|
|
// run during normal operation). Never fatal: a corrupt record must not
|
|
// keep atelet from serving actors. Gated means nothing was attempted
|
|
// (orphans persist until the named file is repaired and the next
|
|
// restart); per-item errors mean the scan ran and reclaimed what it
|
|
// could.
|
|
if _, err := s.RecoverOrphans(context.Background()); err != nil {
|
|
if errors.Is(err, ErrIncompleteEnumeration) {
|
|
slog.Error("Image cache startup orphan scan skipped", slog.Any("err", err))
|
|
} else {
|
|
slog.Warn("Image cache startup orphan recovery incomplete", slog.Any("err", err))
|
|
}
|
|
}
|
|
return s, nil
|
|
}
|
|
|
|
func (s *Store) layersDir() string { return filepath.Join(s.root, "layers", "sha256") }
|
|
func (s *Store) manifestsDir() string { return filepath.Join(s.root, "manifests", "sha256") }
|
|
|
|
func (s *Store) layerDir(diffID v1.Hash) string {
|
|
return filepath.Join(s.root, "layers", diffID.Algorithm, diffID.Hex)
|
|
}
|
|
|
|
func (s *Store) recordPath(digest v1.Hash) string {
|
|
return filepath.Join(s.root, "manifests", digest.Algorithm, digest.Hex+".json")
|
|
}
|
|
|
|
// sweepTempDirs removes unpack temp dirs, retired layer dirs, and
|
|
// manifest-record temp files orphaned by a crash. A layer dir without the
|
|
// temp/retired prefix and a record without a leading dot are always
|
|
// complete (both are moved into place with a single rename), so this is
|
|
// the only recovery the pool needs.
|
|
func (s *Store) sweepTempDirs() error {
|
|
entries, err := os.ReadDir(s.layersDir())
|
|
if err != nil {
|
|
return fmt.Errorf("while listing layer pool: %w", err)
|
|
}
|
|
swept := 0
|
|
for _, e := range entries {
|
|
// ".tmp-": an unpack in flight at crash. ".rm-": a layer retirement
|
|
// renamed aside but not yet removed. Either way, unreferenced
|
|
// garbage.
|
|
if !strings.HasPrefix(e.Name(), ".tmp-") && !strings.HasPrefix(e.Name(), retiredPrefix) {
|
|
continue
|
|
}
|
|
p := filepath.Join(s.layersDir(), e.Name())
|
|
slog.Info("Image cache sweeping orphaned layer dir", slog.String("dir", e.Name()))
|
|
if err := RemoveAllWritable(p); err != nil {
|
|
return fmt.Errorf("while sweeping orphaned layer dir %q: %w", p, err)
|
|
}
|
|
swept++
|
|
}
|
|
if swept > 0 {
|
|
slog.Info("Image cache startup sweep removed orphaned layer dirs", slog.Int("count", swept))
|
|
}
|
|
|
|
// writeRecord's temp files are ".<hex>.json.tmp-<rand>"; finished records
|
|
// are "<hex>.json", so a leading dot alone identifies an orphan.
|
|
records, err := os.ReadDir(s.manifestsDir())
|
|
if err != nil {
|
|
return fmt.Errorf("while listing manifest records: %w", err)
|
|
}
|
|
for _, e := range records {
|
|
if !strings.HasPrefix(e.Name(), ".") {
|
|
continue
|
|
}
|
|
p := filepath.Join(s.manifestsDir(), e.Name())
|
|
if err := os.Remove(p); err != nil {
|
|
return fmt.Errorf("while sweeping orphaned manifest temp file %q: %w", p, err)
|
|
}
|
|
}
|
|
return nil
|
|
}
|
|
|
|
// EnsureImage makes ref's image available in the pool and returns its config
|
|
// and ordered layer directories. Digest refs hit the cache with no network
|
|
// I/O; tag refs cost one HEAD request to resolve the tag to a manifest
|
|
// digest (so tag refs are cacheable, and a moved tag is picked up on the
|
|
// next call).
|
|
func (s *Store) EnsureImage(ctx context.Context, ref string) (_ *Image, err error) {
|
|
// A miss until a complete record proves otherwise; recordRequest
|
|
// reclassifies a failure onto its own outcome.
|
|
outcome := ateattr.ImageCacheOutcomeMiss
|
|
defer func() { s.recordRequest(ctx, outcome, err) }()
|
|
|
|
parsedRef, err := s.parseRef(ref)
|
|
if err != nil {
|
|
return nil, fmt.Errorf("while parsing reference: %w", err)
|
|
}
|
|
|
|
var digest v1.Hash
|
|
if d, ok := parsedRef.(name.Digest); ok {
|
|
digest, err = v1.NewHash(d.DigestStr())
|
|
if err != nil {
|
|
return nil, fmt.Errorf("while parsing digest of %q: %w", ref, err)
|
|
}
|
|
} else {
|
|
// Tag ref: one small HEAD request pins it to an immutable manifest
|
|
// digest, which is the only safe cache key for mutable tags.
|
|
desc, headErr := remote.Head(parsedRef, s.remoteOpts(ctx, parsedRef)...)
|
|
if headErr != nil {
|
|
err = fmt.Errorf("while resolving tag %q to a digest: %w", ref, headErr)
|
|
return nil, err
|
|
}
|
|
digest = desc.Digest
|
|
}
|
|
|
|
img, err := s.cachedImageHit(digest)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
if img != nil {
|
|
outcome = ateattr.ImageCacheOutcomeHit
|
|
slog.InfoContext(ctx, "Image cache hit", slog.String("ref", ref), slog.String("digest", digest.String()))
|
|
return img, nil
|
|
}
|
|
slog.InfoContext(ctx, "Image cache miss", slog.String("ref", ref), slog.String("digest", digest.String()))
|
|
|
|
// Collapse concurrent pulls of the same digest (e.g. several containers of
|
|
// one actor, or several actors landing at once). Callers with very
|
|
// different lifetimes share these flights — a serving Restore, a
|
|
// best-effort prewarm bounded by its own deadline, a caller whose daemon
|
|
// is draining — so the pull runs on a context detached from whichever
|
|
// caller happened to start the flight, bounded only by the store's pull
|
|
// timeout: a waiter must only ever see a real pull failure, never another
|
|
// caller's cancellation. Each caller stops waiting when its own ctx ends,
|
|
// while the pull runs on to warm the cache for the next attempt.
|
|
ch := s.imageSF.DoChan(digest.String(), func() (any, error) {
|
|
pullCtx, cancel := context.WithTimeout(context.WithoutCancel(ctx), s.pullTimeout)
|
|
defer cancel()
|
|
return s.pull(pullCtx, parsedRef, digest)
|
|
})
|
|
select {
|
|
case res := <-ch:
|
|
if res.Err != nil {
|
|
return nil, res.Err
|
|
}
|
|
return res.Val.(*Image), nil
|
|
case <-ctx.Done():
|
|
return nil, fmt.Errorf("while waiting for pull of %s: %w", digest, context.Cause(ctx))
|
|
}
|
|
}
|
|
|
|
// cachedImageHit is the hit side of the hitMu contract: it verifies the
|
|
// cached image and records last-use for eviction's LRU ordering, atomic
|
|
// with respect to eviction's record removal (removeStaleRecord holds
|
|
// hitMu exclusive). Refreshing the mtime also renews the min-age veto, so
|
|
// an image in active use cannot age into eviction between this stat and
|
|
// the ateom's mount.
|
|
func (s *Store) cachedImageHit(digest v1.Hash) (*Image, error) {
|
|
s.hitMu.RLock()
|
|
defer s.hitMu.RUnlock()
|
|
img, err := s.cachedImage(digest)
|
|
if err == nil && img != nil {
|
|
s.touchRecord(digest)
|
|
}
|
|
return img, err
|
|
}
|
|
|
|
// cachedImage returns the cached image for digest, or nil if the record or
|
|
// any of its layer dirs is missing (in which case the caller re-pulls; only
|
|
// the missing layers cost anything).
|
|
func (s *Store) cachedImage(digest v1.Hash) (*Image, error) {
|
|
b, err := os.ReadFile(s.recordPath(digest))
|
|
if errors.Is(err, os.ErrNotExist) {
|
|
return nil, nil
|
|
} else if err != nil {
|
|
return nil, fmt.Errorf("while reading image record for %s: %w", digest, err)
|
|
}
|
|
var rec imageRecord
|
|
if err := json.Unmarshal(b, &rec); err != nil {
|
|
return nil, fmt.Errorf("while decoding image record for %s: %w", digest, err)
|
|
}
|
|
|
|
layerDirs := make([]string, len(rec.DiffIDs))
|
|
for i, d := range rec.DiffIDs {
|
|
diffID, err := v1.NewHash(d)
|
|
if err != nil {
|
|
return nil, fmt.Errorf("invalid diffID %q in image record for %s: %w", d, digest, err)
|
|
}
|
|
dir := s.layerDir(diffID)
|
|
if _, err := os.Stat(filepath.Join(dir, layerFSDirName)); err != nil {
|
|
return nil, nil
|
|
}
|
|
layerDirs[i] = dir
|
|
}
|
|
return &Image{Digest: digest, Config: rec.Config, LayerDirs: layerDirs}, nil
|
|
}
|
|
|
|
// pull fetches the image (by its resolved digest, so what is unpacked is
|
|
// exactly what was recorded), writes the image record, and then unpacks
|
|
// every missing layer into the pool.
|
|
//
|
|
// The record comes first so that every layer is referenced — and thereby
|
|
// safe from eviction — before it can exist on disk. A record therefore
|
|
// means "known image, possibly partially present", which readers already
|
|
// handle: cachedImage verifies every layer and re-pulls what is missing,
|
|
// so an interrupted pull's record is just resumable progress.
|
|
func (s *Store) pull(ctx context.Context, parsedRef name.Reference, digest v1.Hash) (*Image, error) {
|
|
// Re-check under the flight lock: a racing EnsureImage may have completed
|
|
// the pull between our cache miss and winning the singleflight slot.
|
|
if img, err := s.cachedImage(digest); err != nil {
|
|
return nil, err
|
|
} else if img != nil {
|
|
return img, nil
|
|
}
|
|
|
|
tStart := time.Now()
|
|
digestRef := parsedRef.Context().Digest(digest.String())
|
|
img, err := remote.Image(digestRef, s.remoteOpts(ctx, parsedRef)...)
|
|
if err != nil {
|
|
return nil, fmt.Errorf("in remote.Image: %w", err)
|
|
}
|
|
|
|
cfgFile, err := img.ConfigFile()
|
|
if err != nil {
|
|
return nil, fmt.Errorf("while reading image config: %w", err)
|
|
}
|
|
|
|
layers, err := img.Layers()
|
|
if err != nil {
|
|
return nil, fmt.Errorf("while listing image layers: %w", err)
|
|
}
|
|
|
|
if len(cfgFile.RootFS.DiffIDs) != len(layers) {
|
|
return nil, fmt.Errorf("image %s config lists %d diffIDs but manifest has %d layers", digest, len(cfgFile.RootFS.DiffIDs), len(layers))
|
|
}
|
|
diffIDs := make([]string, len(layers))
|
|
for i, d := range cfgFile.RootFS.DiffIDs {
|
|
diffIDs[i] = d.String()
|
|
}
|
|
|
|
// The full diffID list is known from the config before any unpack; make
|
|
// every layer of this pull referenced before it can exist.
|
|
rec := imageRecord{Version: 1, Config: cfgFile.Config, DiffIDs: diffIDs}
|
|
if err := s.writeRecord(digest, rec); err != nil {
|
|
return nil, err
|
|
}
|
|
// For a multi-arch ref the requested digest is the index digest, but the
|
|
// layers unpacked belong to the per-platform child manifest. Record the
|
|
// image under the child digest too, so refs pinned either way hit.
|
|
var actualDigest *v1.Hash
|
|
if actual, err := img.Digest(); err == nil && actual != digest {
|
|
actualDigest = &actual
|
|
if err := s.writeRecord(actual, rec); err != nil {
|
|
slog.WarnContext(ctx, "Failed to record image under platform manifest digest",
|
|
slog.String("digest", actual.String()), slog.Any("err", err))
|
|
}
|
|
}
|
|
|
|
layerDirs := make([]string, len(layers))
|
|
g, gctx := errgroup.WithContext(ctx)
|
|
g.SetLimit(layerPullConcurrency)
|
|
for i, layer := range layers {
|
|
g.Go(func() error {
|
|
diffID, err := layer.DiffID()
|
|
if err != nil {
|
|
return fmt.Errorf("while reading layer diffID: %w", err)
|
|
}
|
|
if diffID.String() != diffIDs[i] {
|
|
// The record must reference exactly what lands on disk.
|
|
return fmt.Errorf("layer %d diffID %s does not match config rootfs diffID %s", i, diffID, diffIDs[i])
|
|
}
|
|
dir, err := s.ensureLayer(gctx, diffID, layer)
|
|
if err != nil {
|
|
return fmt.Errorf("while unpacking layer %s: %w", diffID, err)
|
|
}
|
|
layerDirs[i] = dir
|
|
// Each completed layer refreshes the record's mtime: a pull
|
|
// making progress stays fresh indefinitely; a wedged one ages
|
|
// into ordinary LRU eviction. The twin is not touched — the
|
|
// primary keeps the layers referenced, and the final rewrite
|
|
// recreates the twin if it ages out mid-pull.
|
|
s.touchRecord(digest)
|
|
return nil
|
|
})
|
|
}
|
|
if err := g.Wait(); err != nil {
|
|
return nil, err
|
|
}
|
|
|
|
// Rewrite the record: eviction may legitimately remove it mid-pull
|
|
// (no progress for min-age reads as wedged), and success must never
|
|
// leave the just-unpacked layers unreferenced.
|
|
if err := s.writeRecord(digest, rec); err != nil {
|
|
return nil, fmt.Errorf("while rewriting image record after unpack: %w", err)
|
|
}
|
|
if actualDigest != nil {
|
|
if err := s.writeRecord(*actualDigest, rec); err != nil {
|
|
slog.WarnContext(ctx, "Failed to rewrite platform manifest record after unpack",
|
|
slog.String("digest", actualDigest.String()), slog.Any("err", err))
|
|
}
|
|
}
|
|
|
|
// Never return LayerDirs that are not on disk right now: a vanished
|
|
// dir fails the pull into a clean RPC retry instead of a bundle spec
|
|
// naming a missing lowerdir. Ordered after the rewrite so even this
|
|
// failure path leaves the surviving layers referenced.
|
|
for _, dir := range layerDirs {
|
|
if _, err := os.Stat(filepath.Join(dir, layerFSDirName)); err != nil {
|
|
return nil, fmt.Errorf("layer dir vanished during pull (evicted?): %w", err)
|
|
}
|
|
}
|
|
|
|
slog.InfoContext(ctx, "Image pulled into layer cache",
|
|
slog.String("digest", digest.String()),
|
|
slog.Int("layers", len(layers)),
|
|
slog.Duration("took", time.Since(tStart)))
|
|
|
|
return &Image{Digest: digest, Config: cfgFile.Config, LayerDirs: layerDirs}, nil
|
|
}
|
|
|
|
// ensureLayer makes the unpacked tree for diffID present in the pool,
|
|
// collapsing concurrent requests for the same layer across images.
|
|
func (s *Store) ensureLayer(ctx context.Context, diffID v1.Hash, layer v1.Layer) (string, error) {
|
|
dir := s.layerDir(diffID)
|
|
_, err, _ := s.layerSF.Do(layerFlightKey(diffID.Hex), func() (any, error) {
|
|
if _, err := os.Stat(filepath.Join(dir, layerFSDirName)); err == nil {
|
|
// Refresh the dir mtime inside the flight: retireLayer re-checks
|
|
// the mtime in this same flight, so a layer reused here can
|
|
// never be renamed away between this stat and the image record
|
|
// that will re-reference it.
|
|
now := time.Now()
|
|
if err := os.Chtimes(dir, now, now); err != nil {
|
|
slog.WarnContext(ctx, "Failed to refresh layer mtime on reuse", slog.String("diffid", diffID.String()), slog.Any("err", err))
|
|
}
|
|
return nil, nil
|
|
}
|
|
return nil, s.unpackLayerToPool(ctx, diffID, layer)
|
|
})
|
|
if err != nil {
|
|
return "", err
|
|
}
|
|
return dir, nil
|
|
}
|
|
|
|
// unpackLayerToPool streams the layer (download → decompress → untar) into a
|
|
// temp dir and renames it into place, so a layer dir either exists complete
|
|
// or not at all.
|
|
func (s *Store) unpackLayerToPool(ctx context.Context, diffID v1.Hash, layer v1.Layer) (retErr error) {
|
|
tmp, err := os.MkdirTemp(filepath.Dir(s.layerDir(diffID)), ".tmp-"+diffID.Hex[:12]+"-")
|
|
if err != nil {
|
|
return fmt.Errorf("while creating layer temp dir: %w", err)
|
|
}
|
|
defer func() {
|
|
// No-op once the rename has moved tmp into place.
|
|
if _, err := os.Stat(tmp); err == nil {
|
|
if rmErr := RemoveAllWritable(tmp); rmErr != nil && retErr == nil {
|
|
retErr = fmt.Errorf("while cleaning up layer temp dir: %w", rmErr)
|
|
}
|
|
}
|
|
}()
|
|
|
|
fsDir := filepath.Join(tmp, layerFSDirName)
|
|
if err := os.Mkdir(fsDir, 0o755); err != nil {
|
|
return fmt.Errorf("while creating layer fs dir: %w", err)
|
|
}
|
|
|
|
rc, err := layer.Uncompressed()
|
|
if err != nil {
|
|
return fmt.Errorf("while opening layer stream: %w", err)
|
|
}
|
|
defer rc.Close()
|
|
|
|
root, err := os.OpenRoot(fsDir)
|
|
if err != nil {
|
|
return fmt.Errorf("while opening layer fs dir as os.Root: %w", err)
|
|
}
|
|
defer root.Close()
|
|
|
|
// The uncompressed tar stream is the recorded size: an optimistic
|
|
// estimate (tar framing vs. block rounding), cheap to capture here.
|
|
cr := &countingReader{r: rc}
|
|
wh, err := unpackLayer(ctx, cr, root)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
// Non-fatal: a missing size file is recovered by layerSize's backfill,
|
|
// so a metadata write must not discard a successful unpack.
|
|
if err := os.WriteFile(filepath.Join(tmp, layerSizeFileName), []byte(strconv.FormatInt(cr.n, 10)+"\n"), 0o600); err != nil {
|
|
slog.WarnContext(ctx, "Failed to record layer size; will backfill lazily",
|
|
slog.String("diffid", diffID.String()), slog.Any("err", err))
|
|
}
|
|
|
|
whBytes, err := json.Marshal(wh)
|
|
if err != nil {
|
|
return fmt.Errorf("while encoding layer whiteouts: %w", err)
|
|
}
|
|
if err := os.WriteFile(filepath.Join(tmp, layerWhiteoutsFileName), whBytes, 0o600); err != nil {
|
|
return fmt.Errorf("while writing layer whiteouts: %w", err)
|
|
}
|
|
|
|
if err := os.Rename(tmp, s.layerDir(diffID)); err != nil {
|
|
// A concurrent unpack (another process sharing the pool) may have won;
|
|
// its layer is as good as ours.
|
|
if _, statErr := os.Stat(filepath.Join(s.layerDir(diffID), layerFSDirName)); statErr == nil {
|
|
return nil
|
|
}
|
|
return fmt.Errorf("while moving layer into pool: %w", err)
|
|
}
|
|
return nil
|
|
}
|
|
|
|
func (s *Store) writeRecord(digest v1.Hash, rec imageRecord) error {
|
|
b, err := json.MarshalIndent(rec, "", " ")
|
|
if err != nil {
|
|
return fmt.Errorf("while encoding image record: %w", err)
|
|
}
|
|
path := s.recordPath(digest)
|
|
tmp, err := os.CreateTemp(filepath.Dir(path), "."+filepath.Base(path)+".tmp-*")
|
|
if err != nil {
|
|
return fmt.Errorf("while creating image record temp file: %w", err)
|
|
}
|
|
defer os.Remove(tmp.Name()) // no-op once the rename succeeds
|
|
if _, err := tmp.Write(b); err != nil {
|
|
tmp.Close()
|
|
return fmt.Errorf("while writing image record: %w", err)
|
|
}
|
|
if err := tmp.Close(); err != nil {
|
|
return fmt.Errorf("while closing image record: %w", err)
|
|
}
|
|
if err := os.Rename(tmp.Name(), path); err != nil {
|
|
return fmt.Errorf("while moving image record into place: %w", err)
|
|
}
|
|
return nil
|
|
}
|
|
|
|
// remoteOpts assembles the go-containerregistry options for pulls from
|
|
// parsedRef's registry.
|
|
func (s *Store) remoteOpts(ctx context.Context, parsedRef name.Reference) []remote.Option {
|
|
platform := v1.Platform{
|
|
Architecture: runtime.GOARCH,
|
|
OS: "linux",
|
|
}
|
|
if s.platform != nil {
|
|
platform = *s.platform
|
|
}
|
|
registry := parsedRef.Context().Registry.RegistryStr()
|
|
opts := []remote.Option{
|
|
// Propagate caller ctx into go-containerregistry so cancellation tears
|
|
// down in-flight layer-blob HTTP requests instead of letting them run
|
|
// to completion in background goroutines.
|
|
remote.WithContext(ctx),
|
|
remote.WithPlatform(platform),
|
|
remote.WithRetryBackoff(retryBackoffFor(registry)),
|
|
}
|
|
if s.authenticator != nil && registryUsesGCPAuth(registry) {
|
|
opts = append(opts, remote.WithAuth(s.authenticator))
|
|
}
|
|
return opts
|
|
}
|
|
|
|
func registryUsesGCPAuth(registry string) bool {
|
|
return registry == "gcr.io" || strings.HasSuffix(registry, ".gcr.io") ||
|
|
registry == "pkg.dev" || strings.HasSuffix(registry, ".pkg.dev")
|
|
}
|
|
|
|
// parseRef applies the localhost-registry rewrite (kind local registries) and
|
|
// permits plain-HTTP pulls for localhost/loopback registries, matching docker
|
|
// behavior so local development needs no TLS certs.
|
|
func (s *Store) parseRef(ref string) (name.Reference, error) {
|
|
rewritten := false
|
|
if s.localhostRegistryReplacement != "" {
|
|
if newRef := s.rewriteLocalRegistry(ref); newRef != ref {
|
|
ref = newRef
|
|
rewritten = true
|
|
}
|
|
}
|
|
var nameOpts []name.Option
|
|
if rewritten || isLocalRegistry(ref) {
|
|
nameOpts = append(nameOpts, name.Insecure)
|
|
}
|
|
return name.ParseReference(ref, nameOpts...)
|
|
}
|
|
|
|
func registryHost(ref string) string {
|
|
parts := strings.SplitN(ref, "/", 2)
|
|
reg, err := name.NewRegistry(parts[0], name.Insecure)
|
|
if err != nil {
|
|
return ""
|
|
}
|
|
hostPart := reg.Name()
|
|
if h, _, err := net.SplitHostPort(hostPart); err == nil {
|
|
return h
|
|
}
|
|
return hostPart
|
|
}
|
|
|
|
func isLocalhostOrLoopback(host string) bool {
|
|
if host == "localhost" {
|
|
return true
|
|
}
|
|
if ip := net.ParseIP(host); ip != nil && ip.IsLoopback() {
|
|
return true
|
|
}
|
|
return false
|
|
}
|
|
|
|
func isLocalRegistry(ref string) bool {
|
|
// By default docker permits localhost and 127.0.0.0/8; we also permit the
|
|
// IPv6 loopback here.
|
|
return isLocalhostOrLoopback(registryHost(ref))
|
|
}
|
|
|
|
func (s *Store) rewriteLocalRegistry(ref string) string {
|
|
if isLocalRegistry(ref) {
|
|
parts := strings.SplitN(ref, "/", 2)
|
|
return s.localhostRegistryReplacement + "/" + parts[1]
|
|
}
|
|
return ref
|
|
}
|