This repository has no description
1package engine
2
3import (
4 "context"
5 "errors"
6 "testing"
7 "time"
8
9 "tangled.org/core/spindle/models"
10)
11
12func TestSemaphoreSlotterDisabledDoesNotBlock(t *testing.T) {
13 t.Parallel()
14
15 slotter := NewSemaphoreSlotter(0)
16
17 for range 10 {
18 slot, err := slotter.AcquireWorkflowSlot(context.Background(), zeroWorkflowID(), nil, Wait)
19 if err != nil {
20 t.Fatalf("AcquireWorkflowSlot() error = %v", err)
21 }
22 slot.Release()
23 }
24}
25
26func TestSemaphoreSlotterBlocksUntilRelease(t *testing.T) {
27 t.Parallel()
28
29 slotter := NewSemaphoreSlotter(1)
30
31 first, err := slotter.AcquireWorkflowSlot(context.Background(), zeroWorkflowID(), nil, Wait)
32 if err != nil {
33 t.Fatalf("first AcquireWorkflowSlot() error = %v", err)
34 }
35 releasedFirst := false
36 defer func() {
37 if !releasedFirst {
38 first.Release()
39 }
40 }()
41
42 acquired := make(chan WorkflowSlot, 1)
43 errs := make(chan error, 1)
44 go func() {
45 slot, err := slotter.AcquireWorkflowSlot(context.Background(), zeroWorkflowID(), nil, Wait)
46 if err != nil {
47 errs <- err
48 return
49 }
50 acquired <- slot
51 }()
52
53 assertNotAcquired(t, acquired, errs)
54
55 first.Release()
56 releasedFirst = true
57
58 second := waitForSlot(t, acquired, errs)
59 second.Release()
60}
61
62func TestSemaphoreSlotterHonorsContextCancellation(t *testing.T) {
63 t.Parallel()
64
65 slotter := NewSemaphoreSlotter(1)
66
67 first, err := slotter.AcquireWorkflowSlot(context.Background(), zeroWorkflowID(), nil, Wait)
68 if err != nil {
69 t.Fatalf("first AcquireWorkflowSlot() error = %v", err)
70 }
71 defer first.Release()
72
73 ctx, cancel := context.WithCancel(context.Background())
74 cancel()
75
76 _, err = slotter.AcquireWorkflowSlot(ctx, zeroWorkflowID(), nil, Wait)
77 if !errors.Is(err, context.Canceled) {
78 t.Fatalf("AcquireWorkflowSlot() error = %v, want context.Canceled", err)
79 }
80}
81
82func TestSemaphoreSlotterTryDisabledDoesNotReject(t *testing.T) {
83 t.Parallel()
84
85 slotter := NewSemaphoreSlotter(0)
86
87 for range 10 {
88 slot, err := slotter.AcquireWorkflowSlot(context.Background(), zeroWorkflowID(), nil, NoWait)
89 if err != nil {
90 t.Fatalf("AcquireWorkflowSlot(NoWait) error = %v", err)
91 }
92 slot.Release()
93 }
94}
95
96func TestSemaphoreSlotterTryRejectsWhenFull(t *testing.T) {
97 t.Parallel()
98
99 slotter := NewSemaphoreSlotter(1)
100
101 first, err := slotter.AcquireWorkflowSlot(context.Background(), zeroWorkflowID(), nil, NoWait)
102 if err != nil {
103 t.Fatalf("first AcquireWorkflowSlot(NoWait) error = %v", err)
104 }
105
106 if _, err := slotter.AcquireWorkflowSlot(context.Background(), zeroWorkflowID(), nil, NoWait); !errors.Is(err, ErrNoWorkflowSlots) {
107 t.Fatalf("AcquireWorkflowSlot(NoWait) error = %v, want ErrNoWorkflowSlots", err)
108 }
109
110 first.Release()
111
112 second, err := slotter.AcquireWorkflowSlot(context.Background(), zeroWorkflowID(), nil, NoWait)
113 if err != nil {
114 t.Fatalf("AcquireWorkflowSlot(NoWait) after release error = %v", err)
115 }
116 second.Release()
117}
118
119func assertNotAcquired(t *testing.T, acquired <-chan WorkflowSlot, errs <-chan error) {
120 t.Helper()
121
122 select {
123 case slot := <-acquired:
124 slot.Release()
125 t.Fatal("AcquireWorkflowSlot() acquired a slot before one was released")
126 case err := <-errs:
127 t.Fatalf("AcquireWorkflowSlot() returned unexpected error: %v", err)
128 case <-time.After(25 * time.Millisecond):
129 }
130}
131
132func waitForSlot(t *testing.T, acquired <-chan WorkflowSlot, errs <-chan error) WorkflowSlot {
133 t.Helper()
134
135 select {
136 case slot := <-acquired:
137 return slot
138 case err := <-errs:
139 t.Fatalf("AcquireWorkflowSlot() returned error: %v", err)
140 case <-time.After(time.Second):
141 t.Fatal("timed out waiting for slot acquisition")
142 }
143
144 return nil
145}
146
147func zeroWorkflowID() models.WorkflowId {
148 return models.WorkflowId{}
149}