control: base of status reporting

Signed-off-by: Tonis Tiigi <tonistiigi@gmail.com>
This commit is contained in:
Tonis Tiigi
2017-06-13 14:42:51 -07:00
parent 94d62fd432
commit 62b7d04d01
10 changed files with 2042 additions and 80 deletions

File diff suppressed because it is too large Load Diff

View File

@@ -5,7 +5,7 @@ package control;
service Control {
rpc DiskUsage(DiskUsageRequest) returns (DiskUsageResponse);
rpc Solve(SolveRequest) returns (SolveResponse);
// rpc Status() returns ();
rpc Status(StatusRequest) returns (stream StatusResponse);
}
message DiskUsageRequest {
@@ -28,8 +28,40 @@ message SolveRequest {
}
message SolveResponse {
repeated VertexStatus vertex = 1;
repeated Vertex vtx = 1;
}
message StatusRequest {
string Ref = 1;
}
message StatusResponse {
repeated Vertex vtx = 1;
repeated VertexStatus status = 2;
repeated VertexLog log = 3;
}
message Vertex {
string ID = 1;
repeated string inputs = 2;
string name = 3;
repeated VertexStatus status = 4;
bool cached = 5;
int64 started = 6; // relative, add abolute google.protobuf.Timestamp as well?
int64 completed = 7;
}
message VertexStatus {
string ID = 1;
string vertex = 2;
string name = 3;
int64 current = 4;
int64 total = 5;
}
message VertexLog {
int64 inc = 1;
int64 timestamp = 2;
int64 stream = 3;
bytes msg = 4;
}

49
client/graph.go Normal file
View File

@@ -0,0 +1,49 @@
package client
import (
"time"
digest "github.com/opencontainers/go-digest"
)
type Vertex struct {
ID digest.Digest
Inputs []digest.Digest
Name string
Started time.Time
Completed time.Time
Cached bool
Error string
}
type VertexStatus struct {
ID digest.Digest
Vertex digest.Digest
Name string
Total int
Current int
Timestamp time.Time
}
type VertexLog struct {
Vertex digest.Digest
Stream int
Data []byte
Timestamp time.Time
}
type SolveStatus struct {
Vertexes []*Vertex
Statuses []*VertexStatus
Logs []*VertexLog
}
//
// type VertexEvent struct {
// ID digest.Digest
// Vertex digest.Digest
// Name string
// Total int
// Current int
// Timestamp int64
// }

View File

@@ -64,8 +64,12 @@ func (c *Controller) Solve(ctx context.Context, req *controlapi.SolveRequest) (*
if err != nil {
return nil, errors.Wrap(err, "failed to load")
}
if err := c.solver.Solve(ctx, v); err != nil {
if err := c.solver.Solve(ctx, req.Ref, v); err != nil {
return nil, err
}
return &controlapi.SolveResponse{}, nil
}
func (c *Controller) Status(*controlapi.StatusRequest, controlapi.Control_StatusServer) error {
return errors.Errorf("not implemented")
}

View File

@@ -3,7 +3,6 @@
package control
import (
"context"
"io"
"io/ioutil"
"os"
@@ -18,6 +17,7 @@ import (
ocispec "github.com/opencontainers/image-spec/specs-go/v1"
"github.com/pkg/errors"
"github.com/tonistiigi/buildkit_poc/worker/runcworker"
"golang.org/x/net/context"
)
func NewStandalone(root string) (*Controller, error) {

View File

@@ -3,15 +3,19 @@ package solver
import (
"context"
"os"
"strings"
"sync"
"time"
"golang.org/x/sync/errgroup"
digest "github.com/opencontainers/go-digest"
"github.com/pkg/errors"
"github.com/tonistiigi/buildkit_poc/cache"
"github.com/tonistiigi/buildkit_poc/client"
"github.com/tonistiigi/buildkit_poc/solver/pb"
"github.com/tonistiigi/buildkit_poc/source"
"github.com/tonistiigi/buildkit_poc/util/progress"
"github.com/tonistiigi/buildkit_poc/worker"
)
@@ -22,6 +26,7 @@ type opVertex struct {
refs []cache.ImmutableRef
err error
dgst digest.Digest
vtx client.Vertex
}
func Load(ops [][]byte) (*opVertex, error) {
@@ -63,17 +68,25 @@ func loadReqursive(dgst digest.Digest, op *pb.Op, inputs map[digest.Digest]*pb.O
return v, nil
}
vtx := &opVertex{op: op, dgst: dgst}
inputDigests := make([]digest.Digest, 0, len(op.Inputs))
for _, in := range op.Inputs {
op, ok := inputs[digest.Digest(in.Digest)]
dgst := digest.Digest(in.Digest)
inputDigests = append(inputDigests, dgst)
op, ok := inputs[dgst]
if !ok {
return nil, errors.Errorf("failed to find %s", in)
}
sub, err := loadReqursive(digest.Digest(in.Digest), op, inputs, cache)
sub, err := loadReqursive(dgst, op, inputs, cache)
if err != nil {
return nil, err
}
vtx.inputs = append(vtx.inputs, sub)
}
vtx.vtx = client.Vertex{
Inputs: inputDigests,
Name: vtx.name(),
ID: dgst,
}
cache[dgst] = vtx
return vtx, nil
}
@@ -89,19 +102,44 @@ func (g *opVertex) inputRequiresExport(i int) bool {
}
type Solver struct {
opt Opt
opt Opt
jobs jobs
}
func New(opt Opt) *Solver {
return &Solver{opt: opt}
}
func (s *Solver) Solve(ctx context.Context, g *opVertex) error {
err := g.solve(ctx, s.opt) // TODO: separate exporting
func (s *Solver) Solve(ctx context.Context, id string, g *opVertex) error {
// ctx, cancel := context.WithCancel(ctx)
// defer cancel()
pr, ctx, closeProgressWriter := progress.NewContext(ctx)
_, err := s.jobs.new(id, g, pr)
if err != nil {
return err
}
err = g.solve(ctx, s.opt) // TODO: separate exporting
closeProgressWriter()
if err != nil {
return err
}
g.release(ctx)
// TODO: export final vertex state
return err
}
func (s *Solver) Status(ctx context.Context, id string, statusChan chan *client.SolveStatus) error {
j, err := s.jobs.get(id)
if err != nil {
return nil
}
return j.pipe(ctx, statusChan)
}
func (g *opVertex) release(ctx context.Context) (retErr error) {
for _, i := range g.inputs {
if err := i.release(ctx); err != nil {
@@ -150,8 +188,7 @@ func (g *opVertex) solve(ctx context.Context, opt Opt) (retErr error) {
for _, in := range g.inputs {
eg.Go(func() error {
err := in.solve(ctx, opt)
if err != nil {
if err := in.solve(ctx, opt); err != nil {
return err
}
return nil
@@ -163,6 +200,12 @@ func (g *opVertex) solve(ctx context.Context, opt Opt) (retErr error) {
}
}
pw, _, ctx := progress.FromContext(ctx, g.dgst.String())
defer pw.Done()
g.notifyStarted(pw)
defer g.notifyComplete(pw)
switch op := g.op.Op.(type) {
case *pb.Op_Source:
id, err := source.FromString(op.Source.Identifier)
@@ -232,3 +275,28 @@ func (g *opVertex) solve(ctx context.Context, opt Opt) (retErr error) {
}
return nil
}
func (g *opVertex) notifyStarted(pw progress.ProgressWriter) {
g.vtx.Started = time.Now()
pw.Write(g.vtx)
}
func (g *opVertex) notifyComplete(pw progress.ProgressWriter) {
g.vtx.Completed = time.Now()
pw.Write(g.vtx)
}
func (g *opVertex) name() string {
switch op := g.op.Op.(type) {
case *pb.Op_Source:
return op.Source.Identifier
case *pb.Op_Exec:
name := strings.Join(op.Exec.Meta.Args, " ")
if len(name) > 22 { // TODO: const
name = name[:20] + "..."
}
return name
default:
return "unknown"
}
}

102
solver/run.go Normal file
View File

@@ -0,0 +1,102 @@
package solver
import (
"context"
"sync"
digest "github.com/opencontainers/go-digest"
"github.com/pkg/errors"
"github.com/tonistiigi/buildkit_poc/client"
"github.com/tonistiigi/buildkit_poc/util/progress"
)
type jobs struct {
mu sync.RWMutex
refs map[string]*job
}
func (j *jobs) new(id string, g *opVertex, pr progress.ProgressReader) (*job, error) {
j.mu.Lock()
defer j.mu.Unlock()
if j.refs == nil {
j.refs = make(map[string]*job)
}
if _, ok := j.refs[id]; ok {
return nil, errors.Errorf("id %s exists", id)
}
nj := &job{g: g, pr: progress.NewMultiReader(pr)}
j.refs[id] = nj
go func() {
j.mu.Lock()
defer j.mu.Unlock()
delete(j.refs, id)
}()
return j.refs[id], nil
}
func (j *jobs) get(id string) (*job, error) {
j.mu.RLock()
defer j.mu.RUnlock()
nj, ok := j.refs[id]
if !ok {
return nil, errors.Errorf("no such job %s", id)
}
return nj, nil
}
type job struct {
mu sync.Mutex
g *opVertex
pr *progress.MultiReader
}
func (j *job) pipe(ctx context.Context, ch chan *client.SolveStatus) error {
pr := j.pr.Reader(ctx)
for v := range flatten(j.g) {
ss := &client.SolveStatus{
Vertexes: []*client.Vertex{&v.vtx},
}
select {
case <-ctx.Done():
return ctx.Err()
case ch <- ss:
}
}
for {
p, err := pr.Read(ctx) // add cancelling
if err != nil {
return err
}
switch v := p.Sys.(type) {
case *client.Vertex:
ss := &client.SolveStatus{Vertexes: []*client.Vertex{v}}
select {
case <-ctx.Done():
return ctx.Err()
case ch <- ss:
}
}
}
return nil
}
func flatten(op *opVertex) chan *opVertex {
cache := make(map[digest.Digest]struct{})
ch := make(chan *opVertex, 32)
go sendVertex(ch, op, cache)
return ch
}
func sendVertex(ch chan *opVertex, op *opVertex, cache map[digest.Digest]struct{}) {
for _, v := range op.inputs {
sendVertex(ch, v, cache)
}
if _, ok := cache[op.dgst]; !ok {
ch <- op
cache[op.dgst] = struct{}{}
}
}

View File

@@ -0,0 +1,67 @@
package progress
import (
"context"
"sync"
)
type MultiReader struct {
mu sync.Mutex
main ProgressReader
initialized bool
done chan struct{}
writers map[*progressWriter]struct{}
}
func NewMultiReader(pr ProgressReader) *MultiReader {
mr := &MultiReader{
main: pr,
done: make(chan struct{}),
}
return mr
}
func (mr *MultiReader) Reader(ctx context.Context) ProgressReader {
mr.mu.Lock()
defer mr.mu.Unlock()
pr, ctx, _ := NewContext(ctx)
pw, _, _ := FromContext(ctx, "")
w := pw.(*progressWriter)
mr.writers[w] = struct{}{}
go func() {
select {
case <-ctx.Done():
case <-mr.done:
}
mr.mu.Lock()
defer mr.mu.Unlock()
delete(mr.writers, w)
}()
if !mr.initialized {
go mr.handle()
mr.initialized = true
}
return pr
}
func (mr *MultiReader) handle() error {
for {
p, err := mr.main.Read(context.TODO())
if err != nil {
return err
}
if p == nil {
return nil
}
mr.mu.Lock()
for w := range mr.writers {
w.write(*p)
}
mr.mu.Unlock()
}
}

View File

@@ -20,7 +20,7 @@ func FromContext(ctx context.Context, name string) (ProgressWriter, bool, contex
}
pw = newWriter(pw, name)
ctx = context.WithValue(ctx, contextKey, pw)
return pw, false, ctx
return pw, true, ctx
}
func NewContext(ctx context.Context) (ProgressReader, context.Context, func()) {
@@ -30,7 +30,7 @@ func NewContext(ctx context.Context) (ProgressReader, context.Context, func()) {
}
type ProgressWriter interface {
Write(Progress) error
Write(interface{}) error
Done() error
}
@@ -39,17 +39,17 @@ type ProgressReader interface {
}
type Progress struct {
ID string
// Progress contains a Message or...
Message string
// ...progress of an action
Action string
Current int
Total int
ID string
Timestamp time.Time
Done bool
Sys interface{}
}
type Status struct {
// ...progress of an action
Action string
Current int
Total int
}
type progressReader struct {
@@ -95,7 +95,7 @@ func (pr *progressReader) Read(ctx context.Context) (*Progress, error) {
default:
}
open := false
for _, sh := range pr.handles { // could be more efficient but unlikely that this array will be very big, maybe random ordering?
for _, sh := range pr.handles { // could be more efficient but unlikely that this array will be very big, maybe random ordering? at least remove the completed handlers.
p, ok := sh.next()
if ok {
pr.mu.Unlock()
@@ -148,7 +148,11 @@ func pipe() (*progressReader, *progressWriter, func()) {
func newWriter(pw *progressWriter, name string) *progressWriter {
if pw.id != "" {
name = pw.id + "." + name
if name == "" {
name = pw.id
} else {
name = pw.id + "." + name
}
}
pw = &progressWriter{
id: name,
@@ -163,20 +167,27 @@ type progressWriter struct {
lastP atomic.Value
done bool
reader *progressReader
byKey map[string]atomic.Value
items []atomic.Value
}
func (pw *progressWriter) Write(p Progress) error {
func (pw *progressWriter) Write(s interface{}) error {
if pw.done {
return errors.Errorf("writing to closed progresswriter %s", pw.id)
}
var p Progress
p.ID = pw.id
if p.Timestamp.IsZero() {
p.Timestamp = time.Now()
}
pw.lastP.Store(&p)
p.Timestamp = time.Now()
p.Sys = s
return pw.write(p)
}
func (pw *progressWriter) write(p Progress) error {
if p.Done {
pw.done = true
}
pw.lastP.Store(&p)
pw.reader.cond.Broadcast()
return nil
}
@@ -184,21 +195,24 @@ func (pw *progressWriter) Write(p Progress) error {
func (pw *progressWriter) Done() error {
var p Progress
lastP := pw.lastP.Load().(*Progress)
p.ID = pw.id
p.Timestamp = time.Now()
if lastP != nil {
p = *lastP
if p.Done {
return nil
}
} else {
p = Progress{}
p.Sys = lastP.Sys
}
p.Done = true
return pw.Write(p)
pw.done = true
return pw.write(p)
}
type noOpWriter struct{}
func (pw *noOpWriter) Write(p Progress) error {
func (pw *noOpWriter) Write(p interface{}) error {
return nil
}

View File

@@ -30,7 +30,8 @@ func TestProgress(t *testing.T) {
err = eg.Wait()
assert.NoError(t, err)
assert.Equal(t, 6, len(trace.items))
assert.True(t, len(trace.items) > 5)
assert.True(t, len(trace.items) <= 7)
assert.Equal(t, trace.items[len(trace.items)-1].Done, true)
}
@@ -51,7 +52,7 @@ func TestProgressNested(t *testing.T) {
assert.NoError(t, err)
assert.True(t, len(trace.items) > 9) // usually 14
assert.True(t, len(trace.items) <= 14)
assert.True(t, len(trace.items) <= 15)
streams := 0
for _, t := range trace.items {
if t.Done {
@@ -66,7 +67,7 @@ func calc(ctx context.Context, total int, name string) (int, error) {
defer pw.Done()
sum := 0
pw.Write(Progress{Action: "starting", Total: total})
pw.Write(Status{Action: "starting", Total: total})
for i := 1; i <= total; i++ {
select {
case <-ctx.Done():
@@ -74,12 +75,13 @@ func calc(ctx context.Context, total int, name string) (int, error) {
case <-time.After(10 * time.Millisecond):
}
if i == total {
pw.Write(Progress{Action: "done", Total: total, Current: total, Done: true})
pw.Write(Status{Action: "done", Total: total, Current: total})
} else {
pw.Write(Progress{Action: "calculating", Total: total, Current: i})
pw.Write(Status{Action: "calculating", Total: total, Current: i})
}
sum += i
}
pw.Done()
return sum, nil
}
@@ -90,7 +92,7 @@ func reduceCalc(ctx context.Context, total int) (int, error) {
pw, _, ctx := FromContext(ctx, "reduce")
defer pw.Done()
pw.Write(Progress{Action: "starting"})
pw.Write(Status{Action: "starting"})
// sync step
sum, err := calc(ctx, total, "synccalc")