Skip to repository content312 lines · 12.5 KB · rust
tenant.openagents/omega
No repository description is available.
OpenAgents Git authority 2026-07-28T04:01:14.534Z 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
folding_ranges.rs
1use std::ops::Range;
2use std::sync::Arc;
3use std::time::Duration;
4
5use anyhow::Context as _;
6use collections::{HashMap, HashSet};
7use futures::FutureExt as _;
8use futures::future::{Shared, join_all};
9use gpui::{AppContext as _, AsyncApp, Context, Entity, SharedString, Task};
10use itertools::Itertools;
11use language::Buffer;
12use lsp::LanguageServerId;
13use rpc::{TypedEnvelope, proto};
14use settings::Settings as _;
15use text::Anchor;
16use util::ResultExt as _;
17
18use crate::lsp_command::{GetFoldingRanges, LspCommand as _};
19use crate::lsp_store::{
20 LspStore, LspStoreEvent, RunningFetch, missing_servers_to_query, next_lsp_fetch_id,
21 upstream_lsp_query_server_filter,
22};
23use crate::project_settings::ProjectSettings;
24
25#[derive(Clone, Debug)]
26pub struct LspFoldingRange {
27 pub range: Range<Anchor>,
28 pub collapsed_text: Option<SharedString>,
29}
30
31pub(super) type FoldingRangeTask =
32 Shared<Task<std::result::Result<Vec<LspFoldingRange>, Arc<anyhow::Error>>>>;
33
34#[derive(Debug, Default)]
35pub(super) struct FoldingRangeData {
36 pub(super) ranges: HashMap<LanguageServerId, Vec<LspFoldingRange>>,
37 fetched_servers: HashSet<LanguageServerId>,
38 ranges_update: Option<RunningFetch<FoldingRangeTask>>,
39}
40
41impl FoldingRangeData {
42 pub(super) fn remove_server_data(&mut self, server_id: LanguageServerId) {
43 self.ranges.remove(&server_id);
44 self.fetched_servers.remove(&server_id);
45 RunningFetch::discard_if_queried(&mut self.ranges_update, server_id);
46 }
47
48 fn evict(&mut self, for_server: Option<LanguageServerId>) {
49 match for_server {
50 Some(server_id) => self.remove_server_data(server_id),
51 None => {
52 self.ranges.clear();
53 self.fetched_servers.clear();
54 }
55 }
56 self.ranges_update = None;
57 }
58}
59
60impl LspStore {
61 pub(super) fn refresh_folding_ranges(
62 &mut self,
63 for_server: Option<LanguageServerId>,
64 cx: &mut Context<Self>,
65 ) {
66 for lsp_data in self.lsp_data.values_mut() {
67 if let Some(folding_ranges) = &mut lsp_data.folding_ranges {
68 folding_ranges.evict(for_server);
69 }
70 }
71
72 cx.emit(LspStoreEvent::RefreshFoldingRanges {
73 server_id: for_server,
74 });
75 if let Some((downstream_client, project_id)) = self.downstream_client.as_ref() {
76 downstream_client
77 .send(proto::RefreshFoldingRanges {
78 project_id: *project_id,
79 server_id: for_server.map(|server_id| server_id.to_proto()),
80 })
81 .context("sending refresh folding ranges downstream")
82 .log_err();
83 }
84 }
85
86 pub(super) async fn handle_refresh_folding_ranges(
87 lsp_store: Entity<Self>,
88 envelope: TypedEnvelope<proto::RefreshFoldingRanges>,
89 mut cx: AsyncApp,
90 ) -> anyhow::Result<proto::Ack> {
91 lsp_store.update(&mut cx, |lsp_store, cx| {
92 let server_id = envelope.payload.server_id.map(LanguageServerId::from_proto);
93 lsp_store.refresh_folding_ranges(server_id, cx);
94 });
95 Ok(proto::Ack {})
96 }
97
98 /// Returns a task that resolves to the folding ranges for the given buffer.
99 ///
100 /// Caches results per buffer version so repeated calls for the same version
101 /// return immediately. Deduplicates concurrent in-flight requests.
102 pub fn fetch_folding_ranges(
103 &mut self,
104 buffer: &Entity<Buffer>,
105 cx: &mut Context<Self>,
106 ) -> Task<Vec<LspFoldingRange>> {
107 let version_queried_for = buffer.read(cx).version();
108 let buffer_id = buffer.read(cx).remote_id();
109
110 let current_servers = self.relevant_server_ids_for_capability_check(buffer, cx);
111
112 let mut servers_to_query = None;
113 if let Some(lsp_data) = self.current_lsp_data(buffer_id) {
114 if !version_queried_for.changed_since(&lsp_data.buffer_version)
115 && let Some(cached) = &mut lsp_data.folding_ranges
116 {
117 match missing_servers_to_query(
118 &mut cached.ranges,
119 &mut cached.fetched_servers,
120 ¤t_servers,
121 ) {
122 Some(missing_servers) => servers_to_query = Some(missing_servers),
123 None => {
124 let snapshot = buffer.read(cx).snapshot();
125 return Task::ready(
126 cached
127 .ranges
128 .values()
129 .flatten()
130 .cloned()
131 .sorted_by(|a, b| a.range.start.cmp(&b.range.start, &snapshot))
132 .collect(),
133 );
134 }
135 }
136 }
137 if let Some(folding_ranges) = &lsp_data.folding_ranges
138 && let Some(running) = &folding_ranges.ranges_update
139 && !version_queried_for.changed_since(&running.version)
140 && servers_to_query
141 .as_ref()
142 .is_none_or(|missing| missing.is_subset(&running.servers))
143 {
144 let running = running.task.clone();
145 return cx.background_spawn(async move { running.await.unwrap_or_default() });
146 }
147 }
148
149 let folding_lsp_data = self
150 .latest_lsp_data(buffer, cx)
151 .folding_ranges
152 .get_or_insert_default();
153 let fetch_id = next_lsp_fetch_id();
154 let queried_servers = servers_to_query
155 .clone()
156 .unwrap_or_else(|| current_servers.clone());
157 let buffer = buffer.clone();
158 let query_version = version_queried_for.clone();
159 let new_task = cx
160 .spawn({
161 let queried_servers = queried_servers.clone();
162 async move |lsp_store, cx| {
163 cx.background_executor()
164 .timer(Duration::from_millis(30))
165 .await;
166
167 let fetched = lsp_store
168 .update(cx, |lsp_store, cx| {
169 lsp_store.fetch_folding_ranges_for_buffer(&buffer, servers_to_query, cx)
170 })
171 .map_err(Arc::new)?
172 .await
173 .context("fetching folding ranges")
174 .map_err(Arc::new);
175
176 let fetched = match fetched {
177 Ok(fetched) => fetched,
178 Err(e) => {
179 lsp_store
180 .update(cx, |lsp_store, _| {
181 if let Some(lsp_data) = lsp_store.lsp_data.get_mut(&buffer_id)
182 && let Some(folding_ranges) = &mut lsp_data.folding_ranges
183 {
184 RunningFetch::take_finished(
185 &mut folding_ranges.ranges_update,
186 fetch_id,
187 );
188 }
189 })
190 .ok();
191 return Err(e);
192 }
193 };
194
195 lsp_store
196 .update(cx, |lsp_store, cx| {
197 let lsp_data = lsp_store.latest_lsp_data(&buffer, cx);
198 let folding = lsp_data.folding_ranges.get_or_insert_default();
199
200 if RunningFetch::take_finished(&mut folding.ranges_update, fetch_id)
201 && let Some(fetched_ranges) = fetched
202 {
203 if lsp_data.buffer_version == query_version {
204 folding.ranges.extend(fetched_ranges);
205 folding.fetched_servers.extend(queried_servers);
206 } else if !lsp_data.buffer_version.changed_since(&query_version) {
207 lsp_data.buffer_version = query_version;
208 folding.ranges = fetched_ranges;
209 folding.fetched_servers = queried_servers;
210 }
211 }
212 let snapshot = buffer.read(cx).snapshot();
213 folding
214 .ranges
215 .values()
216 .flatten()
217 .cloned()
218 .sorted_by(|a, b| a.range.start.cmp(&b.range.start, &snapshot))
219 .collect()
220 })
221 .map_err(Arc::new)
222 }
223 })
224 .shared();
225
226 folding_lsp_data.ranges_update = Some(RunningFetch {
227 id: fetch_id,
228 version: version_queried_for,
229 servers: queried_servers,
230 task: new_task.clone(),
231 });
232
233 cx.background_spawn(async move { new_task.await.unwrap_or_default() })
234 }
235
236 fn fetch_folding_ranges_for_buffer(
237 &mut self,
238 buffer: &Entity<Buffer>,
239 for_servers: Option<HashSet<LanguageServerId>>,
240 cx: &mut Context<Self>,
241 ) -> Task<anyhow::Result<Option<HashMap<LanguageServerId, Vec<LspFoldingRange>>>>> {
242 if let Some((client, project_id)) = self.upstream_client() {
243 let request = GetFoldingRanges;
244 if !self.is_capable_for_proto_request(buffer, &request, cx) {
245 return Task::ready(Ok(None));
246 }
247
248 let request_timeout = ProjectSettings::get_global(cx)
249 .global_lsp_settings
250 .get_request_timeout();
251 let request_task = client.request_lsp(
252 project_id,
253 upstream_lsp_query_server_filter(for_servers.as_ref()),
254 request_timeout,
255 cx.background_executor().clone(),
256 request.to_proto(project_id, buffer.read(cx)),
257 );
258 let buffer = buffer.clone();
259 cx.spawn(async move |weak_lsp_store, cx| {
260 let Some(lsp_store) = weak_lsp_store.upgrade() else {
261 return Ok(None);
262 };
263 let Some(responses) = request_task.await? else {
264 return Ok(None);
265 };
266
267 let folding_ranges = join_all(responses.payload.into_iter().map(|response| {
268 let lsp_store = lsp_store.clone();
269 let buffer = buffer.clone();
270 let cx = cx.clone();
271 async move {
272 (
273 LanguageServerId::from_proto(response.server_id),
274 GetFoldingRanges
275 .response_from_proto(response.response, lsp_store, buffer, cx)
276 .await,
277 )
278 }
279 }))
280 .await;
281
282 let mut has_errors = false;
283 let result = folding_ranges
284 .into_iter()
285 .filter_map(|(server_id, ranges)| match ranges {
286 Ok(ranges) => Some((server_id, ranges)),
287 Err(e) => {
288 has_errors = true;
289 log::error!("Failed to fetch folding ranges: {e:#}");
290 None
291 }
292 })
293 .collect::<HashMap<_, _>>();
294 anyhow::ensure!(
295 !has_errors || !result.is_empty(),
296 "Failed to fetch folding ranges"
297 );
298 Ok(Some(result))
299 })
300 } else {
301 let folding_task = self.request_filtered_lsp_locally(
302 buffer,
303 None::<usize>,
304 GetFoldingRanges,
305 for_servers.as_ref(),
306 cx,
307 );
308 cx.background_spawn(async move { Ok(Some(folding_task.await.into_iter().collect())) })
309 }
310 }
311}
312