Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
34 changes: 14 additions & 20 deletions crates/goose/src/agents/reply_parts.rs
Original file line number Diff line number Diff line change
Expand Up @@ -526,13 +526,10 @@ impl Agent {
let manager = self.config.session_manager.clone();
let session = manager.get_session(session_id, false).await?;

let accumulated_usage = session.accumulated_usage + usage.usage;

let accumulated_cost = session
let cost_delta = session
.provider_name
.as_deref()
.and_then(|pn| self.accumulate_cost(session.accumulated_cost, usage, pn))
.or(session.accumulated_cost);
.and_then(|pn| self.estimate_chunk_cost(usage, pn));

let current_usage = if is_compaction_usage {
// After compaction: summary output becomes new input context
Expand All @@ -542,30 +539,27 @@ impl Agent {
usage.usage
};

// Set the context window outright (single writer) while incrementing the
// accumulated totals, so a subagent rolling up concurrently is not
// clobbered - both in one statement.
manager
.update(session_id)
.schedule_id(schedule_id)
.usage(current_usage)
.accumulated_usage(accumulated_usage)
.accumulated_cost(accumulated_cost)
.apply()
.update_usage_metrics(
session_id,
schedule_id,
current_usage,
usage.usage,
cost_delta,
)
.await?;

Ok(())
}

fn accumulate_cost(
&self,
existing: Option<f64>,
usage: &ProviderUsage,
provider_name: &str,
) -> Option<f64> {
fn estimate_chunk_cost(&self, usage: &ProviderUsage, provider_name: &str) -> Option<f64> {
let canonical =
crate::providers::canonical::maybe_get_canonical_model(provider_name, &usage.model)?;

let chunk_cost = canonical.cost.estimate_cost(&usage.usage)?;

Some(existing.unwrap_or(0.0) + chunk_cost)
canonical.cost.estimate_cost(&usage.usage)
}
}

Expand Down
43 changes: 42 additions & 1 deletion crates/goose/src/agents/subagent_handler.rs
Original file line number Diff line number Diff line change
Expand Up @@ -47,7 +47,15 @@ pub struct SubagentRunParams {

pub async fn run_subagent_task(params: SubagentRunParams) -> Result<String, anyhow::Error> {
let return_last_only = params.return_last_only;
let (messages, final_output) = get_agent_messages(params).await.map_err(|e| {
let session_manager = params.config.session_manager.clone();
let parent_session_id = params.task_config.parent_session_id.clone();
let subagent_session_id = params.session_id.clone();

let result = get_agent_messages(params).await;

roll_up_usage_to_parent(&session_manager, &parent_session_id, &subagent_session_id).await;
Comment on lines +54 to +56

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

P2 Badge Make subagent usage roll-up abort-safe

Because the roll-up is sequenced after awaiting the entire subagent run, an async delegate that is cancelled and then force-aborted by handle_load_task_result after its 5-second grace period drops this future while it is still inside get_agent_messages(params).await, so execution never reaches the roll-up call. Any usage already persisted in the subagent session before that abort remains absent from the parent, preserving under-reporting for stuck/cancelled background delegates; move the roll-up to an abort-safe cleanup path or have the abort path add the persisted subagent totals.

Useful? React with 👍 / 👎.


let (messages, final_output) = result.map_err(|e| {
ErrorData::new(
ErrorCode::INTERNAL_ERROR,
format!("Failed to execute task: {}", e),
Expand All @@ -62,6 +70,39 @@ pub async fn run_subagent_task(params: SubagentRunParams) -> Result<String, anyh
Ok(extract_response_text(&messages, return_last_only))
}

/// Descendants already rolled up through this same path, so the subagent's
/// accumulated totals cover the whole subtree.
async fn roll_up_usage_to_parent(
session_manager: &crate::session::SessionManager,
parent_session_id: &str,
subagent_session_id: &str,
) {
if parent_session_id.is_empty() || parent_session_id == subagent_session_id {
return;
}
let Ok(subagent) = session_manager
.get_session(subagent_session_id, false)
.await
else {
return;
};
if let Err(e) = session_manager
.add_accumulated_usage(
parent_session_id,
subagent.accumulated_usage,
subagent.accumulated_cost,
)
.await
{
tracing::warn!(
"Failed to roll up subagent {} usage into parent {}: {}",
subagent_session_id,
parent_session_id,
e
);
}
}

fn extract_response_text(messages: &Conversation, return_last_only: bool) -> String {
if return_last_only {
messages
Expand Down
Loading
Loading