Skip to main content

a3s_code_core/code_intelligence/
local_provider.rs

1//! Local manifest-backed Code Intelligence provider.
2
3use 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
37/// Native local provider sharing the workspace manifest's existing watcher.
38pub 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    /// Create a provider and acquire a cheap, lazily-started runtime generation.
64    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        // Do not block interactive TUI takeover on the first layout resolve.
113        // The manifest stays inactive until after the first flushed frame; the
114        // update task below refreshes as soon as activation/snapshots arrive.
115        // A best-effort background pass covers hosts that never activate.
116        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    /// Stop manifest forwarding and all language processes owned by this provider.
139    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}