Skip to repository content

tenant.openagents/omega

No repository description is available.

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

tasks.rs

590 lines · 21.3 KB · rust
1use std::process::ExitStatus;
2
3use anyhow::Result;
4use collections::HashSet;
5use gpui::{AppContext, AsyncWindowContext, Context, Entity, Task, TaskExt, WeakEntity};
6use language::Buffer;
7use project::{TaskSourceKind, WorktreeId};
8use remote::ConnectionState;
9use task::{
10    DebugScenario, ResolvedTask, SaveStrategy, SharedTaskContext, SpawnInTerminal, TaskContext,
11    TaskHook, TaskTemplate, TaskVariables, VariableName,
12};
13use ui::Window;
14use util::TryFutureExt;
15
16use crate::{SaveIntent, Toast, Workspace, notifications::NotificationId};
17
18#[derive(Clone, Copy, Debug, PartialEq, Eq)]
19pub enum ScheduledTaskResult {
20    Success,
21    Failure,
22    SpawnFailed,
23    Cancelled,
24}
25
26type TaskCompletionHandler = Box<dyn FnOnce(ScheduledTaskResult, &mut AsyncWindowContext)>;
27
28impl Workspace {
29    pub fn schedule_task(
30        self: &mut Workspace,
31        task_source_kind: TaskSourceKind,
32        task_to_resolve: &TaskTemplate,
33        task_cx: &TaskContext,
34        omit_history: bool,
35        window: &mut Window,
36        cx: &mut Context<Self>,
37    ) {
38        match self.project.read(cx).remote_connection_state(cx) {
39            None | Some(ConnectionState::Connected) => {}
40            Some(
41                ConnectionState::Connecting
42                | ConnectionState::Disconnected
43                | ConnectionState::HeartbeatMissed
44                | ConnectionState::Reconnecting,
45            ) => {
46                log::warn!("Cannot schedule tasks when disconnected from a remote host");
47                return;
48            }
49        }
50
51        if let Some(spawn_in_terminal) =
52            task_to_resolve.resolve_task(&task_source_kind.to_id_base(), task_cx)
53        {
54            self.schedule_resolved_task(
55                task_source_kind,
56                spawn_in_terminal,
57                omit_history,
58                window,
59                cx,
60            );
61        }
62    }
63
64    pub fn schedule_resolved_task(
65        self: &mut Workspace,
66        task_source_kind: TaskSourceKind,
67        resolved_task: ResolvedTask,
68        omit_history: bool,
69        window: &mut Window,
70        cx: &mut Context<Workspace>,
71    ) {
72        self.schedule_resolved_task_internal(
73            task_source_kind,
74            resolved_task,
75            omit_history,
76            None,
77            window,
78            cx,
79        );
80    }
81
82    pub fn schedule_resolved_task_with_completion(
83        self: &mut Workspace,
84        task_source_kind: TaskSourceKind,
85        resolved_task: ResolvedTask,
86        omit_history: bool,
87        on_complete: impl FnOnce(ScheduledTaskResult, &mut AsyncWindowContext) + 'static,
88        window: &mut Window,
89        cx: &mut Context<Workspace>,
90    ) {
91        self.schedule_resolved_task_internal(
92            task_source_kind,
93            resolved_task,
94            omit_history,
95            Some(Box::new(on_complete)),
96            window,
97            cx,
98        );
99    }
100
101    fn schedule_resolved_task_internal(
102        self: &mut Workspace,
103        task_source_kind: TaskSourceKind,
104        resolved_task: ResolvedTask,
105        omit_history: bool,
106        on_complete: Option<TaskCompletionHandler>,
107        window: &mut Window,
108        cx: &mut Context<Workspace>,
109    ) {
110        let spawn_in_terminal = resolved_task.resolved.clone();
111        if !omit_history {
112            if let Some(debugger_provider) = self.debugger_provider.as_ref() {
113                debugger_provider.task_scheduled(cx);
114            }
115
116            self.project().update(cx, |project, cx| {
117                if let Some(task_inventory) =
118                    project.task_store().read(cx).task_inventory().cloned()
119                {
120                    task_inventory.update(cx, |inventory, _| {
121                        inventory.task_scheduled(task_source_kind, resolved_task);
122                    })
123                }
124            });
125        }
126
127        if self.terminal_provider.is_some() {
128            let task = cx.spawn_in(window, async move |workspace, cx| {
129                Self::save_for_task(&workspace, spawn_in_terminal.save, cx).await;
130
131                let spawn_task = workspace.update_in(cx, |workspace, window, cx| {
132                    workspace
133                        .terminal_provider
134                        .as_ref()
135                        .map(|terminal_provider| {
136                            terminal_provider.spawn(spawn_in_terminal, window, cx)
137                        })
138                });
139                if let Some(spawn_task) = spawn_task.ok().flatten() {
140                    let res = cx.background_spawn(spawn_task).await;
141                    let result = match res {
142                        Some(Ok(status)) => {
143                            if status.success() {
144                                log::debug!("Task spawn succeeded");
145                                ScheduledTaskResult::Success
146                            } else {
147                                log::debug!("Task spawn failed, code: {:?}", status.code());
148                                ScheduledTaskResult::Failure
149                            }
150                        }
151                        Some(Err(e)) => {
152                            log::error!("Task spawn failed: {e:#}");
153                            _ = workspace.update(cx, |w, cx| {
154                                let id = NotificationId::unique::<ResolvedTask>();
155                                w.show_toast(Toast::new(id, format!("Task spawn failed: {e}")), cx);
156                            });
157                            ScheduledTaskResult::SpawnFailed
158                        }
159                        None => {
160                            log::debug!("Task spawn got cancelled");
161                            ScheduledTaskResult::Cancelled
162                        }
163                    };
164                    if let Some(on_complete) = on_complete {
165                        on_complete(result, cx);
166                    }
167                } else if let Some(on_complete) = on_complete {
168                    on_complete(ScheduledTaskResult::Cancelled, cx);
169                }
170            });
171            self.scheduled_tasks.push(task);
172        }
173    }
174
175    pub async fn save_for_task(
176        workspace: &WeakEntity<Self>,
177        save_strategy: SaveStrategy,
178        cx: &mut AsyncWindowContext,
179    ) {
180        let save_action = match save_strategy {
181            SaveStrategy::All => {
182                let save_all = workspace.update_in(cx, |workspace, window, cx| {
183                    let task = workspace.save_all_internal(SaveIntent::SaveAll, true, window, cx);
184                    cx.background_spawn(async { task.await.map(|_| ()) })
185                });
186                save_all.ok()
187            }
188            SaveStrategy::Current => {
189                let save_current = workspace.update_in(cx, |workspace, window, cx| {
190                    workspace.save_active_item(SaveIntent::SaveAll, window, cx)
191                });
192                save_current.ok()
193            }
194            SaveStrategy::None => None,
195        };
196        if let Some(save_action) = save_action {
197            save_action.log_err().await;
198        }
199    }
200
201    pub fn start_debug_session(
202        &mut self,
203        scenario: DebugScenario,
204        task_context: SharedTaskContext,
205        active_buffer: Option<Entity<Buffer>>,
206        worktree_id: Option<WorktreeId>,
207        window: &mut Window,
208        cx: &mut Context<Self>,
209    ) {
210        if let Some(provider) = self.debugger_provider.as_mut() {
211            provider.start_session(
212                scenario,
213                task_context,
214                active_buffer,
215                worktree_id,
216                window,
217                cx,
218            )
219        }
220    }
221
222    pub fn spawn_in_terminal(
223        self: &mut Workspace,
224        spawn_in_terminal: SpawnInTerminal,
225        window: &mut Window,
226        cx: &mut Context<Workspace>,
227    ) -> Task<Option<Result<ExitStatus>>> {
228        if let Some(terminal_provider) = self.terminal_provider.as_ref() {
229            terminal_provider.spawn(spawn_in_terminal, window, cx)
230        } else {
231            Task::ready(None)
232        }
233    }
234
235    pub fn run_create_worktree_tasks(&mut self, window: &mut Window, cx: &mut Context<Self>) {
236        let project = self.project().clone();
237        let hooks = HashSet::from_iter([TaskHook::CreateWorktree]);
238
239        let worktree_tasks: Vec<(WorktreeId, TaskContext, Vec<TaskTemplate>)> = {
240            let project = project.read(cx);
241            let task_store = project.task_store();
242            let Some(inventory) = task_store.read(cx).task_inventory().cloned() else {
243                return;
244            };
245
246            let git_store = project.git_store().read(cx);
247
248            let mut worktree_tasks = Vec::new();
249            for worktree in project.worktrees(cx) {
250                let worktree = worktree.read(cx);
251                let worktree_id = worktree.id();
252                let worktree_abs_path = worktree.abs_path();
253
254                let templates: Vec<TaskTemplate> = inventory
255                    .read(cx)
256                    .templates_with_hooks(&hooks, worktree_id)
257                    .into_iter()
258                    .map(|(_, template)| template)
259                    .collect();
260
261                if templates.is_empty() {
262                    continue;
263                }
264
265                let mut task_variables = TaskVariables::default();
266                task_variables.insert(
267                    VariableName::WorktreeRoot,
268                    worktree_abs_path.to_string_lossy().into_owned(),
269                );
270
271                if let Some(path) = git_store.original_repo_path_for_worktree(worktree_id, cx) {
272                    task_variables.insert(
273                        VariableName::MainGitWorktree,
274                        path.to_string_lossy().into_owned(),
275                    );
276                }
277
278                let task_context = TaskContext {
279                    cwd: Some(worktree_abs_path.to_path_buf()),
280                    task_variables,
281                    project_env: Default::default(),
282                };
283
284                worktree_tasks.push((worktree_id, task_context, templates));
285            }
286            worktree_tasks
287        };
288
289        if worktree_tasks.is_empty() {
290            return;
291        }
292
293        let task = cx.spawn_in(window, async move |workspace, cx| {
294            let mut tasks = Vec::new();
295            for (worktree_id, task_context, templates) in worktree_tasks {
296                let id_base = format!("worktree_setup_{worktree_id}");
297
298                tasks.push(cx.spawn({
299                    let workspace = workspace.clone();
300                    async move |cx| {
301                        for task_template in templates {
302                            let Some(resolved) =
303                                task_template.resolve_task(&id_base, &task_context)
304                            else {
305                                continue;
306                            };
307
308                            let status = workspace.update_in(cx, |workspace, window, cx| {
309                                workspace.spawn_in_terminal(resolved.resolved, window, cx)
310                            })?;
311
312                            if let Some(result) = status.await {
313                                match result {
314                                    Ok(exit_status) if !exit_status.success() => {
315                                        log::error!(
316                                            "Git worktree setup task failed with status: {:?}",
317                                            exit_status.code()
318                                        );
319                                        break;
320                                    }
321                                    Err(error) => {
322                                        log::error!("Git worktree setup task error: {error:#}");
323                                        break;
324                                    }
325                                    _ => {}
326                                }
327                            }
328                        }
329                        anyhow::Ok(())
330                    }
331                }));
332            }
333
334            futures::future::join_all(tasks).await;
335            anyhow::Ok(())
336        });
337        task.detach_and_log_err(cx);
338    }
339}
340
341#[cfg(test)]
342mod tests {
343    use super::*;
344    use crate::{
345        TerminalProvider,
346        item::test::{TestItem, TestProjectItem},
347        register_serializable_item,
348    };
349    use gpui::{App, TestAppContext};
350    use parking_lot::Mutex;
351    use project::{FakeFs, Project, TaskSourceKind};
352    use serde_json::json;
353    use std::sync::Arc;
354    use task::TaskTemplate;
355
356    struct Fixture {
357        workspace: Entity<Workspace>,
358        item: Entity<TestItem>,
359        task: ResolvedTask,
360        dirty_before_spawn: Arc<Mutex<Option<bool>>>,
361    }
362
363    #[gpui::test]
364    async fn test_schedule_resolved_task_save_all(cx: &mut TestAppContext) {
365        let (fixture, cx) = create_fixture(cx, SaveStrategy::All).await;
366        fixture.workspace.update_in(cx, |workspace, window, cx| {
367            workspace.schedule_resolved_task(
368                TaskSourceKind::UserInput,
369                fixture.task,
370                false,
371                window,
372                cx,
373            );
374        });
375        cx.executor().run_until_parked();
376
377        assert_eq!(*fixture.dirty_before_spawn.lock(), Some(false));
378        assert!(cx.read(|cx| !fixture.item.read(cx).is_dirty));
379    }
380
381    #[gpui::test]
382    async fn test_schedule_resolved_task_save_current(cx: &mut TestAppContext) {
383        let (fixture, cx) = create_fixture(cx, SaveStrategy::Current).await;
384        // Add a second inactive dirty item
385        let inactive = add_test_item(&fixture.workspace, "file2.txt", false, cx);
386        fixture.workspace.update_in(cx, |workspace, window, cx| {
387            workspace.schedule_resolved_task(
388                TaskSourceKind::UserInput,
389                fixture.task,
390                false,
391                window,
392                cx,
393            );
394        });
395        cx.executor().run_until_parked();
396
397        // The active item (fixture.item) should be saved
398        assert_eq!(*fixture.dirty_before_spawn.lock(), Some(false));
399        assert!(cx.read(|cx| !fixture.item.read(cx).is_dirty));
400        // The inactive item should not be saved
401        assert!(cx.read(|cx| inactive.read(cx).is_dirty));
402    }
403
404    #[gpui::test]
405    async fn test_schedule_resolved_task_save_none(cx: &mut TestAppContext) {
406        let (fixture, cx) = create_fixture(cx, SaveStrategy::None).await;
407        fixture.workspace.update_in(cx, |workspace, window, cx| {
408            workspace.schedule_resolved_task(
409                TaskSourceKind::UserInput,
410                fixture.task,
411                false,
412                window,
413                cx,
414            );
415        });
416        cx.executor().run_until_parked();
417
418        assert_eq!(*fixture.dirty_before_spawn.lock(), Some(true));
419        assert!(cx.read(|cx| fixture.item.read(cx).is_dirty));
420    }
421
422    #[gpui::test]
423    async fn test_schedule_resolved_task_with_completion_reports_success(cx: &mut TestAppContext) {
424        let (fixture, cx) = create_fixture(cx, SaveStrategy::None).await;
425        let task_result = Arc::new(Mutex::new(None));
426        fixture.workspace.update_in(cx, |workspace, window, cx| {
427            workspace.schedule_resolved_task_with_completion(
428                TaskSourceKind::UserInput,
429                fixture.task,
430                false,
431                {
432                    let task_result = task_result.clone();
433                    move |result, _| {
434                        *task_result.lock() = Some(result);
435                    }
436                },
437                window,
438                cx,
439            );
440        });
441        cx.executor().run_until_parked();
442
443        assert_eq!(*task_result.lock(), Some(ScheduledTaskResult::Success));
444    }
445
446    async fn create_fixture(
447        cx: &mut TestAppContext,
448        save_strategy: SaveStrategy,
449    ) -> (Fixture, &mut gpui::VisualTestContext) {
450        cx.update(|cx| {
451            let settings_store = settings::SettingsStore::test(cx);
452            cx.set_global(settings_store);
453            theme_settings::init(theme::LoadThemes::JustBase, cx);
454            register_serializable_item::<TestItem>(cx);
455        });
456        let fs = FakeFs::new(cx.executor());
457        fs.insert_tree("/root", json!({ "file.txt": "dirty" }))
458            .await;
459        let project = Project::test(fs.clone(), ["/root".as_ref()], cx).await;
460        let (workspace, cx) =
461            cx.add_window_view(|window, cx| Workspace::test_new(project.clone(), window, cx));
462
463        // Add a dirty item to the workspace
464        let item = add_test_item(&workspace, "file.txt", true, cx);
465
466        let template = TaskTemplate {
467            label: "test".to_string(),
468            command: "echo".to_string(),
469            save: save_strategy,
470            ..Default::default()
471        };
472        let task = template
473            .resolve_task("test", &task::TaskContext::default())
474            .unwrap();
475        let dirty_before_spawn: Arc<Mutex<Option<bool>>> = Arc::default();
476        let terminal_provider = Box::new(TestTerminalProvider {
477            item: item.clone(),
478            dirty_before_spawn: dirty_before_spawn.clone(),
479        });
480        workspace.update(cx, |workspace, _| {
481            workspace.terminal_provider = Some(terminal_provider);
482        });
483        let fixture = Fixture {
484            workspace,
485            item,
486            task,
487            dirty_before_spawn,
488        };
489        (fixture, cx)
490    }
491
492    fn add_test_item(
493        workspace: &Entity<Workspace>,
494        name: &str,
495        active: bool,
496        cx: &mut gpui::VisualTestContext,
497    ) -> Entity<TestItem> {
498        let item = cx.new(|cx| {
499            TestItem::new(cx)
500                .with_dirty(true)
501                .with_project_items(&[TestProjectItem::new(1, name, cx)])
502        });
503        workspace.update_in(cx, |workspace, window, cx| {
504            let pane = workspace.active_pane().clone();
505            workspace.add_item(pane, Box::new(item.clone()), None, true, active, window, cx);
506        });
507        item
508    }
509
510    #[gpui::test]
511    async fn test_save_for_task_all(cx: &mut TestAppContext) {
512        let (fixture, cx) = create_fixture(cx, SaveStrategy::All).await;
513        let workspace = fixture.workspace.downgrade();
514        cx.run_until_parked();
515
516        assert!(cx.read(|cx| fixture.item.read(cx).is_dirty));
517        fixture.workspace.update_in(cx, |_workspace, window, cx| {
518            cx.spawn_in(window, {
519                let workspace = workspace.clone();
520                async move |_this, cx| {
521                    Workspace::save_for_task(&workspace, SaveStrategy::All, cx).await;
522                }
523            })
524            .detach();
525        });
526        cx.run_until_parked();
527        assert!(cx.read(|cx| !fixture.item.read(cx).is_dirty));
528    }
529
530    #[gpui::test]
531    async fn test_save_for_task_none(cx: &mut TestAppContext) {
532        let (fixture, cx) = create_fixture(cx, SaveStrategy::None).await;
533        let workspace = fixture.workspace.downgrade();
534        cx.run_until_parked();
535
536        assert!(cx.read(|cx| fixture.item.read(cx).is_dirty));
537        fixture.workspace.update_in(cx, |_workspace, window, cx| {
538            cx.spawn_in(window, {
539                let workspace = workspace.clone();
540                async move |_this, cx| {
541                    Workspace::save_for_task(&workspace, SaveStrategy::None, cx).await;
542                }
543            })
544            .detach();
545        });
546        cx.run_until_parked();
547        assert!(cx.read(|cx| fixture.item.read(cx).is_dirty));
548    }
549
550    #[gpui::test]
551    async fn test_save_for_task_current(cx: &mut TestAppContext) {
552        let (fixture, cx) = create_fixture(cx, SaveStrategy::Current).await;
553        let inactive = add_test_item(&fixture.workspace, "file2.txt", false, cx);
554        let workspace = fixture.workspace.downgrade();
555        cx.run_until_parked();
556
557        assert!(cx.read(|cx| fixture.item.read(cx).is_dirty));
558        assert!(cx.read(|cx| inactive.read(cx).is_dirty));
559        fixture.workspace.update_in(cx, |_workspace, window, cx| {
560            cx.spawn_in(window, {
561                let workspace = workspace.clone();
562                async move |_this, cx| {
563                    Workspace::save_for_task(&workspace, SaveStrategy::Current, cx).await;
564                }
565            })
566            .detach();
567        });
568        cx.run_until_parked();
569        assert!(cx.read(|cx| !fixture.item.read(cx).is_dirty));
570        assert!(cx.read(|cx| inactive.read(cx).is_dirty));
571    }
572
573    struct TestTerminalProvider {
574        item: Entity<TestItem>,
575        dirty_before_spawn: Arc<Mutex<Option<bool>>>,
576    }
577
578    impl TerminalProvider for TestTerminalProvider {
579        fn spawn(
580            &self,
581            _task: task::SpawnInTerminal,
582            _window: &mut ui::Window,
583            cx: &mut App,
584        ) -> Task<Option<Result<ExitStatus>>> {
585            *self.dirty_before_spawn.lock() = Some(cx.read_entity(&self.item, |e, _| e.is_dirty));
586            Task::ready(Some(Ok(ExitStatus::default())))
587        }
588    }
589}
590
Served at tenant.openagents/omega Member data and write actions are omitted.