Compare commits

...
Author SHA1 Message Date
Alice RyhlandGitHub 38fd42acba chore: prepare bytes v1.2.0 (#556) 2022-07-19 13:39:58 +02:00
7553a67be2 Fix amortized asymptotics of BytesMut (#555)
Signed-off-by: Jiahao XU <[email protected]>
Co-authored-by: Frank Steffahn <[email protected]>
2022-07-19 13:17:53 +02:00
cd188cbd67 Add conversion from Bytes to Vec<u8> (#547)
Signed-off-by: Jiahao XU <[email protected]>
Co-authored-by: Alice Ryhl <[email protected]>
2022-07-13 09:04:23 +02:00
Jiahao XUandGitHub 10d1f6ec5c Fix: From<BytesMut> fo Vec<u8> implementation (#554)
Signed-off-by: Jiahao XU <[email protected]>
2022-07-13 08:59:54 +02:00
Jiahao XUandGitHub 068ed41bc0 Add conversion from BytesMut to Vec<u8> (#543) 2022-07-10 12:44:29 +02:00
Alice RyhlandGitHub f514bd38da miri: don't use int2ptr casts for invalid pointers (#553) 2022-07-09 21:54:34 +02:00
Lucio FrancoandGitHub 28a1eab1e0 chore: Fix unused warnings (#551) 2022-06-23 09:16:03 -07:00
ZettrokeandGitHub 3536017ccf Fix chain remaining_mut(), allowing to chain growing buffer (#488) 2022-06-11 13:12:53 +09:00
Erick TryzelaarandGitHub b8d27c016f Add UninitSlice::as_uninit_slice_mut() (#548)
This adds an unsafe method to convert a `&mut UninitSlice` into a
`&mut [MaybeUninit<u8>]`. This method is unsafe because some of the
bytes in the slice may be initialized, and the caller should not
overwrite them with uninitialized bytes.

This came about when auditing [tokio-util's udp frame], where they want
to pass the unitialized portion of a `BytesMut` to [ReadBuf::uninit].
They need to do this unsafe pointer casting in a few places, which
complicates audits. This method lets us document the safety invariants
the caller needs to maintain when doing this conversion.

[tokio-util's udp frame]: https://github.com/tokio-rs/tokio/blob/master/tokio-util/src/udp/frame.rs#L87
[ReadBuf::uninit]: https://docs.rs/tokio/latest/tokio/io/struct.ReadBuf.html#method.uninit
2022-05-10 11:22:19 +02:00
Taiki EndoandGitHub 0ce4fe3c91 Update actions/checkout action to v3 (#546) 2022-05-01 14:50:55 +02:00
Alice RyhlandGitHub 716a0b189e Only avoid pointer casts when using miri (#545) 2022-04-29 22:01:26 +02:00
Alice Ryhl b4b2c18c27 Revert accidental push directly to master
This reverts commit 89061c3238.

Why am I even able to push to master?
2022-04-29 19:49:28 +02:00
Alice Ryhl 89061c3238 Only avoid pointer casts when using miri 2022-04-29 19:47:34 +02:00
Jiahao XUandGitHub 0a2c43af88 Fix bugs in BytesMut::reserve_inner (#544) 2022-04-28 11:37:33 +02:00
Alice RyhlandGitHub 8198f9e28e Make strict provenance compatible (#542) 2022-04-16 00:26:58 +02:00
Alice RyhlandGitHub 547a32033e Add TSAN support (#541) 2022-04-15 22:46:40 +02:00
Ben KimockandGitHub 724476982b Fix aliasing in Clone by using a raw pointer (#523)
Previously, this code produced a &mut[u8] and a Box<[u8]> to the shared
allocation upon cloning it. If the underlying allocation were actually
shared, such as through a &[u8] from the Deref impl, creating either of
these types incorrectly asserted uniqueness of the allocation.

This fixes the example in #522, but Miri still does not pass on this
test suite with -Zmiri-tag-raw-pointers because Miri does not currently
understand int to pointer casts.
2022-04-06 16:59:20 +02:00
Evan CameronandGitHub 9e6edd18d2 Clarify BytesMut::unsplit docs (#535) 2022-04-06 15:19:00 +02:00
Anthony DeschampsandGitHub e4c723697d docs: redraw layout diagram with box drawing characters. (#539)
I find this diagram very helpful, but a little hard to distinguish
between the boxes and the lines that connect them. This commit redraws
the boxes with line drawing characters so that the boxes appear a
little more solid, and stand out from the other lines.
2022-03-25 10:55:13 +01:00
Jiahao XUandGitHub d4f5023383 Optimize BytesMut::reserve: Reuse vec if possible (#529)
* Optimize `BytesMut::reserve`: Reuse vec if possible

If the `BytesMut` holds a unqiue reference to `KIND_ARC` while the
capacity of the `Vec` is not big enough , reuse the existing `Vec`
instead of allocating a new one.

Signed-off-by: Jiahao XU <[email protected]>
2022-03-16 10:11:42 -04:00
Ralf JungandGitHub 88f5e12350 update Miri CI config (#534) 2022-03-09 10:29:04 +01:00
Rob EdeandGitHub 131dae161f Implement Extend<Bytes> for BytesMut (#527) 2022-01-24 09:58:18 +01:00
Rob EdeandGitHub 0e3b2466f1 Address various clippy warnings (#528) 2022-01-24 09:58:05 +01:00
Donough LiuandGitHub 68afb40df9 Add BytesMut::zeroed (#517) 2021-11-24 10:20:21 +01:00
Cyborus04andGitHub d946ef2e91 const-ify Bytes::len and Bytes::is_empty (#514) 2021-11-09 11:41:40 +01:00
Noah KennedyandGitHub ba5c5c93af Appease miri (#515)
Rewrote the ledger in test_bytes_vec_alloc.rs to not piss off miri. The ledger is now a table within the allocator, which seems to satisfy miri. The old solution was to bundle an extra usize into the beginning of each allocation and then index past the start when deallocating data to get the size.
2021-11-07 10:23:12 -06:00
Alice RyhlandGitHub ebc61e5af1 chore: prepare bytes v1.1.0 (#509) 2021-08-25 17:48:41 +02:00
Christopher HotchkissandGitHub 55e296850d Clarifying actions of clear and truncate. (#508) 2021-08-24 16:12:20 +02:00
Ian JacksonandGitHub 0e9fa0b602 impl From<Box<[u8]>> for Bytes (#504) 2021-08-24 12:42:22 +02:00
ee24be7fa0 ci: fetch cargo hack from github release (#507)
Co-authored-by: Taiki Endo <[email protected]>
2021-08-24 12:33:16 +02:00
Stepan KoltsovandGitHub 2697fa7a9d BufMut::put_bytes(self, val, cnt) (#487)
Equivalent to

```
for _ in 0..cnt {
    self.put_u8(val);
}
```

but may work faster.

Name and signature is chosen to be consistent with `ptr::write_bytes`.

Include three specializations:
* `Vec<u8>`
* `&mut [u8]`
* `BytesMut`

`BytesMut` and `&mut [u8]` specializations use `ptr::write`, `Vec<u8>`
specialization uses `Vec::resize`.
2021-08-09 02:43:53 +09:00
Stepan KoltsovandGitHub fa9cbf1258 Clarify BufPut::put_int behavior (#486)
* writes low bytes, discards high bytes
* panics if `nbytes` is greater than 8
2021-08-09 02:43:09 +09:00
Alice RyhlandGitHub ab8e3c01a8 Clarify BufMut allocation guarantees (#501) 2021-08-07 08:22:18 +02:00
Taiki EndoandGitHub f34dc5c3f9 Remove doc URLs (#498) 2021-08-07 02:06:57 +09:00
GbillouandGitHub baaf12d22a Keep capacity when unsplit on empty other buf (#502) 2021-07-05 16:46:17 +02:00
Taiki EndoandGitHub ed1d24e570 Use ubuntu-latest instead of ubuntu-16.04 (#497) 2021-05-23 23:09:26 +09:00
Taiki EndoandGitHub b89247c713 Update loom to 0.5 (#494) 2021-04-13 23:58:18 +09:00
NoahandGitHub 9c770188fd Fully inline BytesMut::new (#493) 2021-04-11 05:49:33 +09:00
Stepan KoltsovandGitHub b9eade12a5 Specialize copy_to_bytes for Chain and Take (#481)
Avoid allocation when `Take` or `Chain` is composed of `Bytes`
objects.

This works now for `Take`.

`Chain` it works if the requested bytes does not cross boundary
between `Chain` members.
2021-04-11 03:56:35 +09:00
Dan BurkertandGitHub 3d5624a452 Add inline tags to UninitSlice methods (#443)
This appears to be the primary cause of significant performance
regressions in the `prost` test suite in the 0.5 to 0.6 transition.  See
danburkert/prost#381.
2021-04-11 03:19:30 +09:00
Stepan KoltsovandGitHub 2428c152a6 Panic on integer overflow in Chain::remaining (#482)
Make it safer.
2021-02-16 12:44:33 -08:00
ZettrokeandGitHub 268f6f80b4 override put_slice for &mut [u8] (#483) 2021-02-15 16:35:01 -08:00
Alice RyhlandGitHub e4182808df Make bytes_mut -> chunk_mut rename more easily discoverable (#471) 2021-01-23 15:22:43 +09:00
Christopher BunnandGitHub 8daf43e9bd docs: fix broken Take link (#466) 2021-01-20 14:39:23 +09:00
Alice RyhlandGitHub 7b18c1c076 prepare 1.0.1 release (#460) 2021-01-11 18:07:46 +01:00
Ralf JungandGitHub df20a68356 use Box::into_raw instead of mem-forget-in-disguise (#458) 2020-12-31 15:07:28 +01:00
laizyandGitHub 8758a1aba5 add inline for Vec::put_slice (#459) 2020-12-31 12:34:39 +01:00
Ralf JungandGitHub 27a0f9ca6e CI: run test suite in Miri (#456) 2020-12-29 22:54:48 +01:00
Alice RyhlandGitHub ed71a7beb3 Fix deprecation warning (#457) 2020-12-29 22:46:39 +01:00
Carl LercheandGitHub 064ad9a1a0 chore: prepare v1.0.0 release (#453) 2020-12-22 15:30:20 -08:00
Taiki EndoandGitHub ed1d194660 deps: update loom to 0.4 (#452) 2020-12-22 10:30:50 +01:00
Arve KnudsenandGitHub e398b0a209 Update readme / changelog (#451) 2020-12-20 07:31:22 -08:00
Carl LercheandGitHub 06907f3e7b Rename Buf/BufMut, methods to chunk/chunk_mut (#450)
The `bytes()` / `bytes_mut()` name implies the method returns the full
set of bytes represented by `Buf`/`BufMut`. To rectify this, the methods
are renamed to `chunk()` and `chunk_mut()` to reflect the partial nature
of the returned byte slice.

`bytes_vectored()` is renamed `chunks_vectored()`.

Closes #447
2020-12-18 11:04:31 -08:00
Carl LercheandGitHub 54f5ced6c5 remove unused Buf implementation. (#449)
The implementation of `Buf` for `Option<[u8; 1]>` was added to support
`IntoBuf`. The `IntoBuf` trait has since been removed.

Closes #444
2020-12-16 21:51:13 -08:00
Carl LercheandGitHub bd78f19393 chore: prepare for v1.0.0 work (#448) 2020-12-12 08:25:14 -08:00
30 changed files with 1244 additions and 349 deletions
+25 -19
View File
@@ -11,6 +11,7 @@ on:
env: env:
RUSTFLAGS: -Dwarnings RUSTFLAGS: -Dwarnings
RUST_BACKTRACE: 1 RUST_BACKTRACE: 1
nightly: nightly-2021-11-05
defaults: defaults:
run: run:
@@ -22,7 +23,7 @@ jobs:
name: rustfmt name: rustfmt
runs-on: ubuntu-latest runs-on: ubuntu-latest
steps: steps:
- uses: actions/checkout@v2 - uses: actions/checkout@v3
- name: Install Rust - name: Install Rust
run: rustup update stable && rustup default stable run: rustup update stable && rustup default stable
- name: Check formatting - name: Check formatting
@@ -34,7 +35,7 @@ jobs:
# name: clippy # name: clippy
# runs-on: ubuntu-latest # runs-on: ubuntu-latest
# steps: # steps:
# - uses: actions/checkout@v2 # - uses: actions/checkout@v3
# - name: Apply clippy lints # - name: Apply clippy lints
# run: cargo clippy --all-features # run: cargo clippy --all-features
@@ -47,7 +48,7 @@ jobs:
name: minrust name: minrust
runs-on: ubuntu-latest runs-on: ubuntu-latest
steps: steps:
- uses: actions/checkout@v2 - uses: actions/checkout@v3
- name: Install Rust - name: Install Rust
run: rustup update 1.39.0 && rustup default 1.39.0 run: rustup update 1.39.0 && rustup default 1.39.0
- name: Check - name: Check
@@ -64,7 +65,7 @@ jobs:
- windows-latest - windows-latest
runs-on: ${{ matrix.os }} runs-on: ${{ matrix.os }}
steps: steps:
- uses: actions/checkout@v2 - uses: actions/checkout@v3
- name: Install Rust - name: Install Rust
# --no-self-update is necessary because the windows environment cannot self-update rustup.exe. # --no-self-update is necessary because the windows environment cannot self-update rustup.exe.
run: rustup update stable --no-self-update && rustup default stable run: rustup update stable --no-self-update && rustup default stable
@@ -74,14 +75,11 @@ jobs:
# Nightly # Nightly
nightly: nightly:
name: nightly name: nightly
env:
# Pin nightly to avoid being impacted by breakage
RUST_VERSION: nightly-2019-09-25
runs-on: ubuntu-latest runs-on: ubuntu-latest
steps: steps:
- uses: actions/checkout@v2 - uses: actions/checkout@v3
- name: Install Rust - name: Install Rust
run: rustup update $RUST_VERSION && rustup default $RUST_VERSION run: rustup update $nightly && rustup default $nightly
- name: Test - name: Test
run: . ci/test-stable.sh test run: . ci/test-stable.sh test
@@ -96,9 +94,9 @@ jobs:
- powerpc-unknown-linux-gnu - powerpc-unknown-linux-gnu
- powerpc64-unknown-linux-gnu - powerpc64-unknown-linux-gnu
- wasm32-unknown-unknown - wasm32-unknown-unknown
runs-on: ubuntu-16.04 runs-on: ubuntu-latest
steps: steps:
- uses: actions/checkout@v2 - uses: actions/checkout@v3
- name: Install Rust - name: Install Rust
run: rustup update stable && rustup default stable run: rustup update stable && rustup default stable
- name: cross build --target ${{ matrix.target }} - name: cross build --target ${{ matrix.target }}
@@ -118,25 +116,29 @@ jobs:
name: tsan name: tsan
runs-on: ubuntu-latest runs-on: ubuntu-latest
steps: steps:
- uses: actions/checkout@v2 - uses: actions/checkout@v3
- name: Install Rust - name: Install Rust
run: rustup update nightly && rustup default nightly run: rustup update $nightly && rustup default $nightly
- name: Install rust-src - name: Install rust-src
run: rustup component add rust-src run: rustup component add rust-src
- name: ASAN / TSAN - name: ASAN / TSAN
run: . ci/tsan.sh run: . ci/tsan.sh
miri:
name: miri
runs-on: ubuntu-latest
steps:
- uses: actions/checkout@v3
- name: Miri
run: ci/miri.sh
# Loom # Loom
loom: loom:
name: loom name: loom
env:
# Pin nightly to avoid being impacted by breakage
RUST_VERSION: nightly-2020-05-19
runs-on: ubuntu-latest runs-on: ubuntu-latest
steps: steps:
- uses: actions/checkout@v2 - uses: actions/checkout@v3
- name: Install Rust - name: Install Rust
run: rustup update $RUST_VERSION && rustup default $RUST_VERSION run: rustup update $nightly && rustup default $nightly
- name: Loom tests - name: Loom tests
run: RUSTFLAGS="--cfg loom -Dwarnings" cargo test --lib run: RUSTFLAGS="--cfg loom -Dwarnings" cargo test --lib
@@ -149,13 +151,17 @@ jobs:
- nightly - nightly
- minrust - minrust
- cross - cross
- tsan
- loom
runs-on: ubuntu-latest runs-on: ubuntu-latest
steps: steps:
- uses: actions/checkout@v2 - uses: actions/checkout@v3
- name: Install Rust - name: Install Rust
run: rustup update stable && rustup default stable run: rustup update stable && rustup default stable
- name: Build documentation - name: Build documentation
run: cargo doc --no-deps --all-features run: cargo doc --no-deps --all-features
env:
RUSTDOCFLAGS: --cfg docsrs
- name: Publish documentation - name: Publish documentation
run: | run: |
cd target/doc cd target/doc
+65
View File
@@ -1,3 +1,68 @@
# 1.2.0 (July 19, 2022)
### Added
- Add `BytesMut::zeroed` (#517)
- Implement `Extend<Bytes>` for `BytesMut` (#527)
- Add conversion from `BytesMut` to `Vec<u8>` (#543, #554)
- Add conversion from `Bytes` to `Vec<u8>` (#547)
- Add `UninitSlice::as_uninit_slice_mut()` (#548)
- Add const to `Bytes::{len,is_empty}` (#514)
### Changed
- Reuse vector in `BytesMut::reserve` (#539, #544)
### Fixed
- Make miri happy (#515, #523, #542, #545, #553)
- Make tsan happy (#541)
- Fix `remaining_mut()` on chain (#488)
- Fix amortized asymptotics of `BytesMut` (#555)
### Documented
- Redraw layout diagram with box drawing characters (#539)
- Clarify `BytesMut::unsplit` docs (#535)
# 1.1.0 (August 25, 2021)
### Added
- `BufMut::put_bytes(self, val, cnt)` (#487)
- Implement `From<Box<[u8]>>` for `Bytes` (#504)
### Changed
- Override `put_slice` for `&mut [u8]` (#483)
- Panic on integer overflow in `Chain::remaining` (#482)
- Add inline tags to `UninitSlice` methods (#443)
- Override `copy_to_bytes` for Chain and Take (#481)
- Keep capacity when unsplit on empty other buf (#502)
### Documented
- Clarify `BufMut` allocation guarantees (#501)
- Clarify `BufMut::put_int` behavior (#486)
- Clarify actions of `clear` and `truncate`. (#508)
# 1.0.1 (January 11, 2021)
### Changed
- mark `Vec::put_slice` with `#[inline]` (#459)
### Fixed
- Fix deprecation warning (#457)
- use `Box::into_raw` instead of `mem::forget`-in-disguise (#458)
# 1.0.0 (December 22, 2020)
### Changed
- Rename `Buf`/`BufMut` methods `bytes()` and `bytes_mut()` to `chunk()` and `chunk_mut()` (#450)
### Removed
- remove unused Buf implementation. (#449)
# 0.6.0 (October 21, 2020) # 0.6.0 (October 21, 2020)
API polish in preparation for a 1.0 release. API polish in preparation for a 1.0 release.
+6 -6
View File
@@ -2,18 +2,15 @@
name = "bytes" name = "bytes"
# When releasing to crates.io: # When releasing to crates.io:
# - Update html_root_url.
# - Update CHANGELOG.md. # - Update CHANGELOG.md.
# - Update doc URL. # - Create "v1.x.y" git tag.
# - Create "v0.6.x" git tag. version = "1.2.0"
version = "0.6.0"
license = "MIT" license = "MIT"
authors = [ authors = [
"Carl Lerche <[email protected]>", "Carl Lerche <[email protected]>",
"Sean McArthur <[email protected]>", "Sean McArthur <[email protected]>",
] ]
description = "Types and traits for working with bytes" description = "Types and traits for working with bytes"
documentation = "https://docs.rs/bytes/0.6.0/bytes/"
repository = "https://github.com/tokio-rs/bytes" repository = "https://github.com/tokio-rs/bytes"
readme = "README.md" readme = "README.md"
keywords = ["buffers", "zero-copy", "io"] keywords = ["buffers", "zero-copy", "io"]
@@ -31,4 +28,7 @@ serde = { version = "1.0.60", optional = true, default-features = false, feature
serde_test = "1.0" serde_test = "1.0"
[target.'cfg(loom)'.dev-dependencies] [target.'cfg(loom)'.dev-dependencies]
loom = "0.3" loom = "0.5"
[package.metadata.docs.rs]
rustdoc-args = ["--cfg", "docsrs"]
+2 -2
View File
@@ -18,7 +18,7 @@ To use `bytes`, first add this to your `Cargo.toml`:
```toml ```toml
[dependencies] [dependencies]
bytes = "0.6" bytes = "1"
``` ```
Next, add this to your crate: Next, add this to your crate:
@@ -33,7 +33,7 @@ Serde support is optional and disabled by default. To enable use the feature `se
```toml ```toml
[dependencies] [dependencies]
bytes = { version = "0.6", features = ["serde"] } bytes = { version = "1", features = ["serde"] }
``` ```
## License ## License
+4 -5
View File
@@ -46,14 +46,14 @@ impl TestBuf {
} }
impl Buf for TestBuf { impl Buf for TestBuf {
fn remaining(&self) -> usize { fn remaining(&self) -> usize {
return self.buf.len() - self.pos; self.buf.len() - self.pos
} }
fn advance(&mut self, cnt: usize) { fn advance(&mut self, cnt: usize) {
self.pos += cnt; self.pos += cnt;
assert!(self.pos <= self.buf.len()); assert!(self.pos <= self.buf.len());
self.next_readlen(); self.next_readlen();
} }
fn bytes(&self) -> &[u8] { fn chunk(&self) -> &[u8] {
if self.readlen == 0 { if self.readlen == 0 {
Default::default() Default::default()
} else { } else {
@@ -87,8 +87,8 @@ impl Buf for TestBufC {
self.inner.advance(cnt) self.inner.advance(cnt)
} }
#[inline(never)] #[inline(never)]
fn bytes(&self) -> &[u8] { fn chunk(&self) -> &[u8] {
self.inner.bytes() self.inner.chunk()
} }
} }
@@ -159,7 +159,6 @@ macro_rules! bench_group {
mod get_u8 { mod get_u8 {
use super::*; use super::*;
bench_group!(get_u8); bench_group!(get_u8);
bench!(option, option);
} }
mod get_u16 { mod get_u16 {
use super::*; use super::*;
+1
View File
@@ -88,6 +88,7 @@ fn from_long_slice(b: &mut Bencher) {
#[bench] #[bench]
fn slice_empty(b: &mut Bencher) { fn slice_empty(b: &mut Bencher) {
b.iter(|| { b.iter(|| {
// `clone` is to convert to ARC
let b = Bytes::from(vec![17; 1024]).clone(); let b = Bytes::from(vec![17; 1024]).clone();
for i in 0..1000 { for i in 0..1000 {
test::black_box(b.slice(i % 100..i % 100)); test::black_box(b.slice(i % 100..i % 100));
Executable
+11
View File
@@ -0,0 +1,11 @@
#!/bin/bash
set -e
rustup toolchain install nightly --component miri
rustup override set nightly
cargo miri setup
export MIRIFLAGS="-Zmiri-strict-provenance"
cargo miri test
cargo miri test --target mips64-unknown-linux-gnuabi64
+2 -1
View File
@@ -5,7 +5,8 @@ set -ex
cmd="${1:-test}" cmd="${1:-test}"
# Install cargo-hack for feature flag test # Install cargo-hack for feature flag test
cargo install cargo-hack host=$(rustc -Vv | grep host | sed 's/host: //')
curl -LsSf https://github.com/taiki-e/cargo-hack/releases/latest/download/cargo-hack-$host.tar.gz | tar xzf - -C ~/.cargo/bin
# Run with each feature # Run with each feature
# * --each-feature includes both default/no-default features # * --each-feature includes both default/no-default features
+1
View File
@@ -0,0 +1 @@
msrv = "1.39"
+27 -53
View File
@@ -16,7 +16,7 @@ macro_rules! buf_get_impl {
// this Option<ret> trick is to avoid keeping a borrow on self // this Option<ret> trick is to avoid keeping a borrow on self
// when advance() is called (mut borrow) and to call bytes() only once // when advance() is called (mut borrow) and to call bytes() only once
let ret = $this let ret = $this
.bytes() .chunk()
.get(..SIZE) .get(..SIZE)
.map(|src| unsafe { $typ::$conv(*(src as *const _ as *const [_; SIZE])) }); .map(|src| unsafe { $typ::$conv(*(src as *const _ as *const [_; SIZE])) });
@@ -78,7 +78,7 @@ pub trait Buf {
/// the buffer. /// the buffer.
/// ///
/// This value is greater than or equal to the length of the slice returned /// This value is greater than or equal to the length of the slice returned
/// by `bytes`. /// by `chunk()`.
/// ///
/// # Examples /// # Examples
/// ///
@@ -115,31 +115,34 @@ pub trait Buf {
/// ///
/// let mut buf = &b"hello world"[..]; /// let mut buf = &b"hello world"[..];
/// ///
/// assert_eq!(buf.bytes(), &b"hello world"[..]); /// assert_eq!(buf.chunk(), &b"hello world"[..]);
/// ///
/// buf.advance(6); /// buf.advance(6);
/// ///
/// assert_eq!(buf.bytes(), &b"world"[..]); /// assert_eq!(buf.chunk(), &b"world"[..]);
/// ``` /// ```
/// ///
/// # Implementer notes /// # Implementer notes
/// ///
/// This function should never panic. Once the end of the buffer is reached, /// This function should never panic. Once the end of the buffer is reached,
/// i.e., `Buf::remaining` returns 0, calls to `bytes` should return an /// i.e., `Buf::remaining` returns 0, calls to `chunk()` should return an
/// empty slice. /// empty slice.
fn bytes(&self) -> &[u8]; // The `chunk` method was previously called `bytes`. This alias makes the rename
// more easily discoverable.
#[cfg_attr(docsrs, doc(alias = "bytes"))]
fn chunk(&self) -> &[u8];
/// Fills `dst` with potentially multiple slices starting at `self`'s /// Fills `dst` with potentially multiple slices starting at `self`'s
/// current position. /// current position.
/// ///
/// If the `Buf` is backed by disjoint slices of bytes, `bytes_vectored` enables /// If the `Buf` is backed by disjoint slices of bytes, `chunk_vectored` enables
/// fetching more than one slice at once. `dst` is a slice of `IoSlice` /// fetching more than one slice at once. `dst` is a slice of `IoSlice`
/// references, enabling the slice to be directly used with [`writev`] /// references, enabling the slice to be directly used with [`writev`]
/// without any further conversion. The sum of the lengths of all the /// without any further conversion. The sum of the lengths of all the
/// buffers in `dst` will be less than or equal to `Buf::remaining()`. /// buffers in `dst` will be less than or equal to `Buf::remaining()`.
/// ///
/// The entries in `dst` will be overwritten, but the data **contained** by /// The entries in `dst` will be overwritten, but the data **contained** by
/// the slices **will not** be modified. If `bytes_vectored` does not fill every /// the slices **will not** be modified. If `chunk_vectored` does not fill every
/// entry in `dst`, then `dst` is guaranteed to contain all remaining slices /// entry in `dst`, then `dst` is guaranteed to contain all remaining slices
/// in `self. /// in `self.
/// ///
@@ -149,7 +152,7 @@ pub trait Buf {
/// # Implementer notes /// # Implementer notes
/// ///
/// This function should never panic. Once the end of the buffer is reached, /// This function should never panic. Once the end of the buffer is reached,
/// i.e., `Buf::remaining` returns 0, calls to `bytes_vectored` must return 0 /// i.e., `Buf::remaining` returns 0, calls to `chunk_vectored` must return 0
/// without mutating `dst`. /// without mutating `dst`.
/// ///
/// Implementations should also take care to properly handle being called /// Implementations should also take care to properly handle being called
@@ -157,13 +160,13 @@ pub trait Buf {
/// ///
/// [`writev`]: http://man7.org/linux/man-pages/man2/readv.2.html /// [`writev`]: http://man7.org/linux/man-pages/man2/readv.2.html
#[cfg(feature = "std")] #[cfg(feature = "std")]
fn bytes_vectored<'a>(&'a self, dst: &mut [IoSlice<'a>]) -> usize { fn chunks_vectored<'a>(&'a self, dst: &mut [IoSlice<'a>]) -> usize {
if dst.is_empty() { if dst.is_empty() {
return 0; return 0;
} }
if self.has_remaining() { if self.has_remaining() {
dst[0] = IoSlice::new(self.bytes()); dst[0] = IoSlice::new(self.chunk());
1 1
} else { } else {
0 0
@@ -172,7 +175,7 @@ pub trait Buf {
/// Advance the internal cursor of the Buf /// Advance the internal cursor of the Buf
/// ///
/// The next call to `bytes` will return a slice starting `cnt` bytes /// The next call to `chunk()` will return a slice starting `cnt` bytes
/// further into the underlying buffer. /// further into the underlying buffer.
/// ///
/// # Examples /// # Examples
@@ -182,11 +185,11 @@ pub trait Buf {
/// ///
/// let mut buf = &b"hello world"[..]; /// let mut buf = &b"hello world"[..];
/// ///
/// assert_eq!(buf.bytes(), &b"hello world"[..]); /// assert_eq!(buf.chunk(), &b"hello world"[..]);
/// ///
/// buf.advance(6); /// buf.advance(6);
/// ///
/// assert_eq!(buf.bytes(), &b"world"[..]); /// assert_eq!(buf.chunk(), &b"world"[..]);
/// ``` /// ```
/// ///
/// # Panics /// # Panics
@@ -253,7 +256,7 @@ pub trait Buf {
let cnt; let cnt;
unsafe { unsafe {
let src = self.bytes(); let src = self.chunk();
cnt = cmp::min(src.len(), dst.len() - off); cnt = cmp::min(src.len(), dst.len() - off);
ptr::copy_nonoverlapping(src.as_ptr(), dst[off..].as_mut_ptr(), cnt); ptr::copy_nonoverlapping(src.as_ptr(), dst[off..].as_mut_ptr(), cnt);
@@ -283,7 +286,7 @@ pub trait Buf {
/// This function panics if there is no more remaining data in `self`. /// This function panics if there is no more remaining data in `self`.
fn get_u8(&mut self) -> u8 { fn get_u8(&mut self) -> u8 {
assert!(self.remaining() >= 1); assert!(self.remaining() >= 1);
let ret = self.bytes()[0]; let ret = self.chunk()[0];
self.advance(1); self.advance(1);
ret ret
} }
@@ -306,7 +309,7 @@ pub trait Buf {
/// This function panics if there is no more remaining data in `self`. /// This function panics if there is no more remaining data in `self`.
fn get_i8(&mut self) -> i8 { fn get_i8(&mut self) -> i8 {
assert!(self.remaining() >= 1); assert!(self.remaining() >= 1);
let ret = self.bytes()[0] as i8; let ret = self.chunk()[0] as i8;
self.advance(1); self.advance(1);
ret ret
} }
@@ -861,7 +864,7 @@ pub trait Buf {
/// let mut chain = b"hello "[..].chain(&b"world"[..]); /// let mut chain = b"hello "[..].chain(&b"world"[..]);
/// ///
/// let full = chain.copy_to_bytes(11); /// let full = chain.copy_to_bytes(11);
/// assert_eq!(full.bytes(), b"hello world"); /// assert_eq!(full.chunk(), b"hello world");
/// ``` /// ```
fn chain<U: Buf>(self, next: U) -> Chain<Self, U> fn chain<U: Buf>(self, next: U) -> Chain<Self, U>
where where
@@ -908,13 +911,13 @@ macro_rules! deref_forward_buf {
(**self).remaining() (**self).remaining()
} }
fn bytes(&self) -> &[u8] { fn chunk(&self) -> &[u8] {
(**self).bytes() (**self).chunk()
} }
#[cfg(feature = "std")] #[cfg(feature = "std")]
fn bytes_vectored<'b>(&'b self, dst: &mut [IoSlice<'b>]) -> usize { fn chunks_vectored<'b>(&'b self, dst: &mut [IoSlice<'b>]) -> usize {
(**self).bytes_vectored(dst) (**self).chunks_vectored(dst)
} }
fn advance(&mut self, cnt: usize) { fn advance(&mut self, cnt: usize) {
@@ -1022,7 +1025,7 @@ impl Buf for &[u8] {
} }
#[inline] #[inline]
fn bytes(&self) -> &[u8] { fn chunk(&self) -> &[u8] {
self self
} }
@@ -1032,35 +1035,6 @@ impl Buf for &[u8] {
} }
} }
impl Buf for Option<[u8; 1]> {
fn remaining(&self) -> usize {
if self.is_some() {
1
} else {
0
}
}
fn bytes(&self) -> &[u8] {
self.as_ref()
.map(AsRef::as_ref)
.unwrap_or(Default::default())
}
fn advance(&mut self, cnt: usize) {
if cnt == 0 {
return;
}
if self.is_none() {
panic!("overflow");
} else {
assert_eq!(1, cnt);
*self = None;
}
}
}
#[cfg(feature = "std")] #[cfg(feature = "std")]
impl<T: AsRef<[u8]>> Buf for std::io::Cursor<T> { impl<T: AsRef<[u8]>> Buf for std::io::Cursor<T> {
fn remaining(&self) -> usize { fn remaining(&self) -> usize {
@@ -1074,7 +1048,7 @@ impl<T: AsRef<[u8]>> Buf for std::io::Cursor<T> {
len - pos as usize len - pos as usize
} }
fn bytes(&self) -> &[u8] { fn chunk(&self) -> &[u8] {
let len = self.get_ref().as_ref().len(); let len = self.get_ref().as_ref().len();
let pos = self.position(); let pos = self.position();
+97 -30
View File
@@ -31,7 +31,11 @@ pub unsafe trait BufMut {
/// position until the end of the buffer is reached. /// position until the end of the buffer is reached.
/// ///
/// This value is greater than or equal to the length of the slice returned /// This value is greater than or equal to the length of the slice returned
/// by `bytes_mut`. /// by `chunk_mut()`.
///
/// Writing to a `BufMut` may involve allocating more memory on the fly.
/// Implementations may fail before reaching the number of bytes indicated
/// by this method if they encounter an allocation failure.
/// ///
/// # Examples /// # Examples
/// ///
@@ -52,11 +56,15 @@ pub unsafe trait BufMut {
/// Implementations of `remaining_mut` should ensure that the return value /// Implementations of `remaining_mut` should ensure that the return value
/// does not change unless a call is made to `advance_mut` or any other /// does not change unless a call is made to `advance_mut` or any other
/// function that is documented to change the `BufMut`'s current position. /// function that is documented to change the `BufMut`'s current position.
///
/// # Note
///
/// `remaining_mut` may return value smaller than actual available space.
fn remaining_mut(&self) -> usize; fn remaining_mut(&self) -> usize;
/// Advance the internal cursor of the BufMut /// Advance the internal cursor of the BufMut
/// ///
/// The next call to `bytes_mut` will return a slice starting `cnt` bytes /// The next call to `chunk_mut` will return a slice starting `cnt` bytes
/// further into the underlying buffer. /// further into the underlying buffer.
/// ///
/// This function is unsafe because there is no guarantee that the bytes /// This function is unsafe because there is no guarantee that the bytes
@@ -70,11 +78,11 @@ pub unsafe trait BufMut {
/// let mut buf = Vec::with_capacity(16); /// let mut buf = Vec::with_capacity(16);
/// ///
/// // Write some data /// // Write some data
/// buf.bytes_mut()[0..2].copy_from_slice(b"he"); /// buf.chunk_mut()[0..2].copy_from_slice(b"he");
/// unsafe { buf.advance_mut(2) }; /// unsafe { buf.advance_mut(2) };
/// ///
/// // write more bytes /// // write more bytes
/// buf.bytes_mut()[0..3].copy_from_slice(b"llo"); /// buf.chunk_mut()[0..3].copy_from_slice(b"llo");
/// ///
/// unsafe { buf.advance_mut(3); } /// unsafe { buf.advance_mut(3); }
/// ///
@@ -135,14 +143,14 @@ pub unsafe trait BufMut {
/// ///
/// unsafe { /// unsafe {
/// // MaybeUninit::as_mut_ptr /// // MaybeUninit::as_mut_ptr
/// buf.bytes_mut()[0..].as_mut_ptr().write(b'h'); /// buf.chunk_mut()[0..].as_mut_ptr().write(b'h');
/// buf.bytes_mut()[1..].as_mut_ptr().write(b'e'); /// buf.chunk_mut()[1..].as_mut_ptr().write(b'e');
/// ///
/// buf.advance_mut(2); /// buf.advance_mut(2);
/// ///
/// buf.bytes_mut()[0..].as_mut_ptr().write(b'l'); /// buf.chunk_mut()[0..].as_mut_ptr().write(b'l');
/// buf.bytes_mut()[1..].as_mut_ptr().write(b'l'); /// buf.chunk_mut()[1..].as_mut_ptr().write(b'l');
/// buf.bytes_mut()[2..].as_mut_ptr().write(b'o'); /// buf.chunk_mut()[2..].as_mut_ptr().write(b'o');
/// ///
/// buf.advance_mut(3); /// buf.advance_mut(3);
/// } /// }
@@ -153,12 +161,18 @@ pub unsafe trait BufMut {
/// ///
/// # Implementer notes /// # Implementer notes
/// ///
/// This function should never panic. `bytes_mut` should return an empty /// This function should never panic. `chunk_mut` should return an empty
/// slice **if and only if** `remaining_mut` returns 0. In other words, /// slice **if and only if** `remaining_mut()` returns 0. In other words,
/// `bytes_mut` returning an empty slice implies that `remaining_mut` will /// `chunk_mut()` returning an empty slice implies that `remaining_mut()` will
/// return 0 and `remaining_mut` returning 0 implies that `bytes_mut` will /// return 0 and `remaining_mut()` returning 0 implies that `chunk_mut()` will
/// return an empty slice. /// return an empty slice.
fn bytes_mut(&mut self) -> &mut UninitSlice; ///
/// This function may trigger an out-of-memory abort if it tries to allocate
/// memory and fails to do so.
// The `chunk_mut` method was previously called `bytes_mut`. This alias makes the
// rename more easily discoverable.
#[cfg_attr(docsrs, doc(alias = "bytes_mut"))]
fn chunk_mut(&mut self) -> &mut UninitSlice;
/// Transfer bytes into `self` from `src` and advance the cursor by the /// Transfer bytes into `self` from `src` and advance the cursor by the
/// number of bytes written. /// number of bytes written.
@@ -190,8 +204,8 @@ pub unsafe trait BufMut {
let l; let l;
unsafe { unsafe {
let s = src.bytes(); let s = src.chunk();
let d = self.bytes_mut(); let d = self.chunk_mut();
l = cmp::min(s.len(), d.len()); l = cmp::min(s.len(), d.len());
ptr::copy_nonoverlapping(s.as_ptr(), d.as_mut_ptr() as *mut u8, l); ptr::copy_nonoverlapping(s.as_ptr(), d.as_mut_ptr() as *mut u8, l);
@@ -237,7 +251,7 @@ pub unsafe trait BufMut {
let cnt; let cnt;
unsafe { unsafe {
let dst = self.bytes_mut(); let dst = self.chunk_mut();
cnt = cmp::min(dst.len(), src.len() - off); cnt = cmp::min(dst.len(), src.len() - off);
ptr::copy_nonoverlapping(src[off..].as_ptr(), dst.as_mut_ptr() as *mut u8, cnt); ptr::copy_nonoverlapping(src[off..].as_ptr(), dst.as_mut_ptr() as *mut u8, cnt);
@@ -251,6 +265,37 @@ pub unsafe trait BufMut {
} }
} }
/// Put `cnt` bytes `val` into `self`.
///
/// Logically equivalent to calling `self.put_u8(val)` `cnt` times, but may work faster.
///
/// `self` must have at least `cnt` remaining capacity.
///
/// ```
/// use bytes::BufMut;
///
/// let mut dst = [0; 6];
///
/// {
/// let mut buf = &mut dst[..];
/// buf.put_bytes(b'a', 4);
///
/// assert_eq!(2, buf.remaining_mut());
/// }
///
/// assert_eq!(b"aaaa\0\0", &dst);
/// ```
///
/// # Panics
///
/// This function panics if there is not enough remaining capacity in
/// `self`.
fn put_bytes(&mut self, val: u8, cnt: usize) {
for _ in 0..cnt {
self.put_u8(val);
}
}
/// Writes an unsigned 8 bit integer to `self`. /// Writes an unsigned 8 bit integer to `self`.
/// ///
/// The current position is advanced by 1. /// The current position is advanced by 1.
@@ -693,7 +738,7 @@ pub unsafe trait BufMut {
self.put_slice(&n.to_le_bytes()[0..nbytes]); self.put_slice(&n.to_le_bytes()[0..nbytes]);
} }
/// Writes a signed n-byte integer to `self` in big-endian byte order. /// Writes low `nbytes` of a signed integer to `self` in big-endian byte order.
/// ///
/// The current position is advanced by `nbytes`. /// The current position is advanced by `nbytes`.
/// ///
@@ -703,19 +748,19 @@ pub unsafe trait BufMut {
/// use bytes::BufMut; /// use bytes::BufMut;
/// ///
/// let mut buf = vec![]; /// let mut buf = vec![];
/// buf.put_int(0x010203, 3); /// buf.put_int(0x0504010203, 3);
/// assert_eq!(buf, b"\x01\x02\x03"); /// assert_eq!(buf, b"\x01\x02\x03");
/// ``` /// ```
/// ///
/// # Panics /// # Panics
/// ///
/// This function panics if there is not enough remaining capacity in /// This function panics if there is not enough remaining capacity in
/// `self`. /// `self` or if `nbytes` is greater than 8.
fn put_int(&mut self, n: i64, nbytes: usize) { fn put_int(&mut self, n: i64, nbytes: usize) {
self.put_slice(&n.to_be_bytes()[mem::size_of_val(&n) - nbytes..]); self.put_slice(&n.to_be_bytes()[mem::size_of_val(&n) - nbytes..]);
} }
/// Writes a signed n-byte integer to `self` in little-endian byte order. /// Writes low `nbytes` of a signed integer to `self` in little-endian byte order.
/// ///
/// The current position is advanced by `nbytes`. /// The current position is advanced by `nbytes`.
/// ///
@@ -725,14 +770,14 @@ pub unsafe trait BufMut {
/// use bytes::BufMut; /// use bytes::BufMut;
/// ///
/// let mut buf = vec![]; /// let mut buf = vec![];
/// buf.put_int_le(0x010203, 3); /// buf.put_int_le(0x0504010203, 3);
/// assert_eq!(buf, b"\x03\x02\x01"); /// assert_eq!(buf, b"\x03\x02\x01");
/// ``` /// ```
/// ///
/// # Panics /// # Panics
/// ///
/// This function panics if there is not enough remaining capacity in /// This function panics if there is not enough remaining capacity in
/// `self`. /// `self` or if `nbytes` is greater than 8.
fn put_int_le(&mut self, n: i64, nbytes: usize) { fn put_int_le(&mut self, n: i64, nbytes: usize) {
self.put_slice(&n.to_le_bytes()[0..nbytes]); self.put_slice(&n.to_le_bytes()[0..nbytes]);
} }
@@ -913,8 +958,8 @@ macro_rules! deref_forward_bufmut {
(**self).remaining_mut() (**self).remaining_mut()
} }
fn bytes_mut(&mut self) -> &mut UninitSlice { fn chunk_mut(&mut self) -> &mut UninitSlice {
(**self).bytes_mut() (**self).chunk_mut()
} }
unsafe fn advance_mut(&mut self, cnt: usize) { unsafe fn advance_mut(&mut self, cnt: usize) {
@@ -998,7 +1043,7 @@ unsafe impl BufMut for &mut [u8] {
} }
#[inline] #[inline]
fn bytes_mut(&mut self) -> &mut UninitSlice { fn chunk_mut(&mut self) -> &mut UninitSlice {
// UninitSlice is repr(transparent), so safe to transmute // UninitSlice is repr(transparent), so safe to transmute
unsafe { &mut *(*self as *mut [u8] as *mut _) } unsafe { &mut *(*self as *mut [u8] as *mut _) }
} }
@@ -1009,12 +1054,29 @@ unsafe impl BufMut for &mut [u8] {
let (_, b) = core::mem::replace(self, &mut []).split_at_mut(cnt); let (_, b) = core::mem::replace(self, &mut []).split_at_mut(cnt);
*self = b; *self = b;
} }
#[inline]
fn put_slice(&mut self, src: &[u8]) {
self[..src.len()].copy_from_slice(src);
unsafe {
self.advance_mut(src.len());
}
}
fn put_bytes(&mut self, val: u8, cnt: usize) {
assert!(self.remaining_mut() >= cnt);
unsafe {
ptr::write_bytes(self.as_mut_ptr(), val, cnt);
self.advance_mut(cnt);
}
}
} }
unsafe impl BufMut for Vec<u8> { unsafe impl BufMut for Vec<u8> {
#[inline] #[inline]
fn remaining_mut(&self) -> usize { fn remaining_mut(&self) -> usize {
usize::MAX - self.len() // A vector can never have more than isize::MAX bytes
core::isize::MAX as usize - self.len()
} }
#[inline] #[inline]
@@ -1033,7 +1095,7 @@ unsafe impl BufMut for Vec<u8> {
} }
#[inline] #[inline]
fn bytes_mut(&mut self) -> &mut UninitSlice { fn chunk_mut(&mut self) -> &mut UninitSlice {
if self.capacity() == self.len() { if self.capacity() == self.len() {
self.reserve(64); // Grow the vec self.reserve(64); // Grow the vec
} }
@@ -1047,7 +1109,6 @@ unsafe impl BufMut for Vec<u8> {
// Specialize these methods so they can skip checking `remaining_mut` // Specialize these methods so they can skip checking `remaining_mut`
// and `advance_mut`. // and `advance_mut`.
fn put<T: super::Buf>(&mut self, mut src: T) fn put<T: super::Buf>(&mut self, mut src: T)
where where
Self: Sized, Self: Sized,
@@ -1060,7 +1121,7 @@ unsafe impl BufMut for Vec<u8> {
// a block to contain the src.bytes() borrow // a block to contain the src.bytes() borrow
{ {
let s = src.bytes(); let s = src.chunk();
l = s.len(); l = s.len();
self.extend_from_slice(s); self.extend_from_slice(s);
} }
@@ -1069,9 +1130,15 @@ unsafe impl BufMut for Vec<u8> {
} }
} }
#[inline]
fn put_slice(&mut self, src: &[u8]) { fn put_slice(&mut self, src: &[u8]) {
self.extend_from_slice(src); self.extend_from_slice(src);
} }
fn put_bytes(&mut self, val: u8, cnt: usize) {
let new_len = self.len().checked_add(cnt).unwrap();
self.resize(new_len, val);
}
} }
// The existence of this function makes the compiler catch if the BufMut // The existence of this function makes the compiler catch if the BufMut
+32 -12
View File
@@ -1,5 +1,5 @@
use crate::buf::{IntoIter, UninitSlice}; use crate::buf::{IntoIter, UninitSlice};
use crate::{Buf, BufMut}; use crate::{Buf, BufMut, Bytes};
#[cfg(feature = "std")] #[cfg(feature = "std")]
use std::io::IoSlice; use std::io::IoSlice;
@@ -135,14 +135,14 @@ where
U: Buf, U: Buf,
{ {
fn remaining(&self) -> usize { fn remaining(&self) -> usize {
self.a.remaining() + self.b.remaining() self.a.remaining().checked_add(self.b.remaining()).unwrap()
} }
fn bytes(&self) -> &[u8] { fn chunk(&self) -> &[u8] {
if self.a.has_remaining() { if self.a.has_remaining() {
self.a.bytes() self.a.chunk()
} else { } else {
self.b.bytes() self.b.chunk()
} }
} }
@@ -165,11 +165,29 @@ where
} }
#[cfg(feature = "std")] #[cfg(feature = "std")]
fn bytes_vectored<'a>(&'a self, dst: &mut [IoSlice<'a>]) -> usize { fn chunks_vectored<'a>(&'a self, dst: &mut [IoSlice<'a>]) -> usize {
let mut n = self.a.bytes_vectored(dst); let mut n = self.a.chunks_vectored(dst);
n += self.b.bytes_vectored(&mut dst[n..]); n += self.b.chunks_vectored(&mut dst[n..]);
n n
} }
fn copy_to_bytes(&mut self, len: usize) -> Bytes {
let a_rem = self.a.remaining();
if a_rem >= len {
self.a.copy_to_bytes(len)
} else if a_rem == 0 {
self.b.copy_to_bytes(len)
} else {
assert!(
len - a_rem <= self.b.remaining(),
"`len` greater than remaining"
);
let mut ret = crate::BytesMut::with_capacity(len);
ret.put(&mut self.a);
ret.put((&mut self.b).take(len - a_rem));
ret.freeze()
}
}
} }
unsafe impl<T, U> BufMut for Chain<T, U> unsafe impl<T, U> BufMut for Chain<T, U>
@@ -178,14 +196,16 @@ where
U: BufMut, U: BufMut,
{ {
fn remaining_mut(&self) -> usize { fn remaining_mut(&self) -> usize {
self.a.remaining_mut() + self.b.remaining_mut() self.a
.remaining_mut()
.saturating_add(self.b.remaining_mut())
} }
fn bytes_mut(&mut self) -> &mut UninitSlice { fn chunk_mut(&mut self) -> &mut UninitSlice {
if self.a.has_remaining_mut() { if self.a.has_remaining_mut() {
self.a.bytes_mut() self.a.chunk_mut()
} else { } else {
self.b.bytes_mut() self.b.chunk_mut()
} }
} }
+1 -1
View File
@@ -117,7 +117,7 @@ impl<T: Buf> Iterator for IntoIter<T> {
return None; return None;
} }
let b = self.inner.bytes()[0]; let b = self.inner.chunk()[0];
self.inner.advance(1); self.inner.advance(1);
Some(b) Some(b)
+2 -2
View File
@@ -61,8 +61,8 @@ unsafe impl<T: BufMut> BufMut for Limit<T> {
cmp::min(self.inner.remaining_mut(), self.limit) cmp::min(self.inner.remaining_mut(), self.limit)
} }
fn bytes_mut(&mut self) -> &mut UninitSlice { fn chunk_mut(&mut self) -> &mut UninitSlice {
let bytes = self.inner.bytes_mut(); let bytes = self.inner.chunk_mut();
let end = cmp::min(bytes.len(), self.limit); let end = cmp::min(bytes.len(), self.limit);
&mut bytes[..end] &mut bytes[..end]
} }
+1 -1
View File
@@ -73,7 +73,7 @@ impl<B: Buf + Sized> io::Read for Reader<B> {
impl<B: Buf + Sized> io::BufRead for Reader<B> { impl<B: Buf + Sized> io::BufRead for Reader<B> {
fn fill_buf(&mut self) -> io::Result<&[u8]> { fn fill_buf(&mut self) -> io::Result<&[u8]> {
Ok(self.buf.bytes()) Ok(self.buf.chunk())
} }
fn consume(&mut self, amt: usize) { fn consume(&mut self, amt: usize) {
self.buf.advance(amt) self.buf.advance(amt)
+12 -4
View File
@@ -1,11 +1,11 @@
use crate::Buf; use crate::{Buf, Bytes};
use core::cmp; use core::cmp;
/// A `Buf` adapter which limits the bytes read from an underlying buffer. /// A `Buf` adapter which limits the bytes read from an underlying buffer.
/// ///
/// This struct is generally created by calling `take()` on `Buf`. See /// This struct is generally created by calling `take()` on `Buf`. See
/// documentation of [`take()`](trait.BufExt.html#method.take) for more details. /// documentation of [`take()`](trait.Buf.html#method.take) for more details.
#[derive(Debug)] #[derive(Debug)]
pub struct Take<T> { pub struct Take<T> {
inner: T, inner: T,
@@ -134,8 +134,8 @@ impl<T: Buf> Buf for Take<T> {
cmp::min(self.inner.remaining(), self.limit) cmp::min(self.inner.remaining(), self.limit)
} }
fn bytes(&self) -> &[u8] { fn chunk(&self) -> &[u8] {
let bytes = self.inner.bytes(); let bytes = self.inner.chunk();
&bytes[..cmp::min(bytes.len(), self.limit)] &bytes[..cmp::min(bytes.len(), self.limit)]
} }
@@ -144,4 +144,12 @@ impl<T: Buf> Buf for Take<T> {
self.inner.advance(cnt); self.inner.advance(cnt);
self.limit -= cnt; self.limit -= cnt;
} }
fn copy_to_bytes(&mut self, len: usize) -> Bytes {
assert!(len <= self.remaining(), "`len` greater than remaining");
let r = self.inner.copy_to_bytes(len);
self.limit -= len;
r
}
} }
+36 -3
View File
@@ -6,7 +6,7 @@ use core::ops::{
/// Uninitialized byte slice. /// Uninitialized byte slice.
/// ///
/// Returned by `BufMut::bytes_mut()`, the referenced byte slice may be /// Returned by `BufMut::chunk_mut()`, the referenced byte slice may be
/// uninitialized. The wrapper provides safe access without introducing /// uninitialized. The wrapper provides safe access without introducing
/// undefined behavior. /// undefined behavior.
/// ///
@@ -40,6 +40,7 @@ impl UninitSlice {
/// ///
/// let slice = unsafe { UninitSlice::from_raw_parts_mut(ptr, len) }; /// let slice = unsafe { UninitSlice::from_raw_parts_mut(ptr, len) };
/// ``` /// ```
#[inline]
pub unsafe fn from_raw_parts_mut<'a>(ptr: *mut u8, len: usize) -> &'a mut UninitSlice { pub unsafe fn from_raw_parts_mut<'a>(ptr: *mut u8, len: usize) -> &'a mut UninitSlice {
let maybe_init: &mut [MaybeUninit<u8>] = let maybe_init: &mut [MaybeUninit<u8>] =
core::slice::from_raw_parts_mut(ptr as *mut _, len); core::slice::from_raw_parts_mut(ptr as *mut _, len);
@@ -64,6 +65,7 @@ impl UninitSlice {
/// ///
/// assert_eq!(b"boo", &data[..]); /// assert_eq!(b"boo", &data[..]);
/// ``` /// ```
#[inline]
pub fn write_byte(&mut self, index: usize, byte: u8) { pub fn write_byte(&mut self, index: usize, byte: u8) {
assert!(index < self.len()); assert!(index < self.len());
@@ -90,6 +92,7 @@ impl UninitSlice {
/// ///
/// assert_eq!(b"bar", &data[..]); /// assert_eq!(b"bar", &data[..]);
/// ``` /// ```
#[inline]
pub fn copy_from_slice(&mut self, src: &[u8]) { pub fn copy_from_slice(&mut self, src: &[u8]) {
use core::ptr; use core::ptr;
@@ -114,12 +117,39 @@ impl UninitSlice {
/// ///
/// let mut data = [0, 1, 2]; /// let mut data = [0, 1, 2];
/// let mut slice = &mut data[..]; /// let mut slice = &mut data[..];
/// let ptr = BufMut::bytes_mut(&mut slice).as_mut_ptr(); /// let ptr = BufMut::chunk_mut(&mut slice).as_mut_ptr();
/// ``` /// ```
#[inline]
pub fn as_mut_ptr(&mut self) -> *mut u8 { pub fn as_mut_ptr(&mut self) -> *mut u8 {
self.0.as_mut_ptr() as *mut _ self.0.as_mut_ptr() as *mut _
} }
/// Return a `&mut [MaybeUninit<u8>]` to this slice's buffer.
///
/// # Safety
///
/// The caller **must not** read from the referenced memory and **must not** write
/// **uninitialized** bytes to the slice either. This is because `BufMut` implementation
/// that created the `UninitSlice` knows which parts are initialized. Writing uninitalized
/// bytes to the slice may cause the `BufMut` to read those bytes and trigger undefined
/// behavior.
///
/// # Examples
///
/// ```
/// use bytes::BufMut;
///
/// let mut data = [0, 1, 2];
/// let mut slice = &mut data[..];
/// unsafe {
/// let uninit_slice = BufMut::chunk_mut(&mut slice).as_uninit_slice_mut();
/// };
/// ```
#[inline]
pub unsafe fn as_uninit_slice_mut<'a>(&'a mut self) -> &'a mut [MaybeUninit<u8>] {
&mut *(self as *mut _ as *mut [MaybeUninit<u8>])
}
/// Returns the number of bytes in the slice. /// Returns the number of bytes in the slice.
/// ///
/// # Examples /// # Examples
@@ -129,10 +159,11 @@ impl UninitSlice {
/// ///
/// let mut data = [0, 1, 2]; /// let mut data = [0, 1, 2];
/// let mut slice = &mut data[..]; /// let mut slice = &mut data[..];
/// let len = BufMut::bytes_mut(&mut slice).len(); /// let len = BufMut::chunk_mut(&mut slice).len();
/// ///
/// assert_eq!(len, 3); /// assert_eq!(len, 3);
/// ``` /// ```
#[inline]
pub fn len(&self) -> usize { pub fn len(&self) -> usize {
self.0.len() self.0.len()
} }
@@ -150,6 +181,7 @@ macro_rules! impl_index {
impl Index<$t> for UninitSlice { impl Index<$t> for UninitSlice {
type Output = UninitSlice; type Output = UninitSlice;
#[inline]
fn index(&self, index: $t) -> &UninitSlice { fn index(&self, index: $t) -> &UninitSlice {
let maybe_uninit: &[MaybeUninit<u8>] = &self.0[index]; let maybe_uninit: &[MaybeUninit<u8>] = &self.0[index];
unsafe { &*(maybe_uninit as *const [MaybeUninit<u8>] as *const UninitSlice) } unsafe { &*(maybe_uninit as *const [MaybeUninit<u8>] as *const UninitSlice) }
@@ -157,6 +189,7 @@ macro_rules! impl_index {
} }
impl IndexMut<$t> for UninitSlice { impl IndexMut<$t> for UninitSlice {
#[inline]
fn index_mut(&mut self, index: $t) -> &mut UninitSlice { fn index_mut(&mut self, index: $t) -> &mut UninitSlice {
let maybe_uninit: &mut [MaybeUninit<u8>] = &mut self.0[index]; let maybe_uninit: &mut [MaybeUninit<u8>] = &mut self.0[index];
unsafe { &mut *(maybe_uninit as *mut [MaybeUninit<u8>] as *mut UninitSlice) } unsafe { &mut *(maybe_uninit as *mut [MaybeUninit<u8>] as *mut UninitSlice) }
+1 -1
View File
@@ -7,7 +7,7 @@ impl Buf for VecDeque<u8> {
self.len() self.len()
} }
fn bytes(&self) -> &[u8] { fn chunk(&self) -> &[u8] {
let (s1, s2) = self.as_slices(); let (s1, s2) = self.as_slices();
if s1.is_empty() { if s1.is_empty() {
s2 s2
+210 -71
View File
@@ -2,12 +2,18 @@ use core::iter::FromIterator;
use core::ops::{Deref, RangeBounds}; use core::ops::{Deref, RangeBounds};
use core::{cmp, fmt, hash, mem, ptr, slice, usize}; use core::{cmp, fmt, hash, mem, ptr, slice, usize};
use alloc::{borrow::Borrow, boxed::Box, string::String, vec::Vec}; use alloc::{
alloc::{dealloc, Layout},
borrow::Borrow,
boxed::Box,
string::String,
vec::Vec,
};
use crate::buf::IntoIter; use crate::buf::IntoIter;
#[allow(unused)] #[allow(unused)]
use crate::loom::sync::atomic::AtomicMut; use crate::loom::sync::atomic::AtomicMut;
use crate::loom::sync::atomic::{self, AtomicPtr, AtomicUsize, Ordering}; use crate::loom::sync::atomic::{AtomicPtr, AtomicUsize, Ordering};
use crate::Buf; use crate::Buf;
/// A cheaply cloneable and sliceable chunk of contiguous memory. /// A cheaply cloneable and sliceable chunk of contiguous memory.
@@ -55,7 +61,7 @@ use crate::Buf;
/// # Sharing /// # Sharing
/// ///
/// `Bytes` contains a vtable, which allows implementations of `Bytes` to define /// `Bytes` contains a vtable, which allows implementations of `Bytes` to define
/// how sharing/cloneing is implemented in detail. /// how sharing/cloning is implemented in detail.
/// When `Bytes::clone()` is called, `Bytes` will call the vtable function for /// When `Bytes::clone()` is called, `Bytes` will call the vtable function for
/// cloning the backing storage in order to share it behind between multiple /// cloning the backing storage in order to share it behind between multiple
/// `Bytes` instances. /// `Bytes` instances.
@@ -78,18 +84,18 @@ use crate::Buf;
/// ///
/// ```text /// ```text
/// ///
/// Arc ptrs +---------+ /// Arc ptrs ┌─────────┐
/// ________________________ / | Bytes 2 | /// ________________________ / Bytes 2
/// / +---------+ /// / └─────────┘
/// / +-----------+ | | /// / ┌───────────┐ | |
/// |_________/ | Bytes 1 | | | /// |_________/ Bytes 1 | |
/// | +-----------+ | | /// | └───────────┘ | |
/// | | | ___/ data | tail /// | | | ___/ data | tail
/// | data | tail |/ | /// | data | tail |/ |
/// v v v v /// v v v v
/// +-----+---------------------------------+-----+ /// ┌─────┬─────┬───────────┬───────────────┬─────┐
/// | Arc | | | | | /// Arc
/// +-----+---------------------------------+-----+ /// └─────┴─────┴───────────┴───────────────┴─────┘
/// ``` /// ```
pub struct Bytes { pub struct Bytes {
ptr: *const u8, ptr: *const u8,
@@ -103,6 +109,10 @@ pub(crate) struct Vtable {
/// fn(data, ptr, len) /// fn(data, ptr, len)
pub clone: unsafe fn(&AtomicPtr<()>, *const u8, usize) -> Bytes, pub clone: unsafe fn(&AtomicPtr<()>, *const u8, usize) -> Bytes,
/// fn(data, ptr, len) /// fn(data, ptr, len)
///
/// takes `Bytes` to value
pub to_vec: unsafe fn(&AtomicPtr<()>, *const u8, usize) -> Vec<u8>,
/// fn(data, ptr, len)
pub drop: unsafe fn(&mut AtomicPtr<()>, *const u8, usize), pub drop: unsafe fn(&mut AtomicPtr<()>, *const u8, usize),
} }
@@ -179,7 +189,7 @@ impl Bytes {
/// assert_eq!(b.len(), 5); /// assert_eq!(b.len(), 5);
/// ``` /// ```
#[inline] #[inline]
pub fn len(&self) -> usize { pub const fn len(&self) -> usize {
self.len self.len
} }
@@ -194,7 +204,7 @@ impl Bytes {
/// assert!(b.is_empty()); /// assert!(b.is_empty());
/// ``` /// ```
#[inline] #[inline]
pub fn is_empty(&self) -> bool { pub const fn is_empty(&self) -> bool {
self.len == 0 self.len == 0
} }
@@ -262,7 +272,7 @@ impl Bytes {
let mut ret = self.clone(); let mut ret = self.clone();
ret.len = end - begin; ret.len = end - begin;
ret.ptr = unsafe { ret.ptr.offset(begin as isize) }; ret.ptr = unsafe { ret.ptr.add(begin) };
ret ret
} }
@@ -308,15 +318,15 @@ impl Bytes {
assert!( assert!(
sub_p >= bytes_p, sub_p >= bytes_p,
"subset pointer ({:p}) is smaller than self pointer ({:p})", "subset pointer ({:p}) is smaller than self pointer ({:p})",
sub_p as *const u8, subset.as_ptr(),
bytes_p as *const u8, self.as_ptr(),
); );
assert!( assert!(
sub_p + sub_len <= bytes_p + bytes_len, sub_p + sub_len <= bytes_p + bytes_len,
"subset is out of bounds: self = ({:p}, {}), subset = ({:p}, {})", "subset is out of bounds: self = ({:p}, {}), subset = ({:p}, {})",
bytes_p as *const u8, self.as_ptr(),
bytes_len, bytes_len,
sub_p as *const u8, subset.as_ptr(),
sub_len, sub_len,
); );
@@ -501,7 +511,7 @@ impl Bytes {
// should already be asserted, but debug assert for tests // should already be asserted, but debug assert for tests
debug_assert!(self.len >= by, "internal: inc_start out of bounds"); debug_assert!(self.len >= by, "internal: inc_start out of bounds");
self.len -= by; self.len -= by;
self.ptr = self.ptr.offset(by as isize); self.ptr = self.ptr.add(by);
} }
} }
@@ -530,7 +540,7 @@ impl Buf for Bytes {
} }
#[inline] #[inline]
fn bytes(&self) -> &[u8] { fn chunk(&self) -> &[u8] {
self.as_slice() self.as_slice()
} }
@@ -604,7 +614,7 @@ impl<'a> IntoIterator for &'a Bytes {
type IntoIter = core::slice::Iter<'a, u8>; type IntoIter = core::slice::Iter<'a, u8>;
fn into_iter(self) -> Self::IntoIter { fn into_iter(self) -> Self::IntoIter {
self.as_slice().into_iter() self.as_slice().iter()
} }
} }
@@ -686,7 +696,7 @@ impl PartialOrd<Bytes> for str {
impl PartialEq<Vec<u8>> for Bytes { impl PartialEq<Vec<u8>> for Bytes {
fn eq(&self, other: &Vec<u8>) -> bool { fn eq(&self, other: &Vec<u8>) -> bool {
*self == &other[..] *self == other[..]
} }
} }
@@ -710,7 +720,7 @@ impl PartialOrd<Bytes> for Vec<u8> {
impl PartialEq<String> for Bytes { impl PartialEq<String> for Bytes {
fn eq(&self, other: &String) -> bool { fn eq(&self, other: &String) -> bool {
*self == &other[..] *self == other[..]
} }
} }
@@ -797,31 +807,36 @@ impl From<&'static str> for Bytes {
impl From<Vec<u8>> for Bytes { impl From<Vec<u8>> for Bytes {
fn from(vec: Vec<u8>) -> Bytes { fn from(vec: Vec<u8>) -> Bytes {
// into_boxed_slice doesn't return a heap allocation for empty vectors, let slice = vec.into_boxed_slice();
slice.into()
}
}
impl From<Box<[u8]>> for Bytes {
fn from(slice: Box<[u8]>) -> Bytes {
// Box<[u8]> doesn't contain a heap allocation for empty slices,
// so the pointer isn't aligned enough for the KIND_VEC stashing to // so the pointer isn't aligned enough for the KIND_VEC stashing to
// work. // work.
if vec.is_empty() { if slice.is_empty() {
return Bytes::new(); return Bytes::new();
} }
let slice = vec.into_boxed_slice();
let len = slice.len(); let len = slice.len();
let ptr = slice.as_ptr(); let ptr = Box::into_raw(slice) as *mut u8;
drop(Box::into_raw(slice));
if ptr as usize & 0x1 == 0 { if ptr as usize & 0x1 == 0 {
let data = ptr as usize | KIND_VEC; let data = ptr_map(ptr, |addr| addr | KIND_VEC);
Bytes { Bytes {
ptr, ptr,
len, len,
data: AtomicPtr::new(data as *mut _), data: AtomicPtr::new(data.cast()),
vtable: &PROMOTABLE_EVEN_VTABLE, vtable: &PROMOTABLE_EVEN_VTABLE,
} }
} else { } else {
Bytes { Bytes {
ptr, ptr,
len, len,
data: AtomicPtr::new(ptr as *mut _), data: AtomicPtr::new(ptr.cast()),
vtable: &PROMOTABLE_ODD_VTABLE, vtable: &PROMOTABLE_ODD_VTABLE,
} }
} }
@@ -834,6 +849,13 @@ impl From<String> for Bytes {
} }
} }
impl From<Bytes> for Vec<u8> {
fn from(bytes: Bytes) -> Vec<u8> {
let bytes = mem::ManuallyDrop::new(bytes);
unsafe { (bytes.vtable.to_vec)(&bytes.data, bytes.ptr, bytes.len) }
}
}
// ===== impl Vtable ===== // ===== impl Vtable =====
impl fmt::Debug for Vtable { impl fmt::Debug for Vtable {
@@ -849,6 +871,7 @@ impl fmt::Debug for Vtable {
const STATIC_VTABLE: Vtable = Vtable { const STATIC_VTABLE: Vtable = Vtable {
clone: static_clone, clone: static_clone,
to_vec: static_to_vec,
drop: static_drop, drop: static_drop,
}; };
@@ -857,6 +880,11 @@ unsafe fn static_clone(_: &AtomicPtr<()>, ptr: *const u8, len: usize) -> Bytes {
Bytes::from_static(slice) Bytes::from_static(slice)
} }
unsafe fn static_to_vec(_: &AtomicPtr<()>, ptr: *const u8, len: usize) -> Vec<u8> {
let slice = slice::from_raw_parts(ptr, len);
slice.to_vec()
}
unsafe fn static_drop(_: &mut AtomicPtr<()>, _: *const u8, _: usize) { unsafe fn static_drop(_: &mut AtomicPtr<()>, _: *const u8, _: usize) {
// nothing to drop for &'static [u8] // nothing to drop for &'static [u8]
} }
@@ -865,11 +893,13 @@ unsafe fn static_drop(_: &mut AtomicPtr<()>, _: *const u8, _: usize) {
static PROMOTABLE_EVEN_VTABLE: Vtable = Vtable { static PROMOTABLE_EVEN_VTABLE: Vtable = Vtable {
clone: promotable_even_clone, clone: promotable_even_clone,
to_vec: promotable_even_to_vec,
drop: promotable_even_drop, drop: promotable_even_drop,
}; };
static PROMOTABLE_ODD_VTABLE: Vtable = Vtable { static PROMOTABLE_ODD_VTABLE: Vtable = Vtable {
clone: promotable_odd_clone, clone: promotable_odd_clone,
to_vec: promotable_odd_to_vec,
drop: promotable_odd_drop, drop: promotable_odd_drop,
}; };
@@ -878,25 +908,57 @@ unsafe fn promotable_even_clone(data: &AtomicPtr<()>, ptr: *const u8, len: usize
let kind = shared as usize & KIND_MASK; let kind = shared as usize & KIND_MASK;
if kind == KIND_ARC { if kind == KIND_ARC {
shallow_clone_arc(shared as _, ptr, len) shallow_clone_arc(shared.cast(), ptr, len)
} else { } else {
debug_assert_eq!(kind, KIND_VEC); debug_assert_eq!(kind, KIND_VEC);
let buf = (shared as usize & !KIND_MASK) as *mut u8; let buf = ptr_map(shared.cast(), |addr| addr & !KIND_MASK);
shallow_clone_vec(data, shared, buf, ptr, len) shallow_clone_vec(data, shared, buf, ptr, len)
} }
} }
unsafe fn promotable_to_vec(
data: &AtomicPtr<()>,
ptr: *const u8,
len: usize,
f: fn(*mut ()) -> *mut u8,
) -> Vec<u8> {
let shared = data.load(Ordering::Acquire);
let kind = shared as usize & KIND_MASK;
if kind == KIND_ARC {
shared_to_vec_impl(shared.cast(), ptr, len)
} else {
// If Bytes holds a Vec, then the offset must be 0.
debug_assert_eq!(kind, KIND_VEC);
let buf = f(shared);
let cap = (ptr as usize - buf as usize) + len;
// Copy back buffer
ptr::copy(ptr, buf, len);
Vec::from_raw_parts(buf, len, cap)
}
}
unsafe fn promotable_even_to_vec(data: &AtomicPtr<()>, ptr: *const u8, len: usize) -> Vec<u8> {
promotable_to_vec(data, ptr, len, |shared| {
ptr_map(shared.cast(), |addr| addr & !KIND_MASK)
})
}
unsafe fn promotable_even_drop(data: &mut AtomicPtr<()>, ptr: *const u8, len: usize) { unsafe fn promotable_even_drop(data: &mut AtomicPtr<()>, ptr: *const u8, len: usize) {
data.with_mut(|shared| { data.with_mut(|shared| {
let shared = *shared; let shared = *shared;
let kind = shared as usize & KIND_MASK; let kind = shared as usize & KIND_MASK;
if kind == KIND_ARC { if kind == KIND_ARC {
release_shared(shared as *mut Shared); release_shared(shared.cast());
} else { } else {
debug_assert_eq!(kind, KIND_VEC); debug_assert_eq!(kind, KIND_VEC);
let buf = (shared as usize & !KIND_MASK) as *mut u8; let buf = ptr_map(shared.cast(), |addr| addr & !KIND_MASK);
drop(rebuild_boxed_slice(buf, ptr, len)); free_boxed_slice(buf, ptr, len);
} }
}); });
} }
@@ -909,38 +971,49 @@ unsafe fn promotable_odd_clone(data: &AtomicPtr<()>, ptr: *const u8, len: usize)
shallow_clone_arc(shared as _, ptr, len) shallow_clone_arc(shared as _, ptr, len)
} else { } else {
debug_assert_eq!(kind, KIND_VEC); debug_assert_eq!(kind, KIND_VEC);
shallow_clone_vec(data, shared, shared as *mut u8, ptr, len) shallow_clone_vec(data, shared, shared.cast(), ptr, len)
} }
} }
unsafe fn promotable_odd_to_vec(data: &AtomicPtr<()>, ptr: *const u8, len: usize) -> Vec<u8> {
promotable_to_vec(data, ptr, len, |shared| shared.cast())
}
unsafe fn promotable_odd_drop(data: &mut AtomicPtr<()>, ptr: *const u8, len: usize) { unsafe fn promotable_odd_drop(data: &mut AtomicPtr<()>, ptr: *const u8, len: usize) {
data.with_mut(|shared| { data.with_mut(|shared| {
let shared = *shared; let shared = *shared;
let kind = shared as usize & KIND_MASK; let kind = shared as usize & KIND_MASK;
if kind == KIND_ARC { if kind == KIND_ARC {
release_shared(shared as *mut Shared); release_shared(shared.cast());
} else { } else {
debug_assert_eq!(kind, KIND_VEC); debug_assert_eq!(kind, KIND_VEC);
drop(rebuild_boxed_slice(shared as *mut u8, ptr, len)); free_boxed_slice(shared.cast(), ptr, len);
} }
}); });
} }
unsafe fn rebuild_boxed_slice(buf: *mut u8, offset: *const u8, len: usize) -> Box<[u8]> { unsafe fn free_boxed_slice(buf: *mut u8, offset: *const u8, len: usize) {
let cap = (offset as usize - buf as usize) + len; let cap = (offset as usize - buf as usize) + len;
Box::from_raw(slice::from_raw_parts_mut(buf, cap)) dealloc(buf, Layout::from_size_align(cap, 1).unwrap())
} }
// ===== impl SharedVtable ===== // ===== impl SharedVtable =====
struct Shared { struct Shared {
// holds vec for drop, but otherwise doesnt access it // Holds arguments to dealloc upon Drop, but otherwise doesn't use them
_vec: Vec<u8>, buf: *mut u8,
cap: usize,
ref_cnt: AtomicUsize, ref_cnt: AtomicUsize,
} }
impl Drop for Shared {
fn drop(&mut self) {
unsafe { dealloc(self.buf, Layout::from_size_align(self.cap, 1).unwrap()) }
}
}
// Assert that the alignment of `Shared` is divisible by 2. // Assert that the alignment of `Shared` is divisible by 2.
// This is a necessary invariant since we depend on allocating `Shared` a // This is a necessary invariant since we depend on allocating `Shared` a
// shared object to implicitly carry the `KIND_ARC` flag in its pointer. // shared object to implicitly carry the `KIND_ARC` flag in its pointer.
@@ -949,6 +1022,7 @@ const _: [(); 0 - mem::align_of::<Shared>() % 2] = []; // Assert that the alignm
static SHARED_VTABLE: Vtable = Vtable { static SHARED_VTABLE: Vtable = Vtable {
clone: shared_clone, clone: shared_clone,
to_vec: shared_to_vec,
drop: shared_drop, drop: shared_drop,
}; };
@@ -961,9 +1035,42 @@ unsafe fn shared_clone(data: &AtomicPtr<()>, ptr: *const u8, len: usize) -> Byte
shallow_clone_arc(shared as _, ptr, len) shallow_clone_arc(shared as _, ptr, len)
} }
unsafe fn shared_to_vec_impl(shared: *mut Shared, ptr: *const u8, len: usize) -> Vec<u8> {
// Check that the ref_cnt is 1 (unique).
//
// If it is unique, then it is set to 0 with AcqRel fence for the same
// reason in release_shared.
//
// Otherwise, we take the other branch and call release_shared.
if (*shared)
.ref_cnt
.compare_exchange(1, 0, Ordering::AcqRel, Ordering::Relaxed)
.is_ok()
{
let buf = (*shared).buf;
let cap = (*shared).cap;
// Deallocate Shared
drop(Box::from_raw(shared as *mut mem::ManuallyDrop<Shared>));
// Copy back buffer
ptr::copy(ptr, buf, len);
Vec::from_raw_parts(buf, len, cap)
} else {
let v = slice::from_raw_parts(ptr, len).to_vec();
release_shared(shared);
v
}
}
unsafe fn shared_to_vec(data: &AtomicPtr<()>, ptr: *const u8, len: usize) -> Vec<u8> {
shared_to_vec_impl(data.load(Ordering::Relaxed).cast(), ptr, len)
}
unsafe fn shared_drop(data: &mut AtomicPtr<()>, _ptr: *const u8, _len: usize) { unsafe fn shared_drop(data: &mut AtomicPtr<()>, _ptr: *const u8, _len: usize) {
data.with_mut(|shared| { data.with_mut(|shared| {
release_shared(*shared as *mut Shared); release_shared(shared.cast());
}); });
} }
@@ -1001,9 +1108,9 @@ unsafe fn shallow_clone_vec(
// updated and since the buffer hasn't been promoted to an // updated and since the buffer hasn't been promoted to an
// `Arc`, those three fields still are the components of the // `Arc`, those three fields still are the components of the
// vector. // vector.
let vec = rebuild_boxed_slice(buf, offset, len).into_vec();
let shared = Box::new(Shared { let shared = Box::new(Shared {
_vec: vec, buf,
cap: (offset as usize - buf as usize) + len,
// Initialize refcount to 2. One for this reference, and one // Initialize refcount to 2. One for this reference, and one
// for the new clone that will be returned from // for the new clone that will be returned from
// `shallow_clone`. // `shallow_clone`.
@@ -1023,33 +1130,35 @@ unsafe fn shallow_clone_vec(
// `Release` is used synchronize with other threads that // `Release` is used synchronize with other threads that
// will load the `arc` field. // will load the `arc` field.
// //
// If the `compare_and_swap` fails, then the thread lost the // If the `compare_exchange` fails, then the thread lost the
// race to promote the buffer to shared. The `Acquire` // race to promote the buffer to shared. The `Acquire`
// ordering will synchronize with the `compare_and_swap` // ordering will synchronize with the `compare_exchange`
// that happened in the other thread and the `Shared` // that happened in the other thread and the `Shared`
// pointed to by `actual` will be visible. // pointed to by `actual` will be visible.
let actual = atom.compare_and_swap(ptr as _, shared as _, Ordering::AcqRel); match atom.compare_exchange(ptr as _, shared as _, Ordering::AcqRel, Ordering::Acquire) {
Ok(actual) => {
debug_assert!(actual as usize == ptr as usize);
// The upgrade was successful, the new handle can be
// returned.
Bytes {
ptr: offset,
len,
data: AtomicPtr::new(shared as _),
vtable: &SHARED_VTABLE,
}
}
Err(actual) => {
// The upgrade failed, a concurrent clone happened. Release
// the allocation that was made in this thread, it will not
// be needed.
let shared = Box::from_raw(shared);
mem::forget(*shared);
if actual as usize == ptr as usize { // Buffer already promoted to shared storage, so increment ref
// The upgrade was successful, the new handle can be // count.
// returned. shallow_clone_arc(actual as _, offset, len)
return Bytes { }
ptr: offset,
len,
data: AtomicPtr::new(shared as _),
vtable: &SHARED_VTABLE,
};
} }
// The upgrade failed, a concurrent clone happened. Release
// the allocation that was made in this thread, it will not
// be needed.
let shared = Box::from_raw(shared);
mem::forget(*shared);
// Buffer already promoted to shared storage, so increment ref
// count.
shallow_clone_arc(actual as _, offset, len)
} }
unsafe fn release_shared(ptr: *mut Shared) { unsafe fn release_shared(ptr: *mut Shared) {
@@ -1075,10 +1184,40 @@ unsafe fn release_shared(ptr: *mut Shared) {
// > "acquire" operation before deleting the object. // > "acquire" operation before deleting the object.
// //
// [1]: (www.boost.org/doc/libs/1_55_0/doc/html/atomic/usage_examples.html) // [1]: (www.boost.org/doc/libs/1_55_0/doc/html/atomic/usage_examples.html)
atomic::fence(Ordering::Acquire); //
// Thread sanitizer does not support atomic fences. Use an atomic load
// instead.
(*ptr).ref_cnt.load(Ordering::Acquire);
// Drop the data // Drop the data
Box::from_raw(ptr); drop(Box::from_raw(ptr));
}
// Ideally we would always use this version of `ptr_map` since it is strict
// provenance compatible, but it results in worse codegen. We will however still
// use it on miri because it gives better diagnostics for people who test bytes
// code with miri.
//
// See https://github.com/tokio-rs/bytes/pull/545 for more info.
#[cfg(miri)]
fn ptr_map<F>(ptr: *mut u8, f: F) -> *mut u8
where
F: FnOnce(usize) -> usize,
{
let old_addr = ptr as usize;
let new_addr = f(old_addr);
let diff = new_addr.wrapping_sub(old_addr);
ptr.wrapping_add(diff)
}
#[cfg(not(miri))]
fn ptr_map<F>(ptr: *mut u8, f: F) -> *mut u8
where
F: FnOnce(usize) -> usize,
{
let old_addr = ptr as usize;
let new_addr = f(old_addr);
new_addr as *mut u8
} }
// compile-fails // compile-fails
+248 -58
View File
@@ -8,6 +8,7 @@ use alloc::{
borrow::{Borrow, BorrowMut}, borrow::{Borrow, BorrowMut},
boxed::Box, boxed::Box,
string::String, string::String,
vec,
vec::Vec, vec::Vec,
}; };
@@ -15,7 +16,7 @@ use crate::buf::{IntoIter, UninitSlice};
use crate::bytes::Vtable; use crate::bytes::Vtable;
#[allow(unused)] #[allow(unused)]
use crate::loom::sync::atomic::AtomicMut; use crate::loom::sync::atomic::AtomicMut;
use crate::loom::sync::atomic::{self, AtomicPtr, AtomicUsize, Ordering}; use crate::loom::sync::atomic::{AtomicPtr, AtomicUsize, Ordering};
use crate::{Buf, BufMut, Bytes}; use crate::{Buf, BufMut, Bytes};
/// A unique reference to a contiguous slice of memory. /// A unique reference to a contiguous slice of memory.
@@ -252,12 +253,28 @@ impl BytesMut {
let ptr = self.ptr.as_ptr(); let ptr = self.ptr.as_ptr();
let len = self.len; let len = self.len;
let data = AtomicPtr::new(self.data as _); let data = AtomicPtr::new(self.data.cast());
mem::forget(self); mem::forget(self);
unsafe { Bytes::with_vtable(ptr, len, data, &SHARED_VTABLE) } unsafe { Bytes::with_vtable(ptr, len, data, &SHARED_VTABLE) }
} }
} }
/// Creates a new `BytesMut`, which is initialized with zero.
///
/// # Examples
///
/// ```
/// use bytes::BytesMut;
///
/// let zeros = BytesMut::zeroed(42);
///
/// assert_eq!(zeros.len(), 42);
/// zeros.into_iter().for_each(|x| assert_eq!(x, 0));
/// ```
pub fn zeroed(len: usize) -> BytesMut {
BytesMut::from_vec(vec![0; len])
}
/// Splits the bytes into two at the given index. /// Splits the bytes into two at the given index.
/// ///
/// Afterwards `self` contains elements `[0, at)`, and the returned /// Afterwards `self` contains elements `[0, at)`, and the returned
@@ -380,6 +397,8 @@ impl BytesMut {
/// If `len` is greater than the buffer's current length, this has no /// If `len` is greater than the buffer's current length, this has no
/// effect. /// effect.
/// ///
/// Existing underlying capacity is preserved.
///
/// The [`split_off`] method can emulate `truncate`, but this causes the /// The [`split_off`] method can emulate `truncate`, but this causes the
/// excess bytes to be returned instead of dropped. /// excess bytes to be returned instead of dropped.
/// ///
@@ -402,7 +421,7 @@ impl BytesMut {
} }
} }
/// Clears the buffer, removing all data. /// Clears the buffer, removing all data. Existing capacity is preserved.
/// ///
/// # Examples /// # Examples
/// ///
@@ -445,7 +464,7 @@ impl BytesMut {
let additional = new_len - len; let additional = new_len - len;
self.reserve(additional); self.reserve(additional);
unsafe { unsafe {
let dst = self.bytes_mut().as_mut_ptr(); let dst = self.chunk_mut().as_mut_ptr();
ptr::write_bytes(dst, value, additional); ptr::write_bytes(dst, value, additional);
self.set_len(new_len); self.set_len(new_len);
} }
@@ -492,11 +511,20 @@ impl BytesMut {
/// reallocations. A call to `reserve` may result in an allocation. /// reallocations. A call to `reserve` may result in an allocation.
/// ///
/// Before allocating new buffer space, the function will attempt to reclaim /// Before allocating new buffer space, the function will attempt to reclaim
/// space in the existing buffer. If the current handle references a small /// space in the existing buffer. If the current handle references a view
/// view in the original buffer and all other handles have been dropped, /// into a larger original buffer, and all other handles referencing part
/// and the requested capacity is less than or equal to the existing /// of the same original buffer have been dropped, then the current view
/// buffer's capacity, then the current view will be copied to the front of /// can be copied/shifted to the front of the buffer and the handle can take
/// the buffer and the handle will take ownership of the full buffer. /// ownership of the full buffer, provided that the full buffer is large
/// enough to fit the requested additional capacity.
///
/// This optimization will only happen if shifting the data from the current
/// view to the front of the buffer is not too expensive in terms of the
/// (amortized) time required. The precise condition is subject to change;
/// as of now, the length of the data being shifted needs to be at least as
/// large as the distance that it's shifted by. If the current view is empty
/// and the original buffer is large enough to fit the requested additional
/// capacity, then reallocations will never happen.
/// ///
/// # Examples /// # Examples
/// ///
@@ -560,17 +588,34 @@ impl BytesMut {
// space. // space.
// //
// Otherwise, since backed by a vector, use `Vec::reserve` // Otherwise, since backed by a vector, use `Vec::reserve`
//
// We need to make sure that this optimization does not kill the
// amortized runtimes of BytesMut's operations.
unsafe { unsafe {
let (off, prev) = self.get_vec_pos(); let (off, prev) = self.get_vec_pos();
// Only reuse space if we can satisfy the requested additional space. // Only reuse space if we can satisfy the requested additional space.
if self.capacity() - self.len() + off >= additional { //
// There's space - reuse it // Also check if the value of `off` suggests that enough bytes
// have been read to account for the overhead of shifting all
// the data (in an amortized analysis).
// Hence the condition `off >= self.len()`.
//
// This condition also already implies that the buffer is going
// to be (at least) half-empty in the end; so we do not break
// the (amortized) runtime with future resizes of the underlying
// `Vec`.
//
// [For more details check issue #524, and PR #525.]
if self.capacity() - self.len() + off >= additional && off >= self.len() {
// There's enough space, and it's not too much overhead:
// reuse the space!
// //
// Just move the pointer back to the start after copying // Just move the pointer back to the start after copying
// data back. // data back.
let base_ptr = self.ptr.as_ptr().offset(-(off as isize)); let base_ptr = self.ptr.as_ptr().offset(-(off as isize));
ptr::copy(self.ptr.as_ptr(), base_ptr, self.len); // Since `off >= self.len()`, the two regions don't overlap.
ptr::copy_nonoverlapping(self.ptr.as_ptr(), base_ptr, self.len);
self.ptr = vptr(base_ptr); self.ptr = vptr(base_ptr);
self.set_vec_pos(0, prev); self.set_vec_pos(0, prev);
@@ -578,13 +623,14 @@ impl BytesMut {
// can gain capacity back. // can gain capacity back.
self.cap += off; self.cap += off;
} else { } else {
// No space - allocate more // Not enough space, or reusing might be too much overhead:
// allocate more space!
let mut v = let mut v =
ManuallyDrop::new(rebuild_vec(self.ptr.as_ptr(), self.len, self.cap, off)); ManuallyDrop::new(rebuild_vec(self.ptr.as_ptr(), self.len, self.cap, off));
v.reserve(additional); v.reserve(additional);
// Update the info // Update the info
self.ptr = vptr(v.as_mut_ptr().offset(off as isize)); self.ptr = vptr(v.as_mut_ptr().add(off));
self.len = v.len() - off; self.len = v.len() - off;
self.cap = v.capacity() - off; self.cap = v.capacity() - off;
} }
@@ -594,7 +640,7 @@ impl BytesMut {
} }
debug_assert_eq!(kind, KIND_ARC); debug_assert_eq!(kind, KIND_ARC);
let shared: *mut Shared = self.data as _; let shared: *mut Shared = self.data;
// Reserving involves abandoning the currently shared buffer and // Reserving involves abandoning the currently shared buffer and
// allocating a new vector with the requested capacity. // allocating a new vector with the requested capacity.
@@ -617,29 +663,53 @@ impl BytesMut {
// sure that the vector has enough capacity. // sure that the vector has enough capacity.
let v = &mut (*shared).vec; let v = &mut (*shared).vec;
if v.capacity() >= new_cap { let v_capacity = v.capacity();
// The capacity is sufficient, reclaim the buffer let ptr = v.as_mut_ptr();
let ptr = v.as_mut_ptr();
ptr::copy(self.ptr.as_ptr(), ptr, len); let offset = offset_from(self.ptr.as_ptr(), ptr);
// Compare the condition in the `kind == KIND_VEC` case above
// for more details.
if v_capacity >= new_cap && offset >= len {
// The capacity is sufficient, and copying is not too much
// overhead: reclaim the buffer!
// `offset >= len` means: no overlap
ptr::copy_nonoverlapping(self.ptr.as_ptr(), ptr, len);
self.ptr = vptr(ptr); self.ptr = vptr(ptr);
self.cap = v.capacity(); self.cap = v.capacity();
} else {
// calculate offset
let off = (self.ptr.as_ptr() as usize) - (v.as_ptr() as usize);
return; // new_cap is calculated in terms of `BytesMut`, not the underlying
// `Vec`, so it does not take the offset into account.
//
// Thus we have to manually add it here.
new_cap = new_cap.checked_add(off).expect("overflow");
// The vector capacity is not sufficient. The reserve request is
// asking for more than the initial buffer capacity. Allocate more
// than requested if `new_cap` is not much bigger than the current
// capacity.
//
// There are some situations, using `reserve_exact` that the
// buffer capacity could be below `original_capacity`, so do a
// check.
let double = v.capacity().checked_shl(1).unwrap_or(new_cap);
new_cap = cmp::max(double, new_cap);
// No space - allocate more
v.reserve(new_cap - v.len());
// Update the info
self.ptr = vptr(v.as_mut_ptr().add(off));
self.cap = v.capacity() - off;
} }
// The vector capacity is not sufficient. The reserve request is return;
// asking for more than the initial buffer capacity. Allocate more
// than requested if `new_cap` is not much bigger than the current
// capacity.
//
// There are some situations, using `reserve_exact` that the
// buffer capacity could be below `original_capacity`, so do a
// check.
let double = v.capacity().checked_shl(1).unwrap_or(new_cap);
new_cap = cmp::max(cmp::max(double, new_cap), original_capacity);
} else { } else {
new_cap = cmp::max(new_cap, original_capacity); new_cap = cmp::max(new_cap, original_capacity);
} }
@@ -657,7 +727,7 @@ impl BytesMut {
// Update self // Update self
let data = (original_capacity_repr << ORIGINAL_CAPACITY_OFFSET) | KIND_VEC; let data = (original_capacity_repr << ORIGINAL_CAPACITY_OFFSET) | KIND_VEC;
self.data = data as _; self.data = invalid_ptr(data);
self.ptr = vptr(v.as_mut_ptr()); self.ptr = vptr(v.as_mut_ptr());
self.len = v.len(); self.len = v.len();
self.cap = v.capacity(); self.cap = v.capacity();
@@ -688,7 +758,7 @@ impl BytesMut {
// Reserved above // Reserved above
debug_assert!(dst.len() >= cnt); debug_assert!(dst.len() >= cnt);
ptr::copy_nonoverlapping(extend.as_ptr(), dst.as_mut_ptr() as *mut u8, cnt); ptr::copy_nonoverlapping(extend.as_ptr(), dst.as_mut_ptr(), cnt);
} }
unsafe { unsafe {
@@ -698,10 +768,11 @@ impl BytesMut {
/// Absorbs a `BytesMut` that was previously split off. /// Absorbs a `BytesMut` that was previously split off.
/// ///
/// If the two `BytesMut` objects were previously contiguous, i.e., if /// If the two `BytesMut` objects were previously contiguous and not mutated
/// `other` was created by calling `split_off` on this `BytesMut`, then /// in a way that causes re-allocation i.e., if `other` was created by
/// this is an `O(1)` operation that just decreases a reference /// calling `split_off` on this `BytesMut`, then this is an `O(1)` operation
/// count and sets a few indices. Otherwise this method degenerates to /// that just decreases a reference count and sets a few indices.
/// Otherwise this method degenerates to
/// `self.extend_from_slice(other.as_ref())`. /// `self.extend_from_slice(other.as_ref())`.
/// ///
/// # Examples /// # Examples
@@ -752,7 +823,7 @@ impl BytesMut {
ptr, ptr,
len, len,
cap, cap,
data: data as *mut _, data: invalid_ptr(data),
} }
} }
@@ -799,7 +870,7 @@ impl BytesMut {
// Updating the start of the view is setting `ptr` to point to the // Updating the start of the view is setting `ptr` to point to the
// new start and updating the `len` field to reflect the new length // new start and updating the `len` field to reflect the new length
// of the view. // of the view.
self.ptr = vptr(self.ptr.as_ptr().offset(start as isize)); self.ptr = vptr(self.ptr.as_ptr().add(start));
if self.len >= start { if self.len >= start {
self.len -= start; self.len -= start;
@@ -819,11 +890,11 @@ impl BytesMut {
} }
fn try_unsplit(&mut self, other: BytesMut) -> Result<(), BytesMut> { fn try_unsplit(&mut self, other: BytesMut) -> Result<(), BytesMut> {
if other.is_empty() { if other.capacity() == 0 {
return Ok(()); return Ok(());
} }
let ptr = unsafe { self.ptr.as_ptr().offset(self.len as isize) }; let ptr = unsafe { self.ptr.as_ptr().add(self.len) };
if ptr == other.ptr.as_ptr() if ptr == other.ptr.as_ptr()
&& self.kind() == KIND_ARC && self.kind() == KIND_ARC
&& other.kind() == KIND_ARC && other.kind() == KIND_ARC
@@ -873,7 +944,7 @@ impl BytesMut {
// always succeed. // always succeed.
debug_assert_eq!(shared as usize & KIND_MASK, KIND_ARC); debug_assert_eq!(shared as usize & KIND_MASK, KIND_ARC);
self.data = shared as _; self.data = shared;
} }
/// Makes an exact shallow clone of `self`. /// Makes an exact shallow clone of `self`.
@@ -906,13 +977,13 @@ impl BytesMut {
debug_assert_eq!(self.kind(), KIND_VEC); debug_assert_eq!(self.kind(), KIND_VEC);
debug_assert!(pos <= MAX_VEC_POS); debug_assert!(pos <= MAX_VEC_POS);
self.data = ((pos << VEC_POS_OFFSET) | (prev & NOT_VEC_POS_MASK)) as *mut _; self.data = invalid_ptr((pos << VEC_POS_OFFSET) | (prev & NOT_VEC_POS_MASK));
} }
#[inline] #[inline]
fn uninit_slice(&mut self) -> &mut UninitSlice { fn uninit_slice(&mut self) -> &mut UninitSlice {
unsafe { unsafe {
let ptr = self.ptr.as_ptr().offset(self.len as isize); let ptr = self.ptr.as_ptr().add(self.len);
let len = self.cap - self.len; let len = self.cap - self.len;
UninitSlice::from_raw_parts_mut(ptr, len) UninitSlice::from_raw_parts_mut(ptr, len)
@@ -932,7 +1003,7 @@ impl Drop for BytesMut {
let _ = rebuild_vec(self.ptr.as_ptr(), self.len, self.cap, off); let _ = rebuild_vec(self.ptr.as_ptr(), self.len, self.cap, off);
} }
} else if kind == KIND_ARC { } else if kind == KIND_ARC {
unsafe { release_shared(self.data as _) }; unsafe { release_shared(self.data) };
} }
} }
} }
@@ -944,7 +1015,7 @@ impl Buf for BytesMut {
} }
#[inline] #[inline]
fn bytes(&self) -> &[u8] { fn chunk(&self) -> &[u8] {
self.as_slice() self.as_slice()
} }
@@ -985,7 +1056,7 @@ unsafe impl BufMut for BytesMut {
} }
#[inline] #[inline]
fn bytes_mut(&mut self) -> &mut UninitSlice { fn chunk_mut(&mut self) -> &mut UninitSlice {
if self.capacity() == self.len() { if self.capacity() == self.len() {
self.reserve(64); self.reserve(64);
} }
@@ -1000,7 +1071,7 @@ unsafe impl BufMut for BytesMut {
Self: Sized, Self: Sized,
{ {
while src.has_remaining() { while src.has_remaining() {
let s = src.bytes(); let s = src.chunk();
let l = s.len(); let l = s.len();
self.extend_from_slice(s); self.extend_from_slice(s);
src.advance(l); src.advance(l);
@@ -1010,6 +1081,19 @@ unsafe impl BufMut for BytesMut {
fn put_slice(&mut self, src: &[u8]) { fn put_slice(&mut self, src: &[u8]) {
self.extend_from_slice(src); self.extend_from_slice(src);
} }
fn put_bytes(&mut self, val: u8, cnt: usize) {
self.reserve(cnt);
unsafe {
let dst = self.uninit_slice();
// Reserved above
debug_assert!(dst.len() >= cnt);
ptr::write_bytes(dst.as_mut_ptr(), val, cnt);
self.advance_mut(cnt);
}
}
} }
impl AsRef<[u8]> for BytesMut { impl AsRef<[u8]> for BytesMut {
@@ -1146,7 +1230,7 @@ impl<'a> IntoIterator for &'a BytesMut {
type IntoIter = core::slice::Iter<'a, u8>; type IntoIter = core::slice::Iter<'a, u8>;
fn into_iter(self) -> Self::IntoIter { fn into_iter(self) -> Self::IntoIter {
self.as_ref().into_iter() self.as_ref().iter()
} }
} }
@@ -1175,7 +1259,18 @@ impl<'a> Extend<&'a u8> for BytesMut {
where where
T: IntoIterator<Item = &'a u8>, T: IntoIterator<Item = &'a u8>,
{ {
self.extend(iter.into_iter().map(|b| *b)) self.extend(iter.into_iter().copied())
}
}
impl Extend<Bytes> for BytesMut {
fn extend<T>(&mut self, iter: T)
where
T: IntoIterator<Item = Bytes>,
{
for bytes in iter {
self.extend_from_slice(&bytes)
}
} }
} }
@@ -1187,7 +1282,7 @@ impl FromIterator<u8> for BytesMut {
impl<'a> FromIterator<&'a u8> for BytesMut { impl<'a> FromIterator<&'a u8> for BytesMut {
fn from_iter<T: IntoIterator<Item = &'a u8>>(into_iter: T) -> Self { fn from_iter<T: IntoIterator<Item = &'a u8>>(into_iter: T) -> Self {
BytesMut::from_iter(into_iter.into_iter().map(|b| *b)) BytesMut::from_iter(into_iter.into_iter().copied())
} }
} }
@@ -1228,10 +1323,13 @@ unsafe fn release_shared(ptr: *mut Shared) {
// > "acquire" operation before deleting the object. // > "acquire" operation before deleting the object.
// //
// [1]: (www.boost.org/doc/libs/1_55_0/doc/html/atomic/usage_examples.html) // [1]: (www.boost.org/doc/libs/1_55_0/doc/html/atomic/usage_examples.html)
atomic::fence(Ordering::Acquire); //
// Thread sanitizer does not support atomic fences. Use an atomic load
// instead.
(*ptr).ref_count.load(Ordering::Acquire);
// Drop the data // Drop the data
Box::from_raw(ptr); drop(Box::from_raw(ptr));
} }
impl Shared { impl Shared {
@@ -1250,6 +1348,7 @@ impl Shared {
} }
} }
#[inline]
fn original_capacity_to_repr(cap: usize) -> usize { fn original_capacity_to_repr(cap: usize) -> usize {
let width = PTR_WIDTH - ((cap >> MIN_ORIGINAL_CAPACITY_WIDTH).leading_zeros() as usize); let width = PTR_WIDTH - ((cap >> MIN_ORIGINAL_CAPACITY_WIDTH).leading_zeros() as usize);
cmp::min( cmp::min(
@@ -1376,7 +1475,7 @@ impl PartialOrd<BytesMut> for str {
impl PartialEq<Vec<u8>> for BytesMut { impl PartialEq<Vec<u8>> for BytesMut {
fn eq(&self, other: &Vec<u8>) -> bool { fn eq(&self, other: &Vec<u8>) -> bool {
*self == &other[..] *self == other[..]
} }
} }
@@ -1400,7 +1499,7 @@ impl PartialOrd<BytesMut> for Vec<u8> {
impl PartialEq<String> for BytesMut { impl PartialEq<String> for BytesMut {
fn eq(&self, other: &String) -> bool { fn eq(&self, other: &String) -> bool {
*self == &other[..] *self == other[..]
} }
} }
@@ -1466,16 +1565,55 @@ impl PartialOrd<BytesMut> for &str {
impl PartialEq<BytesMut> for Bytes { impl PartialEq<BytesMut> for Bytes {
fn eq(&self, other: &BytesMut) -> bool { fn eq(&self, other: &BytesMut) -> bool {
&other[..] == &self[..] other[..] == self[..]
} }
} }
impl PartialEq<Bytes> for BytesMut { impl PartialEq<Bytes> for BytesMut {
fn eq(&self, other: &Bytes) -> bool { fn eq(&self, other: &Bytes) -> bool {
&other[..] == &self[..] other[..] == self[..]
} }
} }
impl From<BytesMut> for Vec<u8> {
fn from(mut bytes: BytesMut) -> Self {
let kind = bytes.kind();
let mut vec = if kind == KIND_VEC {
unsafe {
let (off, _) = bytes.get_vec_pos();
rebuild_vec(bytes.ptr.as_ptr(), bytes.len, bytes.cap, off)
}
} else if kind == KIND_ARC {
let shared = bytes.data as *mut Shared;
if unsafe { (*shared).is_unique() } {
let vec = mem::replace(unsafe { &mut (*shared).vec }, Vec::new());
unsafe { release_shared(shared) };
vec
} else {
return bytes.deref().to_vec();
}
} else {
return bytes.deref().to_vec();
};
let len = bytes.len;
unsafe {
ptr::copy(bytes.ptr.as_ptr(), vec.as_mut_ptr(), len);
vec.set_len(len);
}
mem::forget(bytes);
vec
}
}
#[inline]
fn vptr(ptr: *mut u8) -> NonNull<u8> { fn vptr(ptr: *mut u8) -> NonNull<u8> {
if cfg!(debug_assertions) { if cfg!(debug_assertions) {
NonNull::new(ptr).expect("Vec pointer should be non-null") NonNull::new(ptr).expect("Vec pointer should be non-null")
@@ -1484,6 +1622,35 @@ fn vptr(ptr: *mut u8) -> NonNull<u8> {
} }
} }
/// Returns a dangling pointer with the given address. This is used to store
/// integer data in pointer fields.
///
/// It is equivalent to `addr as *mut T`, but this fails on miri when strict
/// provenance checking is enabled.
#[inline]
fn invalid_ptr<T>(addr: usize) -> *mut T {
let ptr = core::ptr::null_mut::<u8>().wrapping_add(addr);
debug_assert_eq!(ptr as usize, addr);
ptr.cast::<T>()
}
/// Precondition: dst >= original
///
/// The following line is equivalent to:
///
/// ```rust,ignore
/// self.ptr.as_ptr().offset_from(ptr) as usize;
/// ```
///
/// But due to min rust is 1.39 and it is only stablised
/// in 1.47, we cannot use it.
#[inline]
fn offset_from(dst: *mut u8, original: *mut u8) -> usize {
debug_assert!(dst >= original);
dst as usize - original as usize
}
unsafe fn rebuild_vec(ptr: *mut u8, mut len: usize, mut cap: usize, off: usize) -> Vec<u8> { unsafe fn rebuild_vec(ptr: *mut u8, mut len: usize, mut cap: usize, off: usize) -> Vec<u8> {
let ptr = ptr.offset(-(off as isize)); let ptr = ptr.offset(-(off as isize));
len += off; len += off;
@@ -1496,6 +1663,7 @@ unsafe fn rebuild_vec(ptr: *mut u8, mut len: usize, mut cap: usize, off: usize)
static SHARED_VTABLE: Vtable = Vtable { static SHARED_VTABLE: Vtable = Vtable {
clone: shared_v_clone, clone: shared_v_clone,
to_vec: shared_v_to_vec,
drop: shared_v_drop, drop: shared_v_drop,
}; };
@@ -1503,10 +1671,32 @@ unsafe fn shared_v_clone(data: &AtomicPtr<()>, ptr: *const u8, len: usize) -> By
let shared = data.load(Ordering::Relaxed) as *mut Shared; let shared = data.load(Ordering::Relaxed) as *mut Shared;
increment_shared(shared); increment_shared(shared);
let data = AtomicPtr::new(shared as _); let data = AtomicPtr::new(shared as *mut ());
Bytes::with_vtable(ptr, len, data, &SHARED_VTABLE) Bytes::with_vtable(ptr, len, data, &SHARED_VTABLE)
} }
unsafe fn shared_v_to_vec(data: &AtomicPtr<()>, ptr: *const u8, len: usize) -> Vec<u8> {
let shared: *mut Shared = data.load(Ordering::Relaxed).cast();
if (*shared).is_unique() {
let shared = &mut *shared;
// Drop shared
let mut vec = mem::replace(&mut shared.vec, Vec::new());
release_shared(shared);
// Copy back buffer
ptr::copy(ptr, vec.as_mut_ptr(), len);
vec.set_len(len);
vec
} else {
let v = slice::from_raw_parts(ptr, len).to_vec();
release_shared(shared);
v
}
}
unsafe fn shared_v_drop(data: &mut AtomicPtr<()>, _ptr: *const u8, _len: usize) { unsafe fn shared_v_drop(data: &mut AtomicPtr<()>, _ptr: *const u8, _len: usize) {
data.with_mut(|shared| { data.with_mut(|shared| {
release_shared(*shared as *mut Shared); release_shared(*shared as *mut Shared);
+3 -3
View File
@@ -25,7 +25,7 @@ impl Debug for BytesRef<'_> {
} else if b == b'\0' { } else if b == b'\0' {
write!(f, "\\0")?; write!(f, "\\0")?;
// ASCII printable // ASCII printable
} else if b >= 0x20 && b < 0x7f { } else if (0x20..0x7f).contains(&b) {
write!(f, "{}", b as char)?; write!(f, "{}", b as char)?;
} else { } else {
write!(f, "\\x{:02x}", b)?; write!(f, "\\x{:02x}", b)?;
@@ -38,12 +38,12 @@ impl Debug for BytesRef<'_> {
impl Debug for Bytes { impl Debug for Bytes {
fn fmt(&self, f: &mut Formatter<'_>) -> Result { fn fmt(&self, f: &mut Formatter<'_>) -> Result {
Debug::fmt(&BytesRef(&self.as_ref()), f) Debug::fmt(&BytesRef(self.as_ref()), f)
} }
} }
impl Debug for BytesMut { impl Debug for BytesMut {
fn fmt(&self, f: &mut Formatter<'_>) -> Result { fn fmt(&self, f: &mut Formatter<'_>) -> Result {
Debug::fmt(&BytesRef(&self.as_ref()), f) Debug::fmt(&BytesRef(self.as_ref()), f)
} }
} }
-1
View File
@@ -3,7 +3,6 @@
no_crate_inject, no_crate_inject,
attr(deny(warnings, rust_2018_idioms), allow(dead_code, unused_variables)) attr(deny(warnings, rust_2018_idioms), allow(dead_code, unused_variables))
))] ))]
#![doc(html_root_url = "https://docs.rs/bytes/0.6.0")]
#![no_std] #![no_std]
//! Provides abstractions for working with bytes. //! Provides abstractions for working with bytes.
+2 -2
View File
@@ -1,7 +1,7 @@
#[cfg(not(all(test, loom)))] #[cfg(not(all(test, loom)))]
pub(crate) mod sync { pub(crate) mod sync {
pub(crate) mod atomic { pub(crate) mod atomic {
pub(crate) use core::sync::atomic::{fence, AtomicPtr, AtomicUsize, Ordering}; pub(crate) use core::sync::atomic::{AtomicPtr, AtomicUsize, Ordering};
pub(crate) trait AtomicMut<T> { pub(crate) trait AtomicMut<T> {
fn with_mut<F, R>(&mut self, f: F) -> R fn with_mut<F, R>(&mut self, f: F) -> R
@@ -23,7 +23,7 @@ pub(crate) mod sync {
#[cfg(all(test, loom))] #[cfg(all(test, loom))]
pub(crate) mod sync { pub(crate) mod sync {
pub(crate) mod atomic { pub(crate) mod atomic {
pub(crate) use loom::sync::atomic::{fence, AtomicPtr, AtomicUsize, Ordering}; pub(crate) use loom::sync::atomic::{AtomicPtr, AtomicUsize, Ordering};
pub(crate) trait AtomicMut<T> {} pub(crate) trait AtomicMut<T> {}
} }
+8 -8
View File
@@ -9,17 +9,17 @@ fn test_fresh_cursor_vec() {
let mut buf = &b"hello"[..]; let mut buf = &b"hello"[..];
assert_eq!(buf.remaining(), 5); assert_eq!(buf.remaining(), 5);
assert_eq!(buf.bytes(), b"hello"); assert_eq!(buf.chunk(), b"hello");
buf.advance(2); buf.advance(2);
assert_eq!(buf.remaining(), 3); assert_eq!(buf.remaining(), 3);
assert_eq!(buf.bytes(), b"llo"); assert_eq!(buf.chunk(), b"llo");
buf.advance(3); buf.advance(3);
assert_eq!(buf.remaining(), 0); assert_eq!(buf.remaining(), 0);
assert_eq!(buf.bytes(), b""); assert_eq!(buf.chunk(), b"");
} }
#[test] #[test]
@@ -53,7 +53,7 @@ fn test_bufs_vec() {
let mut dst = [IoSlice::new(b1), IoSlice::new(b2)]; let mut dst = [IoSlice::new(b1), IoSlice::new(b2)];
assert_eq!(1, buf.bytes_vectored(&mut dst[..])); assert_eq!(1, buf.chunks_vectored(&mut dst[..]));
} }
#[test] #[test]
@@ -63,9 +63,9 @@ fn test_vec_deque() {
let mut buffer: VecDeque<u8> = VecDeque::new(); let mut buffer: VecDeque<u8> = VecDeque::new();
buffer.extend(b"hello world"); buffer.extend(b"hello world");
assert_eq!(11, buffer.remaining()); assert_eq!(11, buffer.remaining());
assert_eq!(b"hello world", buffer.bytes()); assert_eq!(b"hello world", buffer.chunk());
buffer.advance(6); buffer.advance(6);
assert_eq!(b"world", buffer.bytes()); assert_eq!(b"world", buffer.chunk());
buffer.extend(b" piece"); buffer.extend(b" piece");
let mut out = [0; 11]; let mut out = [0; 11];
buffer.copy_to_slice(&mut out); buffer.copy_to_slice(&mut out);
@@ -81,8 +81,8 @@ fn test_deref_buf_forwards() {
unreachable!("remaining"); unreachable!("remaining");
} }
fn bytes(&self) -> &[u8] { fn chunk(&self) -> &[u8] {
unreachable!("bytes"); unreachable!("chunk");
} }
fn advance(&mut self, _: usize) { fn advance(&mut self, _: usize) {
+54 -5
View File
@@ -9,15 +9,15 @@ use core::usize;
fn test_vec_as_mut_buf() { fn test_vec_as_mut_buf() {
let mut buf = Vec::with_capacity(64); let mut buf = Vec::with_capacity(64);
assert_eq!(buf.remaining_mut(), usize::MAX); assert_eq!(buf.remaining_mut(), isize::MAX as usize);
assert!(buf.bytes_mut().len() >= 64); assert!(buf.chunk_mut().len() >= 64);
buf.put(&b"zomg"[..]); buf.put(&b"zomg"[..]);
assert_eq!(&buf, b"zomg"); assert_eq!(&buf, b"zomg");
assert_eq!(buf.remaining_mut(), usize::MAX - 4); assert_eq!(buf.remaining_mut(), isize::MAX as usize - 4);
assert_eq!(buf.capacity(), 64); assert_eq!(buf.capacity(), 64);
for _ in 0..16 { for _ in 0..16 {
@@ -27,6 +27,14 @@ fn test_vec_as_mut_buf() {
assert_eq!(buf.len(), 68); assert_eq!(buf.len(), 68);
} }
#[test]
fn test_vec_put_bytes() {
let mut buf = Vec::new();
buf.push(17);
buf.put_bytes(19, 2);
assert_eq!([17, 19, 19], &buf[..]);
}
#[test] #[test]
fn test_put_u8() { fn test_put_u8() {
let mut buf = Vec::with_capacity(8); let mut buf = Vec::with_capacity(8);
@@ -45,6 +53,34 @@ fn test_put_u16() {
assert_eq!(b"\x54\x21", &buf[..]); assert_eq!(b"\x54\x21", &buf[..]);
} }
#[test]
fn test_put_int() {
let mut buf = Vec::with_capacity(8);
buf.put_int(0x1020304050607080, 3);
assert_eq!(b"\x60\x70\x80", &buf[..]);
}
#[test]
#[should_panic]
fn test_put_int_nbytes_overflow() {
let mut buf = Vec::with_capacity(8);
buf.put_int(0x1020304050607080, 9);
}
#[test]
fn test_put_int_le() {
let mut buf = Vec::with_capacity(8);
buf.put_int_le(0x1020304050607080, 3);
assert_eq!(b"\x80\x70\x60", &buf[..]);
}
#[test]
#[should_panic]
fn test_put_int_le_nbytes_overflow() {
let mut buf = Vec::with_capacity(8);
buf.put_int_le(0x1020304050607080, 9);
}
#[test] #[test]
#[should_panic(expected = "cannot advance")] #[should_panic(expected = "cannot advance")]
fn test_vec_advance_mut() { fn test_vec_advance_mut() {
@@ -70,6 +106,19 @@ fn test_mut_slice() {
let mut v = vec![0, 0, 0, 0]; let mut v = vec![0, 0, 0, 0];
let mut s = &mut v[..]; let mut s = &mut v[..];
s.put_u32(42); s.put_u32(42);
assert_eq!(s.len(), 0);
assert_eq!(&v, &[0, 0, 0, 42]);
}
#[test]
fn test_slice_put_bytes() {
let mut v = [0, 0, 0, 0];
let mut s = &mut v[..];
s.put_u8(17);
s.put_bytes(19, 2);
assert_eq!(1, s.remaining_mut());
assert_eq!(&[17, 19, 19, 0], &v[..]);
} }
#[test] #[test]
@@ -81,8 +130,8 @@ fn test_deref_bufmut_forwards() {
unreachable!("remaining_mut"); unreachable!("remaining_mut");
} }
fn bytes_mut(&mut self) -> &mut UninitSlice { fn chunk_mut(&mut self) -> &mut UninitSlice {
unreachable!("bytes_mut"); unreachable!("chunk_mut");
} }
unsafe fn advance_mut(&mut self, _: usize) { unsafe fn advance_mut(&mut self, _: usize) {
+190 -15
View File
@@ -4,8 +4,8 @@ use bytes::{Buf, BufMut, Bytes, BytesMut};
use std::usize; use std::usize;
const LONG: &'static [u8] = b"mary had a little lamb, little lamb, little lamb"; const LONG: &[u8] = b"mary had a little lamb, little lamb, little lamb";
const SHORT: &'static [u8] = b"hello world"; const SHORT: &[u8] = b"hello world";
fn is_sync<T: Sync>() {} fn is_sync<T: Sync>() {}
fn is_send<T: Send>() {} fn is_send<T: Send>() {}
@@ -411,8 +411,8 @@ fn freeze_after_split_off() {
fn fns_defined_for_bytes_mut() { fn fns_defined_for_bytes_mut() {
let mut bytes = BytesMut::from(&b"hello world"[..]); let mut bytes = BytesMut::from(&b"hello world"[..]);
bytes.as_ptr(); let _ = bytes.as_ptr();
bytes.as_mut_ptr(); let _ = bytes.as_mut_ptr();
// Iterator // Iterator
let v: Vec<u8> = bytes.as_ref().iter().cloned().collect(); let v: Vec<u8> = bytes.as_ref().iter().cloned().collect();
@@ -443,7 +443,7 @@ fn reserve_growth() {
let _ = bytes.split(); let _ = bytes.split();
bytes.reserve(65); bytes.reserve(65);
assert_eq!(bytes.capacity(), 128); assert_eq!(bytes.capacity(), 117);
} }
#[test] #[test]
@@ -461,6 +461,7 @@ fn reserve_allocates_at_least_original_capacity() {
} }
#[test] #[test]
#[cfg_attr(miri, ignore)] // Miri is too slow
fn reserve_max_original_capacity_value() { fn reserve_max_original_capacity_value() {
const SIZE: usize = 128 * 1024; const SIZE: usize = 128 * 1024;
@@ -526,6 +527,25 @@ fn reserve_in_arc_nonunique_does_not_overallocate() {
assert_eq!(2001, bytes.capacity()); assert_eq!(2001, bytes.capacity());
} }
/// This function tests `BytesMut::reserve_inner`, where `BytesMut` holds
/// a unique reference to the shared vector and decide to reuse it
/// by reallocating the `Vec`.
#[test]
fn reserve_shared_reuse() {
let mut bytes = BytesMut::with_capacity(1000);
bytes.put_slice(b"Hello, World!");
drop(bytes.split());
bytes.put_slice(b"!123ex123,sadchELLO,_wORLD!");
// Use split_off so that v.capacity() - self.cap != off
drop(bytes.split_off(9));
assert_eq!(&*bytes, b"!123ex123");
bytes.reserve(2000);
assert_eq!(&*bytes, b"!123ex123");
assert_eq!(bytes.capacity(), 2009);
}
#[test] #[test]
fn extend_mut() { fn extend_mut() {
let mut bytes = BytesMut::with_capacity(0); let mut bytes = BytesMut::with_capacity(0);
@@ -543,6 +563,13 @@ fn extend_from_slice_mut() {
} }
} }
#[test]
fn extend_mut_from_bytes() {
let mut bytes = BytesMut::with_capacity(0);
bytes.extend([Bytes::from(LONG)]);
assert_eq!(*bytes, LONG[..]);
}
#[test] #[test]
fn extend_mut_without_size_hint() { fn extend_mut_without_size_hint() {
let mut bytes = BytesMut::with_capacity(0); let mut bytes = BytesMut::with_capacity(0);
@@ -608,15 +635,15 @@ fn advance_past_len() {
#[test] #[test]
// Only run these tests on little endian systems. CI uses qemu for testing // Only run these tests on little endian systems. CI uses qemu for testing
// little endian... and qemu doesn't really support threading all that well. // big endian... and qemu doesn't really support threading all that well.
#[cfg(target_endian = "little")] #[cfg(any(miri, target_endian = "little"))]
fn stress() { fn stress() {
// Tests promoting a buffer from a vec -> shared in a concurrent situation // Tests promoting a buffer from a vec -> shared in a concurrent situation
use std::sync::{Arc, Barrier}; use std::sync::{Arc, Barrier};
use std::thread; use std::thread;
const THREADS: usize = 8; const THREADS: usize = 8;
const ITERS: usize = 1_000; const ITERS: usize = if cfg!(miri) { 100 } else { 1_000 };
for i in 0..ITERS { for i in 0..ITERS {
let data = [i as u8; 256]; let data = [i as u8; 256];
@@ -783,6 +810,31 @@ fn bytes_mut_unsplit_empty_self() {
assert_eq!(b"aaabbbcccddd", &buf[..]); assert_eq!(b"aaabbbcccddd", &buf[..]);
} }
#[test]
fn bytes_mut_unsplit_other_keeps_capacity() {
let mut buf = BytesMut::with_capacity(64);
buf.extend_from_slice(b"aabb");
// non empty other created "from" buf
let mut other = buf.split_off(buf.len());
other.extend_from_slice(b"ccddee");
buf.unsplit(other);
assert_eq!(buf.capacity(), 64);
}
#[test]
fn bytes_mut_unsplit_empty_other_keeps_capacity() {
let mut buf = BytesMut::with_capacity(64);
buf.extend_from_slice(b"aabbccddee");
// empty other created "from" buf
let other = buf.split_off(buf.len());
buf.unsplit(other);
assert_eq!(buf.capacity(), 64);
}
#[test] #[test]
fn bytes_mut_unsplit_arc_different() { fn bytes_mut_unsplit_arc_different() {
let mut buf = BytesMut::with_capacity(64); let mut buf = BytesMut::with_capacity(64);
@@ -848,7 +900,7 @@ fn from_iter_no_size_hint() {
fn test_slice_ref(bytes: &Bytes, start: usize, end: usize, expected: &[u8]) { fn test_slice_ref(bytes: &Bytes, start: usize, end: usize, expected: &[u8]) {
let slice = &(bytes.as_ref()[start..end]); let slice = &(bytes.as_ref()[start..end]);
let sub = bytes.slice_ref(&slice); let sub = bytes.slice_ref(slice);
assert_eq!(&sub[..], expected); assert_eq!(&sub[..], expected);
} }
@@ -868,7 +920,7 @@ fn slice_ref_empty() {
let bytes = Bytes::from(&b""[..]); let bytes = Bytes::from(&b""[..]);
let slice = &(bytes.as_ref()[0..0]); let slice = &(bytes.as_ref()[0..0]);
let sub = bytes.slice_ref(&slice); let sub = bytes.slice_ref(slice);
assert_eq!(&sub[..], b""); assert_eq!(&sub[..], b"");
} }
@@ -912,20 +964,20 @@ fn bytes_buf_mut_advance() {
let mut bytes = BytesMut::with_capacity(1024); let mut bytes = BytesMut::with_capacity(1024);
unsafe { unsafe {
let ptr = bytes.bytes_mut().as_mut_ptr(); let ptr = bytes.chunk_mut().as_mut_ptr();
assert_eq!(1024, bytes.bytes_mut().len()); assert_eq!(1024, bytes.chunk_mut().len());
bytes.advance_mut(10); bytes.advance_mut(10);
let next = bytes.bytes_mut().as_mut_ptr(); let next = bytes.chunk_mut().as_mut_ptr();
assert_eq!(1024 - 10, bytes.bytes_mut().len()); assert_eq!(1024 - 10, bytes.chunk_mut().len());
assert_eq!(ptr.offset(10), next); assert_eq!(ptr.offset(10), next);
// advance to the end // advance to the end
bytes.advance_mut(1024 - 10); bytes.advance_mut(1024 - 10);
// The buffer size is doubled // The buffer size is doubled
assert_eq!(1024, bytes.bytes_mut().len()); assert_eq!(1024, bytes.chunk_mut().len());
} }
} }
@@ -960,3 +1012,126 @@ fn bytes_with_capacity_but_empty() {
let vec = Vec::with_capacity(1); let vec = Vec::with_capacity(1);
let _ = Bytes::from(vec); let _ = Bytes::from(vec);
} }
#[test]
fn bytes_put_bytes() {
let mut bytes = BytesMut::new();
bytes.put_u8(17);
bytes.put_bytes(19, 2);
assert_eq!([17, 19, 19], bytes.as_ref());
}
#[test]
fn box_slice_empty() {
// See https://github.com/tokio-rs/bytes/issues/340
let empty: Box<[u8]> = Default::default();
let b = Bytes::from(empty);
assert!(b.is_empty());
}
#[test]
fn bytes_into_vec() {
// Test kind == KIND_VEC
let content = b"helloworld";
let mut bytes = BytesMut::new();
bytes.put_slice(content);
let vec: Vec<u8> = bytes.into();
assert_eq!(&vec, content);
// Test kind == KIND_ARC, shared.is_unique() == True
let mut bytes = BytesMut::new();
bytes.put_slice(b"abcdewe23");
bytes.put_slice(content);
// Overwrite the bytes to make sure only one reference to the underlying
// Vec exists.
bytes = bytes.split_off(9);
let vec: Vec<u8> = bytes.into();
assert_eq!(&vec, content);
// Test kind == KIND_ARC, shared.is_unique() == False
let prefix = b"abcdewe23";
let mut bytes = BytesMut::new();
bytes.put_slice(prefix);
bytes.put_slice(content);
let vec: Vec<u8> = bytes.split_off(prefix.len()).into();
assert_eq!(&vec, content);
let vec: Vec<u8> = bytes.into();
assert_eq!(&vec, prefix);
}
#[test]
fn test_bytes_into_vec() {
// Test STATIC_VTABLE.to_vec
let bs = b"1b23exfcz3r";
let vec: Vec<u8> = Bytes::from_static(bs).into();
assert_eq!(&*vec, bs);
// Test bytes_mut.SHARED_VTABLE.to_vec impl
eprintln!("1");
let mut bytes_mut: BytesMut = bs[..].into();
// Set kind to KIND_ARC so that after freeze, Bytes will use bytes_mut.SHARED_VTABLE
eprintln!("2");
drop(bytes_mut.split_off(bs.len()));
eprintln!("3");
let b1 = bytes_mut.freeze();
eprintln!("4");
let b2 = b1.clone();
eprintln!("{:#?}", (&*b1).as_ptr());
// shared.is_unique() = False
eprintln!("5");
assert_eq!(&*Vec::from(b2), bs);
// shared.is_unique() = True
eprintln!("6");
assert_eq!(&*Vec::from(b1), bs);
// Test bytes_mut.SHARED_VTABLE.to_vec impl where offset != 0
let mut bytes_mut1: BytesMut = bs[..].into();
let bytes_mut2 = bytes_mut1.split_off(9);
let b1 = bytes_mut1.freeze();
let b2 = bytes_mut2.freeze();
assert_eq!(Vec::from(b2), bs[9..]);
assert_eq!(Vec::from(b1), bs[..9]);
}
#[test]
fn test_bytes_into_vec_promotable_even() {
let vec = vec![33u8; 1024];
// Test cases where kind == KIND_VEC
let b1 = Bytes::from(vec.clone());
assert_eq!(Vec::from(b1), vec);
// Test cases where kind == KIND_ARC, ref_cnt == 1
let b1 = Bytes::from(vec.clone());
drop(b1.clone());
assert_eq!(Vec::from(b1), vec);
// Test cases where kind == KIND_ARC, ref_cnt == 2
let b1 = Bytes::from(vec.clone());
let b2 = b1.clone();
assert_eq!(Vec::from(b1), vec);
// Test cases where vtable = SHARED_VTABLE, kind == KIND_ARC, ref_cnt == 1
assert_eq!(Vec::from(b2), vec);
// Test cases where offset != 0
let mut b1 = Bytes::from(vec.clone());
let b2 = b1.split_off(20);
assert_eq!(Vec::from(b2), vec[20..]);
assert_eq!(Vec::from(b1), vec[..20]);
}
+32 -2
View File
@@ -1,6 +1,8 @@
//! Test using `Bytes` with an allocator that hands out "odd" pointers for //! Test using `Bytes` with an allocator that hands out "odd" pointers for
//! vectors (pointers where the LSB is set). //! vectors (pointers where the LSB is set).
#![cfg(not(miri))] // Miri does not support custom allocators (also, Miri is "odd" by default with 50% chance)
use std::alloc::{GlobalAlloc, Layout, System}; use std::alloc::{GlobalAlloc, Layout, System};
use std::ptr; use std::ptr;
@@ -22,8 +24,7 @@ unsafe impl GlobalAlloc for Odd {
}; };
let ptr = System.alloc(new_layout); let ptr = System.alloc(new_layout);
if !ptr.is_null() { if !ptr.is_null() {
let ptr = ptr.offset(1); ptr.offset(1)
ptr
} else { } else {
ptr ptr
} }
@@ -65,3 +66,32 @@ fn test_bytes_clone_drop() {
let b1 = Bytes::from(vec); let b1 = Bytes::from(vec);
let _b2 = b1.clone(); let _b2 = b1.clone();
} }
#[test]
fn test_bytes_into_vec() {
let vec = vec![33u8; 1024];
// Test cases where kind == KIND_VEC
let b1 = Bytes::from(vec.clone());
assert_eq!(Vec::from(b1), vec);
// Test cases where kind == KIND_ARC, ref_cnt == 1
let b1 = Bytes::from(vec.clone());
drop(b1.clone());
assert_eq!(Vec::from(b1), vec);
// Test cases where kind == KIND_ARC, ref_cnt == 2
let b1 = Bytes::from(vec.clone());
let b2 = b1.clone();
assert_eq!(Vec::from(b1), vec);
// Test cases where vtable = SHARED_VTABLE, kind == KIND_ARC, ref_cnt == 1
assert_eq!(Vec::from(b2), vec);
// Test cases where offset != 0
let mut b1 = Bytes::from(vec.clone());
let b2 = b1.split_off(20);
assert_eq!(Vec::from(b2), vec[20..]);
assert_eq!(Vec::from(b1), vec[..20]);
}
+103 -39
View File
@@ -1,61 +1,87 @@
use std::alloc::{GlobalAlloc, Layout, System}; use std::alloc::{GlobalAlloc, Layout, System};
use std::{mem, ptr}; use std::ptr::null_mut;
use std::sync::atomic::{AtomicPtr, AtomicUsize, Ordering};
use bytes::{Buf, Bytes}; use bytes::{Buf, Bytes};
#[global_allocator] #[global_allocator]
static LEDGER: Ledger = Ledger; static LEDGER: Ledger = Ledger::new();
struct Ledger; const LEDGER_LENGTH: usize = 2048;
const USIZE_SIZE: usize = mem::size_of::<usize>(); struct Ledger {
alloc_table: [(AtomicPtr<u8>, AtomicUsize); LEDGER_LENGTH],
}
unsafe impl GlobalAlloc for Ledger { impl Ledger {
unsafe fn alloc(&self, layout: Layout) -> *mut u8 { const fn new() -> Self {
if layout.align() == 1 && layout.size() > 0 { const ELEM: (AtomicPtr<u8>, AtomicUsize) =
// Allocate extra space to stash a record of (AtomicPtr::new(null_mut()), AtomicUsize::new(0));
// how much space there was. let alloc_table = [ELEM; LEDGER_LENGTH];
let orig_size = layout.size();
let size = orig_size + USIZE_SIZE; Self { alloc_table }
let new_layout = match Layout::from_size_align(size, 1) { }
Ok(layout) => layout,
Err(_err) => return ptr::null_mut(), /// Iterate over our table until we find an open entry, then insert into said entry
}; fn insert(&self, ptr: *mut u8, size: usize) {
let ptr = System.alloc(new_layout); for (entry_ptr, entry_size) in self.alloc_table.iter() {
if !ptr.is_null() { // SeqCst is good enough here, we don't care about perf, i just want to be correct!
(ptr as *mut usize).write(orig_size); if entry_ptr
let ptr = ptr.offset(USIZE_SIZE as isize); .compare_exchange(null_mut(), ptr, Ordering::SeqCst, Ordering::SeqCst)
ptr .is_ok()
} else { {
ptr entry_size.store(size, Ordering::SeqCst);
break;
} }
} else {
System.alloc(layout)
} }
} }
unsafe fn dealloc(&self, ptr: *mut u8, layout: Layout) { fn remove(&self, ptr: *mut u8) -> usize {
if layout.align() == 1 && layout.size() > 0 { for (entry_ptr, entry_size) in self.alloc_table.iter() {
let off_ptr = (ptr as *mut usize).offset(-1); // set the value to be something that will never try and be deallocated, so that we
let orig_size = off_ptr.read(); // don't have any chance of a race condition
if orig_size != layout.size() { //
panic!( // dont worry, LEDGER_LENGTH is really long to compensate for us not reclaiming space
"bad dealloc: alloc size was {}, dealloc size is {}", if entry_ptr
orig_size, .compare_exchange(
layout.size() ptr,
); invalid_ptr(usize::MAX),
Ordering::SeqCst,
Ordering::SeqCst,
)
.is_ok()
{
return entry_size.load(Ordering::SeqCst);
} }
}
let new_layout = match Layout::from_size_align(layout.size() + USIZE_SIZE, 1) { panic!("Couldn't find a matching entry for {:x?}", ptr);
Ok(layout) => layout, }
Err(_err) => std::process::abort(), }
};
System.dealloc(off_ptr as *mut u8, new_layout); unsafe impl GlobalAlloc for Ledger {
unsafe fn alloc(&self, layout: Layout) -> *mut u8 {
let size = layout.size();
let ptr = System.alloc(layout);
self.insert(ptr, size);
ptr
}
unsafe fn dealloc(&self, ptr: *mut u8, layout: Layout) {
let orig_size = self.remove(ptr);
if orig_size != layout.size() {
panic!(
"bad dealloc: alloc size was {}, dealloc size is {}",
orig_size,
layout.size()
);
} else { } else {
System.dealloc(ptr, layout); System.dealloc(ptr, layout);
} }
} }
} }
#[test] #[test]
fn test_bytes_advance() { fn test_bytes_advance() {
let mut bytes = Bytes::from(vec![10, 20, 30]); let mut bytes = Bytes::from(vec![10, 20, 30]);
@@ -77,3 +103,41 @@ fn test_bytes_truncate_and_advance() {
bytes.advance(1); bytes.advance(1);
drop(bytes); drop(bytes);
} }
/// Returns a dangling pointer with the given address. This is used to store
/// integer data in pointer fields.
#[inline]
fn invalid_ptr<T>(addr: usize) -> *mut T {
let ptr = std::ptr::null_mut::<u8>().wrapping_add(addr);
debug_assert_eq!(ptr as usize, addr);
ptr.cast::<T>()
}
#[test]
fn test_bytes_into_vec() {
let vec = vec![33u8; 1024];
// Test cases where kind == KIND_VEC
let b1 = Bytes::from(vec.clone());
assert_eq!(Vec::from(b1), vec);
// Test cases where kind == KIND_ARC, ref_cnt == 1
let b1 = Bytes::from(vec.clone());
drop(b1.clone());
assert_eq!(Vec::from(b1), vec);
// Test cases where kind == KIND_ARC, ref_cnt == 2
let b1 = Bytes::from(vec.clone());
let b2 = b1.clone();
assert_eq!(Vec::from(b1), vec);
// Test cases where vtable = SHARED_VTABLE, kind == KIND_ARC, ref_cnt == 1
assert_eq!(Vec::from(b2), vec);
// Test cases where offset != 0
let mut b1 = Bytes::from(vec.clone());
let b2 = b1.split_off(20);
assert_eq!(Vec::from(b2), vec[20..]);
assert_eq!(Vec::from(b1), vec[..20]);
}
+47 -4
View File
@@ -62,7 +62,7 @@ fn vectored_read() {
IoSlice::new(b4), IoSlice::new(b4),
]; ];
assert_eq!(2, buf.bytes_vectored(&mut iovecs)); assert_eq!(2, buf.chunks_vectored(&mut iovecs));
assert_eq!(iovecs[0][..], b"hello"[..]); assert_eq!(iovecs[0][..], b"hello"[..]);
assert_eq!(iovecs[1][..], b"world"[..]); assert_eq!(iovecs[1][..], b"world"[..]);
assert_eq!(iovecs[2][..], b""[..]); assert_eq!(iovecs[2][..], b""[..]);
@@ -83,7 +83,7 @@ fn vectored_read() {
IoSlice::new(b4), IoSlice::new(b4),
]; ];
assert_eq!(2, buf.bytes_vectored(&mut iovecs)); assert_eq!(2, buf.chunks_vectored(&mut iovecs));
assert_eq!(iovecs[0][..], b"llo"[..]); assert_eq!(iovecs[0][..], b"llo"[..]);
assert_eq!(iovecs[1][..], b"world"[..]); assert_eq!(iovecs[1][..], b"world"[..]);
assert_eq!(iovecs[2][..], b""[..]); assert_eq!(iovecs[2][..], b""[..]);
@@ -104,7 +104,7 @@ fn vectored_read() {
IoSlice::new(b4), IoSlice::new(b4),
]; ];
assert_eq!(1, buf.bytes_vectored(&mut iovecs)); assert_eq!(1, buf.chunks_vectored(&mut iovecs));
assert_eq!(iovecs[0][..], b"world"[..]); assert_eq!(iovecs[0][..], b"world"[..]);
assert_eq!(iovecs[1][..], b""[..]); assert_eq!(iovecs[1][..], b""[..]);
assert_eq!(iovecs[2][..], b""[..]); assert_eq!(iovecs[2][..], b""[..]);
@@ -125,10 +125,53 @@ fn vectored_read() {
IoSlice::new(b4), IoSlice::new(b4),
]; ];
assert_eq!(1, buf.bytes_vectored(&mut iovecs)); assert_eq!(1, buf.chunks_vectored(&mut iovecs));
assert_eq!(iovecs[0][..], b"ld"[..]); assert_eq!(iovecs[0][..], b"ld"[..]);
assert_eq!(iovecs[1][..], b""[..]); assert_eq!(iovecs[1][..], b""[..]);
assert_eq!(iovecs[2][..], b""[..]); assert_eq!(iovecs[2][..], b""[..]);
assert_eq!(iovecs[3][..], b""[..]); assert_eq!(iovecs[3][..], b""[..]);
} }
} }
#[test]
fn chain_growing_buffer() {
let mut buff = [' ' as u8; 10];
let mut vec = b"wassup".to_vec();
let mut chained = (&mut buff[..]).chain_mut(&mut vec).chain_mut(Vec::new()); // Required for potential overflow because remaining_mut for Vec is isize::MAX - vec.len(), but for chain_mut is usize::MAX
chained.put_slice(b"hey there123123");
assert_eq!(&buff, b"hey there1");
assert_eq!(&vec, b"wassup23123");
}
#[test]
fn chain_overflow_remaining_mut() {
let mut chained = Vec::<u8>::new().chain_mut(Vec::new()).chain_mut(Vec::new());
assert_eq!(chained.remaining_mut(), usize::MAX);
chained.put_slice(&[0; 256]);
assert_eq!(chained.remaining_mut(), usize::MAX);
}
#[test]
fn chain_get_bytes() {
let mut ab = Bytes::copy_from_slice(b"ab");
let mut cd = Bytes::copy_from_slice(b"cd");
let ab_ptr = ab.as_ptr();
let cd_ptr = cd.as_ptr();
let mut chain = (&mut ab).chain(&mut cd);
let a = chain.copy_to_bytes(1);
let bc = chain.copy_to_bytes(2);
let d = chain.copy_to_bytes(1);
assert_eq!(Bytes::copy_from_slice(b"a"), a);
assert_eq!(Bytes::copy_from_slice(b"bc"), bc);
assert_eq!(Bytes::copy_from_slice(b"d"), d);
// assert `get_bytes` did not allocate
assert_eq!(ab_ptr, a.as_ptr());
// assert `get_bytes` did not allocate
assert_eq!(cd_ptr.wrapping_offset(1), d.as_ptr());
}
+21 -1
View File
@@ -1,6 +1,7 @@
#![warn(rust_2018_idioms)] #![warn(rust_2018_idioms)]
use bytes::buf::Buf; use bytes::buf::Buf;
use bytes::Bytes;
#[test] #[test]
fn long_take() { fn long_take() {
@@ -8,5 +9,24 @@ fn long_take() {
// overrun the buffer. Regression test for #138. // overrun the buffer. Regression test for #138.
let buf = b"hello world".take(100); let buf = b"hello world".take(100);
assert_eq!(11, buf.remaining()); assert_eq!(11, buf.remaining());
assert_eq!(b"hello world", buf.bytes()); assert_eq!(b"hello world", buf.chunk());
}
#[test]
fn take_copy_to_bytes() {
let mut abcd = Bytes::copy_from_slice(b"abcd");
let abcd_ptr = abcd.as_ptr();
let mut take = (&mut abcd).take(2);
let a = take.copy_to_bytes(1);
assert_eq!(Bytes::copy_from_slice(b"a"), a);
// assert `to_bytes` did not allocate
assert_eq!(abcd_ptr, a.as_ptr());
assert_eq!(Bytes::copy_from_slice(b"bcd"), abcd);
}
#[test]
#[should_panic]
fn take_copy_to_bytes_panics() {
let abcd = Bytes::copy_from_slice(b"abcd");
abcd.take(2).copy_to_bytes(3);
} }