Skip to repository content570 lines · 23.1 KB · rust
tenant.openagents/omega
No repository description is available.
OpenAgents Git authority 2026-07-28T02:50:51.426Z 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
code_lens.rs
1use std::ops::Range;
2use std::sync::Arc;
3
4use anyhow::{Context as _, Result};
5use collections::{HashMap, HashSet};
6use futures::{
7 FutureExt as _,
8 future::{Shared, join_all},
9};
10use gpui::{AppContext as _, AsyncApp, Context, Entity, Task};
11use language::{Anchor, Buffer};
12use lsp::LanguageServerId;
13use rpc::{TypedEnvelope, proto};
14use settings::Settings as _;
15use std::time::Duration;
16use text::OffsetRangeExt as _;
17use util::ResultExt as _;
18
19use crate::{
20 CodeAction, LspAction, LspStore, LspStoreEvent, Project,
21 lsp_command::{GetCodeLens, LspCommand as _},
22 lsp_store::{
23 RunningFetch, missing_servers_to_query, next_lsp_fetch_id, upstream_lsp_query_server_filter,
24 },
25 project_settings::ProjectSettings,
26};
27
28/// Opaque per-action identifier issued by [`LspStore`] at fetch time.
29///
30/// LSP `CodeLens.data` is the server's private payload for resolve
31/// round-trips, so we can't use it (or anything derived from it) to
32/// disambiguate sibling lenses that share the same buffer `range`
33/// (TypeScript's references + implementations is the canonical case).
34/// We tag every cached action with this id and require it back on resolve
35/// so each lens routes to its own request and slot.
36///
37/// Ids are issued in fetch order; sorting by id reproduces server-emit
38/// order, which is how callers recover a stable render order without
39/// paying for an ordered map.
40#[derive(Copy, Clone, Hash, PartialEq, Eq, PartialOrd, Ord, Debug)]
41pub struct CodeLensActionId(u64);
42
43pub type CodeLensActions = HashMap<CodeLensActionId, CodeAction>;
44
45pub(super) type CodeLensTask =
46 Shared<Task<std::result::Result<Option<CodeLensActions>, Arc<anyhow::Error>>>>;
47
48pub type CodeLensResolveTask = Shared<Task<Option<(CodeLensActionId, CodeAction)>>>;
49
50#[derive(Debug, Default)]
51pub(super) struct CodeLensData {
52 pub(super) lens: HashMap<LanguageServerId, CodeLensActions>,
53 fetched_servers: HashSet<LanguageServerId>,
54 pub(super) next_id: u64,
55 pub(super) update: Option<RunningFetch<CodeLensTask>>,
56 pub(super) resolving: HashMap<(LanguageServerId, CodeLensActionId), CodeLensResolveTask>,
57}
58
59impl CodeLensData {
60 pub(super) fn remove_server_data(&mut self, server_id: LanguageServerId) {
61 self.lens.remove(&server_id);
62 self.fetched_servers.remove(&server_id);
63 self.resolving.retain(|(s, _), _| *s != server_id);
64 RunningFetch::discard_if_queried(&mut self.update, server_id);
65 }
66
67 fn evict(&mut self, for_server: Option<LanguageServerId>) {
68 match for_server {
69 Some(server_id) => self.remove_server_data(server_id),
70 None => {
71 self.lens.clear();
72 self.fetched_servers.clear();
73 self.resolving.clear();
74 }
75 }
76 self.update = None;
77 }
78}
79
80fn flatten_cache(lens: &HashMap<LanguageServerId, CodeLensActions>) -> CodeLensActions {
81 let mut out = CodeLensActions::default();
82 out.reserve(lens.values().map(|per_server| per_server.len()).sum());
83 for per_server in lens.values() {
84 for (id, action) in per_server {
85 out.insert(*id, action.clone());
86 }
87 }
88 out
89}
90
91impl LspStore {
92 pub(super) fn refresh_code_lens(
93 &mut self,
94 for_server: Option<LanguageServerId>,
95 cx: &mut Context<Self>,
96 ) {
97 for lsp_data in self.lsp_data.values_mut() {
98 if let Some(code_lens) = &mut lsp_data.code_lens {
99 code_lens.evict(for_server);
100 }
101 }
102
103 cx.emit(LspStoreEvent::RefreshCodeLens {
104 server_id: for_server,
105 });
106 if let Some((downstream_client, project_id)) = self.downstream_client.as_ref() {
107 downstream_client
108 .send(proto::RefreshCodeLens {
109 project_id: *project_id,
110 server_id: for_server.map(|server_id| server_id.to_proto()),
111 })
112 .context("sending refresh code lens downstream")
113 .log_err();
114 }
115 }
116
117 /// Fetches all code lenses for the buffer, each tagged with the
118 /// [`CodeLensActionId`] that callers must pass back to
119 /// [`Self::resolve_code_lens`]. Resolution is the caller's job.
120 pub fn code_lens_actions(
121 &mut self,
122 buffer: &Entity<Buffer>,
123 cx: &mut Context<Self>,
124 ) -> Task<Result<Option<CodeLensActions>>> {
125 let buffer_id = buffer.read(cx).remote_id();
126 let fetch_task = self.fetch_code_lenses(buffer, cx);
127
128 cx.spawn(async move |lsp_store, cx| {
129 fetch_task
130 .await
131 .map_err(|e| anyhow::anyhow!("code lens fetch failed: {e:#}"))?;
132
133 let actions = lsp_store.read_with(cx, |lsp_store, _| {
134 lsp_store
135 .lsp_data
136 .get(&buffer_id)
137 .and_then(|data| data.code_lens.as_ref())
138 .map(|code_lens| flatten_cache(&code_lens.lens))
139 })?;
140 Ok(actions)
141 })
142 }
143
144 fn fetch_code_lenses(
145 &mut self,
146 buffer: &Entity<Buffer>,
147 cx: &mut Context<Self>,
148 ) -> CodeLensTask {
149 let version_queried_for = buffer.read(cx).version();
150 let buffer_id = buffer.read(cx).remote_id();
151 let current_servers = self.relevant_server_ids_for_capability_check(buffer, cx);
152
153 let mut servers_to_query = None;
154 if let Some(lsp_data) = self.current_lsp_data(buffer_id) {
155 if let Some(cached_lens) = &mut lsp_data.code_lens {
156 if !version_queried_for.changed_since(&lsp_data.buffer_version) {
157 match missing_servers_to_query(
158 &mut cached_lens.lens,
159 &mut cached_lens.fetched_servers,
160 ¤t_servers,
161 ) {
162 Some(missing_servers) => servers_to_query = Some(missing_servers),
163 None => {
164 return Task::ready(Ok(Some(flatten_cache(&cached_lens.lens))))
165 .shared();
166 }
167 }
168 }
169 if let Some(running) = cached_lens.update.as_ref()
170 && !version_queried_for.changed_since(&running.version)
171 && servers_to_query
172 .as_ref()
173 .is_none_or(|missing| missing.is_subset(&running.servers))
174 {
175 return running.task.clone();
176 }
177 }
178 }
179
180 let lens_lsp_data = self
181 .latest_lsp_data(buffer, cx)
182 .code_lens
183 .get_or_insert_default();
184 let fetch_id = next_lsp_fetch_id();
185 let queried_servers = servers_to_query
186 .clone()
187 .unwrap_or_else(|| current_servers.clone());
188 let buffer = buffer.clone();
189 let query_version_queried_for = version_queried_for.clone();
190 let new_task = cx
191 .spawn({
192 let queried_servers = queried_servers.clone();
193 async move |lsp_store, cx| {
194 cx.background_executor()
195 .timer(Duration::from_millis(30))
196 .await;
197 let fetched_lens = lsp_store
198 .update(cx, |lsp_store, cx| {
199 lsp_store.fetch_code_lens_for_buffer(&buffer, servers_to_query, cx)
200 })
201 .map_err(Arc::new)?
202 .await
203 .context("fetching code lens")
204 .map_err(Arc::new);
205 let fetched_lens = match fetched_lens {
206 Ok(fetched_lens) => fetched_lens,
207 Err(e) => {
208 lsp_store
209 .update(cx, |lsp_store, _| {
210 if let Some(lens_lsp_data) = lsp_store
211 .lsp_data
212 .get_mut(&buffer_id)
213 .and_then(|lsp_data| lsp_data.code_lens.as_mut())
214 {
215 RunningFetch::take_finished(
216 &mut lens_lsp_data.update,
217 fetch_id,
218 );
219 }
220 })
221 .ok();
222 return Err(e);
223 }
224 };
225
226 lsp_store
227 .update(cx, |lsp_store, _| {
228 let lsp_data = lsp_store.current_lsp_data(buffer_id)?;
229 let code_lens = lsp_data.code_lens.as_mut()?;
230 if !RunningFetch::take_finished(&mut code_lens.update, fetch_id) {
231 return Some(flatten_cache(&code_lens.lens));
232 }
233 if let Some(fetched_lens) = fetched_lens {
234 let mut tagged: HashMap<LanguageServerId, CodeLensActions> =
235 HashMap::default();
236 for (server_id, actions) in fetched_lens {
237 let mut cache = CodeLensActions::default();
238 cache.reserve(actions.len());
239 for action in actions {
240 let id = CodeLensActionId(code_lens.next_id);
241 code_lens.next_id += 1;
242 cache.insert(id, action);
243 }
244 tagged.insert(server_id, cache);
245 }
246 if lsp_data.buffer_version == query_version_queried_for {
247 code_lens.lens.extend(tagged);
248 code_lens.fetched_servers.extend(queried_servers);
249 } else if !lsp_data
250 .buffer_version
251 .changed_since(&query_version_queried_for)
252 {
253 lsp_data.buffer_version = query_version_queried_for;
254 code_lens.lens = tagged;
255 code_lens.fetched_servers = queried_servers;
256 }
257 }
258 Some(flatten_cache(&code_lens.lens))
259 })
260 .map_err(Arc::new)
261 }
262 })
263 .shared();
264 lens_lsp_data.update = Some(RunningFetch {
265 id: fetch_id,
266 version: version_queried_for,
267 servers: queried_servers,
268 task: new_task.clone(),
269 });
270 new_task
271 }
272
273 fn fetch_code_lens_for_buffer(
274 &mut self,
275 buffer: &Entity<Buffer>,
276 for_servers: Option<HashSet<LanguageServerId>>,
277 cx: &mut Context<Self>,
278 ) -> Task<Result<Option<HashMap<LanguageServerId, Vec<CodeAction>>>>> {
279 if let Some((upstream_client, project_id)) = self.upstream_client() {
280 let request = GetCodeLens;
281 if !self.is_capable_for_proto_request(buffer, &request, cx) {
282 return Task::ready(Ok(None));
283 }
284 let request_timeout = ProjectSettings::get_global(cx)
285 .global_lsp_settings
286 .get_request_timeout();
287 let request_task = upstream_client.request_lsp(
288 project_id,
289 upstream_lsp_query_server_filter(for_servers.as_ref()),
290 request_timeout,
291 cx.background_executor().clone(),
292 request.to_proto(project_id, buffer.read(cx)),
293 );
294 let buffer = buffer.clone();
295 cx.spawn(async move |weak_lsp_store, cx| {
296 let Some(lsp_store) = weak_lsp_store.upgrade() else {
297 return Ok(None);
298 };
299 let Some(responses) = request_task.await? else {
300 return Ok(None);
301 };
302
303 let code_lens_actions = join_all(responses.payload.into_iter().map(|response| {
304 let lsp_store = lsp_store.clone();
305 let buffer = buffer.clone();
306 let cx = cx.clone();
307 async move {
308 (
309 LanguageServerId::from_proto(response.server_id),
310 GetCodeLens
311 .response_from_proto(response.response, lsp_store, buffer, cx)
312 .await,
313 )
314 }
315 }))
316 .await;
317
318 let mut has_errors = false;
319 let code_lens_actions = code_lens_actions
320 .into_iter()
321 .filter_map(|(server_id, code_lens)| match code_lens {
322 Ok(code_lens) => Some((server_id, code_lens)),
323 Err(e) => {
324 has_errors = true;
325 log::error!("{e:#}");
326 None
327 }
328 })
329 .collect::<HashMap<_, _>>();
330 anyhow::ensure!(
331 !has_errors || !code_lens_actions.is_empty(),
332 "Failed to fetch code lens"
333 );
334 Ok(Some(code_lens_actions))
335 })
336 } else {
337 let code_lens_actions_task = self.request_filtered_lsp_locally(
338 buffer,
339 None::<usize>,
340 GetCodeLens,
341 for_servers.as_ref(),
342 cx,
343 );
344 cx.background_spawn(async move {
345 Ok(Some(code_lens_actions_task.await.into_iter().collect()))
346 })
347 }
348 }
349
350 /// Resolves a single code lens via `codeLens/resolve`, identified by
351 /// the [`CodeLensActionId`] returned from [`Self::code_lens_actions`].
352 /// The returned task is shared and cached on [`CodeLensData::resolving`]
353 /// keyed by `(server, lens_id)`, so concurrent callers awaiting the
354 /// same lens only drive a single LSP request.
355 ///
356 /// `None` is yielded when the lens cannot be resolved (id no longer
357 /// cached, server gone, no `resolveProvider`, request failure, etc.).
358 /// On success, the cached entry is updated in place before the
359 /// `(id, resolved_action)` pair is returned.
360 ///
361 /// All visibility / batching policy lives in the caller. Remote (proto)
362 /// resolves are forwarded to the host via [`Self::resolve_code_action`].
363 pub fn resolve_code_lens(
364 &mut self,
365 buffer: &Entity<Buffer>,
366 server_id: LanguageServerId,
367 lens_id: CodeLensActionId,
368 cx: &mut Context<Self>,
369 ) -> CodeLensResolveTask {
370 let buffer_id = buffer.read(cx).remote_id();
371
372 let Some(code_lens) = self
373 .lsp_data
374 .get_mut(&buffer_id)
375 .and_then(|data| data.code_lens.as_mut())
376 else {
377 return Task::ready(None).shared();
378 };
379 let key = (server_id, lens_id);
380 if let Some(existing) = code_lens.resolving.get(&key) {
381 return existing.clone();
382 }
383 let Some(cached) = code_lens
384 .lens
385 .get(&server_id)
386 .and_then(|cache| cache.get(&lens_id))
387 else {
388 return Task::ready(None).shared();
389 };
390 if cached.resolved {
391 return Task::ready(Some((lens_id, cached.clone()))).shared();
392 }
393 let LspAction::CodeLens(lens) = &cached.lsp_action else {
394 return Task::ready(None).shared();
395 };
396 let lens = lens.clone();
397 let action = cached.clone();
398
399 if self.upstream_client().is_some() {
400 if !self.check_if_capable_for_proto_request(buffer, GetCodeLens::can_resolve_lens, cx) {
401 return Task::ready(None).shared();
402 }
403 let resolve = self.resolve_code_action(buffer, action, cx);
404 let task = cx
405 .spawn(async move |lsp_store, cx| {
406 let resolved = resolve
407 .await
408 .context("resolving remote code lens")
409 .log_err()?;
410 lsp_store
411 .update(cx, |lsp_store, _| {
412 let code_lens = lsp_store
413 .lsp_data
414 .get_mut(&buffer_id)
415 .and_then(|data| data.code_lens.as_mut())?;
416 code_lens.resolving.remove(&key);
417 let action = code_lens
418 .lens
419 .get_mut(&server_id)
420 .and_then(|cache| cache.get_mut(&lens_id))?;
421 action.resolved = true;
422 action.lsp_action = resolved.lsp_action;
423 Some((lens_id, action.clone()))
424 })
425 .ok()
426 .flatten()
427 })
428 .shared();
429 if let Some(code_lens) = self
430 .lsp_data
431 .get_mut(&buffer_id)
432 .and_then(|data| data.code_lens.as_mut())
433 {
434 code_lens.resolving.insert(key, task.clone());
435 }
436 return task;
437 }
438
439 let Some(server) = self.language_server_for_id(server_id) else {
440 return Task::ready(None).shared();
441 };
442 if !GetCodeLens::can_resolve_lens(&server.capabilities()) {
443 return Task::ready(None).shared();
444 }
445 let request_timeout = ProjectSettings::get_global(cx)
446 .global_lsp_settings
447 .get_request_timeout();
448
449 let task = cx
450 .spawn({
451 async move |lsp_store, cx| {
452 let response = server
453 .request::<lsp::request::CodeLensResolve>(lens, request_timeout)
454 .await
455 .into_response();
456 lsp_store
457 .update(cx, |lsp_store, _| {
458 let code_lens = lsp_store
459 .lsp_data
460 .get_mut(&buffer_id)
461 .and_then(|data| data.code_lens.as_mut())?;
462 code_lens.resolving.remove(&key);
463 let resolved_lens = match response {
464 Ok(resolved_lens) => resolved_lens,
465 Err(e) => {
466 log::warn!("Failed to resolve code lens: {e:#}");
467 return None;
468 }
469 };
470 let action = code_lens
471 .lens
472 .get_mut(&server_id)
473 .and_then(|cache| cache.get_mut(&lens_id))?;
474 action.resolved = true;
475 action.lsp_action = LspAction::CodeLens(resolved_lens);
476 Some((lens_id, action.clone()))
477 })
478 .ok()
479 .flatten()
480 }
481 })
482 .shared();
483
484 if let Some(code_lens) = self
485 .lsp_data
486 .get_mut(&buffer_id)
487 .and_then(|data| data.code_lens.as_mut())
488 {
489 code_lens.resolving.insert(key, task.clone());
490 }
491 task
492 }
493
494 #[cfg(any(test, feature = "test-support"))]
495 pub fn forget_code_lens_task(&mut self, buffer_id: text::BufferId) -> Option<CodeLensTask> {
496 Some(
497 self.lsp_data
498 .get_mut(&buffer_id)?
499 .code_lens
500 .take()?
501 .update
502 .take()?
503 .task,
504 )
505 }
506
507 pub(super) async fn handle_refresh_code_lens(
508 lsp_store: Entity<Self>,
509 envelope: TypedEnvelope<proto::RefreshCodeLens>,
510 mut cx: AsyncApp,
511 ) -> Result<proto::Ack> {
512 lsp_store.update(&mut cx, |lsp_store, cx| {
513 let server_id = envelope.payload.server_id.map(LanguageServerId::from_proto);
514 lsp_store.refresh_code_lens(server_id, cx);
515 });
516 Ok(proto::Ack {})
517 }
518}
519
520impl Project {
521 pub fn code_lens_actions(
522 &mut self,
523 buffer: &Entity<Buffer>,
524 range: Range<Anchor>,
525 cx: &mut Context<Self>,
526 ) -> Task<Result<Option<Vec<CodeAction>>>> {
527 let snapshot = buffer.read(cx).snapshot();
528 let range = range.to_point(&snapshot);
529 let range_start = snapshot.anchor_before(range.start);
530 let range_end = if range.start == range.end {
531 range_start
532 } else {
533 snapshot.anchor_after(range.end)
534 };
535 let range = range_start..range_end;
536 let lsp_store = self.lsp_store();
537 let fetch_task =
538 lsp_store.update(cx, |lsp_store, cx| lsp_store.code_lens_actions(buffer, cx));
539 let buffer = buffer.clone();
540 cx.spawn(async move |_, cx| {
541 let Some(mut tagged) = fetch_task.await? else {
542 return Ok(None);
543 };
544 let snapshot = buffer.read_with(cx, |buffer, _| buffer.snapshot());
545 tagged.retain(|_, action| {
546 range.start.cmp(&action.range.start, &snapshot).is_ge()
547 && range.end.cmp(&action.range.end, &snapshot).is_le()
548 });
549 let resolve_tasks = lsp_store.update(cx, |lsp_store, cx| {
550 tagged
551 .iter()
552 .filter(|(_, action)| !action.resolved)
553 .map(|(id, action)| {
554 lsp_store.resolve_code_lens(&buffer, action.server_id, *id, cx)
555 })
556 .collect::<Vec<_>>()
557 });
558 for (resolved_id, resolved) in join_all(resolve_tasks).await.into_iter().flatten() {
559 if let Some(slot) = tagged.get_mut(&resolved_id) {
560 *slot = resolved;
561 }
562 }
563 // Sort by id to recover server-emit order at the menu boundary.
564 let mut entries: Vec<_> = tagged.into_iter().collect();
565 entries.sort_by_key(|(id, _)| *id);
566 Ok(Some(entries.into_iter().map(|(_, a)| a).collect()))
567 })
568 }
569}
570