All checks were successful
release-tag / release-image (push) Successful in 2m32s
97 lines
2.4 KiB
Go
97 lines
2.4 KiB
Go
package workqueue
|
|
|
|
import (
|
|
"context"
|
|
"errors"
|
|
"testing"
|
|
"time"
|
|
)
|
|
|
|
func TestLimiterBoundsQueue(t *testing.T) {
|
|
l := New(1, 1)
|
|
release, err := l.Acquire(context.Background())
|
|
if err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
defer release()
|
|
|
|
ctx, cancel := context.WithCancel(context.Background())
|
|
defer cancel()
|
|
waiterDone := make(chan error, 1)
|
|
go func() {
|
|
r, err := l.Acquire(ctx)
|
|
if err == nil {
|
|
r()
|
|
}
|
|
waiterDone <- err
|
|
}()
|
|
|
|
deadline := time.Now().Add(time.Second)
|
|
for l.Status().Waiting != 1 && time.Now().Before(deadline) {
|
|
time.Sleep(time.Millisecond)
|
|
}
|
|
if l.Status().Waiting != 1 {
|
|
t.Fatalf("expected one waiter: %+v", l.Status())
|
|
}
|
|
if _, err := l.Acquire(context.Background()); !errors.Is(err, ErrQueueFull) {
|
|
t.Fatalf("expected ErrQueueFull, got %v", err)
|
|
}
|
|
release()
|
|
if err := <-waiterDone; err != nil {
|
|
t.Fatalf("waiter failed: %v", err)
|
|
}
|
|
}
|
|
|
|
func TestLimiterTracksKindsWithoutDuplicateCounters(t *testing.T) {
|
|
l := New(2, 2)
|
|
releaseSearch, err := l.AcquireKind(context.Background(), "searxng.search")
|
|
if err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
releaseChat, err := l.AcquireKind(context.Background(), "ollama.chat")
|
|
if err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
status := l.Status()
|
|
if status.Active != 2 || status.Kinds["searxng.search"].Active != 1 || status.Kinds["ollama.chat"].Active != 1 {
|
|
t.Fatalf("unexpected typed queue status: %+v", status)
|
|
}
|
|
releaseSearch()
|
|
releaseSearch() // release must be idempotent
|
|
releaseChat()
|
|
status = l.Status()
|
|
if status.Active != 0 || status.Kinds["searxng.search"].Active != 0 || status.Kinds["ollama.chat"].Active != 0 {
|
|
t.Fatalf("release duplicated or leaked active counters: %+v", status)
|
|
}
|
|
}
|
|
|
|
func TestLimiterCanRaiseRuntimeLimit(t *testing.T) {
|
|
l := New(1, 4)
|
|
release1, err := l.Acquire(context.Background())
|
|
if err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
acquired := make(chan func(), 1)
|
|
go func() {
|
|
release, err := l.Acquire(context.Background())
|
|
if err == nil {
|
|
acquired <- release
|
|
}
|
|
}()
|
|
time.Sleep(10 * time.Millisecond)
|
|
if got := l.Status(); got.Active != 1 || got.Waiting != 1 {
|
|
t.Fatalf("unexpected before resize: %+v", got)
|
|
}
|
|
l.SetLimits(2, 8)
|
|
select {
|
|
case release2 := <-acquired:
|
|
release2()
|
|
case <-time.After(time.Second):
|
|
t.Fatal("raising limiter did not wake waiter")
|
|
}
|
|
release1()
|
|
if got := l.Status(); got.MaxInflight != 2 || got.QueueSize != 8 || got.Active != 0 {
|
|
t.Fatalf("unexpected resized status: %+v", got)
|
|
}
|
|
}
|