Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
53 changes: 53 additions & 0 deletions internal/pb/pb.go
Original file line number Diff line number Diff line change
Expand Up @@ -60,6 +60,9 @@ type progressBar struct {
size int64
msg atomic.Value // stores string; accessed by mpb render goroutine
startTime time.Time
// indeterminate marks a spinner bar created by Spinner. It has no byte
// counter and no completion trigger until Complete, Reset or Abort.
indeterminate bool
}

// NewProgressBar creates a new progress bar.
Expand Down Expand Up @@ -141,6 +144,51 @@ func (p *ProgressBar) Add(prompt, name string, size int64, reader io.Reader) io.
return reader
}

// Spinner adds an indeterminate bar for a phase that transfers no bytes, such
// as a remote existence check. It renders a spinner, the prompt, the total
// size, and the elapsed time. It has no byte counter and no transfer rate, so
// a slow phase does not show a misleading "0 B / N B" and "0 B/s". Call Reset
// (or Add) with the same name to replace it with a transfer bar once bytes
// start flowing, or Complete / Abort to finish it.
func (p *ProgressBar) Spinner(prompt, name string, size int64) {
if disableProgress.Load() {
return
}

p.mu.Lock()
defer p.mu.Unlock()

// Same replacement semantics as Add.
if oldBar := p.bars[name]; oldBar != nil {
oldBar.Abort(true)
}

newBar := &progressBar{
size: size,
startTime: time.Now(),
indeterminate: true,
}
newBar.msg.Store(fmt.Sprintf("%s %s", prompt, name))

// A zero total disables mpb's completion trigger, so the bar keeps
// spinning until it is replaced, completed, or aborted.
newBar.Bar = p.mpb.New(0,
mpbv8.SpinnerStyle(),
mpbv8.PrependDecorators(
decor.Any(func(_ decor.Statistics) string {
return newBar.msg.Load().(string)
}, decor.WCSyncSpaceR),
),
mpbv8.AppendDecorators(
decor.Name(humanize.Bytes(uint64(size)), decor.WCSyncWidthR),
decor.Name(" | ", decor.WCSyncWidthR),
decor.Elapsed(decor.ET_STYLE_GO, decor.WCSyncWidthR),
),
)

p.bars[name] = newBar
}

// Get returns the progress bar.
func (p *ProgressBar) Get(name string) *progressBar {
p.mu.RLock()
Expand All @@ -158,6 +206,11 @@ func (p *ProgressBar) Complete(name string, msg string) {

if ok {
bar.msg.Store(msg)
if bar.indeterminate {
// A spinner has no total; give it one so mpb marks it complete.
bar.SetTotal(bar.size, true)
return
}
bar.Bar.SetCurrent(bar.size)
}
}
Expand Down
124 changes: 124 additions & 0 deletions internal/pb/pb_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -23,11 +23,47 @@ import (
"strings"
"sync"
"testing"
"time"

"github.com/stretchr/testify/assert"
"github.com/stretchr/testify/require"
)

// lockedBuffer is a bytes.Buffer that is safe for concurrent use. The mpb
// render goroutine writes frames while the test reads them.
type lockedBuffer struct {
mu sync.Mutex
buf bytes.Buffer
}

func (b *lockedBuffer) Write(p []byte) (int, error) {
b.mu.Lock()
defer b.mu.Unlock()
return b.buf.Write(p)
}

func (b *lockedBuffer) String() string {
b.mu.Lock()
defer b.mu.Unlock()
return b.buf.String()
}

// waitForOutput polls the buffer until it contains want or the timeout
// expires, and returns the final contents.
func waitForOutput(t *testing.T, out *lockedBuffer, want string, timeout time.Duration) string {
t.Helper()
deadline := time.Now().Add(timeout)
for time.Now().Before(deadline) {
if s := out.String(); strings.Contains(s, want) {
return s
}
time.Sleep(50 * time.Millisecond)
}
s := out.String()
require.Contains(t, s, want, "progress output never rendered %q", want)
return s
}

// --- Functional tests ---

func TestAdd_WrapsReader(t *testing.T) {
Expand Down Expand Up @@ -113,6 +149,94 @@ func TestReset_NoExistingBar(t *testing.T) {
assert.Equal(t, "Building => new-file", bar.msg.Load().(string))
}

func TestSpinner_CreatesIndeterminateBar(t *testing.T) {
pb := NewProgressBar(io.Discard)
defer pb.Stop()

pb.Spinner("Checking =>", "test-file", 100)

bar := pb.Get("test-file")
require.NotNil(t, bar)
assert.True(t, bar.indeterminate)
assert.Equal(t, int64(100), bar.size)
assert.Equal(t, "Checking => test-file", bar.msg.Load().(string))
assert.False(t, bar.Completed(), "spinner must not complete on its own")
}

func TestSpinner_ResetSwitchesToTransferBar(t *testing.T) {
pb := NewProgressBar(io.Discard)
defer pb.Stop()

pb.Spinner("Checking =>", "test-file", 100)
spinner := pb.Get("test-file")
require.NotNil(t, spinner)

reader := pb.Reset("Pushing =>", "test-file", 100, strings.NewReader("data"))
require.NotNil(t, reader)

bar := pb.Get("test-file")
require.NotNil(t, bar)
assert.NotSame(t, spinner, bar)
assert.False(t, bar.indeterminate)
assert.Equal(t, "Pushing => test-file", bar.msg.Load().(string))
assert.True(t, spinner.Aborted(), "replaced spinner must be aborted")
}

func TestSpinner_CompleteFinishesBar(t *testing.T) {
pb := NewProgressBar(io.Discard)
defer pb.Stop()

pb.Spinner("Checking =>", "test-file", 100)
pb.Complete("test-file", "Skipped => test-file")

bar := pb.Get("test-file")
require.NotNil(t, bar)
assert.Equal(t, "Skipped => test-file", bar.msg.Load().(string))
assert.Eventually(t, bar.Completed, 2*time.Second, 20*time.Millisecond,
"Complete must finish an indeterminate bar")
}

func TestSpinner_DisabledProgress(t *testing.T) {
SetDisableProgress(true)
defer SetDisableProgress(false)

pb := NewProgressBar(io.Discard)
defer pb.Stop()

pb.Spinner("Checking =>", "test-file", 100)
assert.Nil(t, pb.Get("test-file"), "no bar must be tracked when progress is disabled")
}

func TestSpinner_RendersWithoutTransferRate(t *testing.T) {
out := &lockedBuffer{}
pb := NewProgressBar(out)

pb.Spinner("Checking =>", "test-file", 100)
rendered := waitForOutput(t, out, "Checking => test-file", 3*time.Second)
pb.Stop()

// The spinner shows the size and elapsed time, but no byte counter and
// no transfer rate. A transfer bar with no bytes flowing renders
// "0.00 b / 100.00 b | 0.00 b/s" (see TestAdd_RendersTransferRate).
assert.Contains(t, rendered, "100 B")
assert.NotContains(t, rendered, "b/s")
assert.NotContains(t, rendered, "0.00 b / 100.00 b")
}

func TestAdd_RendersTransferRate(t *testing.T) {
out := &lockedBuffer{}
pb := NewProgressBar(out)

pb.Add("Pushing =>", "test-file", 100, nil)
rendered := waitForOutput(t, out, "Pushing => test-file", 3*time.Second)
pb.Stop()

// Sanity check for the assertion above: a transfer bar does render the
// counter and the rate, so the spinner test is not vacuous.
assert.Contains(t, rendered, "0.00 b / 100.00 b")
assert.Contains(t, rendered, "0.00 b/s")
}

// --- Concurrency tests (must pass go test -race) ---

func TestAdd_ConcurrentSameName(t *testing.T) {
Expand Down
7 changes: 4 additions & 3 deletions pkg/backend/push.go
Original file line number Diff line number Diff line change
Expand Up @@ -152,10 +152,11 @@ func pushIfNotExist(ctx context.Context, pb *internalpb.ProgressBar, src storage

// Phase 1: show "Checking" during the existence check. The manifest is
// excluded since its payload is already in memory and pushed together
// with the tag right after. Bar is created with a nil reader so it
// indicates waiting state without transferring bytes.
// with the tag right after. The check transfers no bytes, so use an
// indeterminate spinner instead of a transfer bar that would sit at
// "0 B / N B" and "0 B/s" for the whole HEAD request.
if desc.MediaType != ocispec.MediaTypeImageManifest {
pb.Add(internalpb.NormalizePrompt("Checking "+kind), desc.Digest.String(), desc.Size, nil)
pb.Spinner(internalpb.NormalizePrompt("Checking "+kind), desc.Digest.String(), desc.Size)
}

// check whether the content exists in the destination storage.
Expand Down
Loading