aster_forge_cloud_files_linux/
dispatch.rs1use std::{
4 future::Future,
5 panic::AssertUnwindSafe,
6 sync::{
7 Arc,
8 atomic::{AtomicBool, AtomicU64, Ordering},
9 },
10};
11
12use futures::FutureExt;
13use tokio::{runtime::Handle, sync::OwnedSemaphorePermit};
14
15use crate::{LinuxCloudFilesError, LinuxErrorCode, Result};
16
17#[derive(Default)]
18struct MetricsInner {
19 accepted: AtomicU64,
20 saturated: AtomicU64,
21 closing: AtomicU64,
22 completed: AtomicU64,
23 panicked: AtomicU64,
24}
25
26#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)]
28pub struct LinuxDispatchMetrics {
29 pub accepted: u64,
31 pub saturated: u64,
33 pub closing: u64,
35 pub completed: u64,
37 pub panicked: u64,
39}
40
41#[derive(Debug, Clone, Copy, PartialEq, Eq)]
43pub enum LinuxDispatchRejection {
44 Saturated,
46 Closing,
48}
49
50impl LinuxDispatchRejection {
51 #[must_use]
53 pub const fn error_code(self) -> LinuxErrorCode {
54 match self {
55 Self::Saturated => LinuxErrorCode::TryAgain,
56 Self::Closing => LinuxErrorCode::Io,
57 }
58 }
59}
60
61struct DispatcherInner {
62 runtime: Handle,
63 permits: Arc<tokio::sync::Semaphore>,
64 closing: AtomicBool,
65 metrics: MetricsInner,
66}
67
68#[derive(Clone)]
70pub struct LinuxRequestDispatcher {
71 inner: Arc<DispatcherInner>,
72}
73
74impl LinuxRequestDispatcher {
75 pub fn new(runtime: Handle, max_in_flight: usize) -> Result<Self> {
81 if max_in_flight == 0 {
82 return Err(LinuxCloudFilesError::InvalidConfiguration {
83 reason: "max in-flight FUSE requests must be greater than zero",
84 });
85 }
86 Ok(Self {
87 inner: Arc::new(DispatcherInner {
88 runtime,
89 permits: Arc::new(tokio::sync::Semaphore::new(max_in_flight)),
90 closing: AtomicBool::new(false),
91 metrics: MetricsInner::default(),
92 }),
93 })
94 }
95
96 pub fn reserve(&self) -> std::result::Result<LinuxDispatchReservation, LinuxDispatchRejection> {
102 if self.inner.closing.load(Ordering::Acquire) {
103 self.inner.metrics.closing.fetch_add(1, Ordering::Relaxed);
104 return Err(LinuxDispatchRejection::Closing);
105 }
106 match self.inner.permits.clone().try_acquire_owned() {
107 Ok(permit) => {
108 self.inner.metrics.accepted.fetch_add(1, Ordering::Relaxed);
109 Ok(LinuxDispatchReservation {
110 inner: self.inner.clone(),
111 permit,
112 })
113 }
114 Err(tokio::sync::TryAcquireError::NoPermits) => {
115 self.inner.metrics.saturated.fetch_add(1, Ordering::Relaxed);
116 Err(LinuxDispatchRejection::Saturated)
117 }
118 Err(tokio::sync::TryAcquireError::Closed) => {
119 self.inner.metrics.closing.fetch_add(1, Ordering::Relaxed);
120 Err(LinuxDispatchRejection::Closing)
121 }
122 }
123 }
124
125 pub fn spawn_cleanup<F>(&self, future: F)
132 where
133 F: Future<Output = ()> + Send + 'static,
134 {
135 self.inner.metrics.accepted.fetch_add(1, Ordering::Relaxed);
136 let inner = self.inner.clone();
137 let runtime = inner.runtime.clone();
138 let task = async move {
139 let outcome = AssertUnwindSafe(future).catch_unwind().await;
140 if outcome.is_err() {
141 inner.metrics.panicked.fetch_add(1, Ordering::Relaxed);
142 }
143 inner.metrics.completed.fetch_add(1, Ordering::Release);
144 };
145 std::mem::drop(runtime.spawn(task));
146 }
147
148 pub fn close(&self) {
150 if !self.inner.closing.swap(true, Ordering::AcqRel) {
151 self.inner.permits.close();
152 }
153 }
154
155 #[must_use]
157 pub fn is_closing(&self) -> bool {
158 self.inner.closing.load(Ordering::Acquire)
159 }
160
161 #[must_use]
163 pub fn metrics(&self) -> LinuxDispatchMetrics {
164 LinuxDispatchMetrics {
165 accepted: self.inner.metrics.accepted.load(Ordering::Relaxed),
166 saturated: self.inner.metrics.saturated.load(Ordering::Relaxed),
167 closing: self.inner.metrics.closing.load(Ordering::Relaxed),
168 completed: self.inner.metrics.completed.load(Ordering::Relaxed),
169 panicked: self.inner.metrics.panicked.load(Ordering::Relaxed),
170 }
171 }
172}
173
174#[must_use = "dropping a reservation releases its bounded capacity without dispatching work"]
176pub struct LinuxDispatchReservation {
177 inner: Arc<DispatcherInner>,
178 permit: OwnedSemaphorePermit,
179}
180
181impl LinuxDispatchReservation {
182 pub fn spawn<F>(self, future: F)
184 where
185 F: Future<Output = ()> + Send + 'static,
186 {
187 let inner = self.inner;
188 let runtime = inner.runtime.clone();
189 let permit = self.permit;
190 let task = async move {
191 let outcome = AssertUnwindSafe(future).catch_unwind().await;
192 if outcome.is_err() {
193 inner.metrics.panicked.fetch_add(1, Ordering::Relaxed);
194 }
195 inner.metrics.completed.fetch_add(1, Ordering::Release);
196 drop(permit);
197 };
198 std::mem::drop(runtime.spawn(task));
199 }
200}