history api: support for logs of completed builds

Signed-off-by: Tonis Tiigi <tonistiigi@gmail.com>
This commit is contained in:
Tonis Tiigi
2022-11-21 23:25:55 -08:00
parent 6b004d3536
commit be6e1cf55e
10 changed files with 395 additions and 131 deletions

View File

@@ -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)

125
client/status.go Normal file
View File

@@ -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
}

View File

@@ -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(),
})
}

View File

@@ -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")
}

View File

@@ -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()
}

View File

@@ -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()

View File

@@ -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
}

View File

@@ -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)

View File

@@ -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:
}
}

View File

@@ -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()