use std::collections::HashMap;
use std::time::Duration;
use lsp_types::{ProgressParams, ProgressToken};
use serde::{Deserialize, Serialize};
use tokio::time::Instant;
use crate::config::ServerId;
use crate::lsp::types::ProgressKind;
#[derive(Debug, Clone, Copy, PartialEq, Eq, Default)]
pub enum IndexingState {
#[default]
Unknown,
Loading,
Ready,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Default, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum IndexingPolicy {
#[default]
Auto,
Disabled,
}
impl IndexingPolicy {
#[allow(clippy::trivially_copy_pass_by_ref)]
#[must_use]
pub(crate) const fn is_auto(&self) -> bool {
matches!(self, Self::Auto)
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
enum IndexingSignalSource {
ServerStatus,
Progress,
}
#[derive(Debug, Clone)]
struct IndexingEntry {
state: IndexingState,
source: IndexingSignalSource,
last_updated: Instant,
open: HashMap<ProgressToken, Instant>,
empty_since: Option<Instant>,
latched: bool,
}
impl IndexingEntry {
fn fresh_progress() -> Self {
Self {
state: IndexingState::Unknown,
source: IndexingSignalSource::Progress,
last_updated: Instant::now(),
open: HashMap::new(),
empty_since: None,
latched: false,
}
}
}
const SERVER_STATUS_METHOD: &str = "experimental/serverStatus";
const QUIESCENT_FIELD: &str = "quiescent";
pub const INDEXING_STALENESS_BOUND: Duration = Duration::from_secs(60);
pub const DEFAULT_INDEXING_READY_TIMEOUT_SECS: u64 = 30;
pub const PROGRESS_SETTLE: Duration = Duration::from_secs(3);
pub const PROGRESS_LATCH_IDLE: Duration = Duration::from_secs(30);
const PROGRESS_OPEN_CAP: usize = 64;
const PROGRESS_TOKEN_MAX_LEN: usize = 256;
#[derive(Debug, Default)]
pub struct IndexingTracker {
entries: HashMap<ServerId, IndexingEntry>,
policies: HashMap<ServerId, IndexingPolicy>,
}
impl IndexingTracker {
pub(crate) fn new() -> Self {
Self::default()
}
pub(crate) fn set_policy(&mut self, server_id: ServerId, policy: IndexingPolicy) {
self.policies.insert(server_id, policy);
}
fn is_disabled(&self, server_id: &ServerId) -> bool {
self.policies.get(server_id) == Some(&IndexingPolicy::Disabled)
}
pub(crate) fn observe_server_status(
&mut self,
server_id: &ServerId,
method: &str,
params: Option<&serde_json::Value>,
) {
if self.is_disabled(server_id) {
return;
}
if method != SERVER_STATUS_METHOD {
return;
}
let Some(quiescent) = params
.and_then(|p| p.get(QUIESCENT_FIELD))
.and_then(serde_json::Value::as_bool)
else {
return;
};
let sticky_ready = self.entries.get(server_id).is_some_and(|entry| {
entry.source == IndexingSignalSource::ServerStatus
&& entry.state == IndexingState::Ready
});
if sticky_ready {
return;
}
self.entries.insert(
server_id.clone(),
IndexingEntry {
state: if quiescent {
IndexingState::Ready
} else {
IndexingState::Loading
},
source: IndexingSignalSource::ServerStatus,
last_updated: Instant::now(),
open: HashMap::new(),
empty_since: None,
latched: false,
},
);
}
pub(crate) fn observe_progress(&mut self, server_id: &ServerId, params: &ProgressParams) {
if self.is_disabled(server_id) {
return;
}
let Some(kind) = ProgressKind::from_value(¶ms.value) else {
return;
};
if self
.entries
.get(server_id)
.is_some_and(|entry| entry.source == IndexingSignalSource::ServerStatus)
{
return;
}
if let ProgressToken::String(token) = ¶ms.token
&& token.len() > PROGRESS_TOKEN_MAX_LEN
{
return;
}
let entry = self
.entries
.entry(server_id.clone())
.or_insert_with(IndexingEntry::fresh_progress);
if entry.latched {
return;
}
let now = Instant::now();
entry
.open
.retain(|_, began| began.elapsed() < INDEXING_STALENESS_BOUND);
match kind {
ProgressKind::Begin => {
if entry.open.is_empty()
&& entry
.empty_since
.is_some_and(|since| since.elapsed() >= PROGRESS_LATCH_IDLE)
{
entry.latched = true;
entry.open.clear();
return;
}
if entry.open.len() < PROGRESS_OPEN_CAP {
entry.open.insert(params.token.clone(), now);
}
entry.empty_since = None;
}
ProgressKind::End => {
if entry.open.remove(¶ms.token).is_some() && entry.open.is_empty() {
entry.empty_since = Some(now);
}
}
}
entry.last_updated = now;
}
pub(crate) fn state(&self, server_id: &ServerId) -> IndexingState {
if self.is_disabled(server_id) {
return IndexingState::Unknown;
}
let Some(entry) = self.entries.get(server_id) else {
return IndexingState::Unknown;
};
match entry.source {
IndexingSignalSource::ServerStatus => {
if entry.state == IndexingState::Loading
&& entry.last_updated.elapsed() >= INDEXING_STALENESS_BOUND
{
IndexingState::Unknown
} else {
entry.state
}
}
IndexingSignalSource::Progress => {
if entry.latched {
return IndexingState::Ready;
}
if entry.last_updated.elapsed() >= INDEXING_STALENESS_BOUND {
return IndexingState::Unknown;
}
if entry.open.len() >= PROGRESS_OPEN_CAP {
return IndexingState::Unknown;
}
let has_live_open = entry
.open
.values()
.any(|began| began.elapsed() < INDEXING_STALENESS_BOUND);
if has_live_open {
return IndexingState::Loading;
}
if entry.open.is_empty() {
match entry.empty_since {
Some(since) if since.elapsed() < PROGRESS_SETTLE => IndexingState::Loading,
Some(_) => IndexingState::Ready,
None => IndexingState::Unknown,
}
} else {
IndexingState::Unknown
}
}
}
}
pub(crate) fn reset(&mut self, server_id: &ServerId) {
self.entries.remove(server_id);
}
}
#[cfg(test)]
#[allow(clippy::unwrap_used)]
mod tests {
use super::*;
fn test_server() -> ServerId {
ServerId::from("test-server")
}
fn progress(kind: &str, token: i32) -> ProgressParams {
ProgressParams {
token: ProgressToken::Int(token),
value: serde_json::json!({ "kind": kind }),
}
}
#[test]
fn test_state_defaults_unknown() {
let tracker = IndexingTracker::new();
assert_eq!(tracker.state(&test_server()), IndexingState::Unknown);
}
#[test]
fn test_observe_server_status_quiescent_false_marks_loading() {
let mut tracker = IndexingTracker::new();
let server = test_server();
tracker.observe_server_status(
&server,
"experimental/serverStatus",
Some(&serde_json::json!({"quiescent": false})),
);
assert_eq!(tracker.state(&server), IndexingState::Loading);
}
#[test]
fn test_observe_server_status_quiescent_true_marks_ready() {
let mut tracker = IndexingTracker::new();
let server = test_server();
tracker.observe_server_status(
&server,
"experimental/serverStatus",
Some(&serde_json::json!({"quiescent": true})),
);
assert_eq!(tracker.state(&server), IndexingState::Ready);
}
#[test]
fn test_observe_server_status_ignores_unrecognized_method() {
let mut tracker = IndexingTracker::new();
let server = test_server();
tracker.observe_server_status(
&server,
"window/logMessage",
Some(&serde_json::json!({"quiescent": false})),
);
assert_eq!(tracker.state(&server), IndexingState::Unknown);
}
#[test]
fn test_malformed_server_status_does_not_mark_source_or_disable_progress() {
let mut tracker = IndexingTracker::new();
let server = test_server();
tracker.observe_server_status(&server, "experimental/serverStatus", None);
assert_eq!(tracker.state(&server), IndexingState::Unknown);
tracker.observe_progress(&server, &progress("begin", 1));
assert_eq!(
tracker.state(&server),
IndexingState::Loading,
"the progress path must still be live after a malformed serverStatus payload"
);
}
#[test]
fn test_observe_server_status_ready_is_sticky() {
let mut tracker = IndexingTracker::new();
let server = test_server();
tracker.observe_server_status(
&server,
"experimental/serverStatus",
Some(&serde_json::json!({"quiescent": true})),
);
tracker.observe_server_status(
&server,
"experimental/serverStatus",
Some(&serde_json::json!({"quiescent": false})),
);
assert_eq!(tracker.state(&server), IndexingState::Ready);
}
#[test]
fn test_observe_server_status_tracks_servers_independently() {
let mut tracker = IndexingTracker::new();
let rust = ServerId::from("rust");
let python = ServerId::from("python");
tracker.observe_server_status(
&rust,
"experimental/serverStatus",
Some(&serde_json::json!({"quiescent": false})),
);
assert_eq!(tracker.state(&rust), IndexingState::Loading);
assert_eq!(tracker.state(&python), IndexingState::Unknown);
}
#[tokio::test(start_paused = true)]
async fn test_state_treats_stale_server_status_loading_as_unknown() {
let mut tracker = IndexingTracker::new();
let server = test_server();
tracker.observe_server_status(
&server,
"experimental/serverStatus",
Some(&serde_json::json!({"quiescent": false})),
);
assert_eq!(tracker.state(&server), IndexingState::Loading);
tokio::time::advance(INDEXING_STALENESS_BOUND.saturating_sub(Duration::from_secs(1))).await;
assert_eq!(tracker.state(&server), IndexingState::Loading);
tokio::time::advance(Duration::from_secs(2)).await;
assert_eq!(tracker.state(&server), IndexingState::Unknown);
}
#[tokio::test(start_paused = true)]
async fn test_observe_server_status_refreshes_staleness_clock() {
let mut tracker = IndexingTracker::new();
let server = test_server();
tracker.observe_server_status(
&server,
"experimental/serverStatus",
Some(&serde_json::json!({"quiescent": false})),
);
tokio::time::advance(INDEXING_STALENESS_BOUND.saturating_sub(Duration::from_secs(1))).await;
tracker.observe_server_status(
&server,
"experimental/serverStatus",
Some(&serde_json::json!({"quiescent": false})),
);
tokio::time::advance(INDEXING_STALENESS_BOUND.saturating_sub(Duration::from_secs(1))).await;
assert_eq!(tracker.state(&server), IndexingState::Loading);
}
#[test]
fn test_first_begin_on_fresh_entry_yields_loading_never_ready() {
let mut tracker = IndexingTracker::new();
let server = test_server();
tracker.observe_progress(&server, &progress("begin", 1));
assert_eq!(tracker.state(&server), IndexingState::Loading);
}
#[tokio::test(start_paused = true)]
async fn test_unmatched_end_does_not_create_empty_since() {
let mut tracker = IndexingTracker::new();
let server = test_server();
tracker.observe_progress(&server, &progress("begin", 1));
tracker.observe_progress(&server, &progress("end", 1));
tokio::time::advance(PROGRESS_SETTLE + Duration::from_secs(1)).await;
assert_eq!(
tracker.state(&server),
IndexingState::Ready,
"settled after the real end"
);
tracker.observe_progress(&server, &progress("end", 99));
assert_eq!(
tracker.state(&server),
IndexingState::Ready,
"an unmatched end must not manufacture a fresh settle window"
);
}
#[tokio::test(start_paused = true)]
async fn test_end_as_first_ever_frame_reads_unknown_not_ready() {
let mut tracker = IndexingTracker::new();
let server = test_server();
tracker.observe_progress(&server, &progress("end", 1));
assert_eq!(
tracker.state(&server),
IndexingState::Unknown,
"an end with no observed prior begin must not read Ready immediately"
);
tokio::time::advance(PROGRESS_SETTLE + Duration::from_secs(1)).await;
assert_eq!(
tracker.state(&server),
IndexingState::Unknown,
"must still read Unknown once the settle window would have expired -- there was \
never a real empty_since to settle from"
);
}
#[tokio::test(start_paused = true)]
async fn test_begin_after_latch_idle_gap_latches_permanently() {
let mut tracker = IndexingTracker::new();
let server = test_server();
tracker.observe_progress(&server, &progress("begin", 1));
tracker.observe_progress(&server, &progress("end", 1));
tokio::time::advance(PROGRESS_LATCH_IDLE + Duration::from_secs(1)).await;
tracker.observe_progress(&server, &progress("begin", 2));
assert_eq!(
tracker.state(&server),
IndexingState::Ready,
"a begin after PROGRESS_LATCH_IDLE of quiet must latch, not re-gate"
);
tracker.observe_progress(&server, &progress("end", 2));
tracker.observe_progress(&server, &progress("begin", 3));
assert_eq!(tracker.state(&server), IndexingState::Ready);
}
#[tokio::test(start_paused = true)]
async fn test_begin_after_mid_gap_regates_and_does_not_latch() {
let mut tracker = IndexingTracker::new();
let server = test_server();
tracker.observe_progress(&server, &progress("begin", 1));
tracker.observe_progress(&server, &progress("end", 1));
tokio::time::advance(Duration::from_secs(10)).await; tracker.observe_progress(&server, &progress("begin", 2));
assert_eq!(
tracker.state(&server),
IndexingState::Loading,
"a mid-gap begin must re-gate the next phase, not latch it away"
);
}
#[tokio::test(start_paused = true)]
async fn test_phase_gap_shorter_than_settle_does_not_read_ready() {
let mut tracker = IndexingTracker::new();
let server = test_server();
tracker.observe_progress(&server, &progress("begin", 1));
tracker.observe_progress(&server, &progress("end", 1));
tokio::time::advance(PROGRESS_SETTLE.checked_sub(Duration::from_secs(1)).unwrap()).await;
assert_eq!(tracker.state(&server), IndexingState::Loading);
}
#[test]
fn test_open_non_empty_reads_loading() {
let mut tracker = IndexingTracker::new();
let server = test_server();
tracker.observe_progress(&server, &progress("begin", 1));
tracker.observe_progress(&server, &progress("begin", 2));
tracker.observe_progress(&server, &progress("end", 1));
assert_eq!(
tracker.state(&server),
IndexingState::Loading,
"token 2 is still open"
);
}
#[test]
fn test_observe_progress_tracks_servers_independently() {
let mut tracker = IndexingTracker::new();
let rust = ServerId::from("rust");
let go = ServerId::from("go");
tracker.observe_progress(&rust, &progress("begin", 1));
assert_eq!(
tracker.state(&rust),
IndexingState::Loading,
"rust's token is still open"
);
assert_eq!(
tracker.state(&go),
IndexingState::Unknown,
"go has received no progress signal of its own"
);
}
#[tokio::test(start_paused = true)]
async fn test_progress_last_updated_past_staleness_bound_reads_unknown() {
let mut tracker = IndexingTracker::new();
let server = test_server();
tracker.observe_progress(&server, &progress("begin", 1));
tokio::time::advance(INDEXING_STALENESS_BOUND + Duration::from_secs(1)).await;
assert_eq!(
tracker.state(&server),
IndexingState::Unknown,
"a single op with no phase boundary for over a minute must fail open"
);
}
#[tokio::test(start_paused = true)]
async fn test_multiphase_load_past_staleness_bound_with_boundaries_stays_loading() {
let mut tracker = IndexingTracker::new();
let server = test_server();
tracker.observe_progress(&server, &progress("begin", 1));
tokio::time::advance(Duration::from_secs(40)).await;
tracker.observe_progress(&server, &progress("end", 1));
tracker.observe_progress(&server, &progress("begin", 2));
tokio::time::advance(Duration::from_secs(40)).await;
assert_eq!(tracker.state(&server), IndexingState::Loading);
}
#[tokio::test(start_paused = true)]
async fn test_lost_end_self_heals_after_staleness_bound() {
let mut tracker = IndexingTracker::new();
let server = test_server();
tracker.observe_progress(&server, &progress("begin", 1));
tokio::time::advance(Duration::from_secs(30)).await;
tracker.observe_progress(&server, &progress("begin", 2));
tracker.observe_progress(&server, &progress("end", 2));
tokio::time::advance(Duration::from_secs(25)).await; tracker.observe_progress(&server, &progress("begin", 3));
tracker.observe_progress(&server, &progress("end", 3));
assert_eq!(
tracker.state(&server),
IndexingState::Loading,
"token 1 is still within its own staleness window at t=55s"
);
tokio::time::advance(Duration::from_secs(6)).await; assert_eq!(
tracker.state(&server),
IndexingState::Unknown,
"a token whose own begin exceeded INDEXING_STALENESS_BOUND must stop pinning \
Loading even while unrelated later frames keep the entry-wide last_updated clock \
fresh"
);
}
#[test]
fn test_open_token_cap_forces_unknown() {
let mut tracker = IndexingTracker::new();
let server = test_server();
for i in 0..(PROGRESS_OPEN_CAP - 1) {
tracker.observe_progress(&server, &progress("begin", i32::try_from(i).unwrap()));
}
assert_eq!(
tracker.state(&server),
IndexingState::Loading,
"PROGRESS_OPEN_CAP - 1 distinct open tokens must still read Loading"
);
tracker.observe_progress(
&server,
&progress("begin", i32::try_from(PROGRESS_OPEN_CAP - 1).unwrap()),
);
assert_eq!(
tracker.state(&server),
IndexingState::Unknown,
"reaching PROGRESS_OPEN_CAP must fail open instead of staying Loading forever"
);
}
#[test]
fn test_oversized_string_token_is_rejected() {
let mut tracker = IndexingTracker::new();
let server = test_server();
let oversized = ProgressParams {
token: ProgressToken::String("x".repeat(PROGRESS_TOKEN_MAX_LEN + 1)),
value: serde_json::json!({ "kind": "begin" }),
};
tracker.observe_progress(&server, &oversized);
assert_eq!(
tracker.state(&server),
IndexingState::Unknown,
"an oversized token must be dropped outright, not tracked"
);
}
#[test]
fn test_malformed_progress_kind_ignored() {
let mut tracker = IndexingTracker::new();
let server = test_server();
tracker.observe_progress(&server, &progress("report", 1));
assert_eq!(tracker.state(&server), IndexingState::Unknown);
let missing_kind = ProgressParams {
token: ProgressToken::Int(1),
value: serde_json::json!({}),
};
tracker.observe_progress(&server, &missing_kind);
assert_eq!(tracker.state(&server), IndexingState::Unknown);
}
#[test]
fn test_disabled_policy_pins_unknown() {
let mut tracker = IndexingTracker::new();
let server = test_server();
tracker.set_policy(server.clone(), IndexingPolicy::Disabled);
tracker.observe_server_status(
&server,
"experimental/serverStatus",
Some(&serde_json::json!({"quiescent": false})),
);
tracker.observe_progress(&server, &progress("begin", 1));
assert_eq!(tracker.state(&server), IndexingState::Unknown);
}
#[tokio::test(start_paused = true)]
async fn test_progress_settle_then_quiescent_false_yields_loading() {
let mut tracker = IndexingTracker::new();
let server = test_server();
tracker.observe_progress(&server, &progress("begin", 1));
tracker.observe_progress(&server, &progress("end", 1));
tokio::time::advance(PROGRESS_SETTLE + Duration::from_secs(1)).await;
assert_eq!(tracker.state(&server), IndexingState::Ready);
tracker.observe_server_status(
&server,
"experimental/serverStatus",
Some(&serde_json::json!({"quiescent": false})),
);
assert_eq!(tracker.state(&server), IndexingState::Loading);
}
#[test]
fn test_server_status_entry_ignores_later_progress_frames() {
let mut tracker = IndexingTracker::new();
let server = test_server();
tracker.observe_server_status(
&server,
"experimental/serverStatus",
Some(&serde_json::json!({"quiescent": true})),
);
tracker.observe_progress(&server, &progress("begin", 1));
assert_eq!(tracker.state(&server), IndexingState::Ready);
}
#[test]
fn test_reset_clears_entry_but_not_policy() {
let mut tracker = IndexingTracker::new();
let server = test_server();
tracker.set_policy(server.clone(), IndexingPolicy::Disabled);
tracker.observe_progress(&server, &progress("begin", 1));
tracker.reset(&server);
assert_eq!(
tracker.state(&server),
IndexingState::Unknown,
"policy stays Disabled, so this reads Unknown regardless of the reset entry"
);
let mut auto_tracker = IndexingTracker::new();
auto_tracker.observe_progress(&server, &progress("begin", 1));
auto_tracker.reset(&server);
assert_eq!(auto_tracker.state(&server), IndexingState::Unknown);
}
}