151 lines
3.5 KiB
Go
151 lines
3.5 KiB
Go
// Copyright 2019 Drone.IO Inc. All rights reserved.
|
|
// Use of this source code is governed by the Drone Non-Commercial License
|
|
// that can be found in the LICENSE file.
|
|
|
|
package queue
|
|
|
|
import (
|
|
"context"
|
|
"sync"
|
|
"testing"
|
|
"time"
|
|
|
|
"github.com/drone/drone/core"
|
|
"github.com/drone/drone/mock"
|
|
|
|
"github.com/golang/mock/gomock"
|
|
)
|
|
|
|
func TestQueue(t *testing.T) {
|
|
controller := gomock.NewController(t)
|
|
defer controller.Finish()
|
|
|
|
items := []*core.Stage{
|
|
{ID: 3, OS: "linux", Arch: "amd64"},
|
|
{ID: 2, OS: "linux", Arch: "amd64"},
|
|
{ID: 1, OS: "linux", Arch: "amd64"},
|
|
}
|
|
|
|
ctx := context.Background()
|
|
store := mock.NewMockStageStore(controller)
|
|
store.EXPECT().ListIncomplete(ctx).Return(items, nil).Times(1)
|
|
store.EXPECT().ListIncomplete(ctx).Return(items[1:], nil).Times(1)
|
|
store.EXPECT().ListIncomplete(ctx).Return(items[2:], nil).Times(1)
|
|
|
|
q := newQueue(store)
|
|
for _, item := range items {
|
|
next, err := q.Request(ctx, core.Filter{OS: "linux", Arch: "amd64"})
|
|
if err != nil {
|
|
t.Error(err)
|
|
return
|
|
}
|
|
if got, want := next, item; got != want {
|
|
t.Errorf("Want build %d, got %d", item.ID, item.ID)
|
|
}
|
|
}
|
|
}
|
|
|
|
func TestQueueCancel(t *testing.T) {
|
|
controller := gomock.NewController(t)
|
|
defer controller.Finish()
|
|
|
|
ctx, cancel := context.WithCancel(context.Background())
|
|
store := mock.NewMockStageStore(controller)
|
|
store.EXPECT().ListIncomplete(ctx).Return(nil, nil)
|
|
|
|
q := newQueue(store)
|
|
q.ctx = ctx
|
|
|
|
var wg sync.WaitGroup
|
|
wg.Add(1)
|
|
|
|
go func() {
|
|
build, err := q.Request(ctx, core.Filter{OS: "linux/amd64", Arch: "amd64"})
|
|
if err != context.Canceled {
|
|
t.Errorf("Expected context.Canceled error, got %s", err)
|
|
}
|
|
if build != nil {
|
|
t.Errorf("Expect nil build when subscribe canceled")
|
|
}
|
|
wg.Done()
|
|
}()
|
|
<-time.After(10 * time.Millisecond)
|
|
|
|
q.Lock()
|
|
count := len(q.workers)
|
|
q.Unlock()
|
|
|
|
if got, want := count, 1; got != want {
|
|
t.Errorf("Want %d listener, got %d", want, got)
|
|
}
|
|
|
|
cancel()
|
|
wg.Wait()
|
|
}
|
|
|
|
func TestQueuePush(t *testing.T) {
|
|
controller := gomock.NewController(t)
|
|
defer controller.Finish()
|
|
|
|
item1 := &core.Stage{
|
|
ID: 1,
|
|
OS: "linux",
|
|
Arch: "amd64",
|
|
}
|
|
item2 := &core.Stage{
|
|
ID: 2,
|
|
OS: "linux",
|
|
Arch: "amd64",
|
|
}
|
|
|
|
ctx := context.Background()
|
|
store := mock.NewMockStageStore(controller)
|
|
|
|
q := &queue{
|
|
store: store,
|
|
ready: make(chan struct{}, 1),
|
|
}
|
|
q.Schedule(ctx, item1)
|
|
q.Schedule(ctx, item2)
|
|
select {
|
|
case <-q.ready:
|
|
case <-time.After(time.Millisecond):
|
|
t.Errorf("Expect queue signaled on push")
|
|
}
|
|
}
|
|
|
|
func TestWithinLimits(t *testing.T) {
|
|
tests := []struct {
|
|
ID int64
|
|
RepoID int64
|
|
Name string
|
|
Limit int
|
|
Want bool
|
|
}{
|
|
{Want: true, ID: 1, RepoID: 1, Name: "foo"},
|
|
{Want: true, ID: 2, RepoID: 2, Name: "bar", Limit: 1},
|
|
{Want: true, ID: 3, RepoID: 1, Name: "bar", Limit: 1},
|
|
{Want: false, ID: 4, RepoID: 1, Name: "bar", Limit: 1},
|
|
{Want: false, ID: 5, RepoID: 1, Name: "bar", Limit: 1},
|
|
{Want: true, ID: 6, RepoID: 1, Name: "baz", Limit: 2},
|
|
{Want: true, ID: 7, RepoID: 1, Name: "baz", Limit: 2},
|
|
{Want: false, ID: 8, RepoID: 1, Name: "baz", Limit: 2},
|
|
{Want: false, ID: 9, RepoID: 1, Name: "baz", Limit: 2},
|
|
{Want: true, ID: 10, RepoID: 1, Name: "baz", Limit: 0},
|
|
}
|
|
var stages []*core.Stage
|
|
for _, test := range tests {
|
|
stages = append(stages, &core.Stage{
|
|
ID: test.ID,
|
|
RepoID: test.RepoID,
|
|
Name: test.Name,
|
|
Limit: test.Limit,
|
|
})
|
|
}
|
|
for i, test := range tests {
|
|
stage := stages[i]
|
|
if got, want := withinLimits(stage, stages), test.Want; got != want {
|
|
t.Errorf("Unexpectd results at index %d", i)
|
|
}
|
|
}
|
|
}
|