Skip to repository content590 lines · 21.3 KB · rust
tenant.openagents/omega
No repository description is available.
OpenAgents Git authority 2026-07-28T01:50:01.436Z 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
tasks.rs
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