Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
40 changes: 27 additions & 13 deletions volo-http/Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -42,6 +42,7 @@ pin-project.workspace = true
simdutf8.workspace = true
thiserror.workspace = true
tokio = { workspace = true, features = [
"sync",
"fs",
"time",
"macros",
Expand All @@ -56,14 +57,14 @@ url.workspace = true
# =====optional=====

# server optional
ipnet = { workspace = true, optional = true } # client ip
matchit = { workspace = true, optional = true } # route matching
memchr = { workspace = true, optional = true } # sse
ipnet = { workspace = true, optional = true } # client ip
matchit = { workspace = true, optional = true } # route matching
memchr = { workspace = true, optional = true } # sse
scopeguard = { workspace = true, optional = true } # defer

# client optional
async-broadcast = { workspace = true, optional = true } # service discover
chrono = { workspace = true, optional = true } # stat
async-broadcast = { workspace = true, optional = true } # service discover
chrono = { workspace = true, optional = true } # stat
hickory-resolver = { workspace = true, optional = true } # dns resolver
mime_guess = { workspace = true, optional = true }

Expand Down Expand Up @@ -101,26 +102,39 @@ default-client = ["client", "http1", "json"]
default-server = ["server", "http1", "query", "form", "json", "multipart"]

full = [
"client", "server", # core
"http1", "http2", # protocol
"query", "form", "json", # serde
"tls", # https
"cookie", "multipart", "ws", # exts
"client",
"server", # core
"http1",
"http2", # protocol
"query",
"form",
"json", # serde
"tls", # https
"cookie",
"multipart",
"ws", # exts
]

http1 = ["hyper/http1", "hyper-util/http1"]
http2 = ["hyper/http2", "hyper-util/http2"]

client = [
"hyper/client",
"dep:async-broadcast", "dep:chrono", "dep:hickory-resolver",
"dep:async-broadcast",
"dep:chrono",
"dep:hickory-resolver",
] # client core
server = [
"hyper-util/server",
"dep:ipnet", "dep:matchit", "dep:memchr", "dep:scopeguard", "dep:mime_guess", "dep:chrono",
"dep:ipnet",
"dep:matchit",
"dep:memchr",
"dep:scopeguard",
"dep:mime_guess",
"dep:chrono",
] # server core

__serde = ["dep:serde"] # a private feature for enabling `serde` by `serde_xxx`
__serde = ["dep:serde"] # a private feature for enabling `serde` by `serde_xxx`
query = ["__serde", "dep:serde_urlencoded"]
form = ["__serde", "dep:serde_urlencoded"]
json = ["__serde", "dep:sonic-rs"]
Expand Down
59 changes: 59 additions & 0 deletions volo-http/src/client/transport/pool.rs
Original file line number Diff line number Diff line change
Expand Up @@ -25,6 +25,32 @@ pub struct Pool<K: Key, T> {
inner: Arc<Mutex<PoolInner<K, T>>>,
}

#[cfg(feature = "http1")]
impl<K: Key, T: Poolable> Pool<K, T> {
pub(crate) fn return_handle(&self) -> PoolReturn<K, T> {
PoolReturn {
inner: Arc::downgrade(&self.inner),
}
}
}

#[cfg(feature = "http1")]
pub(crate) struct PoolReturn<K: Key, T: Poolable> {
inner: Weak<Mutex<PoolInner<K, T>>>,
}

#[cfg(feature = "http1")]
impl<K: Key, T: Poolable> PoolReturn<K, T> {
pub(crate) fn put_ready(&self, key: K, value: T) -> Result<(), T> {
let Some(inner) = self.inner.upgrade() else {
return Err(value);
};

inner.lock().put(key, value, &inner);
Ok(())
}
}

// Before using a pooled connection, make sure the sender is not dead.
//
// This is a trait to allow the `client::pool::tests` to work for `i32`.
Expand Down Expand Up @@ -1000,4 +1026,37 @@ mod tests {

assert!(!pool.locked().idle.contains_key(&key));
}

#[tokio::test]
async fn return_handle_put_ready_unparks_checkout() {
let pool = pool_no_timer();
let key = host_key("foo");
let returner = pool.return_handle();
let mut checkout = pool.checkout(key.clone());

// Register the checkout as a waiter before returning the connection.
assert!(PollOnce(&mut checkout).await.is_none());

returner
.put_ready(key, Uniq(41))
.expect("the pool should still exist");

let pooled = checkout.await.expect("the waiter should be notified");
assert_eq!(*pooled, Uniq(41));
}

#[test]
fn return_handle_returns_value_after_pool_drop() {
let key = host_key("foo");
let returner = {
let pool = pool_no_timer::<KeyImpl, Uniq<i32>>();
pool.return_handle()
};

let value = returner
.put_ready(key, Uniq(41))
.expect_err("the weak pool reference should no longer upgrade");

assert_eq!(value, Uniq(41));
}
}
Loading
Loading