1use std::{
4 convert::Infallible,
5 sync::{
6 atomic::{AtomicU64, Ordering},
7 Arc, Mutex as StdMutex, Weak,
8 },
9 time::Duration,
10};
11
12use async_trait::async_trait;
13use tokio::sync::{broadcast, watch, RwLock};
14use tokio_util::sync::CancellationToken;
15
16use super::{
17 project_layout::ProjectLayoutResolver,
18 registry::{
19 LocalCodeIntelligenceRegistry, RegistryAcquireError, RegistryConfig, RegistryKey,
20 RegistryKeyError, RegistryReport, RegistryShutdownError, RegistryShutdownFailure,
21 RuntimeLease,
22 },
23 workspace_runtime::WorkspaceRuntime,
24 CodeDiagnostic, CodeIntelligenceError, CodeIntelligenceResult, CodeIntelligenceState,
25 CodeIntelligenceStatus, CodeLocation, CodePosition, CodeQueryResult, DocumentSymbol,
26 NavigationKind, SymbolInformation, WorkspaceCodeIntelligence,
27};
28use crate::workspace::{
29 LocalWorkspaceManifest, LocalWorkspaceManifestSnapshot, WorkspaceFileChange,
30 WorkspaceFileSystem, WorkspacePath,
31};
32
33const DEFAULT_QUERY_TIMEOUT: Duration = Duration::from_secs(15);
34type RuntimeRegistry = LocalCodeIntelligenceRegistry<WorkspaceRuntime, Infallible>;
35type WorkspaceRuntimeLease = RuntimeLease<WorkspaceRuntime, Infallible>;
36
37pub struct LocalCodeIntelligence {
39 isolation_scope: String,
40 manifest: Arc<LocalWorkspaceManifest>,
41 file_system: Arc<dyn WorkspaceFileSystem>,
42 registry: RuntimeRegistry,
43 current: RwLock<Option<WorkspaceRuntimeLease>>,
44 status: watch::Sender<CodeIntelligenceStatus>,
45 generation: Arc<AtomicU64>,
46 lifetime: CancellationToken,
47 manifest_task: StdMutex<Option<tokio::task::JoinHandle<()>>>,
48 query_timeout: Duration,
49}
50
51impl std::fmt::Debug for LocalCodeIntelligence {
52 fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
53 formatter
54 .debug_struct("LocalCodeIntelligence")
55 .field("isolation_scope", &self.isolation_scope)
56 .field("manifest_root", &self.manifest.snapshot().root)
57 .field("generation", &self.generation.load(Ordering::Acquire))
58 .finish_non_exhaustive()
59 }
60}
61
62impl LocalCodeIntelligence {
63 pub async fn start(
65 isolation_scope: impl Into<String>,
66 manifest: Arc<LocalWorkspaceManifest>,
67 file_system: Arc<dyn WorkspaceFileSystem>,
68 ) -> CodeIntelligenceResult<Arc<Self>> {
69 Self::start_with_timeout(
70 isolation_scope,
71 manifest,
72 file_system,
73 DEFAULT_QUERY_TIMEOUT,
74 )
75 .await
76 }
77
78 pub(crate) async fn start_with_timeout(
79 isolation_scope: impl Into<String>,
80 manifest: Arc<LocalWorkspaceManifest>,
81 file_system: Arc<dyn WorkspaceFileSystem>,
82 query_timeout: Duration,
83 ) -> CodeIntelligenceResult<Arc<Self>> {
84 let isolation_scope = isolation_scope.into();
85 let snapshot_rx = manifest.subscribe();
86 let changes_rx = manifest.subscribe_changes();
87 let registry = RuntimeRegistry::new(
88 RegistryConfig::new(Duration::ZERO, 0),
89 |runtime: Arc<WorkspaceRuntime>| async move {
90 runtime.shutdown().await;
91 Ok(())
92 },
93 );
94 let (status, _) = watch::channel(CodeIntelligenceStatus {
95 state: CodeIntelligenceState::Starting,
96 message: Some("Code Intelligence is preparing the saved workspace".to_owned()),
97 ..CodeIntelligenceStatus::default()
98 });
99 let provider = Arc::new(Self {
100 isolation_scope,
101 manifest,
102 file_system,
103 registry,
104 current: RwLock::new(None),
105 status,
106 generation: Arc::new(AtomicU64::new(0)),
107 lifetime: CancellationToken::new(),
108 manifest_task: StdMutex::new(None),
109 query_timeout,
110 });
111
112 let bootstrap = Arc::clone(&provider);
117 let lifetime = provider.lifetime.clone();
118 tokio::spawn(async move {
119 tokio::select! {
120 _ = lifetime.cancelled() => {}
121 _ = bootstrap.refresh_snapshot(bootstrap.manifest.snapshot()) => {}
122 }
123 });
124 let weak = Arc::downgrade(&provider);
125 let task = tokio::spawn(run_manifest_updates(
126 weak,
127 snapshot_rx,
128 changes_rx,
129 provider.lifetime.clone(),
130 ));
131 *provider
132 .manifest_task
133 .lock()
134 .unwrap_or_else(std::sync::PoisonError::into_inner) = Some(task);
135 Ok(provider)
136 }
137
138 pub async fn shutdown(&self) {
140 self.lifetime.cancel();
141 let task = self
142 .manifest_task
143 .lock()
144 .unwrap_or_else(std::sync::PoisonError::into_inner)
145 .take();
146 if let Some(task) = task {
147 let _ = task.await;
148 }
149 let old = self.current.write().await.take();
150 drop(old);
151 report_registry_cleanup("shutdown", self.registry.shutdown_all().await);
152 self.status.send_replace(CodeIntelligenceStatus {
153 state: CodeIntelligenceState::Unavailable,
154 message: Some("Code Intelligence is shut down".to_owned()),
155 ..CodeIntelligenceStatus::default()
156 });
157 }
158
159 async fn refresh_snapshot(
160 self: &Arc<Self>,
161 snapshot: LocalWorkspaceManifestSnapshot,
162 ) -> CodeIntelligenceResult<()> {
163 let layout = ProjectLayoutResolver::resolve(&snapshot);
164 {
165 let current = self.current.read().await;
166 if let Some(runtime) = current
167 .as_ref()
168 .filter(|runtime| runtime.layout_hash() == layout.layout_hash)
169 {
170 runtime.update_snapshot(&snapshot).await;
171 return Ok(());
172 }
173 }
174
175 let key = RegistryKey::new(
176 self.isolation_scope.clone(),
177 &snapshot.root,
178 layout.layout_hash,
179 )
180 .await
181 .map_err(map_key_error)?;
182 let canonical_root = key.canonical_root().to_path_buf();
183 let runtime_snapshot = snapshot.clone();
184 let file_system = Arc::clone(&self.file_system);
185 let timeout = self.query_timeout;
186 let lease = self
187 .registry
188 .acquire(key, move |_| async move {
189 Ok(WorkspaceRuntime::new(
190 canonical_root,
191 layout,
192 &runtime_snapshot,
193 file_system,
194 timeout,
195 ))
196 })
197 .await
198 .map_err(map_acquire_error)?;
199
200 lease.update_snapshot(&snapshot).await;
201 let mut receiver = lease.subscribe_status();
202 let generation = self.generation.fetch_add(1, Ordering::AcqRel) + 1;
203 self.status.send_replace(receiver.borrow().clone());
204 let old = self.current.write().await.replace(lease);
205 drop(old);
206 report_registry_cleanup("layout refresh", self.registry.cleanup_idle().await);
207 self.spawn_status_forwarder(generation, &mut receiver);
208 Ok(())
209 }
210
211 fn spawn_status_forwarder(
212 &self,
213 generation: u64,
214 receiver: &mut watch::Receiver<CodeIntelligenceStatus>,
215 ) {
216 let mut receiver = receiver.clone();
217 let sender = self.status.clone();
218 let current_generation = Arc::clone(&self.generation);
219 let lifetime = self.lifetime.clone();
220 tokio::spawn(async move {
221 loop {
222 tokio::select! {
223 _ = lifetime.cancelled() => break,
224 changed = receiver.changed() => {
225 if changed.is_err()
226 || current_generation.load(Ordering::Acquire) != generation
227 {
228 break;
229 }
230 sender.send_replace(receiver.borrow().clone());
231 }
232 }
233 }
234 });
235 }
236
237 async fn handle_changes(&self, changes: &[WorkspaceFileChange]) {
238 let current = self.current.read().await;
239 if let Some(runtime) = current.as_ref() {
240 runtime.notify_file_changes(changes).await;
241 }
242 }
243
244 async fn runtime(
245 &self,
246 ) -> CodeIntelligenceResult<tokio::sync::RwLockReadGuard<'_, Option<WorkspaceRuntimeLease>>>
247 {
248 let current = self.current.read().await;
249 if current.is_none() {
250 return Err(CodeIntelligenceError::Unavailable {
251 message: "Code Intelligence has not prepared this workspace yet".to_owned(),
252 });
253 }
254 Ok(current)
255 }
256
257 fn report_update_error(&self, error: &CodeIntelligenceError) {
258 let mut status = self.status.borrow().clone();
259 status.state = CodeIntelligenceState::Degraded;
260 status.message = Some(format!("workspace refresh failed: {error}"));
261 self.status.send_replace(status);
262 }
263}
264
265#[async_trait]
266impl WorkspaceCodeIntelligence for LocalCodeIntelligence {
267 fn subscribe_status(&self) -> watch::Receiver<CodeIntelligenceStatus> {
268 self.status.subscribe()
269 }
270
271 async fn document_symbols(
272 &self,
273 path: &WorkspacePath,
274 cancellation: CancellationToken,
275 ) -> CodeIntelligenceResult<CodeQueryResult<DocumentSymbol>> {
276 let runtime = self.runtime().await?;
277 let Some(runtime) = runtime.as_ref() else {
278 return Err(runtime_unavailable());
279 };
280 runtime.document_symbols(path, cancellation).await
281 }
282
283 async fn search_symbols(
284 &self,
285 query: &str,
286 limit: usize,
287 cancellation: CancellationToken,
288 ) -> CodeIntelligenceResult<CodeQueryResult<SymbolInformation>> {
289 let runtime = self.runtime().await?;
290 let Some(runtime) = runtime.as_ref() else {
291 return Err(runtime_unavailable());
292 };
293 runtime.search_symbols(query, limit, cancellation).await
294 }
295
296 async fn navigate(
297 &self,
298 kind: NavigationKind,
299 path: &WorkspacePath,
300 position: CodePosition,
301 cancellation: CancellationToken,
302 ) -> CodeIntelligenceResult<CodeQueryResult<CodeLocation>> {
303 let runtime = self.runtime().await?;
304 let Some(runtime) = runtime.as_ref() else {
305 return Err(runtime_unavailable());
306 };
307 runtime.navigate(kind, path, position, cancellation).await
308 }
309
310 async fn diagnostics(
311 &self,
312 path: Option<&WorkspacePath>,
313 cancellation: CancellationToken,
314 ) -> CodeIntelligenceResult<CodeQueryResult<CodeDiagnostic>> {
315 let runtime = self.runtime().await?;
316 let Some(runtime) = runtime.as_ref() else {
317 return Err(runtime_unavailable());
318 };
319 runtime.diagnostics(path, cancellation).await
320 }
321}
322
323impl Drop for LocalCodeIntelligence {
324 fn drop(&mut self) {
325 self.lifetime.cancel();
326 if let Some(task) = self
327 .manifest_task
328 .lock()
329 .unwrap_or_else(std::sync::PoisonError::into_inner)
330 .take()
331 {
332 task.abort();
333 }
334 }
335}
336
337async fn run_manifest_updates(
338 provider: Weak<LocalCodeIntelligence>,
339 mut snapshots: broadcast::Receiver<LocalWorkspaceManifestSnapshot>,
340 mut changes: broadcast::Receiver<WorkspaceFileChange>,
341 lifetime: CancellationToken,
342) {
343 loop {
344 tokio::select! {
345 _ = lifetime.cancelled() => break,
346 update = snapshots.recv() => match update {
347 Ok(snapshot) => {
348 let Some(provider) = provider.upgrade() else { break; };
349 if let Err(error) = provider.refresh_snapshot(snapshot).await {
350 provider.report_update_error(&error);
351 }
352 }
353 Err(broadcast::error::RecvError::Lagged(_)) => {
354 let Some(provider) = provider.upgrade() else { break; };
355 if let Err(error) = provider.refresh_snapshot(provider.manifest.snapshot()).await {
356 provider.report_update_error(&error);
357 }
358 }
359 Err(broadcast::error::RecvError::Closed) => break,
360 },
361 update = changes.recv() => match update {
362 Ok(change) => {
363 let mut batch = vec![change];
364 while let Ok(change) = changes.try_recv() {
365 batch.push(change);
366 }
367 let Some(provider) = provider.upgrade() else { break; };
368 provider.handle_changes(&batch).await;
369 }
370 Err(broadcast::error::RecvError::Lagged(count)) => {
371 let Some(provider) = provider.upgrade() else { break; };
372 let mut status = provider.status.borrow().clone();
373 status.state = CodeIntelligenceState::Degraded;
374 status.message = Some(format!(
375 "workspace change stream skipped {count} events; saved documents will resynchronize on query"
376 ));
377 provider.status.send_replace(status);
378 }
379 Err(broadcast::error::RecvError::Closed) => break,
380 },
381 }
382 }
383}
384
385fn map_key_error(error: RegistryKeyError) -> CodeIntelligenceError {
386 CodeIntelligenceError::Unavailable {
387 message: error.to_string(),
388 }
389}
390
391fn map_acquire_error(error: RegistryAcquireError<Infallible>) -> CodeIntelligenceError {
392 match error {
393 RegistryAcquireError::ShuttingDown => CodeIntelligenceError::Unavailable {
394 message: "the Code Intelligence registry is shutting down".to_owned(),
395 },
396 RegistryAcquireError::Factory(error) => match *error {},
397 RegistryAcquireError::FactoryPanicked { message } => CodeIntelligenceError::Unavailable {
398 message: format!("Code Intelligence runtime initialization panicked: {message}"),
399 },
400 RegistryAcquireError::LeaseLimit => CodeIntelligenceError::Unavailable {
401 message: "the Code Intelligence runtime lease limit was exhausted".to_owned(),
402 },
403 }
404}
405
406fn runtime_unavailable() -> CodeIntelligenceError {
407 CodeIntelligenceError::Unavailable {
408 message: "Code Intelligence has not prepared this workspace yet".to_owned(),
409 }
410}
411
412fn report_registry_cleanup(context: &'static str, report: RegistryReport<Infallible>) {
413 if !report.removed.is_empty() {
414 tracing::debug!(
415 context,
416 retired = report.removed.len(),
417 "Code Intelligence retired workspace runtimes"
418 );
419 }
420 for RegistryShutdownError { key, failure } in report.errors {
421 match failure {
422 RegistryShutdownFailure::Runtime(error) => {
423 tracing::error!(
424 context,
425 workspace = ?key.canonical_root(),
426 ?error,
427 "Code Intelligence runtime cleanup returned an impossible error"
428 );
429 }
430 RegistryShutdownFailure::Panicked { message } => {
431 tracing::warn!(
432 context,
433 workspace = ?key.canonical_root(),
434 %message,
435 "Code Intelligence runtime cleanup panicked"
436 );
437 }
438 }
439 }
440}
441
442#[cfg(test)]
443mod tests {
444 use super::*;
445 use crate::workspace::{
446 LocalWorkspaceFile, LocalWorkspaceFileStatus, ManifestWorkspaceBackend,
447 };
448
449 fn file(path: &str) -> LocalWorkspaceFile {
450 LocalWorkspaceFile {
451 path: path.to_owned(),
452 size: 1,
453 modified_ms: Some(1),
454 language: None,
455 status: LocalWorkspaceFileStatus::Tracked,
456 binary: false,
457 generated: false,
458 }
459 }
460
461 async fn scanned_snapshot(manifest: &LocalWorkspaceManifest) -> LocalWorkspaceManifestSnapshot {
462 if manifest.snapshot().version > 0 {
463 return manifest.snapshot();
464 }
465 let mut snapshots = manifest.subscribe();
466 tokio::time::timeout(
467 crate::test_support::external_resource_start_timeout(Duration::from_secs(5)),
468 async {
469 loop {
470 let snapshot = snapshots.recv().await.unwrap();
471 if snapshot.version > 0 {
472 break snapshot;
473 }
474 }
475 },
476 )
477 .await
478 .expect("initial manifest scan should finish")
479 }
480
481 #[tokio::test]
482 async fn provider_reuses_unchanged_layout_and_shutdown_is_idempotent() {
483 let _permit = crate::test_support::resource_intensive_test_permit().await;
484 let workspace = tempfile::tempdir().unwrap();
485 std::fs::create_dir(workspace.path().join("src")).unwrap();
486 std::fs::write(workspace.path().join("Cargo.toml"), "[workspace]\n").unwrap();
487 std::fs::write(workspace.path().join("src/lib.rs"), "pub fn saved() {}\n").unwrap();
488 let backend = ManifestWorkspaceBackend::new(workspace.path());
489 let manifest = backend.manifest();
490 let initial = scanned_snapshot(&manifest).await;
491 let file_system: Arc<dyn WorkspaceFileSystem> = backend;
492 let provider = LocalCodeIntelligence::start_with_timeout(
493 "test-session",
494 manifest,
495 file_system,
496 Duration::from_secs(1),
497 )
498 .await
499 .unwrap();
500 tokio::time::timeout(Duration::from_secs(5), async {
501 loop {
502 if provider.generation.load(Ordering::Acquire) >= 1 {
503 break;
504 }
505 tokio::time::sleep(Duration::from_millis(10)).await;
506 }
507 })
508 .await
509 .expect("background bootstrap refresh should finish");
510 let generation = provider.generation.load(Ordering::Acquire);
511 assert_eq!(generation, 1);
512
513 let mut source_only = initial.clone();
514 source_only.version += 1;
515 source_only.files.push(file("src/new.rs"));
516 provider.refresh_snapshot(source_only).await.unwrap();
517 assert_eq!(provider.generation.load(Ordering::Acquire), generation);
518
519 let mut changed_layout = initial;
520 changed_layout.version += 2;
521 changed_layout
522 .files
523 .retain(|entry| entry.path != "Cargo.toml");
524 changed_layout.files.push(file("package.json"));
525 provider.refresh_snapshot(changed_layout).await.unwrap();
526 assert_eq!(provider.generation.load(Ordering::Acquire), generation + 1);
527
528 provider.shutdown().await;
529 provider.shutdown().await;
530 assert!(provider.current.read().await.is_none());
531 assert!(provider.lifetime.is_cancelled());
532 assert_eq!(
533 provider.status.borrow().state,
534 CodeIntelligenceState::Unavailable
535 );
536 }
537}