use std::collections::VecDeque;
use parking_lot::Mutex;
use serde::{Deserialize, Serialize};
use strop_core::id::{BufferRevision, DocumentId};
use strop_core::worker::Outcome;
use strop_worker_protocol::message::NotifyCoverage;
use strop_worker_protocol::{Event, NotifyHint, Subscription};
use strop_workspace::ResourceLocation;
const MAX_RECORDS: usize = 64;
const MAX_HINTS: usize = 512;
#[derive(Serialize, Deserialize)]
pub(crate) enum Record {
Settled(Outcome<SubscribedScope>),
Hints {
subscription: Subscription,
hints: Vec<NotifyHint>,
},
Overflow { subscription: Subscription },
Boundary { subscription: Subscription },
RemoteSettled {
root: ResourceLocation,
outcome: Outcome<SubscribedScope>,
},
}
#[derive(Serialize, Deserialize)]
pub(crate) struct SubscribedScope {
pub subscription: Subscription,
pub coverage: NotifyCoverage,
pub owned_trace: Option<Vec<u8>>,
}
#[derive(Default)]
struct QueueState {
records: VecDeque<Record>,
rescan: bool,
owned_trace: Option<(Subscription, Vec<u8>)>,
}
pub(crate) struct NotifyQueue {
state: Mutex<QueueState>,
}
impl NotifyQueue {
pub(super) fn new() -> Self {
Self {
state: Mutex::new(QueueState::default()),
}
}
pub(crate) fn set_owned_trace(&self, owned: Option<(Subscription, Vec<u8>)>) {
self.state.lock().owned_trace = owned;
}
pub(crate) fn push_event(&self, event: Event) -> bool {
match event {
Event::Notify {
subscription,
hints,
..
} => self.push_hints(subscription, hints),
Event::NotifyOverflow { subscription, .. } => {
self.push_record(Record::Overflow { subscription });
true
}
Event::ReconcileBoundary { subscription, .. } => {
self.push_record(Record::Boundary { subscription });
true
}
Event::ExecExit { .. } | Event::ExecInput { .. } => false,
}
}
pub(crate) fn push_hints(
&self,
subscription: Subscription,
mut hints: Vec<NotifyHint>,
) -> bool {
let mut state = self.state.lock();
if strop_trace::enabled() {
if let Some((owner, path)) = &state.owned_trace {
if *owner == subscription {
hints.retain(|hint| hint.path.as_slice() != path.as_slice());
}
}
}
if hints.is_empty() {
return false;
}
if let Some(Record::Hints {
subscription: tail_subscription,
hints: tail,
}) = state.records.back_mut()
{
if *tail_subscription == subscription && tail.len() + hints.len() <= MAX_HINTS {
tail.append(&mut hints);
return true;
}
}
if hints.len() > MAX_HINTS {
state.records.push_back(Record::Overflow { subscription });
return true;
}
if state.records.len() >= MAX_RECORDS {
state.rescan = true;
return true;
}
state.records.push_back(Record::Hints {
subscription,
hints,
});
true
}
pub(crate) fn push_record(&self, record: Record) {
let mut state = self.state.lock();
if state.records.len() >= MAX_RECORDS {
state.rescan = true;
return;
}
state.records.push_back(record);
}
pub(super) fn requeue_front(&self, records: Vec<Record>) {
let mut state = self.state.lock();
let room = MAX_RECORDS.saturating_sub(state.records.len());
if records.len() > room {
state.rescan = true;
}
for record in records.into_iter().take(room).rev() {
state.records.push_front(record);
}
}
pub(crate) fn drain(&self) -> (Vec<Record>, bool) {
let mut state = self.state.lock();
(
state.records.drain(..).collect(),
std::mem::take(&mut state.rescan),
)
}
}
#[derive(Debug, Clone, PartialEq, Eq, serde::Serialize, serde::Deserialize)]
pub struct ReloadKey {
pub document: DocumentId,
pub revision: BufferRevision,
#[serde(with = "strop_core::path_serde")]
pub path: std::path::PathBuf,
}