Merge pull request #3294 from tonistiigi/build-history

add build history APIs
This commit is contained in:
Tõnis Tiigi
2022-11-23 13:53:54 -08:00
committed by GitHub
21 changed files with 3435 additions and 161 deletions

File diff suppressed because it is too large Load Diff

View File

@@ -6,6 +6,8 @@ import "github.com/gogo/protobuf/gogoproto/gogo.proto";
import "google/protobuf/timestamp.proto";
import "github.com/moby/buildkit/solver/pb/ops.proto";
import "github.com/moby/buildkit/api/types/worker.proto";
// import "github.com/containerd/containerd/api/types/descriptor.proto";
import "github.com/gogo/googleapis/google/rpc/status.proto";
option (gogoproto.sizer_all) = true;
option (gogoproto.marshaler_all) = true;
@@ -19,6 +21,8 @@ service Control {
rpc Session(stream BytesMessage) returns (stream BytesMessage);
rpc ListWorkers(ListWorkersRequest) returns (ListWorkersResponse);
rpc Info(InfoRequest) returns (InfoResponse);
rpc ListenBuildHistory(BuildHistoryRequest) returns (stream BuildHistoryEvent);
}
message PruneRequest {
@@ -163,3 +167,56 @@ message InfoRequest {}
message InfoResponse {
moby.buildkit.v1.types.BuildkitVersion buildkitVersion = 1;
}
message BuildHistoryRequest {
bool ActiveOnly = 1;
string Ref = 2;
}
enum BuildHistoryEventType {
STARTED = 0;
COMPLETE = 1;
DELETED = 2;
}
message BuildHistoryEvent {
BuildHistoryEventType type = 1;
BuildHistoryRecord record = 2;
}
message BuildHistoryRecord {
string Ref = 1;
string Frontend = 2;
map<string, string> FrontendAttrs = 3;
repeated Exporter Exporters = 4;
google.rpc.Status error = 5;
google.protobuf.Timestamp CreatedAt = 6 [(gogoproto.stdtime) = true];
google.protobuf.Timestamp CompletedAt = 7 [(gogoproto.stdtime) = true];
Descriptor logs = 8;
map<string, string> ExporterResponse = 9;
BuildResultInfo Result = 10;
map<string, BuildResultInfo> Results = 11;
int32 Generation = 12;
// TODO: tags
// TODO: steps/cache summary
// TODO: unclipped logs
// TODO: pinning
}
message Descriptor {
string media_type = 1;
string digest = 2 [(gogoproto.customtype) = "github.com/opencontainers/go-digest.Digest", (gogoproto.nullable) = false];
int64 size = 3;
map<string, string> annotations = 5;
}
message BuildResultInfo {
Descriptor Result = 1;
repeated Descriptor Attestations = 2;
}
message Exporter {
string Type = 1;
map<string, string> Attrs = 2;
}

View File

@@ -168,12 +168,12 @@ func (c *Client) setupDelegatedTracing(ctx context.Context, td TracerDelegate) e
return td.SetSpanExporter(ctx, e)
}
func (c *Client) controlClient() controlapi.ControlClient {
func (c *Client) ControlClient() controlapi.ControlClient {
return controlapi.NewControlClient(c.conn)
}
func (c *Client) Dialer() session.Dialer {
return grpchijack.Dialer(c.controlClient())
return grpchijack.Dialer(c.ControlClient())
}
func (c *Client) Close() error {

View File

@@ -31,7 +31,7 @@ func (c *Client) DiskUsage(ctx context.Context, opts ...DiskUsageOption) ([]*Usa
}
req := &controlapi.DiskUsageRequest{Filter: info.Filter}
resp, err := c.controlClient().DiskUsage(ctx, req)
resp, err := c.ControlClient().DiskUsage(ctx, req)
if err != nil {
return nil, errors.Wrap(err, "failed to call diskusage")
}

View File

@@ -19,7 +19,7 @@ type BuildkitVersion struct {
}
func (c *Client) Info(ctx context.Context) (*Info, error) {
res, err := c.controlClient().Info(ctx, &controlapi.InfoRequest{})
res, err := c.ControlClient().Info(ctx, &controlapi.InfoRequest{})
if err != nil {
return nil, errors.Wrap(err, "failed to call info")
}

View File

@@ -23,7 +23,7 @@ func (c *Client) Prune(ctx context.Context, ch chan UsageInfo, opts ...PruneOpti
if info.All {
req.All = true
}
cl, err := c.controlClient().Prune(ctx, req)
cl, err := c.ControlClient().Prune(ctx, req)
if err != nil {
return errors.Wrap(err, "failed to call prune")
}

View File

@@ -204,7 +204,7 @@ func (c *Client) solve(ctx context.Context, def *llb.Definition, runGateway runG
eg.Go(func() error {
sd := c.sessionDialer
if sd == nil {
sd = grpchijack.Dialer(c.controlClient())
sd = grpchijack.Dialer(c.ControlClient())
}
return s.Run(statusContext, sd)
})
@@ -247,7 +247,7 @@ func (c *Client) solve(ctx context.Context, def *llb.Definition, runGateway runG
frontendInputs[key] = def.ToPB()
}
resp, err := c.controlClient().Solve(ctx, &controlapi.SolveRequest{
resp, err := c.ControlClient().Solve(ctx, &controlapi.SolveRequest{
Ref: ref,
Definition: pbd,
Exporter: ex.Type,
@@ -291,7 +291,7 @@ func (c *Client) solve(ctx context.Context, def *llb.Definition, runGateway runG
}
eg.Go(func() error {
stream, err := c.controlClient().Status(statusContext, &controlapi.StatusRequest{
stream, err := c.ControlClient().Status(statusContext, &controlapi.StatusRequest{
Ref: ref,
})
if err != nil {
@@ -305,52 +305,8 @@ func (c *Client) solve(ctx context.Context, def *llb.Definition, runGateway runG
}
return errors.Wrap(err, "failed to receive status")
}
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,
})
}
if statusChan != nil {
statusChan <- &s
statusChan <- NewSolveStatus(resp)
}
}
})
@@ -393,6 +349,54 @@ 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)

View File

@@ -28,7 +28,7 @@ func (c *Client) ListWorkers(ctx context.Context, opts ...ListWorkersOption) ([]
}
req := &controlapi.ListWorkersRequest{Filter: info.Filter}
resp, err := c.controlClient().ListWorkers(ctx, req)
resp, err := c.ControlClient().ListWorkers(ctx, req)
if err != nil {
return nil, errors.Wrap(err, "failed to list workers")
}

View File

@@ -13,5 +13,7 @@ var debugCommand = cli.Command{
debug.DumpMetadataCommand,
debug.WorkersCommand,
debug.InfoCommand,
debug.MonitorCommand,
debug.LogsCommand,
},
}

View File

@@ -0,0 +1,71 @@
package debug
import (
"fmt"
"io"
"os"
controlapi "github.com/moby/buildkit/api/services/control"
"github.com/moby/buildkit/client"
bccommon "github.com/moby/buildkit/cmd/buildctl/common"
"github.com/moby/buildkit/util/appcontext"
"github.com/moby/buildkit/util/progress/progresswriter"
"github.com/pkg/errors"
"github.com/urfave/cli"
)
var LogsCommand = cli.Command{
Name: "logs",
Usage: "display build logs",
Action: logs,
Flags: []cli.Flag{
cli.StringFlag{
Name: "progress",
Usage: "progress output type",
Value: "auto",
},
},
}
func logs(clicontext *cli.Context) error {
args := clicontext.Args()
if len(args) == 0 {
return fmt.Errorf("build ref must be specified")
}
ref := args[0]
c, err := bccommon.ResolveClient(clicontext)
if err != nil {
return err
}
ctx := appcontext.Context()
cl, err := c.ControlClient().Status(ctx, &controlapi.StatusRequest{
Ref: ref,
})
if err != nil {
return err
}
pw, err := progresswriter.NewPrinter(ctx, os.Stdout, clicontext.String("progress"))
if err != nil {
return err
}
defer func() {
<-pw.Done()
}()
for {
resp, err := cl.Recv()
if err != nil {
close(pw.Status())
if errors.Is(err, io.EOF) {
return nil
}
return err
}
pw.Status() <- client.NewSolveStatus(resp)
}
}

View File

@@ -0,0 +1,52 @@
package debug
import (
"fmt"
controlapi "github.com/moby/buildkit/api/services/control"
bccommon "github.com/moby/buildkit/cmd/buildctl/common"
"github.com/moby/buildkit/util/appcontext"
"github.com/urfave/cli"
)
var MonitorCommand = cli.Command{
Name: "monitor",
Usage: "display build events",
Action: monitor,
Flags: []cli.Flag{
cli.BoolFlag{
Name: "completed",
Usage: "show completed builds",
},
cli.StringFlag{
Name: "ref",
Usage: "show events for a specific build",
},
},
}
func monitor(clicontext *cli.Context) error {
c, err := bccommon.ResolveClient(clicontext)
if err != nil {
return err
}
completed := clicontext.Bool("completed")
ctx := appcontext.Context()
cl, err := c.ControlClient().ListenBuildHistory(ctx, &controlapi.BuildHistoryRequest{
ActiveOnly: !completed,
Ref: clicontext.String("ref"),
})
if err != nil {
return err
}
for {
ev, err := cl.Recv()
if err != nil {
return err
}
fmt.Printf("%s %s\n", ev.Type.String(), ev.Record.Ref)
}
}

View File

@@ -59,6 +59,7 @@ import (
"github.com/pkg/errors"
"github.com/sirupsen/logrus"
"github.com/urfave/cli"
"go.etcd.io/bbolt"
"go.opentelemetry.io/contrib/instrumentation/google.golang.org/grpc/otelgrpc"
"go.opentelemetry.io/otel/propagation"
sdktrace "go.opentelemetry.io/otel/sdk/trace"
@@ -649,6 +650,11 @@ func newController(c *cli.Context, cfg *config.Config) (*control.Controller, err
return nil, err
}
historyDB, err := bbolt.Open(filepath.Join(cfg.Root, "history.db"), 0600, nil)
if err != nil {
return nil, err
}
resolverFn := resolverFunc(cfg)
w, err := wc.GetDefault()
@@ -680,6 +686,8 @@ func newController(c *cli.Context, cfg *config.Config) (*control.Controller, err
CacheKeyStorage: cacheStorage,
Entitlements: cfg.Entitlements,
TraceCollector: tc,
HistoryDB: historyDB,
LeaseManager: w.LeaseManager(),
})
}

View File

@@ -7,6 +7,7 @@ import (
"sync/atomic"
"time"
"github.com/containerd/containerd/leases"
"github.com/docker/distribution/reference"
"github.com/mitchellh/hashstructure/v2"
controlapi "github.com/moby/buildkit/api/services/control"
@@ -31,6 +32,7 @@ import (
"github.com/moby/buildkit/version"
"github.com/moby/buildkit/worker"
"github.com/pkg/errors"
"go.etcd.io/bbolt"
sdktrace "go.opentelemetry.io/otel/sdk/trace"
tracev1 "go.opentelemetry.io/proto/otlp/collector/trace/v1"
"golang.org/x/sync/errgroup"
@@ -48,6 +50,8 @@ type Opt struct {
ResolveCacheImporterFuncs map[string]remotecache.ResolveCacheImporterFunc
Entitlements []string
TraceCollector sdktrace.SpanExporter
HistoryDB *bbolt.DB
LeaseManager leases.Manager
}
type Controller struct { // TODO: ControlService
@@ -55,6 +59,7 @@ type Controller struct { // TODO: ControlService
buildCount int64
opt Opt
solver *llbsolver.Solver
history *llbsolver.HistoryQueue
cache solver.CacheManager
gatewayForwarder *controlgateway.GatewayForwarder
throttledGC func()
@@ -67,6 +72,11 @@ func NewController(opt Opt) (*Controller, error) {
gatewayForwarder := controlgateway.NewGatewayForwarder()
hq := llbsolver.NewHistoryQueue(llbsolver.HistoryQueueOpt{
DB: opt.HistoryDB,
LeaseManager: opt.LeaseManager,
})
s, err := llbsolver.New(llbsolver.Opt{
WorkerController: opt.WorkerController,
Frontends: opt.Frontends,
@@ -75,6 +85,7 @@ func NewController(opt Opt) (*Controller, error) {
GatewayForwarder: gatewayForwarder,
SessionManager: opt.SessionManager,
Entitlements: opt.Entitlements,
HistoryQueue: hq,
})
if err != nil {
return nil, errors.Wrap(err, "failed to create solver")
@@ -83,6 +94,7 @@ func NewController(opt Opt) (*Controller, error) {
c := &Controller{
opt: opt,
solver: s,
history: hq,
cache: cache,
gatewayForwarder: gatewayForwarder,
}
@@ -222,6 +234,15 @@ func (c *Controller) Export(ctx context.Context, req *tracev1.ExportTraceService
return &tracev1.ExportTraceServiceResponse{}, nil
}
func (c *Controller) ListenBuildHistory(req *controlapi.BuildHistoryRequest, srv controlapi.Control_ListenBuildHistoryServer) error {
return c.history.Listen(srv.Context(), req.Ref, req.ActiveOnly, func(h *controlapi.BuildHistoryEvent) error {
if err := srv.Send(h); err != nil {
return err
}
return nil
})
}
func translateLegacySolveRequest(req *controlapi.SolveRequest) error {
// translates ExportRef and ExportAttrs to new Exports (v0.4.0)
if legacyExportRef := req.Cache.ExportRefDeprecated; legacyExportRef != "" {

210
solver/llbsolver/history.go Normal file
View File

@@ -0,0 +1,210 @@
package llbsolver
import (
"context"
"sync"
"github.com/containerd/containerd/leases"
controlapi "github.com/moby/buildkit/api/services/control"
"github.com/pkg/errors"
bolt "go.etcd.io/bbolt"
)
const (
recordsBucket = "_records"
)
type HistoryQueueOpt struct {
DB *bolt.DB
LeaseManager leases.Manager
}
type HistoryQueue struct {
mu sync.Mutex
initOnce sync.Once
HistoryQueueOpt
ps *pubsub[*controlapi.BuildHistoryEvent]
active map[string]*controlapi.BuildHistoryRecord
}
func NewHistoryQueue(opt HistoryQueueOpt) *HistoryQueue {
return &HistoryQueue{
HistoryQueueOpt: opt,
ps: &pubsub[*controlapi.BuildHistoryEvent]{
m: map[*channel[*controlapi.BuildHistoryEvent]]struct{}{},
},
active: map[string]*controlapi.BuildHistoryRecord{},
}
}
func (h *HistoryQueue) init() error {
var err error
h.initOnce.Do(func() {
err = h.DB.Update(func(tx *bolt.Tx) error {
_, err := tx.CreateBucketIfNotExists([]byte(recordsBucket))
return err
})
})
return err
}
func (h *HistoryQueue) leaseID(id string) string {
return "ref_" + id
}
func (h *HistoryQueue) addResource(ctx context.Context, l leases.Lease, desc *controlapi.Descriptor) error {
if desc == nil {
return nil
}
return h.LeaseManager.AddResource(ctx, l, leases.Resource{
ID: string(desc.Digest),
Type: "content",
})
}
func (h *HistoryQueue) Update(ctx context.Context, e *controlapi.BuildHistoryEvent) error {
h.init()
h.mu.Lock()
defer h.mu.Unlock()
if e.Type == controlapi.BuildHistoryEventType_STARTED {
h.active[e.Record.Ref] = e.Record
h.ps.Send(e)
}
if e.Type == controlapi.BuildHistoryEventType_COMPLETE {
delete(h.active, e.Record.Ref)
if err := h.DB.Update(func(tx *bolt.Tx) (err error) {
b := tx.Bucket([]byte(recordsBucket))
if b == nil {
return nil
}
dt, err := e.Record.Marshal()
if err != nil {
return err
}
l, err := h.LeaseManager.Create(ctx, leases.WithID(h.leaseID(e.Record.Ref)))
if err != nil {
return err
}
defer func() {
if err != nil {
h.LeaseManager.Delete(ctx, l)
}
}()
if err := h.addResource(ctx, l, e.Record.Logs); err != nil {
return err
}
return b.Put([]byte(e.Record.Ref), dt)
}); err != nil {
return err
}
h.ps.Send(e)
}
return nil
}
func (h *HistoryQueue) Listen(ctx context.Context, ref string, active bool, f func(*controlapi.BuildHistoryEvent) error) error {
h.init()
h.mu.Lock()
sub := h.ps.Subscribe()
defer sub.close()
for _, e := range h.active {
sub.ps.Send(&controlapi.BuildHistoryEvent{
Type: controlapi.BuildHistoryEventType_STARTED,
Record: e,
})
}
h.mu.Unlock()
if !active {
if err := h.DB.View(func(tx *bolt.Tx) error {
b := tx.Bucket([]byte(recordsBucket))
if b == nil {
return nil
}
return b.ForEach(func(key, dt []byte) error {
var br controlapi.BuildHistoryRecord
if err := br.Unmarshal(dt); err != nil {
return errors.Wrapf(err, "failed to unmarshal build record %s", key)
}
if err := f(&controlapi.BuildHistoryEvent{
Record: &br,
Type: controlapi.BuildHistoryEventType_COMPLETE,
}); err != nil {
return err
}
return nil
})
}); err != nil {
return err
}
}
for {
select {
case <-ctx.Done():
return ctx.Err()
case e := <-sub.ch:
if err := f(e); err != nil {
return err
}
case <-sub.done:
return nil
}
}
}
type pubsub[T any] struct {
mu sync.Mutex
m map[*channel[T]]struct{}
}
func (p *pubsub[T]) Subscribe() *channel[T] {
p.mu.Lock()
c := &channel[T]{
ps: p,
ch: make(chan T, 32),
done: make(chan struct{}),
}
p.m[c] = struct{}{}
p.mu.Unlock()
return c
}
func (p *pubsub[T]) Send(v T) {
p.mu.Lock()
for c := range p.m {
go c.send(v)
}
p.mu.Unlock()
}
type channel[T any] struct {
ps *pubsub[T]
ch chan T
done chan struct{}
closeOnce sync.Once
}
func (p *channel[T]) send(v T) {
select {
case p.ch <- v:
case <-p.done:
}
}
func (p *channel[T]) close() {
p.closeOnce.Do(func() {
p.ps.mu.Lock()
delete(p.ps.m, p)
p.ps.mu.Unlock()
close(p.done)
})
}

View File

@@ -7,6 +7,7 @@ import (
"strings"
"time"
controlapi "github.com/moby/buildkit/api/services/control"
"github.com/moby/buildkit/cache"
cacheconfig "github.com/moby/buildkit/cache/config"
"github.com/moby/buildkit/cache/remotecache"
@@ -24,11 +25,14 @@ import (
"github.com/moby/buildkit/util/buildinfo"
"github.com/moby/buildkit/util/compression"
"github.com/moby/buildkit/util/entitlements"
"github.com/moby/buildkit/util/grpcerrors"
"github.com/moby/buildkit/util/progress"
"github.com/moby/buildkit/worker"
digest "github.com/opencontainers/go-digest"
"github.com/pkg/errors"
"golang.org/x/sync/errgroup"
"google.golang.org/grpc/codes"
"google.golang.org/grpc/status"
)
const keyEntitlements = "llb.entitlements"
@@ -55,6 +59,7 @@ type Opt struct {
GatewayForwarder *controlgateway.GatewayForwarder
SessionManager *session.Manager
WorkerController *worker.Controller
HistoryQueue *HistoryQueue
}
type Solver struct {
@@ -67,6 +72,7 @@ type Solver struct {
gatewayForwarder *controlgateway.GatewayForwarder
sm *session.Manager
entitlements []string
history *HistoryQueue
}
// Processor defines a processing function to be applied after solving, but
@@ -83,6 +89,7 @@ func New(opt Opt) (*Solver, error) {
gatewayForwarder: opt.GatewayForwarder,
sm: opt.SessionManager,
entitlements: opt.Entitlements,
history: opt.HistoryQueue,
}
s.solver = solver.NewSolver(solver.SolverOpt{
@@ -118,7 +125,7 @@ func (s *Solver) Bridge(b solver.Builder) frontend.FrontendLLBBridge {
return s.bridge(b)
}
func (s *Solver) Solve(ctx context.Context, id string, sessionID string, req frontend.SolveRequest, exp ExporterRequest, ent []entitlements.Entitlement, post []Processor) (*client.SolveResponse, error) {
func (s *Solver) Solve(ctx context.Context, id string, sessionID string, req frontend.SolveRequest, exp ExporterRequest, ent []entitlements.Entitlement, post []Processor) (_ *client.SolveResponse, err error) {
j, err := s.solver.NewJob(id)
if err != nil {
return nil, err
@@ -126,6 +133,41 @@ func (s *Solver) Solve(ctx context.Context, id string, sessionID string, req fro
defer j.Discard()
st := time.Now()
rec := &controlapi.BuildHistoryRecord{
Ref: id,
Frontend: req.Frontend,
FrontendAttrs: req.FrontendOpt,
CreatedAt: &st,
}
if err := s.history.Update(ctx, &controlapi.BuildHistoryEvent{
Type: controlapi.BuildHistoryEventType_STARTED,
Record: rec,
}); err != nil {
return nil, err
}
defer func() {
en := time.Now()
rec.CompletedAt = &en
if err != nil {
st, ok := grpcerrors.AsGRPCStatus(grpcerrors.ToGRPC(err))
if !ok {
st = status.New(codes.Unknown, err.Error())
}
rec.Error = grpcerrors.ToRPCStatus(st.Proto())
}
if err1 := s.history.Update(ctx, &controlapi.BuildHistoryEvent{
Type: controlapi.BuildHistoryEventType_COMPLETE,
Record: rec,
}); err1 != nil {
if err == nil {
err = err1
}
}
}()
set, err := entitlements.WhiteList(ent, supportedEntitlements(s.entitlements))
if err != nil {
return nil, err

View File

@@ -1,3 +1,3 @@
package pb
//go:generate protoc -I=. -I=../../vendor/ --gogofaster_out=. ops.proto
//go:generate protoc -I=. -I=../../vendor/ -I=../../vendor/github.com/gogo/protobuf/ --gogofaster_out=. ops.proto

View File

@@ -3,6 +3,7 @@ package solver
import (
"context"
"io"
"sort"
"time"
"github.com/moby/buildkit/util/bklog"
@@ -72,6 +73,22 @@ func (j *Job) Status(ctx context.Context, ch chan *client.SolveStatus) error {
ss.Warnings = append(ss.Warnings, &v)
}
}
sort.Slice(ss.Vertexes, func(i, j int) bool {
if ss.Vertexes[i].Started == nil {
return true
}
if ss.Vertexes[j].Started == nil {
return false
}
return ss.Vertexes[i].Started.Before(*ss.Vertexes[j].Started)
})
sort.Slice(ss.Statuses, func(i, j int) bool {
return ss.Statuses[i].Timestamp.Before(ss.Statuses[j].Timestamp)
})
sort.Slice(ss.Logs, func(i, j int) bool {
return ss.Logs[i].Timestamp.Before(ss.Logs[j].Timestamp)
})
select {
case <-ctx.Done():
return ctx.Err()

View File

@@ -5,6 +5,7 @@ import (
"errors"
"github.com/containerd/typeurl"
rpc "github.com/gogo/googleapis/google/rpc"
gogotypes "github.com/gogo/protobuf/types"
"github.com/golang/protobuf/proto" //nolint:staticcheck
"github.com/golang/protobuf/ptypes/any"
@@ -196,6 +197,20 @@ func FromGRPC(err error) error {
return stack.Enable(err)
}
func ToRPCStatus(st *spb.Status) *rpc.Status {
details := make([]*gogotypes.Any, len(st.Details))
for i, d := range st.Details {
details[i] = gogoAny(d)
}
return &rpc.Status{
Code: int32(st.Code),
Message: st.Message,
Details: details,
}
}
type grpcStatusError struct {
st *status.Status
}

View File

@@ -12,6 +12,7 @@ type MultiReader struct {
initialized bool
done chan struct{}
writers map[*progressWriter]func()
sent []*Progress
}
func NewMultiReader(pr Reader) *MultiReader {
@@ -31,9 +32,59 @@ func (mr *MultiReader) Reader(ctx context.Context) Reader {
pw, _, ctx := NewFromContext(ctx)
w := pw.(*progressWriter)
mr.writers[w] = closeWriter
isBehind := len(mr.sent) > 0
if !isBehind {
mr.writers[w] = closeWriter
}
go func() {
if isBehind {
close := func() {
w.Close()
closeWriter()
}
i := 0
for {
mr.mu.Lock()
sent := mr.sent
count := len(sent) - i
if count == 0 {
select {
case <-ctx.Done():
close()
mr.mu.Unlock()
return
case <-mr.done:
close()
mr.mu.Unlock()
return
default:
}
mr.writers[w] = closeWriter
mr.mu.Unlock()
break
}
mr.mu.Unlock()
for i, p := range sent[i:] {
w.writeRawProgress(p)
if i%100 == 0 {
select {
case <-ctx.Done():
close()
return
case <-mr.done:
close()
return
default:
}
}
}
i += count
}
}
select {
case <-ctx.Done():
case <-mr.done:
@@ -61,6 +112,7 @@ func (mr *MultiReader) handle() error {
w.Close()
c()
}
close(mr.done)
mr.mu.Unlock()
return nil
}
@@ -72,6 +124,7 @@ func (mr *MultiReader) handle() error {
w.writeRawProgress(p)
}
}
mr.sent = append(mr.sent, p...)
mr.mu.Unlock()
}
}

View File

@@ -220,6 +220,10 @@ func (w *Worker) ContentStore() content.Store {
return w.WorkerOpt.ContentStore
}
func (w *Worker) LeaseManager() leases.Manager {
return w.WorkerOpt.LeaseManager
}
func (w *Worker) ID() string {
return w.WorkerOpt.ID
}
@@ -378,7 +382,7 @@ func (w *Worker) Exporter(name string, sm *session.Manager) (exporter.Exporter,
SessionManager: sm,
ImageWriter: w.imageWriter,
RegistryHosts: w.RegistryHosts,
LeaseManager: w.LeaseManager,
LeaseManager: w.LeaseManager(),
})
case client.ExporterLocal:
return localexporter.New(localexporter.Opt{
@@ -393,14 +397,14 @@ func (w *Worker) Exporter(name string, sm *session.Manager) (exporter.Exporter,
SessionManager: sm,
ImageWriter: w.imageWriter,
Variant: ociexporter.VariantOCI,
LeaseManager: w.LeaseManager,
LeaseManager: w.LeaseManager(),
})
case client.ExporterDocker:
return ociexporter.New(ociexporter.Opt{
SessionManager: sm,
ImageWriter: w.imageWriter,
Variant: ociexporter.VariantDocker,
LeaseManager: w.LeaseManager,
LeaseManager: w.LeaseManager(),
})
default:
return nil, errors.Errorf("exporter %q could not be found", name)

View File

@@ -5,6 +5,7 @@ import (
"io"
"github.com/containerd/containerd/content"
"github.com/containerd/containerd/leases"
"github.com/moby/buildkit/cache"
"github.com/moby/buildkit/client"
"github.com/moby/buildkit/client/llb"
@@ -38,6 +39,7 @@ type Worker interface {
ContentStore() content.Store
Executor() executor.Executor
CacheManager() cache.Manager
LeaseManager() leases.Manager
}
type Infos interface {