Skip to main content

lspkit_live/
session.rs

1//! Session lifecycle.
2//!
3//! Wires a [`crate::watcher::FileWatcher`] to a [`crate::scheduler`] and an
4//! analyzer callback. The session owns the generation counter that
5//! [`lspkit::EngineApi`] implementations advance on every successful pass.
6
7use std::sync::atomic::{AtomicU64, Ordering};
8use std::sync::Arc;
9
10use async_trait::async_trait;
11use lspkit::{Cause, Generation, GenerationEvent};
12use tokio::sync::broadcast;
13
14use crate::scheduler::Trigger;
15
16/// Result of a single analysis pass.
17pub type AnalyzerResult = Result<(), AnalyzerError>;
18
19/// Errors emitted by an analyzer callback.
20#[non_exhaustive]
21#[derive(Debug, thiserror::Error)]
22pub enum AnalyzerError {
23    /// The analyzer encountered an unrecoverable problem; the session
24    /// keeps running but the generation is not advanced.
25    #[error("analyzer failed: {0}")]
26    Failed(String),
27}
28
29/// A pluggable analyzer driven by the session.
30#[async_trait]
31pub trait Analyzer: Send + Sync + 'static {
32    /// Called once per coalesced [`Trigger`]. Implementations should rebuild
33    /// or incrementally update their internal state.
34    async fn analyze(&self, trigger: Trigger) -> AnalyzerResult;
35}
36
37/// Session that drives an [`Analyzer`] from a scheduler stream.
38#[derive(Clone)]
39pub struct Session {
40    generation: Arc<AtomicU64>,
41    events: broadcast::Sender<GenerationEvent>,
42}
43
44impl Session {
45    /// New session at [`Generation::ZERO`].
46    #[must_use]
47    pub fn new() -> Self {
48        let (events, _) = broadcast::channel(32);
49        Self {
50            generation: Arc::new(AtomicU64::new(0)),
51            events,
52        }
53    }
54
55    /// Current generation.
56    #[must_use]
57    pub fn generation(&self) -> Generation {
58        Generation::new(self.generation.load(Ordering::SeqCst))
59    }
60
61    /// Subscribe to generation events.
62    #[must_use]
63    pub fn subscribe(&self) -> broadcast::Receiver<GenerationEvent> {
64        self.events.subscribe()
65    }
66
67    /// Bump the generation and broadcast a `Rescan` event.
68    #[must_use]
69    pub fn advance(&self, cause: Cause) -> Generation {
70        let next = self.generation.fetch_add(1, Ordering::SeqCst) + 1;
71        let event = GenerationEvent::new(Generation::new(next), cause);
72        let _ = self.events.send(event);
73        Generation::new(next)
74    }
75
76    /// Drive the session: read triggers, run the analyzer, advance generation
77    /// on success. Returns when `triggers` closes.
78    pub async fn run<A: Analyzer>(
79        self,
80        mut triggers: tokio::sync::mpsc::UnboundedReceiver<Trigger>,
81        analyzer: A,
82    ) {
83        while let Some(trigger) = triggers.recv().await {
84            match analyzer.analyze(trigger).await {
85                Ok(()) => {
86                    let _ = self.advance(Cause::Rescan);
87                }
88                Err(err) => {
89                    tracing::warn!(error = %err, "analyzer pass failed");
90                }
91            }
92        }
93    }
94}
95
96impl Default for Session {
97    fn default() -> Self {
98        Self::new()
99    }
100}
101
102impl std::fmt::Debug for Session {
103    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
104        f.debug_struct("Session")
105            .field("generation", &self.generation())
106            .finish()
107    }
108}