Something went wrong. Try again.
Monorepo for Tangled tangled.org
Something went wrong. Try again.
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581package engine
import ( "context" "crypto/sha256" "database/sql" "encoding/hex" "errors" "fmt" "hash" "io" "log/slog" "net/url" "os/exec" "path/filepath" "strings" "sync" "time"
"github.com/google/uuid" "tangled.org/core/api/tangled" "tangled.org/core/spindle/db" "tangled.org/core/spindle/models" "tangled.org/core/spindle/storage")
type CacheUpdateAction int
const ( CacheUsed CacheUpdateAction = iota + 1 CacheStored CacheDiscarded CacheMissing)
type CacheUpdate struct { Action CacheUpdateAction ID string Ref string SizeBytes int64 Checksum string}
type CacheController interface { Plan(ctx context.Context, pipeline *models.Pipeline, workflow *models.Workflow) ([]models.CacheBinding, error) Apply(ctx context.Context, update CacheUpdate) error}
type localCacheController struct { index *db.DB store storage.Storage repoDir string maxBytesPerOwner int64 maxEntriesPerOwner int64 logger *slog.Logger}
func NewLocalCacheController(index *db.DB, store storage.Storage, repoDir string, maxBytesPerOwner, maxEntriesPerOwner int64, logger *slog.Logger) CacheController { return &localCacheController{ index: index, store: store, repoDir: repoDir, maxBytesPerOwner: maxBytesPerOwner, maxEntriesPerOwner: maxEntriesPerOwner, logger: logger, }}
// on a hash miss, the newest older generation still warms the build// unusable entries degrade to a plain missfunc resolveCaches(ctx context.Context, l *slog.Logger, index *db.DB, repoDID, engine, repoPath, rev string, entries []models.CacheEntry) []models.CacheBinding { resolved := make([]models.CacheBinding, 0, len(entries)) for entryIndex, entry := range entries { hash := "" if len(entry.Hash) > 0 && repoPath != "" { sum, missing := hashKeyFilesLegacy(ctx, repoPath, rev, entry.Hash) for _, m := range missing { l.Warn("cache hash file not in repo", "key", entry.Key, "path", m) } hash = sum } binding := models.CacheBinding{ EntryIndex: entryIndex, Paths: entry.Paths, Key: entry.Key, Hash: hash, CompressionLevel: entry.CompressionLevel, When: entry.When, }
found, err := index.FindCacheEntry(ctx, repoDID, engine, entry.Key, hash) if errors.Is(err, sql.ErrNoRows) && hash != "" { found, err = index.FindFallbackCacheEntry(ctx, repoDID, engine, entry.Key, hash) if err == nil { binding.RestoreName = entry.Key + "-" + found.CacheHash } } if err != nil && !errors.Is(err, sql.ErrNoRows) { l.Warn("cache lookup failed; entry will save but not restore", "key", entry.Key, "err", err) } else if err == nil { binding.RestoreID = found.ID binding.RestoreKey = found.StorageKey } resolved = append(resolved, binding) } return resolved}
func hashKeyFilesLegacy(ctx context.Context, repoPath, rev string, paths []string) (string, []string) { h := sha256.New() var missing []string hashed := 0 for _, p := range paths { blob, err := gitBlobId(ctx, repoPath, rev, p) if err != nil { missing = append(missing, p) continue } fmt.Fprintf(h, "%s=%s\n", p, blob) hashed++ } if hashed == 0 { return "", missing } return hex.EncodeToString(h.Sum(nil))[:12], missing}
func hashKeyFiles(ctx context.Context, repoPath, rev string, paths []string) (string, []string) { h := sha256.New() var missing []string hashed := 0 for _, p := range paths { blob, err := gitBlobId(ctx, repoPath, rev, p) if err != nil { missing = append(missing, p) continue } fmt.Fprintf(h, "%s=%s\n", p, blob) hashed++ } if hashed == 0 { return "", missing } return hex.EncodeToString(h.Sum(nil)), missing}
// sparse checkouts might not have the file, but the object database doesfunc gitBlobId(ctx context.Context, repoPath, rev, path string) (string, error) { out, err := exec.CommandContext(ctx, "git", "-C", repoPath, "rev-parse", rev+":"+path).Output() if err != nil { return "", err } return strings.TrimSpace(string(out)), nil}func cacheIdentityKey(workflowName, ref, rev, key string) string { return url.PathEscape(workflowName) + ":" + url.PathEscape(ref) + ":" + url.PathEscape(rev) + ":" + key}
func pipelineRef(metadata *tangled.Pipeline_TriggerMetadata) string { if metadata == nil { return "" } if metadata.Push != nil { return metadata.Push.Ref } if metadata.PullRequest != nil { return "refs/heads/" + metadata.PullRequest.SourceBranch } if metadata.Manual != nil && metadata.Manual.Ref != nil { return *metadata.Manual.Ref } if metadata.Schedule != nil { return metadata.Schedule.Ref } return ""}
func resolveCachesIdentity(ctx context.Context, l *slog.Logger, index *db.DB, repoDID, engine, repoPath, rev, workflowName, ref string, entries []models.CacheEntry) []models.CacheBinding { resolved := make([]models.CacheBinding, 0, len(entries)) for i, entry := range entries { hash := "" if len(entry.Hash) > 0 && repoPath != "" { hash, _ = hashKeyFiles(ctx, repoPath, rev, entry.Hash) } lookupKey := cacheIdentityKey(workflowName, ref, rev, entry.Key) b := models.CacheBinding{EntryIndex: i, Paths: entry.Paths, Key: entry.Key, Hash: hash, CompressionLevel: entry.CompressionLevel, When: entry.When} found, err := index.FindCacheEntry(ctx, repoDID, engine, lookupKey, hash) if errors.Is(err, sql.ErrNoRows) && entry.Fallback && hash != "" { found, err = index.FindFallbackCacheEntryForIdentity(ctx, repoDID, engine, url.PathEscape(workflowName)+":"+url.PathEscape(ref)+":", hash) if err == nil { b.RestoreName = entry.Key + "-" + found.CacheHash } } if err == nil { b.RestoreID, b.RestoreKey, b.Checksum, b.SizeBytes = found.ID, found.StorageKey, found.Checksum, found.SizeBytes } if err != nil && !errors.Is(err, sql.ErrNoRows) { l.Warn("cache lookup failed; entry will save but not restore", "key", entry.Key, "err", err) } resolved = append(resolved, b) } return resolved}
func (c *localCacheController) Plan(ctx context.Context, pipeline *models.Pipeline, workflow *models.Workflow) ([]models.CacheBinding, error) { repo, err := c.index.GetRepoByDid(pipeline.RepoDid) if err != nil { return nil, fmt.Errorf("cache owner lookup: %w", err) }
repoPath, rev := "", "" if metadata := pipeline.TriggerMetadata; metadata != nil { if resolvedRev, err := models.ExtractCommitSHA(*metadata); err != nil { c.logger.Warn("cannot resolve pipeline commit; cache hashing disabled", "err", err) } else { did := pipeline.RepoDid.String() if metadata.SourceRepo != nil && *metadata.SourceRepo != "" { did = *metadata.SourceRepo } repoPath, rev = filepath.Join(c.repoDir, did), resolvedRev } }
cacheNamespace := workflow.CacheNamespace if cacheNamespace == "" { cacheNamespace = models.CacheNamespace() } cacheEngine := workflow.Engine + "/" + cacheNamespace var bindings []models.CacheBinding if workflow.Name == "" && pipeline.TriggerMetadata == nil { bindings = resolveCaches(ctx, c.logger, c.index, pipeline.RepoDid.String(), cacheEngine, repoPath, rev, workflow.Caches) } else { bindings = resolveCachesIdentity(ctx, c.logger, c.index, pipeline.RepoDid.String(), cacheEngine, repoPath, rev, workflow.Name, pipelineRef(pipeline.TriggerMetadata), workflow.Caches) } now := time.Now() pendingUntil := now // pending rows stay live through the workflow deadline if deadline, ok := ctx.Deadline(); ok { pendingUntil = deadline } for i := range bindings { id := uuid.NewString() cacheKey := bindings[i].Key if workflow.Name != "" || pipeline.TriggerMetadata != nil { cacheKey = cacheIdentityKey(workflow.Name, pipelineRef(pipeline.TriggerMetadata), rev, cacheKey) } saveKey, err := models.CacheObjectKey(pipeline.RepoDid.String(), id) if err != nil { return nil, err } bindings[i].SaveID = id bindings[i].SaveKey = saveKey inserted, err := c.index.InsertCacheEntryWithinQuota(ctx, db.CacheEntry{ ID: id, StorageKey: bindings[i].SaveKey, OwnerDID: repo.Owner.String(), RepoDID: pipeline.RepoDid.String(), Engine: cacheEngine, CacheKey: cacheKey, CacheHash: bindings[i].Hash, State: "pending", CreatedAt: now, LastUsedAt: pendingUntil, }, c.maxEntriesPerOwner) if err != nil { for _, planned := range bindings[:i] { _ = c.Apply(context.WithoutCancel(ctx), CacheUpdate{ Action: CacheDiscarded, ID: planned.SaveID, Ref: planned.SaveKey, }) } for j := range bindings { bindings[j].SaveID = "" bindings[j].SaveKey = "" } c.logger.Warn("cache save reservations failed; restoring only", "err", err) return bindings, nil } if !inserted { bindings[i].SaveID = "" bindings[i].SaveKey = "" } } return bindings, nil}
type cacheByteReserver interface { reserveBytes(context.Context, models.CacheBinding) (int64, error)}
func (c *localCacheController) reserveBytes(ctx context.Context, binding models.CacheBinding) (int64, error) { return c.index.ReserveCacheEntryBytes(ctx, binding.SaveID, c.maxBytesPerOwner)}
func (c *localCacheController) Apply(ctx context.Context, update CacheUpdate) error { now := time.Now() switch update.Action { case CacheUsed: return c.index.TouchCacheEntry(ctx, update.ID, now) case CacheStored: superseded, err := c.index.MarkCacheEntryReadyWithChecksum(ctx, update.ID, update.SizeBytes, update.Checksum, c.maxBytesPerOwner, now) if err != nil { _ = c.store.Delete(context.WithoutCancel(ctx), update.Ref) return err } for _, old := range superseded { if err := c.deleteEntry(context.WithoutCancel(ctx), old.StorageKey, old.ID); err != nil { c.logger.Warn("delete replaced cache failed", "id", old.ID, "err", err) } } case CacheDiscarded, CacheMissing: if update.Action == CacheDiscarded { claimed, err := c.index.ClaimPendingCacheEntry(ctx, update.ID) if err != nil || !claimed { return err } } return c.deleteEntry(ctx, update.Ref, update.ID) default: return fmt.Errorf("unknown cache update action %d", update.Action) } return nil}
func (c *localCacheController) deleteEntry(ctx context.Context, ref, id string) error { pinned, err := c.index.CacheEntryPinned(ctx, id) if err != nil { return err } if pinned { return c.index.ApplyEventBatch(nil, func(tx *db.EventBatchTx) error { return tx.QueueCacheObjectDeletion(ctx, ref, time.Now()) }) } if err := c.store.Delete(ctx, ref); err != nil { return err } return c.index.DeleteCacheEntry(ctx, id)}
type preplannedCacheController struct { apply func(context.Context, CacheUpdate) error}
func NewPreplannedCacheController(apply func(context.Context, CacheUpdate) error) CacheController { return &preplannedCacheController{apply: apply}}
func (c *preplannedCacheController) Plan(_ context.Context, _ *models.Pipeline, workflow *models.Workflow) ([]models.CacheBinding, error) { return workflow.CacheBindings, nil}
func (c *preplannedCacheController) Apply(ctx context.Context, update CacheUpdate) error { return c.apply(ctx, update)}
type trackedCacheStore struct { storage.Storage controller CacheController logger *slog.Logger mu sync.Mutex restores map[string]models.CacheBinding saves map[string]models.CacheBinding}
func newTrackedCacheStore(base storage.Storage, controller CacheController, logger *slog.Logger, bindings []models.CacheBinding) *trackedCacheStore { store := &trackedCacheStore{ Storage: base, controller: controller, logger: logger, restores: make(map[string]models.CacheBinding, len(bindings)), saves: make(map[string]models.CacheBinding, len(bindings)), } for _, binding := range bindings { if binding.RestoreKey != "" { store.restores[binding.RestoreKey] = binding } if binding.SaveKey != "" { store.saves[binding.SaveKey] = binding } } return store}
func (s *trackedCacheStore) beginRestore(ctx context.Context, id string) (bool, *db.DB, error) { if c, ok := s.controller.(*localCacheController); ok { pinned, err := c.index.BeginCacheRestore(ctx, id) return pinned, c.index, err } return true, nil, nil}func (s *trackedCacheStore) endRestore(ctx context.Context, id string) error { if c, ok := s.controller.(*localCacheController); ok { return c.index.EndCacheRestore(ctx, id) } return nil}
func (s *trackedCacheStore) Get(ctx context.Context, key string) (io.ReadCloser, error) { s.mu.Lock() binding, tracked := s.restores[key] s.mu.Unlock() if !tracked { return nil, fmt.Errorf("cache metadata missing for %q", key) }
pinned, index, err := s.beginRestore(ctx, binding.RestoreID) if err != nil { return nil, err } if !pinned { return nil, storage.ErrNotExist } reader, err := s.Storage.Get(ctx, key) if err != nil { _ = s.endRestore(context.WithoutCancel(ctx), binding.RestoreID) if errors.Is(err, storage.ErrNotExist) { _ = s.controller.Apply(context.WithoutCancel(ctx), CacheUpdate{Action: CacheMissing, ID: binding.RestoreID, Ref: binding.RestoreKey}) } return nil, err } if err := s.controller.Apply(ctx, CacheUpdate{ Action: CacheUsed, ID: binding.RestoreID, Ref: binding.RestoreKey, }); err != nil { s.logger.Warn("cache usage update failed", "id", binding.RestoreID, "err", err) } return &verifiedCacheReader{ReadCloser: reader, expected: binding.Checksum, expectedSize: binding.SizeBytes, id: binding.RestoreID, index: index}, nil}
type verifiedCacheReader struct { io.ReadCloser expected string expectedSize int64 id string index *db.DB h hash.Hash n int64 checked bool}
func (r *verifiedCacheReader) verify() error { if r.expectedSize > 0 && r.n != r.expectedSize { return fmt.Errorf("cache object size mismatch: got %d, want %d", r.n, r.expectedSize) } if r.expected != "" && hex.EncodeToString(r.h.Sum(nil)) != r.expected { return fmt.Errorf("cache object checksum mismatch") } return nil}
func (r *verifiedCacheReader) Read(p []byte) (int, error) { if r.h == nil { r.h = sha256.New() } n, err := r.ReadCloser.Read(p) if n > 0 { _, _ = r.h.Write(p[:n]) r.n += int64(n) } if err == io.EOF { r.checked = true if verifyErr := r.verify(); verifyErr != nil { return n, verifyErr } } return n, err}
func (r *verifiedCacheReader) Close() error { var err error if !r.checked { if r.h == nil { r.h = sha256.New() } _, err = io.Copy(io.Discard, r) if err == nil { r.checked = true err = r.verify() } } if closeErr := r.ReadCloser.Close(); err == nil { err = closeErr } if r.index != nil { if endErr := r.index.EndCacheRestore(context.Background(), r.id); err == nil { err = endErr } } return err}
func (s *trackedCacheStore) Put(ctx context.Context, key string, reader io.Reader) error { s.mu.Lock() binding, tracked := s.saves[key] s.mu.Unlock() if !tracked { return fmt.Errorf("cache metadata missing for %q", key) }
limit := int64(0) if reserver, ok := s.controller.(cacheByteReserver); ok { var err error limit, err = reserver.reserveBytes(ctx, binding) if err != nil { return err } } counted := &countingReader{r: reader, limit: limit, h: sha256.New()} put := s.Storage.Put if conditional, ok := s.Storage.(storage.ConditionalStorage); ok { put = func(ctx context.Context, key string, r io.Reader) error { return conditional.PutIfAbsent(ctx, key, r) } } if err := put(ctx, key, counted); err != nil { _ = s.Storage.Delete(context.WithoutCancel(ctx), key) return err } if err := s.controller.Apply(ctx, CacheUpdate{ Action: CacheStored, ID: binding.SaveID, Ref: binding.SaveKey, SizeBytes: counted.n, Checksum: hex.EncodeToString(counted.h.Sum(nil)), }); err != nil { _ = s.Storage.Delete(context.WithoutCancel(ctx), key) return fmt.Errorf("record cache upload: %w", err) }
s.mu.Lock() delete(s.saves, key) s.mu.Unlock() return nil}
func (s *trackedCacheStore) cleanup(ctx context.Context) { s.mu.Lock() pending := make([]models.CacheBinding, 0, len(s.saves)) for _, binding := range s.saves { pending = append(pending, binding) } clear(s.saves) s.mu.Unlock()
for _, binding := range pending { if err := s.controller.Apply(ctx, CacheUpdate{ Action: CacheDiscarded, ID: binding.SaveID, Ref: binding.SaveKey, }); err != nil { s.logger.Warn("discard incomplete cache failed", "id", binding.SaveID, "err", err) } }}
var ErrCacheObjectTooLarge = errors.New("cache object exceeds owner byte quota")
type countingReader struct { r io.Reader n int64 limit int64 h hash.Hash}
func (r *countingReader) Read(p []byte) (int, error) { if r.limit > 0 { remaining := r.limit - r.n if remaining <= 0 { return 0, ErrCacheObjectTooLarge } if int64(len(p)) > remaining { p = p[:remaining] } } n, err := r.r.Read(p) r.n += int64(n) if n > 0 && r.h != nil { _, _ = r.h.Write(p[:n]) } return n, err}