Skip to content
File

Blob: src/workerd/io/io-gate-test.c++

14.1 KB
1// Copyright (c) 2017-2022 Cloudflare, Inc.
2// Licensed under the Apache 2.0 license found in the LICENSE file or at:
3// https://opensource.org/licenses/Apache-2.0
4 
5#include "io-gate.h"
6 
7#include <kj/test.h>
8 
9namespace workerd {
10namespace {
11 
12KJ_TEST("InputGate basics") {
13 kj::EventLoop loop;
14 kj::WaitScope ws(loop);
15 
16 InputGate gate;
17 
18 kj::Promise<InputGate::Lock> promise1 = gate.wait(nullptr);
19 kj::Promise<InputGate::Lock> promise2 = gate.wait(nullptr);
20 kj::Promise<InputGate::Lock> promise3 = gate.wait(nullptr);
21 
22 KJ_ASSERT(promise1.poll(ws));
23 KJ_EXPECT(!promise2.poll(ws));
24 KJ_EXPECT(!promise3.poll(ws));
25 
26 {
27 auto lock = promise1.wait(ws);
28 
29 KJ_EXPECT(!promise2.poll(ws));
30 KJ_EXPECT(!promise3.poll(ws));
31 
32 auto lock2 = lock.addRef(nullptr);
33 { auto drop = kj::mv(lock); }
34 
35 KJ_EXPECT(!promise2.poll(ws));
36 KJ_EXPECT(!promise3.poll(ws));
37 }
38 
39 KJ_EXPECT(promise2.poll(ws));
40 KJ_EXPECT(!promise3.poll(ws)); // we'll cancel this waiter to make sure that works
41 
42 KJ_EXPECT(!gate.onBroken().poll(ws));
43}
44 
45KJ_TEST("InputGate critical section") {
46 kj::EventLoop loop;
47 kj::WaitScope ws(loop);
48 
49 InputGate gate;
50 
51 kj::Own<InputGate::CriticalSection> cs;
52 
53 {
54 auto lock = gate.wait(nullptr).wait(ws);
55 cs = lock.startCriticalSection();
56 }
57 
58 {
59 // Take the first lock.
60 auto firstLock = cs->wait(nullptr).wait(ws);
61 
62 // Other locks are blocked.
63 auto wait1 = cs->wait(nullptr);
64 auto wait2 = cs->wait(nullptr);
65 KJ_EXPECT(!wait1.poll(ws));
66 KJ_EXPECT(!wait2.poll(ws));
67 
68 // Drop it.
69 { auto drop = kj::mv(firstLock); }
70 
71 // Now other locks make progress.
72 {
73 auto lock = wait1.wait(ws);
74 KJ_EXPECT(!wait2.poll(ws));
75 }
76 wait2.wait(ws);
77 }
78 
79 // Can't lock the top-level gate while CriticalSection still exists.
80 auto outerWait = gate.wait(nullptr);
81 KJ_EXPECT(!outerWait.poll(ws));
82 
83 {
84 auto lock = cs->wait(nullptr).wait(ws);
85 cs->succeeded();
86 KJ_EXPECT(!outerWait.poll(ws));
87 }
88 
89 outerWait.wait(ws);
90}
91 
92KJ_TEST("InputGate multiple critical sections start together") {
93 kj::EventLoop loop;
94 kj::WaitScope ws(loop);
95 
96 InputGate gate;
97 
98 kj::Own<InputGate::CriticalSection> cs1;
99 kj::Own<InputGate::CriticalSection> cs2;
100 
101 {
102 auto lock = gate.wait(nullptr).wait(ws);
103 cs1 = lock.startCriticalSection();
104 cs2 = lock.startCriticalSection();
105 }
106 
107 // Start cs1.
108 cs1->wait(nullptr).wait(ws);
109 
110 // Can't start cs2 yet.
111 auto cs2Wait = cs2->wait(nullptr);
112 KJ_EXPECT(!cs2Wait.poll(ws));
113 
114 cs1->succeeded();
115 
116 cs2Wait.wait(ws);
117}
118 
119KJ_TEST("InputGate nested critical sections") {
120 kj::EventLoop loop;
121 kj::WaitScope ws(loop);
122 
123 InputGate gate;
124 
125 kj::Own<InputGate::CriticalSection> cs1;
126 kj::Own<InputGate::CriticalSection> cs2;
127 
128 {
129 auto lock = gate.wait(nullptr).wait(ws);
130 cs1 = lock.startCriticalSection();
131 }
132 
133 {
134 auto lock = cs1->wait(nullptr).wait(ws);
135 cs2 = lock.startCriticalSection();
136 }
137 
138 // Start cs2.
139 cs2->wait(nullptr).wait(ws);
140 
141 // Can't start new tasks in cs1 until cs2 finishes.
142 auto cs1Wait = cs1->wait(nullptr);
143 KJ_EXPECT(!cs1Wait.poll(ws));
144 
145 cs2->succeeded();
146 
147 cs1Wait.wait(ws);
148}
149 
150KJ_TEST("InputGate nested critical section outlives parent") {
151 kj::EventLoop loop;
152 kj::WaitScope ws(loop);
153 
154 InputGate gate;
155 
156 kj::Own<InputGate::CriticalSection> cs1;
157 kj::Own<InputGate::CriticalSection> cs2;
158 
159 {
160 auto lock = gate.wait(nullptr).wait(ws);
161 cs1 = lock.startCriticalSection();
162 }
163 
164 {
165 auto lock = cs1->wait(nullptr).wait(ws);
166 cs2 = lock.startCriticalSection();
167 }
168 
169 // Start cs2.
170 cs2->wait(nullptr).wait(ws);
171 
172 // Mark cs1 done. (Note that, in a real program, this probably can't happen like this, because a
173 // lock would be taken on cs1 before marking it done, and that lock would wait for cs2 to
174 // finish. But I want to make sure it works anyway.)
175 cs1->succeeded();
176 
177 // Can't start new tasks in at root until cs2 finishes.
178 auto rootWait = gate.wait(nullptr);
179 KJ_EXPECT(!rootWait.poll(ws));
180 
181 cs2->succeeded();
182 
183 rootWait.wait(ws);
184}
185 
186KJ_TEST("InputGate deeply nested critical sections") {
187 kj::EventLoop loop;
188 kj::WaitScope ws(loop);
189 
190 InputGate gate;
191 
192 kj::Own<InputGate::CriticalSection> cs1;
193 kj::Own<InputGate::CriticalSection> cs2;
194 kj::Own<InputGate::CriticalSection> cs3;
195 kj::Own<InputGate::CriticalSection> cs4;
196 
197 {
198 auto lock = gate.wait(nullptr).wait(ws);
199 cs1 = lock.startCriticalSection();
200 }
201 
202 {
203 auto lock = cs1->wait(nullptr).wait(ws);
204 cs2 = lock.startCriticalSection();
205 }
206 
207 {
208 auto lock = cs2->wait(nullptr).wait(ws);
209 cs3 = lock.startCriticalSection();
210 cs4 = lock.startCriticalSection();
211 }
212 
213 // Start cs2
214 cs2->wait(nullptr).wait(ws);
215 
216 // Add some waiters to cs2, some of which are waiting to start more nested critical sections
217 auto lock = cs2->wait(nullptr).wait(ws);
218 auto waiter1 = cs2->wait(nullptr);
219 auto waiter2 = cs2->wait(nullptr);
220 
221 // Both of these wait on cs2 indirectly, as they are nested under cs2
222 auto waiter3 = cs3->wait(nullptr);
223 auto waiter4 = cs4->wait(nullptr);
224 
225 KJ_EXPECT(!waiter1.poll(ws));
226 KJ_EXPECT(!waiter2.poll(ws));
227 KJ_EXPECT(!waiter3.poll(ws));
228 KJ_EXPECT(!waiter4.poll(ws));
229 
230 // Mark cs2 as complete with outstanding waiters, and drop our reference to it.
231 cs2->succeeded();
232 cs2 = nullptr;
233 
234 // Our waiters should still be outstanding as we have not released the lock
235 KJ_EXPECT(!waiter1.poll(ws));
236 KJ_EXPECT(!waiter2.poll(ws));
237 KJ_EXPECT(!waiter3.poll(ws));
238 KJ_EXPECT(!waiter4.poll(ws));
239 
240 // Drop some outstanding waiters
241 { auto drop = kj::mv(waiter2); }
242 { auto drop = kj::mv(waiter4); }
243 
244 // Release the lock on cs2
245 { auto drop = kj::mv(lock); }
246 
247 // cs3 should have started
248 KJ_ASSERT(!waiter1.poll(ws));
249 KJ_ASSERT(waiter3.poll(ws));
250 auto lock2 = waiter3.wait(ws);
251 
252 // Add a waiter on cs3
253 auto waiter5 = cs3->wait(nullptr);
254 KJ_ASSERT(!waiter5.poll(ws));
255 
256 // Can't start new tasks on the root until both cs1 and cs3 have succeeded, and all outstanding
257 // tasks have either been dropped or completed.
258 auto waiter6 = gate.wait(nullptr);
259 KJ_ASSERT(!waiter6.poll(ws));
260 
261 cs1->succeeded();
262 cs3->succeeded();
263 
264 // drop waiter5
265 { auto drop = kj::mv(waiter5); }
266 
267 // Release the lock on cs3
268 { auto drop = kj::mv(lock2); }
269 
270 // Our root task should be ready now.
271 KJ_ASSERT(waiter6.poll(ws));
272 waiter6.wait(ws);
273}
274 
275KJ_TEST("InputGate critical section lock outlives critical section") {
276 kj::EventLoop loop;
277 kj::WaitScope ws(loop);
278 
279 InputGate gate;
280 
281 kj::Own<InputGate::CriticalSection> cs;
282 
283 {
284 auto lock = gate.wait(nullptr).wait(ws);
285 cs = lock.startCriticalSection();
286 }
287 
288 // Start critical section.
289 auto lock = cs->wait(nullptr).wait(ws);
290 KJ_ASSERT(lock.isFor(gate));
291 
292 // Mark it done, even though a lock is still outstanding.
293 cs->succeeded();
294 
295 // Drop our reference.
296 cs = nullptr;
297 
298 // Lock should have been reparented, so should still work.
299 KJ_ASSERT(lock.isFor(gate));
300 
301 // Adding a ref and dropping it shouldn't cause trouble.
302 lock.addRef(nullptr);
303 
304 // The gate should still be locked
305 auto waiter = gate.wait(nullptr);
306 KJ_EXPECT(!waiter.poll(ws));
307 
308 // Drop the outstanding lock
309 { auto drop = kj::mv(lock); }
310 
311 // Our waiter should resolve now
312 KJ_ASSERT(waiter.poll(ws));
313 KJ_EXPECT(waiter.wait(ws).isFor(gate));
314}
315 
316KJ_TEST("InputGate broken") {
317 kj::EventLoop loop;
318 kj::WaitScope ws(loop);
319 
320 InputGate gate;
321 
322 auto brokenPromise = gate.onBroken();
323 
324 kj::Own<InputGate::CriticalSection> cs1;
325 kj::Own<InputGate::CriticalSection> cs2;
326 kj::Own<InputGate::CriticalSection> cs3;
327 
328 {
329 auto lock = gate.wait(nullptr).wait(ws);
330 cs1 = lock.startCriticalSection();
331 cs3 = lock.startCriticalSection();
332 }
333 
334 {
335 auto lock = cs1->wait(nullptr).wait(ws);
336 cs2 = lock.startCriticalSection();
337 }
338 
339 // start cs2
340 cs2->wait(nullptr).wait(ws);
341 
342 auto cs1Wait = cs1->wait(nullptr);
343 KJ_EXPECT(!cs1Wait.poll(ws));
344 
345 auto cs3Wait = cs3->wait(nullptr);
346 KJ_EXPECT(!cs3Wait.poll(ws));
347 
348 auto rootWait = gate.wait(nullptr);
349 KJ_EXPECT(!rootWait.poll(ws));
350 
351 cs2->failed(KJ_EXCEPTION(FAILED, "foobar"));
352 
353 KJ_EXPECT_THROW_MESSAGE("foobar", cs1Wait.wait(ws));
354 KJ_EXPECT_THROW_MESSAGE("foobar", cs3Wait.wait(ws));
355 KJ_EXPECT_THROW_MESSAGE("foobar", rootWait.wait(ws));
356 KJ_EXPECT_THROW_MESSAGE("foobar", cs2->wait(nullptr).wait(ws));
357 KJ_EXPECT_THROW_MESSAGE("foobar", brokenPromise.wait(ws));
358 KJ_EXPECT_THROW_MESSAGE("foobar", gate.onBroken().wait(ws));
359}
360 
361// =======================================================================================
362 
363KJ_TEST("OutputGate basics") {
364 kj::EventLoop loop;
365 kj::WaitScope ws(loop);
366 
367 OutputGate gate;
368 
369 KJ_EXPECT(gate.wait(nullptr).poll(ws));
370 
371 auto paf1 = kj::newPromiseAndFulfiller<void>();
372 auto blocker1 = gate.lockWhile(kj::mv(paf1.promise), nullptr);
373 
374 auto promise1 = gate.wait(nullptr);
375 auto promise2 = gate.wait(nullptr);
376 
377 auto paf2 = kj::newPromiseAndFulfiller<void>();
378 auto blocker2 = gate.lockWhile(kj::mv(paf2.promise), nullptr);
379 
380 auto promise3 = gate.wait(nullptr);
381 
382 KJ_EXPECT(!promise1.poll(ws));
383 KJ_EXPECT(!promise2.poll(ws));
384 KJ_EXPECT(!promise3.poll(ws));
385 
386 KJ_EXPECT(!blocker1.poll(ws));
387 paf1.fulfiller->fulfill();
388 KJ_EXPECT(blocker1.poll(ws));
389 blocker1.wait(ws);
390 
391 KJ_EXPECT(promise1.poll(ws));
392 promise1.wait(ws);
393 KJ_EXPECT(promise2.poll(ws));
394 promise2.wait(ws);
395 KJ_EXPECT(!promise3.poll(ws));
396 
397 KJ_EXPECT(!blocker2.poll(ws));
398 paf2.fulfiller->fulfill();
399 KJ_EXPECT(blocker2.poll(ws));
400 blocker2.wait(ws);
401 
402 KJ_EXPECT(promise3.poll(ws));
403 promise3.wait(ws);
404 
405 KJ_EXPECT(!gate.onBroken().poll(ws));
406}
407 
408KJ_TEST("OutputGate out-of-order") {
409 kj::EventLoop loop;
410 kj::WaitScope ws(loop);
411 
412 OutputGate gate;
413 
414 KJ_EXPECT(gate.wait(nullptr).poll(ws));
415 
416 auto paf1 = kj::newPromiseAndFulfiller<void>();
417 auto blocker1 = gate.lockWhile(kj::mv(paf1.promise), nullptr);
418 
419 auto promise1 = gate.wait(nullptr);
420 auto promise2 = gate.wait(nullptr);
421 
422 auto paf2 = kj::newPromiseAndFulfiller<void>();
423 auto blocker2 = gate.lockWhile(kj::mv(paf2.promise), nullptr);
424 
425 auto promise3 = gate.wait(nullptr);
426 
427 KJ_EXPECT(!promise1.poll(ws));
428 KJ_EXPECT(!promise2.poll(ws));
429 KJ_EXPECT(!promise3.poll(ws));
430 
431 // Fulfill second blocker first.
432 KJ_EXPECT(!blocker2.poll(ws));
433 paf2.fulfiller->fulfill();
434 KJ_EXPECT(blocker2.poll(ws));
435 blocker2.wait(ws);
436 
437 // Everything is still blocked.
438 KJ_EXPECT(!promise1.poll(ws));
439 KJ_EXPECT(!promise2.poll(ws));
440 KJ_EXPECT(!promise3.poll(ws));
441 
442 // Fulfill the first one.
443 KJ_EXPECT(!blocker1.poll(ws));
444 paf1.fulfiller->fulfill();
445 KJ_EXPECT(blocker1.poll(ws));
446 blocker1.wait(ws);
447 
448 // Everything unblocked.
449 KJ_EXPECT(promise1.poll(ws));
450 promise1.wait(ws);
451 KJ_EXPECT(promise2.poll(ws));
452 promise2.wait(ws);
453 KJ_EXPECT(promise3.poll(ws));
454 promise3.wait(ws);
455 
456 KJ_EXPECT(!gate.onBroken().poll(ws));
457}
458 
459KJ_TEST("OutputGate exception") {
460 kj::EventLoop loop;
461 kj::WaitScope ws(loop);
462 
463 OutputGate gate;
464 auto onBroken = gate.onBroken();
465 
466 KJ_EXPECT(gate.wait(nullptr).poll(ws));
467 
468 auto paf1 = kj::newPromiseAndFulfiller<void>();
469 auto blocker1 = gate.lockWhile(kj::mv(paf1.promise), nullptr);
470 
471 auto promise1 = gate.wait(nullptr);
472 auto promise2 = gate.wait(nullptr);
473 
474 auto paf2 = kj::newPromiseAndFulfiller<void>();
475 auto blocker2 = gate.lockWhile(kj::mv(paf2.promise), nullptr);
476 
477 auto promise3 = gate.wait(nullptr);
478 
479 KJ_EXPECT(!promise1.poll(ws));
480 KJ_EXPECT(!promise2.poll(ws));
481 KJ_EXPECT(!promise3.poll(ws));
482 
483 // Let's have the second blocker fail first.
484 paf2.fulfiller->reject(KJ_EXCEPTION(FAILED, "foo"));
485 KJ_EXPECT(blocker2.poll(ws));
486 KJ_EXPECT_THROW_MESSAGE("foo", blocker2.wait(ws));
487 
488 // Promises are all still waiting. TECHNICALLY, it would be OK to fail-fast the third promise,
489 // but for now we don't.
490 KJ_EXPECT(!promise1.poll(ws));
491 KJ_EXPECT(!promise2.poll(ws));
492 KJ_EXPECT(!promise3.poll(ws));
493 
494 // We are marked broken at this point, though.
495 KJ_ASSERT(onBroken.poll(ws));
496 KJ_EXPECT_THROW_MESSAGE("foo", onBroken.wait(ws));
497 
498 // Fulfill the first blocker (normally, not with an exception).
499 KJ_EXPECT(!blocker1.poll(ws));
500 paf1.fulfiller->fulfill();
501 KJ_EXPECT(blocker1.poll(ws));
502 blocker1.wait(ws);
503 
504 // Everything unblocked, but only the third promise fails.
505 KJ_EXPECT(promise1.poll(ws));
506 promise1.wait(ws);
507 KJ_EXPECT(promise2.poll(ws));
508 promise2.wait(ws);
509 KJ_EXPECT(promise3.poll(ws));
510 KJ_EXPECT_THROW_MESSAGE("foo", promise3.wait(ws));
511 
512 // Still broken.
513 onBroken = gate.onBroken();
514 KJ_ASSERT(onBroken.poll(ws));
515 KJ_EXPECT_THROW_MESSAGE("foo", onBroken.wait(ws));
516}
517 
518KJ_TEST("OutputGate canceled") {
519 kj::EventLoop loop;
520 kj::WaitScope ws(loop);
521 
522 OutputGate gate;
523 auto onBroken = gate.onBroken();
524 
525 KJ_EXPECT(gate.wait(nullptr).poll(ws));
526 
527 auto paf1 = kj::newPromiseAndFulfiller<void>();
528 auto blocker1 = gate.lockWhile(kj::mv(paf1.promise), nullptr);
529 
530 auto promise1 = gate.wait(nullptr);
531 auto promise2 = gate.wait(nullptr);
532 
533 auto blocker2 = gate.lockWhile(kj::Promise<void>(kj::NEVER_DONE), nullptr);
534 
535 auto promise3 = gate.wait(nullptr);
536 
537 KJ_EXPECT(!promise1.poll(ws));
538 KJ_EXPECT(!promise2.poll(ws));
539 KJ_EXPECT(!promise3.poll(ws));
540 
541 // Let's cancel the second blocker first.
542 blocker2 = nullptr;
543 
544 // Promises are all still waiting. TECHNICALLY, it would be OK to fail-fast the third promise,
545 // but for now we don't.
546 KJ_EXPECT(!promise1.poll(ws));
547 KJ_EXPECT(!promise2.poll(ws));
548 KJ_EXPECT(!promise3.poll(ws));
549 
550 // We are marked broken at this point, though.
551 KJ_ASSERT(onBroken.poll(ws));
552 KJ_EXPECT_THROW_MESSAGE("output lock was canceled before completion", onBroken.wait(ws));
553 
554 // Fulfill the first blocker (normally, not with an exception).
555 KJ_EXPECT(!blocker1.poll(ws));
556 paf1.fulfiller->fulfill();
557 KJ_EXPECT(blocker1.poll(ws));
558 blocker1.wait(ws);
559 
560 // Everything unblocked, but only the third promise fails.
561 KJ_EXPECT(promise1.poll(ws));
562 promise1.wait(ws);
563 KJ_EXPECT(promise2.poll(ws));
564 promise2.wait(ws);
565 KJ_EXPECT(promise3.poll(ws));
566 KJ_EXPECT_THROW_MESSAGE("output lock was canceled before completion", promise3.wait(ws));
567 
568 // Still broken.
569 onBroken = gate.onBroken();
570 KJ_ASSERT(onBroken.poll(ws));
571 KJ_EXPECT_THROW_MESSAGE("output lock was canceled before completion", onBroken.wait(ws));
572}
573 
574} // namespace
575} // namespace workerd