diff options
Diffstat (limited to 'internal/pkg/pkg.go')
| -rw-r--r-- | internal/pkg/pkg.go | 2947 |
1 files changed, 0 insertions, 2947 deletions
diff --git a/internal/pkg/pkg.go b/internal/pkg/pkg.go deleted file mode 100644 index 268d0248..00000000 --- a/internal/pkg/pkg.go +++ /dev/null @@ -1,2947 +0,0 @@ -// Package pkg provides utilities for packaging software. -package pkg - -import ( - "bufio" - "bytes" - "cmp" - "context" - "crypto/rand" - "crypto/sha512" - "encoding/base64" - "encoding/binary" - "errors" - "fmt" - "hash" - "io" - "io/fs" - "iter" - "maps" - "math" - "os" - "path/filepath" - "runtime" - "slices" - "strconv" - "strings" - "sync" - "sync/atomic" - "syscall" - "testing" - "time" - "unique" - "unsafe" - - "hakurei.app/check" - "hakurei.app/internal/info" - "hakurei.app/internal/lockedfile" - "hakurei.app/message" -) - -const ( - // programName is the string identifying this build system. - programName = "internal/pkg" -) - -type ( - // A Checksum is a SHA-384 checksum computed for a cured [Artifact]. - Checksum = [sha512.Size384]byte - - // An ID is a unique identifier returned by [KnownIdent.ID]. This value must - // be deterministically determined ahead of time. - ID Checksum -) - -// Encode is abbreviation for base64.URLEncoding.EncodeToString(checksum[:]). -func Encode(checksum Checksum) string { - return base64.URLEncoding.EncodeToString(checksum[:]) -} - -// Decode is abbreviation for base64.URLEncoding.Decode(checksum[:], []byte(s)). -func Decode(buf *Checksum, s string) (err error) { - var n int - n, err = base64.URLEncoding.Decode(buf[:], []byte(s)) - if err == nil && n != len(buf) { - err = io.ErrUnexpectedEOF - } - return -} - -// MustDecode decodes a string representation of [Checksum] and panics if there -// is a decoding error or the resulting data is too short. -func MustDecode(s string) (checksum Checksum) { - if err := Decode(&checksum, s); err != nil { - panic(err) - } - return -} - -var ( - // extension is a string uniquely identifying a set of custom [Artifact] - // implementations registered by calling [Register]. - extension string - - // openMu synchronises access to global state for initialisation. - openMu sync.Mutex - // opened is false if [Open] was never called. - opened bool -) - -// Extension returns a string uniquely identifying the currently registered set -// of custom [Artifact], or the zero value if none was registered. -func Extension() string { return extension } - -// ValidExtension returns whether s is valid for use in a call to SetExtension. -func ValidExtension(s string) bool { - if l := len(s); l == 0 || l > 128 { - return false - } - for _, v := range s { - if v < 'a' || v > 'z' { - return false - } - } - return true -} - -// ErrInvalidExtension is returned for a variant identification string for which -// [ValidExtension] returns false. -var ErrInvalidExtension = errors.New("invalid extension variant identification string") - -// SetExtension sets the extension variant identification string. SetExtension -// must be called before [Open] if custom [Artifact] implementations had been -// recorded by calling [Register]. -// -// The variant identification string must be between 1 and 128 bytes long and -// consists of only bytes between 'a' and 'z'. -// -// SetExtension is not safe for concurrent use. SetExtension is called at most -// once and must not be called after the first instance of Cache has been opened. -func SetExtension(s string) { - openMu.Lock() - defer openMu.Unlock() - - if opened { - panic("attempting to set extension after open") - } - if extension != "" { - panic("attempting to set extension twice") - } - if !ValidExtension(s) { - panic(ErrInvalidExtension) - } - extension = s - statusHeader = makeStatusHeader(s) -} - -// common holds elements and receives methods shared between different contexts. -type common struct { - // Context specific to this [Artifact]. The toplevel context in [Cache] must - // not be exposed directly. - ctx context.Context - - // Address of underlying [Cache], should be zeroed or made unusable after - // Cure returns and must not be exposed directly. - cache *Cache -} - -// TContext is passed to [TrivialArtifact.Cure] and provides information and -// methods required for curing the [TrivialArtifact]. -// -// Methods of TContext are safe for concurrent use. TContext is valid -// until [TrivialArtifact.Cure] returns. -type TContext struct { - // Populated during [Cache.Cure]. - work, temp *check.Absolute - - // Target [Artifact] encoded identifier. - ids string - // Pathname status was created at. - statusPath, statusSPath *check.Absolute - // File statusHeader and logs are written to. - status *os.File - // Error value during prepareStatus. - statusErr error - - common -} - -// makeStatusHeader creates the header written to every status file. This should -// not be called directly, its result is stored in statusHeader and will not -// change after the first [Cache] is opened. -func makeStatusHeader(extension string) string { - s := programName - if v := info.Version(); v != info.FallbackVersion { - s += " " + v - } - if extension != "" { - s += " with " + extension + " extensions" - } - s += " (" + runtime.GOARCH + ")" - if name, err := os.Hostname(); err == nil { - s += " on " + name - } - s += "\n\n" - return s -} - -// statusHeader is the header written to all status files in dirStatus. -var statusHeader = makeStatusHeader("") - -// prepareStatus initialises the status file once. -func (t *TContext) prepareStatus(writeHeader bool) error { - if t.statusPath != nil || t.status != nil { - return t.statusErr - } - - t.statusPath = t.cache.base.Append( - dirStatus, - t.ids, - ) - if t.status, t.statusErr = os.OpenFile( - t.statusPath.String(), - syscall.O_CREAT|syscall.O_EXCL|syscall.O_WRONLY, - 0400, - ); t.statusErr != nil { - return t.statusErr - } - - if writeHeader { - _, t.statusErr = t.status.WriteString(statusHeader) - } - return t.statusErr -} - -// GetStatusWriter returns a [io.Writer] for build logs. The caller must not -// seek this writer before the position it was first returned in. -func (t *TContext) GetStatusWriter() (io.Writer, error) { - err := t.prepareStatus(true) - return t.status, err -} - -// destroy destroys the temporary directory and joins its errors with the error -// referred to by errP. If the error referred to by errP is non-nil, the work -// directory is removed similarly. [Cache] is responsible for making sure work -// is never left behind for a successful [Cache.Cure]. -// -// If implementation had requested status, it is closed with error joined with -// the error referred to by errP. If the error referred to by errP is non-nil, -// the status file is removed from the filesystem. -// -// destroy must be deferred by [Cache.Cure] if [TContext] is passed to any Cure -// implementation. It should not be called prior to that point. -func (t *TContext) destroy(errP *error) { - if chmodErr, removeErr := removeAll(t.temp); chmodErr != nil || removeErr != nil { - *errP = errors.Join(*errP, chmodErr, removeErr) - } - - if *errP != nil { - chmodErr, removeErr := removeAll(t.work) - if chmodErr != nil || removeErr != nil { - *errP = errors.Join(*errP, chmodErr, removeErr) - } else if errors.Is(*errP, os.ErrExist) { - if linkError, ok := errors.AsType[*os.LinkError](*errP); ok && - linkError != nil && - linkError.Op == "rename" { - // two artifacts may be backed by the same file - *errP = nil - } - } - } - - if t.status != nil { - if err := t.status.Close(); err != nil { - *errP = errors.Join(*errP, err) - } - if *errP != nil { - *errP = errors.Join(*errP, os.Rename( - t.statusPath.String(), t.cache.base.Append( - dirFault, - t.ids+"."+strconv.FormatUint(uint64( - time.Now().UnixNano(), - ), 10), - ).String(), - )) - if t.statusSPath != nil { - t.cache.checksumMu.Lock() - *errP = errors.Join(*errP, os.Remove(t.statusSPath.String())) - t.cache.checksumMu.Unlock() - } - } - t.status = nil - } -} - -// Unwrap returns the underlying [context.Context]. -func (c *common) Unwrap() context.Context { return c.ctx } - -// GetMessage returns [message.Msg] held by the underlying [Cache]. -func (c *common) GetMessage() message.Msg { return c.cache.msg } - -// GetJobs returns the preferred number of jobs to run, when applicable. Its -// value must not affect cure outcome. -func (c *common) GetJobs() int { return c.cache.attr.Jobs } - -// GetLoad returns the preferred load average target, when applicable. Its -// value must not affect cure outcome. -func (c *common) GetLoad() int { return c.cache.attr.Load } - -// GetWorkDir returns a pathname to a directory which [Artifact] is expected to -// write its output to. This is not the final resting place of the [Artifact] -// and this pathname should not be directly referred to in the final contents. -func (t *TContext) GetWorkDir() *check.Absolute { return t.work } - -// GetTempDir returns a pathname which implementations may use as scratch space. -// A directory is not created automatically, implementations are expected to -// create it if they wish to use it, using [os.MkdirAll]. -func (t *TContext) GetTempDir() *check.Absolute { return t.temp } - -// Open tries to open [Artifact] for reading. If a implements [FileArtifact], -// its reader might be used directly, eliminating the roundtrip to vfs. -// Otherwise, it must cure into a directory containing a single regular file. -// -// If err is nil, the caller must close the resulting [io.ReadCloser] and return -// its error, if any. Failure to read r to EOF may result in a spurious -// [ChecksumMismatchError], or the underlying implementation may block on Close. -func (c *common) Open(a Artifact) (r io.ReadCloser, err error) { - if f, ok := a.(FileArtifact); ok { - return c.cache.openFile(c.ctx, f) - } - - var pathname *check.Absolute - if pathname, _, _, err = c.cache.cure(a, true, false); err != nil { - return - } - - var entries []os.DirEntry - if entries, err = os.ReadDir(pathname.String()); err != nil { - return - } - - if len(entries) != 1 || !entries[0].Type().IsRegular() { - err = errors.New( - "input directory does not contain a single regular file", - ) - return - } else { - return os.Open(pathname.Append(entries[0].Name()).String()) - } -} - -// FContext is passed to [FloodArtifact.Cure] and provides information and -// methods required for curing the [FloodArtifact]. -// -// Methods of FContext are safe for concurrent use. FContext is valid -// until [FloodArtifact.Cure] returns. -type FContext struct { - TContext - - // Cured top-level inputs looked up by Pathname. - inputs map[Artifact]cureRes -} - -// linkSubstitute links status for substitute if populated. -func (f *FContext) linkSubstitute(ids, substitutes string) (err error) { - if f.status == nil || ids == substitutes { - return - } - - statusS := f.cache.base.Append( - dirStatus, - substitutes, - ) - f.cache.checksumMu.Lock() - err = os.Link(f.cache.base.Append( - dirStatus, - ids, - ).String(), statusS.String()) - f.cache.checksumMu.Unlock() - if err == nil { - f.statusSPath = statusS - } - return -} - -// InvalidLookupError is the identifier of non-input [Artifact] looked up -// via [FContext.GetArtifact] by a misbehaving [Artifact] implementation. -type InvalidLookupError ID - -func (e InvalidLookupError) Error() string { - return "attempting to look up non-input artifact " + Encode(e) -} - -var _ error = InvalidLookupError{} - -// GetArtifact returns the identifier pathname and checksum of an [Artifact]. -// Calling Pathname with an [Artifact] not part of the slice returned by -// [Artifact.Inputs] panics. -func (f *FContext) GetArtifact(a Artifact) ( - pathname *check.Absolute, - checksum unique.Handle[Checksum], -) { - if res, ok := f.inputs[a]; ok { - return res.pathname, res.checksum - } - panic(InvalidLookupError(f.cache.Ident(a).Value())) -} - -// RContext is passed to [FileArtifact.Cure] and provides helper methods useful -// for curing the [FileArtifact]. -// -// Methods of RContext are safe for concurrent use. RContext is valid -// until [FileArtifact.Cure] returns. -type RContext struct{ common } - -// An Artifact is a read-only reference to a piece of data that may be created -// deterministically but might not currently be available in memory or on the -// filesystem. -type Artifact interface { - // Kind returns the [Kind] of artifact. This is usually unique to the - // concrete type but two functionally identical implementations of - // [Artifact] is allowed to return the same [Kind] value. - Kind() Kind - - // Params writes deterministic values describing [Artifact]. Implementations - // must guarantee that these values are unique among differing instances - // of the same implementation with identical dependencies and conveys enough - // information to create another instance of [Artifact] identical to the - // instance emitting these values. The new instance created via [IRReadFunc] - // from these values must then produce identical IR values. - // - // Result must remain identical across multiple invocations. - Params(ctx *IContext) - - // Inputs returns a slice of [Artifact] the current instance has access to - // while producing its output. - // - // Callers must not modify the retuned slice. - // - // Result must remain identical across multiple invocations. - Inputs() []Artifact - - // IsExclusive returns whether the [Artifact] is exclusive. Exclusive - // artifacts might not run in parallel with each other, and are still - // subject to the cures limit. - // - // Some implementations may saturate the CPU for a nontrivial amount of - // time. Curing multiple such implementations simultaneously causes - // significant CPU scheduler overhead. An exclusive artifact will generally - // not be cured alongside another exclusive artifact, thus alleviating this - // overhead. - // - // Note that [Cache] reserves the right to still cure exclusive - // artifacts concurrently as this is not a synchronisation primitive but - // an optimisation one. Implementations are forbidden from accessing global - // state regardless of exclusivity. - // - // Result must remain identical across multiple invocations. - IsExclusive() bool -} - -// FloodArtifact refers to an [Artifact] requiring its entire dependency graph -// to be cured prior to curing itself. -type FloodArtifact interface { - // Cure cures the current [Artifact] to the working directory obtained via - // [TContext.GetWorkDir] embedded in [FContext]. - // - // Implementations must not retain c. - Cure(f *FContext) (err error) - - Artifact -} - -// TrivialArtifact refers to an [Artifact] that cures without requiring that -// any other [Artifact] is cured before it. Its dependency tree is ignored after -// computing its identifier. -// -// TrivialArtifact is unable to cure any other [Artifact] and it cannot access -// pathnames. This type of [Artifact] is primarily intended for dependency-less -// artifacts or direct dependencies that only consists of [FileArtifact]. -type TrivialArtifact interface { - // Cure cures the current [Artifact] to the working directory obtained via - // [TContext.GetWorkDir]. - // - // Implementations must not retain c. - Cure(t *TContext) (err error) - - Artifact -} - -// KnownIdent is optionally implemented by [Artifact] and is used instead of -// [Cache.Ident] when it is available. -// -// This is very subtle to use correctly. The implementation must ensure that -// this value is globally unique, otherwise [Cache] can enter an inconsistent -// state. This should not be implemented outside of testing. -type KnownIdent interface { - // ID returns a globally unique identifier referring to the current - // [Artifact]. This value must be known ahead of time and guaranteed to be - // unique without having obtained the full contents of the [Artifact]. - ID() ID - - Artifact -} - -// KnownChecksum is optionally implemented by [Artifact] for an artifact with -// output known ahead of time. -type KnownChecksum interface { - // Checksum returns the address of a known checksum. - // - // Callers must not modify the [Checksum]. - // - // Result must remain identical across multiple invocations. - Checksum() Checksum - - Artifact -} - -// CuresExempt is optionally implemented for an artifact exempt to the -// cache-wide cures counter and limit. -type CuresExempt interface { - Artifact - - // CuresExempt is a no-op function but serves to distinguish implementations - // that are cures-exempt. - CuresExempt() -} - -// FileArtifact refers to an [Artifact] backed by a single file. -// -// FileArtifact does not support fine-grained cancellation. Its context is -// inherited from the first [TrivialArtifact] or [FloodArtifact] that opens it. -type FileArtifact interface { - // IsExecutable returns whether the resulting filesystem entry should be made - // executable, if the [FileArtifact] is cured to the on-disk cache. - // - // Result must remain identical across multiple invocations. - IsExecutable() bool - - // Cure returns [io.ReadCloser] of the full contents of [FileArtifact]. If - // [FileArtifact] implements [KnownChecksum], Cure is responsible for - // validating any data it produces and must return [ChecksumMismatchError] - // if validation fails. This error is conventionally returned during the - // first call to Close, but may be returned during any call to Read before - // EOF, or by Cure itself. - // - // Callers are responsible for closing the resulting [io.ReadCloser]. - // - // The resulting [io.ReadCloser] across multiple invocations must have - // identical behaviour. - Cure(r *RContext) (io.ReadCloser, error) - - Artifact -} - -// RevisionArtifact is optionally implemented by an artifact that had undergone -// internal changes affecting its behaviour while retaining its IR structure. -type RevisionArtifact interface { - // Revision returns the revision number of [Artifact]. This value is always - // represented in the IR. An IR stream produced from the same [Kind] with - // differing revision is rejected. - // - // Result must remain identical across multiple invocations. - Revision() uint64 - - Artifact -} - -// GetRevision returns the revision number of an [Artifact]. -func GetRevision(a Artifact) (revision uint64) { - revision = math.MaxUint64 - if r, ok := a.(RevisionArtifact); ok { - revision = r.Revision() - } - return -} - -// reportName returns a string describing [Artifact] presented to the user. -func reportName(a Artifact, id unique.Handle[ID]) string { - r := Encode(id.Value()) - if s, ok := a.(fmt.Stringer); ok { - if name := s.String(); name != "" { - r += "-" + name - } - } - return r -} - -// Kind corresponds to the concrete type of [Artifact] and is used to create -// identifier for an [Artifact] with dependencies. -type Kind uint64 - -const ( - // KindHTTPGet is the kind of [Artifact] returned by [NewHTTPGet]. - KindHTTPGet Kind = iota - // KindTar is the kind of [Artifact] returned by [NewTar]. - KindTar - // KindExec is the kind of [Artifact] returned by [NewExec]. - KindExec - // KindExecNet is the kind of [Artifact] returned by [NewExec] but with a - // non-nil checksum. - KindExecNet - // KindFile is the kind of [Artifact] returned by [NewFile]. - KindFile - // KindDecompress is the kind of [Artifact] returned by [NewDecompress]. - KindDecompress - // KindArchive is the kind of [Artifact] returned by [NewArchive]. - KindArchive - - // _kindEnd is the total number of kinds and does not denote a kind. - _kindEnd - - // KindCustomOffset is the first [Kind] value reserved for implementations - // not from this package. - KindCustomOffset = 1 << 31 -) - -const ( - // kindCollection is the kind of [Collect]. It never cures successfully. - kindCollection Kind = KindCustomOffset - 1 - iota -) - -const ( - // fileLock is the lock file for exclusive access to the cache directory. - fileLock = "lock" - // fileVariant is a file holding the variant identification string set by a - // prior call to [SetExtension]. - fileVariant = "variant" - - // dirSubstitute holds symlinks to artifacts by checksum, named after their - // substitute identifier. - dirSubstitute = "substitute" - // dirIdentifier holds symlinks to artifacts by checksum, named after their - // IR-based identifier. - dirIdentifier = "identifier" - // dirChecksum holds artifacts named after their [Checksum]. - dirChecksum = "checksum" - // dirStatus holds artifact metadata and logs named after their IR-based - // identifier. For [FloodArtifact], the same file is also available under - // its substitute identifier. - dirStatus = "status" - // dirFault holds status files of faulted cures. - dirFault = "fault" - - // dirWork holds working pathnames set up during [Cache.Cure]. - dirWork = "work" - // dirTemp holds scratch space allocated during [Cache.Cure]. - dirTemp = "temp" - - // dirExecScratch is scratch space set up for the container started by - // [Cache.EnterExec]. Exclusivity via Cache.inExec. - dirExecScratch = "scratch" - - // checksumLinknamePrefix is prepended to the encoded [Checksum] value - // of an [Artifact] when creating a symbolic link to dirChecksum. - checksumLinknamePrefix = "../" + dirChecksum + "/" -) - -// cureRes are the non-error results returned by [Cache.Cure]. -type cureRes struct { - pathname *check.Absolute - checksum unique.Handle[Checksum] -} - -// A pendingArtifactDep is an input [Artifact] pending concurrent curing, -// subject to the cures limit. Values pointed to by result addresses are safe -// to access after the [sync.WaitGroup] associated with this pendingArtifactDep -// is done. pendingArtifactDep must not be reused or modified after it is sent -// to cure. -type pendingArtifactDep struct { - // Dependency artifact populated during [Cache.Cure]. - a Artifact - - // Address of result pathname populated during [Cache.Cure] and dereferenced - // if curing succeeds. - resP *cureRes - - // Address of result error map populated during [Cache.Cure], dereferenced - // after acquiring errsMu if curing fails. No additional action is taken, - // [Cache] and its caller are responsible for further error handling. - errs InputError - // Address of mutex synchronising access to errs. - errsMu *sync.Mutex - - // For synchronising access to result buffer. - *sync.WaitGroup -} - -const ( - // CValidateKnown arranges for [KnownChecksum] outcomes to be validated to - // match its intended checksum. - // - // A correct implementation of [KnownChecksum] does not successfully cure - // with output not matching its intended checksum. When an implementation - // fails to perform this validation correctly, the on-disk format enters - // an inconsistent state (correctable by [Cache.Scrub]). - // - // This flag causes [Cache.Cure] to always compute the checksum, and reject - // a cure if it does not match the intended checksum. - // - // This behaviour significantly reduces performance and is not recommended - // outside of testing a custom [Artifact] implementation. - CValidateKnown = 1 << iota - - // CSchedIdle arranges for the [ext.SCHED_IDLE] scheduling priority to be - // set for [KindExec] and [KindExecNet] containers. - CSchedIdle - - // CAssumeChecksum enables the use of [KnownChecksum] for duplicate function - // call suppression via the on-disk cache. - // - // This may cause incorrect cure outcome if an impossible checksum is - // specified that matches an output already present in the on-disk cache. - // This may be avoided by purposefully specifying a statistically - // unattainable checksum, like the zero value. - // - // While this optimisation might seem appealing, it is almost never - // applicable in real world use. Almost every time this path was taken, it - // was caused by an incorrect checksum accidentally left behind while - // bumping a package. Only enable this if you are really sure you need it. - CAssumeChecksum - - // CHostAbstract disables restriction of sandboxed processes from connecting - // to an abstract UNIX socket created by a host process. - // - // This is considered less secure in some systems, but does not introduce - // impurity due to [KindExecNet] being [KnownChecksum]. This flag exists - // to support kernels without Landlock LSM enabled. - CHostAbstract - - // CPromoteVariant allows [pkg.Open] to promote an unextended on-disk cache - // to the current extension variant. This is a one-way operation. - CPromoteVariant - - // CSuppressInit arranges for verbose output of the container init to be - // suppressed regardless of [message.Msg] state. - CSuppressInit - - // CIgnoreSubstitutes disables content-based input substitution. - CIgnoreSubstitutes - - // CExternShallow arranges for only non-flood inputs to be fetched when - // curing an [Artifact] available via the external cache. - CExternShallow - - // CColourOutput enables output colouring via ANSI control sequences. - CColourOutput -) - -// toplevel holds [context.WithCancel] over caller-supplied context, where all -// [Artifact] context are derived from. -type toplevel struct { - ctx context.Context - cancel context.CancelFunc -} - -// newToplevel returns the address of a new toplevel via ctx. -func newToplevel(ctx context.Context) *toplevel { - var t toplevel - t.ctx, t.cancel = context.WithCancel(ctx) - return &t -} - -// pendingCure provides synchronisation and cancellation for pending cures. -type pendingCure struct { - // Closed on cure completion. - done <-chan struct{} - // Error outcome, safe to access after done is closed. - err error - // Cancels the corresponding cure. - cancel context.CancelFunc -} - -// An External cache provides prepared [Artifact] cure outcomes. -type External interface { - // Artifact returns the address of the [Checksum] of the cure outcome of - // an [Artifact] corresponding to id, or nil if this [Artifact] is not - // available in the external cache. - Artifact(ctx context.Context, id unique.Handle[ID]) (*Checksum, error) - // Checksum returns an [Artifact] producing the specified checksum. - Checksum(checksum unique.Handle[Checksum]) Artifact - // Status returns [io.ReadCloser] of the status file of an [Artifact] - // corresponding to id, or nil if this [Artifact] is not available or a - // status file is not present. - Status(r *RContext, id unique.Handle[ID]) (io.ReadCloser, error) -} - -// Cache is a support layer that implementations of [Artifact] can use to store -// cured [Artifact] data in a content addressed fashion. -type Cache struct { - // Cures of any variant of [Artifact] sends to cures before entering the - // implementation and receives an equal amount of elements after. - cures chan struct{} - - // Parent context which toplevel was derived from. - parent context.Context - // For deriving curing context, must not be accessed directly. - toplevel atomic.Pointer[toplevel] - // For waiting on input curing goroutines. - wg sync.WaitGroup - // Reports new cures and passed to [Artifact]. - msg message.Msg - // Select graphics rendition sequences, populated by Open. - sgrRes, sgrIdent, sgrWarn, sgrErr string - - // Directory where all [Cache] related files are placed. - base *check.Absolute - // Immutable [CacheAttr] populated by [Open]. - attr CacheAttr - - // Must not be exposed directly. - irCache - - // Synchronises access to dirChecksum. - checksumMu sync.RWMutex - - // Presence of an alternative in the cache. Keys are not valid identifiers - // and must not be used as such. - substitute map[unique.Handle[ID]]unique.Handle[Checksum] - // Synchronises access to substitute and corresponding filesystem entries. - substituteMu sync.RWMutex - // Identifier to content pair cache. - ident map[unique.Handle[ID]]unique.Handle[Checksum] - // Identifier to error pair for unrecoverably faulted [Artifact]. - identErr map[unique.Handle[ID]]error - // Pending identifiers, accessed through Cure for entries not in ident. - identPending map[unique.Handle[ID]]*pendingCure - // Synchronises access to ident and corresponding filesystem entries. - identMu sync.RWMutex - // Synchronises entry into Abort and Cure. - abortMu sync.RWMutex - - // Synchronises entry into exclusive artifacts for the cure method. - exclMu sync.Mutex - // Buffered I/O free list, must not be accessed directly. - brPool, bwPool sync.Pool - - // Optional external cache implementation. - extern External - // Caches responses from extern. - externCache map[unique.Handle[ID]]unique.Handle[Checksum] - // Synchronises access to extern. - externMu sync.RWMutex - - // Unlocks the on-filesystem cache. Must only be called from Close. - unlock func() - // Whether [Cache] is considered closed. - closed bool - // Synchronises calls to Abort and Close. - closeMu sync.Mutex - - // Whether EnterExec has not yet returned. - inExec atomic.Bool -} - -// extIdent is a [Kind] concatenated with [ID]. -type extIdent [wordSize + len(ID{})]byte - -// getIdentBuf returns the address of an extIdent for Ident. -func (ic *irCache) getIdentBuf() *extIdent { return ic.identPool.Get().(*extIdent) } - -// putIdentBuf adds buf to identPool. -func (ic *irCache) putIdentBuf(buf *extIdent) { ic.identPool.Put(buf) } - -// storeIdent adds an [Artifact] to the artifact cache. -func (ic *irCache) storeIdent(a Artifact, buf *extIdent) unique.Handle[ID] { - idu := unique.Make(ID(buf[wordSize:])) - ic.artifact.Store(a, idu) - return idu -} - -// Ident returns the identifier of an [Artifact]. -func (ic *irCache) Ident(a Artifact) unique.Handle[ID] { - buf, idu := ic.unsafeIdent(a, false) - if buf != nil { - idu = ic.storeIdent(a, buf) - ic.putIdentBuf(buf) - } - return idu -} - -// unsafeIdent implements Ident but returns the underlying buffer for a newly -// computed identifier. Callers must return this buffer to identPool. encodeKind -// is only a hint, kind may still be encoded in the buffer. -func (ic *irCache) unsafeIdent(a Artifact, encodeKind bool) ( - buf *extIdent, - idu unique.Handle[ID], -) { - if id, ok := ic.artifact.Load(a); ok { - idu = id.(unique.Handle[ID]) - return - } - - if ki, ok := a.(KnownIdent); ok { - buf = ic.getIdentBuf() - if encodeKind { - binary.LittleEndian.PutUint64(buf[:], uint64(a.Kind())) - } - *(*ID)(buf[wordSize:]) = ki.ID() - return - } - - buf = ic.getIdentBuf() - h := sha512.New384() - if err := ic.Encode(h, a); err != nil { - // unreachable - panic(err) - } - binary.LittleEndian.PutUint64(buf[:], uint64(a.Kind())) - h.Sum(buf[wordSize:wordSize]) - return -} - -// getReader is like [bufio.NewReader] but for brPool. -func (c *Cache) getReader(r io.Reader) *bufio.Reader { - br := c.brPool.Get().(*bufio.Reader) - br.Reset(r) - return br -} - -// putReader adds br to brPool. -func (c *Cache) putReader(br *bufio.Reader) { c.brPool.Put(br) } - -// bufioReadCloser is the concrete type of value returned by Cache.getReaderRC. -type bufioReadCloser struct { - // Saved close error. - closeErr error - // Synchronises calls to Close. - closeOnce sync.Once - - // For backing freelist. - c *Cache - // Underlying reader. - r io.ReadCloser - // Allocated from c. - *bufio.Reader -} - -// Close closes the underlying reader, saves its return value, and returns the -// [bufio.Reader] instance to the backing [Cache]. -func (brc *bufioReadCloser) Close() error { - brc.closeOnce.Do(func() { - br := brc.Reader - brc.Reader = nil - brc.c.putReader(br) - brc.closeErr = brc.r.Close() - }) - return brc.closeErr -} - -// getReaderRC is like getReader, but returns an [io.ReadCloser]. -func (c *Cache) getReaderRC(r io.ReadCloser) io.ReadCloser { - return &bufioReadCloser{c: c, r: r, Reader: c.getReader(r)} -} - -// getWriter is like [bufio.NewWriter] but for bwPool. -func (c *Cache) getWriter(w io.Writer) *bufio.Writer { - bw := c.bwPool.Get().(*bufio.Writer) - bw.Reset(w) - return bw -} - -// putWriter adds bw to bwPool. -func (c *Cache) putWriter(bw *bufio.Writer) { c.bwPool.Put(bw) } - -// A ChecksumMismatchError describes an [Artifact] with unexpected content. -type ChecksumMismatchError struct { - // Actual and expected checksums. - Got, Want Checksum -} - -func (e *ChecksumMismatchError) Error() string { - return "got " + Encode(e.Got) + - " instead of " + Encode(e.Want) -} - -// LinknamePrefixError describes a malformed linkname to a [Checksum]. -type LinknamePrefixError string - -func (e LinknamePrefixError) Error() string { - return "linkname " + strconv.Quote(string(e)) + " missing prefix" -} - -// readlinkChecksum reads a symbolic link to a dirChecksum entry and saves the -// decoded [Checksum] to the value pointed to by buf. The checksumLinknamePrefix -// is required. -func readlinkChecksum(a *check.Absolute, buf *Checksum) error { - linkname, err := os.Readlink(a.String()) - if err != nil { - return nil - } - - if !strings.HasPrefix(linkname, checksumLinknamePrefix) { - return LinknamePrefixError(linkname) - } - return Decode(buf, linkname[len(checksumLinknamePrefix):]) -} - -// SetExternal sets e as the [External] implementation of c. -func (c *Cache) SetExternal(e External) { - c.externMu.Lock() - c.externCache = make(map[unique.Handle[ID]]unique.Handle[Checksum]) - c.extern = e - c.externMu.Unlock() -} - -// ScrubError describes the outcome of a [Cache.Scrub] call where errors were -// found and removed from the underlying storage of [Cache]. -type ScrubError struct { - // Content-addressed entries not matching their checksum. This can happen - // if an incorrect [FileArtifact] implementation was cured against - // a non-strict [Cache]. - ChecksumMismatches []ChecksumMismatchError - // Dangling identifier symlinks. This can happen if the content-addressed - // entry was removed while scrubbing due to a checksum mismatch. - DanglingIdentifiers []ID - // Dangling status files. This can happen if a dangling status symlink was - // removed while scrubbing. - DanglingStatus []ID - // Miscellaneous errors, including [os.ReadDir] on checksum and identifier - // directories, [Decode] on entry names and [os.RemoveAll] on inconsistent - // entries. - Errs map[unique.Handle[string]][]error -} - -// errs is a deterministic iterator over Errs. -func (e *ScrubError) errs(yield func(unique.Handle[string], []error) bool) { - keys := slices.AppendSeq( - make([]unique.Handle[string], 0, len(e.Errs)), - maps.Keys(e.Errs), - ) - slices.SortFunc(keys, func(a, b unique.Handle[string]) int { - return strings.Compare(a.Value(), b.Value()) - }) - for _, key := range keys { - if !yield(key, e.Errs[key]) { - break - } - } -} - -// Unwrap returns a concatenation of ChecksumMismatches and Errs. -func (e *ScrubError) Unwrap() []error { - s := make([]error, 0, len(e.ChecksumMismatches)+len(e.Errs)) - for _, err := range e.ChecksumMismatches { - s = append(s, &err) - } - for _, errs := range e.errs { - s = append(s, errs...) - } - return s -} - -// Error returns a multi-line representation of [ScrubError]. -func (e *ScrubError) Error() string { - var segments []string - var buf strings.Builder - - if len(e.ChecksumMismatches) > 0 { - buf.Reset() - buf.WriteString("checksum mismatches:\n") - for _, m := range e.ChecksumMismatches { - buf.WriteString(m.Error() + "\n") - } - segments = append(segments, buf.String()) - } - if len(e.DanglingIdentifiers) > 0 { - buf.Reset() - buf.WriteString("dangling identifiers:\n") - for _, id := range e.DanglingIdentifiers { - buf.WriteString(Encode(id) + "\n") - } - segments = append(segments, buf.String()) - } - if len(e.DanglingStatus) > 0 { - buf.Reset() - buf.WriteString("dangling status:\n") - for _, id := range e.DanglingStatus { - buf.WriteString(Encode(id) + "\n") - } - segments = append(segments, buf.String()) - } - if len(e.Errs) > 0 { - buf.Reset() - buf.WriteString("errors during scrub:\n") - for pathname, errs := range e.errs { - buf.WriteString(" " + pathname.Value() + ":\n") - for _, err := range errs { - buf.WriteString(" " + err.Error() + "\n") - } - } - segments = append(segments, buf.String()) - } - return strings.Join(segments, "\n") -} - -// Scrub frees internal in-memory identifier to content pair cache, verifies all -// cached artifacts against their checksums, checks for dangling identifier -// symlinks and removes them if found. -// -// This method is not safe for concurrent use with any other method. -func (c *Cache) Scrub(checks int) error { - if checks <= 0 { - checks = runtime.NumCPU() - } - - c.substituteMu.Lock() - defer c.substituteMu.Unlock() - c.identMu.Lock() - defer c.identMu.Unlock() - c.checksumMu.Lock() - defer c.checksumMu.Unlock() - - c.substitute = make(map[unique.Handle[ID]]unique.Handle[Checksum]) - c.ident = make(map[unique.Handle[ID]]unique.Handle[Checksum]) - c.identErr = make(map[unique.Handle[ID]]error) - c.artifact.Clear() - - var ( - se = ScrubError{Errs: make(map[unique.Handle[string]][]error)} - seMu sync.Mutex - - addErr = func(pathname *check.Absolute, err error) { - seMu.Lock() - se.Errs[pathname.Handle()] = append(se.Errs[pathname.Handle()], err) - seMu.Unlock() - } - ) - - type checkEntry struct { - ent os.DirEntry - check func(ent os.DirEntry, want *Checksum) bool - } - var ( - dir *check.Absolute - wg sync.WaitGroup - w = make(chan checkEntry, checks) - p = sync.Pool{New: func() any { return new(Checksum) }} - ) - condemn := func(ent os.DirEntry) { - pathname := dir.Append(ent.Name()) - chmodErr, removeErr := removeAll(pathname) - if chmodErr != nil { - addErr(pathname, chmodErr) - } - if removeErr != nil { - addErr(pathname, removeErr) - } - } - for i := 0; i < checks; i++ { - go func() { - for ce := range w { - want := p.Get().(*Checksum) - ent := ce.ent - if err := Decode(want, ent.Name()); err != nil { - addErr(dir.Append(ent.Name()), err) - wg.Go(func() { condemn(ent) }) - } else if !ce.check(ent, want) { - wg.Go(func() { condemn(ent) }) - } else { - c.msg.Verbosef( - "%s%s%s is consistent", - c.sgrIdent, ent.Name(), c.sgrRes, - ) - } - p.Put(want) - wg.Done() - } - }() - } - defer close(w) - - dir = c.base.Append(dirChecksum) - if entries, readdirErr := os.ReadDir(dir.String()); readdirErr != nil { - addErr(dir, readdirErr) - } else { - wg.Add(len(entries)) - for _, ent := range entries { - w <- checkEntry{ent, func(ent os.DirEntry, want *Checksum) bool { - got := p.Get().(*Checksum) - defer p.Put(got) - - pathname := dir.Append(ent.Name()) - if ent.IsDir() { - if err := SumDir(got, pathname); err != nil { - addErr(pathname, err) - return true - } - } else if ent.Type().IsRegular() { - h := sha512.New384() - - if r, err := os.Open(pathname.String()); err != nil { - addErr(pathname, err) - return true - } else { - _, err = io.Copy(h, r) - closeErr := r.Close() - if closeErr != nil { - addErr(pathname, closeErr) - } - if err != nil { - addErr(pathname, err) - } - } - h.Sum(got[:0]) - } else { - addErr(pathname, InvalidFileModeError(ent.Type())) - return false - } - - if *got != *want { - seMu.Lock() - se.ChecksumMismatches = append(se.ChecksumMismatches, - ChecksumMismatchError{Got: *got, Want: *want}, - ) - seMu.Unlock() - return false - } - return true - }} - } - wg.Wait() - } - - for _, suffix := range []string{ - dirSubstitute, - dirIdentifier, - } { - dir = c.base.Append(suffix) - if entries, readdirErr := os.ReadDir(dir.String()); readdirErr != nil { - addErr(dir, readdirErr) - } else { - wg.Add(len(entries)) - for _, ent := range entries { - w <- checkEntry{ent, func(ent os.DirEntry, want *Checksum) bool { - got := p.Get().(*Checksum) - defer p.Put(got) - - pathname := dir.Append(ent.Name()) - if linkname, err := os.Readlink( - pathname.String(), - ); err != nil { - seMu.Lock() - se.Errs[pathname.Handle()] = append(se.Errs[pathname.Handle()], err) - se.DanglingIdentifiers = append(se.DanglingIdentifiers, *want) - seMu.Unlock() - return false - } else if err = Decode(got, filepath.Base(linkname)); err != nil { - seMu.Lock() - lnp := dir.Append(linkname) - se.Errs[lnp.Handle()] = append(se.Errs[lnp.Handle()], err) - se.DanglingIdentifiers = append(se.DanglingIdentifiers, *want) - seMu.Unlock() - return false - } - - if _, err := os.Stat(pathname.String()); err != nil { - if !errors.Is(err, os.ErrNotExist) { - addErr(pathname, err) - } - seMu.Lock() - se.DanglingIdentifiers = append(se.DanglingIdentifiers, *want) - seMu.Unlock() - return false - } - return true - }} - } - wg.Wait() - } - } - - dir = c.base.Append(dirStatus) - if entries, readdirErr := os.ReadDir(dir.String()); readdirErr != nil { - if !errors.Is(readdirErr, os.ErrNotExist) { - addErr(dir, readdirErr) - } - } else { - wg.Add(len(entries)) - for _, ent := range entries { - w <- checkEntry{ent, func(ent os.DirEntry, want *Checksum) bool { - got := p.Get().(*Checksum) - defer p.Put(got) - - var ok bool - for _, name := range [...]string{ - dirIdentifier, - dirSubstitute, - } { - if _, err := os.Stat(c.base.Append( - name, - ent.Name(), - ).String()); err != nil { - if !errors.Is(err, os.ErrNotExist) { - addErr(dir.Append(ent.Name()), err) - } - continue - } - ok = true - } - if !ok { - seMu.Lock() - se.DanglingStatus = append(se.DanglingStatus, *want) - seMu.Unlock() - } - return ok - }} - } - wg.Wait() - } - - if len(c.identPending) > 0 { - addErr(c.base, errors.New( - "scrub began with pending artifacts", - )) - } else { - pathname := c.base.Append(dirWork) - chmodErr, removeErr := removeAll(pathname) - if chmodErr != nil { - addErr(pathname, chmodErr) - } - if removeErr != nil { - addErr(pathname, removeErr) - } - - if err := os.Mkdir(pathname.String(), 0700); err != nil { - addErr(pathname, err) - } - - pathname = c.base.Append(dirTemp) - chmodErr, removeErr = removeAll(pathname) - if chmodErr != nil { - addErr(pathname, chmodErr) - } - if removeErr != nil { - addErr(pathname, removeErr) - } - } - - if len(se.ChecksumMismatches) > 0 || - len(se.DanglingIdentifiers) > 0 || - len(se.DanglingStatus) > 0 || - len(se.Errs) > 0 { - slices.SortFunc(se.ChecksumMismatches, func(a, b ChecksumMismatchError) int { - return bytes.Compare(a.Want[:], b.Want[:]) - }) - slices.SortFunc(se.DanglingIdentifiers, func(a, b ID) int { - return bytes.Compare(a[:], b[:]) - }) - slices.SortFunc(se.DanglingStatus, func(a, b ID) int { - return bytes.Compare(a[:], b[:]) - }) - return &se - } else { - return nil - } -} - -// loadOrStoreIdent attempts to load a cached [Artifact] by its identifier or -// wait for a pending [Artifact] to cure. If neither is possible, the current -// identifier is stored in identPending and a non-nil channel is returned. -// -// Since identErr is treated as grow-only, loadOrStoreIdent must not be entered -// without holding a read lock on abortMu. -func (c *Cache) loadOrStoreIdent(id unique.Handle[ID]) ( - ctx context.Context, - done chan<- struct{}, - checksum unique.Handle[Checksum], - err error, -) { - var ok bool - - c.identMu.Lock() - if checksum, ok = c.ident[id]; ok { - c.identMu.Unlock() - return - } - if err, ok = c.identErr[id]; ok { - c.identMu.Unlock() - return - } - - var pending *pendingCure - if pending, ok = c.identPending[id]; ok { - c.identMu.Unlock() - <-pending.done - c.identMu.RLock() - if checksum, ok = c.ident[id]; !ok { - err = pending.err - } - c.identMu.RUnlock() - return - } - - d := make(chan struct{}) - pending = &pendingCure{done: d} - ctx, pending.cancel = context.WithCancel(c.toplevel.Load().ctx) - c.wg.Add(1) - c.identPending[id] = pending - c.identMu.Unlock() - done = d - return -} - -// finaliseIdent commits a checksum or error to ident for an identifier -// previously submitted to identPending. -func (c *Cache) finaliseIdent( - done chan<- struct{}, - id unique.Handle[ID], - checksum unique.Handle[Checksum], - err error, -) { - c.identMu.Lock() - if err != nil { - c.identPending[id].err = err - c.identErr[id] = err - } else { - c.ident[id] = checksum - } - delete(c.identPending, id) - c.identMu.Unlock() - c.wg.Done() - - close(done) -} - -// zeroChecksum is a zero [Checksum] handle, used for comparison only. -var zeroChecksum unique.Handle[Checksum] - -// loadSubstitute returns a checksum corresponding to a substitute identifier, -// or zeroChecksum if an alternative is not available. -func (c *Cache) loadSubstitute( - substitute unique.Handle[ID], -) (unique.Handle[Checksum], error) { - c.substituteMu.RLock() - if checksum, ok := c.substitute[substitute]; ok { - c.substituteMu.RUnlock() - return checksum, nil - } - - linkname, err := os.Readlink(c.base.Append( - dirSubstitute, - Encode(substitute.Value()), - ).String()) - c.substituteMu.RUnlock() - - if err != nil { - if !errors.Is(err, os.ErrNotExist) { - return zeroChecksum, err - } - - c.substituteMu.Lock() - c.substitute[substitute] = zeroChecksum - c.substituteMu.Unlock() - return zeroChecksum, nil - } - - var checksum unique.Handle[Checksum] - buf := c.getIdentBuf() - err = Decode((*Checksum)(buf[:]), filepath.Base(linkname)) - if err == nil { - checksum = unique.Make(Checksum(buf[:])) - - c.substituteMu.Lock() - c.substitute[substitute] = checksum - c.substituteMu.Unlock() - } - c.putIdentBuf(buf) - - return checksum, err -} - -// Done returns a channel that is closed when the ongoing cure of an [Artifact] -// referred to by the specified identifier completes. Done may return nil if -// no ongoing cure of the specified identifier exists. -func (c *Cache) Done(id unique.Handle[ID]) <-chan struct{} { - c.identMu.RLock() - pending, ok := c.identPending[id] - c.identMu.RUnlock() - if !ok || pending == nil { - return nil - } - return pending.done -} - -// Cancel cancels the ongoing cure of an [Artifact] referred to by the specified -// identifier. Cancel returns whether the [context.CancelFunc] has been killed. -// Cancel returns after the cure is complete. -func (c *Cache) Cancel(id unique.Handle[ID]) bool { - c.identMu.RLock() - pending, ok := c.identPending[id] - c.identMu.RUnlock() - if !ok || pending == nil || pending.cancel == nil { - return false - } - pending.cancel() - <-pending.done - - c.abortMu.Lock() - c.identMu.Lock() - delete(c.identErr, id) - c.identMu.Unlock() - c.abortMu.Unlock() - return true -} - -// openFile tries to load [FileArtifact] from [Cache], and if that fails, -// obtains it via [FileArtifact.Cure] instead. Notably, it does not cure -// [FileArtifact] to the filesystem. If err is nil, the caller is responsible -// for closing the resulting [io.ReadCloser]. -// -// The context must originate from loadOrStoreIdent to enable cancellation. -func (c *Cache) openFile( - ctx context.Context, - f FileArtifact, -) (r io.ReadCloser, err error) { - if kc, ok := f.(KnownChecksum); c.attr.Flags&CAssumeChecksum != 0 && ok { - c.checksumMu.RLock() - r, err = os.Open(c.base.Append( - dirChecksum, - Encode(kc.Checksum()), - ).String()) - c.checksumMu.RUnlock() - } else { - c.identMu.RLock() - r, err = os.Open(c.base.Append( - dirIdentifier, - Encode(c.Ident(f).Value()), - ).String()) - c.identMu.RUnlock() - } - - if err != nil { - if !errors.Is(err, os.ErrNotExist) { - return - } - id := c.Ident(f) - if c.msg.IsVerbose() { - rn := reportName(f, id) - c.msg.Verbosef("curing %s%s%s in memory...", c.sgrIdent, rn, c.sgrRes) - defer func() { - if err == nil { - c.msg.Verbosef("opened %s%s%s for reading", c.sgrIdent, rn, c.sgrRes) - } - }() - } - return f.Cure(&RContext{common{ctx, c}}) - } - return -} - -// InvalidFileModeError describes a [FloodArtifact.Cure] or -// [TrivialArtifact.Cure] that did not result in a regular file or directory -// located at the work pathname. -type InvalidFileModeError fs.FileMode - -// Error returns a constant string. -func (e InvalidFileModeError) Error() string { - return "artifact did not produce a regular file or directory" -} - -// NoOutputError describes a [FloodArtifact.Cure] or [TrivialArtifact.Cure] -// that did not populate its work pathname despite completing successfully. -type NoOutputError struct{} - -// Unwrap returns [os.ErrNotExist]. -func (NoOutputError) Unwrap() error { return os.ErrNotExist } - -// Error returns a constant string. -func (NoOutputError) Error() string { - return "artifact cured successfully but did not produce any output" -} - -// removeAll is similar to [os.RemoveAll] but is robust against any permissions. -func removeAll(pathname *check.Absolute) (chmodErr, removeErr error) { - chmodErr = filepath.WalkDir(pathname.String(), func( - path string, - d fs.DirEntry, - err error, - ) error { - if err != nil { - return err - } - if d.IsDir() { - return os.Chmod(path, 0700) - } - return nil - }) - if errors.Is(chmodErr, os.ErrNotExist) { - chmodErr = nil - } - removeErr = os.RemoveAll(pathname.String()) - return -} - -// zeroTimes zeroes atime and mtime for the named file. -func zeroTimes(path string) (err error) { - // include/uapi/linux/fcntl.h - const ( - AT_FDCWD = -100 - AT_SYMLINK_NOFOLLOW = 0x100 - ) - _AT_FDCWD := AT_FDCWD - - var _p0 *byte - _p0, err = syscall.BytePtrFromString(path) - if err != nil { - return - } - if _, _, errno := syscall.Syscall6( - syscall.SYS_UTIMENSAT, - uintptr(_AT_FDCWD), - uintptr(unsafe.Pointer(_p0)), - uintptr(unsafe.Pointer(new([2]syscall.Timespec))), - AT_SYMLINK_NOFOLLOW, - 0, 0, - ); errno != 0 { - return os.NewSyscallError("utimensat", errno) - } - return -} - -// overrideFileInfo overrides the permission bits of [fs.FileInfo] to 0500 and -// is the concrete type returned by overrideFile.Stat. -type overrideFileInfo struct{ fs.FileInfo } - -// Mode returns [fs.FileMode] with its permission bits set to 0500. -func (fi overrideFileInfo) Mode() fs.FileMode { - return fi.FileInfo.Mode()&(^fs.FileMode(0777)) | 0500 -} - -// Sys returns nil to avoid passing the original permission bits. -func (fi overrideFileInfo) Sys() any { return nil } - -// overrideFile overrides the permission bits of [fs.File] to 0500 and is the -// concrete type returned by dotOverrideFS for calls with "." passed as name. -type overrideFile struct{ fs.File } - -func (f overrideFile) Stat() (fi fs.FileInfo, err error) { - fi, err = f.File.Stat() - if err != nil { - return - } - fi = overrideFileInfo{fi} - return -} - -// dirFS is implemented by the concrete type of the return value of [os.DirFS]. -type dirFS interface { - fs.StatFS - fs.ReadFileFS - fs.ReadDirFS - fs.ReadLinkFS -} - -// dotOverrideFS overrides the permission bits of "." to 0500 to avoid the extra -// system calls to add and remove write bit from the target directory. -type dotOverrideFS struct{ dirFS } - -// Open wraps the underlying [fs.FS] with "." special case. -func (fsys dotOverrideFS) Open(name string) (f fs.File, err error) { - f, err = fsys.dirFS.Open(name) - if err != nil || name != "." { - return - } - f = overrideFile{f} - return -} - -// Stat wraps the underlying [fs.FS] with "." special case. -func (fsys dotOverrideFS) Stat(name string) (fi fs.FileInfo, err error) { - fi, err = fsys.dirFS.Stat(name) - if err != nil || name != "." { - return - } - fi = overrideFileInfo{fi} - return -} - -// InvalidArtifactError describes an artifact that does not implement a -// supported Cure method. -type InvalidArtifactError ID - -func (e InvalidArtifactError) Error() string { - return "artifact " + Encode(e) + " cannot be cured" -} - -// Cure cures the [Artifact] and returns its pathname and [Checksum]. Direct -// calls to Cure are not subject to the cures limit. -func (c *Cache) Cure(a Artifact) ( - pathname *check.Absolute, - checksum unique.Handle[Checksum], - err error, -) { - c.abortMu.RLock() - defer c.abortMu.RUnlock() - - if err = c.toplevel.Load().ctx.Err(); err != nil { - return - } - - pathname, checksum, _, err = c.cure(a, true, false) - return -} - -// CureWhence is like Cure, but returns the whence value. -func (c *Cache) CureWhence(a Artifact) ( - pathname *check.Absolute, - checksum unique.Handle[Checksum], - whence int, - err error, -) { - c.abortMu.RLock() - defer c.abortMu.RUnlock() - - if err = c.toplevel.Load().ctx.Err(); err != nil { - return - } - - return c.cure(a, true, false) -} - -// CureNew is like Cure, but always enters the implementation. -func (c *Cache) CureNew(a Artifact) ( - pathname *check.Absolute, - checksum unique.Handle[Checksum], - err error, -) { - c.abortMu.RLock() - defer c.abortMu.RUnlock() - - if err = c.toplevel.Load().ctx.Err(); err != nil { - return - } - - var whence int -retry: - pathname, checksum, whence, err = c.cure(a, true, true) - if err != nil || whence == WNew { - return - } - goto retry -} - -// An InputError describes inputs of a [FloodArtifact] which had failed to cure. -type InputError map[Artifact]error - -// unwrap returns an iterator over sorted, deduplicated [Artifact] and their -// corresponding identifier. -func (e InputError) unwrap() iter.Seq2[Artifact, unique.Handle[ID]] { - ir := NewIR() - - type input struct { - a Artifact - id unique.Handle[ID] - } - p := make([]input, 0, len(e)) - for a := range e { - p = append(p, input{a, ir.Ident(a)}) - } - - var identBuf [2]ID - slices.SortFunc(p, func(a, b input) int { - identBuf[0], identBuf[1] = a.id.Value(), b.id.Value() - return slices.Compare(identBuf[0][:], identBuf[1][:]) - }) - p = slices.CompactFunc(p, func(a, b input) bool { return a.id == b.id }) - - return func(yield func(Artifact, unique.Handle[ID]) bool) { - for _, i := range p { - if !yield(i.a, i.id) { - return - } - } - } -} - -// Error returns a user-facing, deterministic text representation of e. -func (e InputError) Error() string { - var buf strings.Builder - buf.WriteString("errors curing inputs:") - for a, id := range e.unwrap() { - buf.WriteString("\n\t") - buf.WriteString(reportName(a, id)) - buf.WriteString(": ") - buf.WriteString(e[a].Error()) - } - return buf.String() -} - -// Unwrap returns a slice of underlying errors sorted by identifier. -func (e InputError) Unwrap() []error { - errs := make([]error, 0, len(e)) - for a := range e.unwrap() { - errs = append(errs, e[a]) - } - return errs -} - -// enterCure must be called before entering an [Artifact] implementation. -func (c *Cache) enterCure(a Artifact, curesExempt bool) error { - if c.attr.Notify != nil { - c.attr.Notify <- true - } - - if a.IsExclusive() { - c.exclMu.Lock() - } - if curesExempt { - return nil - } - - ctx := c.toplevel.Load().ctx - select { - case c.cures <- struct{}{}: - return nil - - case <-ctx.Done(): - if a.IsExclusive() { - c.exclMu.Unlock() - } - return ctx.Err() - } -} - -// exitCure must be called after exiting an [Artifact] implementation. -func (c *Cache) exitCure(a Artifact, curesExempt bool) { - if c.attr.Notify != nil { - c.attr.Notify <- false - } - - if a.IsExclusive() { - c.exclMu.Unlock() - } - if curesExempt { - return - } - - <-c.cures -} - -// measuredReader implements [io.ReadCloser] and measures the checksum during -// Close. If the underlying reader is not read to EOF, Close blocks until all -// remaining data is consumed and validated. -type measuredReader struct { - // Underlying reader. Never exposed directly. - r io.ReadCloser - // For validating checksum. Never exposed directly. - h hash.Hash - // Buffers writes to h, initialised by [Cache]. Never exposed directly. - hbw *bufio.Writer - // Expected checksum, compared during Close. - want unique.Handle[Checksum] - - // For accessing free lists. - c *Cache - - // Set up via [io.TeeReader] by [Cache]. - io.Reader -} - -// Close reads the underlying [io.ReadCloser] to EOF, closes it and measures its -// outcome. It returns a [ChecksumMismatchError] for an unexpected checksum. -func (mr *measuredReader) Close() (err error) { - if mr.hbw == nil || mr.Reader == nil { - return os.ErrInvalid - } - err = mr.hbw.Flush() - mr.c.putWriter(mr.hbw) - mr.hbw, mr.Reader = nil, nil - if err != nil { - _ = mr.r.Close() - return - } - var n int64 - if n, err = io.Copy(mr.h, mr.r); err != nil { - _ = mr.r.Close() - return - } - - if n > 0 { - mr.c.msg.Verbosef( - "%smissed %d bytes on measured reader%s", - mr.c.sgrWarn, n, mr.c.sgrRes, - ) - } - - if err = mr.r.Close(); err != nil { - return - } - - buf := mr.c.getIdentBuf() - mr.h.Sum(buf[:0]) - - if got := Checksum(buf[:]); got != mr.want.Value() { - err = &ChecksumMismatchError{ - Got: got, - Want: mr.want.Value(), - } - } - - mr.c.putIdentBuf(buf) - return -} - -// newMeasuredReader implements [RContext.NewMeasuredReader]. -func (c *Cache) newMeasuredReader( - r io.ReadCloser, - checksum unique.Handle[Checksum], -) io.ReadCloser { - mr := measuredReader{r: r, h: sha512.New384(), want: checksum, c: c} - mr.hbw = c.getWriter(mr.h) - mr.Reader = io.TeeReader(r, mr.hbw) - return &mr -} - -// NewMeasuredReader returns an [io.ReadCloser] implementing behaviour required -// by [FileArtifact]. The resulting [io.ReadCloser] holds a buffer originating -// from [Cache] and must be closed to return this buffer. -func (r *RContext) NewMeasuredReader( - rc io.ReadCloser, - checksum unique.Handle[Checksum], -) io.ReadCloser { - return r.cache.newMeasuredReader(rc, checksum) -} - -// tryChecksum dereferences a symlink to a cure outcome. -func (c *Cache) tryChecksum(pathname *check.Absolute) ( - checksum unique.Handle[Checksum], - err error, -) { - _, err = os.Lstat(pathname.String()) - if err == nil { - var name string - if name, err = os.Readlink(pathname.String()); err != nil { - return - } - buf := c.getIdentBuf() - err = Decode((*Checksum)(buf[:]), filepath.Base(name)) - if err == nil { - checksum = unique.Make(Checksum(buf[:])) - } - c.putIdentBuf(buf) - } - return -} - -// tryLocal attempts to obtain an [Artifact] outcome from the filesystem. -func (c *Cache) tryLocal(id unique.Handle[ID]) (unique.Handle[Checksum], error) { - return c.tryChecksum(c.base.Append( - dirIdentifier, - Encode(id.Value()), - )) -} - -// tryExtern attempts to obtain an [Artifact] outcome from extern. -func (c *Cache) tryExtern(ctx context.Context, id unique.Handle[ID]) ( - unique.Handle[Checksum], - error, -) { - c.externMu.RLock() - defer c.externMu.RUnlock() - - checksum, ok := c.externCache[id] - if !ok { - if c.extern == nil { - return zeroChecksum, nil - } - - v, err := c.extern.Artifact(ctx, id) - if err != nil { - return zeroChecksum, err - } - if v == nil { - return zeroChecksum, nil - } - checksum = unique.Make(*v) - - var got unique.Handle[Checksum] - if _, got, err = c.Cure(c.extern.Checksum(checksum)); err != nil { - return checksum, err - } else if got != checksum { - return zeroChecksum, &ChecksumMismatchError{got.Value(), checksum.Value()} - } - } - return checksum, nil -} - -// cureMany concurrently collects outcome of multiple [Artifact]. -func (c *Cache) cureMany( - inputs []Artifact, - r map[Artifact]cureRes, - shallow bool, -) ([]bool, error) { - var wg sync.WaitGroup - wg.Add(len(inputs)) - var mask []bool - res := make([]cureRes, len(inputs)) - errs := make(InputError) - var errsMu sync.Mutex - if shallow { - mask = make([]bool, len(inputs)) - } - for i, d := range inputs { - if shallow { - if _, ok := d.(FloodArtifact); ok { - mask[i] = true - wg.Done() - continue - } - - if kc, ok := d.(KnownChecksum); ok { - res[i].checksum = unique.Make(kc.Checksum()) - wg.Done() - continue - } - } - pending := pendingArtifactDep{d, &res[i], errs, &errsMu, &wg} - go pending.cure(c) - } - wg.Wait() - - if len(errs) > 0 { - return mask, errs - } - for i, p := range res { - if shallow && mask[i] { - continue - } - r[inputs[i]] = p - } - return mask, nil -} - -// HangingInputError describes an input of an [Artifact] on a sparse cache -// without an outcome available locally or via [External]. -type HangingInputError unique.Handle[ID] - -func (e HangingInputError) Error() string { - return Encode(unique.Handle[ID](e).Value()) + " is unavailable" -} - -const ( - // WNew indicates a cure entering the implementation. - WNew = iota - // WCache indicates an [Artifact] present in the cache. - WCache - // WSubstitute indicates an [Artifact] hitting a content-based input - // substitution. - WSubstitute - // WExternal indicates a cure entering the external cache. - WExternal -) - -// WhenceString returns a printable string for a whence value. -func WhenceString(whence int) string { - switch whence { - case WNew: - return "new" - case WCache: - return "cache" - case WSubstitute: - return "substitute" - case WExternal: - return "external" - - default: - return "invalid whence " + strconv.Itoa(whence) - } -} - -// cure implements Cure without acquiring a read lock on abortMu. cure must not -// be entered during Abort. -func (c *Cache) cure(a Artifact, curesExempt, rebuild bool) ( - pathname *check.Absolute, - checksum unique.Handle[Checksum], - whence int, - err error, -) { - id := c.Ident(a) - if rebuild { - var v ID - _, _ = rand.Read(v[:]) - id = unique.Make(v) - } - - ids := Encode(id.Value()) - pathname = c.base.Append( - dirIdentifier, - ids, - ) - defer func() { - if err != nil { - pathname = nil - checksum = unique.Handle[Checksum]{} - } - }() - - if _, ok := a.(CuresExempt); ok { - curesExempt = true - } - - var ( - ctx context.Context - done chan<- struct{} - ) - ctx, done, checksum, err = c.loadOrStoreIdent(id) - if done == nil { - whence = WCache - return - } else { - defer func() { c.finaliseIdent(done, id, checksum, err) }() - } - - checksum, err = c.tryChecksum(pathname) - if err == nil || !errors.Is(err, os.ErrNotExist) { - whence = WCache - return - } - - var ( - checksums string - substitute unique.Handle[ID] - alternative *check.Absolute - ) - defer func() { - if err == nil && checksums != "" { - linkname := checksumLinknamePrefix + checksums - - err = os.Symlink( - linkname, - pathname.String(), - ) - if err == nil { - err = zeroTimes(pathname.String()) - } - - if err == nil && alternative != nil && substitute != id { - c.substituteMu.Lock() - err = os.Symlink( - linkname, - alternative.String(), - ) - if errors.Is(err, os.ErrExist) { - c.msg.Verbosef( - "creating alternative over %s%s%s for artifact %s%s%s", - c.sgrIdent, Encode(substitute.Value()), c.sgrRes, - c.sgrIdent, ids, c.sgrRes, - ) - err = nil - } - if err == nil { - err = zeroTimes(alternative.String()) - } - if err == nil && checksum != zeroChecksum { - c.substitute[substitute] = checksum - } - c.substituteMu.Unlock() - } - } - }() - - var checksumPathname *check.Absolute - var checksumFi os.FileInfo - if kc, ok := a.(KnownChecksum); ok { - checksum = unique.Make(kc.Checksum()) - checksums = Encode(checksum.Value()) - checksumPathname = c.base.Append( - dirChecksum, - checksums, - ) - - if c.attr.Flags&CAssumeChecksum != 0 { - c.checksumMu.RLock() - checksumFi, err = os.Stat(checksumPathname.String()) - c.checksumMu.RUnlock() - - if err != nil { - if !errors.Is(err, os.ErrNotExist) { - return - } - - checksumFi, err = nil, nil - } - } - } - - whence = WNew - if c.msg.IsVerbose() { - rn := reportName(a, id) - c.msg.Verbosef("curing %s%s%s...", c.sgrIdent, rn, c.sgrRes) - defer func() { - if err != nil { - return - } - if checksums != "" { - c.msg.Verbosef( - "cured %s%s%s checksum %s%s%s", - c.sgrIdent, rn, c.sgrRes, - c.sgrIdent, checksums, c.sgrRes) - } else { - c.msg.Verbosef("cured %s%s%s", c.sgrIdent, rn, c.sgrRes) - } - }() - } - - // cure FileArtifact outside type switch to skip TContext initialisation - if f, ok := a.(FileArtifact); ok { - if checksumFi != nil { - whence = WCache - if !checksumFi.Mode().IsRegular() { - // unreachable - err = InvalidFileModeError(checksumFi.Mode()) - } - return - } - - perm := os.FileMode(0400) - if f.IsExecutable() { - perm = 0500 - } - - work := c.base.Append(dirWork, ids) - var w *os.File - if w, err = os.OpenFile( - work.String(), - os.O_CREATE|os.O_EXCL|os.O_WRONLY, - perm, - ); err != nil { - return - } - defer func() { - closeErr := w.Close() - if err == nil { - err = closeErr - } - - removeErr := os.Remove(work.String()) - if err == nil && !errors.Is(removeErr, os.ErrNotExist) { - err = removeErr - } - }() - - var r io.ReadCloser - if err = c.enterCure(a, curesExempt); err != nil { - return - } - r, err = f.Cure(&RContext{common{ctx, c}}) - if err == nil { - if checksumPathname == nil || c.attr.Flags&CValidateKnown != 0 { - h := sha512.New384() - hbw := c.getWriter(h) - _, err = io.Copy(w, io.TeeReader(r, hbw)) - flushErr := hbw.Flush() - c.putWriter(hbw) - if err == nil { - err = flushErr - } - - if err == nil { - buf := c.getIdentBuf() - h.Sum(buf[:0]) - - if checksumPathname == nil { - checksum = unique.Make(Checksum(buf[:])) - checksums = Encode(Checksum(buf[:])) - } else if c.attr.Flags&CValidateKnown != 0 { - if got := Checksum(buf[:]); got != checksum.Value() { - err = &ChecksumMismatchError{ - Got: got, - Want: checksum.Value(), - } - } - } - - c.putIdentBuf(buf) - - if checksumPathname == nil { - checksumPathname = c.base.Append( - dirChecksum, - checksums, - ) - } - } - } else { - _, err = io.Copy(w, r) - } - - closeErr := r.Close() - if err == nil { - err = closeErr - } - } - c.exitCure(a, curesExempt) - if err != nil { - if c.msg.IsVerbose() { - c.msg.Verbosef( - "cure file %s%s%s: %s%v%s", - c.sgrIdent, reportName(f, id), c.sgrRes, - c.sgrErr, err, c.sgrRes, - ) - } - return - } - - c.checksumMu.Lock() - if err = os.Rename( - work.String(), - checksumPathname.String(), - ); err != nil { - c.checksumMu.Unlock() - return - } - timeErr := zeroTimes(checksumPathname.String()) - c.checksumMu.Unlock() - - if err == nil { - err = timeErr - } - return - } - - if checksumFi != nil { - whence = WCache - if !checksumFi.Mode().IsDir() { - // unreachable - err = InvalidFileModeError(checksumFi.Mode()) - } - return - } - - t := TContext{ - c.base.Append(dirWork, ids), - c.base.Append(dirTemp, ids), - ids, nil, nil, nil, nil, - common{ctx, c}, - } - switch ca := a.(type) { - case TrivialArtifact: - defer t.destroy(&err) - if err = c.enterCure(a, curesExempt); err != nil { - return - } - err = ca.Cure(&t) - c.exitCure(a, curesExempt) - if err != nil { - if c.msg.IsVerbose() { - c.msg.Verbosef( - "cure trivial %s%s%s: %s%v%s", - c.sgrIdent, reportName(ca, id), c.sgrRes, - c.sgrErr, err, c.sgrRes, - ) - } - return - } - break - - case FloodArtifact: - var externChecksum unique.Handle[Checksum] - if externChecksum, err = c.tryExtern(ctx, id); err != nil { - if c.msg.IsVerbose() { - c.msg.Verbosef( - "extern %s%s%s: %s%v%s", - c.sgrIdent, reportName(ca, id), c.sgrRes, - c.sgrErr, err, c.sgrRes, - ) - } - return - } - extern := externChecksum != zeroChecksum - shallow := extern && c.attr.Flags&CExternShallow != 0 - - inputs := a.Inputs() - f := FContext{t, make(map[Artifact]cureRes, len(inputs))} - var mask []bool - if mask, err = c.cureMany(inputs, f.inputs, shallow); err != nil { - return - } - - if shallow { - for i, d := range inputs { - if !mask[i] { - continue - } - - if kc, ok := d.(KnownChecksum); ok { - f.inputs[d] = cureRes{checksum: unique.Make(kc.Checksum())} - continue - } - - var sum unique.Handle[Checksum] - did := c.Ident(d) - sum, err = c.tryLocal(did) - if err != nil { - if !errors.Is(err, os.ErrNotExist) { - return - } - - sum, err = c.tryExtern(ctx, did) - if err != nil { - return - } - if sum == zeroChecksum { - if c.msg.IsVerbose() { - c.msg.Verbosef( - "input %s%s%s not available", - c.sgrIdent, reportName(d, did), c.sgrRes, - ) - } - err = HangingInputError(did) - return - } - } - - f.inputs[d] = cureRes{checksum: sum} - } - } - - sh := sha512.New384() - err = c.encode(sh, a, f.inputs) - if err != nil { - return - } - - buf := c.getIdentBuf() - sh.Sum(buf[wordSize:wordSize]) - substitute = unique.Make(ID(buf[wordSize:])) - substitutes := Encode(substitute.Value()) - c.putIdentBuf(buf) - alternative = c.base.Append( - dirSubstitute, - substitutes, - ) - - if !rebuild && c.attr.Flags&CIgnoreSubstitutes == 0 { - var substituteChecksum unique.Handle[Checksum] - substituteChecksum, err = c.loadSubstitute(substitute) - if err != nil { - return - } - if substituteChecksum != zeroChecksum { - whence = WSubstitute - checksum = substituteChecksum - checksums = Encode(checksum.Value()) - checksumPathname = c.base.Append( - dirChecksum, - checksums, - ) - if _, err = os.Lstat(c.base.Append( - dirStatus, - substitutes, - ).String()); err == nil { - err = os.Symlink(substitutes, c.base.Append( - dirStatus, - ids, - ).String()) - } else if errors.Is(err, os.ErrNotExist) { - err = nil - } - return - } - } - - defer f.destroy(&err) - - if extern { - whence = WExternal - if checksum != zeroChecksum && externChecksum != checksum { - err = &ChecksumMismatchError{externChecksum.Value(), checksum.Value()} - if c.msg.IsVerbose() { - c.msg.Verbosef( - "extern %s%s%s: %s%v%s", - c.sgrIdent, reportName(ca, id), c.sgrRes, - c.sgrErr, err, c.sgrRes, - ) - } - return - } - - var externStatus io.ReadCloser - c.externMu.RLock() - externStatus, err = c.extern.Status(&RContext{common{ctx, c}}, id) - c.externMu.RUnlock() - if err != nil { - return - } - - checksum = externChecksum - checksums = Encode(checksum.Value()) - checksumPathname = c.base.Append( - dirChecksum, - checksums, - ) - - if externStatus != nil { - if err = f.prepareStatus(false); err != nil { - _ = externStatus.Close() - return - } else if _, err = io.Copy(f.status, externStatus); err != nil { - _ = externStatus.Close() - return - } else if err = externStatus.Close(); err != nil { - return - } else if !rebuild { - if err = f.linkSubstitute(ids, substitutes); err != nil { - return - } - } - } - return - } - - if err = c.enterCure(a, curesExempt); err != nil { - return - } - err = ca.Cure(&f) - c.exitCure(a, curesExempt) - - if !rebuild && err == nil { - err = f.linkSubstitute(ids, substitutes) - } - if err != nil { - if c.msg.IsVerbose() { - c.msg.Verbosef( - "cure %s%s%s: %s%v%s", - c.sgrIdent, reportName(ca, id), c.sgrRes, - c.sgrErr, err, c.sgrRes, - ) - } - return - } - break - - default: - err = InvalidArtifactError(id.Value()) - return - } - t.cache = nil - - var fi os.FileInfo - if fi, err = os.Lstat(t.work.String()); err != nil { - if errors.Is(err, os.ErrNotExist) { - err = NoOutputError{} - } - return - } - - if !fi.IsDir() { - if !fi.Mode().IsRegular() { - err = InvalidFileModeError(fi.Mode()) - } else { - err = errors.New("non-file artifact produced regular file") - } - return - } - - var gotChecksum Checksum - if err = SumFS( - &gotChecksum, - dotOverrideFS{os.DirFS(t.work.String()).(dirFS)}, - ".", - ); err != nil { - return - } - - if checksumPathname == nil { - checksum = unique.Make(gotChecksum) - checksums = Encode(gotChecksum) - checksumPathname = c.base.Append( - dirChecksum, - checksums, - ) - } else if gotChecksum != checksum.Value() { - err = &ChecksumMismatchError{ - Got: gotChecksum, - Want: checksum.Value(), - } - if c.msg.IsVerbose() { - c.msg.Verbosef( - "validate %s%s%s: %s%v%s", - c.sgrIdent, reportName(a, id), c.sgrRes, - c.sgrErr, err, c.sgrRes, - ) - } - return - } - - if err = os.Chmod(t.work.String(), 0700); err != nil { - return - } - if err = filepath.WalkDir(t.work.String(), func(path string, _ fs.DirEntry, err error) error { - if err != nil { - return err - } - return zeroTimes(path) - }); err != nil { - return - } - c.checksumMu.Lock() - if err = os.Rename( - t.work.String(), - checksumPathname.String(), - ); err != nil { - if !errors.Is(err, os.ErrExist) { - c.checksumMu.Unlock() - return - } - // err is zeroed during deferred cleanup - } else { - err = os.Chmod(checksumPathname.String(), 0500) - } - c.checksumMu.Unlock() - return -} - -// cure cures the pending [Artifact], stores its result and notifies the caller. -func (pending *pendingArtifactDep) cure(c *Cache) { - defer pending.Done() - - var err error - pending.resP.pathname, pending.resP.checksum, _, err = c.cure(pending.a, false, false) - if err == nil { - return - } - - pending.errsMu.Lock() - if errs, ok := err.(InputError); ok { - maps.Copy(pending.errs, errs) - } else { - pending.errs[pending.a] = err - } - pending.errsMu.Unlock() -} - -// OpenStatus attempts to open the status file associated to an [Artifact]. If -// err is nil, the caller must close the resulting reader. -func (c *Cache) OpenStatus(a Artifact) (r io.ReadSeekCloser, err error) { - c.identMu.RLock() - r, err = os.Open(c.base.Append( - dirStatus, - Encode(c.Ident(a).Value())).String(), - ) - c.identMu.RUnlock() - return -} - -// Fault holds the pathname and termination time of an [Artifact] fault entry. -type Fault struct { - *check.Absolute - t uint64 -} - -// Time returns the instant in time where the fault occurred. -func (f Fault) Time() time.Time { return time.Unix(0, int64(f.t)) } - -// Open opens the underlying entry for reading. -func (f Fault) Open() (io.ReadCloser, error) { return os.Open(f.Absolute.String()) } - -// Destroy removes the underlying fault entry. -func (f Fault) Destroy() error { return os.Remove(f.Absolute.String()) } - -// ReadFaults returns fault entries for an [Artifact]. -func (c *Cache) ReadFaults(a Artifact) (faults []Fault, err error) { - prefix := Encode(c.Ident(a).Value()) + "." - var dents []os.DirEntry - if dents, err = os.ReadDir(c.base.Append(dirFault).String()); err != nil { - return - } - - for _, dent := range dents { - name := dent.Name() - if !strings.HasPrefix(name, prefix) { - continue - } - var t uint64 - t, err = strconv.ParseUint(name[len(prefix):], 10, 64) - if err != nil { - return - } - - faults = append(faults, Fault{c.base.Append( - dirFault, - name, - ), t}) - } - - slices.SortFunc(faults, func(a, b Fault) int { - return cmp.Compare(a.t, b.t) - }) - return -} - -// Abort cancels all pending cures and waits for them to clean up, but does not -// close the cache. -func (c *Cache) Abort() { - c.closeMu.Lock() - defer c.closeMu.Unlock() - - if c.closed { - return - } - - c.toplevel.Load().cancel() - c.abortMu.Lock() - defer c.abortMu.Unlock() - - // holding abortMu, identPending stays empty - c.wg.Wait() - c.identMu.Lock() - c.toplevel.Store(newToplevel(c.parent)) - clear(c.identErr) - c.identMu.Unlock() -} - -// Close cancels all pending cures and waits for them to clean up. -func (c *Cache) Close() { - c.closeMu.Lock() - defer c.closeMu.Unlock() - - if c.closed { - return - } - - c.closed = true - c.toplevel.Load().cancel() - c.wg.Wait() - close(c.cures) - c.unlock() - - if c.attr.Notify != nil { - close(c.attr.Notify) - } -} - -// UnsupportedVariantError describes an on-disk cache with an extension variant -// identification string that differs from the value returned by [Extension]. -type UnsupportedVariantError string - -func (e UnsupportedVariantError) Error() string { - return "unsupported variant " + strconv.Quote(string(e)) -} - -var ( - // ErrWouldPromote is returned by [Open] if the [CPromoteVariant] bit is not - // set and the on-disk cache requires variant promotion. - ErrWouldPromote = errors.New("operation would promote unextended cache") -) - -// CacheAttr holds the attributes that will be applied to a new [Cache] opened -// by [Open]. -type CacheAttr struct { - // Concurrent cures of a [FloodArtifact] dependency graph. - Cures int - // Options affecting [Cache] behaviour. - Flags int - // Preferred job count, when applicable. - Jobs int - // Preferred loadavg target, when applicable. - Load int - // Optional cure entry and exit notification. - Notify chan<- bool - - // Omit the [lockedfile] lock. - skipLock bool -} - -// Open returns the address of a newly opened instance of [Cache]. -// -// Concurrent cures of a [FloodArtifact] dependency graph is limited to the -// caller-supplied value, however direct calls to [Cache.Cure] is not subject -// to this limitation. -// -// A cures or jobs value of 0 or lower is equivalent to the value returned by -// [runtime.NumCPU]. -// -// A successful call to Open guarantees exclusive access to the on-filesystem -// cache for the resulting instance of [Cache]. The [Cache.Close] method cancels -// and waits for pending cures on [Cache] before releasing this lock and must be -// called once the [Cache] is no longer needed. -func Open( - ctx context.Context, - msg message.Msg, - base *check.Absolute, - attr *CacheAttr, -) (*Cache, error) { - openMu.Lock() - defer openMu.Unlock() - opened = true - - if extension == "" && len(irArtifact) != int(_kindEnd) { - panic("attempting to open cache with incomplete variant setup") - } - - var a CacheAttr - if attr != nil { - a = *attr - } - - if a.Cures < 1 { - a.Cures = runtime.NumCPU() - } - if a.Jobs < 1 { - a.Jobs = runtime.NumCPU() - } - if a.Load < 1 { - a.Load = runtime.NumCPU() + 2 - } - - for _, name := range []string{ - dirSubstitute, - dirIdentifier, - dirChecksum, - dirStatus, - dirFault, - dirWork, - } { - if err := os.MkdirAll( - base.Append(name).String(), - 0700, - ); err != nil && !errors.Is(err, os.ErrExist) { - return nil, err - } - } - - c := Cache{ - parent: ctx, - - cures: make(chan struct{}, a.Cures), - attr: a, - - msg: msg, - base: base, - - irCache: zeroIRCache(), - - substitute: make(map[unique.Handle[ID]]unique.Handle[Checksum]), - ident: make(map[unique.Handle[ID]]unique.Handle[Checksum]), - identErr: make(map[unique.Handle[ID]]error), - identPending: make(map[unique.Handle[ID]]*pendingCure), - - brPool: sync.Pool{New: func() any { return new(bufio.Reader) }}, - bwPool: sync.Pool{New: func() any { return new(bufio.Writer) }}, - } - c.toplevel.Store(newToplevel(ctx)) - - if !a.skipLock || !testing.Testing() { - if unlock, err := lockedfile.MutexAt( - base.Append(fileLock).String(), - ).Lock(); err != nil { - return nil, err - } else { - c.unlock = unlock - } - } else { - c.unlock = func() {} - } - - for _, name := range []string{ - dirWork, - dirTemp, - } { - dents, err := os.ReadDir(base.Append(name).String()) - if err != nil { - if errors.Is(err, os.ErrNotExist) { - continue - } - c.unlock() - return nil, err - } - if len(dents) != 0 { - c.unlock() - return nil, fmt.Errorf( - "%s is not empty, scrub likely required", - name, - ) - } - } - - if _, err := os.ReadDir(base.Append( - dirExecScratch, - ).String()); !errors.Is(err, os.ErrNotExist) { - c.unlock() - if err != nil { - return nil, err - } - return nil, errors.New(dirExecScratch + " is present, scrub likely required") - } - - variantPath := base.Append(fileVariant).String() - if p, err := os.ReadFile(variantPath); err != nil { - if !errors.Is(err, os.ErrNotExist) { - c.unlock() - return nil, err - } - // nonexistence implies newly created cache, or a cache predating - // variant identification strings, in which case it is silently promoted - if err = os.WriteFile( - variantPath, - []byte(extension), - 0400, - ); err != nil { - c.unlock() - return nil, err - } - } else if s := string(p); s == "" { - if extension != "" { - if a.Flags&CPromoteVariant == 0 { - c.unlock() - return nil, ErrWouldPromote - } - if err = os.WriteFile( - variantPath, - []byte(extension), - 0400, - ); err != nil { - c.unlock() - return nil, err - } - } - } else if !ValidExtension(s) { - c.unlock() - return nil, ErrInvalidExtension - } else if s != extension { - c.unlock() - return nil, UnsupportedVariantError(s) - } - - if a.Flags&CColourOutput != 0 { - c.sgrRes = "\x1b[0m" - c.sgrIdent = "\x1b[1m" - c.sgrWarn = "\x1b[35m" - c.sgrErr = "\x1b[1;31m" - } - - return &c, nil -} - -// Collected is returned by [Collect.Cure] to indicate a successful collection. -type Collected struct{} - -// Error returns a constant string to satisfy error, but should never be seen -// by the user. -func (Collected) Error() string { return "artifacts successfully collected" } - -// IsCollected returns whether the underlying error contains that of the result -// of curing a [Collect] helper. -func IsCollected(err error) bool { return errors.As(err, new(Collected)) } - -// Collect implements [pkg.FloodArtifact] to concurrently cure multiple -// [pkg.Artifact]. It returns [Collected]. -type Collect []Artifact - -var _ Artifact = new(Collect) - -// Cure returns [Collected]. -func (*Collect) Cure(*FContext) error { return Collected{} } - -// Kind returns the hardcoded [pkg.Kind] value. -func (*Collect) Kind() Kind { return kindCollection } - -// Params is a noop: dependencies are already represented in the header. -func (*Collect) Params(*IContext) {} - -// Inputs returns [Collect] as is. -func (c *Collect) Inputs() []Artifact { return *c } - -// IsExclusive returns false: Cure is a noop. -func (*Collect) IsExclusive() bool { return false } |
