From be6e1cf55eb49d6035c1d37bae5ea472470b59d1 Mon Sep 17 00:00:00 2001 From: Tonis Tiigi Date: Mon, 21 Nov 2022 23:25:55 -0800 Subject: [PATCH] history api: support for logs of completed builds Signed-off-by: Tonis Tiigi --- client/solve.go | 48 ---------- client/status.go | 125 ++++++++++++++++++++++++++ cmd/buildkitd/main.go | 1 + control/control.go | 91 +++++++------------ control/init.go | 10 --- solver/jobs.go | 7 +- solver/llbsolver/history.go | 167 +++++++++++++++++++++++++++++++++++ solver/llbsolver/solver.go | 43 +++++++++ util/progress/multireader.go | 12 +-- util/progress/progress.go | 22 +++-- 10 files changed, 395 insertions(+), 131 deletions(-) create mode 100644 client/status.go delete mode 100644 control/init.go diff --git a/client/solve.go b/client/solve.go index b1d78b536..37a808664 100644 --- a/client/solve.go +++ b/client/solve.go @@ -350,54 +350,6 @@ func (c *Client) solve(ctx context.Context, def *llb.Definition, runGateway runG return res, nil } -func NewSolveStatus(resp *controlapi.StatusResponse) *SolveStatus { - s := &SolveStatus{} - for _, v := range resp.Vertexes { - s.Vertexes = append(s.Vertexes, &Vertex{ - Digest: v.Digest, - Inputs: v.Inputs, - Name: v.Name, - Started: v.Started, - Completed: v.Completed, - Error: v.Error, - Cached: v.Cached, - ProgressGroup: v.ProgressGroup, - }) - } - for _, v := range resp.Statuses { - s.Statuses = append(s.Statuses, &VertexStatus{ - ID: v.ID, - Vertex: v.Vertex, - Name: v.Name, - Total: v.Total, - Current: v.Current, - Timestamp: v.Timestamp, - Started: v.Started, - Completed: v.Completed, - }) - } - for _, v := range resp.Logs { - s.Logs = append(s.Logs, &VertexLog{ - Vertex: v.Vertex, - Stream: int(v.Stream), - Data: v.Msg, - Timestamp: v.Timestamp, - }) - } - for _, v := range resp.Warnings { - s.Warnings = append(s.Warnings, &VertexWarning{ - Vertex: v.Vertex, - Level: int(v.Level), - Short: v.Short, - Detail: v.Detail, - URL: v.Url, - SourceInfo: v.Info, - Range: v.Ranges, - }) - } - return s -} - func prepareSyncedDirs(def *llb.Definition, localDirs map[string]string) (filesync.StaticDirSource, error) { for _, d := range localDirs { fi, err := os.Stat(d) diff --git a/client/status.go b/client/status.go new file mode 100644 index 000000000..d692094af --- /dev/null +++ b/client/status.go @@ -0,0 +1,125 @@ +package client + +import ( + controlapi "github.com/moby/buildkit/api/services/control" +) + +var emptyLogVertexSize int + +func init() { + emptyLogVertex := controlapi.VertexLog{} + emptyLogVertexSize = emptyLogVertex.Size() +} + +func NewSolveStatus(resp *controlapi.StatusResponse) *SolveStatus { + s := &SolveStatus{} + for _, v := range resp.Vertexes { + s.Vertexes = append(s.Vertexes, &Vertex{ + Digest: v.Digest, + Inputs: v.Inputs, + Name: v.Name, + Started: v.Started, + Completed: v.Completed, + Error: v.Error, + Cached: v.Cached, + ProgressGroup: v.ProgressGroup, + }) + } + for _, v := range resp.Statuses { + s.Statuses = append(s.Statuses, &VertexStatus{ + ID: v.ID, + Vertex: v.Vertex, + Name: v.Name, + Total: v.Total, + Current: v.Current, + Timestamp: v.Timestamp, + Started: v.Started, + Completed: v.Completed, + }) + } + for _, v := range resp.Logs { + s.Logs = append(s.Logs, &VertexLog{ + Vertex: v.Vertex, + Stream: int(v.Stream), + Data: v.Msg, + Timestamp: v.Timestamp, + }) + } + for _, v := range resp.Warnings { + s.Warnings = append(s.Warnings, &VertexWarning{ + Vertex: v.Vertex, + Level: int(v.Level), + Short: v.Short, + Detail: v.Detail, + URL: v.Url, + SourceInfo: v.Info, + Range: v.Ranges, + }) + } + return s +} + +func (ss *SolveStatus) Marshal() (out []*controlapi.StatusResponse) { + logSize := 0 + for { + retry := false + sr := controlapi.StatusResponse{} + for _, v := range ss.Vertexes { + sr.Vertexes = append(sr.Vertexes, &controlapi.Vertex{ + Digest: v.Digest, + Inputs: v.Inputs, + Name: v.Name, + Started: v.Started, + Completed: v.Completed, + Error: v.Error, + Cached: v.Cached, + ProgressGroup: v.ProgressGroup, + }) + } + for _, v := range ss.Statuses { + sr.Statuses = append(sr.Statuses, &controlapi.VertexStatus{ + ID: v.ID, + Vertex: v.Vertex, + Name: v.Name, + Current: v.Current, + Total: v.Total, + Timestamp: v.Timestamp, + Started: v.Started, + Completed: v.Completed, + }) + } + for i, v := range ss.Logs { + sr.Logs = append(sr.Logs, &controlapi.VertexLog{ + Vertex: v.Vertex, + Stream: int64(v.Stream), + Msg: v.Data, + Timestamp: v.Timestamp, + }) + logSize += len(v.Data) + emptyLogVertexSize + // avoid logs growing big and split apart if they do + if logSize > 1024*1024 { + ss.Vertexes = nil + ss.Statuses = nil + ss.Logs = ss.Logs[i+1:] + retry = true + break + } + } + for _, v := range ss.Warnings { + sr.Warnings = append(sr.Warnings, &controlapi.VertexWarning{ + Vertex: v.Vertex, + Level: int64(v.Level), + Short: v.Short, + Detail: v.Detail, + Info: v.SourceInfo, + Ranges: v.Range, + Url: v.URL, + }) + } + out = append(out, &sr) + if !retry { + break + } + } + return +} diff --git a/cmd/buildkitd/main.go b/cmd/buildkitd/main.go index d57b17414..beaae0953 100644 --- a/cmd/buildkitd/main.go +++ b/cmd/buildkitd/main.go @@ -688,6 +688,7 @@ func newController(c *cli.Context, cfg *config.Config) (*control.Controller, err TraceCollector: tc, HistoryDB: historyDB, LeaseManager: w.LeaseManager(), + ContentStore: w.ContentStore(), }) } diff --git a/control/control.go b/control/control.go index eb53827e9..ce7957f80 100644 --- a/control/control.go +++ b/control/control.go @@ -7,7 +7,10 @@ import ( "sync/atomic" "time" + contentapi "github.com/containerd/containerd/api/services/content/v1" + "github.com/containerd/containerd/content" "github.com/containerd/containerd/leases" + "github.com/containerd/containerd/services/content/contentserver" "github.com/docker/distribution/reference" "github.com/mitchellh/hashstructure/v2" controlapi "github.com/moby/buildkit/api/services/control" @@ -31,6 +34,7 @@ import ( "github.com/moby/buildkit/util/tracing/transform" "github.com/moby/buildkit/version" "github.com/moby/buildkit/worker" + digest "github.com/opencontainers/go-digest" "github.com/pkg/errors" "go.etcd.io/bbolt" sdktrace "go.opentelemetry.io/otel/sdk/trace" @@ -52,6 +56,7 @@ type Opt struct { TraceCollector sdktrace.SpanExporter HistoryDB *bbolt.DB LeaseManager leases.Manager + ContentStore content.Store } type Controller struct { // TODO: ControlService @@ -75,6 +80,7 @@ func NewController(opt Opt) (*Controller, error) { hq := llbsolver.NewHistoryQueue(llbsolver.HistoryQueueOpt{ DB: opt.HistoryDB, LeaseManager: opt.LeaseManager, + ContentStore: opt.ContentStore, }) s, err := llbsolver.New(llbsolver.Opt{ @@ -115,6 +121,9 @@ func (c *Controller) Register(server *grpc.Server) { controlapi.RegisterControlServer(server, c) c.gatewayForwarder.Register(server) tracev1.RegisterTraceServiceServer(server, c) + + store := &roContentStore{c.opt.ContentStore} + contentapi.RegisterContentServer(server, contentserver.New(store)) } func (c *Controller) DiskUsage(ctx context.Context, r *controlapi.DiskUsageRequest) (*controlapi.DiskUsageResponse, error) { @@ -407,68 +416,10 @@ func (c *Controller) Status(req *controlapi.StatusRequest, stream controlapi.Con if !ok { return nil } - logSize := 0 - for { - retry := false - sr := controlapi.StatusResponse{} - for _, v := range ss.Vertexes { - sr.Vertexes = append(sr.Vertexes, &controlapi.Vertex{ - Digest: v.Digest, - Inputs: v.Inputs, - Name: v.Name, - Started: v.Started, - Completed: v.Completed, - Error: v.Error, - Cached: v.Cached, - ProgressGroup: v.ProgressGroup, - }) - } - for _, v := range ss.Statuses { - sr.Statuses = append(sr.Statuses, &controlapi.VertexStatus{ - ID: v.ID, - Vertex: v.Vertex, - Name: v.Name, - Current: v.Current, - Total: v.Total, - Timestamp: v.Timestamp, - Started: v.Started, - Completed: v.Completed, - }) - } - for i, v := range ss.Logs { - sr.Logs = append(sr.Logs, &controlapi.VertexLog{ - Vertex: v.Vertex, - Stream: int64(v.Stream), - Msg: v.Data, - Timestamp: v.Timestamp, - }) - logSize += len(v.Data) + emptyLogVertexSize - // avoid logs growing big and split apart if they do - if logSize > 1024*1024 { - ss.Vertexes = nil - ss.Statuses = nil - ss.Logs = ss.Logs[i+1:] - retry = true - break - } - } - for _, v := range ss.Warnings { - sr.Warnings = append(sr.Warnings, &controlapi.VertexWarning{ - Vertex: v.Vertex, - Level: int64(v.Level), - Short: v.Short, - Detail: v.Detail, - Info: v.SourceInfo, - Ranges: v.Range, - Url: v.URL, - }) - } - if err := stream.SendMsg(&sr); err != nil { + for _, sr := range ss.Marshal() { + if err := stream.SendMsg(sr); err != nil { return err } - if !retry { - break - } } } }) @@ -633,3 +584,23 @@ func cacheOptKey(opt controlapi.CacheOptionsEntry) (string, error) { } return fmt.Sprint(opt.Type, ":", hash), nil } + +type roContentStore struct { + content.Store +} + +func (cs *roContentStore) Writer(ctx context.Context, opts ...content.WriterOpt) (content.Writer, error) { + return nil, errors.Errorf("read-only content store") +} + +func (cs *roContentStore) Delete(ctx context.Context, dgst digest.Digest) error { + return errors.Errorf("read-only content store") +} + +func (cs *roContentStore) Update(ctx context.Context, info content.Info, fieldpaths ...string) (content.Info, error) { + return content.Info{}, errors.Errorf("read-only content store") +} + +func (cs *roContentStore) Abort(ctx context.Context, ref string) error { + return errors.Errorf("read-only content store") +} diff --git a/control/init.go b/control/init.go deleted file mode 100644 index 2e86133e4..000000000 --- a/control/init.go +++ /dev/null @@ -1,10 +0,0 @@ -package control - -import controlapi "github.com/moby/buildkit/api/services/control" - -var emptyLogVertexSize int - -func init() { - emptyLogVertex := controlapi.VertexLog{} - emptyLogVertexSize = emptyLogVertex.Size() -} diff --git a/solver/jobs.go b/solver/jobs.go index 070b4020b..5a4f2ba6b 100644 --- a/solver/jobs.go +++ b/solver/jobs.go @@ -557,9 +557,12 @@ func (j *Job) walkProvenance(ctx context.Context, e Edge, f func(ProvenanceProvi return nil } -func (j *Job) Discard() error { - defer j.progressCloser() +func (j *Job) CloseProgress() { + j.progressCloser() + j.pw.Close() +} +func (j *Job) Discard() error { j.list.mu.Lock() defer j.list.mu.Unlock() diff --git a/solver/llbsolver/history.go b/solver/llbsolver/history.go index 75f333192..72159e27b 100644 --- a/solver/llbsolver/history.go +++ b/solver/llbsolver/history.go @@ -1,11 +1,22 @@ package llbsolver import ( + "bufio" "context" + "encoding/binary" + "io" + "os" "sync" + "time" + "github.com/containerd/containerd/content" + "github.com/containerd/containerd/errdefs" "github.com/containerd/containerd/leases" controlapi "github.com/moby/buildkit/api/services/control" + "github.com/moby/buildkit/client" + "github.com/moby/buildkit/util/leaseutil" + digest "github.com/opencontainers/go-digest" + ocispecs "github.com/opencontainers/image-spec/specs-go/v1" "github.com/pkg/errors" bolt "go.etcd.io/bbolt" ) @@ -17,6 +28,7 @@ const ( type HistoryQueueOpt struct { DB *bolt.DB LeaseManager leases.Manager + ContentStore content.Store } type HistoryQueue struct { @@ -62,6 +74,70 @@ func (h *HistoryQueue) addResource(ctx context.Context, l leases.Lease, desc *co }) } +func (h *HistoryQueue) Status(ctx context.Context, ref string, st chan<- *client.SolveStatus) error { + var br controlapi.BuildHistoryRecord + if err := h.DB.View(func(tx *bolt.Tx) error { + b := tx.Bucket([]byte(recordsBucket)) + if b == nil { + return nil + } + dt := b.Get([]byte(ref)) + if dt == nil { + return os.ErrNotExist + } + + if err := br.Unmarshal(dt); err != nil { + return errors.Wrapf(err, "failed to unmarshal build record %s", ref) + } + return nil + }); err != nil { + return err + } + + if br.Logs == nil { + return nil + } + + ra, err := h.ContentStore.ReaderAt(ctx, ocispecs.Descriptor{ + Digest: br.Logs.Digest, + Size: br.Logs.Size_, + MediaType: br.Logs.MediaType, + }) + if err != nil { + return err + } + defer ra.Close() + + brdr := bufio.NewReader(&reader{ReaderAt: ra}) + + buf := make([]byte, 32*1024) + + for { + _, err := io.ReadAtLeast(brdr, buf[:4], 4) + if err != nil { + if errors.Is(err, io.EOF) { + break + } + return err + } + sz := binary.LittleEndian.Uint32(buf[:4]) + if sz > uint32(len(buf)) { + buf = make([]byte, sz) + } + _, err = io.ReadAtLeast(brdr, buf[:sz], int(sz)) + if err != nil { + return err + } + var sr controlapi.StatusResponse + if err := sr.Unmarshal(buf[:sz]); err != nil { + return err + } + st <- client.NewSolveStatus(&sr) + } + + return nil +} + func (h *HistoryQueue) Update(ctx context.Context, e *controlapi.BuildHistoryEvent) error { h.init() h.mu.Lock() @@ -108,6 +184,86 @@ func (h *HistoryQueue) Update(ctx context.Context, e *controlapi.BuildHistoryEve return nil } +func (h *HistoryQueue) ImportStatus(ctx context.Context, ch chan *client.SolveStatus) (_ *ocispecs.Descriptor, _ func(), err error) { + defer func() { + if ch == nil { + return + } + for range ch { + } + }() + + l, err := h.LeaseManager.Create(ctx, leases.WithRandomID(), leases.WithExpiration(5*time.Minute), leaseutil.MakeTemporary) + if err != nil { + return nil, nil, err + } + defer func() { + if err != nil { + h.LeaseManager.Delete(ctx, l) + } + }() + ctx = leases.WithLease(ctx, l.ID) + + w, err := content.OpenWriter(ctx, h.ContentStore, content.WithRef("status-"+h.leaseID(l.ID))) + if err != nil { + return nil, nil, err + } + bufW := bufio.NewWriter(w) + + defer func() { + if err != nil && w != nil { + w.Close() + } + }() + + dgst := digest.Canonical.Digester() + total := 0 + + buf := make([]byte, 32*1024) + for st := range ch { + hdr := make([]byte, 4) + for _, pst := range st.Marshal() { + sz := pst.Size() + if len(buf) < sz { + buf = make([]byte, sz) + } + n, err := pst.MarshalTo(buf) + if err != nil { + return nil, nil, err + } + binary.LittleEndian.PutUint32(hdr, uint32(n)) + if _, err := bufW.Write(hdr); err != nil { + return nil, nil, err + } + if _, err := bufW.Write(buf[:n]); err != nil { + return nil, nil, err + } + dgst.Hash().Write(hdr) + dgst.Hash().Write(buf[:n]) + total += 4 + n + } + } + if err := bufW.Flush(); err != nil { + return nil, nil, err + } + + if err := w.Commit(ctx, int64(total), dgst.Digest()); err != nil { + if !errdefs.IsAlreadyExists(err) { + return nil, nil, err + } + } + w = nil + + return &ocispecs.Descriptor{ + MediaType: "application/vnd.buildkit.status.v0", + Digest: dgst.Digest(), + Size: int64(total), + }, + func() { + h.LeaseManager.Delete(context.TODO(), l) + }, nil +} + func (h *HistoryQueue) Listen(ctx context.Context, ref string, active bool, f func(*controlapi.BuildHistoryEvent) error) error { h.init() @@ -208,3 +364,14 @@ func (p *channel[T]) close() { close(p.done) }) } + +type reader struct { + io.ReaderAt + pos int64 +} + +func (r *reader) Read(p []byte) (int, error) { + n, err := r.ReaderAt.ReadAt(p, r.pos) + r.pos += int64(len(p)) + return n, err +} diff --git a/solver/llbsolver/solver.go b/solver/llbsolver/solver.go index 1a4e803cb..0dec3be71 100644 --- a/solver/llbsolver/solver.go +++ b/solver/llbsolver/solver.go @@ -4,6 +4,7 @@ import ( "context" "encoding/base64" "fmt" + "os" "strings" "time" @@ -29,6 +30,7 @@ import ( "github.com/moby/buildkit/util/progress" "github.com/moby/buildkit/worker" digest "github.com/opencontainers/go-digest" + ocispecs "github.com/opencontainers/image-spec/specs-go/v1" "github.com/pkg/errors" "golang.org/x/sync/errgroup" "google.golang.org/grpc/codes" @@ -151,6 +153,38 @@ func (s *Solver) Solve(ctx context.Context, id string, sessionID string, req fro en := time.Now() rec.CompletedAt = &en + j.CloseProgress() + + ch := make(chan *client.SolveStatus) + eg, ctx2 := errgroup.WithContext(ctx) + var releaseStatus func() + eg.Go(func() error { + var err error + var desc *ocispecs.Descriptor + desc, releaseStatus, err = s.history.ImportStatus(ctx2, ch) + if err != nil { + return err + } + rec.Logs = &controlapi.Descriptor{ + Digest: desc.Digest, + Size_: desc.Size, + MediaType: desc.MediaType, + } + return nil + }) + eg.Go(func() error { + return j.Status(ctx2, ch) + }) + if err1 := eg.Wait(); err == nil { + err = err1 + } + + defer func() { + if releaseStatus != nil { + releaseStatus() + } + }() + if err != nil { st, ok := grpcerrors.AsGRPCStatus(grpcerrors.ToGRPC(err)) if !ok { @@ -592,6 +626,15 @@ func withDescHandlerCacheOpts(ctx context.Context, ref cache.ImmutableRef) conte } func (s *Solver) Status(ctx context.Context, id string, statusChan chan *client.SolveStatus) error { + if err := s.history.Status(ctx, id, statusChan); err != nil { + if !errors.Is(err, os.ErrNotExist) { + close(statusChan) + return err + } + } else { + close(statusChan) + return nil + } j, err := s.solver.Get(id) if err != nil { close(statusChan) diff --git a/util/progress/multireader.go b/util/progress/multireader.go index 3f9d37816..b0d92dde8 100644 --- a/util/progress/multireader.go +++ b/util/progress/multireader.go @@ -35,8 +35,13 @@ func (mr *MultiReader) Reader(ctx context.Context) Reader { isBehind := len(mr.sent) > 0 - if !isBehind { - mr.writers[w] = closeWriter + select { + case <-mr.done: + isBehind = true + default: + if !isBehind { + mr.writers[w] = closeWriter + } } go func() { @@ -74,9 +79,6 @@ func (mr *MultiReader) Reader(ctx context.Context) Reader { case <-ctx.Done(): close() return - case <-mr.done: - close() - return default: } } diff --git a/util/progress/progress.go b/util/progress/progress.go index 4fabf3769..fbbb22de0 100644 --- a/util/progress/progress.go +++ b/util/progress/progress.go @@ -118,12 +118,22 @@ func (pr *progressReader) Read(ctx context.Context) ([]*Progress, error) { done := make(chan struct{}) defer close(done) go func() { - select { - case <-done: - case <-ctx.Done(): - pr.mu.Lock() - pr.cond.Broadcast() - pr.mu.Unlock() + prdone := pr.ctx.Done() + for { + select { + case <-done: + return + case <-ctx.Done(): + pr.mu.Lock() + pr.cond.Broadcast() + pr.mu.Unlock() + return + case <-prdone: + pr.mu.Lock() + pr.cond.Broadcast() + pr.mu.Unlock() + prdone = nil + } } }() pr.mu.Lock()