File size: 1,443 Bytes
52a9af3 | 1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 30 31 32 33 34 35 36 37 38 39 40 41 42 43 44 45 46 47 48 49 50 51 52 53 54 55 56 57 | use std::sync::atomic::AtomicBool;
use std::sync::atomic::Ordering;
use tokio::sync::AcquireError;
use tokio::sync::Semaphore;
use tokio::sync::SemaphorePermit;
/// Owns MCP invalidation and the single gate used to publish runtime updates.
pub(super) struct McpRefresh {
pending: AtomicBool,
gate: Semaphore,
}
impl McpRefresh {
pub(super) fn new() -> Self {
Self {
pending: AtomicBool::new(false),
gate: Semaphore::new(/*permits*/ 1),
}
}
pub(super) fn invalidate(&self) {
self.pending.store(true, Ordering::Release);
}
#[cfg(test)]
pub(super) fn is_pending(&self) -> bool {
self.pending.load(Ordering::Acquire)
}
pub(super) fn claim(&self) -> bool {
self.pending.swap(false, Ordering::AcqRel)
}
#[tracing::instrument(name = "mcp.runtime.refresh_wait", skip_all)]
pub(super) async fn acquire(&self) -> Result<SemaphorePermit<'_>, AcquireError> {
self.gate.acquire().await
}
pub(super) fn close(&self) {
self.gate.close();
}
}
/// Restores a claimed refresh when its task is cancelled before publication.
pub(super) struct McpRefreshInvalidationGuard<'a> {
pub(super) refresh: &'a McpRefresh,
pub(super) published: bool,
}
impl Drop for McpRefreshInvalidationGuard<'_> {
fn drop(&mut self) {
if !self.published {
self.refresh.invalidate();
}
}
}
|