use crate::response::Response;
use hyper::upgrade::{OnUpgrade as HyperOnUpgrade, Upgraded};
use std::sync::Arc;
use tokio::sync::{OwnedSemaphorePermit, Semaphore};
#[derive(Clone)]
pub(crate) struct UpgradeBudget(pub Arc<Semaphore>);
pub(crate) struct PendingUpgrade {
pub(crate) on_upgrade: HyperOnUpgrade,
pub(crate) budget: UpgradeBudget,
}
#[allow(dead_code)]
pub struct UpgradePermit(OwnedSemaphorePermit);
pub struct OnUpgrade {
inner: HyperOnUpgrade,
permit: OwnedSemaphorePermit,
}
impl OnUpgrade {
pub async fn upgrade(self) -> Result<(Upgraded, UpgradePermit), hyper::Error> {
let io = self.inner.await?;
Ok((io, UpgradePermit(self.permit)))
}
}
pub(crate) fn take_upgrade(
pending: PendingUpgrade,
) -> Result<OnUpgrade, Box<Response>> {
let permit = match pending.budget.0.clone().try_acquire_owned() {
Ok(p) => p,
Err(_) => {
return Err(Box::new(
Response::text("Service Unavailable")
.status(503)
.header("retry-after", "5"),
));
}
};
Ok(OnUpgrade {
inner: pending.on_upgrade,
permit,
})
}