1use crate::response::Response;
4use hyper::upgrade::{OnUpgrade as HyperOnUpgrade, Upgraded};
5use std::sync::Arc;
6use tokio::sync::{OwnedSemaphorePermit, Semaphore};
7
8#[derive(Clone)]
10pub(crate) struct UpgradeBudget(pub Arc<Semaphore>);
11
12pub(crate) struct PendingUpgrade {
14 pub(crate) on_upgrade: HyperOnUpgrade,
15 pub(crate) budget: UpgradeBudget,
16}
17
18#[allow(dead_code)]
20pub struct UpgradePermit(OwnedSemaphorePermit);
21
22pub struct OnUpgrade {
26 inner: HyperOnUpgrade,
27 permit: OwnedSemaphorePermit,
28}
29
30impl OnUpgrade {
31 pub async fn upgrade(self) -> Result<(Upgraded, UpgradePermit), hyper::Error> {
36 let io = self.inner.await?;
37 Ok((io, UpgradePermit(self.permit)))
38 }
39}
40
41pub(crate) fn take_upgrade(pending: PendingUpgrade) -> Result<OnUpgrade, Box<Response>> {
43 let permit = match pending.budget.0.clone().try_acquire_owned() {
44 Ok(p) => p,
45 Err(_) => {
46 return Err(Box::new(
47 Response::text("Service Unavailable")
48 .status(503)
49 .header("retry-after", "5"),
50 ));
51 }
52 };
53 Ok(OnUpgrade {
54 inner: pending.on_upgrade,
55 permit,
56 })
57}