Skip to main content

clt_database/io/
completions.rs

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
16/// Callback for read completions. Returns `Some(error)` if the callback detects an error
17/// (e.g., short read), which will be stored in the completion and propagated to VDBE.
18pub 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    /// Optional completion state. If None, it means we are Yield in order to not allocate anything
28    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    // TODO: add abort signal
51}
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            // only call and change waker if it would awake a different task
72            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    /// None means we completed successfully
107    // Thread safe with OnceLock
108    pub(super) result: crate::sync::OnceLock<Option<CompletionError>>,
109    context: Context,
110    /// Optional parent group this completion belongs to
111    parent: OnceLock<Arc<GroupCompletionInner>>,
112    /// Keeps the write buffer alive for async I/O backends (io_uring, VFS)
113    /// where pwrite returns before the kernel has consumed the buffer.
114    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    /// The children added so far. Used by error paths that need to
147    /// wait on the kernel side via `IO::drain_completions` after
148    /// cancelling the group.
149    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        // Store the group completion reference for later callback
169        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 the completion has not completed, link it to the group.
175            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            // Return early if there was an error.
184            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            // Mark the successful completion as done.
191            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            // Set result to Some(None) on success so succeeded() returns true
200            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    /// Number of completions that need to finish
224    outstanding: AtomicUsize,
225    /// Callback to invoke when all completions finish
226    complete: Box<dyn Fn(Result<i32, CompletionError>) + Send + Sync>,
227    /// Cached result after all completions finish
228    result: OnceLock<Option<CompletionError>>,
229    /// Reference to the group's own Completion for notifying parents
230    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    /// Stores a write buffer reference in the completion to keep it alive
306    /// until the I/O completes. Required for async backends (io_uring, VFS)
307    /// where pwrite returns before the kernel has consumed the buffer.
308    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    /// Create a yield completion. These are completed by default allowing to yield control without
355    /// allocating memory.
356    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    /// Fire a "progress wake" without consuming the completion's result —
367    /// wakes the waker on this completion *and* on its parent group, if any.
368    /// Used by IO backends to nudge the future when an operation made progress
369    /// but isn't fully done (e.g. an io_uring writev that completed only the
370    /// first chunk and was resubmitted internally). Without this, a poll
371    /// whose drained CQEs are all intermediate-chunk completions returns
372    /// without waking anything, and since `step()` is the only thing that
373    /// drains the CQ, the resubmitted chunks pile up and the task deadlocks.
374    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                        // For groups, check the group's cached result field
419                        // (set when the last completion finishes)
420                        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    /// Checks if the Completion completed or errored
430    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    /// Returns true if this completion is an explicit yield — a signal to
441    /// return control to the cooperative scheduler so other connections can make
442    /// progress. Unlike real I/O completions that happen to be finished,
443    /// yield completions must not be treated as "ready to continue immediately"
444    /// because the yielding operation is waiting on external state (e.g. a lock
445    /// held by another fiber) that can only change when other fibers are stepped.
446    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            // Run the type-specific callback. For ReadCompletion, this returns
468            // an optional error detected by the callback (e.g., short read).
469            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            // Use callback error if present, otherwise use the original IO error
491            let final_error = callback_error.or_else(|| result.err());
492
493            if let Some(group) = inner.parent.get() {
494                // Capture first error in group
495                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                    // progress wake so the waiter keeps driving io.step,
501                    // If prev > 1, there are still children outstanding after this one.
502                    if let Some(group_completion) = group.self_completion.get() {
503                        group_completion.wake();
504                    }
505                }
506                // If this was the last completion in the group, trigger the group's callback
507                // which will recursively call this same callback() method to notify parents
508                if prev == 1 {
509                    // Set result to Some(None) on success so succeeded() returns true
510                    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        // call the waker regardless
521        inner.context.wake();
522    }
523
524    /// only call this method if you are sure that the completion is
525    /// a ReadCompletion, panics otherwise
526    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    /// Link this completion to a group completion (internal use only)
535    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        // Set the parent (can only be set once)
542        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        // Verify the callback was actually called
635        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        // Test some completions added to group, then finish before build()
732        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        // Add all to group while pending
738        group.add(&c1);
739        group.add(&c2);
740        group.add(&c3);
741
742        // Complete c1 and c2 AFTER adding but BEFORE build()
743        c1.complete(0);
744        c2.complete(0);
745
746        let group = group.build();
747
748        // c1 and c2 finished before build(), so outstanding should account for them
749        // Only c3 should be pending
750        assert!(!group.finished());
751        assert!(!group.succeeded());
752
753        // Complete c3
754        c3.complete(0);
755
756        // Now the group should be finished
757        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        // Test when all completions are already finished before build()
765        let mut group = CompletionGroup::new(|_| {});
766        let c1 = Completion::new_write(|_| {});
767        let c2 = Completion::new_write(|_| {});
768
769        // Complete both before adding to group
770        c1.complete(0);
771        c2.complete(0);
772
773        group.add(&c1);
774        group.add(&c2);
775
776        let group = group.build();
777
778        // All completions were already complete, so group should be finished immediately
779        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        // Complete c1 and c3 before adding to group
800        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        // Only c2 and c4 should be pending
811        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        // Test when a completion finishes with error before build()
827        let mut group = CompletionGroup::new(|_| {});
828        let c1 = Completion::new_write(|_| {});
829        let c2 = Completion::new_write(|_| {});
830
831        // Complete c1 with error before adding to group
832        c1.error(CompletionError::Aborted);
833
834        group.add(&c1);
835        group.add(&c2);
836
837        let group = group.build();
838
839        // Group should immediately fail with the error
840        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        // This test verifies the fix for the bug where CompletionGroup::add()
848        // would skip successfully-finished completions. This caused problems
849        // when code used drain() to move completions into a group, because
850        // finished completions would be removed from the source but not tracked
851        // by the group, effectively losing them.
852        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        // Simulate the pattern: create multiple completions, complete some,
858        // then add ALL of them to a group (like drain() would do)
859        let mut completions = Vec::new();
860
861        // Create 4 completions
862        for _ in 0..4 {
863            completions.push(Completion::new_write(|_| {}));
864        }
865
866        // Complete 2 of them before adding to group (simulate async completion)
867        completions[0].complete(0);
868        completions[2].complete(0);
869
870        // Now create a group and add ALL completions (like drain() would do)
871        let mut group = CompletionGroup::new(move |_| {
872            callback_count_clone.fetch_add(1, Ordering::SeqCst);
873        });
874
875        // Add all completions to the group
876        for c in &completions {
877            group.add(c);
878        }
879
880        let group = group.build();
881
882        // The group should track all 4 completions:
883        // - c[0] and c[2] are already finished
884        // - c[1] and c[3] are still pending
885        // So the group should not be finished yet
886        assert!(!group.finished());
887        assert_eq!(callback_count.load(Ordering::SeqCst), 0);
888
889        // Complete the first pending completion
890        completions[1].complete(0);
891        assert!(!group.finished());
892        assert_eq!(callback_count.load(Ordering::SeqCst), 0);
893
894        // Complete the last pending completion - now group should finish
895        completions[3].complete(0);
896        assert!(group.finished());
897        assert!(group.succeeded());
898        assert_eq!(callback_count.load(Ordering::SeqCst), 1);
899
900        // Verify no errors
901        assert!(group.get_error().is_none());
902    }
903
904    #[test]
905    fn test_completion_group_with_all_finished_successfully() {
906        // Edge case: all completions are already successfully finished
907        // when added to the group. The group should complete immediately.
908        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        // Create and immediately complete 3 completions
916        for _ in 0..3 {
917            let c = Completion::new_write(|_| {});
918            c.complete(0);
919            completions.push(c);
920        }
921
922        // Add all already-completed completions to group
923        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        // Group should be immediately finished since all completions were done
934        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        // Track callbacks at different levels
945        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        // Create child group 1 with 2 completions
950        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        // Create child group 2 with 2 completions
961        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        // Create parent group containing both child groups
972        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        // Initially nothing should be finished
981        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        // Complete first completion in child group 1
989        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        // Complete second completion in child group 1 - should finish child group 1
996        c2.complete(0);
997        assert!(child_group1.finished());
998        assert!(child_group1.succeeded());
999        assert_eq!(child1_called.load(Ordering::SeqCst), 1);
1000
1001        // Parent should not be finished yet because child group 2 is still pending
1002        assert!(!parent_group.finished());
1003        assert_eq!(parent_called.load(Ordering::SeqCst), 0);
1004
1005        // Complete first completion in child group 2
1006        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        // Complete second completion in child group 2 - should finish everything
1013        c4.complete(0);
1014        assert!(child_group2.finished());
1015        assert!(child_group2.succeeded());
1016        assert_eq!(child2_called.load(Ordering::SeqCst), 1);
1017
1018        // Parent should now be finished
1019        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        // Create child group with 2 completions
1033        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        // Create parent group containing child group and another completion
1044        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        // Complete child group with success
1054        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        // Parent still pending
1061        assert!(!parent_group.finished());
1062        assert!(!parent_called.load(Ordering::SeqCst));
1063
1064        // Complete c3 with error
1065        c3.error(CompletionError::Aborted);
1066
1067        // Parent should finish with error
1068        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    // Tests for individual completion success/failure status
1075
1076    #[test]
1077    fn test_write_completion_pending_status() {
1078        let c = Completion::new_write(|_| {});
1079
1080        // Pending completion should not be finished, succeeded, or failed
1081        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        // Yield completions are always considered finished and succeeded
1221        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        // Completing a completion multiple times should only trigger the callback once
1280        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        // Callback should only be called once
1294        assert_eq!(call_count.load(Ordering::SeqCst), 1);
1295        assert!(c.succeeded());
1296    }
1297
1298    #[test]
1299    fn test_completion_idempotent_error() {
1300        // Erroring a completion multiple times should only trigger the callback once
1301        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); // Try completing after error
1313
1314        // Callback should only be called once
1315        assert_eq!(call_count.load(Ordering::SeqCst), 1);
1316        assert!(c.failed());
1317    }
1318}