mirror of
https://github.com/tokio-rs/tokio.git
synced 2026-08-21 00:00:10 +02:00
chore: fix formatting, remove old rustfmt.toml (#2007)
`cargo fmt` has a bug where it does not format modules scoped with feature flags.
This commit is contained in:
committed by
Carl Lerche
parent
f309b295bb
commit
8656b7b8eb
@@ -13,5 +13,6 @@ jobs:
|
||||
cargo fmt --version
|
||||
displayName: Install rustfmt
|
||||
- script: |
|
||||
cargo fmt --all -- --check
|
||||
# Workaround for rust-lang/cargo#7732
|
||||
rustfmt --check --edition 2018 $(find . -name '*.rs' -print)
|
||||
displayName: Check formatting
|
||||
|
||||
@@ -1 +0,0 @@
|
||||
edition = "2018"
|
||||
@@ -3,7 +3,10 @@ use crate::codec::encoder::Encoder;
|
||||
use crate::codec::framed_read::{framed_read2, framed_read2_with_buffer, FramedRead2};
|
||||
use crate::codec::framed_write::{framed_write2, framed_write2_with_buffer, FramedWrite2};
|
||||
|
||||
use tokio::{io::{AsyncBufRead, AsyncRead, AsyncWrite}, stream::Stream};
|
||||
use tokio::{
|
||||
io::{AsyncBufRead, AsyncRead, AsyncWrite},
|
||||
stream::Stream,
|
||||
};
|
||||
|
||||
use bytes::BytesMut;
|
||||
use futures_sink::Sink;
|
||||
|
||||
@@ -2,7 +2,10 @@ use crate::codec::decoder::Decoder;
|
||||
use crate::codec::encoder::Encoder;
|
||||
use crate::codec::framed::{Fuse, ProjectFuse};
|
||||
|
||||
use tokio::{io::{AsyncBufRead, AsyncRead, AsyncWrite}, stream::Stream};
|
||||
use tokio::{
|
||||
io::{AsyncBufRead, AsyncRead, AsyncWrite},
|
||||
stream::Stream,
|
||||
};
|
||||
|
||||
use bytes::BytesMut;
|
||||
use futures_core::ready;
|
||||
|
||||
@@ -1,6 +1,6 @@
|
||||
use sdt::pin::Pin;
|
||||
use std::future::Future;
|
||||
use std::marker;
|
||||
use sdt::pin::Pin;
|
||||
use std::task::{Context, Poll};
|
||||
|
||||
/// Future for the [`pending()`] function.
|
||||
@@ -29,7 +29,8 @@ struct Pending<T> {
|
||||
pub async fn pending() -> ! {
|
||||
Pending {
|
||||
_data: marker::PhantomData,
|
||||
}.await
|
||||
}
|
||||
.await
|
||||
}
|
||||
|
||||
impl<T> Future for Pending<T> {
|
||||
@@ -40,5 +41,4 @@ impl<T> Future for Pending<T> {
|
||||
}
|
||||
}
|
||||
|
||||
impl<T> Unpin for Pending<T> {
|
||||
}
|
||||
impl<T> Unpin for Pending<T> {}
|
||||
|
||||
@@ -12,8 +12,8 @@ use std::cell::RefCell;
|
||||
use std::fmt;
|
||||
use std::io;
|
||||
use std::marker::PhantomData;
|
||||
use std::sync::{Arc, Weak};
|
||||
use std::sync::atomic::Ordering::SeqCst;
|
||||
use std::sync::{Arc, Weak};
|
||||
use std::task::Waker;
|
||||
use std::time::Duration;
|
||||
|
||||
|
||||
@@ -3,7 +3,7 @@ use crate::loom::sync::atomic::AtomicUsize;
|
||||
use crate::util::bit;
|
||||
use crate::util::slab::{Address, Entry, Generation};
|
||||
|
||||
use std::sync::atomic::Ordering::{Acquire, AcqRel, SeqCst};
|
||||
use std::sync::atomic::Ordering::{AcqRel, Acquire, SeqCst};
|
||||
|
||||
#[derive(Debug)]
|
||||
pub(crate) struct ScheduledIo {
|
||||
@@ -29,12 +29,10 @@ impl Entry for ScheduledIo {
|
||||
|
||||
let next = PACK.pack(generation.next().to_usize(), 0);
|
||||
|
||||
match self.readiness.compare_exchange(
|
||||
current,
|
||||
next,
|
||||
AcqRel,
|
||||
Acquire,
|
||||
) {
|
||||
match self
|
||||
.readiness
|
||||
.compare_exchange(current, next, AcqRel, Acquire)
|
||||
{
|
||||
Ok(_) => break,
|
||||
Err(actual) => current = actual,
|
||||
}
|
||||
|
||||
@@ -1,5 +1,5 @@
|
||||
use crate::io::driver::platform;
|
||||
use crate::io::{AsyncRead, AsyncWrite, Registration};
|
||||
use crate::io::driver::{platform};
|
||||
|
||||
use mio::event::Evented;
|
||||
use std::fmt;
|
||||
|
||||
@@ -1,9 +1,9 @@
|
||||
use crate::io::driver::{Direction, Handle, platform};
|
||||
use crate::io::driver::{platform, Direction, Handle};
|
||||
use crate::util::slab::Address;
|
||||
|
||||
use mio::{self, Evented};
|
||||
use std::task::{Context, Poll};
|
||||
use std::io;
|
||||
use std::task::{Context, Poll};
|
||||
|
||||
cfg_io_driver! {
|
||||
/// Associates an I/O resource with the reactor instance that drives it.
|
||||
|
||||
@@ -2,8 +2,8 @@ use crate::io::util::chain::{chain, Chain};
|
||||
use crate::io::util::read::{read, Read};
|
||||
use crate::io::util::read_buf::{read_buf, ReadBuf};
|
||||
use crate::io::util::read_exact::{read_exact, ReadExact};
|
||||
use crate::io::util::read_int::{ReadU8, ReadU16, ReadU32, ReadU64, ReadU128};
|
||||
use crate::io::util::read_int::{ReadI8, ReadI16, ReadI32, ReadI64, ReadI128};
|
||||
use crate::io::util::read_int::{ReadI128, ReadI16, ReadI32, ReadI64, ReadI8};
|
||||
use crate::io::util::read_int::{ReadU128, ReadU16, ReadU32, ReadU64, ReadU8};
|
||||
use crate::io::util::read_to_end::{read_to_end, ReadToEnd};
|
||||
use crate::io::util::read_to_string::{read_to_string, ReadToString};
|
||||
use crate::io::util::take::{take, Take};
|
||||
|
||||
@@ -35,4 +35,4 @@ pub trait AsyncSeekExt: AsyncSeek {
|
||||
}
|
||||
}
|
||||
|
||||
impl<S: AsyncSeek + ?Sized> AsyncSeekExt for S {}
|
||||
impl<S: AsyncSeek + ?Sized> AsyncSeekExt for S {}
|
||||
|
||||
@@ -3,8 +3,8 @@ use crate::io::util::shutdown::{shutdown, Shutdown};
|
||||
use crate::io::util::write::{write, Write};
|
||||
use crate::io::util::write_all::{write_all, WriteAll};
|
||||
use crate::io::util::write_buf::{write_buf, WriteBuf};
|
||||
use crate::io::util::write_int::{WriteU8, WriteU16, WriteU32, WriteU64, WriteU128};
|
||||
use crate::io::util::write_int::{WriteI8, WriteI16, WriteI32, WriteI64, WriteI128};
|
||||
use crate::io::util::write_int::{WriteI128, WriteI16, WriteI32, WriteI64, WriteI8};
|
||||
use crate::io::util::write_int::{WriteU128, WriteU16, WriteU32, WriteU64, WriteU8};
|
||||
use crate::io::AsyncWrite;
|
||||
|
||||
use bytes::Buf;
|
||||
|
||||
@@ -48,7 +48,8 @@ macro_rules! reader {
|
||||
}
|
||||
|
||||
while *me.read < $bytes as u8 {
|
||||
*me.read += match me.src
|
||||
*me.read += match me
|
||||
.src
|
||||
.as_mut()
|
||||
.poll_read(cx, &mut me.buf[*me.read as usize..])
|
||||
{
|
||||
|
||||
@@ -49,7 +49,8 @@ macro_rules! writer {
|
||||
}
|
||||
|
||||
while *me.written < $bytes as u8 {
|
||||
*me.written += match me.dst
|
||||
*me.written += match me
|
||||
.dst
|
||||
.as_mut()
|
||||
.poll_write(cx, &me.buf[*me.written as usize..])
|
||||
{
|
||||
@@ -77,10 +78,7 @@ macro_rules! writer8 {
|
||||
|
||||
impl<W> $name<W> {
|
||||
pub(crate) fn new(dst: W, byte: $ty) -> Self {
|
||||
Self {
|
||||
dst,
|
||||
byte,
|
||||
}
|
||||
Self { dst, byte }
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
+66
-18
@@ -674,22 +674,70 @@ impl TcpStream {
|
||||
|
||||
// IoSlice isn't Copy, so we must expand this manually ;_;
|
||||
let mut slices: [IoSlice<'_>; MAX_BUFS] = [
|
||||
IoSlice::new(S), IoSlice::new(S), IoSlice::new(S), IoSlice::new(S),
|
||||
IoSlice::new(S), IoSlice::new(S), IoSlice::new(S), IoSlice::new(S),
|
||||
IoSlice::new(S), IoSlice::new(S), IoSlice::new(S), IoSlice::new(S),
|
||||
IoSlice::new(S), IoSlice::new(S), IoSlice::new(S), IoSlice::new(S),
|
||||
IoSlice::new(S), IoSlice::new(S), IoSlice::new(S), IoSlice::new(S),
|
||||
IoSlice::new(S), IoSlice::new(S), IoSlice::new(S), IoSlice::new(S),
|
||||
IoSlice::new(S), IoSlice::new(S), IoSlice::new(S), IoSlice::new(S),
|
||||
IoSlice::new(S), IoSlice::new(S), IoSlice::new(S), IoSlice::new(S),
|
||||
IoSlice::new(S), IoSlice::new(S), IoSlice::new(S), IoSlice::new(S),
|
||||
IoSlice::new(S), IoSlice::new(S), IoSlice::new(S), IoSlice::new(S),
|
||||
IoSlice::new(S), IoSlice::new(S), IoSlice::new(S), IoSlice::new(S),
|
||||
IoSlice::new(S), IoSlice::new(S), IoSlice::new(S), IoSlice::new(S),
|
||||
IoSlice::new(S), IoSlice::new(S), IoSlice::new(S), IoSlice::new(S),
|
||||
IoSlice::new(S), IoSlice::new(S), IoSlice::new(S), IoSlice::new(S),
|
||||
IoSlice::new(S), IoSlice::new(S), IoSlice::new(S), IoSlice::new(S),
|
||||
IoSlice::new(S), IoSlice::new(S), IoSlice::new(S), IoSlice::new(S),
|
||||
IoSlice::new(S),
|
||||
IoSlice::new(S),
|
||||
IoSlice::new(S),
|
||||
IoSlice::new(S),
|
||||
IoSlice::new(S),
|
||||
IoSlice::new(S),
|
||||
IoSlice::new(S),
|
||||
IoSlice::new(S),
|
||||
IoSlice::new(S),
|
||||
IoSlice::new(S),
|
||||
IoSlice::new(S),
|
||||
IoSlice::new(S),
|
||||
IoSlice::new(S),
|
||||
IoSlice::new(S),
|
||||
IoSlice::new(S),
|
||||
IoSlice::new(S),
|
||||
IoSlice::new(S),
|
||||
IoSlice::new(S),
|
||||
IoSlice::new(S),
|
||||
IoSlice::new(S),
|
||||
IoSlice::new(S),
|
||||
IoSlice::new(S),
|
||||
IoSlice::new(S),
|
||||
IoSlice::new(S),
|
||||
IoSlice::new(S),
|
||||
IoSlice::new(S),
|
||||
IoSlice::new(S),
|
||||
IoSlice::new(S),
|
||||
IoSlice::new(S),
|
||||
IoSlice::new(S),
|
||||
IoSlice::new(S),
|
||||
IoSlice::new(S),
|
||||
IoSlice::new(S),
|
||||
IoSlice::new(S),
|
||||
IoSlice::new(S),
|
||||
IoSlice::new(S),
|
||||
IoSlice::new(S),
|
||||
IoSlice::new(S),
|
||||
IoSlice::new(S),
|
||||
IoSlice::new(S),
|
||||
IoSlice::new(S),
|
||||
IoSlice::new(S),
|
||||
IoSlice::new(S),
|
||||
IoSlice::new(S),
|
||||
IoSlice::new(S),
|
||||
IoSlice::new(S),
|
||||
IoSlice::new(S),
|
||||
IoSlice::new(S),
|
||||
IoSlice::new(S),
|
||||
IoSlice::new(S),
|
||||
IoSlice::new(S),
|
||||
IoSlice::new(S),
|
||||
IoSlice::new(S),
|
||||
IoSlice::new(S),
|
||||
IoSlice::new(S),
|
||||
IoSlice::new(S),
|
||||
IoSlice::new(S),
|
||||
IoSlice::new(S),
|
||||
IoSlice::new(S),
|
||||
IoSlice::new(S),
|
||||
IoSlice::new(S),
|
||||
IoSlice::new(S),
|
||||
IoSlice::new(S),
|
||||
IoSlice::new(S),
|
||||
];
|
||||
let cnt = buf.bytes_vectored(&mut slices);
|
||||
|
||||
@@ -703,11 +751,11 @@ impl TcpStream {
|
||||
Ok(n) => {
|
||||
buf.advance(n);
|
||||
Poll::Ready(Ok(n))
|
||||
},
|
||||
}
|
||||
Err(ref e) if e.kind() == io::ErrorKind::WouldBlock => {
|
||||
self.io.clear_write_ready(cx)?;
|
||||
Poll::Pending
|
||||
},
|
||||
}
|
||||
Err(e) => Poll::Ready(Err(e)),
|
||||
}
|
||||
}
|
||||
|
||||
@@ -4,4 +4,4 @@ pub(crate) mod socket;
|
||||
pub(crate) use socket::UdpSocket;
|
||||
|
||||
mod split;
|
||||
pub use split::{RecvHalf, SendHalf, ReuniteError};
|
||||
pub use split::{RecvHalf, ReuniteError, SendHalf};
|
||||
|
||||
@@ -556,7 +556,7 @@ impl Command {
|
||||
imp::spawn_child(&mut self.std).map(|spawned_child| Child {
|
||||
child: ChildDropGuard {
|
||||
inner: spawned_child.child,
|
||||
kill_on_drop: self.kill_on_drop
|
||||
kill_on_drop: self.kill_on_drop,
|
||||
},
|
||||
stdin: spawned_child.stdin.map(|inner| ChildStdin { inner }),
|
||||
stdout: spawned_child.stdout.map(|inner| ChildStdout { inner }),
|
||||
|
||||
@@ -2,10 +2,10 @@
|
||||
|
||||
use crate::loom::sync::{Arc, Condvar, Mutex};
|
||||
use crate::loom::thread;
|
||||
use crate::runtime::{self, io, time, Builder, Callback};
|
||||
use crate::runtime::blocking::shutdown;
|
||||
use crate::runtime::blocking::schedule::NoopSchedule;
|
||||
use crate::runtime::blocking::shutdown;
|
||||
use crate::runtime::blocking::task::BlockingTask;
|
||||
use crate::runtime::{self, io, time, Builder, Callback};
|
||||
use crate::task::{self, JoinHandle};
|
||||
|
||||
use std::cell::Cell;
|
||||
@@ -55,7 +55,6 @@ struct Inner {
|
||||
clock: time::Clock,
|
||||
|
||||
thread_cap: usize,
|
||||
|
||||
}
|
||||
|
||||
struct Shared {
|
||||
|
||||
+26
-11
@@ -2,8 +2,8 @@
|
||||
//!
|
||||
//! A combination of the various resource driver park handles.
|
||||
|
||||
use crate::loom::sync::{Arc, Mutex, Condvar};
|
||||
use crate::loom::sync::atomic::AtomicUsize;
|
||||
use crate::loom::sync::{Arc, Condvar, Mutex};
|
||||
use crate::loom::thread;
|
||||
use crate::park::{Park, Unpark};
|
||||
use crate::runtime::time;
|
||||
@@ -84,7 +84,9 @@ impl Park for Parker {
|
||||
type Error = ();
|
||||
|
||||
fn unpark(&self) -> Unparker {
|
||||
Unparker { inner: self.inner.clone() }
|
||||
Unparker {
|
||||
inner: self.inner.clone(),
|
||||
}
|
||||
}
|
||||
|
||||
fn park(&mut self) -> Result<(), Self::Error> {
|
||||
@@ -97,8 +99,7 @@ impl Park for Parker {
|
||||
assert_eq!(duration, Duration::from_millis(0));
|
||||
|
||||
if let Some(mut driver) = self.inner.shared.driver.try_lock() {
|
||||
driver.park_timeout(duration)
|
||||
.map_err(|_| ())
|
||||
driver.park_timeout(duration).map_err(|_| ())
|
||||
} else {
|
||||
Ok(())
|
||||
}
|
||||
@@ -117,7 +118,11 @@ impl Inner {
|
||||
for _ in 0..3 {
|
||||
// If we were previously notified then we consume this notification and
|
||||
// return quickly.
|
||||
if self.state.compare_exchange(NOTIFIED, EMPTY, SeqCst, SeqCst).is_ok() {
|
||||
if self
|
||||
.state
|
||||
.compare_exchange(NOTIFIED, EMPTY, SeqCst, SeqCst)
|
||||
.is_ok()
|
||||
{
|
||||
return;
|
||||
}
|
||||
|
||||
@@ -135,7 +140,10 @@ impl Inner {
|
||||
// Otherwise we need to coordinate going to sleep
|
||||
let mut m = self.mutex.lock().unwrap();
|
||||
|
||||
match self.state.compare_exchange(EMPTY, PARKED_CONDVAR, SeqCst, SeqCst) {
|
||||
match self
|
||||
.state
|
||||
.compare_exchange(EMPTY, PARKED_CONDVAR, SeqCst, SeqCst)
|
||||
{
|
||||
Ok(_) => {}
|
||||
Err(NOTIFIED) => {
|
||||
// We must read here, even though we know it will be `NOTIFIED`.
|
||||
@@ -155,7 +163,11 @@ impl Inner {
|
||||
loop {
|
||||
m = self.condvar.wait(m).unwrap();
|
||||
|
||||
if self.state.compare_exchange(NOTIFIED, EMPTY, SeqCst, SeqCst).is_ok() {
|
||||
if self
|
||||
.state
|
||||
.compare_exchange(NOTIFIED, EMPTY, SeqCst, SeqCst)
|
||||
.is_ok()
|
||||
{
|
||||
// got a notification
|
||||
return;
|
||||
}
|
||||
@@ -165,7 +177,10 @@ impl Inner {
|
||||
}
|
||||
|
||||
fn park_driver(&self, driver: &mut time::Driver) {
|
||||
match self.state.compare_exchange(EMPTY, PARKED_DRIVER, SeqCst, SeqCst) {
|
||||
match self
|
||||
.state
|
||||
.compare_exchange(EMPTY, PARKED_DRIVER, SeqCst, SeqCst)
|
||||
{
|
||||
Ok(_) => {}
|
||||
Err(NOTIFIED) => {
|
||||
// We must read here, even though we know it will be `NOTIFIED`.
|
||||
@@ -186,7 +201,7 @@ impl Inner {
|
||||
driver.park().unwrap();
|
||||
|
||||
match self.state.swap(EMPTY, SeqCst) {
|
||||
NOTIFIED => {} // got a notification, hurray!
|
||||
NOTIFIED => {} // got a notification, hurray!
|
||||
PARKED_DRIVER => {} // no notification, alas
|
||||
n => panic!("inconsistent park_timeout state: {}", n),
|
||||
}
|
||||
@@ -199,8 +214,8 @@ impl Inner {
|
||||
// is already `NOTIFIED`. That is why this must be a swap rather than a
|
||||
// compare-and-swap that returns if it reads `NOTIFIED` on failure.
|
||||
match self.state.swap(NOTIFIED, SeqCst) {
|
||||
EMPTY => {}, // no one was waiting
|
||||
NOTIFIED => {}, // already unparked
|
||||
EMPTY => {} // no one was waiting
|
||||
NOTIFIED => {} // already unparked
|
||||
PARKED_CONDVAR => self.unpark_condvar(),
|
||||
PARKED_DRIVER => self.unpark_driver(),
|
||||
actual => panic!("inconsistent state in unpark; actual = {}", actual),
|
||||
|
||||
@@ -54,20 +54,12 @@ pub(crate) struct Workers {
|
||||
}
|
||||
|
||||
impl ThreadPool {
|
||||
pub(crate) fn new(
|
||||
pool_size: usize,
|
||||
parker: Parker,
|
||||
) -> (ThreadPool, Workers) {
|
||||
let (pool, workers) = worker::create_set(
|
||||
pool_size,
|
||||
parker,
|
||||
);
|
||||
pub(crate) fn new(pool_size: usize, parker: Parker) -> (ThreadPool, Workers) {
|
||||
let (pool, workers) = worker::create_set(pool_size, parker);
|
||||
|
||||
let spawner = Spawner::new(pool);
|
||||
|
||||
let pool = ThreadPool {
|
||||
spawner,
|
||||
};
|
||||
let pool = ThreadPool { spawner };
|
||||
|
||||
(pool, Workers { workers })
|
||||
}
|
||||
|
||||
@@ -4,8 +4,8 @@
|
||||
|
||||
use crate::loom::rand::seed;
|
||||
use crate::park::Park;
|
||||
use crate::runtime::Parker;
|
||||
use crate::runtime::thread_pool::{current, queue, Idle, Owned, Shared};
|
||||
use crate::runtime::Parker;
|
||||
use crate::task::{self, JoinHandle, Task};
|
||||
use crate::util::{CachePadded, FastRand};
|
||||
|
||||
|
||||
@@ -1,5 +1,5 @@
|
||||
use crate::runtime::{self, Runtime};
|
||||
use crate::runtime::tests::loom_oneshot as oneshot;
|
||||
use crate::runtime::{self, Runtime};
|
||||
use crate::spawn;
|
||||
|
||||
use loom::sync::atomic::{AtomicBool, AtomicUsize};
|
||||
|
||||
@@ -1,9 +1,9 @@
|
||||
use crate::loom::cell::CausalCell;
|
||||
use crate::loom::sync::Arc;
|
||||
use crate::park::Park;
|
||||
use crate::runtime::{self, blocking};
|
||||
use crate::runtime::park::Parker;
|
||||
use crate::runtime::thread_pool::{current, slice, Owned, Shared, Spawner};
|
||||
use crate::runtime::{self, blocking};
|
||||
use crate::task::Task;
|
||||
|
||||
use std::cell::Cell;
|
||||
@@ -78,10 +78,7 @@ struct GenerationGuard<'a> {
|
||||
struct WorkerGone;
|
||||
|
||||
// TODO: Move into slices
|
||||
pub(super) fn create_set(
|
||||
pool_size: usize,
|
||||
parker: Parker,
|
||||
) -> (Arc<slice::Set>, Vec<Worker>) {
|
||||
pub(super) fn create_set(pool_size: usize, parker: Parker) -> (Arc<slice::Set>, Vec<Worker>) {
|
||||
// Create the parks...
|
||||
let parkers: Vec<_> = (0..pool_size).map(|_| parker.clone()).collect();
|
||||
|
||||
@@ -95,13 +92,7 @@ pub(super) fn create_set(
|
||||
let workers = parkers
|
||||
.into_iter()
|
||||
.enumerate()
|
||||
.map(|(index, parker)| {
|
||||
Worker::new(
|
||||
slices.clone(),
|
||||
index,
|
||||
parker,
|
||||
)
|
||||
})
|
||||
.map(|(index, parker)| Worker::new(slices.clone(), index, parker))
|
||||
.collect();
|
||||
|
||||
(slices, workers)
|
||||
@@ -116,11 +107,7 @@ const GLOBAL_POLL_INTERVAL: u16 = 61;
|
||||
impl Worker {
|
||||
// Safe as aquiring a lock is required before doing anything potentially
|
||||
// dangerous.
|
||||
pub(super) fn new(
|
||||
slices: Arc<slice::Set>,
|
||||
index: usize,
|
||||
park: Parker,
|
||||
) -> Self {
|
||||
pub(super) fn new(slices: Arc<slice::Set>, index: usize, park: Parker) -> Self {
|
||||
Worker {
|
||||
inner: Arc::new(Inner {
|
||||
park: CausalCell::new(park),
|
||||
|
||||
@@ -109,8 +109,8 @@
|
||||
|
||||
use crate::loom::cell::CausalCell;
|
||||
use crate::loom::future::AtomicWaker;
|
||||
use crate::loom::sync::{Mutex, Arc, Condvar};
|
||||
use crate::loom::sync::atomic::{AtomicBool, AtomicPtr, AtomicUsize, spin_loop_hint};
|
||||
use crate::loom::sync::atomic::{spin_loop_hint, AtomicBool, AtomicPtr, AtomicUsize};
|
||||
use crate::loom::sync::{Arc, Condvar, Mutex};
|
||||
|
||||
use std::fmt;
|
||||
use std::ptr;
|
||||
@@ -387,10 +387,7 @@ pub fn channel<T>(mut capacity: usize) -> (Sender<T>, Receiver<T>) {
|
||||
let shared = Arc::new(Shared {
|
||||
buffer: buffer.into_boxed_slice(),
|
||||
mask: capacity - 1,
|
||||
tail: Mutex::new(Tail {
|
||||
pos: 0,
|
||||
rx_cnt: 1,
|
||||
}),
|
||||
tail: Mutex::new(Tail { pos: 0, rx_cnt: 1 }),
|
||||
condvar: Condvar::new(),
|
||||
wait_stack: AtomicPtr::new(ptr::null_mut()),
|
||||
num_tx: AtomicUsize::new(1),
|
||||
@@ -406,9 +403,7 @@ pub fn channel<T>(mut capacity: usize) -> (Sender<T>, Receiver<T>) {
|
||||
}),
|
||||
};
|
||||
|
||||
let tx = Sender {
|
||||
shared,
|
||||
};
|
||||
let tx = Sender { shared };
|
||||
|
||||
(tx, rx)
|
||||
}
|
||||
@@ -852,7 +847,10 @@ where
|
||||
// access to `self.wait.next`.
|
||||
self.wait.next.with_mut(|ptr| unsafe { *ptr = curr });
|
||||
|
||||
let res = self.shared.wait_stack.compare_exchange(curr, node, SeqCst, SeqCst);
|
||||
let res = self
|
||||
.shared
|
||||
.wait_stack
|
||||
.compare_exchange(curr, node, SeqCst, SeqCst);
|
||||
|
||||
match res {
|
||||
Ok(_) => return,
|
||||
|
||||
@@ -60,7 +60,9 @@ impl Semaphore {
|
||||
sem: &self,
|
||||
ll_permit: ll::Permit::new(),
|
||||
};
|
||||
poll_fn(|cx| permit.ll_permit.poll_acquire(cx, &self.ll_sem)).await.unwrap();
|
||||
poll_fn(|cx| permit.ll_permit.poll_acquire(cx, &self.ll_sem))
|
||||
.await
|
||||
.unwrap();
|
||||
permit
|
||||
}
|
||||
|
||||
@@ -68,10 +70,12 @@ impl Semaphore {
|
||||
pub fn try_acquire(&self) -> Result<SemaphorePermit<'_>, TryAcquireError> {
|
||||
let mut ll_permit = ll::Permit::new();
|
||||
match ll_permit.try_acquire(&self.ll_sem) {
|
||||
Ok(_) => Ok(SemaphorePermit { sem: self, ll_permit }),
|
||||
Ok(_) => Ok(SemaphorePermit {
|
||||
sem: self,
|
||||
ll_permit,
|
||||
}),
|
||||
Err(_) => Err(TryAcquireError(())),
|
||||
}
|
||||
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
@@ -66,9 +66,7 @@ fn broadcast_wrap() {
|
||||
match rx1.recv().await {
|
||||
Ok(_) => num += 1,
|
||||
Err(Closed) => break,
|
||||
Err(Lagged(n)) => {
|
||||
num += n as usize
|
||||
},
|
||||
Err(Lagged(n)) => num += n as usize,
|
||||
}
|
||||
}
|
||||
|
||||
@@ -84,9 +82,7 @@ fn broadcast_wrap() {
|
||||
match rx2.recv().await {
|
||||
Ok(_) => num += 1,
|
||||
Err(Closed) => break,
|
||||
Err(Lagged(n)) => {
|
||||
num += n as usize
|
||||
}
|
||||
Err(Lagged(n)) => num += n as usize,
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
@@ -94,7 +94,7 @@ pub use interval::{interval, interval_at, Interval};
|
||||
|
||||
mod timeout;
|
||||
#[doc(inline)]
|
||||
pub use timeout::{timeout, timeout_at, Timeout, Elapsed};
|
||||
pub use timeout::{timeout, timeout_at, Elapsed, Timeout};
|
||||
|
||||
cfg_stream! {
|
||||
mod throttle;
|
||||
|
||||
@@ -97,16 +97,11 @@ impl<T: Stream> Stream for Throttle<T> {
|
||||
|
||||
fn poll_next(mut self: Pin<&mut Self>, cx: &mut task::Context<'_>) -> Poll<Option<Self::Item>> {
|
||||
if !self.has_delayed && self.delay.is_some() {
|
||||
ready!(Pin::new(self.as_mut()
|
||||
.project().delay.as_mut().unwrap())
|
||||
.poll(cx));
|
||||
ready!(Pin::new(self.as_mut().project().delay.as_mut().unwrap()).poll(cx));
|
||||
*self.as_mut().project().has_delayed = true;
|
||||
}
|
||||
|
||||
let value = ready!(self
|
||||
.as_mut()
|
||||
.project().stream
|
||||
.poll_next(cx));
|
||||
let value = ready!(self.as_mut().project().stream.poll_next(cx));
|
||||
|
||||
if value.is_some() {
|
||||
let dur = self.duration;
|
||||
|
||||
@@ -21,10 +21,7 @@ impl Pack {
|
||||
pub(crate) const fn least_significant(width: u32) -> Pack {
|
||||
let mask = mask_for(width);
|
||||
|
||||
Pack {
|
||||
mask,
|
||||
shift: 0,
|
||||
}
|
||||
Pack { mask, shift: 0 }
|
||||
}
|
||||
|
||||
/// Value is packed in the `width` more-significant bits.
|
||||
@@ -32,10 +29,7 @@ impl Pack {
|
||||
let shift = pointer_width() - self.mask.leading_zeros();
|
||||
let mask = mask_for(width) << shift;
|
||||
|
||||
Pack {
|
||||
mask,
|
||||
shift,
|
||||
}
|
||||
Pack { mask, shift }
|
||||
}
|
||||
|
||||
/// Mask used to unpack value
|
||||
@@ -65,7 +59,11 @@ impl Pack {
|
||||
|
||||
impl fmt::Debug for Pack {
|
||||
fn fmt(&self, fmt: &mut fmt::Formatter<'_>) -> fmt::Result {
|
||||
write!(fmt, "Pack {{ mask: {:b}, shift: {} }}", self.mask, self.shift)
|
||||
write!(
|
||||
fmt,
|
||||
"Pack {{ mask: {:b}, shift: {} }}",
|
||||
self.mask, self.shift
|
||||
)
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
@@ -50,7 +50,7 @@
|
||||
//! ```
|
||||
|
||||
use crate::util::bit;
|
||||
use crate::util::slab::{Generation, MAX_PAGES, MAX_THREADS, INITIAL_PAGE_SIZE};
|
||||
use crate::util::slab::{Generation, INITIAL_PAGE_SIZE, MAX_PAGES, MAX_THREADS};
|
||||
|
||||
use std::usize;
|
||||
|
||||
@@ -61,15 +61,14 @@ pub(crate) struct Address(usize);
|
||||
const PAGE_INDEX_SHIFT: u32 = INITIAL_PAGE_SIZE.trailing_zeros() + 1;
|
||||
|
||||
/// Address in the shard
|
||||
const SLOT: bit::Pack = bit::Pack::least_significant(
|
||||
MAX_PAGES as u32 + PAGE_INDEX_SHIFT);
|
||||
const SLOT: bit::Pack = bit::Pack::least_significant(MAX_PAGES as u32 + PAGE_INDEX_SHIFT);
|
||||
|
||||
/// Masks the thread identifier
|
||||
const THREAD: bit::Pack = SLOT.then(MAX_THREADS.trailing_zeros() + 1);
|
||||
|
||||
/// Masks the generation
|
||||
const GENERATION: bit::Pack = THREAD.then(
|
||||
bit::pointer_width().wrapping_sub(RESERVED.width() + THREAD.width() + SLOT.width()));
|
||||
const GENERATION: bit::Pack = THREAD
|
||||
.then(bit::pointer_width().wrapping_sub(RESERVED.width() + THREAD.width() + SLOT.width()));
|
||||
|
||||
// Chosen arbitrarily
|
||||
const RESERVED: bit::Pack = bit::Pack::most_significant(5);
|
||||
|
||||
@@ -102,8 +102,6 @@ impl<T: Entry> Slab<T> {
|
||||
|
||||
impl<T> fmt::Debug for Slab<T> {
|
||||
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
|
||||
f.debug_struct("Slab")
|
||||
.field("shard", &self.shard)
|
||||
.finish()
|
||||
f.debug_struct("Slab").field("shard", &self.shard).finish()
|
||||
}
|
||||
}
|
||||
|
||||
@@ -1,4 +1,4 @@
|
||||
use crate::util::slab::{Address, Entry, page, MAX_PAGES};
|
||||
use crate::util::slab::{page, Address, Entry, MAX_PAGES};
|
||||
|
||||
use std::fmt;
|
||||
|
||||
@@ -49,10 +49,7 @@ impl<T: Entry> Shard<T> {
|
||||
|
||||
let local = (0..MAX_PAGES).map(|_| page::Local::new()).collect();
|
||||
|
||||
Shard {
|
||||
local,
|
||||
shared,
|
||||
}
|
||||
Shard { local, shared }
|
||||
}
|
||||
|
||||
pub(super) fn alloc(&self) -> Option<Address> {
|
||||
|
||||
@@ -1,5 +1,5 @@
|
||||
use crate::loom::cell::CausalCell;
|
||||
use crate::util::slab::{Generation, Entry};
|
||||
use crate::util::slab::{Entry, Generation};
|
||||
|
||||
/// Stores an entry in the slab.
|
||||
pub(super) struct Slot<T> {
|
||||
|
||||
@@ -15,10 +15,10 @@ pub(crate) struct LockGuard<'a, T> {
|
||||
_p: PhantomData<std::rc::Rc<()>>,
|
||||
}
|
||||
|
||||
unsafe impl<T: Send> Send for TryLock<T> { }
|
||||
unsafe impl<T: Send> Sync for TryLock<T> { }
|
||||
unsafe impl<T: Send> Send for TryLock<T> {}
|
||||
unsafe impl<T: Send> Sync for TryLock<T> {}
|
||||
|
||||
unsafe impl<T: Sync> Sync for LockGuard<'_, T> { }
|
||||
unsafe impl<T: Sync> Sync for LockGuard<'_, T> {}
|
||||
|
||||
impl<T> TryLock<T> {
|
||||
/// Create a new `TryLock`
|
||||
@@ -31,7 +31,11 @@ impl<T> TryLock<T> {
|
||||
|
||||
/// Attempt to acquire lock
|
||||
pub(crate) fn try_lock(&self) -> Option<LockGuard<'_, T>> {
|
||||
if self.locked.compare_exchange(false, true, SeqCst, SeqCst).is_err() {
|
||||
if self
|
||||
.locked
|
||||
.compare_exchange(false, true, SeqCst, SeqCst)
|
||||
.is_err()
|
||||
{
|
||||
return None;
|
||||
}
|
||||
|
||||
|
||||
Reference in New Issue
Block a user