use std::sync::{Arc, Mutex};
use anyhow::Result;
use salsa::plumbing::current_revision;
use self::task::BackgroundTaskBuilder;
use self::thread::{JoinHandle, ThreadPriority};
use crate::config::Config;
use crate::server::client::{Client, Notifier, Requester, Responder};
use crate::server::connection::ClientSender;
use crate::server::schedule::task::BackgroundFnBuilder;
use crate::state::{MetaState, MetaStateInner, State};
mod task;
pub mod thread;
pub(super) use self::task::BackgroundSchedule;
pub use self::task::{Handler, RetryTaskInfo, SyncMutTask, Task};
use crate::server::schedule::task::SyncTask;
pub fn event_loop_thread(
func: impl FnOnce() -> Result<()> + Send + 'static,
) -> Result<JoinHandle<Result<()>>> {
const MAIN_THREAD_STACK_SIZE: usize = 2 * 1024 * 1024;
const MAIN_THREAD_NAME: &str = "cairols:main";
Ok(thread::Builder::new(ThreadPriority::LatencySensitive)
.name(MAIN_THREAD_NAME.into())
.stack_size(MAIN_THREAD_STACK_SIZE)
.spawn(func)?)
}
type SyncTaskHook = Box<dyn Fn(&mut State, MetaState, Notifier)>;
pub struct Scheduler<'s> {
state: &'s mut State,
client: Client<'s>,
background_pool: thread::Pool,
fmt_pool: thread::Pool,
sync_mut_task_hooks: Vec<SyncTaskHook>,
pub meta_state: MetaState,
}
impl<'s> Scheduler<'s> {
pub fn new(state: &'s mut State, sender: ClientSender) -> Self {
let analysis_event_sender =
state.analysis_progress_controller.server_tracker().events_sender();
let meta_state = Arc::new(Mutex::new(MetaStateInner::new(analysis_event_sender)));
Self {
state,
client: Client::new(sender),
background_pool: thread::Pool::new(usize::MAX, "worker"),
fmt_pool: thread::Pool::new(1, "fmt"),
sync_mut_task_hooks: Default::default(),
meta_state,
}
}
pub fn response(&mut self, response: lsp_server::Response) -> Task<'s> {
self.client.requester.pop_response_task(response)
}
pub fn dispatch(&mut self, task: Task<'s>) {
let build_task_fn = |func: BackgroundFnBuilder| {
let static_func = func(self.state, self.meta_state.clone());
let notifier = self.client.notifier();
let responder = self.client.responder();
move || static_func(notifier, responder)
};
match task {
Task::SyncMut(SyncMutTask { func }) => {
let notifier = self.client.notifier();
let responder = self.client.responder();
let old_revision = current_revision(&self.state.db);
let old_config: &Config = &self.state.config.snapshot();
func(self.state, notifier.clone(), &mut self.client.requester, responder);
let new_revision = current_revision(&self.state.db);
let new_config: &Config = &self.state.config.snapshot();
if old_revision != new_revision || old_config != new_config {
for hook in &self.sync_mut_task_hooks {
hook(self.state, self.meta_state.clone(), notifier.clone());
}
}
}
Task::Sync(SyncTask { func }) => {
let notifier = self.client.notifier();
let responder = self.client.responder();
func(
self.state,
self.meta_state.clone(),
notifier.clone(),
&mut self.client.requester,
responder,
);
}
Task::Background(BackgroundTaskBuilder { schedule, builder: func }) => {
let task = build_task_fn(func);
match schedule {
BackgroundSchedule::Worker => {
self.background_pool.spawn(ThreadPriority::Worker, task);
}
BackgroundSchedule::LatencySensitive => {
self.background_pool.spawn(ThreadPriority::LatencySensitive, task);
}
}
}
Task::Fmt(func) => {
let task = build_task_fn(func);
self.fmt_pool.spawn(ThreadPriority::LatencySensitive, task);
}
}
}
pub fn local_mut(
&mut self,
func: impl FnOnce(&mut State, Notifier, &mut Requester<'_>, Responder) + 's,
) {
self.dispatch(Task::local_mut(func));
}
pub fn local(
&mut self,
func: impl FnOnce(&State, MetaState, Notifier, &mut Requester<'_>, Responder) + 's,
) {
self.dispatch(Task::local(func));
}
pub fn on_sync_mut_task(&mut self, hook: impl Fn(&mut State, MetaState, Notifier) + 'static) {
self.sync_mut_task_hooks.push(Box::new(hook));
}
}