Skip to main content

ntex_io/
buf.rs

1use std::{cell::Cell, cell::Ref, cell::RefCell, fmt, io, iter, mem, task::Poll};
2
3use ntex_bytes::{BytePageSize, BytePages, BytesMut};
4
5use crate::IoRef;
6
7/// Buffers of the filter chain, ordered from the application toward the
8/// transport.
9///
10/// Without a filter layer the application and the transport share `app`.
11/// Each added layer becomes the outermost one: it gets a new `app` buffer, and
12/// the previous `app` buffer, with any data buffered in it, moves one position
13/// toward the transport. The first one becomes `wire`, further ones go to
14/// `mid`. The layer at position `i` uses the buffers at positions `i` and
15/// `i + 1`, the innermost layer writes to and reads from `wire`.
16///
17/// Adding a layer moves buffers, so it needs exclusive access to the layers.
18/// Everything that holds references into them, including while it runs code
19/// it does not control such as closures, codecs or filters, holds a shared
20/// borrow.
21pub(crate) struct Stack(RefCell<Layers>);
22
23pub(crate) struct Layers {
24    app: Buffer,
25    mid: Vec<Buffer>,
26    wire: Option<Buffer>,
27}
28
29struct Buffer {
30    read: Cell<Option<BytesMut>>,
31    write: RefCell<BytePages>,
32}
33
34impl fmt::Debug for Stack {
35    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
36        f.debug_struct("Stack")
37            .field("layers", &self.borrow().count())
38            .finish()
39    }
40}
41
42impl fmt::Debug for Layers {
43    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
44        f.debug_struct("Layers")
45            .field("layers", &self.count())
46            .finish()
47    }
48}
49
50impl Layers {
51    /// Returns the number of installed filter layers.
52    fn count(&self) -> usize {
53        self.wire.as_ref().map_or(0, |_| self.mid.len() + 1)
54    }
55
56    fn buffers(&self) -> impl Iterator<Item = &Buffer> {
57        iter::once(&self.app)
58            .chain(self.mid.iter())
59            .chain(self.wire.iter())
60    }
61
62    /// Returns the buffer at position `idx`, counted from the application.
63    fn get(&self, idx: usize) -> Option<&Buffer> {
64        if idx == 0 {
65            Some(&self.app)
66        } else if let Some(buf) = self.mid.get(idx - 1) {
67            Some(buf)
68        } else if idx == self.mid.len() + 1 {
69            self.wire.as_ref()
70        } else {
71            None
72        }
73    }
74
75    /// Returns the transport-facing buffer.
76    fn transport(&self) -> &Buffer {
77        self.wire.as_ref().unwrap_or(&self.app)
78    }
79
80    /// Returns the size of the transport-facing write buffer.
81    fn write_dst_size(&self) -> usize {
82        self.transport().write_len()
83    }
84}
85
86impl Stack {
87    pub(crate) fn new(size: BytePageSize) -> Self {
88        Self(RefCell::new(Layers {
89            app: Buffer::new(size),
90            mid: Vec::new(),
91            wire: None,
92        }))
93    }
94
95    /// Borrows the layers, the stack cannot change until the borrow is
96    /// dropped.
97    #[inline]
98    pub(crate) fn borrow(&self) -> Ref<'_, Layers> {
99        self.0.borrow()
100    }
101
102    /// Whether references into the stack may be held by running code.
103    pub(crate) fn is_borrowed(&self) -> bool {
104        self.0.try_borrow_mut().is_err()
105    }
106
107    pub(crate) fn set_page_size(&self, size: BytePageSize) {
108        for b in self.borrow().buffers() {
109            b.with_write_if_free(|b| b.set_page_size(size));
110        }
111    }
112
113    /// Adds a buffer for a new outermost layer.
114    ///
115    /// # Panics
116    ///
117    /// Panics if the stack is borrowed.
118    pub(crate) fn add_layer(&self, page_size: BytePageSize) {
119        let Ok(mut layers) = self.0.try_borrow_mut() else {
120            panic!("filter buffers are in use");
121        };
122        let outer = mem::replace(&mut layers.app, Buffer::new(page_size));
123        if layers.wire.is_none() {
124            layers.wire = Some(outer);
125        } else {
126            layers.mid.insert(0, outer);
127        }
128    }
129
130    pub(crate) fn with_read_src<F, R>(&self, io: &IoRef, f: F) -> R
131    where
132        F: FnOnce(&mut BytesMut) -> R,
133    {
134        self.borrow().transport().with_read(io, f)
135    }
136
137    pub(crate) fn with_read_dst<F, R>(&self, io: &IoRef, f: F) -> R
138    where
139        F: FnOnce(&mut BytesMut) -> R,
140    {
141        self.borrow().app.with_read(io, f)
142    }
143
144    pub(crate) fn write_buf_size(&self) -> usize {
145        // check size for first level because delayed filter processing
146        let layers = self.borrow();
147        if let Some(wire) = &layers.wire {
148            layers.app.write_len() + wire.write_len()
149        } else {
150            layers.app.write_len()
151        }
152    }
153
154    /// Returns the size of the transport-facing write buffer.
155    pub(crate) fn write_dst_size(&self) -> usize {
156        self.borrow().write_dst_size()
157    }
158
159    pub(crate) fn with_write_src<F, R>(&self, f: F) -> R
160    where
161        F: FnOnce(&mut BytePages) -> R,
162    {
163        self.borrow().app.with_write(f)
164    }
165
166    pub(crate) fn with_write_dst<F, R>(&self, f: F) -> R
167    where
168        F: FnOnce(&mut BytePages) -> R,
169    {
170        self.borrow().transport().with_write(f)
171    }
172
173    pub(crate) fn read_dst_size(&self) -> usize {
174        self.borrow().app.read_len()
175    }
176
177    pub(crate) fn with_filter<F, R>(&self, io: &IoRef, f: F) -> R
178    where
179        F: FnOnce(&mut FilterCtx<'_>) -> R,
180    {
181        let layers = self.borrow();
182        let mut ctx = FilterCtx {
183            io,
184            idx: 0,
185            layers: &layers,
186            st: FilterUpdates { wants_write: false },
187        };
188        f(&mut ctx)
189    }
190
191    pub(crate) fn get_read_buf(&self) -> Option<BytesMut> {
192        self.borrow().transport().read.take()
193    }
194
195    pub(crate) fn set_read_buf(&self, buf: BytesMut) {
196        let layers = self.borrow();
197        let buffer = layers.transport();
198        if let Some(mut first_buf) = buffer.read.take() {
199            first_buf.extend_from_slice(&buf);
200            buffer.read.set(Some(first_buf));
201        } else if !buf.is_empty() {
202            buffer.read.set(Some(buf));
203        }
204    }
205
206    pub(crate) fn process_read_buf(&self, io: &IoRef) -> io::Result<FilterUpdates> {
207        let layers = self.borrow();
208        let mut ctx = FilterCtx {
209            io,
210            idx: 0,
211            layers: &layers,
212            st: FilterUpdates { wants_write: false },
213        };
214        io.with_callbacks(|cb| cb.before_processing(io));
215        let result = io.filter().process_read_buf(&mut ctx);
216        io.with_callbacks(|cb| cb.after_processing(io));
217
218        result.map(|()| ctx.st)
219    }
220
221    pub(crate) fn process_read_buf_no_cb(&self, io: &IoRef) -> io::Result<FilterUpdates> {
222        let layers = self.borrow();
223        let mut ctx = FilterCtx {
224            io,
225            idx: 0,
226            layers: &layers,
227            st: FilterUpdates { wants_write: false },
228        };
229        io.filter().process_read_buf(&mut ctx).map(|()| ctx.st)
230    }
231
232    pub(crate) fn process_write_buf(&self, io: &IoRef) -> io::Result<()> {
233        let layers = self.borrow();
234        if layers.app.is_write_empty() {
235            Ok(())
236        } else {
237            let mut ctx = FilterCtx {
238                io,
239                idx: 0,
240                layers: &layers,
241                st: FilterUpdates { wants_write: true },
242            };
243            io.with_callbacks(|cb| cb.before_processing(io));
244            let res = io.filter().process_write_buf(&mut ctx);
245            io.with_callbacks(|cb| cb.after_processing(io));
246
247            res
248        }
249    }
250
251    pub(crate) fn process_write_buf_no_cb(&self, io: &IoRef) -> io::Result<()> {
252        let layers = self.borrow();
253        if layers.app.is_write_empty() {
254            Ok(())
255        } else {
256            let mut ctx = FilterCtx {
257                io,
258                idx: 0,
259                layers: &layers,
260                st: FilterUpdates { wants_write: true },
261            };
262            io.filter().process_write_buf(&mut ctx)
263        }
264    }
265
266    pub(crate) fn process_write_buf_force(&self, io: &IoRef) -> io::Result<()> {
267        let layers = self.borrow();
268        let mut ctx = FilterCtx {
269            io,
270            idx: 0,
271            layers: &layers,
272            st: FilterUpdates { wants_write: true },
273        };
274        io.with_callbacks(|cb| cb.before_processing(io));
275        let res = io.filter().process_write_buf(&mut ctx);
276        io.with_callbacks(|cb| cb.after_processing(io));
277
278        res
279    }
280
281    pub(crate) fn process_shutdown(&self, io: &IoRef) -> io::Result<Poll<()>> {
282        self.process_write_buf(io)?;
283        io.with_callbacks(|cb| cb.before_processing(io));
284        let res = self.with_filter(io, |ctx| io.filter().shutdown(ctx));
285        io.with_callbacks(|cb| cb.after_processing(io));
286
287        res
288    }
289
290    /// Releases the data of every buffer once nothing can consume it anymore.
291    ///
292    /// Read buffers go back to the cache and write pages are freed, the
293    /// buffers themselves stay usable.
294    pub(crate) fn release(&self) {
295        for b in self.borrow().buffers() {
296            drop(b.read.take());
297            b.with_write_if_free(BytePages::clear);
298        }
299    }
300}
301
302impl Buffer {
303    fn new(size: BytePageSize) -> Self {
304        Buffer {
305            read: Cell::new(None),
306            write: RefCell::new(BytePages::new(size)),
307        }
308    }
309
310    /// Calls `f` unless the write buffer is in use.
311    fn with_write_if_free(&self, f: impl FnOnce(&mut BytePages)) {
312        if let Ok(mut wb) = self.write.try_borrow_mut() {
313            f(&mut wb);
314        }
315    }
316
317    fn is_write_empty(&self) -> bool {
318        self.write.borrow().is_empty()
319    }
320
321    fn read_len(&self) -> usize {
322        if let Some(rb) = self.read.take() {
323            let l = rb.len();
324            self.read.set(Some(rb));
325            l
326        } else {
327            0
328        }
329    }
330
331    fn write_len(&self) -> usize {
332        self.write.borrow().len()
333    }
334
335    fn with_read<F, R>(&self, io: &IoRef, f: F) -> R
336    where
337        F: FnOnce(&mut BytesMut) -> R,
338    {
339        let mut rb = self.read.take().unwrap_or_else(|| io.0.get_read_buf());
340        let result = f(&mut rb);
341
342        #[cfg(debug_assertions)]
343        // check nested updates
344        if self.read.take().is_some() {
345            log::error!("Nested read io operation is detected");
346            io.terminate();
347        }
348
349        if !rb.is_empty() {
350            self.read.set(Some(rb));
351        }
352        result
353    }
354
355    fn with_write<F, R>(&self, f: F) -> R
356    where
357        F: FnOnce(&mut BytePages) -> R,
358    {
359        f(&mut self.write.borrow_mut())
360    }
361}
362
363#[derive(Copy, Clone, Debug)]
364pub(crate) struct FilterUpdates {
365    pub(crate) wants_write: bool,
366}
367
368#[derive(Debug)]
369/// Context used while traversing a complete filter chain.
370///
371/// A context tracks the current layer and the write activity accumulated while
372/// traversing it. [`with_next`](Self::with_next) advances to the inner layer,
373/// while [`with_buffer`](Self::with_buffer) exposes the buffers adjacent to the
374/// current layer.
375pub struct FilterCtx<'a> {
376    io: &'a IoRef,
377    idx: usize,
378    layers: &'a Layers,
379    st: FilterUpdates,
380}
381
382impl FilterCtx<'_> {
383    #[inline]
384    /// Gets a reference to the I/O object.
385    pub fn io(&self) -> &IoRef {
386        self.io
387    }
388
389    #[inline]
390    /// Gets the I/O tag.
391    pub fn tag(&self) -> &'static str {
392        self.io.tag()
393    }
394
395    #[inline]
396    /// Invokes `f` with the context advanced to the next inner filter.
397    ///
398    /// The previous layer is restored after `f` returns.
399    pub fn with_next<F, R>(&mut self, f: F) -> R
400    where
401        F: FnOnce(&mut Self) -> R,
402    {
403        self.idx += 1;
404        let res = f(self);
405        self.idx -= 1;
406        res
407    }
408
409    #[inline]
410    /// Invokes `f` with the buffers adjacent to the current filter.
411    pub fn with_buffer<F, R>(&mut self, f: F) -> R
412    where
413        F: FnOnce(&mut FilterBuf<'_>) -> R,
414    {
415        let mut buf = FilterBuf {
416            io: self.io,
417            curr: self.buffer(),
418            next: self.layers.get(self.idx + 1),
419            wants_write: Cell::new(self.st.wants_write),
420        };
421        let result = f(&mut buf);
422        if buf.wants_write.get() {
423            self.st.wants_write = true;
424        }
425        result
426    }
427
428    #[inline]
429    /// Returns the size of the application-facing read buffer.
430    pub fn read_dst_size(&self) -> usize {
431        self.layers.app.read_len()
432    }
433
434    #[inline]
435    /// Returns the size of the transport-facing write buffer.
436    pub fn write_dst_size(&mut self) -> usize {
437        self.layers.write_dst_size()
438    }
439
440    pub(crate) fn clear_write_buf(&mut self) {
441        self.buffer().with_write(BytePages::clear);
442    }
443
444    fn buffer(&self) -> &Buffer {
445        self.layers
446            .get(self.idx)
447            .expect("Filter context is outside of the filter chain")
448    }
449}
450
451#[derive(Debug)]
452/// Buffers and connection state adjacent to one [`FilterLayer`](crate::FilterLayer).
453///
454/// For reads, the source is transport-facing and the destination is
455/// application-facing. For writes, the source is application-facing and the
456/// destination is transport-facing. Buffers are returned to the chain after
457/// each closure completes; empty read buffers may be returned to the
458/// configured cache.
459///
460/// Buffer access is not reentrant. A closure must not access an overlapping
461/// buffer again through this object or its [`IoRef`]. Nested write access
462/// panics; nested read access is unsupported and may terminate the connection
463/// or lose nested buffer changes.
464pub struct FilterBuf<'a> {
465    io: &'a IoRef,
466    curr: &'a Buffer,
467    // `None` below the innermost layer, the transport uses `curr` directly
468    next: Option<&'a Buffer>,
469    wants_write: Cell<bool>,
470}
471
472impl FilterBuf<'_> {
473    #[inline]
474    /// Gets a reference to the I/O object.
475    pub fn io(&self) -> &IoRef {
476        self.io
477    }
478
479    #[inline]
480    /// Gets the I/O tag.
481    pub fn tag(&self) -> &'static str {
482        self.io.tag()
483    }
484
485    /// Provides mutable access to the transport-facing read source.
486    ///
487    /// The source is optional because no bytes may currently be allocated for
488    /// this edge of the filter chain. Leaving an empty buffer in the option
489    /// returns it to the configured cache.
490    ///
491    /// The closure must not access this read source again.
492    pub fn with_read_src<F, R>(&self, f: F) -> R
493    where
494        F: FnOnce(&mut Option<BytesMut>) -> R,
495    {
496        let mut read_src = self.next.and_then(|b| b.read.take());
497        let result = f(&mut read_src);
498        self.put_read_src(read_src);
499        result
500    }
501
502    /// Provides the transport-facing read source and application-facing
503    /// destination.
504    ///
505    /// Implementations normally consume bytes from `src` and append decoded or
506    /// transformed bytes to `dst`. Unconsumed source bytes are retained for the
507    /// next invocation.
508    ///
509    /// The closure must not access either read buffer again.
510    pub fn with_read_buffers<F, R>(&self, f: F) -> R
511    where
512        F: FnOnce(&mut Option<BytesMut>, &mut BytesMut) -> R,
513    {
514        let mut read_src = self.next.and_then(|b| b.read.take());
515        let mut read_dst = self
516            .curr
517            .read
518            .take()
519            .unwrap_or_else(|| self.io.0.get_read_buf());
520
521        let result = f(&mut read_src, &mut read_dst);
522
523        self.put_read_src(read_src);
524        if !read_dst.is_empty() {
525            self.curr.read.set(Some(read_dst));
526        }
527
528        result
529    }
530
531    fn put_read_src(&self, src: Option<BytesMut>) {
532        if let Some(b) = src
533            && !b.is_empty()
534        {
535            // Without a filter layer there is no transport-facing read
536            // source, the transport reads into the application-facing
537            // buffer. Input stored here is never read.
538            debug_assert!(
539                self.next.is_some(),
540                "{}: input stored in the read source of the innermost filter buffer is never read",
541                self.io.tag()
542            );
543            if let Some(next) = self.next {
544                next.read.set(Some(b));
545            }
546        }
547    }
548
549    #[inline]
550    /// Provides the application-facing write source and transport-facing
551    /// destination.
552    ///
553    /// Implementations normally consume bytes from `src` and append encoded or
554    /// transformed bytes to `dst`. Appending destination bytes marks the write
555    /// chain as needing transport progress.
556    ///
557    /// # Panics
558    ///
559    /// Panics if the closure accesses either write buffer again.
560    pub fn with_write_buffers<F, R>(&self, f: F) -> R
561    where
562        F: FnOnce(&mut BytePages, &mut BytePages) -> R,
563    {
564        let mut write_curr = self.curr.write.borrow_mut();
565        let mut next = self.next.map(|b| b.write.borrow_mut());
566        let mut on_demand;
567        let write_next = if let Some(next) = next.as_deref_mut() {
568            next
569        } else {
570            // Without a filter layer there is no transport-facing write
571            // buffer, the transport writes from the application-facing one.
572            // Output written here is never delivered.
573            on_demand = BytePages::new(write_curr.page_size());
574            &mut on_demand
575        };
576        let write_len = if self.wants_write.get() {
577            0
578        } else {
579            write_next.len()
580        };
581
582        let result = f(&mut write_curr, write_next);
583
584        if !self.wants_write.get() && write_next.len() > write_len {
585            self.wants_write.set(true);
586        }
587        debug_assert!(
588            self.next.is_some() || write_next.is_empty(),
589            "{}: output written to the write destination of the innermost filter buffer is never sent",
590            self.io.tag()
591        );
592        result
593    }
594}
595
596impl fmt::Debug for Buffer {
597    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
598        let read = self.read.take();
599        let write = self.write.try_borrow();
600
601        let result = f
602            .debug_struct("Buffer")
603            .field("read", &read)
604            .field("write", &write.as_deref().ok())
605            .finish();
606        self.read.set(read);
607        result
608    }
609}
610
611#[cfg(test)]
612mod tests {
613    use std::ptr;
614
615    use ntex_bytes::BufMut;
616
617    use super::*;
618    use crate::{Io, testing::IoTest};
619
620    #[test]
621    fn miri_add_layer_keeps_buffers() {
622        let stack = Stack::new(BytePageSize::Size8);
623        for i in 0..4u8 {
624            let layers = stack.borrow();
625            let nested = stack.borrow();
626            let mut buf = BytesMut::new();
627            buf.extend_from_slice(&[i]);
628            layers.app.read.set(Some(buf));
629            drop(nested);
630            assert!(stack.is_borrowed());
631            drop(layers);
632            stack.add_layer(BytePageSize::Size8);
633        }
634
635        let layers = stack.borrow();
636        assert_eq!(layers.count(), 4);
637        let reads: Vec<_> = layers
638            .buffers()
639            .map(|b| b.read.take().map(|b| b.to_vec()))
640            .collect();
641        assert_eq!(
642            reads,
643            [
644                None,
645                Some(vec![3]),
646                Some(vec![2]),
647                Some(vec![1]),
648                Some(vec![0])
649            ]
650        );
651    }
652
653    #[test]
654    #[should_panic(expected = "filter buffers are in use")]
655    fn miri_add_layer_while_borrowed() {
656        let stack = Stack::new(BytePageSize::Size8);
657        let _layers = stack.borrow();
658        stack.add_layer(BytePageSize::Size8);
659    }
660
661    #[ntex::test]
662    async fn stack_without_layers() {
663        let (_, server) = IoTest::create();
664        let io = Io::from(server);
665        let ioref = io.get_ref();
666
667        let stack = Stack::new(BytePageSize::Size8);
668        let layers = stack.borrow();
669        assert_eq!(layers.count(), 0);
670        assert!(format!("{stack:?}").contains("layers: 0"));
671        assert!(layers.mid.is_empty() && layers.wire.is_none());
672        assert!(ptr::eq(layers.transport(), &raw const layers.app));
673        assert!(layers.get(1).is_none());
674        assert_eq!(stack.read_dst_size(), 0);
675        assert_eq!(stack.write_buf_size(), 0);
676
677        // the application and the transport share one buffer
678        stack.with_write_src(|buf| buf.put_slice(b"out"));
679        assert_eq!(stack.write_buf_size(), 3);
680        assert_eq!(stack.write_dst_size(), 3);
681        stack.set_read_buf(BytesMut::from(&b"one"[..]));
682        stack.set_read_buf(BytesMut::from(&b"-two"[..]));
683        assert_eq!(stack.read_dst_size(), 7);
684        stack.with_read_dst(&ioref, |buf| assert_eq!(&buf[..], b"one-two"));
685
686        // the base position has no inner buffer, nothing is kept for it
687        stack.with_filter(&ioref, |ctx| {
688            assert_eq!(ctx.read_dst_size(), 7);
689            assert_eq!(ctx.write_dst_size(), 3);
690            ctx.with_buffer(|buf| {
691                assert!(ptr::eq(buf.curr, &raw const layers.app));
692                assert!(buf.next.is_none());
693                buf.with_read_buffers(|src, dst| {
694                    assert!(src.is_none());
695                    assert_eq!(&dst[..], b"one-two");
696                });
697                buf.with_read_src(|src| *src = Some(BytesMut::new()));
698                buf.with_write_buffers(|src, dst| {
699                    assert_eq!(src.len(), 3);
700                    assert!(dst.is_empty());
701                });
702            });
703        });
704        assert!(layers.wire.is_none());
705
706        assert_eq!(
707            stack.with_write_dst(|buf| buf.split_to(3).freeze()),
708            b"out".as_ref()
709        );
710        assert_eq!(stack.get_read_buf().as_deref(), Some(b"one-two".as_ref()));
711        assert!(stack.get_read_buf().is_none());
712        stack.set_read_buf(BytesMut::new());
713        assert!(stack.get_read_buf().is_none());
714
715        stack.set_page_size(BytePageSize::Size32);
716        layers
717            .app
718            .with_write(|buf| assert_eq!(buf.page_size(), BytePageSize::Size32));
719    }
720
721    #[ntex::test]
722    async fn stack_with_one_layer() {
723        let (_, server) = IoTest::create();
724        let io = Io::from(server);
725        let ioref = io.get_ref();
726
727        // data buffered before the layer is added belongs to the transport
728        let stack = Stack::new(BytePageSize::Size8);
729        stack.with_write_src(|buf| buf.put_slice(b"plain"));
730        stack.set_read_buf(BytesMut::from(&b"hello"[..]));
731        stack.add_layer(BytePageSize::Size16);
732
733        let layers = stack.borrow();
734        assert_eq!(layers.count(), 1);
735        assert!(format!("{stack:?}").contains("layers: 1"));
736        assert!(layers.mid.is_empty());
737        let wire = layers.wire.as_ref().unwrap();
738        assert!(ptr::eq(layers.transport(), wire));
739        assert!(ptr::eq(layers.get(1).unwrap(), wire));
740        assert!(layers.get(2).is_none());
741        assert_eq!(stack.read_dst_size(), 0);
742        assert_eq!(stack.write_dst_size(), 5);
743        layers.app.with_write(|buf| {
744            assert!(buf.is_empty());
745            assert_eq!(buf.page_size(), BytePageSize::Size16);
746        });
747
748        // output of both sides is pending, filter processing may be delayed
749        stack.with_write_src(|buf| buf.put_slice(b"app"));
750        assert_eq!(stack.write_buf_size(), 8);
751
752        stack.with_filter(&ioref, |ctx| {
753            ctx.with_buffer(|buf| {
754                assert!(ptr::eq(buf.curr, &raw const layers.app));
755                buf.with_read_buffers(|src, dst| {
756                    dst.extend_from_slice(&src.take().unwrap());
757                });
758                buf.with_write_buffers(|src, dst| {
759                    assert_eq!(dst.len(), 5);
760                    src.move_to(dst);
761                });
762            });
763            ctx.with_next(|ctx| {
764                ctx.with_buffer(|buf| {
765                    assert!(ptr::eq(buf.curr, wire));
766                    assert!(buf.next.is_none());
767                });
768            });
769        });
770
771        assert_eq!(stack.read_dst_size(), 5);
772        assert!(stack.get_read_buf().is_none());
773        stack.with_read_dst(&ioref, |buf| assert_eq!(&buf[..], b"hello"));
774        assert_eq!(stack.write_buf_size(), 8);
775        assert_eq!(
776            stack.with_write_dst(|buf| buf.split_to(8).freeze()),
777            b"plainapp".as_ref()
778        );
779        assert_eq!(stack.write_buf_size(), 0);
780    }
781
782    #[ntex::test]
783    async fn stack_with_more_layers() {
784        type Seen = (Option<Vec<u8>>, Vec<u8>, usize, usize);
785
786        fn visit(ctx: &mut FilterCtx<'_>, seen: &mut Vec<Seen>) {
787            ctx.with_buffer(|buf| {
788                let (src, dst) = buf.with_read_buffers(|src, dst| {
789                    (src.as_deref().map(<[u8]>::to_vec), dst.to_vec())
790                });
791                let (wsrc, wdst) = buf.with_write_buffers(|src, dst| (src.len(), dst.len()));
792                seen.push((src, dst, wsrc, wdst));
793            });
794            if seen.len() < 4 {
795                ctx.with_next(|ctx| visit(ctx, seen));
796            }
797        }
798
799        let (_, server) = IoTest::create();
800        let io = Io::from(server);
801        let ioref = io.get_ref();
802
803        // each layer is added outermost, buffered data moves toward the
804        // transport, the data of the first layer ends up in `wire`
805        let stack = Stack::new(BytePageSize::Size8);
806        for data in ["a", "bb", "cccc"] {
807            stack.with_write_src(|buf| buf.put_slice(data.as_bytes()));
808            stack
809                .borrow()
810                .app
811                .read
812                .set(Some(BytesMut::from(data.as_bytes())));
813            stack.add_layer(BytePageSize::Size16);
814        }
815        let layers = stack.borrow();
816        assert_eq!(layers.count(), 3);
817        assert!(format!("{stack:?}").contains("layers: 3"));
818        assert_eq!(layers.mid.len(), 2);
819        assert!(ptr::eq(layers.transport(), layers.wire.as_ref().unwrap()));
820        assert!(layers.get(4).is_none());
821
822        // only the outermost and the transport-facing buffers are pending
823        stack.with_write_src(|buf| buf.put_slice(b"app"));
824        assert_eq!(stack.write_buf_size(), 4);
825        assert_eq!(stack.write_dst_size(), 1);
826        assert_eq!(stack.read_dst_size(), 0);
827
828        let mut seen = Vec::new();
829        stack.with_filter(&ioref, |ctx| visit(ctx, &mut seen));
830        assert_eq!(
831            seen,
832            [
833                (Some(b"cccc".to_vec()), Vec::new(), 3, 4),
834                (Some(b"bb".to_vec()), b"cccc".to_vec(), 4, 2),
835                (Some(b"a".to_vec()), b"bb".to_vec(), 2, 1),
836                (None, b"a".to_vec(), 1, 0),
837            ]
838        );
839        assert_eq!(stack.get_read_buf().as_deref(), Some(b"a".as_ref()));
840
841        stack.set_page_size(BytePageSize::Size32);
842        for buf in layers.buffers() {
843            buf.with_write(|buf| assert_eq!(buf.page_size(), BytePageSize::Size32));
844        }
845
846        stack.release();
847        assert_eq!(layers.buffers().count(), 4);
848        for buf in layers.buffers() {
849            assert_eq!(buf.read_len(), 0);
850            assert_eq!(buf.write_len(), 0);
851        }
852    }
853
854    #[ntex::test]
855    #[should_panic(expected = "outside of the filter chain")]
856    async fn filter_ctx_outside_of_chain() {
857        let (_, server) = IoTest::create();
858        let io = Io::from(server);
859        let ioref = io.get_ref();
860        let stack = Stack::new(BytePageSize::Size8);
861        stack.add_layer(BytePageSize::Size8);
862
863        stack.with_filter(&ioref, |ctx| {
864            ctx.with_next(|ctx| ctx.with_next(|ctx| ctx.with_buffer(|_| ())));
865        });
866    }
867
868    #[ntex::test]
869    async fn set_read_buf_merges_into_cacheable_buffer() {
870        let high = BytePageSize::Size16.capacity();
871        let stack = Stack::new(BytePageSize::Size8);
872
873        // unconsumed input, most of the buffer is taken by a decoded frame
874        // that is still alive
875        let mut first = BytesMut::with_page_size(BytePageSize::Size16);
876        first.extend_from_slice(&vec![1; high - 100]);
877        let frame = first.split_to(high - 1100);
878        stack.set_read_buf(first);
879
880        // a read into a buffer of its own completes
881        let mut second = BytesMut::with_page_size(BytePageSize::Size16);
882        second.extend_from_slice(&[2; 4000]);
883        stack.set_read_buf(second);
884
885        let merged = stack.get_read_buf().unwrap();
886        assert_eq!(merged.len(), 5000);
887        assert_eq!(&merged[..1000], &[1; 1000][..]);
888        assert_eq!(&merged[1000..], &[2; 4000][..]);
889        assert_eq!(merged.capacity(), high);
890        assert_eq!(frame.len(), high - 1100);
891    }
892
893    #[ntex::test]
894    async fn filter_read_buffers() {
895        let (_, server) = IoTest::create();
896        let io = Io::from(server);
897        let ioref = io.get_ref();
898        let stack = Stack::new(BytePageSize::Size8);
899        stack.add_layer(BytePageSize::Size8);
900        stack.set_read_buf(BytesMut::from(&b"input"[..]));
901
902        stack.with_filter(&ioref, |ctx| {
903            assert_eq!(ctx.io(), &ioref);
904            assert_eq!(ctx.tag(), ioref.tag());
905            assert_eq!(ctx.read_dst_size(), 0);
906
907            ctx.with_buffer(|buf| {
908                assert_eq!(buf.io(), &ioref);
909                assert_eq!(buf.tag(), ioref.tag());
910                buf.with_read_buffers(|src, dst| {
911                    let src = src.as_mut().unwrap();
912                    dst.extend_from_slice(&src.split_to(2));
913                });
914            });
915        });
916
917        assert_eq!(stack.read_dst_size(), 2);
918        stack.with_read_dst(&ioref, |buf| assert_eq!(&buf[..], b"in"));
919        assert_eq!(stack.get_read_buf().as_deref(), Some(b"put".as_ref()));
920
921        stack.with_filter(&ioref, |ctx| {
922            ctx.with_buffer(|buf| {
923                buf.with_read_src(|src| {
924                    *src = Some(BytesMut::from(&b"next"[..]));
925                });
926            });
927        });
928        assert_eq!(stack.get_read_buf().as_deref(), Some(b"next".as_ref()));
929    }
930
931    #[cfg(debug_assertions)]
932    #[ntex::test]
933    #[should_panic(expected = "is never sent")]
934    async fn innermost_write_destination_output_asserts() {
935        let (_, server) = IoTest::create();
936        let io = Io::from(server);
937        let ioref = io.get_ref();
938        let stack = Stack::new(BytePageSize::Size8);
939
940        stack.with_write_src(|buf| buf.put_slice(b"out"));
941        stack.with_filter(&ioref, |ctx| {
942            ctx.with_buffer(|buf| buf.with_write_buffers(BytePages::move_to));
943        });
944    }
945
946    #[cfg(debug_assertions)]
947    #[ntex::test]
948    #[should_panic(expected = "is never read")]
949    async fn innermost_read_source_input_asserts() {
950        let (_, server) = IoTest::create();
951        let io = Io::from(server);
952        let ioref = io.get_ref();
953        let stack = Stack::new(BytePageSize::Size8);
954
955        stack.with_filter(&ioref, |ctx| {
956            ctx.with_buffer(|buf| {
957                buf.with_read_src(|src| *src = Some(BytesMut::from(&b"in"[..])));
958            });
959        });
960    }
961
962    #[ntex::test]
963    async fn filter_write_buffers_and_updates() {
964        let (_, server) = IoTest::create();
965        let io = Io::from(server);
966        let ioref = io.get_ref();
967        let stack = Stack::new(BytePageSize::Size8);
968        stack.add_layer(BytePageSize::Size8);
969        stack.with_write_src(|buf| buf.put_slice(b"output"));
970
971        let updates = stack.with_filter(&ioref, |ctx| {
972            assert_eq!(ctx.write_dst_size(), 0);
973            ctx.with_buffer(|buf| {
974                buf.with_write_buffers(|src, dst| {
975                    assert_eq!(src.len(), 6);
976                    assert_eq!(dst.len(), 0);
977                    src.move_to(dst);
978                });
979            });
980            ctx.st
981        });
982
983        assert!(updates.wants_write);
984        assert_eq!(stack.write_buf_size(), 6);
985        assert_eq!(
986            stack.with_write_dst(|buf| buf.split_to(6).freeze()),
987            b"output".as_ref()
988        );
989        assert_eq!(stack.write_buf_size(), 0);
990    }
991
992    #[ntex::test]
993    async fn buffer_debug_preserves_contents() {
994        let (_, server) = IoTest::create();
995        let io = Io::from(server);
996        let ioref = io.get_ref();
997        let buffer = Buffer::new(BytePageSize::Size8);
998
999        buffer.with_read(&ioref, |buf| buf.extend_from_slice(b"read"));
1000        buffer.with_write(|buf| buf.put_slice(b"write"));
1001
1002        let debug = format!("{buffer:?}");
1003        assert!(debug.contains("Buffer"));
1004        assert_eq!(buffer.read_len(), 4);
1005        assert_eq!(buffer.write_len(), 5);
1006
1007        // a write buffer in use is not shown
1008        let debug = buffer.with_write(|_| format!("{buffer:?}"));
1009        assert!(debug.contains("write: None"));
1010        assert_eq!(buffer.read_len(), 4);
1011    }
1012}