Skip to main content

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}