diff options
| -rw-r--r-- | cmd/mbf/internal/ci/ci.go | 522 | ||||
| -rw-r--r-- | cmd/mbf/internal/ci/ci_test.go | 55 | ||||
| -rw-r--r-- | cmd/mbf/main.go | 37 | ||||
| -rw-r--r-- | internal/rosa/llvm_test.go | 2 | ||||
| -rw-r--r-- | internal/rosa/mirror.go | 1 | ||||
| -rw-r--r-- | internal/rosa/package/gnu.az | 18 | ||||
| -rw-r--r-- | internal/rosa/package/gzip/late-include.patch | 21 | ||||
| -rw-r--r-- | internal/rosa/package/gzip/package.az | 18 | ||||
| -rw-r--r-- | pkg/archive.go | 8 | ||||
| -rw-r--r-- | pkg/pkg.go | 3 |
10 files changed, 660 insertions, 25 deletions
diff --git a/cmd/mbf/internal/ci/ci.go b/cmd/mbf/internal/ci/ci.go new file mode 100644 index 00000000..47c63d67 --- /dev/null +++ b/cmd/mbf/internal/ci/ci.go @@ -0,0 +1,522 @@ +// Package ci implements the CI service and client. +package ci + +import ( + "archive/tar" + "bytes" + "compress/gzip" + "context" + "crypto/ed25519" + "encoding/binary" + "errors" + "io" + "io/fs" + "log" + "net" + "net/http" + "os" + "path" + "path/filepath" + "sync" + "syscall" + "time" + "unique" + "unsafe" + _ "unsafe" // for go:linkname + + "hakurei.app/internal/rosa" + "hakurei.app/message" + "hakurei.app/pkg" +) + +// setSource is made available here to accept prepared hakurei tarballs. +// +//go:linkname setSource hakurei.app/internal/rosa.(*S).setSource +func setSource(s *rosa.S, p []byte, version string) + +// inotifyInit returns a new inotify instance. +func inotifyInit() (*os.File, error) { + fd, err := syscall.InotifyInit1(syscall.IN_NONBLOCK | syscall.IN_CLOEXEC) + if err != nil { + return nil, os.NewSyscallError("inotify_init1", err) + } + return os.NewFile(uintptr(fd), "inotify"), nil +} + +// inotifyAddWatch adds pathname to in. +func inotifyAddWatch( + in *os.File, + pathname string, + mask uint32, +) (watchdesc int, err error) { + sc, _err := in.SyscallConn() + if _err != nil { + return -1, _err + } + if _err = sc.Control(func(fd uintptr) { + watchdesc, err = syscall.InotifyAddWatch(int(fd), pathname, mask) + }); _err != nil { + return -1, _err + } + return +} + +// Follow reads from the file at pathname and writes its contents and any new +// contents to w. follow returns if a read, write or inotify error occurs, or +// the context is canceled. +func Follow(ctx context.Context, pathname string, w io.Writer) (err error) { + var in *os.File + if in, err = inotifyInit(); err != nil { + return + } + defer func() { + if _err := in.Close(); err == nil { + err = _err + } + }() + + if _, err = inotifyAddWatch(in, pathname, syscall.IN_MODIFY); err != nil { + return + } + var r *os.File + if r, err = os.Open(pathname); err != nil { + return + } + defer func() { + if _err := r.Close(); err == nil { + err = _err + } + }() + + done := make(chan struct{}) + defer close(done) + go func() { + select { + case <-ctx.Done(): + now := time.Now() + _ = in.SetDeadline(now) + return + + case <-done: + return + } + }() + + if _, err = io.Copy(w, r); err != nil { + return + } + + buf := make([]byte, os.Getpagesize()) + for { + if _, err = io.Copy(w, r); err != nil { + return + } + if _, err = in.Read(buf); err != nil { + if errors.Is(err, os.ErrDeadlineExceeded) { + err = nil + } + return + } + } +} + +// Path wraps the [pkg.Cache] pathname. +type Path string + +// String returns the value of p. +func (p Path) String() string { return string(p) } + +// append is [filepath.Join] with p as the first element. +func (p Path) append(elem ...string) string { + return filepath.Join(append([]string{p.String()}, elem...)...) +} + +// New returns a new [Path]. +func New(c *pkg.Cache) Path { return Path(c.Path().String()) } + +// name returns the CI socket pathname. +func (p Path) name() string { return p.append("ci") } + +// client returns the CI http client. +func (p Path) client() *http.Client { + var d net.Dialer + addr := net.UnixAddr{ + Net: "unix", + Name: p.name(), + } + + return &http.Client{Transport: &http.Transport{ + DialContext: func(ctx context.Context, _, _ string) (net.Conn, error) { + return d.DialUnix(ctx, "unix", nil, &addr) + }, + }} +} + +// cure writes the identifier of the pending artifact, cures the artifact, and +// writes the cure whence. For an unsuccessful cure, a negative whence is +// written, followed by a user-facing error string. +func cure(c *pkg.Cache, a pkg.Artifact, w http.ResponseWriter) error { + h := w.Header() + h.Set("Content-Type", "application/octet-stream") + h.Set("Cache-Control", "no-cache") + + f, ok := w.(http.Flusher) + if !ok { + _, _ = w.Write([]byte{0}) + return errors.ErrUnsupported + } + + id := c.Ident(a).Value() + if _, err := w.Write(id[:]); err != nil { + return err + } + f.Flush() + + _, _, whence, err := c.CureWhence(a) + if err != nil { + whence = -1 + } + if _, _err := w.Write( + binary.LittleEndian.AppendUint64(nil, uint64(whence)), + ); _err != nil { + return _err + } + + if err != nil { + if _, _err := io.WriteString(w, err.Error()); _err != nil { + return _err + } + } + f.Flush() + return err +} + +// spool holds reusable [rosa.S] instances. +var spool = sync.Pool{New: func() any { return rosa.New() }} + +// getS returns the address of a populated [rosa.S]. Its hakurei-source may be +// clobbered and must be replaced using setSource before use. +func getS() *rosa.S { return spool.Get().(*rosa.S) } + +// putS returns s to spool. +func putS(s *rosa.S) { spool.Put(s) } + +// versionSize is the maximum size of the specified version string, plus its +// deliminator byte. +const versionSize = 8 + 16 + 6 + 2 + +// errBadVersion is returned by readSource if a header does not contain +// the deliminator byte. +var errBadVersion = errors.New("unterminated version string") + +// readSource reads a version string and compressed source tarball from r and +// returns the address of a [rosa.S] with this source tarball. The resulting +// [rosa.S] must be returned via putS. +func readSource(w http.ResponseWriter, r *http.Request) (*rosa.S, error) { + var header [versionSize]byte + _, err := io.ReadFull(r.Body, header[:]) + if err != nil { + _ = r.Body.Close() + http.Error(w, "bad header", http.StatusBadRequest) + return nil, err + } + + var version string + if i := bytes.IndexByte(header[:], 0); i < 0 { + _ = r.Body.Close() + http.Error(w, "unterminated version string", http.StatusBadRequest) + return nil, errBadVersion + } else { + version = unsafe.String(&header[0], i) + } + + var p []byte + if p, err = io.ReadAll(r.Body); err != nil { + _ = r.Body.Close() + http.Error(w, "cannot receive payload", http.StatusInternalServerError) + return nil, err + } + + s := getS() + setSource(s, p, version) + return s, r.Body.Close() +} + +// ErrDaemonError is returned by writeSource generally if cure could not flush +// on the connection to notify completion. +var ErrDaemonError = errors.New("CI service could not process the request") + +// writeSource writes a source tarball to the specified endpoint of the CI +// backend servicing the cache referred to by cm. +func (p Path) writeSource( + ctx context.Context, + w io.Writer, + endpoint, source, version string, +) (*pkg.ID, error) { + if len(version) >= versionSize { + return nil, syscall.ENOMEM + } + var header [versionSize]byte + copy(header[:], version[:]) + + var buf bytes.Buffer + gw := gzip.NewWriter(&buf) + tw := tar.NewWriter(gw) + + if err := filepath.WalkDir(source, func(path string, d fs.DirEntry, err error) error { + if err != nil { + return err + } + + if d.IsDir() && d.Name() == ".git" { + return fs.SkipDir + } + + var fi fs.FileInfo + if fi, err = d.Info(); err != nil { + return err + } + + var linkname string + if fi.Mode()&fs.ModeSymlink != 0 { + if linkname, err = os.Readlink(path); err != nil { + return err + } + } + + var h *tar.Header + if h, err = tar.FileInfoHeader(fi, linkname); err != nil { + return err + } + h.Name = path + + var isVersion bool + if dir, file := filepath.Split(path); filepath.Base(dir) == "dist" && + file == "VERSION" && fi.Mode().IsRegular() { + isVersion = true + h.Size = int64(len(version)) + } + + if err = tw.WriteHeader(h); err != nil { + return err + } + + if isVersion { + _, err = io.WriteString(tw, version) + return err + } + + if fi.Mode().IsRegular() { + var f io.ReadCloser + if f, err = os.Open(path); err != nil { + return err + } + + _, err = io.Copy(tw, f) + if _err := f.Close(); err == nil { + err = _err + } + if err != nil { + return err + } + } + return nil + }); err != nil { + return nil, err + } + if err := tw.Close(); err != nil { + return nil, err + } + if err := gw.Close(); err != nil { + return nil, err + } + + req, err := http.NewRequestWithContext( + ctx, + http.MethodPost, + "http://"+path.Join("_", endpoint), + io.MultiReader(bytes.NewReader(header[:]), bytes.NewReader(buf.Bytes())), + ) + if err != nil { + return nil, err + } + + var resp *http.Response + resp, err = p.client().Do(req) + if err != nil { + return nil, err + } + + var id pkg.ID + _, err = io.ReadFull(resp.Body, id[:]) + if err != nil { + _ = resp.Body.Close() + if errors.Is(err, io.ErrUnexpectedEOF) { + return nil, ErrDaemonError + } + return nil, err + } + + c, cancel := context.WithCancel(ctx) + done := make(chan error, 1) + var whence int + go func() { + defer cancel() + var wbuf [8]byte + _, _err := io.ReadFull(resp.Body, wbuf[:]) + whence = int(binary.LittleEndian.Uint64(wbuf[:])) + done <- _err + }() + + if w != nil { + retry: + err = Follow(c, filepath.Join(p.String(), "status", pkg.Encode(id)), w) + if err != nil { + if ctx.Err() == nil && errors.Is(err, os.ErrNotExist) { + goto retry + } + + _ = resp.Body.Close() + return nil, err + } + } + + err = <-done + if err != nil { + _ = resp.Body.Close() + return nil, err + } + + if whence < 0 { + var m []byte + if m, err = io.ReadAll(resp.Body); err != nil { + _ = resp.Body.Close() + return nil, err + } else if err = resp.Body.Close(); err != nil { + return nil, err + } + return nil, errors.New(unsafe.String(unsafe.SliceData(m), len(m))) + } + return &id, resp.Body.Close() +} + +// The stubKey is used by the mirror service exposed by serve where +// authentication is unnecessary. +var stubKey = ed25519.NewKeyFromSeed(make([]byte, ed25519.SeedSize)) + +// fetch fetches the outcome of id and writes it to the specified directory. +func (p Path) fetch(ctx context.Context, id *pkg.ID, output string) error { + c := p.client() + r, err := rosa.NewRemote( + c, "http://_", + stubKey.Public().(ed25519.PublicKey), + ) + if err != nil { + return err + } + + var sum *pkg.Checksum + if sum, err = r.Artifact(ctx, unique.Make(*id)); err != nil { + return err + } + + var req *http.Request + if req, err = http.NewRequestWithContext( + ctx, + http.MethodGet, + "http://"+path.Join("_", "outcome", pkg.Encode(*sum)), + nil, + ); err != nil { + return err + } + + var resp *http.Response + if resp, err = c.Do(req); err != nil { + return err + } + + err = pkg.Extract(resp.Body, output, nil) + if closeErr := resp.Body.Close(); err == nil { + err = closeErr + } + return err +} + +// MakeDist creates a hakurei distribution using the CI service. +func (p Path) MakeDist( + ctx context.Context, + w io.Writer, + output, source, version string, +) error { + id, err := p.writeSource(ctx, w, "/dist", source, version) + if err != nil { + return err + } + return p.fetch(ctx, id, output) +} + +// Serve services CI workload dispatched to c. +func Serve(ctx context.Context, msg message.Msg, c *pkg.Cache) error { + const shutdownTimeout = 15 * time.Second + p := New(c) + addr := net.UnixAddr{ + Net: "unix", + Name: p.name(), + } + + var mux http.ServeMux + mux.HandleFunc("POST /dist", func(w http.ResponseWriter, r *http.Request) { + s, err := readSource(w, r) + if err != nil { + msg.Verbose(err) + return + } + defer putS(s) + + _, a := s.Std().MustLoad(rosa.H("hakurei-dist")) + if err = cure(c, a, w); err != nil { + msg.Verbose(err) + return + } + if msg.IsVerbose() { + msg.Verbosef( + "satisfied distribution %s", + pkg.Encode(c.Ident(a).Value()), + ) + } + }) + + if r, err := os.OpenRoot(p.String()); err != nil { + return err + } else { + defer func() { + if err = r.Close(); err != nil { + msg.Verbose(err) + } + }() + rosa.NewMirror(msg, r.FS(), stubKey).Register(&mux) + } + + server := http.Server{Handler: &mux} + go func() { + <-ctx.Done() + cc, cancel := context.WithTimeout(context.Background(), shutdownTimeout) + defer cancel() + if _err := server.Shutdown(cc); _err != nil { + log.Fatal(_err) + } + }() + + ul, err := net.ListenUnix("unix", &addr) + if err != nil { + return err + } + ul.SetUnlinkOnClose(true) + msg.Verbosef("listening on %s", addr.Net) + + err = server.Serve(ul) + if errors.Is(err, http.ErrServerClosed) { + err = nil + } + return err +} diff --git a/cmd/mbf/internal/ci/ci_test.go b/cmd/mbf/internal/ci/ci_test.go new file mode 100644 index 00000000..b5ed134e --- /dev/null +++ b/cmd/mbf/internal/ci/ci_test.go @@ -0,0 +1,55 @@ +package ci_test + +import ( + "bytes" + "context" + "os" + "path/filepath" + "testing" + + "hakurei.app/cmd/mbf/internal/ci" +) + +func TestFollow(t *testing.T) { + t.Parallel() + + pathname := filepath.Join(t.TempDir(), "f") + w, err := os.Create(pathname) + if err != nil { + t.Fatal(err) + } + + var buf bytes.Buffer + ctx, cancel := context.WithCancel(t.Context()) + + var want string + go func() { + defer cancel() + for _, s := range []string{ + "\xde\xad\xbe\xef", + "\xff\xff\xff\xff", + "\x00\x00", + } { + want += s + if _, _err := w.WriteString(s); _err != nil { + panic(_err) + } + } + }() + +retry: + if err = ci.Follow(ctx, pathname, &buf); err != nil { + t.Fatal(err) + } + <-ctx.Done() + + // the inotify event takes time to arrive, and there is no way to + // synchronise for this cleanly + if buf.Len() != len(want) { + goto retry + } + + if got := buf.String(); got != want { + t.Fatalf("follow: %q, want %q", got, want) + } +} diff --git a/cmd/mbf/main.go b/cmd/mbf/main.go index 22bfbc3e..014a7f47 100644 --- a/cmd/mbf/main.go +++ b/cmd/mbf/main.go @@ -43,6 +43,7 @@ import ( "hakurei.app/message" "hakurei.app/pkg" + "hakurei.app/cmd/mbf/internal/ci" "hakurei.app/cmd/mbf/internal/pkgsite" "hakurei.app/cmd/mbf/internal/pkgsite/ui" ) @@ -504,7 +505,7 @@ func main() { c.NewCommand( "daemon", "Service artifact IR with Rosa OS extensions", - func(args []string) error { + func([]string) error { ul, err := net.ListenUnix("unix", &addr) if err != nil { return err @@ -514,6 +515,40 @@ func main() { }, ) + _ci := c.New("ci", command.UsageInternal) + _ci.NewCommand( + "daemon", + "Service CI workload dispatched through the socket", + func([]string) error { + return cm.Do(func(cache *pkg.Cache) error { + return ci.Serve(ctx, msg, cache) + }) + }, + ) + { + var flagOutput string + _ci.NewCommand( + "dist", + "Request distribution tarball for the specified source directory", + func(args []string) error { + if len(args) != 2 { + return errors.New("dist requires 2 arguments") + } + + return ci.Path(cm.base).MakeDist( + ctx, + os.Stdout, + flagOutput, + args[0], args[1], + ) + }, + ).Flag( + &flagOutput, + "o", command.StringFlag("."), + "Write the resulting distribution to the named directory", + ) + } + c.NewCommand( "keygen", "Create keypair for local cache", diff --git a/internal/rosa/llvm_test.go b/internal/rosa/llvm_test.go index 3152107c..3fa8fc9d 100644 --- a/internal/rosa/llvm_test.go +++ b/internal/rosa/llvm_test.go @@ -8,7 +8,7 @@ import ( ) func TestLLVMInputs(t *testing.T) { - const wantInputCount = 466 + const wantInputCount = 472 _, llvm := rosa.MustLoad(rosa.H("llvm")) var n int diff --git a/internal/rosa/mirror.go b/internal/rosa/mirror.go index 4af931c7..20dd3231 100644 --- a/internal/rosa/mirror.go +++ b/internal/rosa/mirror.go @@ -451,7 +451,6 @@ func writeF[B any]( return } m.msg.GetLogger().Println(err) - w.WriteHeader(http.StatusInternalServerError) } // Register configures an [http.ServeMux] for servicing mirror requests. diff --git a/internal/rosa/package/gnu.az b/internal/rosa/package/gnu.az index c3f1be64..2c408206 100644 --- a/internal/rosa/package/gnu.az +++ b/internal/rosa/package/gnu.az @@ -125,24 +125,6 @@ package libtool { ]; } -package gzip { - description = "a popular data compression program"; - website = "https://www.gnu.org/software/gzip"; - anitya = 1290; - - version# = "1.15"; - source = remoteTar { - url = "https://mirrors.kernel.org/gnu/gzip/gzip-"+version+".tar.gz"; - checksum = "N8v7cH-Vhmxvk1PYaF-B1vH-r-rMT3Z3VyMdCFXhBd74a8v62wT3dKes-u1JCvUG"; - compress = gzip; - }; - - exec = make { - // dependency loop - check = nil; - }; -} - package sed { description = "a non-interactive command-line text editor"; website = "https://www.gnu.org/software/sed"; diff --git a/internal/rosa/package/gzip/late-include.patch b/internal/rosa/package/gzip/late-include.patch new file mode 100644 index 00000000..4638d578 --- /dev/null +++ b/internal/rosa/package/gzip/late-include.patch @@ -0,0 +1,21 @@ +diff --git a/gzip.c b/gzip.c +index 220f6fc..ef79a89 100644 +--- a/gzip.c ++++ b/gzip.c +@@ -56,6 +56,8 @@ static char const license_msg[] = + + #include <config.h> + ++#include <signal.h> ++ + #include "tailor.h" + + #include "gzip.h" +@@ -79,7 +81,6 @@ static char const license_msg[] = + #include <inttypes.h> + #include <limits.h> + #include <locale.h> +-#include <signal.h> + #include <stdcountof.h> + #include <stddef.h> + #include <stdlib.h> diff --git a/internal/rosa/package/gzip/package.az b/internal/rosa/package/gzip/package.az new file mode 100644 index 00000000..6c6b3b1f --- /dev/null +++ b/internal/rosa/package/gzip/package.az @@ -0,0 +1,18 @@ +package gzip { + description = "a popular data compression program"; + website = "https://www.gnu.org/software/gzip"; + anitya = 1290; + + version# = "1.15"; + source = remoteTar { + url = "https://mirrors.kernel.org/gnu/gzip/gzip-"+version+".tar.gz"; + checksum = "N8v7cH-Vhmxvk1PYaF-B1vH-r-rMT3Z3VyMdCFXhBd74a8v62wT3dKes-u1JCvUG"; + compress = gzip; + }; + patches = [ "late-include.patch" ]; + + exec = make { + // dependency loop + check = nil; + }; +} diff --git a/pkg/archive.go b/pkg/archive.go index 532c0df9..ba7c61b9 100644 --- a/pkg/archive.go +++ b/pkg/archive.go @@ -314,10 +314,10 @@ func (archiveArtifact) IsExclusive() bool { return false } // Revision satisfies [RevisionArtifact] for status behaviour. func (archiveArtifact) Revision() uint64 { return 0 } -// Unpack reads an archive stream from r and unpacks its contents to dir. If +// Extract reads an archive stream from r and writes its contents to dir. If // writeStatus is non-nil, it is called with a user-facing message after -// each successfully unpacked entry. -func Unpack( +// each successfully extracted entry. +func Extract( r io.Reader, dir string, writeStatus func(s string) error, @@ -444,7 +444,7 @@ func (a archiveArtifact) Cure(t *TContext) (err error) { } msg := t.GetMessage() - err = Unpack(r, t.GetWorkDir().String(), func(s string) (err error) { + err = Extract(r, t.GetWorkDir().String(), func(s string) (err error) { msg.Verbose(s) if _, err = io.WriteString(status, s); err != nil { return @@ -2933,6 +2933,9 @@ func Open( return &c, nil } +// Path returns the pathname of the directory holding [Cache] state. +func (c *Cache) Path() *check.Absolute { return c.base } + // Collected is returned by [Collect.Cure] to indicate a successful collection. type Collected struct{} |
