Skip to content
Merged
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
4 changes: 1 addition & 3 deletions .golangci.yaml
Original file line number Diff line number Diff line change
Expand Up @@ -4,7 +4,7 @@ version: "2"

run:
go: "1.26"
build-tags: [ safe ]
build-tags: [ safe, pprof ]
modules-download-mode: readonly

linters:
Expand All @@ -16,8 +16,6 @@ linters:
- paralleltest
- godoclint
- forbidigo
- wsl_v5
- whitespace
# Discouraged linters
- noinlineerr # Disallows inline error handling (`if err := ...; err != nil {`).
- embeddedstructfieldcheck # Embedded types should be at the top of the field list of a struct, and there must be an empty line separating embedded fields from regular fields. [fast]
Expand Down
3 changes: 1 addition & 2 deletions .goreleaser.yaml
Original file line number Diff line number Diff line change
Expand Up @@ -25,8 +25,7 @@ builds:
- '7'
flags:
- -trimpath
- -tags=safe
- -tags=pprof
- -tags=safe,pprof
ldflags:
- -s -w -X github.com/foomo/contentserver/cmd.version={{.Version}}

Expand Down
9 changes: 9 additions & 0 deletions client/client.go
Original file line number Diff line number Diff line change
Expand Up @@ -39,10 +39,12 @@ func (c *Client) Update(ctx context.Context) (*responses.Update, error) {
type serverResponse struct {
Reply *responses.Update
}

resp := serverResponse{}
if err := c.t.Call(ctx, handler.RouteUpdate, &requests.Update{}, &resp); err != nil {
return nil, err
}

return resp.Reply, nil
}

Expand All @@ -51,6 +53,7 @@ func (c *Client) GetContent(ctx context.Context, request *requests.Content) (*co
type serverResponse struct {
Reply *content.SiteContent
}

resp := serverResponse{}
if err := c.t.Call(ctx, handler.RouteGetContent, request, &resp); err != nil {
return nil, err
Expand All @@ -69,6 +72,7 @@ func (c *Client) GetURIs(ctx context.Context, dimension string, ids []string) (m
if err := c.t.Call(ctx, handler.RouteGetURIs, &requests.URIs{Dimension: dimension, IDs: ids}, &resp); err != nil {
return nil, err
}

return resp.Reply, nil
}

Expand All @@ -78,13 +82,16 @@ func (c *Client) GetNodes(ctx context.Context, env *requests.Env, nodes map[stri
Env: env,
Nodes: nodes,
}

type serverResponse struct {
Reply map[string]*content.Node
}

resp := serverResponse{}
if err := c.t.Call(ctx, handler.RouteGetNodes, r, &resp); err != nil {
return nil, err
}

return resp.Reply, nil
}

Expand All @@ -93,10 +100,12 @@ func (c *Client) GetRepo(ctx context.Context) (map[string]*content.RepoNode, err
type serverResponse struct {
Reply map[string]*content.RepoNode
}

resp := serverResponse{}
if err := c.t.Call(ctx, handler.RouteGetRepo, &requests.Repo{}, &resp); err != nil {
return nil, err
}

return resp.Reply, nil
}

Expand Down
19 changes: 19 additions & 0 deletions client/client_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -32,8 +32,10 @@ func TestGetURIs(t *testing.T) {
testWithClients(t, func(t *testing.T, c *client.Client) {
t.Helper()
t.Parallel()

request := mock.MakeValidURIsRequest()
uriMap, err := c.GetURIs(t.Context(), request.Dimension, request.IDs)

time.Sleep(100 * time.Millisecond)
require.NoError(t, err)
assert.Equal(t, "/a", uriMap[request.IDs[0]])
Expand All @@ -46,6 +48,7 @@ func TestGetRepo(t *testing.T) {
t.Parallel()
r, err := c.GetRepo(t.Context())
require.NoError(t, err)

if assert.NotEmpty(t, r, "received empty JSON from GetRepo") {
assert.InDelta(t, 1.0, r["dimension_foo"].Nodes["id-a"].Data["baz"].(float64), 0, "failed to drill deep for data") //nolint:forcetypeassert
}
Expand All @@ -56,17 +59,21 @@ func TestGetNodes(t *testing.T) {
testWithClients(t, func(t *testing.T, c *client.Client) {
t.Helper()
t.Parallel()

nodesRequest := mock.MakeNodesRequest()
nodes, err := c.GetNodes(t.Context(), nodesRequest.Env, nodesRequest.Nodes)
require.NoError(t, err)

testNode, ok := nodes["test"]
if !ok {
t.Fatal("that should be a node")
}

testData, ok := testNode.Item.Data["foo"]
if !ok {
t.Fatal("where is foo")
}

if testData != "bar" {
t.Fatal("testData should have bennd bar not", testData)
}
Expand All @@ -77,6 +84,7 @@ func TestGetContent(t *testing.T) {
testWithClients(t, func(t *testing.T, c *client.Client) {
t.Helper()
t.Parallel()

request := mock.MakeValidContentRequest()
response, err := c.GetContent(t.Context(), request)
require.NoError(t, err)
Expand All @@ -88,9 +96,12 @@ func TestGetContent(t *testing.T) {
func benchmarkServerAndClientGetContent(b *testing.B, numGroups, numCalls int, client GetContentClient) {
b.Helper()
b.ResetTimer()

for i := 0; i < b.N; i++ {
start := time.Now()

benchmarkClientAndServerGetContent(b, numGroups, numCalls, client)

dur := time.Since(start)
totalCalls := numGroups * numCalls
b.Log("requests per second", int(float64(totalCalls)/(float64(dur)/float64(1000000000))), dur, totalCalls)
Expand All @@ -99,18 +110,22 @@ func benchmarkServerAndClientGetContent(b *testing.B, numGroups, numCalls int, c

func benchmarkClientAndServerGetContent(tb testing.TB, numGroups, numCalls int, client GetContentClient) {
tb.Helper()

var wg sync.WaitGroup
wg.Add(numGroups)

for range numGroups {
go func() {
defer wg.Done()

request := mock.MakeValidContentRequest()
for range numCalls {
response, err := client.GetContent(tb.Context(), request)
if err == nil {
if request.URI != response.URI {
tb.Fatal("uri mismatch")
}

if response.Status != content.StatusOk {
tb.Fatal("unexpected status")
}
Expand Down Expand Up @@ -149,17 +164,20 @@ func testWithClients(t *testing.T, testFunc func(t *testing.T, c *client.Client)
func initRepo(tb testing.TB, l *zap.Logger) *repo.Repo {
tb.Helper()
testRepoServer, varDir := mock.GetMockData(tb)

h, err := repo.NewHistory(l,
repo.HistoryWithHistoryDir(varDir),
)
if err != nil {
tb.Fatal(err)
}

r := repo.New(l,
testRepoServer.URL+"/repo-two-dimensions.json",
h,
)
up := make(chan bool, 1)

r.OnLoaded(func() {
up <- true
})
Expand All @@ -168,6 +186,7 @@ func initRepo(tb testing.TB, l *zap.Logger) *repo.Repo {
// preventing race conditions with logging after test completion.
ctx, cancel := context.WithCancel(context.Background())
go r.Start(ctx) //nolint:errcheck

<-up

tb.Cleanup(func() {
Expand Down
13 changes: 13 additions & 0 deletions client/connectionpool.go
Original file line number Diff line number Diff line change
Expand Up @@ -22,6 +22,7 @@ func newConnectionPool(url string, connectionPoolSize int, waitTimeout time.Dura
chanDrainPool: make(chan int),
}
go connPool.run(connectionPoolSize, waitTimeout)

return connPool
}

Expand All @@ -31,6 +32,7 @@ func (c *connectionPool) run(connectionPoolSize int, waitTimeout time.Duration)
err error
conn net.Conn
}

type waitPoolEntry struct {
entryTime time.Time
chanConn chan net.Conn
Expand All @@ -46,6 +48,7 @@ func (c *connectionPool) run(connectionPoolSize int, waitTimeout time.Duration)
busy: false,
}
}

RunLoop:
for {
// fmt.Println("----------------------- run loop ------------------------")
Expand All @@ -55,6 +58,7 @@ RunLoop:
for _, waitPoolEntry := range waitPool {
waitPoolEntry.chanConn <- nil
}

break RunLoop
case <-time.After(waitTimeout):
// fmt.Println("tick", len(connectionPool), len(waitPool))
Expand All @@ -72,6 +76,7 @@ RunLoop:
nextI = i + 1
}
}

waitPool[nextI] = &waitPoolEntry{
chanConn: chanReturnNextConn,
entryTime: time.Now(),
Expand All @@ -94,6 +99,7 @@ RunLoop:
for _, poolEntry := range connectionPool {
if poolEntry.conn == nil {
var d net.Dialer

newConn, errDial := d.DialContext(context.Background(), "tcp", c.url)
poolEntry.err = errDial
poolEntry.conn = newConn
Expand All @@ -104,12 +110,16 @@ RunLoop:
if len(waitPool) == 0 {
break
}

if poolEntry.err == nil && poolEntry.conn != nil && !poolEntry.busy {
for i, waitPoolEntry := range waitPool {
// fmt.Println("---------------------------> serving wait pool", i, waitPoolEntry)
poolEntry.busy = true

delete(waitPool, i)

waitPoolEntry.chanConn <- poolEntry.conn

break
}
}
Expand All @@ -122,13 +132,16 @@ RunLoop:
for i, waitPoolEntry := range waitPool {
if now.Sub(waitPoolEntry.entryTime) > waitTimeout {
waitPoolLoosers = append(waitPoolLoosers, i)

waitPoolEntry.chanConn <- nil
}
}

for _, i := range waitPoolLoosers {
delete(waitPool, i)
}
}

c.chanDrainPool = nil
c.chanConnReturn = nil
c.chanConnGet = nil
Expand Down
5 changes: 5 additions & 0 deletions client/httptransport.go
Original file line number Diff line number Diff line change
Expand Up @@ -72,6 +72,7 @@ func (t *HTTPTransport) Call(ctx context.Context, route handler.Route, request a
if errMarshal != nil {
return errMarshal
}

req, errNewRequest := http.NewRequestWithContext(
ctx,
http.MethodPost,
Expand All @@ -81,6 +82,7 @@ func (t *HTTPTransport) Call(ctx context.Context, route handler.Route, request a
if errNewRequest != nil {
return errNewRequest
}

httpResponse, errDo := t.httpClient.Do(req) // #nosec G704 -- The client transport must call the caller-configured contentserver endpoint.
if errDo != nil {
return errDo
Expand All @@ -90,13 +92,16 @@ func (t *HTTPTransport) Call(ctx context.Context, route handler.Route, request a
if httpResponse.StatusCode != http.StatusOK {
return errors.New("non 200 reply")
}

if httpResponse.Body == nil {
return errors.New("empty response body")
}

responseBytes, errRead := io.ReadAll(httpResponse.Body)
if errRead != nil {
return errRead
}

return json.Unmarshal(responseBytes, response)
}

Expand Down
3 changes: 3 additions & 0 deletions client/httptransport_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -52,13 +52,16 @@ type GetContentClient interface {

func newHTTPClient(tb testing.TB, server *httptest.Server) *client.Client {
tb.Helper()

c, err := client.NewHTTPClient(server.URL + pathContentserver)
require.NoError(tb, err)

return c
}

func initHTTPRepoServer(tb testing.TB, l *zap.Logger) *httptest.Server {
tb.Helper()
r := initRepo(tb, l)

return httptest.NewServer(handler.NewHTTP(l, r))
}
Loading