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(
43 pending: PendingUpgrade,
44) -> Result<OnUpgrade, Box<Response>> {
45 let permit = match pending.budget.0.clone().try_acquire_owned() {
46 Ok(p) => p,
47 Err(_) => {
48 return Err(Box::new(
49 Response::text("Service Unavailable")
50 .status(503)
51 .header("retry-after", "5"),
52 ));
53 }
54 };
55 Ok(OnUpgrade {
56 inner: pending.on_upgrade,
57 permit,
58 })
59}