130 lines
3.7 KiB
Go
130 lines
3.7 KiB
Go
package updater_test
|
|
|
|
import (
|
|
"context"
|
|
"errors"
|
|
"sync"
|
|
"testing"
|
|
"time"
|
|
|
|
dto "github.com/prometheus/client_model/go"
|
|
"github.com/prometheus/client_golang/prometheus"
|
|
"github.com/shcizo/package-updater/internal/discovery"
|
|
"github.com/shcizo/package-updater/internal/metrics"
|
|
"github.com/shcizo/package-updater/internal/updater"
|
|
"github.com/stretchr/testify/require"
|
|
)
|
|
|
|
type fakeExec struct {
|
|
mu sync.Mutex
|
|
calls []discovery.Job
|
|
delay time.Duration
|
|
errForSvc map[string]error
|
|
}
|
|
|
|
func (f *fakeExec) Execute(ctx context.Context, job discovery.Job) error {
|
|
f.mu.Lock()
|
|
f.calls = append(f.calls, job)
|
|
f.mu.Unlock()
|
|
if f.delay > 0 {
|
|
select {
|
|
case <-time.After(f.delay):
|
|
case <-ctx.Done():
|
|
return ctx.Err()
|
|
}
|
|
}
|
|
if err, ok := f.errForSvc[job.Service]; ok {
|
|
return err
|
|
}
|
|
return nil
|
|
}
|
|
|
|
func (f *fakeExec) callsCopy() []discovery.Job {
|
|
f.mu.Lock()
|
|
defer f.mu.Unlock()
|
|
return append([]discovery.Job(nil), f.calls...)
|
|
}
|
|
|
|
func TestQueue_ProcessesFIFO(t *testing.T) {
|
|
exec := &fakeExec{delay: 20 * time.Millisecond}
|
|
q := updater.NewQueue(exec, nil)
|
|
q.Start(context.Background())
|
|
defer q.Stop()
|
|
|
|
job1 := discovery.Job{Project: "a", Service: "svc"}
|
|
job2 := discovery.Job{Project: "b", Service: "svc"}
|
|
job3 := discovery.Job{Project: "c", Service: "svc"}
|
|
|
|
results := make(chan []updater.Result, 3)
|
|
go func() { results <- q.Submit(context.Background(), []discovery.Job{job1}) }()
|
|
time.Sleep(5 * time.Millisecond)
|
|
go func() { results <- q.Submit(context.Background(), []discovery.Job{job2}) }()
|
|
time.Sleep(5 * time.Millisecond)
|
|
go func() { results <- q.Submit(context.Background(), []discovery.Job{job3}) }()
|
|
|
|
for i := 0; i < 3; i++ {
|
|
<-results
|
|
}
|
|
|
|
calls := exec.callsCopy()
|
|
require.Equal(t, []string{"a", "b", "c"}, []string{calls[0].Project, calls[1].Project, calls[2].Project})
|
|
}
|
|
|
|
func TestQueue_ReturnsPerJobResults(t *testing.T) {
|
|
exec := &fakeExec{errForSvc: map[string]error{"failing": errors.New("boom")}}
|
|
q := updater.NewQueue(exec, nil)
|
|
q.Start(context.Background())
|
|
defer q.Stop()
|
|
|
|
jobs := []discovery.Job{
|
|
{Project: "p1", Service: "ok"},
|
|
{Project: "p2", Service: "failing"},
|
|
{Project: "p3", Service: "refused-svc", Refused: true, RefusedReason: "outside root"},
|
|
}
|
|
results := q.Submit(context.Background(), jobs)
|
|
|
|
require.Len(t, results, 3)
|
|
require.Equal(t, updater.StatusUpdated, results[0].Status)
|
|
require.Equal(t, updater.StatusFailed, results[1].Status)
|
|
require.Contains(t, results[1].Error, "boom")
|
|
require.Equal(t, updater.StatusRefused, results[2].Status)
|
|
require.Contains(t, results[2].Error, "outside root")
|
|
}
|
|
|
|
func TestQueue_Timeout(t *testing.T) {
|
|
exec := &fakeExec{delay: 200 * time.Millisecond}
|
|
q := updater.NewQueue(exec, nil)
|
|
q.Start(context.Background())
|
|
defer q.Stop()
|
|
|
|
ctx, cancel := context.WithTimeout(context.Background(), 30*time.Millisecond)
|
|
defer cancel()
|
|
|
|
jobs := []discovery.Job{{Project: "slow", Service: "svc"}}
|
|
results := q.Submit(ctx, jobs)
|
|
|
|
require.Len(t, results, 1)
|
|
require.Equal(t, updater.StatusTimeout, results[0].Status)
|
|
}
|
|
|
|
func TestQueue_RecordsMetrics(t *testing.T) {
|
|
reg := prometheus.NewRegistry()
|
|
m := metrics.New(reg)
|
|
exec := &fakeExec{}
|
|
q := updater.NewQueue(exec, m)
|
|
q.Start(context.Background())
|
|
defer q.Stop()
|
|
|
|
q.Submit(context.Background(), []discovery.Job{{Project: "p", Service: "s"}})
|
|
|
|
// Verify UpdateJobs counter was incremented for the successful job.
|
|
metric := &dto.Metric{}
|
|
require.NoError(t, m.UpdateJobs.WithLabelValues("p", "s", "updated").Write(metric))
|
|
require.Equal(t, 1.0, metric.GetCounter().GetValue())
|
|
|
|
// Verify LastUpdateTime was set (non-zero).
|
|
tsMetric := &dto.Metric{}
|
|
require.NoError(t, m.LastUpdateTime.WithLabelValues("p", "s").Write(tsMetric))
|
|
require.Greater(t, tsMetric.GetGauge().GetValue(), 0.0)
|
|
}
|