diff --git a/.github/workflows/go-test.yml b/.github/workflows/go-test.yml index f1d7507..fabb20f 100644 --- a/.github/workflows/go-test.yml +++ b/.github/workflows/go-test.yml @@ -17,7 +17,7 @@ jobs: go-version: "1.21" - name: Install git. run: | - sudo apt-get install -y libcurl4-openssl-dev + sudo apt-get install -y cargo libcurl4-openssl-dev git clone https://github.com/git/git.git git -C git checkout next sudo make -j 16 -C git prefix=/usr NO_GETTEXT=YesPlease all install @@ -28,4 +28,3 @@ jobs: run: make test - name: go integration test run: make test-integration - diff --git a/internal/spokes/spokes.go b/internal/spokes/spokes.go index c3ccc55..1140520 100644 --- a/internal/spokes/spokes.go +++ b/internal/spokes/spokes.go @@ -15,6 +15,7 @@ import ( "path/filepath" "regexp" "strings" + "sync" "syscall" "time" @@ -30,9 +31,10 @@ import ( const ( // maximum length of a pkt-line's data component - maxPacketDataLength = 65516 - nullSHA1OID = objectformat.NullOIDSHA1 - nullSHA256OID = objectformat.NullOIDSHA256 + maxPacketDataLength = 65516 + advertisementBufferSize = 64 * 1024 + nullSHA1OID = objectformat.NullOIDSHA1 + nullSHA256OID = objectformat.NullOIDSHA256 ) // Exec is similar to a main func for the new version of receive-pack. @@ -148,14 +150,8 @@ func (r *spokesReceivePack) execute(ctx context.Context) error { // We only need to perform the references discovery when we are not using the HTTP protocol or, if we are using it, // we only run the discovery phase when the http-backend-info-refs/advertise-refs option has been set if r.advertiseRefs || !r.statelessRPC { - if sockstat.GetBool("spokes_receive_pack_isolated_reference_discovery") { - if err := r.performReferenceDiscoveryIsolatedPipes(ctx); err != nil { - return err - } - } else { - if err := r.performReferenceDiscovery(ctx); err != nil { - return err - } + if err := r.performBufferedReferenceDiscovery(ctx); err != nil { + return err } } @@ -264,6 +260,27 @@ func (r *spokesReceivePack) execute(ctx context.Context) error { return nil } +func (r *spokesReceivePack) performBufferedReferenceDiscovery(ctx context.Context) error { + output := r.output + buffered := bufio.NewWriterSize(output, advertisementBufferSize) + r.output = buffered + defer func() { r.output = output }() + + var err error + if sockstat.GetBool("spokes_receive_pack_isolated_reference_discovery") { + err = r.performReferenceDiscoveryIsolatedPipes(ctx) + } else { + err = r.performReferenceDiscovery(ctx) + } + if err != nil { + return err + } + if err := buffered.Flush(); err != nil { + return fmt.Errorf("flushing reference advertisement: %w", err) + } + return nil +} + func supportedCapabilities(of objectformat.ObjectFormat) string { return fmt.Sprintf( "report-status report-status-v2 delete-refs side-band-64k ofs-delta atomic object-format=%s quiet", @@ -485,8 +502,14 @@ func (r *spokesReceivePack) performReferenceDiscovery(ctx context.Context) error } } + // The legacy pipeline starts each reference collector concurrently. + // Serialize complete pkt-lines and the one-time capabilities decision. + var advertiseMu sync.Mutex var wroteCapabilities bool advertiseRef := func(line []byte) error { + advertiseMu.Lock() + defer advertiseMu.Unlock() + if len(line) < 41 { return fmt.Errorf("malformed ref line: %q", string(line)) } diff --git a/internal/spokes/spokes_test.go b/internal/spokes/spokes_test.go index ea1e161..1639ebc 100644 --- a/internal/spokes/spokes_test.go +++ b/internal/spokes/spokes_test.go @@ -4,7 +4,9 @@ import ( "bytes" "context" "fmt" + "io" "os" + "strings" "testing" "github.com/github/spokes-receive-pack/internal/config" @@ -12,6 +14,16 @@ import ( "github.com/stretchr/testify/require" ) +type countingWriter struct { + w io.Writer + writes int +} + +func (w *countingWriter) Write(p []byte) (int, error) { + w.writes++ + return w.w.Write(p) +} + func TestCheckHiddenRefs(t *testing.T) { hiddenRefs := []string{"refs/pull/", "refs/gh/", "refs/__gh__", "!refs/__gh__/svn"} for _, p := range []struct { @@ -259,15 +271,53 @@ func TestPerformReferenceDiscovery(t *testing.T) { require.NoError(t, os.Chdir("testdata/lots-of-refs.git")) t.Cleanup(func() { _ = os.Chdir(origwd) }) + for _, isolated := range []bool{false, true} { + t.Run(fmt.Sprintf("isolated=%t", isolated), func(t *testing.T) { + isolatedValue := "" + if isolated { + isolatedValue = "bool:true" + } + t.Setenv("GIT_SOCKSTAT_VAR_spokes_receive_pack_isolated_reference_discovery", isolatedValue) + + var buf bytes.Buffer + output := &countingWriter{w: &buf} + wd, _ := os.Getwd() + r := &spokesReceivePack{ + config: &config.Config{}, + output: output, + repoPath: wd, + capabilities: "anything", + } + + assert.NoError(t, r.performBufferedReferenceDiscovery(context.Background())) + assert.Equal(t, expectedReferenceList, buf.String()) + assert.Equal(t, 1, output.writes) + }) + } +} + +func TestPerformBufferedReferenceDiscoveryConcurrentCollectors(t *testing.T) { + origwd, err := os.Getwd() + require.NoError(t, err) + require.NoError(t, os.Chdir("testdata/lots-of-refs.git")) + t.Cleanup(func() { _ = os.Chdir(origwd) }) + t.Setenv("GIT_SOCKSTAT_VAR_spokes_receive_pack_isolated_reference_discovery", "") + var buf bytes.Buffer + output := &countingWriter{w: &buf} wd, _ := os.Getwd() r := &spokesReceivePack{ - config: &config.Config{}, - output: &buf, + config: &config.Config{Entries: []config.ConfigEntry{ + {Key: "transfer.hiderefs", Value: "refs/tags/"}, + {Key: "transfer.hiderefs", Value: "!refs/tags/"}, + }}, + output: output, repoPath: wd, capabilities: "anything", } - assert.NoError(t, r.performReferenceDiscovery(context.Background())) - assert.Equal(t, expectedReferenceList, buf.String()) + assert.NoError(t, r.performBufferedReferenceDiscovery(context.Background())) + assert.Equal(t, len(expectedReferenceList), buf.Len()) + assert.True(t, strings.HasSuffix(buf.String(), "0000")) + assert.Equal(t, 1, output.writes) } diff --git a/spokes-receive-pack.go b/spokes-receive-pack.go index be7d3d3..7a15587 100644 --- a/spokes-receive-pack.go +++ b/spokes-receive-pack.go @@ -5,6 +5,7 @@ import ( "fmt" "io" "os" + "runtime" "github.com/github/spokes-receive-pack/internal/receivepack" "github.com/github/spokes-receive-pack/internal/sockstat" @@ -14,6 +15,7 @@ import ( var BuildVersion string func main() { + setGOMAXPROCS() exitCode, err := mainImpl(os.Stdin, os.Stdout, os.Stderr, os.Args[1:]) if err != nil { fmt.Fprintf(os.Stderr, "error: %v\n", err) @@ -21,6 +23,11 @@ func main() { os.Exit(exitCode) } +func setGOMAXPROCS() { + // This short-lived command is serial/IO-bound; inheriting host-wide GOMAXPROCS causes large GC/runtime thread fanout. + runtime.GOMAXPROCS(1) +} + func mainImpl(stdin io.Reader, stdout, stderr io.Writer, args []string) (int, error) { ctx := context.Background() if !sockstat.GetBool("spokes_quarantine") { diff --git a/spokes-receive-pack_test.go b/spokes-receive-pack_test.go new file mode 100644 index 0000000..712f20b --- /dev/null +++ b/spokes-receive-pack_test.go @@ -0,0 +1,19 @@ +package main + +import ( + "runtime" + "testing" +) + +func TestSetGOMAXPROCS(t *testing.T) { + previous := runtime.GOMAXPROCS(0) + t.Cleanup(func() { + runtime.GOMAXPROCS(previous) + }) + + setGOMAXPROCS() + + if got := runtime.GOMAXPROCS(0); got != 1 { + t.Fatalf("GOMAXPROCS() = %d, want 1", got) + } +}