feat(summon): add peek mode for async background tasks (#9519)
Signed-off-by: Douwe M Osinga <douwe@sidewalklabs.com> Co-authored-by: Douwe M Osinga <douwe@sidewalklabs.com>
This commit is contained in:
@@ -517,6 +517,11 @@ impl SummonClient {
|
||||
"type": "boolean",
|
||||
"default": false,
|
||||
"description": "For running background tasks: cancel and return output."
|
||||
},
|
||||
"peek": {
|
||||
"type": "boolean",
|
||||
"default": false,
|
||||
"description": "For running background tasks: check progress without blocking. Returns turn count, idle time, and recent tool activity."
|
||||
}
|
||||
}
|
||||
});
|
||||
@@ -527,11 +532,13 @@ impl SummonClient {
|
||||
Call with no arguments to list all available sources (subrecipes, recipes, agents).\n\
|
||||
Call with a source name to load its content into your context.\n\
|
||||
For background tasks: load(source: \"task_id\") waits for the task and returns the result.\n\
|
||||
To cancel a running task: load(source: \"task_id\", cancel: true) stops and returns output.\n\n\
|
||||
To cancel a running task: load(source: \"task_id\", cancel: true) stops and returns output.\n\
|
||||
To check progress: load(source: \"task_id\", peek: true) returns status without blocking.\n\n\
|
||||
Examples:\n\
|
||||
- load() → Lists available sources\n\
|
||||
- load(source: \"deploy\") → Loads the deploy recipe\n\
|
||||
- load(source: \"20260219_1\") → Waits for background task, then returns result"
|
||||
- load(source: \"20260219_1\") → Waits for background task, then returns result\n\
|
||||
- load(source: \"20260219_1\", peek: true) → Check task progress without waiting"
|
||||
.to_string(),
|
||||
schema.as_object().unwrap().clone(),
|
||||
)
|
||||
@@ -815,6 +822,12 @@ impl SummonClient {
|
||||
.and_then(|v| v.as_bool())
|
||||
.unwrap_or(false);
|
||||
|
||||
let peek = arguments
|
||||
.as_ref()
|
||||
.and_then(|args| args.get("peek"))
|
||||
.and_then(|v| v.as_bool())
|
||||
.unwrap_or(false);
|
||||
|
||||
let working_dir = self.get_working_dir(session_id).await;
|
||||
|
||||
if source_name.is_none() {
|
||||
@@ -827,7 +840,7 @@ impl SummonClient {
|
||||
let name = source_name.unwrap();
|
||||
|
||||
if is_session_id(name) {
|
||||
let task_result = self.handle_load_task_result(name, cancel).await?;
|
||||
let task_result = self.handle_load_task_result(name, cancel, peek).await?;
|
||||
let mut meta = Meta::new();
|
||||
meta.0.insert(
|
||||
"subagent_session_id".to_string(),
|
||||
@@ -861,11 +874,32 @@ impl SummonClient {
|
||||
&self,
|
||||
task_id: &str,
|
||||
cancel: bool,
|
||||
peek: bool,
|
||||
) -> Result<TaskLoadResult, String> {
|
||||
let mut completed = self.completed_tasks.lock().await;
|
||||
|
||||
if let Some(task) = completed.remove(task_id) {
|
||||
let status_key = match &task.result {
|
||||
let completed_entry = if peek {
|
||||
completed.get(task_id).map(|task| {
|
||||
(
|
||||
task.result.clone(),
|
||||
task.description.clone(),
|
||||
task.duration,
|
||||
task.turns_taken,
|
||||
)
|
||||
})
|
||||
} else {
|
||||
completed.remove(task_id).map(|task| {
|
||||
(
|
||||
task.result,
|
||||
task.description,
|
||||
task.duration,
|
||||
task.turns_taken,
|
||||
)
|
||||
})
|
||||
};
|
||||
|
||||
if let Some((result, description, duration, turns_taken)) = completed_entry {
|
||||
let status_key = match &result {
|
||||
Ok(_) => "completed",
|
||||
Err(e) if e.starts_with("Task panicked:") => "panicked",
|
||||
Err(_) => "failed",
|
||||
@@ -875,7 +909,7 @@ impl SummonClient {
|
||||
"panicked" => "✗ Panicked",
|
||||
_ => "✗ Failed",
|
||||
};
|
||||
let output = match task.result {
|
||||
let output = match result {
|
||||
Ok(output) => output,
|
||||
Err(error) => format!("Error: {}", error),
|
||||
};
|
||||
@@ -887,15 +921,15 @@ impl SummonClient {
|
||||
**Duration:** {} ({} turns)\n\n\
|
||||
## Output\n\n{}",
|
||||
task_id,
|
||||
task.description,
|
||||
description,
|
||||
status,
|
||||
round_duration(task.duration),
|
||||
task.turns_taken,
|
||||
round_duration(duration),
|
||||
turns_taken,
|
||||
output
|
||||
))],
|
||||
status: status_key,
|
||||
turns: Some(task.turns_taken),
|
||||
duration_secs: Some(task.duration.as_secs()),
|
||||
turns: Some(turns_taken),
|
||||
duration_secs: Some(duration.as_secs()),
|
||||
});
|
||||
}
|
||||
|
||||
@@ -903,6 +937,40 @@ impl SummonClient {
|
||||
|
||||
let mut running = self.background_tasks.lock().await;
|
||||
if running.contains_key(task_id) {
|
||||
if peek {
|
||||
let task = running.get(task_id).unwrap();
|
||||
let elapsed = task.started_at.elapsed();
|
||||
let turns_taken = task.turns.load(Ordering::Relaxed);
|
||||
let now = current_epoch_millis();
|
||||
let idle_ms = now.saturating_sub(task.last_activity.load(Ordering::Relaxed));
|
||||
let description = task.description.clone();
|
||||
|
||||
let buffered_count = task.notification_buffer.lock().await.len();
|
||||
|
||||
drop(running);
|
||||
|
||||
let mut output = format!(
|
||||
"# Background Task Status: {}\n\n**Task:** {}\n**Status:** ⏳ Running\n**Elapsed:** {}\n**Turns taken:** {}\n**Idle:** {}\n**Buffered tool calls:** {}",
|
||||
task_id,
|
||||
description,
|
||||
round_duration(elapsed),
|
||||
turns_taken,
|
||||
round_duration(Duration::from_millis(idle_ms)),
|
||||
buffered_count,
|
||||
);
|
||||
|
||||
if buffered_count == 0 && turns_taken == 0 {
|
||||
output.push_str("\n\n_Task is initialising (no tool activity yet)._");
|
||||
}
|
||||
|
||||
return Ok(TaskLoadResult {
|
||||
content: vec![Content::text(output)],
|
||||
status: "running",
|
||||
turns: Some(turns_taken),
|
||||
duration_secs: Some(elapsed.as_secs()),
|
||||
});
|
||||
}
|
||||
|
||||
if cancel {
|
||||
let task = running.remove(task_id).unwrap();
|
||||
drop(running);
|
||||
@@ -2589,7 +2657,9 @@ You review code."#;
|
||||
let client = SummonClient::new(create_test_context()).unwrap();
|
||||
let temp_dir = TempDir::new().unwrap();
|
||||
|
||||
let result = client.handle_load_task_result("20260204_999", false).await;
|
||||
let result = client
|
||||
.handle_load_task_result("20260204_999", false, false)
|
||||
.await;
|
||||
assert!(result.is_err());
|
||||
assert!(result.unwrap_err().contains("not found"));
|
||||
|
||||
@@ -2631,7 +2701,7 @@ You review code."#;
|
||||
let mut subscriber = client.subscribe().await;
|
||||
|
||||
let result = client
|
||||
.handle_load_task_result("20260204_1", false)
|
||||
.handle_load_task_result("20260204_1", false, false)
|
||||
.await
|
||||
.expect("load should wait and return result");
|
||||
let text = extract_text(&result.content[0]);
|
||||
@@ -2693,7 +2763,7 @@ You review code."#;
|
||||
assert!(discovery_text.contains("20260204_3"));
|
||||
|
||||
let result = client
|
||||
.handle_load_task_result("20260204_2", false)
|
||||
.handle_load_task_result("20260204_2", false, false)
|
||||
.await
|
||||
.unwrap();
|
||||
let text = extract_text(&result.content[0]);
|
||||
@@ -2713,7 +2783,7 @@ You review code."#;
|
||||
.contains_key("20260204_2"));
|
||||
|
||||
let result = client
|
||||
.handle_load_task_result("20260204_3", false)
|
||||
.handle_load_task_result("20260204_3", false, false)
|
||||
.await
|
||||
.unwrap();
|
||||
let text = extract_text(&result.content[0]);
|
||||
@@ -2721,7 +2791,9 @@ You review code."#;
|
||||
assert!(text.contains("Error: Something went wrong"));
|
||||
assert_eq!(result.status, "failed");
|
||||
|
||||
let result = client.handle_load_task_result("20260204_3", false).await;
|
||||
let result = client
|
||||
.handle_load_task_result("20260204_3", false, false)
|
||||
.await;
|
||||
assert!(result.is_err());
|
||||
assert!(result.unwrap_err().contains("not found"));
|
||||
|
||||
@@ -2755,7 +2827,7 @@ You review code."#;
|
||||
}
|
||||
|
||||
let result = client
|
||||
.handle_load_task_result("20260204_1", true)
|
||||
.handle_load_task_result("20260204_1", true, false)
|
||||
.await
|
||||
.unwrap();
|
||||
let text = extract_text(&result.content[0]);
|
||||
@@ -2771,4 +2843,98 @@ You review code."#;
|
||||
.await
|
||||
.contains_key("20260204_1"));
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn test_peek_running_task() {
|
||||
let client = SummonClient::new(create_test_context()).unwrap();
|
||||
|
||||
{
|
||||
let mut running = client.background_tasks.lock().await;
|
||||
running.insert(
|
||||
"20260204_1".to_string(),
|
||||
BackgroundTask {
|
||||
id: "20260204_1".to_string(),
|
||||
description: "Long running analysis".to_string(),
|
||||
started_at: Instant::now(),
|
||||
turns: Arc::new(AtomicU32::new(7)),
|
||||
last_activity: Arc::new(AtomicU64::new(current_epoch_millis())),
|
||||
handle: tokio::spawn(async {
|
||||
tokio::time::sleep(Duration::from_secs(1000)).await;
|
||||
Ok("eventual result".to_string())
|
||||
}),
|
||||
cancellation_token: CancellationToken::new(),
|
||||
notification_buffer: Arc::new(Mutex::new(Vec::new())),
|
||||
},
|
||||
);
|
||||
}
|
||||
|
||||
// Peek should return status without removing the task
|
||||
let result = client
|
||||
.handle_load_task_result("20260204_1", false, true)
|
||||
.await
|
||||
.unwrap();
|
||||
let text = extract_text(&result.content[0]);
|
||||
assert!(text.contains("Running"));
|
||||
assert!(text.contains("Long running analysis"));
|
||||
assert!(text.contains("7")); // turns taken
|
||||
|
||||
// Task should still be in background_tasks (not consumed)
|
||||
assert!(client
|
||||
.background_tasks
|
||||
.lock()
|
||||
.await
|
||||
.contains_key("20260204_1"));
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn test_peek_nonexistent_task() {
|
||||
let client = SummonClient::new(create_test_context()).unwrap();
|
||||
|
||||
let result = client
|
||||
.handle_load_task_result("20260204_999", false, true)
|
||||
.await;
|
||||
assert!(result.is_err());
|
||||
assert!(result.unwrap_err().contains("not found"));
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn test_peek_completed_task_returns_result() {
|
||||
let client = SummonClient::new(create_test_context()).unwrap();
|
||||
|
||||
{
|
||||
let mut completed = client.completed_tasks.lock().await;
|
||||
completed.insert(
|
||||
"20260204_1".to_string(),
|
||||
CompletedTask {
|
||||
id: "20260204_1".to_string(),
|
||||
description: "Finished task".to_string(),
|
||||
result: Ok("final output".to_string()),
|
||||
turns_taken: 4,
|
||||
duration: Duration::from_secs(30),
|
||||
completed_at: Instant::now(),
|
||||
},
|
||||
);
|
||||
}
|
||||
|
||||
// Peek on a completed task should return the full result (same as non-peek)
|
||||
let result = client
|
||||
.handle_load_task_result("20260204_1", false, true)
|
||||
.await
|
||||
.unwrap();
|
||||
let text = extract_text(&result.content[0]);
|
||||
assert!(text.contains("Completed"));
|
||||
assert!(text.contains("final output"));
|
||||
|
||||
// Peek must be non-destructive: the result is still retrievable afterwards.
|
||||
assert!(client
|
||||
.completed_tasks
|
||||
.lock()
|
||||
.await
|
||||
.contains_key("20260204_1"));
|
||||
let result = client
|
||||
.handle_load_task_result("20260204_1", false, false)
|
||||
.await
|
||||
.unwrap();
|
||||
assert!(extract_text(&result.content[0]).contains("final output"));
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user