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
7pub(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 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 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 fn transport(&self) -> &Buffer {
77 self.wire.as_ref().unwrap_or(&self.app)
78 }
79
80 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 #[inline]
98 pub(crate) fn borrow(&self) -> Ref<'_, Layers> {
99 self.0.borrow()
100 }
101
102 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 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 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 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 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 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 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)]
369pub struct FilterCtx<'a> {
376 io: &'a IoRef,
377 idx: usize,
378 layers: &'a Layers,
379 st: FilterUpdates,
380}
381
382impl FilterCtx<'_> {
383 #[inline]
384 pub fn io(&self) -> &IoRef {
386 self.io
387 }
388
389 #[inline]
390 pub fn tag(&self) -> &'static str {
392 self.io.tag()
393 }
394
395 #[inline]
396 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 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 pub fn read_dst_size(&self) -> usize {
431 self.layers.app.read_len()
432 }
433
434 #[inline]
435 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)]
452pub struct FilterBuf<'a> {
465 io: &'a IoRef,
466 curr: &'a Buffer,
467 next: Option<&'a Buffer>,
469 wants_write: Cell<bool>,
470}
471
472impl FilterBuf<'_> {
473 #[inline]
474 pub fn io(&self) -> &IoRef {
476 self.io
477 }
478
479 #[inline]
480 pub fn tag(&self) -> &'static str {
482 self.io.tag()
483 }
484
485 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 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 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 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 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 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 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 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 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 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 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 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 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 let debug = buffer.with_write(|_| format!("{buffer:?}"));
1009 assert!(debug.contains("write: None"));
1010 assert_eq!(buffer.read_len(), 4);
1011 }
1012}