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
3 changes: 1 addition & 2 deletions .github/workflows/go-test.yml
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand All @@ -28,4 +28,3 @@ jobs:
run: make test
- name: go integration test
run: make test-integration

45 changes: 34 additions & 11 deletions internal/spokes/spokes.go
Original file line number Diff line number Diff line change
Expand Up @@ -15,6 +15,7 @@ import (
"path/filepath"
"regexp"
"strings"
"sync"
"syscall"
"time"

Expand All @@ -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.
Expand Down Expand Up @@ -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
}
}

Expand Down Expand Up @@ -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",
Expand Down Expand Up @@ -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))
}
Expand Down
58 changes: 54 additions & 4 deletions internal/spokes/spokes_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -4,14 +4,26 @@ import (
"bytes"
"context"
"fmt"
"io"
"os"
"strings"
"testing"

"github.com/github/spokes-receive-pack/internal/config"
"github.com/stretchr/testify/assert"
"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 {
Expand Down Expand Up @@ -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)
}
7 changes: 7 additions & 0 deletions spokes-receive-pack.go
Original file line number Diff line number Diff line change
Expand Up @@ -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"
Expand All @@ -14,13 +15,19 @@ 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)
}
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") {
Expand Down
19 changes: 19 additions & 0 deletions spokes-receive-pack_test.go
Original file line number Diff line number Diff line change
@@ -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)
}
}
Loading