aster_forge_cloud_files_linux/
dispatch.rs

1//! Bounded non-blocking bridge from FUSE callback threads to an explicit Tokio runtime.
2
3use 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/// Snapshot of bounded callback-to-async dispatch outcomes.
27#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)]
28pub struct LinuxDispatchMetrics {
29    /// Requests accepted into the explicit runtime.
30    pub accepted: u64,
31    /// Requests rejected because every in-flight permit was occupied.
32    pub saturated: u64,
33    /// Requests rejected after dispatcher closing began.
34    pub closing: u64,
35    /// Accepted tasks that reached a normal or panic terminal state.
36    pub completed: u64,
37    /// Accepted task futures that panicked. Dropping their owned FUSE reply produces `EIO`.
38    pub panicked: u64,
39}
40
41/// Reason a callback could not reserve bounded asynchronous execution capacity.
42#[derive(Debug, Clone, Copy, PartialEq, Eq)]
43pub enum LinuxDispatchRejection {
44    /// Every configured in-flight permit was occupied.
45    Saturated,
46    /// The mount session is closing and accepts no new work.
47    Closing,
48}
49
50impl LinuxDispatchRejection {
51    /// Returns the FUSE-visible portable error classification.
52    #[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/// Explicit bounded dispatcher shared by one mount session.
69#[derive(Clone)]
70pub struct LinuxRequestDispatcher {
71    inner: Arc<DispatcherInner>,
72}
73
74impl LinuxRequestDispatcher {
75    /// Creates a dispatcher backed by a caller-owned runtime.
76    /// # Errors
77    ///
78    /// Returns an error when validation fails or an underlying backend, store, or platform
79    /// operation fails.
80    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    /// Reserves one in-flight slot without blocking the callback thread.
97    /// # Errors
98    ///
99    /// Returns an error when validation fails or an underlying backend, store, or platform
100    /// operation fails.
101    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    /// Dispatches mandatory cleanup for an already accepted native handle.
126    ///
127    /// Cleanup bypasses the normal request semaphore and closing fence because FUSE issues one
128    /// final `release` per successful `open`; rejecting that cleanup would leak adapter and
129    /// product-owned session state. Callers must use this only for bounded-by-open-handle terminal
130    /// work, never for new backend requests.
131    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    /// Begins idempotent closing and rejects every subsequent reservation.
149    pub fn close(&self) {
150        if !self.inner.closing.swap(true, Ordering::AcqRel) {
151            self.inner.permits.close();
152        }
153    }
154
155    /// Returns whether new callback work is rejected.
156    #[must_use]
157    pub fn is_closing(&self) -> bool {
158        self.inner.closing.load(Ordering::Acquire)
159    }
160
161    /// Returns an atomic metrics snapshot.
162    #[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/// Owned capacity reservation that transfers one native reply into asynchronous work.
175#[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    /// Spawns one accepted reply-owning future on the caller-provided runtime.
183    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}