use super::*;
mod interrupted;
pub use interrupted::InterruptedToolCalls;
#[derive(Debug, Clone, PartialEq)]
#[non_exhaustive]
pub struct ParkedToolCalls {
pub turn_id: TurnId,
pub tool_calls: Vec<everruns_contracts::tool_types::ToolCall>,
}
pub(super) struct ParkedTurn {
pub(super) calls: ParkedToolCalls,
pub(super) resume: TurnState,
}
pub(super) type ParkedTurns = Arc<Mutex<std::collections::HashMap<SessionId, ParkedTurn>>>;
fn no_parked_turn(session_id: SessionId) -> AgentLoopError {
AgentLoopError::store(format!(
"session {session_id} has no turn waiting for tool results"
))
}
pub(super) fn lock_parked(
parked: &ParkedTurns,
) -> std::sync::MutexGuard<'_, std::collections::HashMap<SessionId, ParkedTurn>> {
parked
.lock()
.unwrap_or_else(|poisoned| poisoned.into_inner())
}
impl InProcessRuntime {
pub fn parked_tool_calls(&self, session_id: SessionId) -> Option<ParkedToolCalls> {
lock_parked(&self.parked_turns)
.get(&session_id)
.map(|parked| parked.calls.clone())
}
pub async fn resume_steerable_turn(
&self,
session_id: SessionId,
results: Vec<ToolCompletedData>,
steering: TurnSteering,
) -> Result<TurnResult> {
let resume = self
.deliver_parked_tool_results(session_id, results)
.await?;
let snapshot = self.resolved_execution_snapshot(session_id).await?;
let drive = TurnDrive {
session_id,
org_id: resume.org_id,
turn_id: resume.turn_id.unwrap_or_default(),
input_message_id: resume.input_message_id,
harness_id: snapshot.harness_id,
agent_id: snapshot.agent_id,
workspace_id: snapshot.workspace_id,
};
let plan = TurnPlan::ScheduleReason(resume.clone());
self.drive_turn_plan(drive, InProcessExecution::new(resume), plan, steering)
.await
}
pub async fn record_parked_tool_results(
&self,
session_id: SessionId,
results: Vec<ToolCompletedData>,
) -> Result<()> {
let (turn_id, input_message_id) = lock_parked(&self.parked_turns)
.get(&session_id)
.map(|parked| (parked.calls.turn_id, parked.resume.input_message_id))
.ok_or_else(|| no_parked_turn(session_id))?;
self.emit_tool_results(session_id, turn_id, input_message_id, results)
.await
}
async fn emit_tool_results(
&self,
session_id: SessionId,
turn_id: TurnId,
input_message_id: MessageId,
results: Vec<ToolCompletedData>,
) -> Result<()> {
for result in results {
self.event_emitter
.emit(EventRequest::new(
session_id,
EventContext::turn(turn_id, input_message_id),
result,
))
.await?;
}
Ok(())
}
#[doc(hidden)]
pub async fn deliver_parked_tool_results(
&self,
session_id: SessionId,
results: Vec<ToolCompletedData>,
) -> Result<TurnState> {
let parked = self
.take_parked_turn(session_id)
.ok_or_else(|| no_parked_turn(session_id))?;
let turn_id = parked.calls.turn_id;
let input_message_id = parked.resume.input_message_id;
self.emit_tool_results(session_id, turn_id, input_message_id, results)
.await?;
let mut resume = parked.resume;
resume.turn_id = Some(turn_id);
Ok(resume)
}
#[doc(hidden)]
pub fn park_turn(
&self,
session_id: SessionId,
turn_id: TurnId,
tool_calls: Vec<everruns_contracts::tool_types::ToolCall>,
resume: TurnState,
) {
lock_parked(&self.parked_turns).insert(
session_id,
ParkedTurn {
calls: ParkedToolCalls {
turn_id,
tool_calls,
},
resume,
},
);
}
#[doc(hidden)]
pub fn supersede_parked_turn(&self, session_id: SessionId) {
self.take_parked_turn(session_id);
}
pub(super) fn take_parked_turn(&self, session_id: SessionId) -> Option<ParkedTurn> {
lock_parked(&self.parked_turns).remove(&session_id)
}
}