Something went wrong. Try again.
Monorepo for Tangled tangled.org
Something went wrong. Try again.
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267package engine
import ( "context" "errors" "fmt" "testing" "time")
// resources for testingtype ru struct{ a, b int64 }
func (r ru) Fits(limit ru) bool { if limit.a > 0 && r.a > limit.a { return false } if limit.b > 0 && r.b > limit.b { return false } return true}func (r ru) Add(o ru) ru { return ru{r.a + o.a, r.b + o.b} }func (r ru) Sub(o ru) ru { return ru{max(0, r.a-o.a), max(0, r.b-o.b)} }func (r ru) String() string { return fmt.Sprintf("a=%d b=%d", r.a, r.b)}
type acquireResult struct { slot WorkflowSlot err error}
func TestResourceSchedulerZeroLimitsDoNotApply(t *testing.T) { t.Parallel()
scheduler := NewResourceScheduler(ru{}, ru{}, 0)
slot, err := scheduler.Acquire(context.Background(), ru{a: 1 << 20, b: 1 << 20}, Wait) if err != nil { t.Fatalf("Acquire() error = %v", err) } slot.Release()}
func TestResourceSchedulerRejectsRequestsThatCanNeverFit(t *testing.T) { t.Parallel()
scheduler := NewResourceScheduler(ru{a: 1024, b: 10_000}, ru{a: 512, b: 5_000}, 0)
_, err := scheduler.Acquire(context.Background(), ru{a: 768, b: 100}, Wait) if !errors.Is(err, ErrNoWorkflowSlots) { t.Fatalf("Acquire() error = %v, want ErrNoWorkflowSlots", err) }
_, err = scheduler.Acquire(context.Background(), ru{a: 128, b: 12_000}, Wait) if !errors.Is(err, ErrNoWorkflowSlots) { t.Fatalf("Acquire() error = %v, want ErrNoWorkflowSlots", err) }}
func TestResourceSchedulerWaitsUntilResourcesAreReleased(t *testing.T) { t.Parallel()
scheduler := NewResourceScheduler(ru{a: 1024}, ru{}, 0)
first, err := scheduler.Acquire(context.Background(), ru{a: 1024}, Wait) if err != nil { t.Fatalf("first Acquire() error = %v", err) } defer first.Release()
ch := acquireAsync(context.Background(), scheduler, ru{a: 1}) assertAcquireBlocked(t, ch)
first.Release() first = NoopSlot{}
second := waitAcquireOK(t, ch) second.Release()}
func TestResourceSchedulerReleaseIsIdempotent(t *testing.T) { t.Parallel()
scheduler := NewResourceScheduler(ru{a: 1}, ru{}, 0)
slot, err := scheduler.Acquire(context.Background(), ru{a: 1}, Wait) if err != nil { t.Fatalf("Acquire() error = %v", err) }
slot.Release() slot.Release()
second, err := scheduler.Acquire(context.Background(), ru{a: 1}, Wait) if err != nil { t.Fatalf("Acquire() after double release error = %v", err) } second.Release()}
func TestResourceSchedulerBackfillsPastBlockedHead(t *testing.T) { t.Parallel()
scheduler := NewResourceScheduler(ru{a: 1024}, ru{}, time.Hour) // disable aging so we test pure backfill
hold, err := scheduler.Acquire(context.Background(), ru{a: 512}, Wait) if err != nil { t.Fatalf("hold Acquire() error = %v", err) } defer hold.Release()
bigCh := acquireAsync(context.Background(), scheduler, ru{a: 768}) assertAcquireBlocked(t, bigCh)
smallCh := acquireAsync(context.Background(), scheduler, ru{a: 256}) small := waitAcquireOK(t, smallCh) small.Release()
assertAcquireBlocked(t, bigCh)}
func TestResourceSchedulerAgingReservesCapacityForBlockedHead(t *testing.T) { t.Parallel()
scheduler := NewResourceScheduler(ru{a: 1024}, ru{}, 10*time.Millisecond) fakeNow := time.Now() scheduler.now = func() time.Time { return fakeNow }
hold, err := scheduler.Acquire(context.Background(), ru{a: 512}, Wait) if err != nil { t.Fatalf("hold Acquire() error = %v", err) }
bigCh := acquireAsync(context.Background(), scheduler, ru{a: 768}) assertAcquireBlocked(t, bigCh)
fakeNow = fakeNow.Add(time.Second)
// big is now aged and reserves its 768. a 256 request would fit // alongside the held 512, but the reservation blocks it. smallCh := acquireAsync(context.Background(), scheduler, ru{a: 256}) assertAcquireBlocked(t, smallCh)
hold.Release()
big := waitAcquireOK(t, bigCh) small := waitAcquireOK(t, smallCh) small.Release() big.Release()}
func TestResourceSchedulerTryRejectsWhenNoRoomNow(t *testing.T) { t.Parallel()
scheduler := NewResourceScheduler(ru{a: 1024}, ru{}, 0)
first, err := scheduler.Acquire(context.Background(), ru{a: 1024}, NoWait) if err != nil { t.Fatalf("first Acquire(NoWait) error = %v", err) }
if _, err := scheduler.Acquire(context.Background(), ru{a: 1}, NoWait); !errors.Is(err, ErrNoWorkflowSlots) { t.Fatalf("Acquire(NoWait) error = %v, want ErrNoWorkflowSlots", err) }
first.Release()
second, err := scheduler.Acquire(context.Background(), ru{a: 1}, NoWait) if err != nil { t.Fatalf("Acquire(NoWait) after release error = %v", err) } second.Release()}
func TestResourceSchedulerTryRejectsRequestsThatCanNeverFit(t *testing.T) { t.Parallel()
scheduler := NewResourceScheduler(ru{a: 1024, b: 10_000}, ru{a: 512, b: 5_000}, 0)
if _, err := scheduler.Acquire(context.Background(), ru{a: 768, b: 100}, NoWait); !errors.Is(err, ErrNoWorkflowSlots) { t.Fatalf("Acquire(NoWait) error = %v, want ErrNoWorkflowSlots", err) }}
func TestResourceSchedulerTryIgnoresQueuedWaiters(t *testing.T) { t.Parallel()
scheduler := NewResourceScheduler(ru{a: 1024}, ru{}, time.Hour)
hold, err := scheduler.Acquire(context.Background(), ru{a: 512}, Wait) if err != nil { t.Fatalf("hold Acquire() error = %v", err) } defer hold.Release()
bigCh := acquireAsync(context.Background(), scheduler, ru{a: 768}) assertAcquireBlocked(t, bigCh)
slot, err := scheduler.Acquire(context.Background(), ru{a: 256}, NoWait) if err != nil { t.Fatalf("Acquire(NoWait) error = %v, want success past queued waiter", err) } slot.Release()}
func TestResourceSchedulerZeroLimitsTryDoesNotReject(t *testing.T) { t.Parallel()
scheduler := NewResourceScheduler(ru{}, ru{}, 0)
slot, err := scheduler.Acquire(context.Background(), ru{a: 1 << 20, b: 1 << 20}, NoWait) if err != nil { t.Fatalf("Acquire(NoWait) error = %v", err) } slot.Release()}
func acquireAsync(ctx context.Context, scheduler *ResourceScheduler[ru], req ru) <-chan acquireResult { ch := make(chan acquireResult, 1) go func() { slot, err := scheduler.Acquire(ctx, req, Wait) ch <- acquireResult{slot: slot, err: err} }() return ch}
func assertAcquireBlocked(t *testing.T, ch <-chan acquireResult) { t.Helper()
select { case res := <-ch: if res.slot != nil { res.slot.Release() } t.Fatalf("Acquire() returned before resources were available: err=%v", res.err) case <-time.After(25 * time.Millisecond): }}
func waitAcquireOK(t *testing.T, ch <-chan acquireResult) WorkflowSlot { t.Helper()
res := waitAcquireResult(t, ch) if res.err != nil { t.Fatalf("Acquire() error = %v", res.err) } if res.slot == nil { t.Fatal("Acquire() returned nil slot") } return res.slot}
func waitAcquireResult(t *testing.T, ch <-chan acquireResult) acquireResult { t.Helper()
select { case res := <-ch: return res case <-time.After(time.Second): t.Fatal("timed out waiting for Acquire() result") }
return acquireResult{}}