Merge pull request #84 from entireio/fix/pack-observer-non-fatal · Entire
Merge pull request #84 from entireio/fix/pack-observer-non-fatal
146fcd5→main·
Soph·4w ago·2 files·+52 added/-1 removed
Keep a pack-observer Scanner error from aborting the upload
Changes
2
internal/strategy/bootstrap
Mpack_observer.go+25/-1
Mpack_observer_test.go+27
77 unmodified lines
78
79
80
81
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
77 unmodified lines
headerReady: make(chan struct{}),
done: make(chan struct{}),
}
obs.tee = io.TeeReader(src, pw)
// Tee through a best-effort writer: once the consume goroutine stops and
// closes pr (on a Scanner error or after a malformed pack), writes to pw
// would fail with io.ErrClosedPipe and TeeReader would surface that from
// Read, aborting the live upload. Observation is non-fatal, so absorb the
// failure and keep the bytes flowing to the server.
obs.tee = io.TeeReader(src, &bestEffortWriter{w: pw})
go obs.consume(pr)
return obs
}
// bestEffortWriter forwards bytes to the observer pipe but never propagates a
// write failure: after the first error it silently drops subsequent writes and
// always reports success, so the wrapping TeeReader cannot turn a stopped
// observer into a failed upload. Touched only by the upload's Read goroutine.
type bestEffortWriter struct {
w io.Writer
broken bool
}
func (b *bestEffortWriter) Write(p []byte) (int, error) {
if b.broken {
return len(p), nil
}
if _, err := b.w.Write(p); err != nil {
b.broken = true
}
return len(p), nil
}
func (o *packStreamObserver) Read(p []byte) (int, error) {
if o.aborted.Load() {
return 0, ErrPackUploadAborted
Minternal/strategy/bootstrap/pack_observer.go+25/-1
115 unmodified lines
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
115 unmodified lines
}
}
// A Scanner error mid-stream must not abort the upload: every byte handed in
// must still come back out of Read (the documented "non-fatal for the upload"
// contract), with the error recorded only for debugging. Previously the
// Scanner's deferred pipe close made the TeeReader's write fail and killed the
// push.
func TestPackStreamObserverScannerErrorDoesNotAbortUpload(t *testing.T) {
t.Parallel()
// Not a valid packfile — the Scanner fails parsing the header almost
// immediately and closes its pipe reader while bytes are still flowing.
payload := bytes.Repeat([]byte("definitely-not-a-packfile\n"), 2048)
o := newPackStreamObserver(io.NopCloser(bytes.NewReader(payload)))
out, err := io.ReadAll(o)
if err != nil {
t.Fatalf("upload Read aborted by a non-fatal observer error: %v", err)
}
if !bytes.Equal(out, payload) {
t.Fatalf("observer dropped bytes: got %d, want %d", len(out), len(payload))
}
if err := o.Close(); err != nil {
t.Fatalf("close: %v", err)
}
if o.ScannerError() == nil {
t.Fatal("expected a scanner error to be recorded for non-pack input")
}
}
// TestPackStreamObserverHeaderReadyEarly verifies that the header is
// observed before all bytes have been pulled — important for callers
// that want to make subdivision decisions partway through the upload.