opendal_core/raw/oio/buf/queue_buf.rs
1// Licensed to the Apache Software Foundation (ASF) under one
2// or more contributor license agreements. See the NOTICE file
3// distributed with this work for additional information
4// regarding copyright ownership. The ASF licenses this file
5// to you under the Apache License, Version 2.0 (the
6// "License"); you may not use this file except in compliance
7// with the License. You may obtain a copy of the License at
8//
9// http://www.apache.org/licenses/LICENSE-2.0
10//
11// Unless required by applicable law or agreed to in writing,
12// software distributed under the License is distributed on an
13// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
14// KIND, either express or implied. See the License for the
15// specific language governing permissions and limitations
16// under the License.
17
18use std::collections::VecDeque;
19use std::mem;
20
21use bytes::Buf;
22
23use crate::*;
24
25/// QueueBuf is a queue of [`Buffer`].
26///
27/// QueueBuf allows:
28///
29/// - Storing multiple buffers without copying bytes.
30/// - Consuming buffers in order.
31///
32/// QueueBuf mainly provides the following operations:
33///
34/// - `push`: Push a new buffer in the queue.
35/// - `collect`: Collect all buffer in the queue as a new [`Buffer`]
36/// - `advance`: Advance the queue by `cnt` bytes.
37#[derive(Clone, Default)]
38pub struct QueueBuf(VecDeque<Buffer>);
39
40impl QueueBuf {
41 /// Create a new buffer queue.
42 #[inline]
43 pub fn new() -> Self {
44 Self::default()
45 }
46
47 /// Push new [`Buffer`] into the queue.
48 #[inline]
49 pub fn push(&mut self, buf: Buffer) {
50 if buf.is_empty() {
51 return;
52 }
53
54 self.0.push_back(buf);
55 }
56
57 /// Total bytes size inside the buffer queue.
58 #[inline]
59 pub fn len(&self) -> usize {
60 self.0.iter().map(|b| b.len()).sum()
61 }
62
63 /// Is the buffer queue empty.
64 #[inline]
65 pub fn is_empty(&self) -> bool {
66 self.len() == 0
67 }
68
69 /// Take the entire buffer queue and leave `self` in empty states.
70 #[inline]
71 pub fn take(&mut self) -> QueueBuf {
72 mem::take(self)
73 }
74
75 /// Build a new [`Buffer`] from the queue.
76 ///
77 /// If the queue is empty, it will return an empty buffer. Otherwise, it will iterate over all
78 /// buffers and collect them into a new buffer.
79 ///
80 /// # Notes
81 ///
82 /// There are allocation overheads when collecting multiple buffers into a new buffer. But
83 /// most of them should be acceptable since we can expect the item length of buffers are slower
84 /// than 4k.
85 #[inline]
86 pub fn collect(mut self) -> Buffer {
87 if self.0.is_empty() {
88 Buffer::new()
89 } else if self.0.len() == 1 {
90 self.0.pop_front().unwrap()
91 } else {
92 self.0.into_iter().flatten().collect()
93 }
94 }
95
96 /// Advance the buffer queue by `cnt` bytes.
97 #[inline]
98 pub fn advance(&mut self, cnt: usize) {
99 assert!(cnt <= self.len(), "cannot advance past {cnt} bytes");
100
101 let mut new_cnt = cnt;
102 while new_cnt > 0 {
103 let buf = self.0.front_mut().expect("buffer must be valid");
104 if new_cnt < buf.remaining() {
105 buf.advance(new_cnt);
106 break;
107 } else {
108 new_cnt -= buf.remaining();
109 self.0.pop_front();
110 }
111 }
112 }
113
114 /// Clear the buffer queue.
115 #[inline]
116 pub fn clear(&mut self) {
117 self.0.clear()
118 }
119}