use std::sync::atomic::{AtomicU64, Ordering};
use std::sync::Arc;
use async_trait::async_trait;
use lspkit::{Cause, Generation, GenerationEvent};
use tokio::sync::broadcast;
use crate::scheduler::Trigger;
pub type AnalyzerResult = Result<(), AnalyzerError>;
#[non_exhaustive]
#[derive(Debug, thiserror::Error)]
pub enum AnalyzerError {
#[error("analyzer failed: {0}")]
Failed(String),
}
#[async_trait]
pub trait Analyzer: Send + Sync + 'static {
async fn analyze(&self, trigger: Trigger) -> AnalyzerResult;
}
#[derive(Clone)]
pub struct Session {
generation: Arc<AtomicU64>,
events: broadcast::Sender<GenerationEvent>,
}
impl Session {
#[must_use]
pub fn new() -> Self {
let (events, _) = broadcast::channel(32);
Self {
generation: Arc::new(AtomicU64::new(0)),
events,
}
}
#[must_use]
pub fn generation(&self) -> Generation {
Generation::new(self.generation.load(Ordering::SeqCst))
}
#[must_use]
pub fn subscribe(&self) -> broadcast::Receiver<GenerationEvent> {
self.events.subscribe()
}
#[must_use]
pub fn advance(&self, cause: Cause) -> Generation {
let next = self.generation.fetch_add(1, Ordering::SeqCst) + 1;
let event = GenerationEvent::new(Generation::new(next), cause);
let _ = self.events.send(event);
Generation::new(next)
}
pub async fn run<A: Analyzer>(
self,
mut triggers: tokio::sync::mpsc::UnboundedReceiver<Trigger>,
analyzer: A,
) {
while let Some(trigger) = triggers.recv().await {
match analyzer.analyze(trigger).await {
Ok(()) => {
let _ = self.advance(Cause::Rescan);
}
Err(err) => {
tracing::warn!(error = %err, "analyzer pass failed");
}
}
}
}
}
impl Default for Session {
fn default() -> Self {
Self::new()
}
}
impl std::fmt::Debug for Session {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("Session")
.field("generation", &self.generation())
.finish()
}
}