Close incremental pack streams on read errors · Entire
Close incremental pack streams on read errors
Sessions
Changes
3
- docs
- Mrewrite-issue-list.md +1/-1
- internal/strategy/incremental
- Mincremental.go +30
- Mincremental_test.go +69
156 unmodified lines
156 unmodified lines
Current rewrite note:
- Ownership of stream lifecycle is clearer than on `main`, and the rewrite now has direct tests for key pack-stream close behavior on success and error paths.
- Direct strategy-level error-path tests now verify that relay bootstrap and incremental paths close source pack streams when pushes fail.
- Direct strategy-level error-path tests now verify that relay bootstrap and incremental paths close source pack streams when pushes fail, and incremental relay now uses the same close-once ownership model as bootstrap instead of relying on the target pusher to close on error.
- Batched integration coverage now also exercises a failed checkpoint pack push followed by a resume-from-temp-ref retry.
- Batched bootstrap strategy tests now also cover an actual mid-read checkpoint stream interruption where the target pusher returns the read error without closing the stream.
- Lower-level `gitproto.PushPack` rejection paths now also close the provided pack stream instead of leaking it on preflight command errors.
Code Snippet
package main
import (
"context"
"fmt"
"io"
"sync"
"github.com/go-git/go-git/v6/plumbing"
)
// closeOnceReadCloser wraps an io.ReadCloser for ensuring it is closed only once.
type closeOnceReadCloser struct {
io.ReadCloser
once sync.Once
}
func (c *closeOnceReadCloser) Close() error {
var err error
c.once.Do(func() {
err = c.ReadCloser.Close()
})
return err
}
func closeOnce(rc io.ReadCloser) io.ReadCloser {
if rc == nil {
return nil
}
if _, ok := rc.(*closeOnceReadCloser); ok {
return rc
}
return &closeOnceReadCloser{ReadCloser: rc}
}
Tests
func TestExecuteIncrementalRelayUsesTargetRefsAsHaves(t *testing.T) {
mainRef := plumbing.NewBranchReferenceName("main")
oldHash := plumbing.NewHash("1111111111111111111111111111111111111111")
}
func TestExecuteIncrementalRelayClosesPackOnReadInterruption(t *testing.T) {
mainRef := plumbing.NewBranchReferenceName("main")
oldHash := plumbing.NewHash("1111111111111111111111111111111111111111")
newHash := plumbing.NewHash("2222222222222222222222222222222222222222")
pack := &interruptedReadCloser{first: []byte("PACK"), err: io.ErrUnexpectedEOF}
_, err := Execute(context.Background(), Params{
SourceService: fakeSourceService{
fetchPack: func(_ context.Context, _ *gitproto.Conn, _ map[plumbing.ReferenceName]gitproto.DesiredRef, _ map[plumbing.ReferenceName]plumbing.Hash) (io.ReadCloser, error) {
return pack, nil
},
},
TargetPusher: fakeTargetPusher{
pushPack: func(_ context.Context, _ []gitproto.PushCommand, pack io.ReadCloser) error {
_, err := io.Copy(io.Discard, pack)
return err
},
},
DesiredRefs: map[plumbing.ReferenceName]planner.DesiredRef{
mainRef: {
SourceRef: mainRef,
TargetRef: mainRef,
SourceHash: newHash,
Kind: planner.RefKindBranch,
},
},
TargetRefs: map[plumbing.ReferenceName]plumbing.Hash{mainRef: oldHash},
PushPlans: []planner.BranchPlan{{
SourceRef: mainRef,
TargetRef: mainRef,
SourceHash: newHash,
TargetHash: oldHash,
Kind: planner.RefKindBranch,
Action: planner.ActionUpdate,
}},
CanRelay: func(bool, bool, bool, []planner.BranchPlan) (bool, string) {
return true, "fast-forward"
},
}, planner.PlanConfig{})
if !errors.Is(err, io.ErrUnexpectedEOF) {
t.Fatalf("expected interrupted read error, got %%v", err)
}
if !pack.closed {
t.Fatal("expected pack to be closed after read interruption")
}
}
func TestToGP(t *testing.T) {
hash1 := plumbing.NewHash("aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa")
hash2 := plumbing.NewHash("bbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbb")
}