add temporal service to builds. (#2842)

This commit is contained in:
Max Novich
2025-06-10 08:17:45 -07:00
committed by GitHub
parent f23fc192e4
commit cb6a819de4
44 changed files with 5145 additions and 180 deletions
+2 -2
View File
@@ -6,7 +6,7 @@ use anyhow::Result;
use etcetera::{choose_app_strategy, AppStrategy};
use goose::agents::Agent;
use goose::config::APP_STRATEGY;
use goose::scheduler::Scheduler as GooseScheduler;
use goose::scheduler_factory::SchedulerFactory;
use tower_http::cors::{Any, CorsLayer};
use tracing::info;
@@ -28,7 +28,7 @@ pub async fn run() -> Result<()> {
.data_dir()
.join("schedules.json");
let scheduler_instance = GooseScheduler::new(schedule_file_path).await?;
let scheduler_instance = SchedulerFactory::create(schedule_file_path).await?;
app_state.set_scheduler(scheduler_instance).await;
let cors = CorsLayer::new()
@@ -472,7 +472,7 @@ mod tests {
.unwrap()
.data_dir()
.join("schedules.json");
let sched = goose::scheduler::Scheduler::new(sched_storage_path)
let sched = goose::scheduler_factory::SchedulerFactory::create_legacy(sched_storage_path)
.await
.unwrap();
test_state.set_scheduler(sched).await;
+4 -3
View File
@@ -541,9 +541,10 @@ mod tests {
let state = AppState::new(Arc::new(agent), "test-secret".to_string()).await;
let scheduler_path = goose::scheduler::get_default_scheduler_storage_path()
.expect("Failed to get default scheduler storage path");
let scheduler = goose::scheduler::Scheduler::new(scheduler_path)
.await
.unwrap();
let scheduler =
goose::scheduler_factory::SchedulerFactory::create_legacy(scheduler_path)
.await
.unwrap();
state.set_scheduler(scheduler).await;
let app = routes(state);
+17 -2
View File
@@ -108,6 +108,11 @@ async fn create_schedule(
.scheduler()
.await
.map_err(|_| StatusCode::INTERNAL_SERVER_ERROR)?;
tracing::info!(
"Server: Calling scheduler.add_scheduled_job() for job '{}'",
req.id
);
let job = ScheduledJob {
id: req.id,
source: req.recipe_source,
@@ -147,7 +152,12 @@ async fn list_schedules(
.scheduler()
.await
.map_err(|_| StatusCode::INTERNAL_SERVER_ERROR)?;
let jobs = scheduler.list_scheduled_jobs().await;
tracing::info!("Server: Calling scheduler.list_scheduled_jobs()");
let jobs = scheduler.list_scheduled_jobs().await.map_err(|e| {
eprintln!("Error listing schedules: {:?}", e);
StatusCode::INTERNAL_SERVER_ERROR
})?;
Ok(Json(ListSchedulesResponse { jobs }))
}
@@ -210,6 +220,8 @@ async fn run_now_handler(
.await
.map_err(|_| StatusCode::INTERNAL_SERVER_ERROR)?;
tracing::info!("Server: Calling scheduler.run_now() for job '{}'", id);
match scheduler.run_now(&id).await {
Ok(session_id) => Ok(Json(RunNowResponse { session_id })),
Err(e) => {
@@ -408,7 +420,10 @@ async fn update_schedule(
})?;
// Return the updated schedule
let jobs = scheduler.list_scheduled_jobs().await;
let jobs = scheduler.list_scheduled_jobs().await.map_err(|e| {
eprintln!("Error listing schedules after update: {:?}", e);
StatusCode::INTERNAL_SERVER_ERROR
})?;
let updated_job = jobs
.into_iter()
.find(|job| job.id == id)
+4 -4
View File
@@ -1,5 +1,5 @@
use goose::agents::Agent;
use goose::scheduler::Scheduler;
use goose::scheduler_trait::SchedulerTrait;
use std::sync::Arc;
use tokio::sync::Mutex;
@@ -9,7 +9,7 @@ pub type AgentRef = Arc<Agent>;
pub struct AppState {
agent: Option<AgentRef>,
pub secret_key: String,
pub scheduler: Arc<Mutex<Option<Arc<Scheduler>>>>,
pub scheduler: Arc<Mutex<Option<Arc<dyn SchedulerTrait>>>>,
}
impl AppState {
@@ -27,12 +27,12 @@ impl AppState {
.ok_or_else(|| anyhow::anyhow!("Agent needs to be created first."))
}
pub async fn set_scheduler(&self, sched: Arc<Scheduler>) {
pub async fn set_scheduler(&self, sched: Arc<dyn SchedulerTrait>) {
let mut guard = self.scheduler.lock().await;
*guard = Some(sched);
}
pub async fn scheduler(&self) -> Result<Arc<Scheduler>, anyhow::Error> {
pub async fn scheduler(&self) -> Result<Arc<dyn SchedulerTrait>, anyhow::Error> {
self.scheduler
.lock()
.await