Skip to content
Merged
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
184 changes: 167 additions & 17 deletions src/crates/assembly/core/src/agentic/coordination/coordinator.rs
Original file line number Diff line number Diff line change
Expand Up @@ -4957,15 +4957,22 @@ Update the persona files and delete BOOTSTRAP.md as soon as bootstrap is complet
session_id: &str,
workspace_path: &Path,
) -> OpenBitFunResult<()> {
let _goal_guard = self.lock_thread_goal_operation(session_id).await;
let goal_guard = self.lock_thread_goal_operation(session_id).await;
let storage_path = self
.resolve_thread_goal_storage_path(session_id, workspace_path)
.await?;
// Dropping the goal drops the work it started: capture its turn before the
// clear releases the goal-to-turn binding.
let goal_turn = self.goal_driven_turn_id(session_id);
self.thread_goal_store()
.clear_thread_goal(session_id, storage_path.as_path())
.await?;
self.thread_goal_runtime(session_id).clear_active_goal(None);
self.emit_thread_goal_updated(session_id, None).await;
drop(goal_guard);
if let Some(turn_id) = goal_turn {
self.stop_goal_driven_turn(session_id, &turn_id).await;
}
Ok(())
}

Expand Down Expand Up @@ -5005,19 +5012,18 @@ Update the persona files and delete BOOTSTRAP.md as soon as bootstrap is complet
"cannot edit goal for session {session_id}: no goal exists"
))
})?;
let status = match existing.status {
ThreadGoalStatus::BudgetLimited | ThreadGoalStatus::Complete => {
Some(ThreadGoalStatus::Active)
}
_ => None,
};
// Editing the objective is an explicit goal action, not a status edit: the
// user is telling the goal what to work on next, so it runs again instead
// of staying parked on the status the edit was made from. A goal that was
// paused, blocked or over quota therefore comes back as active.
let resuming = thread_goal_status_is_resumable(existing.status);
let result = self
.thread_goal_store()
.set_thread_goal(
session_id,
storage_path.as_path(),
Some(objective),
status,
Some(ThreadGoalStatus::Active),
None,
false,
)
Expand All @@ -5029,9 +5035,24 @@ Update the persona files and delete BOOTSTRAP.md as soon as bootstrap is complet
self.emit_thread_goal_updated(session_id, Some(result.goal.clone()))
.await;
drop(goal_guard);
if objective_changed && result.goal.is_active() {
self.apply_objective_updated_steering(session_id, &result.goal)
.await;
// Saving the same text through the edit dialog is still the user asking for
// this goal to run, so a restarting edit delivers its steering either way.
if result.goal.is_active() && (objective_changed || resuming) {
if resuming {
// The edit restarts the goal, so it has to retire what stopped it:
// until the interrupted turn that paused the goal is abandoned, the
// objective steering below is parked and the active goal never moves.
self.release_superseded_interrupted_turn(session_id).await;
clear_thread_goal_continuation_abort(session_id);
}
// The restart is handed over in the background like resume: an idle
// session makes this steering admit the goal's next turn, and the user is
// waiting on the edit dialog, not on that turn being admitted.
self.schedule_thread_goal_steering(
session_id,
&result.goal,
AgentThreadGoalDeliveryKind::ObjectiveUpdated,
);
}
Ok(result.goal)
}
Expand Down Expand Up @@ -5194,6 +5215,13 @@ Update the persona files and delete BOOTSTRAP.md as soon as bootstrap is complet
&& previous
.as_ref()
.is_some_and(|goal| thread_goal_status_is_resumable(goal.status));
// Pausing has to stop the goal's turn. The status write below drops the
// goal-to-turn binding, so capture the turn it currently owns first.
let paused_goal_turn = if status == ThreadGoalStatus::Paused {
self.goal_driven_turn_id(session_id)
} else {
None
};
let result = self
.thread_goal_store()
.set_thread_goal(
Expand All @@ -5213,9 +5241,19 @@ Update the persona files and delete BOOTSTRAP.md as soon as bootstrap is complet
self.emit_thread_goal_updated(session_id, Some(result.goal.clone()))
.await;
drop(goal_guard);
if let Some(turn_id) = paused_goal_turn {
self.stop_goal_driven_turn(session_id, &turn_id).await;
}
if resuming && result.goal.is_active() {
// Resuming is explicit too: the interruption that paused the goal must
// not keep parking the steering that starts its next turn.
self.release_superseded_interrupted_turn(session_id).await;
clear_thread_goal_continuation_abort(session_id);
self.schedule_thread_goal_resumed_steering(session_id, &result.goal);
self.schedule_thread_goal_steering(
session_id,
&result.goal,
AgentThreadGoalDeliveryKind::Resumed,
);
}
Ok(result.goal)
}
Expand Down Expand Up @@ -5258,10 +5296,117 @@ Update the persona files and delete BOOTSTRAP.md as soon as bootstrap is complet
}
}

fn schedule_thread_goal_resumed_steering(&self, session_id: &str, goal: &ThreadGoal) {
/// The in-flight turn that is working toward the session's thread goal.
///
/// The goal runtime only binds a goal to a turn while that goal is active, so
/// a bound turn that is also the session's current turn is exactly the work
/// the goal owns. Everything else in the session belongs to the user.
fn goal_driven_turn_id(&self, session_id: &str) -> Option<String> {
let (turn_id, _) = self.thread_goal_runtime(session_id).current_turn_usage()?;
if turn_id.is_empty() {
return None;
}
let session = self.session_manager.get_session(session_id)?;
match session.state {
SessionState::Processing {
current_turn_id, ..
} if current_turn_id == turn_id => Some(turn_id),
_ => None,
}
}

/// Stop the turn that is working toward a thread goal.
///
/// Pausing or dropping a goal has to stop the work it started: the goal's own
/// turn keeps issuing rounds on the model and calling tools long after the
/// goal stopped being active, which reads as a paused or deleted goal that
/// keeps running. Only the goal's turn is stopped; an unrelated user turn in
/// the same session keeps running.
async fn stop_goal_driven_turn(&self, session_id: &str, turn_id: &str) {
info!(
"Stopping thread goal turn: session_id={}, turn_id={}",
session_id, turn_id
);
// `Cancelled` is the disposition for work the user stopped. It also ends
// the recoverable-interruption fence a goal turn would otherwise leave
// behind, which is what lets the next explicit goal activation start.
//
// Stopping a turn can pause the goal it belonged to, so this call reaches
// back into the goal status write. That cycle is real but shallow, so the
// future is boxed instead of letting the async state machine recurse.
let cancellation = Box::pin(self.cancel_dialog_turn(session_id, turn_id)).await;
if let Err(error) = cancellation {
warn!(
"Failed to stop thread goal turn: session_id={}, turn_id={}, error={}",
session_id, turn_id, error
);
}
}

/// Drop a dispatch hold an interrupted turn left on the session.
///
/// An interrupted turn parks every submission that is not a user submission
/// until the user recovers or sends a new turn. A thread goal never submits a
/// user turn, so an explicit goal activation has to retire that fence first;
/// otherwise the goal is active with a kickoff that is parked forever.
async fn release_superseded_interrupted_turn(&self, session_id: &str) {
let holds_dispatch = match self
.session_manager
.latest_dialog_turn_holds_dispatch(session_id)
.await
{
Ok(holds_dispatch) => holds_dispatch,
Err(error) => {
warn!(
"Failed to inspect the interrupted turn fence for a thread goal: session_id={}, error={}",
session_id, error
);
return;
}
};
if !holds_dispatch {
return;
}
let interrupted_turn_id = self
.session_manager
.get_session(session_id)
.and_then(|session| session.dialog_turn_ids.last().cloned());
match self
.session_manager
.abandon_interrupted_dialog_turn(session_id, interrupted_turn_id.as_deref())
.await
{
Ok(Some(abandoned_turn_id)) => info!(
"Retired interrupted turn for thread goal activation: session_id={}, turn_id={}",
session_id, abandoned_turn_id
),
Ok(None) => {}
Err(error) => warn!(
"Failed to retire interrupted turn for thread goal activation: session_id={}, error={}",
session_id, error
),
}
}

/// Hand goal steering to the runtime without waiting for the delivery to finish.
///
/// The runtime injects steering into a running turn, but an idle session makes
/// it admit the goal's next turn, which costs seconds on a large workspace. The
/// caller is an interactive command that already persisted the goal state, so it
/// must not block on turn admission; the delivery outcome is reported by log only.
fn schedule_thread_goal_steering(
&self,
session_id: &str,
goal: &ThreadGoal,
kind: AgentThreadGoalDeliveryKind,
) {
if !goal.is_active() {
return;
}
let kind_label = match kind {
AgentThreadGoalDeliveryKind::Resumed => "resumed",
AgentThreadGoalDeliveryKind::ObjectiveUpdated => "objective_updated",
};
let agent_type = match self.session_manager.get_session(session_id) {
Some(session) => {
let agent_type = session.agent_type.trim();
Expand Down Expand Up @@ -5295,8 +5440,8 @@ Update the persona files and delete BOOTSTRAP.md as soon as bootstrap is complet
Ok(runtime) => runtime,
Err(error) => {
warn!(
"Agent runtime lifecycle delivery is not available; thread goal resume steering skipped: session_id={}, error={}",
session_id, error
"Agent runtime lifecycle delivery is not available; thread goal steering skipped: session_id={}, kind={}, error={}",
session_id, kind_label, error
);
return;
}
Expand All @@ -5308,14 +5453,15 @@ Update the persona files and delete BOOTSTRAP.md as soon as bootstrap is complet
workspace_path,
remote_connection_id,
remote_ssh_host,
kind: AgentThreadGoalDeliveryKind::Resumed,
kind,
goal,
})
.await
{
warn!(
"Failed to deliver thread goal resume steering: session_id={}, error={}",
"Failed to deliver thread goal steering: session_id={}, kind={}, error={}",
session_id,
kind_label,
CoreServiceAgentRuntime::runtime_error_message(error)
);
}
Expand Down Expand Up @@ -5431,6 +5577,10 @@ Update the persona files and delete BOOTSTRAP.md as soon as bootstrap is complet
.get_thread_goal(&session_id, storage_path.as_path())
.await?;
let replace_existing = existing.is_some();
// An explicit goal supersedes a recoverable interruption: the objective
// steering below is not a user turn, so as long as that interrupted turn
// fences dispatch the goal would stay active with a parked kickoff.
self.release_superseded_interrupted_turn(&session_id).await;
let goal = self
.set_thread_goal_objective(
&session_id,
Expand Down
46 changes: 46 additions & 0 deletions src/crates/assembly/core/src/agentic/coordination/scheduler.rs
Original file line number Diff line number Diff line change
Expand Up @@ -5122,6 +5122,52 @@ mod tests {
assert_eq!(paused.status, ThreadGoalStatus::Paused);
}

#[tokio::test]
async fn editing_a_paused_goal_puts_it_back_to_work() {
let (scheduler, sessions, _, root) = test_scheduler_with_persistence(true);
let session = "goal-edit";
mark_session_processing(&sessions, &root, session, "goal-turn").await;
let storage = sessions
.effective_session_storage_path(session)
.await
.unwrap();
scheduler
.coordinator
.create_thread_goal(session, &storage, "finish the work".into(), None)
.await
.unwrap();

// Saving the objective the goal already carries is a restart too: the user
// opened the edit dialog to put the goal back to work, so it must not come
// back as active-with-nothing-running.
scheduler
.coordinator
.set_thread_goal_status(session, &storage, ThreadGoalStatus::Paused)
.await
.unwrap();
let unchanged = scheduler
.coordinator
.update_thread_goal_objective(session, &storage, "finish the work".into())
.await
.unwrap();
assert_eq!(unchanged.status, ThreadGoalStatus::Active);

// A rewritten objective keeps the goal running instead of inheriting the
// pause it was edited from.
scheduler
.coordinator
.set_thread_goal_status(session, &storage, ThreadGoalStatus::Paused)
.await
.unwrap();
let edited = scheduler
.coordinator
.update_thread_goal_objective(session, &storage, "finish the rest".into())
.await
.unwrap();
assert_eq!(edited.status, ThreadGoalStatus::Active);
assert_eq!(edited.objective, "finish the rest");
}

#[tokio::test]
async fn thread_goal_plain_prompt_steering_activates_without_an_extra_turn() {
let (scheduler, session_manager, _, root) = test_scheduler_with_persistence(true);
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -49,7 +49,7 @@ export const chatInputAppearanceDescriptor: AppearanceSurfaceDescriptor = {
{ id: 'command', attribute: 'data-openbitfun-command', values: ['actions', 'all', 'skills', 'modes'] },
{ id: 'commandItemKind', attribute: 'data-openbitfun-command-item-kind', values: ['action', 'mode', 'skill', 'mcp'] },
{ id: 'action', attribute: 'data-openbitfun-action', values: ['cancel', 'continue-interrupted', 'retry', 'send', 'split'] },
{ id: 'boostItemKind', attribute: 'data-openbitfun-boost-item-kind', values: ['directive', 'agent', 'context', 'skill', 'additional-mode', 'manage'] },
{ id: 'boostItemKind', attribute: 'data-openbitfun-boost-item-kind', values: ['directive', 'agent', 'context', 'skill', 'additional-mode', 'manage', 'thread-goal'] },
],
states: [
{ id: 'processing', selector: { kind: 'ancestorPart', part: 'root', suffix: '[data-openbitfun-state~="processing"]' } },
Expand Down
Loading
Loading