package updater_test import ( "context" "errors" "sync" "testing" "time" "github.com/shcizo/package-updater/internal/discovery" "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) 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) 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) 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) }