Skip to repository content540 lines · 22.9 KB · rust
tenant.openagents/omega
No repository description is available.
OpenAgents Git authority 2026-07-28T02:51:54.082Z Public web read
NIP-34 coordinate
30617:7649603503856e5148d571eac2766b288a8ff1e9e35d380337a1d2b0015b4f92:omegaMaintainersHidden in public view
References2 branches · 1 tag
Read-only clone
git clone https://openagents.com/git/tenant.openagents/omega.gitBrowse files
document_links.rs
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 ¤t_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