feat (acp): custom methods for managing scheduler (#9972)

This commit is contained in:
Lifei Zhou
2026-06-25 09:11:35 +10:00
committed by GitHub
parent 09998b12cf
commit eea698945c
28 changed files with 2050 additions and 484 deletions
@@ -6,6 +6,8 @@ use std::collections::HashMap;
mod recipe;
pub use recipe::*;
mod schedule;
pub use schedule::*;
/// Schema descriptor for a single custom method, produced by the
/// `#[custom_methods]` macro's generated `custom_method_schemas()` function.
@@ -0,0 +1,169 @@
use agent_client_protocol::schema::SessionInfo;
use agent_client_protocol::{JsonRpcRequest, JsonRpcResponse};
use schemars::JsonSchema;
use serde::{Deserialize, Serialize};
use super::{EmptyResponse, RecipeDto};
#[derive(Debug, Default, Clone, Serialize, Deserialize, JsonSchema)]
#[serde(rename_all = "camelCase")]
pub struct ScheduledJobDto {
pub id: String,
pub source: String,
pub cron: String,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub last_run: Option<String>,
pub currently_running: bool,
pub paused: bool,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub current_session_id: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub job_start_time: Option<String>,
}
#[derive(Debug, Default, Clone, Serialize, Deserialize, JsonSchema, JsonRpcRequest)]
#[request(
method = "_goose/unstable/schedules/list",
response = ListSchedulesResponse
)]
pub struct ListSchedulesRequest {}
#[derive(Debug, Default, Clone, Serialize, Deserialize, JsonSchema, JsonRpcResponse)]
pub struct ListSchedulesResponse {
pub jobs: Vec<ScheduledJobDto>,
}
#[derive(Debug, Default, Clone, Serialize, Deserialize, JsonSchema, JsonRpcRequest)]
#[request(
method = "_goose/unstable/schedules/create",
response = CreateScheduleResponse
)]
pub struct CreateScheduleRequest {
pub id: String,
pub recipe: RecipeDto,
pub cron: String,
}
#[derive(Debug, Default, Clone, Serialize, Deserialize, JsonSchema, JsonRpcResponse)]
pub struct CreateScheduleResponse {
pub job: ScheduledJobDto,
}
#[derive(Debug, Default, Clone, Serialize, Deserialize, JsonSchema, JsonRpcRequest)]
#[request(method = "_goose/unstable/schedules/delete", response = EmptyResponse)]
#[serde(rename_all = "camelCase")]
pub struct DeleteScheduleRequest {
pub schedule_id: String,
}
#[derive(Debug, Default, Clone, Serialize, Deserialize, JsonSchema, JsonRpcRequest)]
#[request(
method = "_goose/unstable/schedules/update",
response = UpdateScheduleResponse
)]
#[serde(rename_all = "camelCase")]
pub struct UpdateScheduleRequest {
pub schedule_id: String,
pub cron: String,
}
#[derive(Debug, Default, Clone, Serialize, Deserialize, JsonSchema, JsonRpcResponse)]
pub struct UpdateScheduleResponse {
pub job: ScheduledJobDto,
}
#[derive(Debug, Default, Clone, Serialize, Deserialize, JsonSchema, JsonRpcRequest)]
#[request(
method = "_goose/unstable/schedules/run-now",
response = RunScheduleNowResponse
)]
#[serde(rename_all = "camelCase")]
pub struct RunScheduleNowRequest {
pub schedule_id: String,
}
#[derive(Debug, Default, Clone, Serialize, Deserialize, JsonSchema, JsonRpcResponse)]
#[serde(rename_all = "camelCase")]
pub struct RunScheduleNowResponse {
pub status: RunScheduleNowStatus,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub session_id: Option<String>,
}
#[derive(Debug, Default, Clone, Serialize, Deserialize, JsonSchema)]
#[serde(rename_all = "camelCase")]
pub enum RunScheduleNowStatus {
#[default]
Completed,
Cancelled,
}
#[derive(Debug, Default, Clone, Serialize, Deserialize, JsonSchema, JsonRpcRequest)]
#[request(
method = "_goose/unstable/schedules/sessions/list",
response = ListScheduleSessionsResponse
)]
#[serde(rename_all = "camelCase")]
pub struct ListScheduleSessionsRequest {
pub schedule_id: String,
pub limit: usize,
}
#[derive(Debug, Default, Clone, Serialize, Deserialize, JsonSchema, JsonRpcResponse)]
pub struct ListScheduleSessionsResponse {
pub sessions: Vec<SessionInfo>,
}
#[derive(Debug, Default, Clone, Serialize, Deserialize, JsonSchema, JsonRpcRequest)]
#[request(method = "_goose/unstable/schedules/pause", response = EmptyResponse)]
#[serde(rename_all = "camelCase")]
pub struct PauseScheduleRequest {
pub schedule_id: String,
}
#[derive(Debug, Default, Clone, Serialize, Deserialize, JsonSchema, JsonRpcRequest)]
#[request(
method = "_goose/unstable/schedules/unpause",
response = EmptyResponse
)]
#[serde(rename_all = "camelCase")]
pub struct UnpauseScheduleRequest {
pub schedule_id: String,
}
#[derive(Debug, Default, Clone, Serialize, Deserialize, JsonSchema, JsonRpcRequest)]
#[request(
method = "_goose/unstable/schedules/running-job/kill",
response = KillRunningJobResponse
)]
#[serde(rename_all = "camelCase")]
pub struct KillRunningJobRequest {
pub job_id: String,
}
#[derive(Debug, Default, Clone, Serialize, Deserialize, JsonSchema, JsonRpcResponse)]
pub struct KillRunningJobResponse {
pub message: String,
}
#[derive(Debug, Default, Clone, Serialize, Deserialize, JsonSchema, JsonRpcRequest)]
#[request(
method = "_goose/unstable/schedules/running-job/inspect",
response = InspectRunningJobResponse
)]
#[serde(rename_all = "camelCase")]
pub struct InspectRunningJobRequest {
pub job_id: String,
}
#[derive(Debug, Default, Clone, Serialize, Deserialize, JsonSchema, JsonRpcResponse)]
#[serde(rename_all = "camelCase")]
pub struct InspectRunningJobResponse {
pub running: bool,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub session_id: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub job_start_time: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub running_duration_seconds: Option<i64>,
}
+50
View File
@@ -250,6 +250,56 @@
"requestType": "RecipeToYamlRequest_unstable",
"responseType": "RecipeToYamlResponse_unstable"
},
{
"method": "_goose/unstable/schedules/list",
"requestType": "ListSchedulesRequest_unstable",
"responseType": "ListSchedulesResponse_unstable"
},
{
"method": "_goose/unstable/schedules/sessions/list",
"requestType": "ListScheduleSessionsRequest_unstable",
"responseType": "ListScheduleSessionsResponse_unstable"
},
{
"method": "_goose/unstable/schedules/create",
"requestType": "CreateScheduleRequest_unstable",
"responseType": "CreateScheduleResponse_unstable"
},
{
"method": "_goose/unstable/schedules/delete",
"requestType": "DeleteScheduleRequest_unstable",
"responseType": "EmptyResponse"
},
{
"method": "_goose/unstable/schedules/pause",
"requestType": "PauseScheduleRequest_unstable",
"responseType": "EmptyResponse"
},
{
"method": "_goose/unstable/schedules/unpause",
"requestType": "UnpauseScheduleRequest_unstable",
"responseType": "EmptyResponse"
},
{
"method": "_goose/unstable/schedules/update",
"requestType": "UpdateScheduleRequest_unstable",
"responseType": "UpdateScheduleResponse_unstable"
},
{
"method": "_goose/unstable/schedules/run-now",
"requestType": "RunScheduleNowRequest_unstable",
"responseType": "RunScheduleNowResponse_unstable"
},
{
"method": "_goose/unstable/schedules/running-job/kill",
"requestType": "KillRunningJobRequest_unstable",
"responseType": "KillRunningJobResponse_unstable"
},
{
"method": "_goose/unstable/schedules/running-job/inspect",
"requestType": "InspectRunningJobRequest_unstable",
"responseType": "InspectRunningJobResponse_unstable"
},
{
"method": "_goose/unstable/session/info",
"requestType": "GetSessionInfoRequest_unstable",
+474 -16
View File
@@ -3520,32 +3520,105 @@
"x-side": "agent",
"x-method": "_goose/unstable/recipes/to-yaml"
},
"GetSessionInfoRequest_unstable": {
"ListSchedulesRequest_unstable": {
"type": "object",
"properties": {
"sessionId": {
"type": "string"
}
},
"required": [
"sessionId"
],
"description": "Return list-style metadata for a single session without loading the conversation.",
"x-side": "agent",
"x-method": "_goose/unstable/session/info"
"x-method": "_goose/unstable/schedules/list"
},
"GetSessionInfoResponse_unstable": {
"ListSchedulesResponse_unstable": {
"type": "object",
"properties": {
"session": {
"$ref": "#/$defs/SessionInfo"
"jobs": {
"type": "array",
"items": {
"$ref": "#/$defs/ScheduledJobDto"
}
}
},
"required": [
"session"
"jobs"
],
"x-side": "agent",
"x-method": "_goose/unstable/session/info"
"x-method": "_goose/unstable/schedules/list"
},
"ScheduledJobDto": {
"type": "object",
"properties": {
"id": {
"type": "string"
},
"source": {
"type": "string"
},
"cron": {
"type": "string"
},
"lastRun": {
"type": [
"string",
"null"
]
},
"currentlyRunning": {
"type": "boolean"
},
"paused": {
"type": "boolean"
},
"currentSessionId": {
"type": [
"string",
"null"
]
},
"jobStartTime": {
"type": [
"string",
"null"
]
}
},
"required": [
"id",
"source",
"cron",
"currentlyRunning",
"paused"
]
},
"ListScheduleSessionsRequest_unstable": {
"type": "object",
"properties": {
"scheduleId": {
"type": "string"
},
"limit": {
"type": "integer",
"minimum": 0
}
},
"required": [
"scheduleId",
"limit"
],
"x-side": "agent",
"x-method": "_goose/unstable/schedules/sessions/list"
},
"ListScheduleSessionsResponse_unstable": {
"type": "object",
"properties": {
"sessions": {
"type": "array",
"items": {
"$ref": "#/$defs/SessionInfo"
}
}
},
"required": [
"sessions"
],
"x-side": "agent",
"x-method": "_goose/unstable/schedules/sessions/list"
},
"SessionInfo": {
"type": "object",
@@ -3598,6 +3671,245 @@
"type": "string",
"description": "A unique identifier for a conversation session between a client and agent.\n\nSessions maintain their own context, conversation history, and state,\nallowing multiple independent interactions with the same agent.\n\nSee protocol docs: [Session ID](https://agentclientprotocol.com/protocol/session-setup#session-id)"
},
"CreateScheduleRequest_unstable": {
"type": "object",
"properties": {
"id": {
"type": "string"
},
"recipe": {
"$ref": "#/$defs/RecipeDto"
},
"cron": {
"type": "string"
}
},
"required": [
"id",
"recipe",
"cron"
],
"x-side": "agent",
"x-method": "_goose/unstable/schedules/create"
},
"CreateScheduleResponse_unstable": {
"type": "object",
"properties": {
"job": {
"$ref": "#/$defs/ScheduledJobDto"
}
},
"required": [
"job"
],
"x-side": "agent",
"x-method": "_goose/unstable/schedules/create"
},
"DeleteScheduleRequest_unstable": {
"type": "object",
"properties": {
"scheduleId": {
"type": "string"
}
},
"required": [
"scheduleId"
],
"x-side": "agent",
"x-method": "_goose/unstable/schedules/delete"
},
"PauseScheduleRequest_unstable": {
"type": "object",
"properties": {
"scheduleId": {
"type": "string"
}
},
"required": [
"scheduleId"
],
"x-side": "agent",
"x-method": "_goose/unstable/schedules/pause"
},
"UnpauseScheduleRequest_unstable": {
"type": "object",
"properties": {
"scheduleId": {
"type": "string"
}
},
"required": [
"scheduleId"
],
"x-side": "agent",
"x-method": "_goose/unstable/schedules/unpause"
},
"UpdateScheduleRequest_unstable": {
"type": "object",
"properties": {
"scheduleId": {
"type": "string"
},
"cron": {
"type": "string"
}
},
"required": [
"scheduleId",
"cron"
],
"x-side": "agent",
"x-method": "_goose/unstable/schedules/update"
},
"UpdateScheduleResponse_unstable": {
"type": "object",
"properties": {
"job": {
"$ref": "#/$defs/ScheduledJobDto"
}
},
"required": [
"job"
],
"x-side": "agent",
"x-method": "_goose/unstable/schedules/update"
},
"RunScheduleNowRequest_unstable": {
"type": "object",
"properties": {
"scheduleId": {
"type": "string"
}
},
"required": [
"scheduleId"
],
"x-side": "agent",
"x-method": "_goose/unstable/schedules/run-now"
},
"RunScheduleNowResponse_unstable": {
"type": "object",
"properties": {
"status": {
"$ref": "#/$defs/RunScheduleNowStatus"
},
"sessionId": {
"type": [
"string",
"null"
]
}
},
"required": [
"status"
],
"x-side": "agent",
"x-method": "_goose/unstable/schedules/run-now"
},
"RunScheduleNowStatus": {
"type": "string",
"enum": [
"completed",
"cancelled"
]
},
"KillRunningJobRequest_unstable": {
"type": "object",
"properties": {
"jobId": {
"type": "string"
}
},
"required": [
"jobId"
],
"x-side": "agent",
"x-method": "_goose/unstable/schedules/running-job/kill"
},
"KillRunningJobResponse_unstable": {
"type": "object",
"properties": {
"message": {
"type": "string"
}
},
"required": [
"message"
],
"x-side": "agent",
"x-method": "_goose/unstable/schedules/running-job/kill"
},
"InspectRunningJobRequest_unstable": {
"type": "object",
"properties": {
"jobId": {
"type": "string"
}
},
"required": [
"jobId"
],
"x-side": "agent",
"x-method": "_goose/unstable/schedules/running-job/inspect"
},
"InspectRunningJobResponse_unstable": {
"type": "object",
"properties": {
"running": {
"type": "boolean"
},
"sessionId": {
"type": [
"string",
"null"
]
},
"jobStartTime": {
"type": [
"string",
"null"
]
},
"runningDurationSeconds": {
"type": [
"integer",
"null"
]
}
},
"required": [
"running"
],
"x-side": "agent",
"x-method": "_goose/unstable/schedules/running-job/inspect"
},
"GetSessionInfoRequest_unstable": {
"type": "object",
"properties": {
"sessionId": {
"type": "string"
}
},
"required": [
"sessionId"
],
"description": "Return list-style metadata for a single session without loading the conversation.",
"x-side": "agent",
"x-method": "_goose/unstable/session/info"
},
"GetSessionInfoResponse_unstable": {
"type": "object",
"properties": {
"session": {
"$ref": "#/$defs/SessionInfo"
}
},
"required": [
"session"
],
"x-side": "agent",
"x-method": "_goose/unstable/session/info"
},
"TruncateSessionConversationRequest_unstable": {
"type": "object",
"properties": {
@@ -5195,6 +5507,96 @@
"description": "Params for _goose/unstable/recipes/to-yaml",
"title": "RecipeToYamlRequest_unstable"
},
{
"allOf": [
{
"$ref": "#/$defs/ListSchedulesRequest_unstable"
}
],
"description": "Params for _goose/unstable/schedules/list",
"title": "ListSchedulesRequest_unstable"
},
{
"allOf": [
{
"$ref": "#/$defs/ListScheduleSessionsRequest_unstable"
}
],
"description": "Params for _goose/unstable/schedules/sessions/list",
"title": "ListScheduleSessionsRequest_unstable"
},
{
"allOf": [
{
"$ref": "#/$defs/CreateScheduleRequest_unstable"
}
],
"description": "Params for _goose/unstable/schedules/create",
"title": "CreateScheduleRequest_unstable"
},
{
"allOf": [
{
"$ref": "#/$defs/DeleteScheduleRequest_unstable"
}
],
"description": "Params for _goose/unstable/schedules/delete",
"title": "DeleteScheduleRequest_unstable"
},
{
"allOf": [
{
"$ref": "#/$defs/PauseScheduleRequest_unstable"
}
],
"description": "Params for _goose/unstable/schedules/pause",
"title": "PauseScheduleRequest_unstable"
},
{
"allOf": [
{
"$ref": "#/$defs/UnpauseScheduleRequest_unstable"
}
],
"description": "Params for _goose/unstable/schedules/unpause",
"title": "UnpauseScheduleRequest_unstable"
},
{
"allOf": [
{
"$ref": "#/$defs/UpdateScheduleRequest_unstable"
}
],
"description": "Params for _goose/unstable/schedules/update",
"title": "UpdateScheduleRequest_unstable"
},
{
"allOf": [
{
"$ref": "#/$defs/RunScheduleNowRequest_unstable"
}
],
"description": "Params for _goose/unstable/schedules/run-now",
"title": "RunScheduleNowRequest_unstable"
},
{
"allOf": [
{
"$ref": "#/$defs/KillRunningJobRequest_unstable"
}
],
"description": "Params for _goose/unstable/schedules/running-job/kill",
"title": "KillRunningJobRequest_unstable"
},
{
"allOf": [
{
"$ref": "#/$defs/InspectRunningJobRequest_unstable"
}
],
"description": "Params for _goose/unstable/schedules/running-job/inspect",
"title": "InspectRunningJobRequest_unstable"
},
{
"allOf": [
{
@@ -5721,6 +6123,62 @@
],
"title": "RecipeToYamlResponse_unstable"
},
{
"allOf": [
{
"$ref": "#/$defs/ListSchedulesResponse_unstable"
}
],
"title": "ListSchedulesResponse_unstable"
},
{
"allOf": [
{
"$ref": "#/$defs/ListScheduleSessionsResponse_unstable"
}
],
"title": "ListScheduleSessionsResponse_unstable"
},
{
"allOf": [
{
"$ref": "#/$defs/CreateScheduleResponse_unstable"
}
],
"title": "CreateScheduleResponse_unstable"
},
{
"allOf": [
{
"$ref": "#/$defs/UpdateScheduleResponse_unstable"
}
],
"title": "UpdateScheduleResponse_unstable"
},
{
"allOf": [
{
"$ref": "#/$defs/RunScheduleNowResponse_unstable"
}
],
"title": "RunScheduleNowResponse_unstable"
},
{
"allOf": [
{
"$ref": "#/$defs/KillRunningJobResponse_unstable"
}
],
"title": "KillRunningJobResponse_unstable"
},
{
"allOf": [
{
"$ref": "#/$defs/InspectRunningJobResponse_unstable"
}
],
"title": "InspectRunningJobResponse_unstable"
},
{
"allOf": [
{
+1
View File
@@ -99,6 +99,7 @@ mod onboarding;
mod providers;
mod recipe;
mod resources;
mod schedule;
mod slash_commands;
mod sources;
mod tool_notifications;
@@ -416,6 +416,86 @@ impl GooseAcpAgent {
self.on_recipe_to_yaml(req).await
}
#[custom_method(ListSchedulesRequest)]
async fn dispatch_list_schedules(
&self,
req: ListSchedulesRequest,
) -> Result<ListSchedulesResponse, agent_client_protocol::Error> {
self.on_list_schedules(req).await
}
#[custom_method(ListScheduleSessionsRequest)]
async fn dispatch_list_schedule_sessions(
&self,
req: ListScheduleSessionsRequest,
) -> Result<ListScheduleSessionsResponse, agent_client_protocol::Error> {
self.on_list_schedule_sessions(req).await
}
#[custom_method(CreateScheduleRequest)]
async fn dispatch_create_schedule(
&self,
req: CreateScheduleRequest,
) -> Result<CreateScheduleResponse, agent_client_protocol::Error> {
self.on_create_schedule(req).await
}
#[custom_method(DeleteScheduleRequest)]
async fn dispatch_delete_schedule(
&self,
req: DeleteScheduleRequest,
) -> Result<EmptyResponse, agent_client_protocol::Error> {
self.on_delete_schedule(req).await
}
#[custom_method(PauseScheduleRequest)]
async fn dispatch_pause_schedule(
&self,
req: PauseScheduleRequest,
) -> Result<EmptyResponse, agent_client_protocol::Error> {
self.on_pause_schedule(req).await
}
#[custom_method(UnpauseScheduleRequest)]
async fn dispatch_unpause_schedule(
&self,
req: UnpauseScheduleRequest,
) -> Result<EmptyResponse, agent_client_protocol::Error> {
self.on_unpause_schedule(req).await
}
#[custom_method(UpdateScheduleRequest)]
async fn dispatch_update_schedule(
&self,
req: UpdateScheduleRequest,
) -> Result<UpdateScheduleResponse, agent_client_protocol::Error> {
self.on_update_schedule(req).await
}
#[custom_method(RunScheduleNowRequest)]
async fn dispatch_run_schedule_now(
&self,
req: RunScheduleNowRequest,
) -> Result<RunScheduleNowResponse, agent_client_protocol::Error> {
self.on_run_schedule_now(req).await
}
#[custom_method(KillRunningJobRequest)]
async fn dispatch_kill_running_job(
&self,
req: KillRunningJobRequest,
) -> Result<KillRunningJobResponse, agent_client_protocol::Error> {
self.on_kill_running_job(req).await
}
#[custom_method(InspectRunningJobRequest)]
async fn dispatch_inspect_running_job(
&self,
req: InspectRunningJobRequest,
) -> Result<InspectRunningJobResponse, agent_client_protocol::Error> {
self.on_inspect_running_job(req).await
}
#[custom_method(GetSessionInfoRequest)]
async fn dispatch_get_session_info(
&self,
+343
View File
@@ -0,0 +1,343 @@
use goose_sdk_types::custom_requests::{
CreateScheduleRequest, CreateScheduleResponse, DeleteScheduleRequest, EmptyResponse,
InspectRunningJobRequest, InspectRunningJobResponse, KillRunningJobRequest,
KillRunningJobResponse, ListScheduleSessionsRequest, ListScheduleSessionsResponse,
ListSchedulesRequest, ListSchedulesResponse, PauseScheduleRequest, RunScheduleNowRequest,
RunScheduleNowResponse, RunScheduleNowStatus, ScheduledJobDto, UnpauseScheduleRequest,
UpdateScheduleRequest, UpdateScheduleResponse,
};
use tokio::fs;
use super::{build_session_info, GooseAcpAgent, ResultExt};
use crate::recipe::validate_recipe::validate_recipe_template_from_content;
use crate::recipe::Recipe;
use crate::scheduler::{get_default_scheduled_recipes_dir, ScheduledJob, SchedulerError};
fn validate_schedule_id(id: &str) -> Result<(), agent_client_protocol::Error> {
let is_valid = !id.is_empty()
&& id
.chars()
.all(|c| c.is_ascii_alphanumeric() || c == '-' || c == '_' || c == ' ');
if !is_valid {
return Err(agent_client_protocol::Error::invalid_params().data(
"Schedule name must use only alphanumeric characters, hyphens, underscores, or spaces",
));
}
Ok(())
}
fn validate_schedule_recipe(recipe: &Recipe) -> Result<(), agent_client_protocol::Error> {
let recipe_yaml = recipe
.to_yaml()
.map_err(|e| agent_client_protocol::Error::invalid_params().data(e.to_string()))?;
validate_recipe_template_from_content(&recipe_yaml, None)
.map_err(|e| agent_client_protocol::Error::invalid_params().data(e.to_string()))?;
Ok(())
}
fn schedule_not_found_or_internal(error: SchedulerError) -> agent_client_protocol::Error {
match error {
SchedulerError::JobNotFound(id) => {
agent_client_protocol::Error::resource_not_found(Some(id))
}
error => agent_client_protocol::Error::internal_error().data(error.to_string()),
}
}
fn create_schedule_error(error: SchedulerError) -> agent_client_protocol::Error {
match error {
SchedulerError::CronParseError(message) => agent_client_protocol::Error::invalid_params()
.data(format!("Invalid cron expression: {message}")),
SchedulerError::RecipeLoadError(message) => agent_client_protocol::Error::invalid_params()
.data(format!("Recipe load error: {message}")),
SchedulerError::JobIdExists(id) => agent_client_protocol::Error::invalid_params()
.data(format!("Job ID already exists: {id}")),
error => agent_client_protocol::Error::internal_error()
.data(format!("Error creating schedule: {error}")),
}
}
fn schedule_state_error(error: SchedulerError) -> agent_client_protocol::Error {
match error {
SchedulerError::JobNotFound(id) => {
agent_client_protocol::Error::resource_not_found(Some(id))
}
SchedulerError::AnyhowError(error) => {
agent_client_protocol::Error::invalid_params().data(error.to_string())
}
error => agent_client_protocol::Error::internal_error().data(error.to_string()),
}
}
fn update_schedule_error(error: SchedulerError) -> agent_client_protocol::Error {
match error {
SchedulerError::JobNotFound(id) => {
agent_client_protocol::Error::resource_not_found(Some(id))
}
SchedulerError::AnyhowError(error) => {
agent_client_protocol::Error::invalid_params().data(error.to_string())
}
SchedulerError::CronParseError(message) => agent_client_protocol::Error::invalid_params()
.data(format!("Invalid cron expression: {message}")),
error => agent_client_protocol::Error::internal_error().data(error.to_string()),
}
}
fn run_schedule_now_error(
error: SchedulerError,
) -> Result<RunScheduleNowResponse, agent_client_protocol::Error> {
match error {
SchedulerError::JobNotFound(id) => {
Err(agent_client_protocol::Error::resource_not_found(Some(id)))
}
SchedulerError::AnyhowError(error)
if error.to_string().contains("was successfully cancelled") =>
{
Ok(RunScheduleNowResponse {
status: RunScheduleNowStatus::Cancelled,
session_id: None,
})
}
error => Err(agent_client_protocol::Error::internal_error()
.data(format!("Error running schedule: {error}"))),
}
}
fn scheduled_job_to_dto(job: ScheduledJob) -> ScheduledJobDto {
ScheduledJobDto {
id: job.id,
source: job.source,
cron: job.cron,
last_run: job.last_run.map(|value| value.to_rfc3339()),
currently_running: job.currently_running,
paused: job.paused,
current_session_id: job.current_session_id,
job_start_time: job.process_start_time.map(|value| value.to_rfc3339()),
}
}
impl GooseAcpAgent {
pub(super) async fn on_list_schedules(
&self,
_req: ListSchedulesRequest,
) -> Result<ListSchedulesResponse, agent_client_protocol::Error> {
let jobs = self
.agent_manager
.scheduler()
.list_scheduled_jobs()
.await
.into_iter()
.map(scheduled_job_to_dto)
.collect();
Ok(ListSchedulesResponse { jobs })
}
pub(super) async fn on_list_schedule_sessions(
&self,
req: ListScheduleSessionsRequest,
) -> Result<ListScheduleSessionsResponse, agent_client_protocol::Error> {
let sessions = self
.agent_manager
.scheduler()
.sessions(&req.schedule_id, req.limit)
.await
.internal_err_ctx("Failed to fetch schedule sessions")?
.into_iter()
.map(|(_, session)| build_session_info(session))
.collect();
Ok(ListScheduleSessionsResponse { sessions })
}
pub(super) async fn on_create_schedule(
&self,
req: CreateScheduleRequest,
) -> Result<CreateScheduleResponse, agent_client_protocol::Error> {
let id = req.id.trim().to_string();
validate_schedule_id(&id)?;
let recipe = Recipe::try_from(req.recipe).map_err(|e| {
agent_client_protocol::Error::invalid_params().data(format!("recipe: {e}"))
})?;
if recipe.check_for_security_warnings() {
return Err(agent_client_protocol::Error::invalid_params().data(
"This recipe contains hidden characters that could be malicious. Please remove them before trying to save.",
));
}
validate_schedule_recipe(&recipe)?;
let scheduled_recipes_dir = get_default_scheduled_recipes_dir().map_err(|e| {
agent_client_protocol::Error::internal_error()
.data(format!("Failed to get scheduled recipes directory: {e}"))
})?;
let recipe_path = scheduled_recipes_dir.join(format!("{id}.yaml"));
let yaml_content = recipe.to_yaml().map_err(|e| {
agent_client_protocol::Error::internal_error()
.data(format!("Failed to convert recipe to YAML: {e}"))
})?;
fs::write(&recipe_path, yaml_content).await.map_err(|e| {
agent_client_protocol::Error::internal_error()
.data(format!("Failed to save recipe file: {e}"))
})?;
let job = ScheduledJob {
id,
source: recipe_path.to_string_lossy().into_owned(),
cron: req.cron,
last_run: None,
currently_running: false,
paused: false,
current_session_id: None,
process_start_time: None,
parameters: vec![],
recipe_base_dir: None,
};
self.agent_manager
.scheduler()
.add_scheduled_job(job.clone(), false)
.await
.map_err(create_schedule_error)?;
Ok(CreateScheduleResponse {
job: scheduled_job_to_dto(job),
})
}
pub(super) async fn on_delete_schedule(
&self,
req: DeleteScheduleRequest,
) -> Result<EmptyResponse, agent_client_protocol::Error> {
self.agent_manager
.scheduler()
.remove_scheduled_job(&req.schedule_id, false)
.await
.map_err(schedule_not_found_or_internal)?;
Ok(EmptyResponse {})
}
pub(super) async fn on_pause_schedule(
&self,
req: PauseScheduleRequest,
) -> Result<EmptyResponse, agent_client_protocol::Error> {
self.agent_manager
.scheduler()
.pause_schedule(&req.schedule_id)
.await
.map_err(schedule_state_error)?;
Ok(EmptyResponse {})
}
pub(super) async fn on_unpause_schedule(
&self,
req: UnpauseScheduleRequest,
) -> Result<EmptyResponse, agent_client_protocol::Error> {
self.agent_manager
.scheduler()
.unpause_schedule(&req.schedule_id)
.await
.map_err(schedule_not_found_or_internal)?;
Ok(EmptyResponse {})
}
pub(super) async fn on_update_schedule(
&self,
req: UpdateScheduleRequest,
) -> Result<UpdateScheduleResponse, agent_client_protocol::Error> {
let schedule_id = req.schedule_id;
let cron = req.cron;
let scheduler = self.agent_manager.scheduler();
scheduler
.update_schedule(&schedule_id, cron)
.await
.map_err(update_schedule_error)?;
let job = scheduler
.list_scheduled_jobs()
.await
.into_iter()
.find(|job| job.id == schedule_id)
.ok_or_else(|| {
agent_client_protocol::Error::internal_error()
.data("Schedule not found after update")
})?;
Ok(UpdateScheduleResponse {
job: scheduled_job_to_dto(job),
})
}
pub(super) async fn on_run_schedule_now(
&self,
req: RunScheduleNowRequest,
) -> Result<RunScheduleNowResponse, agent_client_protocol::Error> {
match self
.agent_manager
.scheduler()
.run_now(&req.schedule_id)
.await
{
Ok(session_id) => Ok(RunScheduleNowResponse {
status: RunScheduleNowStatus::Completed,
session_id: Some(session_id),
}),
Err(error) => run_schedule_now_error(error),
}
}
pub(super) async fn on_kill_running_job(
&self,
req: KillRunningJobRequest,
) -> Result<KillRunningJobResponse, agent_client_protocol::Error> {
self.agent_manager
.scheduler()
.kill_running_job(&req.job_id)
.await
.map_err(schedule_state_error)?;
Ok(KillRunningJobResponse {
message: format!("Successfully killed running job '{}'", req.job_id),
})
}
pub(super) async fn on_inspect_running_job(
&self,
req: InspectRunningJobRequest,
) -> Result<InspectRunningJobResponse, agent_client_protocol::Error> {
let job = self
.agent_manager
.scheduler()
.list_scheduled_jobs()
.await
.into_iter()
.find(|job| job.id == req.job_id)
.ok_or_else(|| agent_client_protocol::Error::resource_not_found(Some(req.job_id)))?;
if !job.currently_running {
return Ok(InspectRunningJobResponse::default());
}
let running_duration_seconds = job.process_start_time.map(|start_time| {
chrono::Utc::now()
.signed_duration_since(start_time)
.num_seconds()
});
Ok(InspectRunningJobResponse {
running: true,
session_id: job.current_session_id,
job_start_time: job.process_start_time.map(|value| value.to_rfc3339()),
running_duration_seconds,
})
}
}
+156 -7
View File
@@ -139,6 +139,16 @@ async fn persist_jobs(
Ok(())
}
fn clear_running_state(job: &mut ScheduledJob) -> bool {
let changed = job.currently_running
|| job.current_session_id.is_some()
|| job.process_start_time.is_some();
job.currently_running = false;
job.current_session_id = None;
job.process_start_time = None;
changed
}
pub struct Scheduler {
tokio_scheduler: TokioJobScheduler,
jobs: Arc<Mutex<JobsMap>>,
@@ -432,7 +442,7 @@ impl Scheduler {
return;
}
let list: Vec<ScheduledJob> = match serde_json::from_str(&data) {
let mut list: Vec<ScheduledJob> = match serde_json::from_str(&data) {
Ok(jobs) => jobs,
Err(e) => {
tracing::error!(
@@ -444,6 +454,20 @@ impl Scheduler {
}
};
let reset_stale_running_state = list
.iter_mut()
.fold(false, |changed, job| clear_running_state(job) || changed);
if reset_stale_running_state {
match serde_json::to_string_pretty(&list) {
Ok(data) => {
if let Err(e) = fs::write(&self.storage_path, data) {
tracing::error!("Failed to persist scheduler startup state: {}", e);
}
}
Err(e) => tracing::error!("Failed to serialize scheduler startup state: {}", e),
}
}
for job_to_load in list {
if !Path::new(&job_to_load.source).exists() {
tracing::warn!(
@@ -554,12 +578,15 @@ impl Scheduler {
pub async fn list_scheduled_jobs(&self) -> Vec<ScheduledJob> {
self.sync_from_storage().await;
self.jobs
let mut jobs: Vec<ScheduledJob> = self
.jobs
.lock()
.await
.values()
.map(|(_, j)| j.clone())
.collect()
.collect();
jobs.sort_by(|a, b| a.id.cmp(&b.id));
jobs
}
pub async fn remove_scheduled_job(
@@ -648,6 +675,7 @@ impl Scheduler {
cancel_token.clone(),
)
.await;
let was_cancelled = cancel_token.is_cancelled();
{
let mut tasks = self.running_tasks.lock().await;
@@ -667,6 +695,10 @@ impl Scheduler {
persist_jobs(&self.storage_path, &self.jobs).await?;
match result {
_ if was_cancelled => Err(SchedulerError::AnyhowError(anyhow!(
"Job '{}' was successfully cancelled",
sched_id
))),
Ok(session_id) => Ok(session_id),
Err(e) => Err(SchedulerError::AnyhowError(anyhow!(
"Job '{}' failed: {}",
@@ -770,14 +802,25 @@ impl Scheduler {
}
}
let token = {
let mut tasks = self.running_tasks.lock().await;
tasks.remove(sched_id)
};
if let Some(token) = token {
token.cancel();
}
{
let tasks = self.running_tasks.lock().await;
if let Some(token) = tasks.get(sched_id) {
token.cancel();
let mut jobs_guard = self.jobs.lock().await;
match jobs_guard.get_mut(sched_id) {
Some((_, job)) => {
clear_running_state(job);
}
None => return Err(SchedulerError::JobNotFound(sched_id.to_string())),
}
}
Ok(())
persist_jobs(&self.storage_path, &self.jobs).await
}
pub async fn get_running_job_info(
@@ -1227,6 +1270,112 @@ mod tests {
);
}
#[tokio::test]
async fn test_kill_running_job_clears_state_and_persists() {
let temp_dir = tempdir().unwrap();
let storage_path = temp_dir.path().join("schedule.json");
let recipe_path = create_test_recipe(temp_dir.path(), "running_job");
let session_manager = Arc::new(SessionManager::new(temp_dir.path().to_path_buf()));
let scheduler = Scheduler::new(storage_path.clone(), session_manager)
.await
.unwrap();
let job = ScheduledJob {
id: "running_job".to_string(),
source: recipe_path.to_string_lossy().to_string(),
cron: "0 0 0 1 1 *".to_string(),
last_run: None,
currently_running: false,
paused: false,
current_session_id: None,
process_start_time: None,
parameters: vec![],
recipe_base_dir: None,
};
scheduler.add_scheduled_job(job, false).await.unwrap();
{
let mut jobs_guard = scheduler.jobs.lock().await;
let (_, job) = jobs_guard.get_mut("running_job").unwrap();
job.currently_running = true;
job.current_session_id = Some("session-id".to_string());
job.process_start_time = Some(Utc::now());
}
{
let mut tasks = scheduler.running_tasks.lock().await;
tasks.insert("running_job".to_string(), CancellationToken::new());
}
persist_jobs(&storage_path, &scheduler.jobs).await.unwrap();
scheduler.kill_running_job("running_job").await.unwrap();
let jobs = scheduler.list_scheduled_jobs().await;
let killed_job = jobs.iter().find(|job| job.id == "running_job").unwrap();
assert!(!killed_job.currently_running);
assert!(killed_job.current_session_id.is_none());
assert!(killed_job.process_start_time.is_none());
assert!(scheduler.running_tasks.lock().await.is_empty());
let persisted_jobs: Vec<ScheduledJob> =
serde_json::from_str(&fs::read_to_string(storage_path).unwrap()).unwrap();
let persisted_job = persisted_jobs
.iter()
.find(|job| job.id == "running_job")
.unwrap();
assert!(!persisted_job.currently_running);
assert!(persisted_job.current_session_id.is_none());
assert!(persisted_job.process_start_time.is_none());
}
#[tokio::test]
async fn test_load_jobs_from_storage_clears_stale_running_state() {
let temp_dir = tempdir().unwrap();
let storage_path = temp_dir.path().join("schedule.json");
let recipe_path = create_test_recipe(temp_dir.path(), "stale_running_job");
let started_at = Utc::now();
let stale_job = ScheduledJob {
id: "stale_running_job".to_string(),
source: recipe_path.to_string_lossy().to_string(),
cron: "0 0 0 1 1 *".to_string(),
last_run: None,
currently_running: true,
paused: false,
current_session_id: Some("stale-session-id".to_string()),
process_start_time: Some(started_at),
parameters: vec![],
recipe_base_dir: None,
};
fs::write(
&storage_path,
serde_json::to_string_pretty(&vec![stale_job]).unwrap(),
)
.unwrap();
let session_manager = Arc::new(SessionManager::new(temp_dir.path().to_path_buf()));
let scheduler = Scheduler::new(storage_path.clone(), session_manager)
.await
.unwrap();
let jobs = scheduler.list_scheduled_jobs().await;
let loaded_job = jobs
.iter()
.find(|job| job.id == "stale_running_job")
.unwrap();
assert!(!loaded_job.currently_running);
assert!(loaded_job.current_session_id.is_none());
assert!(loaded_job.process_start_time.is_none());
let persisted_jobs: Vec<ScheduledJob> =
serde_json::from_str(&fs::read_to_string(storage_path).unwrap()).unwrap();
let persisted_job = persisted_jobs
.iter()
.find(|job| job.id == "stale_running_job")
.unwrap();
assert!(!persisted_job.currently_running);
assert!(persisted_job.current_session_id.is_none());
assert!(persisted_job.process_start_time.is_none());
}
#[tokio::test]
async fn test_job_with_no_prompt_does_not_panic() {
let _guard = env_lock::lock_env([