1use crate::turso_assert_eq;
2use core::fmt::{self, Debug};
3use std::{
4 future::Future,
5 sync::{
6 atomic::{AtomicUsize, Ordering},
7 Arc, OnceLock,
8 },
9 task::{Poll, Waker},
10};
11
12use crate::sync::Mutex;
13
14use crate::{Buffer, CompletionError};
15
16pub type ReadComplete =
19 dyn Fn(Result<(Arc<Buffer>, i32), CompletionError>) -> Option<CompletionError> + Send + Sync;
20pub type WriteComplete = dyn Fn(Result<i32, CompletionError>) + Send + Sync;
21pub type SyncComplete = dyn Fn(Result<i32, CompletionError>) + Send + Sync;
22pub type TruncateComplete = dyn Fn(Result<i32, CompletionError>) + Send + Sync;
23
24#[must_use]
25#[derive(Debug, Clone)]
26pub struct Completion {
27 pub(super) inner: Option<Arc<CompletionInner>>,
29}
30
31impl Future for Completion {
32 type Output = Result<(), crate::LimboError>;
33
34 fn poll(self: std::pin::Pin<&mut Self>, cx: &mut std::task::Context<'_>) -> Poll<Self::Output> {
35 self.set_waker(cx.waker());
36 if self.finished() {
37 self.wake();
38 let res = self
39 .get_error()
40 .map_or(Ok(()), |err| Err(crate::LimboError::CompletionError(err)));
41 return Poll::Ready(res);
42 }
43 Poll::Pending
44 }
45}
46
47#[derive(Debug, Default)]
48struct ContextInner {
49 waker: Option<Waker>,
50 }
52
53#[derive(Debug, Clone)]
54pub struct Context {
55 inner: Arc<Mutex<ContextInner>>,
56}
57
58impl ContextInner {
59 pub fn new() -> Self {
60 Self { waker: None }
61 }
62
63 pub fn wake(&mut self) {
64 if let Some(waker) = self.waker.take() {
65 waker.wake();
66 }
67 }
68
69 pub fn set_waker(&mut self, waker: &Waker) {
70 if let Some(curr_waker) = self.waker.as_mut() {
71 if !curr_waker.will_wake(waker) {
73 let prev_waker = std::mem::replace(curr_waker, waker.clone());
74 prev_waker.wake();
75 }
76 } else {
77 self.waker = Some(waker.clone());
78 }
79 }
80}
81
82impl Default for Context {
83 fn default() -> Self {
84 Self::new()
85 }
86}
87
88impl Context {
89 pub fn new() -> Self {
90 Self {
91 inner: Arc::new(Mutex::new(ContextInner::new())),
92 }
93 }
94
95 pub fn wake(&self) {
96 self.inner.lock().wake();
97 }
98
99 pub fn set_waker(&self, waker: &Waker) {
100 self.inner.lock().set_waker(waker);
101 }
102}
103
104pub(super) struct CompletionInner {
105 completion_type: CompletionType,
106 pub(super) result: crate::sync::OnceLock<Option<CompletionError>>,
109 context: Context,
110 parent: OnceLock<Arc<GroupCompletionInner>>,
112 write_buffer: OnceLock<Arc<Buffer>>,
115}
116
117impl fmt::Debug for CompletionInner {
118 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
119 f.debug_struct("CompletionInner")
120 .field("completion_type", &self.completion_type)
121 .field("parent", &self.parent.get().is_some())
122 .finish()
123 }
124}
125
126pub struct CompletionGroup {
127 completions: Vec<Completion>,
128 callback: Box<dyn Fn(Result<i32, CompletionError>) + Send + Sync>,
129}
130
131impl CompletionGroup {
132 pub fn new<F>(callback: F) -> Self
133 where
134 F: Fn(Result<i32, CompletionError>) + Send + Sync + 'static,
135 {
136 Self {
137 completions: Vec::new(),
138 callback: Box::new(callback),
139 }
140 }
141
142 pub fn add(&mut self, completion: &Completion) {
143 self.completions.push(completion.clone());
144 }
145
146 pub fn completions(&self) -> &[Completion] {
150 &self.completions
151 }
152
153 pub fn cancel(&self) {
154 for c in &self.completions {
155 c.abort();
156 }
157 }
158
159 pub fn build(self) -> Completion {
160 let total = self.completions.len();
161 if total == 0 {
162 (self.callback)(Ok(0));
163 return Completion::new_yield();
164 }
165 let group_completion = GroupCompletion::new(self.callback, total);
166 let group = Completion::new(CompletionType::Group(group_completion));
167
168 if let CompletionType::Group(ref g) = group.get_inner().completion_type {
170 let _ = g.inner.self_completion.set(group.clone());
171 }
172
173 for mut c in self.completions {
174 if !c.finished() {
176 c.link_internal(&group);
177 continue;
178 }
179 let group_inner = match &group.get_inner().completion_type {
180 CompletionType::Group(g) => &g.inner,
181 _ => unreachable!(),
182 };
183 if let Some(err) = c.get_error() {
185 let _ = group_inner.result.set(Some(err));
186 group_inner.outstanding.store(0, Ordering::SeqCst);
187 (group_inner.complete)(Err(err));
188 return group;
189 }
190 group_inner.outstanding.fetch_sub(1, Ordering::SeqCst);
192 }
193
194 let group_inner = match &group.get_inner().completion_type {
195 CompletionType::Group(g) => &g.inner,
196 _ => unreachable!(),
197 };
198 if group_inner.outstanding.load(Ordering::SeqCst) == 0 {
199 let _ = group_inner.result.set(None);
201 (group_inner.complete)(Ok(0));
202 }
203 group
204 }
205}
206
207pub struct GroupCompletion {
208 inner: Arc<GroupCompletionInner>,
209}
210
211impl fmt::Debug for GroupCompletion {
212 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
213 f.debug_struct("GroupCompletion")
214 .field(
215 "outstanding",
216 &self.inner.outstanding.load(Ordering::SeqCst),
217 )
218 .finish()
219 }
220}
221
222struct GroupCompletionInner {
223 outstanding: AtomicUsize,
225 complete: Box<dyn Fn(Result<i32, CompletionError>) + Send + Sync>,
227 result: OnceLock<Option<CompletionError>>,
229 self_completion: OnceLock<Completion>,
231}
232
233impl GroupCompletion {
234 pub fn new<F>(complete: F, outstanding: usize) -> Self
235 where
236 F: Fn(Result<i32, CompletionError>) + Send + Sync + 'static,
237 {
238 Self {
239 inner: Arc::new(GroupCompletionInner {
240 outstanding: AtomicUsize::new(outstanding),
241 complete: Box::new(complete),
242 result: OnceLock::new(),
243 self_completion: OnceLock::new(),
244 }),
245 }
246 }
247
248 pub fn callback(&self, result: Result<i32, CompletionError>) {
249 turso_assert_eq!(
250 self.inner.outstanding.load(Ordering::SeqCst),
251 0,
252 "callback called before all completions finished"
253 );
254 (self.inner.complete)(result);
255 }
256}
257
258impl Debug for CompletionType {
259 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
260 match self {
261 Self::Read(..) => f.debug_tuple("Read").finish(),
262 Self::Write(..) => f.debug_tuple("Write").finish(),
263 Self::Sync(..) => f.debug_tuple("Sync").finish(),
264 Self::Truncate(..) => f.debug_tuple("Truncate").finish(),
265 Self::Group(..) => f.debug_tuple("Group").finish(),
266 Self::Yield => f.debug_tuple("Yield").finish(),
267 }
268 }
269}
270
271pub enum CompletionType {
272 Read(ReadCompletion),
273 Write(WriteCompletion),
274 Sync(SyncCompletion),
275 Truncate(TruncateCompletion),
276 Group(GroupCompletion),
277 Yield,
278}
279
280impl CompletionInner {
281 fn new(completion_type: CompletionType) -> Self {
282 Self {
283 completion_type,
284 result: OnceLock::new(),
285 context: Context::new(),
286 parent: OnceLock::new(),
287 write_buffer: OnceLock::new(),
288 }
289 }
290}
291
292impl Completion {
293 pub fn new(completion_type: CompletionType) -> Self {
294 Self {
295 inner: Some(Arc::new(CompletionInner::new(completion_type))),
296 }
297 }
298
299 pub(super) fn get_inner(&self) -> &Arc<CompletionInner> {
300 self.inner
301 .as_ref()
302 .expect("completion inner should be initialized")
303 }
304
305 pub fn keep_write_buffer_alive(&self, buf: Arc<Buffer>) {
309 self.get_inner()
310 .write_buffer
311 .set(buf)
312 .expect("write buffer should only be set once");
313 }
314
315 pub fn new_write<F>(complete: F) -> Self
316 where
317 F: Fn(Result<i32, CompletionError>) + Send + Sync + 'static,
318 {
319 Self::new(CompletionType::Write(WriteCompletion::new(Box::new(
320 complete,
321 ))))
322 }
323
324 pub fn new_read<F>(buf: Arc<Buffer>, complete: F) -> Self
325 where
326 F: Fn(Result<(Arc<Buffer>, i32), CompletionError>) -> Option<CompletionError>
327 + Send
328 + Sync
329 + 'static,
330 {
331 Self::new(CompletionType::Read(ReadCompletion::new(
332 buf,
333 Box::new(complete),
334 )))
335 }
336 pub fn new_sync<F>(complete: F) -> Self
337 where
338 F: Fn(Result<i32, CompletionError>) + Send + Sync + 'static,
339 {
340 Self::new(CompletionType::Sync(SyncCompletion::new(Box::new(
341 complete,
342 ))))
343 }
344
345 pub fn new_trunc<F>(complete: F) -> Self
346 where
347 F: Fn(Result<i32, CompletionError>) + Send + Sync + 'static,
348 {
349 Self::new(CompletionType::Truncate(TruncateCompletion::new(Box::new(
350 complete,
351 ))))
352 }
353
354 pub fn new_yield() -> Self {
357 Self { inner: None }
358 }
359
360 pub fn wake(&self) {
361 if let Some(inner) = &self.inner {
362 inner.context.wake();
363 }
364 }
365
366 pub fn wake_progress(&self) {
375 if let Some(inner) = &self.inner {
376 if let Some(group) = inner.parent.get() {
377 if let Some(group_completion) = group.self_completion.get() {
378 group_completion.wake();
379 }
380 }
381 inner.context.wake();
382 }
383 }
384
385 pub fn set_waker(&self, waker: &Waker) {
386 if self.finished() || self.inner.is_none() {
387 waker.wake_by_ref();
388 } else {
389 self.get_inner().context.set_waker(waker);
390 }
391 }
392
393 pub fn succeeded(&self) -> bool {
394 match &self.inner {
395 Some(inner) => match &inner.completion_type {
396 CompletionType::Group(g) => {
397 g.inner.outstanding.load(Ordering::SeqCst) == 0
398 && g.inner.result.get().is_some_and(|e| e.is_none())
399 }
400 _ => inner.result.get().is_some_and(|e| e.is_none()),
401 },
402 None => true,
403 }
404 }
405
406 pub fn failed(&self) -> bool {
407 match &self.inner {
408 Some(inner) => inner.result.get().is_some_and(|val| val.is_some()),
409 None => false,
410 }
411 }
412
413 pub fn get_error(&self) -> Option<CompletionError> {
414 match &self.inner {
415 Some(inner) => {
416 match &inner.completion_type {
417 CompletionType::Group(g) => {
418 g.inner.result.get().and_then(|res| *res)
421 }
422 _ => inner.result.get().and_then(|res| *res),
423 }
424 }
425 None => None,
426 }
427 }
428
429 pub fn finished(&self) -> bool {
431 match &self.inner {
432 Some(inner) => match &inner.completion_type {
433 CompletionType::Group(g) => g.inner.outstanding.load(Ordering::SeqCst) == 0,
434 _ => inner.result.get().is_some(),
435 },
436 None => true,
437 }
438 }
439
440 pub fn is_explicit_yield(&self) -> bool {
447 self.inner.is_none()
448 }
449
450 pub fn complete(&self, result: i32) {
451 let result = Ok(result);
452 self.callback(result);
453 }
454
455 pub fn error(&self, err: CompletionError) {
456 let result = Err(err);
457 self.callback(result);
458 }
459
460 pub fn abort(&self) {
461 self.error(CompletionError::Aborted);
462 }
463
464 fn callback(&self, result: Result<i32, CompletionError>) {
465 let inner = self.get_inner();
466 inner.result.get_or_init(|| {
467 let callback_error = match &inner.completion_type {
470 CompletionType::Read(r) => r.callback(result),
471 CompletionType::Write(w) => {
472 w.callback(result);
473 None
474 }
475 CompletionType::Sync(s) => {
476 s.callback(result);
477 None
478 }
479 CompletionType::Truncate(t) => {
480 t.callback(result);
481 None
482 }
483 CompletionType::Group(g) => {
484 g.callback(result);
485 None
486 }
487 CompletionType::Yield => None,
488 };
489
490 let final_error = callback_error.or_else(|| result.err());
492
493 if let Some(group) = inner.parent.get() {
494 if let Some(err) = final_error {
496 let _ = group.result.set(Some(err));
497 }
498 let prev = group.outstanding.fetch_sub(1, Ordering::SeqCst);
499 if prev > 1 {
500 if let Some(group_completion) = group.self_completion.get() {
503 group_completion.wake();
504 }
505 }
506 if prev == 1 {
509 let _ = group.result.set(None);
511 if let Some(group_completion) = group.self_completion.get() {
512 let group_result = group.result.get().and_then(|e| *e);
513 group_completion.callback(group_result.map_or(Ok(0), Err));
514 }
515 }
516 }
517
518 final_error
519 });
520 inner.context.wake();
522 }
523
524 pub fn as_read(&self) -> &ReadCompletion {
527 let inner = self.get_inner();
528 match inner.completion_type {
529 CompletionType::Read(ref r) => r,
530 _ => unreachable!(),
531 }
532 }
533
534 fn link_internal(&mut self, group: &Completion) {
536 let group_inner = match &group.get_inner().completion_type {
537 CompletionType::Group(g) => &g.inner,
538 _ => panic!("link_internal() requires a group completion"),
539 };
540
541 if self.get_inner().parent.set(group_inner.clone()).is_err() {
543 panic!("completion can only be linked once");
544 }
545 }
546}
547
548pub struct ReadCompletion {
549 pub buf: Arc<Buffer>,
550 pub complete: Box<ReadComplete>,
551}
552
553impl ReadCompletion {
554 pub fn new(buf: Arc<Buffer>, complete: Box<ReadComplete>) -> Self {
555 Self { buf, complete }
556 }
557
558 pub fn buf(&self) -> &Buffer {
559 &self.buf
560 }
561
562 pub fn callback(&self, bytes_read: Result<i32, CompletionError>) -> Option<CompletionError> {
563 (self.complete)(bytes_read.map(|b| (self.buf.clone(), b)))
564 }
565
566 pub fn buf_arc(&self) -> Arc<Buffer> {
567 self.buf.clone()
568 }
569}
570
571pub struct WriteCompletion {
572 pub complete: Box<WriteComplete>,
573}
574
575impl WriteCompletion {
576 pub fn new(complete: Box<WriteComplete>) -> Self {
577 Self { complete }
578 }
579
580 pub fn callback(&self, bytes_written: Result<i32, CompletionError>) {
581 (self.complete)(bytes_written);
582 }
583}
584
585pub struct SyncCompletion {
586 pub complete: Box<SyncComplete>,
587}
588
589impl SyncCompletion {
590 pub fn new(complete: Box<SyncComplete>) -> Self {
591 Self { complete }
592 }
593
594 pub fn callback(&self, res: Result<i32, CompletionError>) {
595 (self.complete)(res);
596 }
597}
598
599pub struct TruncateCompletion {
600 pub complete: Box<TruncateComplete>,
601}
602
603impl TruncateCompletion {
604 pub fn new(complete: Box<TruncateComplete>) -> Self {
605 Self { complete }
606 }
607
608 pub fn callback(&self, res: Result<i32, CompletionError>) {
609 (self.complete)(res);
610 }
611}
612
613#[cfg(clt_turso_tests)]
614mod tests {
615 use crate::CompletionError;
616
617 use super::*;
618
619 #[test]
620 fn test_completion_group_empty() {
621 use crate::sync::atomic::{AtomicBool, Ordering};
622
623 let callback_called = Arc::new(AtomicBool::new(false));
624 let callback_called_clone = callback_called.clone();
625
626 let group = CompletionGroup::new(move |_| {
627 callback_called_clone.store(true, Ordering::SeqCst);
628 });
629 let group = group.build();
630 assert!(group.finished());
631 assert!(group.succeeded());
632 assert!(group.get_error().is_none());
633
634 assert!(
636 callback_called.load(Ordering::SeqCst),
637 "callback should be called for empty group"
638 );
639 }
640
641 #[test]
642 fn test_completion_group_single_completion() {
643 let mut group = CompletionGroup::new(|_| {});
644 let c = Completion::new_write(|_| {});
645 group.add(&c);
646 let group = group.build();
647
648 assert!(!group.finished());
649 assert!(!group.succeeded());
650
651 c.complete(0);
652
653 assert!(group.finished());
654 assert!(group.succeeded());
655 assert!(group.get_error().is_none());
656 }
657
658 #[test]
659 fn test_completion_group_multiple_completions() {
660 let mut group = CompletionGroup::new(|_| {});
661 let c1 = Completion::new_write(|_| {});
662 let c2 = Completion::new_write(|_| {});
663 let c3 = Completion::new_write(|_| {});
664 group.add(&c1);
665 group.add(&c2);
666 group.add(&c3);
667 let group = group.build();
668
669 assert!(!group.succeeded());
670 assert!(!group.finished());
671
672 c1.complete(0);
673 assert!(!group.succeeded());
674 assert!(!group.finished());
675
676 c2.complete(0);
677 assert!(!group.succeeded());
678 assert!(!group.finished());
679
680 c3.complete(0);
681 assert!(group.succeeded());
682 assert!(group.finished());
683 }
684
685 #[test]
686 fn test_completion_group_with_error() {
687 let mut group = CompletionGroup::new(|_| {});
688 let c1 = Completion::new_write(|_| {});
689 let c2 = Completion::new_write(|_| {});
690 group.add(&c1);
691 group.add(&c2);
692 let group = group.build();
693
694 c1.complete(0);
695 c2.error(CompletionError::Aborted);
696
697 assert!(group.finished());
698 assert!(!group.succeeded());
699 assert_eq!(group.get_error(), Some(CompletionError::Aborted));
700 }
701
702 #[test]
703 fn test_completion_group_callback() {
704 use crate::sync::atomic::{AtomicBool, Ordering};
705 let called = Arc::new(AtomicBool::new(false));
706 let called_clone = called.clone();
707
708 let mut group = CompletionGroup::new(move |_| {
709 called_clone.store(true, Ordering::SeqCst);
710 });
711
712 let c1 = Completion::new_write(|_| {});
713 let c2 = Completion::new_write(|_| {});
714 group.add(&c1);
715 group.add(&c2);
716 let group = group.build();
717
718 assert!(!called.load(Ordering::SeqCst));
719
720 c1.complete(0);
721 assert!(!called.load(Ordering::SeqCst));
722
723 c2.complete(0);
724 assert!(called.load(Ordering::SeqCst));
725 assert!(group.finished());
726 assert!(group.succeeded());
727 }
728
729 #[test]
730 fn test_completion_group_some_already_completed() {
731 let mut group = CompletionGroup::new(|_| {});
733 let c1 = Completion::new_write(|_| {});
734 let c2 = Completion::new_write(|_| {});
735 let c3 = Completion::new_write(|_| {});
736
737 group.add(&c1);
739 group.add(&c2);
740 group.add(&c3);
741
742 c1.complete(0);
744 c2.complete(0);
745
746 let group = group.build();
747
748 assert!(!group.finished());
751 assert!(!group.succeeded());
752
753 c3.complete(0);
755
756 assert!(group.finished());
758 assert!(group.succeeded());
759 assert!(group.get_error().is_none());
760 }
761
762 #[test]
763 fn test_completion_group_all_already_completed() {
764 let mut group = CompletionGroup::new(|_| {});
766 let c1 = Completion::new_write(|_| {});
767 let c2 = Completion::new_write(|_| {});
768
769 c1.complete(0);
771 c2.complete(0);
772
773 group.add(&c1);
774 group.add(&c2);
775
776 let group = group.build();
777
778 assert!(group.finished());
780 assert!(group.succeeded());
781 assert!(group.get_error().is_none());
782 }
783
784 #[test]
785 fn test_completion_group_mixed_finished_and_pending() {
786 use crate::sync::atomic::{AtomicBool, Ordering};
787 let called = Arc::new(AtomicBool::new(false));
788 let called_clone = called.clone();
789
790 let mut group = CompletionGroup::new(move |_| {
791 called_clone.store(true, Ordering::SeqCst);
792 });
793
794 let c1 = Completion::new_write(|_| {});
795 let c2 = Completion::new_write(|_| {});
796 let c3 = Completion::new_write(|_| {});
797 let c4 = Completion::new_write(|_| {});
798
799 c1.complete(0);
801 c3.complete(0);
802
803 group.add(&c1);
804 group.add(&c2);
805 group.add(&c3);
806 group.add(&c4);
807
808 let group = group.build();
809
810 assert!(!group.finished());
812 assert!(!called.load(Ordering::SeqCst));
813
814 c2.complete(0);
815 assert!(!group.finished());
816 assert!(!called.load(Ordering::SeqCst));
817
818 c4.complete(0);
819 assert!(group.finished());
820 assert!(group.succeeded());
821 assert!(called.load(Ordering::SeqCst));
822 }
823
824 #[test]
825 fn test_completion_group_already_completed_with_error() {
826 let mut group = CompletionGroup::new(|_| {});
828 let c1 = Completion::new_write(|_| {});
829 let c2 = Completion::new_write(|_| {});
830
831 c1.error(CompletionError::Aborted);
833
834 group.add(&c1);
835 group.add(&c2);
836
837 let group = group.build();
838
839 assert!(group.finished());
841 assert!(!group.succeeded());
842 assert_eq!(group.get_error(), Some(CompletionError::Aborted));
843 }
844
845 #[test]
846 fn test_completion_group_tracks_all_completions() {
847 use crate::sync::atomic::{AtomicUsize, Ordering};
853
854 let callback_count = Arc::new(AtomicUsize::new(0));
855 let callback_count_clone = callback_count.clone();
856
857 let mut completions = Vec::new();
860
861 for _ in 0..4 {
863 completions.push(Completion::new_write(|_| {}));
864 }
865
866 completions[0].complete(0);
868 completions[2].complete(0);
869
870 let mut group = CompletionGroup::new(move |_| {
872 callback_count_clone.fetch_add(1, Ordering::SeqCst);
873 });
874
875 for c in &completions {
877 group.add(c);
878 }
879
880 let group = group.build();
881
882 assert!(!group.finished());
887 assert_eq!(callback_count.load(Ordering::SeqCst), 0);
888
889 completions[1].complete(0);
891 assert!(!group.finished());
892 assert_eq!(callback_count.load(Ordering::SeqCst), 0);
893
894 completions[3].complete(0);
896 assert!(group.finished());
897 assert!(group.succeeded());
898 assert_eq!(callback_count.load(Ordering::SeqCst), 1);
899
900 assert!(group.get_error().is_none());
902 }
903
904 #[test]
905 fn test_completion_group_with_all_finished_successfully() {
906 use crate::sync::atomic::{AtomicBool, Ordering};
909
910 let callback_called = Arc::new(AtomicBool::new(false));
911 let callback_called_clone = callback_called.clone();
912
913 let mut completions = Vec::new();
914
915 for _ in 0..3 {
917 let c = Completion::new_write(|_| {});
918 c.complete(0);
919 completions.push(c);
920 }
921
922 let mut group = CompletionGroup::new(move |_| {
924 callback_called_clone.store(true, Ordering::SeqCst);
925 });
926
927 for c in &completions {
928 group.add(c);
929 }
930
931 let group = group.build();
932
933 assert!(group.finished());
935 assert!(group.succeeded());
936 assert!(callback_called.load(Ordering::SeqCst));
937 assert!(group.get_error().is_none());
938 }
939
940 #[test]
941 fn test_completion_group_nested() {
942 use crate::sync::atomic::{AtomicUsize, Ordering};
943
944 let parent_called = Arc::new(AtomicUsize::new(0));
946 let child1_called = Arc::new(AtomicUsize::new(0));
947 let child2_called = Arc::new(AtomicUsize::new(0));
948
949 let child1_called_clone = child1_called.clone();
951 let mut child_group1 = CompletionGroup::new(move |_| {
952 child1_called_clone.fetch_add(1, Ordering::SeqCst);
953 });
954 let c1 = Completion::new_write(|_| {});
955 let c2 = Completion::new_write(|_| {});
956 child_group1.add(&c1);
957 child_group1.add(&c2);
958 let child_group1 = child_group1.build();
959
960 let child2_called_clone = child2_called.clone();
962 let mut child_group2 = CompletionGroup::new(move |_| {
963 child2_called_clone.fetch_add(1, Ordering::SeqCst);
964 });
965 let c3 = Completion::new_write(|_| {});
966 let c4 = Completion::new_write(|_| {});
967 child_group2.add(&c3);
968 child_group2.add(&c4);
969 let child_group2 = child_group2.build();
970
971 let parent_called_clone = parent_called.clone();
973 let mut parent_group = CompletionGroup::new(move |_| {
974 parent_called_clone.fetch_add(1, Ordering::SeqCst);
975 });
976 parent_group.add(&child_group1);
977 parent_group.add(&child_group2);
978 let parent_group = parent_group.build();
979
980 assert!(!parent_group.finished());
982 assert!(!child_group1.finished());
983 assert!(!child_group2.finished());
984 assert_eq!(parent_called.load(Ordering::SeqCst), 0);
985 assert_eq!(child1_called.load(Ordering::SeqCst), 0);
986 assert_eq!(child2_called.load(Ordering::SeqCst), 0);
987
988 c1.complete(0);
990 assert!(!child_group1.finished());
991 assert!(!parent_group.finished());
992 assert_eq!(child1_called.load(Ordering::SeqCst), 0);
993 assert_eq!(parent_called.load(Ordering::SeqCst), 0);
994
995 c2.complete(0);
997 assert!(child_group1.finished());
998 assert!(child_group1.succeeded());
999 assert_eq!(child1_called.load(Ordering::SeqCst), 1);
1000
1001 assert!(!parent_group.finished());
1003 assert_eq!(parent_called.load(Ordering::SeqCst), 0);
1004
1005 c3.complete(0);
1007 assert!(!child_group2.finished());
1008 assert!(!parent_group.finished());
1009 assert_eq!(child2_called.load(Ordering::SeqCst), 0);
1010 assert_eq!(parent_called.load(Ordering::SeqCst), 0);
1011
1012 c4.complete(0);
1014 assert!(child_group2.finished());
1015 assert!(child_group2.succeeded());
1016 assert_eq!(child2_called.load(Ordering::SeqCst), 1);
1017
1018 assert!(parent_group.finished());
1020 assert!(parent_group.succeeded());
1021 assert_eq!(parent_called.load(Ordering::SeqCst), 1);
1022 assert!(parent_group.get_error().is_none());
1023 }
1024
1025 #[test]
1026 fn test_completion_group_nested_with_error() {
1027 use crate::sync::atomic::{AtomicBool, Ordering};
1028
1029 let parent_called = Arc::new(AtomicBool::new(false));
1030 let child_called = Arc::new(AtomicBool::new(false));
1031
1032 let child_called_clone = child_called.clone();
1034 let mut child_group = CompletionGroup::new(move |_| {
1035 child_called_clone.store(true, Ordering::SeqCst);
1036 });
1037 let c1 = Completion::new_write(|_| {});
1038 let c2 = Completion::new_write(|_| {});
1039 child_group.add(&c1);
1040 child_group.add(&c2);
1041 let child_group = child_group.build();
1042
1043 let parent_called_clone = parent_called.clone();
1045 let mut parent_group = CompletionGroup::new(move |_| {
1046 parent_called_clone.store(true, Ordering::SeqCst);
1047 });
1048 let c3 = Completion::new_write(|_| {});
1049 parent_group.add(&child_group);
1050 parent_group.add(&c3);
1051 let parent_group = parent_group.build();
1052
1053 c1.complete(0);
1055 c2.complete(0);
1056 assert!(child_group.finished());
1057 assert!(child_group.succeeded());
1058 assert!(child_called.load(Ordering::SeqCst));
1059
1060 assert!(!parent_group.finished());
1062 assert!(!parent_called.load(Ordering::SeqCst));
1063
1064 c3.error(CompletionError::Aborted);
1066
1067 assert!(parent_group.finished());
1069 assert!(!parent_group.succeeded());
1070 assert_eq!(parent_group.get_error(), Some(CompletionError::Aborted));
1071 assert!(parent_called.load(Ordering::SeqCst));
1072 }
1073
1074 #[test]
1077 fn test_write_completion_pending_status() {
1078 let c = Completion::new_write(|_| {});
1079
1080 assert!(!c.finished());
1082 assert!(!c.succeeded());
1083 assert!(!c.failed());
1084 assert!(c.get_error().is_none());
1085 }
1086
1087 #[test]
1088 fn test_write_completion_success() {
1089 let c = Completion::new_write(|_| {});
1090
1091 c.complete(42);
1092
1093 assert!(c.finished());
1094 assert!(c.succeeded());
1095 assert!(!c.failed());
1096 assert!(c.get_error().is_none());
1097 }
1098
1099 #[test]
1100 fn test_write_completion_failure() {
1101 let c = Completion::new_write(|_| {});
1102
1103 c.error(CompletionError::Aborted);
1104
1105 assert!(c.finished());
1106 assert!(!c.succeeded());
1107 assert!(c.failed());
1108 assert_eq!(c.get_error(), Some(CompletionError::Aborted));
1109 }
1110
1111 #[test]
1112 fn test_read_completion_pending_status() {
1113 let buf = Arc::new(crate::Buffer::new_temporary(4096));
1114 let c = Completion::new_read(buf, |_| None);
1115
1116 assert!(!c.finished());
1117 assert!(!c.succeeded());
1118 assert!(!c.failed());
1119 assert!(c.get_error().is_none());
1120 }
1121
1122 #[test]
1123 fn test_read_completion_success() {
1124 let buf = Arc::new(crate::Buffer::new_temporary(4096));
1125 let c = Completion::new_read(buf, |_| None);
1126
1127 c.complete(1024);
1128
1129 assert!(c.finished());
1130 assert!(c.succeeded());
1131 assert!(!c.failed());
1132 assert!(c.get_error().is_none());
1133 }
1134
1135 #[test]
1136 fn test_read_completion_failure() {
1137 let buf = Arc::new(crate::Buffer::new_temporary(4096));
1138 let c = Completion::new_read(buf, |_| None);
1139
1140 c.error(CompletionError::Aborted);
1141
1142 assert!(c.finished());
1143 assert!(!c.succeeded());
1144 assert!(c.failed());
1145 assert_eq!(c.get_error(), Some(CompletionError::Aborted));
1146 }
1147
1148 #[test]
1149 fn test_sync_completion_pending_status() {
1150 let c = Completion::new_sync(|_| {});
1151
1152 assert!(!c.finished());
1153 assert!(!c.succeeded());
1154 assert!(!c.failed());
1155 assert!(c.get_error().is_none());
1156 }
1157
1158 #[test]
1159 fn test_sync_completion_success() {
1160 let c = Completion::new_sync(|_| {});
1161
1162 c.complete(0);
1163
1164 assert!(c.finished());
1165 assert!(c.succeeded());
1166 assert!(!c.failed());
1167 assert!(c.get_error().is_none());
1168 }
1169
1170 #[test]
1171 fn test_sync_completion_failure() {
1172 let c = Completion::new_sync(|_| {});
1173
1174 c.error(CompletionError::Aborted);
1175
1176 assert!(c.finished());
1177 assert!(!c.succeeded());
1178 assert!(c.failed());
1179 assert_eq!(c.get_error(), Some(CompletionError::Aborted));
1180 }
1181
1182 #[test]
1183 fn test_truncate_completion_pending_status() {
1184 let c = Completion::new_trunc(|_| {});
1185
1186 assert!(!c.finished());
1187 assert!(!c.succeeded());
1188 assert!(!c.failed());
1189 assert!(c.get_error().is_none());
1190 }
1191
1192 #[test]
1193 fn test_truncate_completion_success() {
1194 let c = Completion::new_trunc(|_| {});
1195
1196 c.complete(0);
1197
1198 assert!(c.finished());
1199 assert!(c.succeeded());
1200 assert!(!c.failed());
1201 assert!(c.get_error().is_none());
1202 }
1203
1204 #[test]
1205 fn test_truncate_completion_failure() {
1206 let c = Completion::new_trunc(|_| {});
1207
1208 c.error(CompletionError::Aborted);
1209
1210 assert!(c.finished());
1211 assert!(!c.succeeded());
1212 assert!(c.failed());
1213 assert_eq!(c.get_error(), Some(CompletionError::Aborted));
1214 }
1215
1216 #[test]
1217 fn test_yield_completion_status() {
1218 let c = Completion::new_yield();
1219
1220 assert!(c.finished());
1222 assert!(c.succeeded());
1223 assert!(!c.failed());
1224 assert!(c.get_error().is_none());
1225 }
1226
1227 #[test]
1228 fn test_completion_abort() {
1229 let c = Completion::new_write(|_| {});
1230
1231 c.abort();
1232
1233 assert!(c.finished());
1234 assert!(!c.succeeded());
1235 assert!(c.failed());
1236 assert_eq!(c.get_error(), Some(CompletionError::Aborted));
1237 }
1238
1239 #[test]
1240 fn test_completion_callback_receives_success_result() {
1241 use crate::sync::atomic::{AtomicI32, Ordering};
1242
1243 let result_value = Arc::new(AtomicI32::new(-1));
1244 let result_value_clone = result_value.clone();
1245
1246 let c = Completion::new_write(move |res| {
1247 if let Ok(val) = res {
1248 result_value_clone.store(val, Ordering::SeqCst);
1249 }
1250 });
1251
1252 c.complete(42);
1253
1254 assert_eq!(result_value.load(Ordering::SeqCst), 42);
1255 assert!(c.succeeded());
1256 }
1257
1258 #[test]
1259 fn test_completion_callback_receives_error_result() {
1260 use crate::sync::atomic::{AtomicBool, Ordering};
1261
1262 let got_error = Arc::new(AtomicBool::new(false));
1263 let got_error_clone = got_error.clone();
1264
1265 let c = Completion::new_write(move |res| {
1266 if res.is_err() {
1267 got_error_clone.store(true, Ordering::SeqCst);
1268 }
1269 });
1270
1271 c.error(CompletionError::Aborted);
1272
1273 assert!(got_error.load(Ordering::SeqCst));
1274 assert!(c.failed());
1275 }
1276
1277 #[test]
1278 fn test_completion_idempotent_complete() {
1279 use crate::sync::atomic::{AtomicUsize, Ordering};
1281
1282 let call_count = Arc::new(AtomicUsize::new(0));
1283 let call_count_clone = call_count.clone();
1284
1285 let c = Completion::new_write(move |_| {
1286 call_count_clone.fetch_add(1, Ordering::SeqCst);
1287 });
1288
1289 c.complete(1);
1290 c.complete(2);
1291 c.complete(3);
1292
1293 assert_eq!(call_count.load(Ordering::SeqCst), 1);
1295 assert!(c.succeeded());
1296 }
1297
1298 #[test]
1299 fn test_completion_idempotent_error() {
1300 use crate::sync::atomic::{AtomicUsize, Ordering};
1302
1303 let call_count = Arc::new(AtomicUsize::new(0));
1304 let call_count_clone = call_count.clone();
1305
1306 let c = Completion::new_write(move |_| {
1307 call_count_clone.fetch_add(1, Ordering::SeqCst);
1308 });
1309
1310 c.error(CompletionError::Aborted);
1311 c.error(CompletionError::Aborted);
1312 c.complete(0); assert_eq!(call_count.load(Ordering::SeqCst), 1);
1316 assert!(c.failed());
1317 }
1318}