From 6481448d50bbc5fb1223e6b33df8147e3e109c06 Mon Sep 17 00:00:00 2001 From: Henry Park Date: Sat, 28 Feb 2026 19:58:07 -0800 Subject: [PATCH] feat(web): DB-backed Jobs tab + scheduler-dispatched local jobs (#436) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit * feat(web): DB-backed Jobs tab, scheduler-dispatched local jobs, remove active-jobs-bar - Remove active-jobs-bar UI element (HTML, CSS, JS polling) - Move job handlers from server.rs to handlers/jobs.rs - Remove user_id scoping (single-user gateway) - Add list_agent_jobs() and agent_job_summary() to Database trait (both postgres and libsql backends) for non-sandbox job visibility - Wire SchedulerSlot into CreateJobTool so execute_local dispatches via scheduler (persists to DB + spawns worker) instead of creating phantom ContextManager-only jobs - Update /status and /list slash commands to read from DB for consistency with Jobs tab - Fix worker mark_completed: skip if already terminal or stuck - Add agent job cancel via DB update in both web handler and slash cmd - Add Stuck → Completed guard with tracing in worker completion path Co-Authored-By: Claude Sonnet 4.6 * fix: address PR review comments - Log warning when get_context fails in worker completion path - Extract duplicated status-counting logic into AgentJobSummary::add_count() helper, used by both postgres and libsql backends Co-Authored-By: Claude Sonnet 4.6 --------- Co-authored-by: Claude Sonnet 4.6 Co-authored-by: Nick Pismenkov <50764773+nickpismenkov@users.noreply.github.com> --- src/agent/agent_loop.rs | 5 + src/agent/commands.rs | 96 +++++- src/agent/worker.rs | 31 ++ src/channels/web/handlers/jobs.rs | 220 ++++++++----- src/channels/web/server.rs | 514 +----------------------------- src/db/libsql/jobs.rs | 65 +++- src/db/mod.rs | 6 +- src/db/postgres.rs | 12 +- src/history/mod.rs | 4 +- src/history/store.rs | 88 +++++ src/main.rs | 10 + src/tools/builtin/job.rs | 53 ++- src/tools/builtin/mod.rs | 2 +- src/tools/registry.rs | 7 +- 14 files changed, 508 insertions(+), 605 deletions(-) diff --git a/src/agent/agent_loop.rs b/src/agent/agent_loop.rs index b2e5538b..38ae30d3 100644 --- a/src/agent/agent_loop.rs +++ b/src/agent/agent_loop.rs @@ -138,6 +138,11 @@ impl Agent { // Convenience accessors + /// Get the scheduler (for external wiring, e.g. CreateJobTool). + pub fn scheduler(&self) -> Arc { + Arc::clone(&self.scheduler) + } + pub(super) fn store(&self) -> Option<&Arc> { self.deps.store.as_ref() } diff --git a/src/agent/commands.rs b/src/agent/commands.rs index 56778ed4..2aab2e4e 100644 --- a/src/agent/commands.rs +++ b/src/agent/commands.rs @@ -12,6 +12,7 @@ use crate::agent::session::Session; use crate::agent::submission::SubmissionResult; use crate::agent::{Agent, MessageIntent}; use crate::channels::{IncomingMessage, StatusUpdate}; +use crate::context::JobState; use crate::error::Error; use crate::llm::{ChatMessage, Reasoning}; @@ -117,6 +118,22 @@ impl Agent { let uuid = Uuid::parse_str(&id) .map_err(|_| crate::error::JobError::NotFound { id: Uuid::nil() })?; + // Try DB first for persistent state, fall back to ContextManager. + if let Some(store) = self.store() + && let Ok(Some(ctx)) = store.get_job(uuid).await + { + return Ok(format!( + "Job: {}\nStatus: {:?}\nCreated: {}\nStarted: {}\nActual cost: {}", + ctx.title, + ctx.state, + ctx.created_at.format("%Y-%m-%d %H:%M:%S"), + ctx.started_at + .map(|t| t.format("%Y-%m-%d %H:%M:%S").to_string()) + .unwrap_or_else(|| "Not started".to_string()), + ctx.actual_cost + )); + } + let ctx = self.context_manager.get_context(uuid).await?; if ctx.user_id != user_id { return Err(crate::error::JobError::NotFound { id: uuid }.into()); @@ -134,10 +151,38 @@ impl Agent { )) } None => { - // Show summary of all jobs + // Show summary from DB for consistency with Jobs tab. + if let Some(store) = self.store() { + let mut total = 0; + let mut in_progress = 0; + let mut completed = 0; + let mut failed = 0; + let mut stuck = 0; + + if let Ok(s) = store.agent_job_summary().await { + total += s.total; + in_progress += s.in_progress; + completed += s.completed; + failed += s.failed; + stuck += s.stuck; + } + if let Ok(s) = store.sandbox_job_summary().await { + total += s.total; + in_progress += s.running; + completed += s.completed; + failed += s.failed + s.interrupted; + } + + return Ok(format!( + "Jobs summary: Total: {} In Progress: {} Completed: {} Failed: {} Stuck: {}", + total, in_progress, completed, failed, stuck + )); + } + + // Fallback to ContextManager if no DB. let summary = self.context_manager.summary_for(user_id).await; Ok(format!( - "Jobs summary:\n Total: {}\n In Progress: {}\n Completed: {}\n Failed: {}\n Stuck: {}", + "Jobs summary: Total: {} In Progress: {} Completed: {} Failed: {} Stuck: {}", summary.total, summary.in_progress, summary.completed, @@ -159,6 +204,15 @@ impl Agent { self.scheduler.stop(uuid).await?; + // Also update DB so the Jobs tab reflects cancellation immediately. + if let Some(store) = self.store() + && let Err(e) = store + .update_job_status(uuid, JobState::Cancelled, Some("Cancelled by user")) + .await + { + tracing::warn!(job_id = %uuid, "Failed to persist cancellation to DB: {}", e); + } + Ok(format!("Job {} has been cancelled.", job_id)) } @@ -167,21 +221,49 @@ impl Agent { user_id: &str, _filter: Option, ) -> Result { - let jobs = self.context_manager.all_jobs_for(user_id).await; + // List from DB for consistency with Jobs tab. + if let Some(store) = self.store() { + let agent_jobs = match store.list_agent_jobs().await { + Ok(jobs) => jobs, + Err(e) => { + tracing::warn!("Failed to list agent jobs: {}", e); + Vec::new() + } + }; + let sandbox_jobs = match store.list_sandbox_jobs().await { + Ok(jobs) => jobs, + Err(e) => { + tracing::warn!("Failed to list sandbox jobs: {}", e); + Vec::new() + } + }; + if agent_jobs.is_empty() && sandbox_jobs.is_empty() { + return Ok("No jobs found.".to_string()); + } + + let mut output = String::from("Jobs:\n"); + for j in &agent_jobs { + output.push_str(&format!(" {} - {} ({})\n", j.id, j.title, j.status)); + } + for j in &sandbox_jobs { + output.push_str(&format!(" {} - {} ({})\n", j.id, j.task, j.status)); + } + return Ok(output); + } + + // Fallback to ContextManager if no DB. + let jobs = self.context_manager.all_jobs_for(user_id).await; if jobs.is_empty() { return Ok("No jobs found.".to_string()); } let mut output = String::from("Jobs:\n"); for job_id in jobs { - if let Ok(ctx) = self.context_manager.get_context(job_id).await - && ctx.user_id == user_id - { + if let Ok(ctx) = self.context_manager.get_context(job_id).await { output.push_str(&format!(" {} - {} ({:?})\n", job_id, ctx.title, ctx.state)); } } - Ok(output) } diff --git a/src/agent/worker.rs b/src/agent/worker.rs index 87ee0ed4..67374b4d 100644 --- a/src/agent/worker.rs +++ b/src/agent/worker.rs @@ -158,6 +158,37 @@ Report when the job is complete or if you encounter issues you cannot resolve."# match result { Ok(Ok(())) => { tracing::info!("Worker for job {} completed successfully", self.job_id); + // Only mark completed if still in an active, non-stuck state. + // The execution_loop may have already called mark_completed or + // mark_stuck (e.g. "plan completed but work remains"). + let current_state = self + .context_manager() + .get_context(self.job_id) + .await + .map(|ctx| ctx.state); + match current_state { + Ok(state) if state.is_terminal() => { + // Already in a terminal state (e.g. execution_loop + // called mark_completed itself). + } + Ok(JobState::Stuck) => { + // execution_loop marked this as stuck (e.g. "plan + // completed but work remains"); leave for self-repair. + tracing::info!( + "Job {} returned Ok but is Stuck — leaving for self-repair", + self.job_id + ); + } + Ok(_) => { + self.mark_completed().await?; + } + Err(e) => { + tracing::warn!( + job_id = %self.job_id, + "Failed to get job context, cannot mark as completed: {}", e + ); + } + } } Ok(Err(e)) => { tracing::error!("Worker for job {} failed: {}", self.job_id, e); diff --git a/src/channels/web/handlers/jobs.rs b/src/channels/web/handlers/jobs.rs index 567acc7a..c32ccb85 100644 --- a/src/channels/web/handlers/jobs.rs +++ b/src/channels/web/handlers/jobs.rs @@ -1,5 +1,6 @@ //! Job and sandbox API handlers. +use std::collections::HashSet; use std::sync::Arc; use axum::{ @@ -21,32 +22,55 @@ pub async fn jobs_list_handler( "Database not available".to_string(), ))?; - // Fetch sandbox jobs scoped to the authenticated user. - let sandbox_jobs = store - .list_sandbox_jobs_for_user(&state.user_id) - .await - .map_err(|e| (StatusCode::INTERNAL_SERVER_ERROR, e.to_string()))?; + let mut jobs: Vec = Vec::new(); + let mut seen_ids: HashSet = HashSet::new(); - // Scope jobs to the authenticated user. - let mut jobs: Vec = sandbox_jobs - .iter() - .filter(|j| j.user_id == state.user_id) - .map(|j| { - let ui_state = match j.status.as_str() { - "creating" => "pending", - "running" => "in_progress", - s => s, - }; - JobInfo { - id: j.id, - title: j.task.clone(), - state: ui_state.to_string(), - user_id: j.user_id.clone(), - created_at: j.created_at.to_rfc3339(), - started_at: j.started_at.map(|dt| dt.to_rfc3339()), + // Fetch sandbox jobs from database. + match store.list_sandbox_jobs().await { + Ok(sandbox_jobs) => { + for j in &sandbox_jobs { + let ui_state = match j.status.as_str() { + "creating" => "pending", + "running" => "in_progress", + s => s, + }; + seen_ids.insert(j.id); + jobs.push(JobInfo { + id: j.id, + title: j.task.clone(), + state: ui_state.to_string(), + user_id: j.user_id.clone(), + created_at: j.created_at.to_rfc3339(), + started_at: j.started_at.map(|dt| dt.to_rfc3339()), + }); } - }) - .collect(); + } + Err(e) => { + tracing::warn!("Failed to list sandbox jobs: {}", e); + } + } + + // Fetch agent (non-sandbox) jobs from database, deduplicating by ID. + match store.list_agent_jobs().await { + Ok(agent_jobs) => { + for j in &agent_jobs { + if seen_ids.contains(&j.id) { + continue; + } + jobs.push(JobInfo { + id: j.id, + title: j.title.clone(), + state: j.status.clone(), + user_id: j.user_id.clone(), + created_at: j.created_at.to_rfc3339(), + started_at: j.started_at.map(|dt| dt.to_rfc3339()), + }); + } + } + Err(e) => { + tracing::warn!("Failed to list agent jobs: {}", e); + } + } // Most recent first. jobs.sort_by(|a, b| b.created_at.cmp(&a.created_at)); @@ -62,18 +86,49 @@ pub async fn jobs_summary_handler( "Database not available".to_string(), ))?; - let s = store - .sandbox_job_summary_for_user(&state.user_id) - .await - .map_err(|e| (StatusCode::INTERNAL_SERVER_ERROR, e.to_string()))?; + let mut total = 0; + let mut pending = 0; + let mut in_progress = 0; + let mut completed = 0; + let mut failed = 0; + let mut stuck = 0; + + // Sandbox job counts. + match store.sandbox_job_summary().await { + Ok(s) => { + total += s.total; + pending += s.creating; + in_progress += s.running; + completed += s.completed; + failed += s.failed + s.interrupted; + } + Err(e) => { + tracing::warn!("Failed to fetch sandbox job summary: {}", e); + } + } + + // Agent job counts. + match store.agent_job_summary().await { + Ok(s) => { + total += s.total; + pending += s.pending; + in_progress += s.in_progress; + completed += s.completed; + failed += s.failed; + stuck += s.stuck; + } + Err(e) => { + tracing::warn!("Failed to fetch agent job summary: {}", e); + } + } Ok(Json(JobSummaryResponse { - total: s.total, - pending: s.creating, - in_progress: s.running, - completed: s.completed, - failed: s.failed + s.interrupted, - stuck: 0, + total, + pending, + in_progress, + completed, + failed, + stuck, })) } @@ -81,16 +136,16 @@ pub async fn jobs_detail_handler( State(state): State>, Path(id): Path, ) -> Result, (StatusCode, String)> { + let store = state.store.as_ref().ok_or(( + StatusCode::SERVICE_UNAVAILABLE, + "Database not available".to_string(), + ))?; + let job_id = Uuid::parse_str(&id) .map_err(|_| (StatusCode::BAD_REQUEST, "Invalid job ID".to_string()))?; - // Try sandbox job from DB first, scoped to the authenticated user. - if let Some(ref store) = state.store - && let Ok(Some(job)) = store.get_sandbox_job(job_id).await - { - if job.user_id != state.user_id { - return Err((StatusCode::NOT_FOUND, "Job not found".to_string())); - } + // Try sandbox job from DB first. + if let Ok(Some(job)) = store.get_sandbox_job(job_id).await { let browse_id = std::path::Path::new(&job.project_dir) .file_name() .map(|n| n.to_string_lossy().to_string()) @@ -146,6 +201,30 @@ pub async fn jobs_detail_handler( })); } + // Fall back to agent job from DB. + if let Ok(Some(ctx)) = store.get_job(job_id).await { + let elapsed_secs = ctx.started_at.map(|start| { + let end = ctx.completed_at.unwrap_or_else(chrono::Utc::now); + (end - start).num_seconds().max(0) as u64 + }); + + return Ok(Json(JobDetailResponse { + id: ctx.job_id, + title: ctx.title.clone(), + description: ctx.description.clone(), + state: ctx.state.to_string(), + user_id: ctx.user_id.clone(), + created_at: ctx.created_at.to_rfc3339(), + started_at: ctx.started_at.map(|dt| dt.to_rfc3339()), + completed_at: ctx.completed_at.map(|dt| dt.to_rfc3339()), + elapsed_secs, + project_dir: None, + browse_url: None, + job_mode: None, + transitions: Vec::new(), + })); + } + Err((StatusCode::NOT_FOUND, "Job not found".to_string())) } @@ -156,13 +235,10 @@ pub async fn jobs_cancel_handler( let job_id = Uuid::parse_str(&id) .map_err(|_| (StatusCode::BAD_REQUEST, "Invalid job ID".to_string()))?; - // Try sandbox job cancellation, scoped to the authenticated user. + // Try sandbox job cancellation. if let Some(ref store) = state.store && let Ok(Some(job)) = store.get_sandbox_job(job_id).await { - if job.user_id != state.user_id { - return Err((StatusCode::NOT_FOUND, "Job not found".to_string())); - } if job.status == "running" || job.status == "creating" { // Stop the container if we have a job manager. if let Some(ref jm) = state.job_manager @@ -188,6 +264,26 @@ pub async fn jobs_cancel_handler( }))); } + // Fall back to agent job cancellation via DB status update. + if let Some(ref store) = state.store + && let Ok(Some(job)) = store.get_job(job_id).await + { + if job.state.is_active() { + store + .update_job_status( + job_id, + crate::context::JobState::Cancelled, + Some("Cancelled by user"), + ) + .await + .map_err(|e| (StatusCode::INTERNAL_SERVER_ERROR, e.to_string()))?; + } + return Ok(Json(serde_json::json!({ + "status": "cancelled", + "job_id": job_id, + }))); + } + Err((StatusCode::NOT_FOUND, "Job not found".to_string())) } @@ -213,11 +309,6 @@ pub async fn jobs_restart_handler( .map_err(|e| (StatusCode::INTERNAL_SERVER_ERROR, e.to_string()))? .ok_or((StatusCode::NOT_FOUND, "Job not found".to_string()))?; - // Scope to the authenticated user. - if old_job.user_id != state.user_id { - return Err((StatusCode::NOT_FOUND, "Job not found".to_string())); - } - if old_job.status != "interrupted" && old_job.status != "failed" { return Err(( StatusCode::CONFLICT, @@ -310,16 +401,6 @@ pub async fn jobs_prompt_handler( .parse() .map_err(|_| (StatusCode::BAD_REQUEST, "Invalid job ID".to_string()))?; - // Verify user owns this job. - if let Some(ref store) = state.store - && !store - .sandbox_job_belongs_to_user(job_id, &state.user_id) - .await - .unwrap_or(false) - { - return Err((StatusCode::NOT_FOUND, "Job not found".to_string())); - } - let content = body .get("content") .and_then(|v| v.as_str()) @@ -358,15 +439,6 @@ pub async fn jobs_events_handler( .parse() .map_err(|_| (StatusCode::BAD_REQUEST, "Invalid job ID".to_string()))?; - // Verify user owns this job. - if !store - .sandbox_job_belongs_to_user(job_id, &state.user_id) - .await - .unwrap_or(false) - { - return Err((StatusCode::NOT_FOUND, "Job not found".to_string())); - } - let events = store .list_job_events(job_id, None) .await @@ -416,11 +488,6 @@ pub async fn job_files_list_handler( .map_err(|e| (StatusCode::INTERNAL_SERVER_ERROR, e.to_string()))? .ok_or((StatusCode::NOT_FOUND, "Job not found".to_string()))?; - // Verify user owns this job. - if job.user_id != state.user_id { - return Err((StatusCode::NOT_FOUND, "Job not found".to_string())); - } - let base = std::path::PathBuf::from(&job.project_dir); let rel_path = query.path.as_deref().unwrap_or(""); let target = base.join(rel_path); @@ -484,11 +551,6 @@ pub async fn job_files_read_handler( .map_err(|e| (StatusCode::INTERNAL_SERVER_ERROR, e.to_string()))? .ok_or((StatusCode::NOT_FOUND, "Job not found".to_string()))?; - // Verify user owns this job. - if job.user_id != state.user_id { - return Err((StatusCode::NOT_FOUND, "Job not found".to_string())); - } - let path = query.path.as_deref().ok_or(( StatusCode::BAD_REQUEST, "path parameter required".to_string(), diff --git a/src/channels/web/server.rs b/src/channels/web/server.rs index 0964a01b..7bcb30eb 100644 --- a/src/channels/web/server.rs +++ b/src/channels/web/server.rs @@ -29,6 +29,11 @@ use crate::agent::SessionManager; use crate::bootstrap::ironclaw_base_dir; use crate::channels::IncomingMessage; use crate::channels::web::auth::{AuthState, auth_middleware}; +use crate::channels::web::handlers::jobs::{ + job_files_list_handler, job_files_read_handler, jobs_cancel_handler, jobs_detail_handler, + jobs_events_handler, jobs_list_handler, jobs_prompt_handler, jobs_restart_handler, + jobs_summary_handler, +}; use crate::channels::web::handlers::skills::{ skills_install_handler, skills_list_handler, skills_remove_handler, skills_search_handler, }; @@ -1116,514 +1121,7 @@ async fn memory_search_handler( Ok(Json(MemorySearchResponse { results: hits })) } -// --- Jobs handlers --- - -async fn jobs_list_handler( - State(state): State>, -) -> Result, (StatusCode, String)> { - let store = state.store.as_ref().ok_or(( - StatusCode::SERVICE_UNAVAILABLE, - "Database not available".to_string(), - ))?; - - // Fetch sandbox jobs scoped to the authenticated user. - let sandbox_jobs = store - .list_sandbox_jobs_for_user(&state.user_id) - .await - .map_err(|e| (StatusCode::INTERNAL_SERVER_ERROR, e.to_string()))?; - - // Scope jobs to the authenticated user. - let mut jobs: Vec = sandbox_jobs - .iter() - .filter(|j| j.user_id == state.user_id) - .map(|j| { - let ui_state = match j.status.as_str() { - "creating" => "pending", - "running" => "in_progress", - s => s, - }; - JobInfo { - id: j.id, - title: j.task.clone(), - state: ui_state.to_string(), - user_id: j.user_id.clone(), - created_at: j.created_at.to_rfc3339(), - started_at: j.started_at.map(|dt| dt.to_rfc3339()), - } - }) - .collect(); - - // Most recent first. - jobs.sort_by(|a, b| b.created_at.cmp(&a.created_at)); - - Ok(Json(JobListResponse { jobs })) -} - -async fn jobs_summary_handler( - State(state): State>, -) -> Result, (StatusCode, String)> { - let store = state.store.as_ref().ok_or(( - StatusCode::SERVICE_UNAVAILABLE, - "Database not available".to_string(), - ))?; - - let s = store - .sandbox_job_summary_for_user(&state.user_id) - .await - .map_err(|e| (StatusCode::INTERNAL_SERVER_ERROR, e.to_string()))?; - - Ok(Json(JobSummaryResponse { - total: s.total, - pending: s.creating, - in_progress: s.running, - completed: s.completed, - failed: s.failed + s.interrupted, - stuck: 0, - })) -} - -async fn jobs_detail_handler( - State(state): State>, - Path(id): Path, -) -> Result, (StatusCode, String)> { - let job_id = Uuid::parse_str(&id) - .map_err(|_| (StatusCode::BAD_REQUEST, "Invalid job ID".to_string()))?; - - // Try sandbox job from DB first, scoped to the authenticated user. - if let Some(ref store) = state.store - && let Ok(Some(job)) = store.get_sandbox_job(job_id).await - { - if job.user_id != state.user_id { - return Err((StatusCode::NOT_FOUND, "Job not found".to_string())); - } - let browse_id = std::path::Path::new(&job.project_dir) - .file_name() - .map(|n| n.to_string_lossy().to_string()) - .unwrap_or_else(|| job.id.to_string()); - - let ui_state = match job.status.as_str() { - "creating" => "pending", - "running" => "in_progress", - s => s, - }; - - let elapsed_secs = job.started_at.map(|start| { - let end = job.completed_at.unwrap_or_else(chrono::Utc::now); - (end - start).num_seconds().max(0) as u64 - }); - - // Synthesize transitions from timestamps. - let mut transitions = Vec::new(); - if let Some(started) = job.started_at { - transitions.push(TransitionInfo { - from: "creating".to_string(), - to: "running".to_string(), - timestamp: started.to_rfc3339(), - reason: None, - }); - } - if let Some(completed) = job.completed_at { - transitions.push(TransitionInfo { - from: "running".to_string(), - to: job.status.clone(), - timestamp: completed.to_rfc3339(), - reason: job.failure_reason.clone(), - }); - } - - return Ok(Json(JobDetailResponse { - id: job.id, - title: job.task.clone(), - description: String::new(), - state: ui_state.to_string(), - user_id: job.user_id.clone(), - created_at: job.created_at.to_rfc3339(), - started_at: job.started_at.map(|dt| dt.to_rfc3339()), - completed_at: job.completed_at.map(|dt| dt.to_rfc3339()), - elapsed_secs, - project_dir: Some(job.project_dir.clone()), - browse_url: Some(format!("/projects/{}/", browse_id)), - job_mode: { - let mode = store.get_sandbox_job_mode(job.id).await.ok().flatten(); - mode.filter(|m| m != "worker") - }, - transitions, - })); - } - - Err((StatusCode::NOT_FOUND, "Job not found".to_string())) -} - -async fn jobs_cancel_handler( - State(state): State>, - Path(id): Path, -) -> Result, (StatusCode, String)> { - let job_id = Uuid::parse_str(&id) - .map_err(|_| (StatusCode::BAD_REQUEST, "Invalid job ID".to_string()))?; - - // Try sandbox job cancellation, scoped to the authenticated user. - if let Some(ref store) = state.store - && let Ok(Some(job)) = store.get_sandbox_job(job_id).await - { - if job.user_id != state.user_id { - return Err((StatusCode::NOT_FOUND, "Job not found".to_string())); - } - if job.status == "running" || job.status == "creating" { - // Stop the container if we have a job manager. - if let Some(ref jm) = state.job_manager - && let Err(e) = jm.stop_job(job_id).await - { - tracing::warn!(job_id = %job_id, error = %e, "Failed to stop container during cancellation"); - } - store - .update_sandbox_job_status( - job_id, - "failed", - Some(false), - Some("Cancelled by user"), - None, - Some(chrono::Utc::now()), - ) - .await - .map_err(|e| (StatusCode::INTERNAL_SERVER_ERROR, e.to_string()))?; - } - return Ok(Json(serde_json::json!({ - "status": "cancelled", - "job_id": job_id, - }))); - } - - Err((StatusCode::NOT_FOUND, "Job not found".to_string())) -} - -async fn jobs_restart_handler( - State(state): State>, - Path(id): Path, -) -> Result, (StatusCode, String)> { - let store = state.store.as_ref().ok_or(( - StatusCode::SERVICE_UNAVAILABLE, - "Database not available".to_string(), - ))?; - let jm = state.job_manager.as_ref().ok_or(( - StatusCode::SERVICE_UNAVAILABLE, - "Sandbox not enabled".to_string(), - ))?; - - let old_job_id = Uuid::parse_str(&id) - .map_err(|_| (StatusCode::BAD_REQUEST, "Invalid job ID".to_string()))?; - - let old_job = store - .get_sandbox_job(old_job_id) - .await - .map_err(|e| (StatusCode::INTERNAL_SERVER_ERROR, e.to_string()))? - .ok_or((StatusCode::NOT_FOUND, "Job not found".to_string()))?; - - // Scope to the authenticated user. - if old_job.user_id != state.user_id { - return Err((StatusCode::NOT_FOUND, "Job not found".to_string())); - } - - if old_job.status != "interrupted" && old_job.status != "failed" { - return Err(( - StatusCode::CONFLICT, - format!("Cannot restart job in state '{}'", old_job.status), - )); - } - - // Create a new job with the same task and project_dir. - let new_job_id = Uuid::new_v4(); - let now = chrono::Utc::now(); - - let record = crate::history::SandboxJobRecord { - id: new_job_id, - task: old_job.task.clone(), - status: "creating".to_string(), - user_id: old_job.user_id.clone(), - project_dir: old_job.project_dir.clone(), - success: None, - failure_reason: None, - created_at: now, - started_at: None, - completed_at: None, - credential_grants_json: old_job.credential_grants_json.clone(), - }; - store - .save_sandbox_job(&record) - .await - .map_err(|e| (StatusCode::INTERNAL_SERVER_ERROR, e.to_string()))?; - - // Look up the original job's mode so the restart uses the same mode. - let mode = match store.get_sandbox_job_mode(old_job_id).await { - Ok(Some(m)) if m == "claude_code" => crate::orchestrator::job_manager::JobMode::ClaudeCode, - _ => crate::orchestrator::job_manager::JobMode::Worker, - }; - - // Restore credential grants from the original job so the restarted container - // has access to the same secrets. - let credential_grants: Vec = - serde_json::from_str(&old_job.credential_grants_json).unwrap_or_else(|e| { - tracing::warn!( - job_id = %old_job.id, - "Failed to deserialize credential grants from stored job: {}. \ - Restarted job will have no credentials.", - e - ); - vec![] - }); - - let project_dir = std::path::PathBuf::from(&old_job.project_dir); - let _token = jm - .create_job( - new_job_id, - &old_job.task, - Some(project_dir), - mode, - credential_grants, - ) - .await - .map_err(|e| { - ( - StatusCode::INTERNAL_SERVER_ERROR, - format!("Failed to create container: {}", e), - ) - })?; - - store - .update_sandbox_job_status(new_job_id, "running", None, None, Some(now), None) - .await - .map_err(|e| (StatusCode::INTERNAL_SERVER_ERROR, e.to_string()))?; - - Ok(Json(serde_json::json!({ - "status": "restarted", - "old_job_id": old_job_id, - "new_job_id": new_job_id, - }))) -} - -// --- Claude Code prompt and events handlers --- - -/// Submit a follow-up prompt to a running Claude Code sandbox job. -async fn jobs_prompt_handler( - State(state): State>, - Path(id): Path, - Json(body): Json, -) -> Result, (StatusCode, String)> { - let prompt_queue = state.prompt_queue.as_ref().ok_or(( - StatusCode::NOT_IMPLEMENTED, - "Claude Code not configured".to_string(), - ))?; - - let job_id: uuid::Uuid = id - .parse() - .map_err(|_| (StatusCode::BAD_REQUEST, "Invalid job ID".to_string()))?; - - // Verify user owns this job. - if let Some(ref store) = state.store - && !store - .sandbox_job_belongs_to_user(job_id, &state.user_id) - .await - .unwrap_or(false) - { - return Err((StatusCode::NOT_FOUND, "Job not found".to_string())); - } - - let content = body - .get("content") - .and_then(|v| v.as_str()) - .ok_or(( - StatusCode::BAD_REQUEST, - "Missing 'content' field".to_string(), - ))? - .to_string(); - - let done = body.get("done").and_then(|v| v.as_bool()).unwrap_or(false); - - let prompt = crate::orchestrator::api::PendingPrompt { content, done }; - - { - let mut queue = prompt_queue.lock().await; - queue.entry(job_id).or_default().push_back(prompt); - } - - Ok(Json(serde_json::json!({ - "status": "queued", - "job_id": job_id.to_string(), - }))) -} - -/// Load persisted job events for a job (for history replay on page open). -async fn jobs_events_handler( - State(state): State>, - Path(id): Path, -) -> Result, (StatusCode, String)> { - let store = state.store.as_ref().ok_or(( - StatusCode::NOT_IMPLEMENTED, - "Database not available".to_string(), - ))?; - - let job_id: uuid::Uuid = id - .parse() - .map_err(|_| (StatusCode::BAD_REQUEST, "Invalid job ID".to_string()))?; - - // Verify user owns this job. - if !store - .sandbox_job_belongs_to_user(job_id, &state.user_id) - .await - .unwrap_or(false) - { - return Err((StatusCode::NOT_FOUND, "Job not found".to_string())); - } - - let events = store - .list_job_events(job_id, None) - .await - .map_err(|e| (StatusCode::INTERNAL_SERVER_ERROR, e.to_string()))?; - - let events_json: Vec = events - .into_iter() - .map(|e| { - serde_json::json!({ - "id": e.id, - "event_type": e.event_type, - "data": e.data, - "created_at": e.created_at.to_rfc3339(), - }) - }) - .collect(); - - Ok(Json(serde_json::json!({ - "job_id": job_id.to_string(), - "events": events_json, - }))) -} - -// --- Project file handlers for sandbox jobs --- - -#[derive(Deserialize)] -struct FilePathQuery { - path: Option, -} - -async fn job_files_list_handler( - State(state): State>, - Path(id): Path, - Query(query): Query, -) -> Result, (StatusCode, String)> { - let store = state.store.as_ref().ok_or(( - StatusCode::SERVICE_UNAVAILABLE, - "Database not available".to_string(), - ))?; - - let job_id = Uuid::parse_str(&id) - .map_err(|_| (StatusCode::BAD_REQUEST, "Invalid job ID".to_string()))?; - - let job = store - .get_sandbox_job(job_id) - .await - .map_err(|e| (StatusCode::INTERNAL_SERVER_ERROR, e.to_string()))? - .ok_or((StatusCode::NOT_FOUND, "Job not found".to_string()))?; - - // Verify user owns this job. - if job.user_id != state.user_id { - return Err((StatusCode::NOT_FOUND, "Job not found".to_string())); - } - - let base = std::path::PathBuf::from(&job.project_dir); - let rel_path = query.path.as_deref().unwrap_or(""); - let target = base.join(rel_path); - - // Path traversal guard. - let canonical = target - .canonicalize() - .map_err(|_| (StatusCode::NOT_FOUND, "Path not found".to_string()))?; - let base_canonical = base - .canonicalize() - .map_err(|_| (StatusCode::NOT_FOUND, "Project dir not found".to_string()))?; - if !canonical.starts_with(&base_canonical) { - return Err((StatusCode::FORBIDDEN, "Forbidden".to_string())); - } - - let mut entries = Vec::new(); - let mut read_dir = tokio::fs::read_dir(&canonical) - .await - .map_err(|_| (StatusCode::NOT_FOUND, "Cannot read directory".to_string()))?; - - while let Ok(Some(entry)) = read_dir.next_entry().await { - let name = entry.file_name().to_string_lossy().to_string(); - let is_dir = entry - .file_type() - .await - .map(|ft| ft.is_dir()) - .unwrap_or(false); - let rel = if rel_path.is_empty() { - name.clone() - } else { - format!("{}/{}", rel_path, name) - }; - entries.push(ProjectFileEntry { - name, - path: rel, - is_dir, - }); - } - - entries.sort_by(|a, b| b.is_dir.cmp(&a.is_dir).then_with(|| a.name.cmp(&b.name))); - - Ok(Json(ProjectFilesResponse { entries })) -} - -async fn job_files_read_handler( - State(state): State>, - Path(id): Path, - Query(query): Query, -) -> Result, (StatusCode, String)> { - let store = state.store.as_ref().ok_or(( - StatusCode::SERVICE_UNAVAILABLE, - "Database not available".to_string(), - ))?; - - let job_id = Uuid::parse_str(&id) - .map_err(|_| (StatusCode::BAD_REQUEST, "Invalid job ID".to_string()))?; - - let job = store - .get_sandbox_job(job_id) - .await - .map_err(|e| (StatusCode::INTERNAL_SERVER_ERROR, e.to_string()))? - .ok_or((StatusCode::NOT_FOUND, "Job not found".to_string()))?; - - // Verify user owns this job. - if job.user_id != state.user_id { - return Err((StatusCode::NOT_FOUND, "Job not found".to_string())); - } - - let path = query.path.as_deref().ok_or(( - StatusCode::BAD_REQUEST, - "path parameter required".to_string(), - ))?; - - let base = std::path::PathBuf::from(&job.project_dir); - let file_path = base.join(path); - - let canonical = file_path - .canonicalize() - .map_err(|_| (StatusCode::NOT_FOUND, "File not found".to_string()))?; - let base_canonical = base - .canonicalize() - .map_err(|_| (StatusCode::NOT_FOUND, "Project dir not found".to_string()))?; - if !canonical.starts_with(&base_canonical) { - return Err((StatusCode::FORBIDDEN, "Forbidden".to_string())); - } - - let content = tokio::fs::read_to_string(&canonical) - .await - .map_err(|_| (StatusCode::NOT_FOUND, "Cannot read file".to_string()))?; - - Ok(Json(ProjectFileReadResponse { - path: path.to_string(), - content, - })) -} - +// Job handlers moved to handlers/jobs.rs // --- Logs handlers --- async fn logs_events_handler( diff --git a/src/db/libsql/jobs.rs b/src/db/libsql/jobs.rs index 0c46b231..78a55b81 100644 --- a/src/db/libsql/jobs.rs +++ b/src/db/libsql/jobs.rs @@ -12,7 +12,7 @@ use super::{ use crate::context::{ActionRecord, JobContext, JobState}; use crate::db::JobStore; use crate::error::DatabaseError; -use crate::history::LlmCallRecord; +use crate::history::{AgentJobRecord, AgentJobSummary, LlmCallRecord}; use chrono::Utc; @@ -173,6 +173,69 @@ impl JobStore for LibSqlBackend { Ok(ids) } + async fn list_agent_jobs(&self) -> Result, DatabaseError> { + let conn = self.connect().await?; + let mut rows = conn + .query( + r#" + SELECT id, title, status, user_id, failure_reason, + created_at, started_at, completed_at + FROM agent_jobs WHERE source = 'direct' + ORDER BY created_at DESC + "#, + (), + ) + .await + .map_err(|e| DatabaseError::Query(e.to_string()))?; + + let mut jobs = Vec::new(); + while let Some(row) = rows + .next() + .await + .map_err(|e| DatabaseError::Query(e.to_string()))? + { + let id_str = get_text(&row, 0); + let Ok(id) = id_str.parse() else { + tracing::warn!("Skipping agent job with invalid UUID: {}", id_str); + continue; + }; + jobs.push(AgentJobRecord { + id, + title: get_text(&row, 1), + status: get_text(&row, 2), + user_id: get_text(&row, 3), + failure_reason: get_opt_text(&row, 4), + created_at: get_ts(&row, 5), + started_at: get_opt_ts(&row, 6), + completed_at: get_opt_ts(&row, 7), + }); + } + Ok(jobs) + } + + async fn agent_job_summary(&self) -> Result { + let conn = self.connect().await?; + let mut rows = conn + .query( + "SELECT status, COUNT(*) as cnt FROM agent_jobs WHERE source = 'direct' GROUP BY status", + (), + ) + .await + .map_err(|e| DatabaseError::Query(e.to_string()))?; + + let mut summary = AgentJobSummary::default(); + while let Some(row) = rows + .next() + .await + .map_err(|e| DatabaseError::Query(e.to_string()))? + { + let status = get_text(&row, 0); + let count = get_i64(&row, 1) as usize; + summary.add_count(&status, count); + } + Ok(summary) + } + async fn save_action(&self, job_id: Uuid, action: &ActionRecord) -> Result<(), DatabaseError> { let conn = self.connect().await?; let duration_ms = action.duration.as_millis() as i64; diff --git a/src/db/mod.rs b/src/db/mod.rs index 1101a311..7a6b8941 100644 --- a/src/db/mod.rs +++ b/src/db/mod.rs @@ -32,8 +32,8 @@ use crate::context::{ActionRecord, JobContext, JobState}; use crate::error::DatabaseError; use crate::error::WorkspaceError; use crate::history::{ - ConversationMessage, ConversationSummary, JobEventRecord, LlmCallRecord, SandboxJobRecord, - SandboxJobSummary, SettingRow, + AgentJobRecord, AgentJobSummary, ConversationMessage, ConversationSummary, JobEventRecord, + LlmCallRecord, SandboxJobRecord, SandboxJobSummary, SettingRow, }; use crate::workspace::{MemoryChunk, MemoryDocument, WorkspaceEntry}; use crate::workspace::{SearchConfig, SearchResult}; @@ -172,6 +172,8 @@ pub trait JobStore: Send + Sync { ) -> Result<(), DatabaseError>; async fn mark_job_stuck(&self, id: Uuid) -> Result<(), DatabaseError>; async fn get_stuck_jobs(&self) -> Result, DatabaseError>; + async fn list_agent_jobs(&self) -> Result, DatabaseError>; + async fn agent_job_summary(&self) -> Result; async fn save_action(&self, job_id: Uuid, action: &ActionRecord) -> Result<(), DatabaseError>; async fn get_job_actions(&self, job_id: Uuid) -> Result, DatabaseError>; async fn record_llm_call(&self, record: &LlmCallRecord<'_>) -> Result; diff --git a/src/db/postgres.rs b/src/db/postgres.rs index 094b268f..27d5e70b 100644 --- a/src/db/postgres.rs +++ b/src/db/postgres.rs @@ -21,8 +21,8 @@ use crate::db::{ }; use crate::error::{DatabaseError, WorkspaceError}; use crate::history::{ - ConversationMessage, ConversationSummary, JobEventRecord, LlmCallRecord, SandboxJobRecord, - SandboxJobSummary, SettingRow, Store, + AgentJobRecord, AgentJobSummary, ConversationMessage, ConversationSummary, JobEventRecord, + LlmCallRecord, SandboxJobRecord, SandboxJobSummary, SettingRow, Store, }; use crate::workspace::{ MemoryChunk, MemoryDocument, Repository, SearchConfig, SearchResult, WorkspaceEntry, @@ -215,6 +215,14 @@ impl JobStore for PgBackend { self.store.get_stuck_jobs().await } + async fn list_agent_jobs(&self) -> Result, DatabaseError> { + self.store.list_agent_jobs().await + } + + async fn agent_job_summary(&self) -> Result { + self.store.agent_job_summary().await + } + async fn save_action(&self, job_id: Uuid, action: &ActionRecord) -> Result<(), DatabaseError> { self.store.save_action(job_id, action).await } diff --git a/src/history/mod.rs b/src/history/mod.rs index 4f448cfc..3194a680 100644 --- a/src/history/mod.rs +++ b/src/history/mod.rs @@ -14,6 +14,6 @@ pub use analytics::{JobStats, ToolStats}; #[cfg(feature = "postgres")] pub use store::Store; pub use store::{ - ConversationMessage, ConversationSummary, JobEventRecord, LlmCallRecord, SandboxJobRecord, - SandboxJobSummary, SettingRow, + AgentJobRecord, AgentJobSummary, ConversationMessage, ConversationSummary, JobEventRecord, + LlmCallRecord, SandboxJobRecord, SandboxJobSummary, SettingRow, }; diff --git a/src/history/store.rs b/src/history/store.rs index 7e2ff0f3..592c876e 100644 --- a/src/history/store.rs +++ b/src/history/store.rs @@ -487,6 +487,45 @@ pub struct SandboxJobSummary { pub interrupted: usize, } +/// Lightweight record for agent (non-sandbox) jobs, used by the web Jobs tab. +#[derive(Debug, Clone)] +pub struct AgentJobRecord { + pub id: Uuid, + pub title: String, + pub status: String, + pub user_id: String, + pub created_at: DateTime, + pub started_at: Option>, + pub completed_at: Option>, + pub failure_reason: Option, +} + +/// Summary counts for agent (non-sandbox) jobs. +#[derive(Debug, Clone, Default)] +pub struct AgentJobSummary { + pub total: usize, + pub pending: usize, + pub in_progress: usize, + pub completed: usize, + pub failed: usize, + pub stuck: usize, +} + +impl AgentJobSummary { + /// Accumulate a status/count pair into the summary buckets. + pub fn add_count(&mut self, status: &str, count: usize) { + self.total += count; + match status { + "pending" => self.pending += count, + "in_progress" => self.in_progress += count, + "completed" | "submitted" | "accepted" => self.completed += count, + "failed" | "cancelled" => self.failed += count, + "stuck" => self.stuck += count, + _ => {} + } + } +} + #[cfg(feature = "postgres")] impl Store { /// Insert a new sandbox job into `agent_jobs`. @@ -754,6 +793,55 @@ impl Store { } Ok(summary) } + + /// List all agent (non-sandbox) jobs, most recent first. + pub async fn list_agent_jobs(&self) -> Result, DatabaseError> { + let conn = self.conn().await?; + let rows = conn + .query( + r#" + SELECT id, title, status, user_id, failure_reason, + created_at, started_at, completed_at + FROM agent_jobs WHERE source = 'direct' + ORDER BY created_at DESC + "#, + &[], + ) + .await?; + + Ok(rows + .iter() + .map(|r| AgentJobRecord { + id: r.get("id"), + title: r.get("title"), + status: r.get("status"), + user_id: r.get::<_, Option>("user_id").unwrap_or_default(), + created_at: r.get("created_at"), + started_at: r.get("started_at"), + completed_at: r.get("completed_at"), + failure_reason: r.get("failure_reason"), + }) + .collect()) + } + + /// Summary counts for agent (non-sandbox) jobs. + pub async fn agent_job_summary(&self) -> Result { + let conn = self.conn().await?; + let rows = conn + .query( + "SELECT status, COUNT(*) as cnt FROM agent_jobs WHERE source = 'direct' GROUP BY status", + &[], + ) + .await?; + + let mut summary = AgentJobSummary::default(); + for row in &rows { + let status: String = row.get("status"); + let count: i64 = row.get("cnt"); + summary.add_count(&status, count as usize); + } + Ok(summary) + } } // ==================== Job Events ==================== diff --git a/src/main.rs b/src/main.rs index b545159c..55a56b1d 100644 --- a/src/main.rs +++ b/src/main.rs @@ -456,9 +456,16 @@ async fn async_main() -> anyhow::Result<()> { let session_manager = Arc::new(ironclaw::agent::SessionManager::new().with_hooks(components.hooks.clone())); + // Lazy scheduler slot — filled after Agent::new creates the Scheduler. + // Allows CreateJobTool to dispatch local jobs via the Scheduler even though + // the Scheduler is created after tools are registered (chicken-and-egg). + let scheduler_slot: ironclaw::tools::builtin::SchedulerSlot = + Arc::new(tokio::sync::RwLock::new(None)); + // Register job tools (sandbox deps auto-injected when container_job_manager is available) components.tools.register_job_tools( Arc::clone(&components.context_manager), + Some(scheduler_slot.clone()), container_job_manager.clone(), components.db.clone(), job_event_tx.clone(), @@ -648,6 +655,9 @@ async fn async_main() -> anyhow::Result<()> { Some(session_manager), ); + // Fill the scheduler slot now that Agent (and its Scheduler) exist. + *scheduler_slot.write().await = Some(agent.scheduler()); + agent.run().await?; // ── Shutdown ──────────────────────────────────────────────────────── diff --git a/src/tools/builtin/job.rs b/src/tools/builtin/job.rs index 7da26577..a571773e 100644 --- a/src/tools/builtin/job.rs +++ b/src/tools/builtin/job.rs @@ -12,6 +12,7 @@ use std::time::Duration; use async_trait::async_trait; use chrono::Utc; +use tokio::sync::RwLock; use uuid::Uuid; use crate::bootstrap::ironclaw_base_dir; @@ -25,6 +26,12 @@ use crate::orchestrator::job_manager::{ContainerJobManager, JobMode}; use crate::secrets::SecretsStore; use crate::tools::tool::{ApprovalRequirement, Tool, ToolError, ToolOutput, require_str}; +/// Lazy scheduler reference, filled after Agent::new creates the Scheduler. +/// +/// Solves the chicken-and-egg: tools are registered before the Scheduler exists +/// (Scheduler needs the ToolRegistry). Created empty, filled after Agent::new. +pub type SchedulerSlot = Arc>>>; + /// Resolve a job ID from a full UUID or a short prefix (like git short SHAs). /// /// Tries full UUID parse first. If that fails, treats the input as a hex prefix @@ -73,6 +80,8 @@ async fn resolve_job_id(input: &str, context_manager: &ContextManager) -> Result /// job via the ContextManager. The LLM never needs to know the difference. pub struct CreateJobTool { context_manager: Arc, + /// Lazy scheduler for dispatching local (non-sandbox) jobs. + scheduler_slot: Option, job_manager: Option>, store: Option>, /// Broadcast sender for job events (used to subscribe a monitor). @@ -87,6 +96,7 @@ impl CreateJobTool { pub fn new(context_manager: Arc) -> Self { Self { context_manager, + scheduler_slot: None, job_manager: None, store: None, event_tx: None, @@ -118,6 +128,12 @@ impl CreateJobTool { self } + /// Inject a lazy scheduler slot for dispatching local (non-sandbox) jobs. + pub fn with_scheduler_slot(mut self, slot: SchedulerSlot) -> Self { + self.scheduler_slot = Some(slot); + self + } + /// Inject secrets store for credential validation. pub fn with_secrets(mut self, secrets: Arc) -> Self { self.secrets_store = Some(secrets); @@ -239,7 +255,8 @@ impl CreateJobTool { } } - /// Execute via in-memory ContextManager (no sandbox). + /// Execute via Scheduler (persists to DB + spawns worker), or fall back to + /// ContextManager-only if the scheduler isn't available yet. async fn execute_local( &self, title: &str, @@ -247,6 +264,38 @@ impl CreateJobTool { ctx: &JobContext, ) -> Result { let start = std::time::Instant::now(); + + // Use the scheduler if available — creates in ContextManager, persists + // to DB, transitions to InProgress, and spawns a worker. The new job + // runs independently with its own Worker and LLM context (not inheriting + // the parent conversation). MaxJobsExceeded is returned as error JSON + // so the LLM can report it to the user. + if let Some(ref slot) = self.scheduler_slot + && let Some(ref scheduler) = *slot.read().await + { + return match scheduler + .dispatch_job(&ctx.user_id, title, description, None) + .await + { + Ok(job_id) => { + let result = serde_json::json!({ + "job_id": job_id.to_string(), + "title": title, + "status": "in_progress", + "message": format!("Created and scheduled job '{}'", title) + }); + Ok(ToolOutput::success(result, start.elapsed())) + } + Err(e) => { + let result = serde_json::json!({ + "error": e.to_string() + }); + Ok(ToolOutput::success(result, start.elapsed())) + } + }; + } + + // Fallback: ContextManager-only (scheduler not yet initialized). match self .context_manager .create_job_for_user(&ctx.user_id, title, description) @@ -257,7 +306,7 @@ impl CreateJobTool { "job_id": job_id.to_string(), "title": title, "status": "pending", - "message": format!("Created job '{}'", title) + "message": format!("Created job '{}' (not scheduled — scheduler unavailable)", title) }); Ok(ToolOutput::success(result, start.elapsed())) } diff --git a/src/tools/builtin/mod.rs b/src/tools/builtin/mod.rs index 1092ae57..98037495 100644 --- a/src/tools/builtin/mod.rs +++ b/src/tools/builtin/mod.rs @@ -22,7 +22,7 @@ pub use file::{ApplyPatchTool, ListDirTool, ReadFileTool, WriteFileTool}; pub use http::HttpTool; pub use job::{ CancelJobTool, CreateJobTool, JobEventsTool, JobPromptTool, JobStatusTool, ListJobsTool, - PromptQueue, + PromptQueue, SchedulerSlot, }; pub use json::JsonTool; pub use memory::{MemoryReadTool, MemorySearchTool, MemoryTreeTool, MemoryWriteTool}; diff --git a/src/tools/registry.rs b/src/tools/registry.rs index 2ed639db..78b064c1 100644 --- a/src/tools/registry.rs +++ b/src/tools/registry.rs @@ -286,11 +286,13 @@ impl ToolRegistry { /// /// Job tools allow the LLM to create, list, check status, and cancel jobs. /// When sandbox deps are provided, `create_job` automatically delegates to - /// Docker containers. Otherwise it creates in-memory jobs via ContextManager. + /// Docker containers. Otherwise it dispatches via the Scheduler (which + /// persists to DB and spawns a worker). #[allow(clippy::too_many_arguments)] pub fn register_job_tools( &self, context_manager: Arc, + scheduler_slot: Option, job_manager: Option>, store: Option>, job_event_tx: Option< @@ -301,6 +303,9 @@ impl ToolRegistry { secrets_store: Option>, ) { let mut create_tool = CreateJobTool::new(Arc::clone(&context_manager)); + if let Some(slot) = scheduler_slot { + create_tool = create_tool.with_scheduler_slot(slot); + } if let Some(jm) = job_manager { create_tool = create_tool.with_sandbox(jm, store.clone()); }