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>>>;
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 parked = self.take_parked_turn(session_id).ok_or_else(|| {
AgentLoopError::store(format!(
"session {session_id} has no turn waiting for tool results"
))
})?;
let snapshot = self.resolved_execution_snapshot(session_id).await?;
let turn_id = parked.calls.turn_id;
let input_message_id = parked.resume.input_message_id;
for result in results {
self.event_emitter
.emit(EventRequest::new(
session_id,
EventContext::turn(turn_id, input_message_id),
result,
))
.await?;
}
let drive = TurnDrive {
session_id,
org_id: parked.resume.org_id,
turn_id,
input_message_id,
harness_id: snapshot.harness_id,
agent_id: snapshot.agent_id,
workspace_id: snapshot.workspace_id,
};
let plan = TurnPlan::ScheduleReason(parked.resume.clone());
self.drive_turn_plan(
drive,
InProcessExecution::new(parked.resume),
plan,
steering,
)
.await
}
pub(super) fn take_parked_turn(&self, session_id: SessionId) -> Option<ParkedTurn> {
lock_parked(&self.parked_turns).remove(&session_id)
}
}