Skip to repository content

tenant.openagents/omega

No repository description is available.

OpenAgents Git authority 2026-07-28T03:36:13.371Z 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

headless_host.rs

353 lines · 12.0 KB · rust
1use std::{path::PathBuf, sync::Arc};
2
3use anyhow::{Context as _, Result};
4use client::{TypedEnvelope, proto};
5use collections::{HashMap, HashSet};
6use extension::{
7    Extension, ExtensionDebugAdapterProviderProxy, ExtensionHostProxy, ExtensionLanguageProxy,
8    ExtensionLanguageServerProxy, ExtensionManifest,
9};
10use fs::{Fs, RemoveOptions, RenameOptions};
11use futures::future::join_all;
12use gpui::{App, AppContext as _, AsyncApp, Context, Entity, Task, WeakEntity};
13use http_client::HttpClient;
14use language::{LanguageConfig, LanguageName, LanguageQueries, LoadedLanguage};
15use lsp::LanguageServerName;
16use node_runtime::NodeRuntime;
17
18use crate::wasm_host::{WasmExtension, WasmHost};
19
20#[derive(Clone, Debug)]
21pub struct ExtensionVersion {
22    pub id: String,
23    pub version: String,
24    pub dev: bool,
25}
26
27pub struct HeadlessExtensionStore {
28    pub fs: Arc<dyn Fs>,
29    pub extension_dir: PathBuf,
30    pub proxy: Arc<ExtensionHostProxy>,
31    pub wasm_host: Arc<WasmHost>,
32    pub loaded_extensions: HashMap<Arc<str>, Arc<str>>,
33    pub loaded_languages: HashMap<Arc<str>, Vec<LanguageName>>,
34    pub loaded_language_servers: HashMap<Arc<str>, Vec<(LanguageServerName, LanguageName)>>,
35}
36
37impl HeadlessExtensionStore {
38    pub fn new(
39        fs: Arc<dyn Fs>,
40        http_client: Arc<dyn HttpClient>,
41        extension_dir: PathBuf,
42        extension_host_proxy: Arc<ExtensionHostProxy>,
43        node_runtime: NodeRuntime,
44        cx: &mut App,
45    ) -> Entity<Self> {
46        cx.new(|cx| Self {
47            fs: fs.clone(),
48            wasm_host: WasmHost::new(
49                fs.clone(),
50                http_client.clone(),
51                node_runtime,
52                extension_host_proxy.clone(),
53                extension_dir.join("work"),
54                cx,
55            ),
56            extension_dir,
57            proxy: extension_host_proxy,
58            loaded_extensions: Default::default(),
59            loaded_languages: Default::default(),
60            loaded_language_servers: Default::default(),
61        })
62    }
63
64    pub fn sync_extensions(
65        &mut self,
66        extensions: Vec<ExtensionVersion>,
67        cx: &Context<Self>,
68    ) -> Task<Result<Vec<ExtensionVersion>>> {
69        let on_client = HashSet::from_iter(extensions.iter().map(|e| e.id.as_str()));
70        let to_remove: Vec<Arc<str>> = self
71            .loaded_extensions
72            .keys()
73            .filter(|id| !on_client.contains(id.as_ref()))
74            .cloned()
75            .collect();
76        let to_load: Vec<ExtensionVersion> = extensions
77            .into_iter()
78            .filter(|e| {
79                if e.dev {
80                    return true;
81                }
82                self.loaded_extensions
83                    .get(e.id.as_str())
84                    .is_none_or(|loaded| loaded.as_ref() != e.version.as_str())
85            })
86            .collect();
87
88        cx.spawn(async move |this, cx| {
89            let mut missing = Vec::new();
90
91            for extension_id in to_remove {
92                log::info!("removing extension: {}", extension_id);
93                this.update(cx, |this, cx| this.uninstall_extension(&extension_id, cx))?
94                    .await?;
95            }
96
97            for extension in to_load {
98                if let Err(e) = Self::load_extension(this.clone(), extension.clone(), cx).await {
99                    log::info!("failed to load extension: {}, {:#}", extension.id, e);
100                    missing.push(extension)
101                } else if extension.dev {
102                    missing.push(extension)
103                }
104            }
105
106            Ok(missing)
107        })
108    }
109
110    pub async fn load_extension(
111        this: WeakEntity<Self>,
112        extension: ExtensionVersion,
113        cx: &mut AsyncApp,
114    ) -> Result<()> {
115        let (fs, wasm_host, extension_dir) = this.update(cx, |this, _cx| {
116            this.loaded_extensions.insert(
117                extension.id.clone().into(),
118                extension.version.clone().into(),
119            );
120            (
121                this.fs.clone(),
122                this.wasm_host.clone(),
123                this.extension_dir.join(&extension.id),
124            )
125        })?;
126
127        let manifest = Arc::new(ExtensionManifest::load(fs.clone(), &extension_dir).await?);
128
129        debug_assert!(!manifest.languages.is_empty() || manifest.allow_remote_load());
130
131        if manifest.version.as_ref() != extension.version.as_str() {
132            anyhow::bail!(
133                "mismatched versions: ({}) != ({})",
134                manifest.version,
135                extension.version
136            )
137        }
138
139        for language_path in &manifest.languages {
140            let language_path = extension_dir.join(language_path);
141            let config = fs
142                .load(&language_path.join(LanguageConfig::FILE_NAME))
143                .await?;
144            let mut config = ::toml::from_str::<LanguageConfig>(&config)?;
145
146            this.update(cx, |this, _cx| {
147                this.loaded_languages
148                    .entry(manifest.id.clone())
149                    .or_default()
150                    .push(config.name.clone());
151
152                config.grammar = None;
153
154                this.proxy.register_language(
155                    config.name.clone(),
156                    None,
157                    config.matcher.clone(),
158                    config.hidden,
159                    Arc::new(move || {
160                        Ok(LoadedLanguage {
161                            config: config.clone(),
162                            queries: LanguageQueries::default(),
163                            context_provider: None,
164                            toolchain_provider: None,
165                            manifest_name: None,
166                        })
167                    }),
168                );
169            })?;
170        }
171
172        if !manifest.allow_remote_load() {
173            return Ok(());
174        }
175
176        let wasm_extension: Arc<dyn Extension> =
177            Arc::new(WasmExtension::load(&extension_dir, &manifest, wasm_host.clone(), cx).await?);
178
179        for (language_server_id, language_server_config) in &manifest.language_servers {
180            for language in language_server_config.languages() {
181                this.update(cx, |this, _cx| {
182                    this.loaded_language_servers
183                        .entry(manifest.id.clone())
184                        .or_default()
185                        .push((language_server_id.clone(), language.clone()));
186                    this.proxy.register_language_server(
187                        wasm_extension.clone(),
188                        language_server_id.clone(),
189                        language.clone(),
190                    );
191                })?;
192            }
193            log::info!("Loaded language server: {}", language_server_id);
194        }
195
196        for (debug_adapter, meta) in &manifest.debug_adapters {
197            let schema_path = extension::build_debug_adapter_schema_path(debug_adapter, meta)?;
198
199            this.update(cx, |this, _cx| {
200                this.proxy.register_debug_adapter(
201                    wasm_extension.clone(),
202                    debug_adapter.clone(),
203                    &extension_dir.join(schema_path),
204                );
205            })?;
206            log::info!("Loaded debug adapter: {}", debug_adapter);
207        }
208
209        for debug_locator in manifest.debug_locators.keys() {
210            this.update(cx, |this, _cx| {
211                this.proxy
212                    .register_debug_locator(wasm_extension.clone(), debug_locator.clone());
213            })?;
214            log::info!("Loaded debug locator: {}", debug_locator);
215        }
216
217        Ok(())
218    }
219
220    fn uninstall_extension(
221        &mut self,
222        extension_id: &Arc<str>,
223        cx: &mut Context<Self>,
224    ) -> Task<Result<()>> {
225        self.loaded_extensions.remove(extension_id);
226
227        let languages_to_remove = self
228            .loaded_languages
229            .remove(extension_id)
230            .unwrap_or_default();
231        self.proxy.remove_languages(&languages_to_remove, &[]);
232
233        let servers_to_remove = self
234            .loaded_language_servers
235            .remove(extension_id)
236            .unwrap_or_default();
237        let proxy = self.proxy.clone();
238        let path = self.extension_dir.join(&extension_id.to_string());
239        let fs = self.fs.clone();
240        cx.spawn(async move |_, cx| {
241            let mut removal_tasks = Vec::with_capacity(servers_to_remove.len());
242            cx.update(|cx| {
243                for (language_server_name, language) in servers_to_remove {
244                    removal_tasks.push(proxy.remove_language_server(
245                        &language,
246                        &language_server_name,
247                        cx,
248                    ));
249                }
250            });
251            let _ = join_all(removal_tasks).await;
252
253            fs.remove_dir(
254                &path,
255                RemoveOptions {
256                    recursive: true,
257                    ignore_if_not_exists: true,
258                },
259            )
260            .await
261            .with_context(|| format!("Removing directory {path:?}"))
262        })
263    }
264
265    pub fn install_extension(
266        &mut self,
267        extension: ExtensionVersion,
268        tmp_path: PathBuf,
269        cx: &mut Context<Self>,
270    ) -> Task<Result<()>> {
271        let path = self.extension_dir.join(&extension.id);
272        let fs = self.fs.clone();
273
274        cx.spawn(async move |this, cx| {
275            if fs.is_dir(&path).await {
276                this.update(cx, |this, cx| {
277                    this.uninstall_extension(&extension.id.clone().into(), cx)
278                })?
279                .await?;
280            }
281
282            fs.rename(&tmp_path, &path, RenameOptions::default())
283                .await
284                .with_context(|| format!("Failed to rename {tmp_path:?} to {path:?}"))?;
285
286            Self::load_extension(this, extension, cx).await
287        })
288    }
289
290    pub async fn handle_sync_extensions(
291        extension_store: Entity<HeadlessExtensionStore>,
292        envelope: TypedEnvelope<proto::SyncExtensions>,
293        mut cx: AsyncApp,
294    ) -> Result<proto::SyncExtensionsResponse> {
295        let requested_extensions =
296            envelope
297                .payload
298                .extensions
299                .into_iter()
300                .map(|p| ExtensionVersion {
301                    id: p.id,
302                    version: p.version,
303                    dev: p.dev,
304                });
305        let missing_extensions = extension_store
306            .update(&mut cx, |extension_store, cx| {
307                extension_store.sync_extensions(requested_extensions.collect(), cx)
308            })
309            .await?;
310
311        Ok(proto::SyncExtensionsResponse {
312            missing_extensions: missing_extensions
313                .into_iter()
314                .map(|e| proto::Extension {
315                    id: e.id,
316                    version: e.version,
317                    dev: e.dev,
318                })
319                .collect(),
320            tmp_dir: paths::remote_extensions_uploads_dir()
321                .to_string_lossy()
322                .to_string(),
323        })
324    }
325
326    pub async fn handle_install_extension(
327        extensions: Entity<HeadlessExtensionStore>,
328        envelope: TypedEnvelope<proto::InstallExtension>,
329        mut cx: AsyncApp,
330    ) -> Result<proto::Ack> {
331        let extension = envelope
332            .payload
333            .extension
334            .context("Invalid InstallExtension request")?;
335
336        extensions
337            .update(&mut cx, |extensions, cx| {
338                extensions.install_extension(
339                    ExtensionVersion {
340                        id: extension.id,
341                        version: extension.version,
342                        dev: extension.dev,
343                    },
344                    PathBuf::from(envelope.payload.tmp_dir),
345                    cx,
346                )
347            })
348            .await?;
349
350        Ok(proto::Ack {})
351    }
352}
353
Served at tenant.openagents/omega Member data and write actions are omitted.