Skip to content

Commit a8a6594

Browse files
committed
BytesCursor
1 parent 75e4f24 commit a8a6594

2 files changed

Lines changed: 149 additions & 72 deletions

File tree

Lines changed: 110 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,110 @@
1+
use crate::transport::Error;
2+
use std::{io, mem};
3+
4+
#[derive(Debug, Default)]
5+
pub(crate) struct BytesCursor<'a> {
6+
slice: &'a mut [u8],
7+
/// position <= slice.len()
8+
position: usize,
9+
}
10+
11+
impl<'a> BytesCursor<'a> {
12+
pub(crate) fn new(slice: &'a mut [u8], position: usize) -> Self {
13+
assert!(position <= slice.len());
14+
BytesCursor { slice, position }
15+
}
16+
17+
pub(crate) fn slice_mut(&mut self) -> &mut [u8] {
18+
self.slice
19+
}
20+
21+
pub(crate) fn position(&self) -> usize {
22+
self.position
23+
}
24+
25+
pub(crate) fn available_bytes(&self) -> usize {
26+
self.slice.len() - self.position
27+
}
28+
29+
pub(crate) fn check_available_space(&self, sz: usize) -> io::Result<()> {
30+
if sz > self.available_bytes() {
31+
Err(io::Error::new(
32+
io::ErrorKind::InvalidData,
33+
format!(
34+
"data out of range, available {} requested {}",
35+
self.available_bytes(),
36+
sz
37+
),
38+
))
39+
} else {
40+
Ok(())
41+
}
42+
}
43+
44+
pub(crate) fn available_slice(&mut self) -> &mut [u8] {
45+
&mut self.slice[self.position..]
46+
}
47+
48+
pub(crate) fn account_written(&mut self, count: usize) {
49+
assert!(self.available_bytes() >= count);
50+
self.position += count;
51+
}
52+
53+
pub(crate) fn written(&self) -> &[u8] {
54+
&self.slice[..self.position]
55+
}
56+
57+
pub(crate) fn split_at(&mut self, offset: usize) -> Result<BytesCursor<'a>, Error> {
58+
if self.slice.len() < offset {
59+
return Err(Error::SplitOutOfBounds(offset));
60+
}
61+
62+
let (len1, len2) = if self.position > offset {
63+
(offset, self.position - offset)
64+
} else {
65+
(self.position, 0)
66+
};
67+
let (slice1, slice2) = mem::take(&mut self.slice).split_at_mut(offset);
68+
*self = BytesCursor::new(slice1, len1);
69+
Ok(BytesCursor::new(slice2, len2))
70+
}
71+
72+
pub(crate) fn extend_from_slice(&mut self, slice: &[u8]) {
73+
self.slice[self.position..][..slice.len()].copy_from_slice(slice);
74+
self.position += slice.len();
75+
}
76+
}
77+
78+
#[cfg(test)]
79+
mod tests {
80+
use crate::transport::fusedev::bytes_cursor::BytesCursor;
81+
82+
#[test]
83+
fn test_split_at() {
84+
let mut array = [0, 1, 2, 3, 4, 5, 6, 7, 8, 9];
85+
86+
let mut b = BytesCursor::new(&mut array, 0);
87+
let mut b1 = b.split_at(4).unwrap();
88+
89+
assert_eq!(&[0, 1, 2, 3], b.slice_mut());
90+
assert_eq!(0, b.position());
91+
assert_eq!(&[4, 5, 6, 7, 8, 9], b1.slice_mut());
92+
assert_eq!(0, b1.position());
93+
94+
let mut b = BytesCursor::new(&mut array, 6);
95+
let mut b1 = b.split_at(4).unwrap();
96+
assert_eq!(&[0, 1, 2, 3], b.slice_mut());
97+
assert_eq!(4, b.position());
98+
assert_eq!(&[4, 5, 6, 7, 8, 9], b1.slice_mut());
99+
assert_eq!(2, b1.position());
100+
}
101+
102+
#[test]
103+
fn test_extend_from_slice() {
104+
let mut array = [0, 1, 2, 3, 4, 5, 6, 7, 8, 9];
105+
let mut b = BytesCursor::new(&mut array, 3);
106+
b.extend_from_slice(&[b'a', b'b']);
107+
assert_eq!(&[0, 1, 2, b'a', b'b', 5, 6, 7, 8, 9], b.slice);
108+
assert_eq!(5, b.position());
109+
}
110+
}

‎src/transport/fusedev/mod.rs‎

Lines changed: 39 additions & 72 deletions
Original file line numberDiff line numberDiff line change
@@ -10,7 +10,6 @@
1010
use std::collections::VecDeque;
1111
use std::io::{self, IoSlice, Write};
1212
use std::marker::PhantomData;
13-
use std::mem::ManuallyDrop;
1413
use std::os::fd::AsRawFd;
1514
use std::os::unix::io::RawFd;
1615

@@ -20,8 +19,11 @@ use vm_memory::{ByteValued, VolatileSlice};
2019

2120
use super::{Error, FileReadWriteVolatile, IoBuffers, Reader, Result, Writer};
2221
use crate::file_buf::FileVolatileSlice;
22+
use crate::transport::fusedev::bytes_cursor::BytesCursor;
2323
use crate::BitmapSlice;
2424

25+
mod bytes_cursor;
26+
2527
#[cfg(target_os = "linux")]
2628
mod linux_session;
2729
#[cfg(target_os = "linux")]
@@ -87,19 +89,19 @@ impl<'a, S: BitmapSlice + Default> Reader<'a, S> {
8789
pub struct FuseDevWriter<'a, S: BitmapSlice = ()> {
8890
fd: RawFd,
8991
buffered: bool,
90-
buf: ManuallyDrop<Vec<u8>>,
92+
buf: BytesCursor<'a>,
9193
bitmapslice: S,
9294
phantom: PhantomData<&'a mut [S]>,
9395
}
9496

9597
impl<'a, S: BitmapSlice + Default> FuseDevWriter<'a, S> {
9698
/// Construct a new [Writer].
9799
pub fn new(fd: RawFd, data_buf: &'a mut [u8]) -> Result<FuseDevWriter<'a, S>> {
98-
let buf = unsafe { Vec::from_raw_parts(data_buf.as_mut_ptr(), 0, data_buf.len()) };
100+
let buf = BytesCursor::new(data_buf, 0);
99101
Ok(FuseDevWriter {
100102
fd,
101103
buffered: false,
102-
buf: ManuallyDrop::new(buf),
104+
buf,
103105
bitmapslice: S::default(),
104106
phantom: PhantomData,
105107
})
@@ -113,27 +115,14 @@ impl<'a, S: BitmapSlice> FuseDevWriter<'a, S> {
113115
/// `Writer` can write up to `available_bytes() - offset` bytes. Returns an error if
114116
/// `offset > self.available_bytes()`.
115117
pub fn split_at(&mut self, offset: usize) -> Result<FuseDevWriter<'a, S>> {
116-
if self.buf.capacity() < offset {
117-
return Err(Error::SplitOutOfBounds(offset));
118-
}
119-
120-
let (len1, len2) = if self.buf.len() > offset {
121-
(offset, self.buf.len() - offset)
122-
} else {
123-
(self.buf.len(), 0)
124-
};
125-
let cap2 = self.buf.capacity() - offset;
126-
let ptr = self.buf.as_mut_ptr();
118+
let buf2 = self.buf.split_at(offset)?;
127119

128-
// Safe because both buffers refer to different parts of the same underlying `data_buf`.
129-
self.buf = unsafe { ManuallyDrop::new(Vec::from_raw_parts(ptr, len1, offset)) };
130120
self.buffered = true;
131-
let buf = unsafe { ManuallyDrop::new(Vec::from_raw_parts(ptr.add(offset), len2, cap2)) };
132121

133122
Ok(FuseDevWriter {
134123
fd: self.fd,
135124
buffered: true,
136-
buf,
125+
buf: buf2,
137126
bitmapslice: self.bitmapslice.clone(),
138127
phantom: PhantomData,
139128
})
@@ -146,15 +135,15 @@ impl<'a, S: BitmapSlice> FuseDevWriter<'a, S> {
146135
}
147136

148137
let o = match other {
149-
Some(Writer::FuseDev(w)) => w.buf.as_slice(),
138+
Some(Writer::FuseDev(w)) => w.buf.written(),
150139
_ => &[],
151140
};
152-
let res = match (self.buf.len(), o.len()) {
141+
let res = match (self.buf.position(), o.len()) {
153142
(0, 0) => Ok(0),
154143
(0, _) => write(self.fd, o),
155-
(_, 0) => write(self.fd, self.buf.as_slice()),
144+
(_, 0) => write(self.fd, self.buf.written()),
156145
(_, _) => {
157-
let bufs = [IoSlice::new(self.buf.as_slice()), IoSlice::new(o)];
146+
let bufs = [IoSlice::new(self.buf.written()), IoSlice::new(o)];
158147
writev(self.fd, &bufs)
159148
}
160149
};
@@ -164,18 +153,12 @@ impl<'a, S: BitmapSlice> FuseDevWriter<'a, S> {
164153

165154
/// Return number of bytes already written to the internal buffer.
166155
pub fn bytes_written(&self) -> usize {
167-
self.buf.len()
156+
self.buf.position()
168157
}
169158

170159
/// Return number of bytes available for writing.
171160
pub fn available_bytes(&self) -> usize {
172-
self.buf.capacity() - self.buf.len()
173-
}
174-
175-
fn account_written(&mut self, count: usize) {
176-
let new_len = self.buf.len() + count;
177-
// Safe because check_avail_space() ensures that `count` is valid.
178-
unsafe { self.buf.set_len(new_len) };
161+
self.buf.available_bytes()
179162
}
180163

181164
/// Write an object to the writer.
@@ -193,21 +176,17 @@ impl<'a, S: BitmapSlice> FuseDevWriter<'a, S> {
193176
) -> io::Result<usize> {
194177
self.check_available_space(count)?;
195178

196-
let cnt = src.read_vectored_volatile(
197-
// Safe because we have made sure buf has at least count capacity above
198-
unsafe {
199-
&[FileVolatileSlice::from_raw_ptr(
200-
self.buf.as_mut_ptr().add(self.buf.len()),
201-
count,
202-
)]
203-
},
204-
)?;
205-
self.account_written(cnt);
179+
let cnt = src.read_vectored_volatile(unsafe {
180+
&[FileVolatileSlice::from_mut_slice(
181+
&mut self.buf.available_slice()[..count],
182+
)]
183+
})?;
184+
self.buf.account_written(cnt);
206185

207186
if self.buffered {
208187
Ok(cnt)
209188
} else {
210-
Self::do_write(self.fd, &self.buf[..cnt])
189+
Self::do_write(self.fd, &self.buf.written()[..cnt])
211190
}
212191
}
213192

@@ -222,21 +201,19 @@ impl<'a, S: BitmapSlice> FuseDevWriter<'a, S> {
222201
self.check_available_space(count)?;
223202

224203
let cnt = src.read_vectored_at_volatile(
225-
// Safe because we have made sure buf has at least count capacity above
226204
unsafe {
227-
&[FileVolatileSlice::from_raw_ptr(
228-
self.buf.as_mut_ptr().add(self.buf.len()),
229-
count,
205+
&[FileVolatileSlice::from_mut_slice(
206+
&mut self.buf.available_slice()[..count],
230207
)]
231208
},
232209
off,
233210
)?;
234-
self.account_written(cnt);
211+
self.buf.account_written(cnt);
235212

236213
if self.buffered {
237214
Ok(cnt)
238215
} else {
239-
Self::do_write(self.fd, &self.buf[..cnt])
216+
Self::do_write(self.fd, &self.buf.slice_mut()[..cnt])
240217
}
241218
}
242219

@@ -266,19 +243,8 @@ impl<'a, S: BitmapSlice> FuseDevWriter<'a, S> {
266243
}
267244

268245
fn check_available_space(&self, sz: usize) -> io::Result<()> {
269-
assert!(self.buffered || self.buf.is_empty());
270-
if sz > self.available_bytes() {
271-
Err(io::Error::new(
272-
io::ErrorKind::InvalidData,
273-
format!(
274-
"data out of range, available {} requested {}",
275-
self.available_bytes(),
276-
sz
277-
),
278-
))
279-
} else {
280-
Ok(())
281-
}
246+
assert!(self.buffered || self.buf.position() == 0);
247+
self.buf.check_available_space(sz)
282248
}
283249

284250
fn do_write(fd: RawFd, data: &[u8]) -> io::Result<usize> {
@@ -298,7 +264,7 @@ impl<S: BitmapSlice> Write for FuseDevWriter<'_, S> {
298264
Ok(data.len())
299265
} else {
300266
Self::do_write(self.fd, data).inspect(|&x| {
301-
self.account_written(x);
267+
self.buf.account_written(x);
302268
})
303269
}
304270
}
@@ -319,7 +285,7 @@ impl<S: BitmapSlice> Write for FuseDevWriter<'_, S> {
319285
}
320286
writev(self.fd, bufs)
321287
.inspect(|&x| {
322-
self.account_written(x);
288+
self.buf.account_written(x);
323289
})
324290
.map_err(|e| {
325291
error! {"fail to write to fuse device on commit: {}", e};
@@ -393,7 +359,7 @@ mod async_io {
393359
} else {
394360
nix::sys::uio::pwrite(self.fd, data, 0)
395361
.map(|x| {
396-
self.account_written(x);
362+
self.buf.account_written(x);
397363
x
398364
})
399365
.map_err(|e| {
@@ -419,7 +385,7 @@ mod async_io {
419385
let bufs = [IoSlice::new(data), IoSlice::new(data2)];
420386
writev(self.fd, &bufs)
421387
.map(|x| {
422-
self.account_written(x);
388+
self.buf.account_written(x);
423389
x
424390
})
425391
.map_err(|e| {
@@ -451,7 +417,7 @@ mod async_io {
451417
let bufs = [IoSlice::new(data), IoSlice::new(data2), IoSlice::new(data3)];
452418
writev(self.fd, &bufs)
453419
.map(|x| {
454-
self.account_written(x);
420+
self.buf.account_written(x);
455421
x
456422
})
457423
.map_err(|e| {
@@ -491,11 +457,12 @@ mod async_io {
491457
) -> io::Result<usize> {
492458
self.check_available_space(count)?;
493459

494-
let buf = unsafe { FileVolatileBuf::from_raw_ptr(self.buf.as_mut_ptr(), 0, count) };
460+
let buf =
461+
unsafe { FileVolatileBuf::new_with_data(&mut self.buf.slice_mut()[..count], 0) };
495462
let (res, _) = src.async_read_at_volatile(buf, off).await;
496463
match res {
497464
Ok(cnt) => {
498-
self.account_written(cnt);
465+
self.buf.account_written(cnt);
499466
if self.buffered {
500467
Ok(cnt)
501468
} else {
@@ -515,22 +482,22 @@ mod async_io {
515482
/// We need this because the lifetime of others is usually shorter than self.
516483
pub async fn async_commit(&mut self, other: Option<&Writer<'a, S>>) -> io::Result<usize> {
517484
let o = match other {
518-
Some(Writer::FuseDev(w)) => w.buf.as_slice(),
485+
Some(Writer::FuseDev(w)) => w.buf.written(),
519486
_ => &[],
520487
};
521488

522-
let res = match (self.buf.len(), o.len()) {
489+
let res = match (self.buf.position(), o.len()) {
523490
(0, 0) => Ok(0),
524491
(0, _) => nix::sys::uio::pwrite(self.fd, o, 0).map_err(|e| {
525492
error! {"fail to write to fuse device fd {}: {}", self.fd, e};
526493
io::Error::other(format!("{}", e))
527494
}),
528-
(_, 0) => nix::sys::uio::pwrite(self.fd, self.buf.as_slice(), 0).map_err(|e| {
495+
(_, 0) => nix::sys::uio::pwrite(self.fd, self.buf.written(), 0).map_err(|e| {
529496
error! {"fail to write to fuse device fd {}: {}", self.fd, e};
530497
io::Error::other(format!("{}", e))
531498
}),
532499
(_, _) => {
533-
let bufs = [IoSlice::new(self.buf.as_slice()), IoSlice::new(o)];
500+
let bufs = [IoSlice::new(self.buf.written()), IoSlice::new(o)];
534501
writev(self.fd, &bufs).map_err(|e| {
535502
error! {"fail to write to fuse device fd {}: {}", self.fd, e};
536503
io::Error::other(format!("{}", e))

0 commit comments

Comments
 (0)