Skip to content

Commit 69ff057

Browse files
committed
Harden staged push polling
1 parent 6437409 commit 69ff057

3 files changed

Lines changed: 109 additions & 50 deletions

File tree

pkg/cmd/push.go

Lines changed: 23 additions & 15 deletions
Original file line numberDiff line numberDiff line change
@@ -95,10 +95,31 @@ func handleLocalPush(ctx context.Context, cmd *cli.Command) error {
9595
targetName = args[1]
9696
}
9797

98-
return pushLocalImage(ctx, cmd, sourceImage, targetName, nil)
98+
return pushLocalImage(ctx, cmd, sourceImage, targetName)
9999
}
100100

101-
func pushLocalImage(ctx context.Context, cmd *cli.Command, sourceImage, targetName string, img v1.Image) error {
101+
func pushLocalImage(ctx context.Context, cmd *cli.Command, sourceImage, targetName string) error {
102+
fmt.Fprintf(os.Stderr, "Loading image %s from Docker...\n", sourceImage)
103+
img, err := loadDockerImage(sourceImage)
104+
if err != nil {
105+
return err
106+
}
107+
return uploadLocalImage(ctx, cmd, targetName, img)
108+
}
109+
110+
func loadDockerImage(image string) (v1.Image, error) {
111+
srcRef, err := name.ParseReference(image)
112+
if err != nil {
113+
return nil, fmt.Errorf("invalid source image: %w", err)
114+
}
115+
img, err := daemon.Image(srcRef)
116+
if err != nil {
117+
return nil, fmt.Errorf("load image: %w", err)
118+
}
119+
return img, nil
120+
}
121+
122+
func uploadLocalImage(ctx context.Context, cmd *cli.Command, targetName string, img v1.Image) error {
102123
baseURL := resolveBaseURL(cmd)
103124

104125
parsedURL, err := url.Parse(baseURL)
@@ -114,11 +135,6 @@ func pushLocalImage(ctx context.Context, cmd *cli.Command, sourceImage, targetNa
114135

115136
registryHost := parsedURL.Host
116137

117-
srcRef, err := name.ParseReference(sourceImage)
118-
if err != nil {
119-
return fmt.Errorf("invalid source image: %w", err)
120-
}
121-
122138
// Build and validate the target before opening the Docker daemon. The
123139
// server computes the image digest from the manifest, while the tag keeps
124140
// the image addressable with Docker-like image names after the push.
@@ -132,14 +148,6 @@ func pushLocalImage(ctx context.Context, cmd *cli.Command, sourceImage, targetNa
132148
return fmt.Errorf("invalid target: %w", err)
133149
}
134150

135-
if img == nil {
136-
fmt.Fprintf(os.Stderr, "Loading image %s from Docker...\n", sourceImage)
137-
img, err = daemon.Image(srcRef)
138-
if err != nil {
139-
return fmt.Errorf("load image: %w", err)
140-
}
141-
}
142-
143151
fmt.Fprintf(os.Stderr, "The push refers to repository [%s]\n", dstRef.Context().Name())
144152

145153
token := resolveAPIKey()

pkg/cmd/push_test.go

Lines changed: 25 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -5,6 +5,7 @@ import (
55
"testing"
66

77
v1 "github.com/google/go-containerregistry/pkg/v1"
8+
"github.com/kernel/hypeman-go"
89
"github.com/stretchr/testify/assert"
910
)
1011

@@ -32,6 +33,30 @@ func TestRenderPushProgressNonInteractive(t *testing.T) {
3233
assert.Empty(t, output.String())
3334
}
3435

36+
func TestPushStatusText(t *testing.T) {
37+
lastBytes := int64(0)
38+
39+
assert.Equal(t, "queued · registry.example.com/app:v1", pushStatusText(&hypeman.Push{
40+
Status: hypeman.PushStatusQueued,
41+
Target: "registry.example.com/app:v1",
42+
}, &lastBytes))
43+
assert.Equal(t, "pushing 2.0 KB · 2 layers · registry.example.com/app:v1", pushStatusText(&hypeman.Push{
44+
Status: hypeman.PushStatusPushing,
45+
Bytes: 2048,
46+
Layers: 2,
47+
Target: "registry.example.com/app:v1",
48+
}, &lastBytes))
49+
assert.Equal(t, int64(2048), lastBytes)
50+
assert.Equal(t, "pushed · digest: sha256:abc", pushStatusText(&hypeman.Push{
51+
Status: hypeman.PushStatusPushed,
52+
Digest: "sha256:abc",
53+
}, &lastBytes))
54+
assert.Equal(t, "failed · registry unavailable", pushStatusText(&hypeman.Push{
55+
Status: hypeman.PushStatusFailed,
56+
Error: "registry unavailable",
57+
}, &lastBytes))
58+
}
59+
3560
func TestPushStatusRenderer(t *testing.T) {
3661
var output bytes.Buffer
3762
renderer := &pushStatusRenderer{output: &output, interactive: true}

pkg/cmd/pushcmd.go

Lines changed: 61 additions & 35 deletions
Original file line numberDiff line numberDiff line change
@@ -10,7 +10,6 @@ import (
1010
"time"
1111

1212
"github.com/google/go-containerregistry/pkg/name"
13-
"github.com/google/go-containerregistry/pkg/v1/daemon"
1413
"github.com/kernel/hypeman-go"
1514
"github.com/kernel/hypeman-go/option"
1615
"github.com/tidwall/gjson"
@@ -95,24 +94,20 @@ func handleRemotePushTarget(ctx context.Context, cmd *cli.Command, target string
9594
// The one-argument form follows Docker's local-tag flow: TARGET must be
9695
// present in the local Docker daemon before it can be staged and pushed.
9796
// Cached Hypeman images use the explicit IMAGE TARGET form instead.
98-
srcRef, err := name.ParseReference(target)
99-
if err != nil {
100-
return err
101-
}
102-
img, err := daemon.Image(srcRef)
97+
img, err := loadDockerImage(target)
10398
if err != nil {
10499
return fmt.Errorf("load local Docker image %q: %w; tag it first or use hypeman push <image> <target> for a cached Hypeman image", target, err)
105100
}
106101

107102
fmt.Fprintf(os.Stderr, "Staging local image %s in Hypeman...\n", target)
108-
if err := pushLocalImage(ctx, cmd, target, target, img); err != nil {
103+
if err := uploadLocalImage(ctx, cmd, target, img); err != nil {
109104
return err
110105
}
111106

112107
client := hypeman.NewClient(getDefaultRequestOptions(cmd)...)
113-
imported, err := client.Images.Get(ctx, url.PathEscape(target))
108+
imported, err := waitForImageRecord(ctx, &client, target)
114109
if err != nil {
115-
return fmt.Errorf("get staged image %s: %w", target, err)
110+
return err
116111
}
117112
if err := waitForImageReady(ctx, &client, imported); err != nil {
118113
return err
@@ -134,6 +129,27 @@ func validateRemotePushTarget(target string) error {
134129
return nil
135130
}
136131

132+
func waitForImageRecord(ctx context.Context, client *hypeman.Client, imageName string) (*hypeman.Image, error) {
133+
ticker := time.NewTicker(300 * time.Millisecond)
134+
defer ticker.Stop()
135+
136+
for {
137+
img, err := client.Images.Get(ctx, url.PathEscape(imageName))
138+
if err == nil {
139+
return img, nil
140+
}
141+
if !isNotFoundError(err) {
142+
return nil, fmt.Errorf("get staged image %s: %w", imageName, err)
143+
}
144+
145+
select {
146+
case <-ctx.Done():
147+
return nil, ctx.Err()
148+
case <-ticker.C:
149+
}
150+
}
151+
}
152+
137153
func handlePushCreate(ctx context.Context, cmd *cli.Command) error {
138154
args := cmd.Args().Slice()
139155
if len(args) != 2 {
@@ -226,41 +242,29 @@ func waitForPush(ctx context.Context, client *hypeman.Client, push *hypeman.Push
226242
output: os.Stderr,
227243
interactive: term.IsTerminal(int(os.Stderr.Fd())),
228244
}
245+
defer renderer.finish()
229246
}
230247

231-
current := push
232248
var lastBytes int64
233-
for {
234-
if renderer != nil {
235-
switch current.Status {
236-
case hypeman.PushStatusQueued:
237-
renderer.update(fmt.Sprintf("queued · %s", current.Target))
238-
case hypeman.PushStatusPushing:
239-
if current.Bytes > lastBytes {
240-
lastBytes = current.Bytes
241-
}
242-
renderer.update(fmt.Sprintf("pushing %s · %d layers · %s", formatBytes(lastBytes), current.Layers, current.Target))
243-
case hypeman.PushStatusPushed:
244-
renderer.update(fmt.Sprintf("pushed · digest: %s", current.Digest))
245-
case hypeman.PushStatusFailed:
246-
message := current.Error
247-
if message == "" {
248-
message = "unknown error"
249-
}
250-
renderer.update("failed · " + message)
251-
}
249+
return pollPush(ctx, client, push, opts, func(current *hypeman.Push) {
250+
if renderer == nil {
251+
return
252252
}
253+
if message := pushStatusText(current, &lastBytes); message != "" {
254+
renderer.update(message)
255+
}
256+
})
257+
}
258+
259+
func pollPush(ctx context.Context, client *hypeman.Client, push *hypeman.Push, opts []option.RequestOption, update func(*hypeman.Push)) (*hypeman.Push, error) {
260+
current := push
261+
for {
262+
update(current)
253263

254264
switch current.Status {
255265
case hypeman.PushStatusPushed:
256-
if renderer != nil {
257-
renderer.finish()
258-
}
259266
return current, nil
260267
case hypeman.PushStatusFailed:
261-
if renderer != nil {
262-
renderer.finish()
263-
}
264268
if current.Error != "" {
265269
return nil, fmt.Errorf("push %s failed: %s", push.ID, current.Error)
266270
}
@@ -283,6 +287,28 @@ func waitForPush(ctx context.Context, client *hypeman.Client, push *hypeman.Push
283287
}
284288
}
285289

290+
func pushStatusText(push *hypeman.Push, lastBytes *int64) string {
291+
switch push.Status {
292+
case hypeman.PushStatusQueued:
293+
return fmt.Sprintf("queued · %s", push.Target)
294+
case hypeman.PushStatusPushing:
295+
if push.Bytes > *lastBytes {
296+
*lastBytes = push.Bytes
297+
}
298+
return fmt.Sprintf("pushing %s · %d layers · %s", formatBytes(*lastBytes), push.Layers, push.Target)
299+
case hypeman.PushStatusPushed:
300+
return fmt.Sprintf("pushed · digest: %s", push.Digest)
301+
case hypeman.PushStatusFailed:
302+
message := push.Error
303+
if message == "" {
304+
message = "unknown error"
305+
}
306+
return "failed · " + message
307+
default:
308+
return ""
309+
}
310+
}
311+
286312
type pushStatusRenderer struct {
287313
output io.Writer
288314
interactive bool

0 commit comments

Comments
 (0)