aboutsummaryrefslogtreecommitdiffhomepage
path: root/pkg/pkg.go
diff options
context:
space:
mode:
Diffstat (limited to 'pkg/pkg.go')
-rw-r--r--pkg/pkg.go2952
1 files changed, 2952 insertions, 0 deletions
diff --git a/pkg/pkg.go b/pkg/pkg.go
new file mode 100644
index 00000000..5d9d8d48
--- /dev/null
+++ b/pkg/pkg.go
@@ -0,0 +1,2952 @@
+// Package pkg provides low-level primitives for packaging software.
+//
+// The public interface, IR format, and IR representation of [Artifact]
+// implementations not satisfying [RevisionArtifact] are covered by the
+// compatibility promise. [RevisionArtifact] revisions may change between any
+// two releases.
+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 = "hakurei.app/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 }