Skip to repository content

tenant.openagents/omega

No repository description is available.

OpenAgents Git authority 2026-07-28T03:59:28.509Z Public web read
NIP-34 coordinate30617:7649603503856e5148d571eac2766b288a8ff1e9e35d380337a1d2b0015b4f92:omega
MaintainersHidden in public view
References2 branches · 1 tag
Read-only clonegit clone https://openagents.com/git/tenant.openagents/omega.git
Browse files

document_links.rs

540 lines · 22.9 KB · rust
1use std::ops::Range;
2use std::str::FromStr as _;
3use std::sync::Arc;
4use std::time::Duration;
5
6use anyhow::Context as _;
7use collections::{HashMap, HashSet};
8use futures::FutureExt as _;
9use futures::future::{Shared, join_all};
10use gpui::{AppContext as _, AsyncApp, Context, Entity, SharedString, Task};
11use language::{Buffer, point_to_lsp};
12use lsp::LanguageServerId;
13use lsp::request::DocumentLinkResolve;
14use rpc::{TypedEnvelope, proto};
15use settings::Settings as _;
16use text::{Anchor, BufferId, ToPointUtf16 as _};
17use util::ResultExt as _;
18
19use crate::lsp_command::{GetDocumentLinks, LspCommand as _};
20use crate::lsp_store::{
21    LspStore, LspStoreEvent, RunningFetch, missing_servers_to_query, next_lsp_fetch_id,
22};
23use crate::project_settings::ProjectSettings;
24
25#[derive(Copy, Clone, Hash, PartialEq, Eq, PartialOrd, Ord, Debug)]
26pub struct DocumentLinkId(u64);
27
28#[derive(Clone, Debug)]
29pub struct LspDocumentLink {
30    pub range: Range<Anchor>,
31    pub target: Option<SharedString>,
32    pub tooltip: Option<SharedString>,
33    pub data: Option<serde_json::Value>,
34    pub resolved: bool,
35}
36
37pub type BufferDocumentLinks = HashMap<LanguageServerId, HashMap<DocumentLinkId, LspDocumentLink>>;
38
39pub(super) type DocumentLinksTask =
40    Shared<Task<std::result::Result<Option<BufferDocumentLinks>, Arc<anyhow::Error>>>>;
41
42pub type DocumentLinkResolveTask = Shared<Task<Option<(DocumentLinkId, LspDocumentLink)>>>;
43
44#[derive(Debug, Default)]
45pub(super) struct DocumentLinksData {
46    pub(super) links: BufferDocumentLinks,
47    fetched_servers: HashSet<LanguageServerId>,
48    pub(super) next_id: u64,
49    links_update: Option<RunningFetch<DocumentLinksTask>>,
50    pub(super) link_resolves: HashMap<(LanguageServerId, DocumentLinkId), DocumentLinkResolveTask>,
51}
52
53impl DocumentLinksData {
54    pub(super) fn remove_server_data(&mut self, server_id: LanguageServerId) {
55        self.links.remove(&server_id);
56        self.fetched_servers.remove(&server_id);
57        self.link_resolves
58            .retain(|(resolved_server, _), _| *resolved_server != server_id);
59        RunningFetch::discard_if_queried(&mut self.links_update, server_id);
60    }
61
62    fn evict(&mut self, for_server: Option<LanguageServerId>) {
63        match for_server {
64            Some(server_id) => self.remove_server_data(server_id),
65            None => {
66                self.links.clear();
67                self.fetched_servers.clear();
68                self.link_resolves.clear();
69            }
70        }
71        self.links_update = None;
72    }
73}
74
75/// Mirror of [`crate::lsp_store::ResolvedHint`] for document links: callers
76/// either get the resolved entry directly, an in-flight `Shared` task to await
77/// (deduplicated across editors), or `None` when the cache no longer contains
78/// a matching link.
79pub enum ResolvedDocumentLink {
80    Resolved(LspDocumentLink),
81    Resolving(DocumentLinkResolveTask),
82}
83
84impl LspStore {
85    pub(super) fn refresh_document_links(
86        &mut self,
87        for_server: Option<LanguageServerId>,
88        cx: &mut Context<Self>,
89    ) {
90        for lsp_data in self.lsp_data.values_mut() {
91            if let Some(document_links) = &mut lsp_data.document_links {
92                document_links.evict(for_server);
93            }
94        }
95
96        cx.emit(LspStoreEvent::RefreshDocumentLinks {
97            server_id: for_server,
98        });
99        if let Some((downstream_client, project_id)) = self.downstream_client.as_ref() {
100            downstream_client
101                .send(proto::RefreshDocumentLinks {
102                    project_id: *project_id,
103                    server_id: for_server.map(|server_id| server_id.to_proto()),
104                })
105                .context("sending refresh document links downstream")
106                .log_err();
107        }
108    }
109
110    pub(super) async fn handle_refresh_document_links(
111        lsp_store: Entity<Self>,
112        envelope: TypedEnvelope<proto::RefreshDocumentLinks>,
113        mut cx: AsyncApp,
114    ) -> anyhow::Result<proto::Ack> {
115        lsp_store.update(&mut cx, |lsp_store, cx| {
116            let server_id = envelope.payload.server_id.map(LanguageServerId::from_proto);
117            lsp_store.refresh_document_links(server_id, cx);
118        });
119        Ok(proto::Ack {})
120    }
121
122    /// `Some(..)` means the underlying state was actually refreshed; `None`
123    /// means the fetch was skipped or failed, and the caller should keep its
124    /// previous data.
125    pub fn fetch_document_links(
126        &mut self,
127        buffer: &Entity<Buffer>,
128        cx: &mut Context<Self>,
129    ) -> Task<Option<BufferDocumentLinks>> {
130        let version_queried_for = buffer.read(cx).version();
131        let buffer_id = buffer.read(cx).remote_id();
132
133        let current_servers = self.relevant_server_ids_for_capability_check(buffer, cx);
134
135        let mut servers_to_query = None;
136        if let Some(lsp_data) = self.current_lsp_data(buffer_id) {
137            if !version_queried_for.changed_since(&lsp_data.buffer_version)
138                && let Some(cached) = &mut lsp_data.document_links
139            {
140                match missing_servers_to_query(
141                    &mut cached.links,
142                    &mut cached.fetched_servers,
143                    &current_servers,
144                ) {
145                    Some(missing_servers) => servers_to_query = Some(missing_servers),
146                    None => return Task::ready(Some(cached.links.clone())),
147                }
148            }
149            if let Some(document_links) = &lsp_data.document_links
150                && let Some(running) = &document_links.links_update
151                && !version_queried_for.changed_since(&running.version)
152                && servers_to_query
153                    .as_ref()
154                    .is_none_or(|missing| missing.is_subset(&running.servers))
155            {
156                let running = running.task.clone();
157                return cx.background_spawn(async move { running.await.ok().flatten() });
158            }
159        }
160
161        let links_lsp_data = self
162            .latest_lsp_data(buffer, cx)
163            .document_links
164            .get_or_insert_default();
165        let fetch_id = next_lsp_fetch_id();
166        let queried_servers = servers_to_query
167            .clone()
168            .unwrap_or_else(|| current_servers.clone());
169        let buffer = buffer.clone();
170        let query_version = version_queried_for.clone();
171        let new_task = cx
172            .spawn({
173                let queried_servers = queried_servers.clone();
174                async move |lsp_store, cx| {
175                    cx.background_executor()
176                        .timer(Duration::from_millis(30))
177                        .await;
178
179                    let fetched = lsp_store
180                        .update(cx, |lsp_store, cx| {
181                            lsp_store.fetch_document_links_for_buffer(&buffer, servers_to_query, cx)
182                        })
183                        .map_err(Arc::new)?
184                        .await
185                        .context("fetching document links")
186                        .map_err(Arc::new);
187
188                    let fetched = match fetched {
189                        Ok(fetched) => fetched,
190                        Err(e) => {
191                            lsp_store
192                                .update(cx, |lsp_store, _| {
193                                    if let Some(lsp_data) = lsp_store.lsp_data.get_mut(&buffer_id)
194                                        && let Some(document_links) = &mut lsp_data.document_links
195                                    {
196                                        RunningFetch::take_finished(
197                                            &mut document_links.links_update,
198                                            fetch_id,
199                                        );
200                                    }
201                                })
202                                .ok();
203                            return Err(e);
204                        }
205                    };
206
207                    lsp_store
208                        .update(cx, |lsp_store, cx| {
209                            let lsp_data = lsp_store.latest_lsp_data(&buffer, cx);
210                            let links_data = lsp_data.document_links.get_or_insert_default();
211                            if !RunningFetch::take_finished(&mut links_data.links_update, fetch_id)
212                            {
213                                return Some(links_data.links.clone());
214                            }
215
216                            let Some(fetched_links) = fetched else {
217                                return None;
218                            };
219
220                            let mut tagged = BufferDocumentLinks::default();
221                            for (server_id, server_links) in fetched_links {
222                                let mut by_id = HashMap::default();
223                                by_id.reserve(server_links.len());
224                                for link in server_links {
225                                    let id = DocumentLinkId(links_data.next_id);
226                                    links_data.next_id += 1;
227                                    by_id.insert(id, link);
228                                }
229                                tagged.insert(server_id, by_id);
230                            }
231
232                            if lsp_data.buffer_version == query_version {
233                                for (server_id, new_links) in &tagged {
234                                    links_data.links.insert(*server_id, new_links.clone());
235                                }
236                                links_data.fetched_servers.extend(queried_servers);
237                                // The newly inserted links are unresolved by definition; drop any
238                                // pending resolves that were keyed against the prior entries for
239                                // those servers so callers re-issue against the fresh ids.
240                                links_data.link_resolves.clear();
241                                Some(links_data.links.clone())
242                            } else if !lsp_data.buffer_version.changed_since(&query_version) {
243                                lsp_data.buffer_version = query_version;
244                                links_data.links = tagged;
245                                links_data.fetched_servers = queried_servers;
246                                links_data.link_resolves.clear();
247                                Some(links_data.links.clone())
248                            } else {
249                                None
250                            }
251                        })
252                        .map_err(Arc::new)
253                }
254            })
255            .shared();
256
257        links_lsp_data.links_update = Some(RunningFetch {
258            id: fetch_id,
259            version: version_queried_for,
260            servers: queried_servers,
261            task: new_task.clone(),
262        });
263
264        cx.background_spawn(async move { new_task.await.ok().flatten() })
265    }
266
267    #[cfg(any(test, feature = "test-support"))]
268    pub fn document_links_for_buffer(&self, buffer_id: BufferId) -> Option<BufferDocumentLinks> {
269        let data = self.lsp_data.get(&buffer_id)?;
270        let document_links = data.document_links.as_ref()?;
271        Some(document_links.links.clone())
272    }
273
274    fn fetch_document_links_for_buffer(
275        &mut self,
276        buffer: &Entity<Buffer>,
277        for_servers: Option<HashSet<LanguageServerId>>,
278        cx: &mut Context<Self>,
279    ) -> Task<anyhow::Result<Option<HashMap<LanguageServerId, Vec<LspDocumentLink>>>>> {
280        if let Some((client, project_id)) = self.upstream_client() {
281            // No `for_servers` filter is forwarded: unlike its siblings, `GetDocumentLinks`
282            // is answered from the host's own `fetch_document_links` cache with the full
283            // per-server map, to spare the LSP request (see collab's
284            // `test_lsp_document_links`), so filtering could not reduce the work anyway.
285            let request = GetDocumentLinks;
286            if !self.is_capable_for_proto_request(buffer, &request, cx) {
287                return Task::ready(Ok(None));
288            }
289
290            let request_timeout = ProjectSettings::get_global(cx)
291                .global_lsp_settings
292                .get_request_timeout();
293            let request_task = client.request_lsp(
294                project_id,
295                None,
296                request_timeout,
297                cx.background_executor().clone(),
298                request.to_proto(project_id, buffer.read(cx)),
299            );
300            let buffer = buffer.clone();
301            cx.spawn(async move |weak_lsp_store, cx| {
302                let Some(lsp_store) = weak_lsp_store.upgrade() else {
303                    return Ok(None);
304                };
305                let Some(responses) = request_task.await? else {
306                    return Ok(None);
307                };
308
309                let document_links = join_all(responses.payload.into_iter().map(|response| {
310                    let lsp_store = lsp_store.clone();
311                    let buffer = buffer.clone();
312                    let cx = cx.clone();
313                    async move {
314                        let server_id = LanguageServerId::from_proto(response.server_id);
315                        let links = GetDocumentLinks
316                            .response_from_proto(response.response, lsp_store, buffer, cx)
317                            .await;
318                        (server_id, links)
319                    }
320                }))
321                .await;
322
323                let mut has_errors = false;
324                let result = document_links
325                    .into_iter()
326                    .filter_map(|(server_id, links)| match links {
327                        Ok(links) => Some((server_id, links)),
328                        Err(e) => {
329                            has_errors = true;
330                            log::error!(
331                                "Failed to fetch document links for server {server_id}: {e:#}"
332                            );
333                            None
334                        }
335                    })
336                    .collect::<HashMap<_, _>>();
337                anyhow::ensure!(
338                    !has_errors || !result.is_empty(),
339                    "Failed to fetch document links"
340                );
341                Ok(Some(result))
342            })
343        } else {
344            let links_task = self.request_filtered_lsp_locally(
345                buffer,
346                None::<usize>,
347                GetDocumentLinks,
348                for_servers.as_ref(),
349                cx,
350            );
351            cx.background_spawn(async move { Ok(Some(links_task.await.into_iter().collect())) })
352        }
353    }
354
355    /// Returns the resolved state for a cached document link, deduplicating
356    /// in-flight `documentLink/resolve` requests across editors via a `Shared`
357    /// task stored on `DocumentLinksData`.
358    ///
359    /// `link_id` is the [`DocumentLinkId`] stamped on the cached link by
360    /// [`Self::fetch_document_links`]; sibling links sharing the same buffer
361    /// range are disambiguated by it. `None` is returned when the cache no
362    /// longer holds a matching link (likely a version bump in between).
363    pub fn resolved_document_link(
364        &mut self,
365        buffer: &Entity<Buffer>,
366        server_id: LanguageServerId,
367        link_id: DocumentLinkId,
368        cx: &mut Context<Self>,
369    ) -> Option<ResolvedDocumentLink> {
370        let buffer_id = buffer.read(cx).remote_id();
371
372        let document_links = self.lsp_data.get(&buffer_id)?.document_links.as_ref()?;
373        let cached_link = document_links.links.get(&server_id)?.get(&link_id)?.clone();
374
375        if cached_link.resolved {
376            return Some(ResolvedDocumentLink::Resolved(cached_link));
377        }
378
379        let key = (server_id, link_id);
380        if let Some(running) = document_links.link_resolves.get(&key) {
381            return Some(ResolvedDocumentLink::Resolving(running.clone()));
382        }
383
384        let resolve_task = self.resolve_document_link_request(buffer, server_id, &cached_link, cx);
385        let query_version = self.lsp_data.get(&buffer_id)?.buffer_version.clone();
386        let resolve_task = cx
387            .spawn(async move |lsp_store, cx| {
388                let resolved = resolve_task.await;
389                lsp_store
390                    .update(cx, |lsp_store, _| {
391                        let lsp_data = lsp_store.lsp_data.get_mut(&buffer_id)?;
392                        if lsp_data.buffer_version != query_version {
393                            return None;
394                        }
395                        let links_data = lsp_data.document_links.as_mut()?;
396                        links_data.link_resolves.remove(&key);
397                        let updated = match resolved {
398                            Some(resolved) => lsp_store
399                                .cache_resolved_link(buffer_id, server_id, link_id, &resolved)?,
400                            None => {
401                                // No further resolution is possible (no capability,
402                                // missing server, or LSP error); mark as resolved so we
403                                // do not keep retrying on every hover, and yield the
404                                // entry as-is so awaiters can still surface it.
405                                let links_data = lsp_data.document_links.as_mut()?;
406                                let link =
407                                    links_data.links.get_mut(&server_id)?.get_mut(&link_id)?;
408                                link.resolved = true;
409                                link.clone()
410                            }
411                        };
412                        Some((link_id, updated))
413                    })
414                    .ok()
415                    .flatten()
416            })
417            .shared();
418
419        let document_links = self.lsp_data.get_mut(&buffer_id)?.document_links.as_mut()?;
420        document_links
421            .link_resolves
422            .insert(key, resolve_task.clone());
423        Some(ResolvedDocumentLink::Resolving(resolve_task))
424    }
425
426    /// Builds the LSP/proto request task for a single unresolved link. Returns
427    /// a task that yields `None` when the resolve request cannot be issued
428    /// (no upstream capability, no local server, or no `resolveProvider`).
429    fn resolve_document_link_request(
430        &self,
431        buffer: &Entity<Buffer>,
432        server_id: LanguageServerId,
433        cached_link: &LspDocumentLink,
434        cx: &mut Context<Self>,
435    ) -> Task<Option<lsp::DocumentLink>> {
436        let snapshot = buffer.read(cx).snapshot();
437        let buffer_id = buffer.read(cx).remote_id();
438        let lsp_link = lsp::DocumentLink {
439            range: lsp::Range {
440                start: point_to_lsp(cached_link.range.start.to_point_utf16(&snapshot)),
441                end: point_to_lsp(cached_link.range.end.to_point_utf16(&snapshot)),
442            },
443            target: cached_link
444                .target
445                .as_ref()
446                .and_then(|s| lsp::Uri::from_str(s).ok()),
447            tooltip: cached_link.tooltip.as_deref().map(str::to_string),
448            data: cached_link.data.clone(),
449        };
450
451        if let Some((upstream_client, project_id)) = self.upstream_client() {
452            if !self.check_if_capable_for_proto_request(buffer, can_resolve_link, cx) {
453                return Task::ready(None);
454            }
455            let request = proto::ResolveDocumentLink {
456                project_id,
457                buffer_id: buffer_id.into(),
458                language_server_id: server_id.0 as u64,
459                lsp_link: serde_json::to_vec(&lsp_link).unwrap_or_default(),
460            };
461            cx.background_spawn(async move {
462                let response = upstream_client.request(request).await.log_err()?;
463                serde_json::from_slice::<lsp::DocumentLink>(&response.lsp_link).log_err()
464            })
465        } else {
466            let Some(server) = self.language_server_for_id(server_id) else {
467                return Task::ready(None);
468            };
469            if !can_resolve_link(&server.capabilities()) {
470                return Task::ready(None);
471            }
472            let request_timeout = ProjectSettings::get_global(cx)
473                .global_lsp_settings
474                .get_request_timeout();
475            cx.background_spawn(async move {
476                server
477                    .request::<DocumentLinkResolve>(lsp_link, request_timeout)
478                    .await
479                    .into_response()
480                    .log_err()
481            })
482        }
483    }
484
485    fn cache_resolved_link(
486        &mut self,
487        buffer_id: BufferId,
488        server_id: LanguageServerId,
489        link_id: DocumentLinkId,
490        resolved: &lsp::DocumentLink,
491    ) -> Option<LspDocumentLink> {
492        let document_links = self.lsp_data.get_mut(&buffer_id)?.document_links.as_mut()?;
493        let link = document_links
494            .links
495            .get_mut(&server_id)?
496            .get_mut(&link_id)?;
497        link.target = resolved.target.as_ref().map(|u| u.to_string().into());
498        if let Some(tooltip) = &resolved.tooltip {
499            link.tooltip = Some(tooltip.clone().into());
500        }
501        link.data = resolved.data.clone();
502        link.resolved = true;
503        Some(link.clone())
504    }
505
506    pub(super) async fn handle_resolve_document_link(
507        lsp_store: Entity<Self>,
508        envelope: TypedEnvelope<proto::ResolveDocumentLink>,
509        mut cx: AsyncApp,
510    ) -> anyhow::Result<proto::ResolveDocumentLinkResponse> {
511        let lsp_link: lsp::DocumentLink = serde_json::from_slice(&envelope.payload.lsp_link)
512            .context("deserializing document link to resolve")?;
513        let server_id = LanguageServerId::from_proto(envelope.payload.language_server_id);
514
515        let resolve_task = lsp_store.update(&mut cx, |lsp_store, cx| {
516            let server = lsp_store
517                .language_server_for_id(server_id)
518                .with_context(|| format!("No language server {server_id}"))?;
519            let timeout = ProjectSettings::get_global(cx)
520                .global_lsp_settings
521                .get_request_timeout();
522            anyhow::Ok(server.request::<DocumentLinkResolve>(lsp_link, timeout))
523        })?;
524        let resolved = resolve_task.await.into_response()?;
525
526        Ok(proto::ResolveDocumentLinkResponse {
527            lsp_link: serde_json::to_vec(&resolved)
528                .context("serializing resolved document link")?,
529        })
530    }
531}
532
533fn can_resolve_link(capabilities: &lsp::ServerCapabilities) -> bool {
534    capabilities
535        .document_link_provider
536        .as_ref()
537        .and_then(|opts| opts.resolve_provider)
538        .unwrap_or(false)
539}
540
Served at tenant.openagents/omega Member data and write actions are omitted.