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) } }