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
109 changes: 93 additions & 16 deletions crates/libsy/src/algorithms/plan_execute.rs
Original file line number Diff line number Diff line change
Expand Up @@ -80,40 +80,45 @@ impl PlanExecute {
})
}

fn phase(&self, request: &Request) -> Phase {
fn phase(&self, request: &Request) -> Result<Phase> {
let signals = ToolSignals::from_request(request, None);
let mutation_seen = signals.edit_count > 0 || signals.write_count > 0;
let Some(identity) = RoutingIdentity::from_request(request) else {
return if mutation_seen {
return Ok(if mutation_seen {
Phase::Handoff
} else {
Phase::Plan
};
});
};
let is_session_final = request
.metadata
.as_ref()
.and_then(|metadata| metadata.session_final)
== Some(true);

let mut sessions = self.executing_sessions.lock();
let phase = if sessions.contains(&identity) {
Phase::Execute
} else if mutation_seen {
if sessions.len() >= MAX_EXECUTING_SESSIONS
&& let Some(evicted) = sessions.iter().next().cloned()
{
sessions.remove(&evicted);
if !is_session_final {
if sessions.len() >= MAX_EXECUTING_SESSIONS {
return Err(LibsyError::AlgorithmError {
message: format!(
"plan_execute reached its limit of {MAX_EXECUTING_SESSIONS} executing sessions. \
Finish an existing session with session_final before retrying the handoff"
),
});
}
sessions.insert(identity.clone());
}
sessions.insert(identity.clone());
Phase::Handoff
} else {
Phase::Plan
};
if request
.metadata
.as_ref()
.and_then(|metadata| metadata.session_final)
== Some(true)
{
if is_session_final {
sessions.remove(&identity);
}
phase
Ok(phase)
}

fn replay_planner_reasoning_as_text(request: &mut Request) -> usize {
Expand Down Expand Up @@ -166,7 +171,7 @@ impl Algorithm for PlanExecute {
driver: Driver,
mut request: Request,
) -> Result<RoutingOutcome> {
match self.phase(&request) {
match self.phase(&request)? {
Phase::Plan => {
prepend_system_prompt(&mut request, &self.config.planning_prompt);
tracing::debug!(phase = "plan", "plan-execute selected capable tier");
Expand Down Expand Up @@ -347,6 +352,78 @@ mod tests {
assert_eq!(selected, CAPABLE);
}

#[tokio::test]
async fn capacity_preserves_execution_after_compaction() {
let algorithm = algorithm(PlanExecuteConfig::default());
let sessions: Vec<_> = (0..MAX_EXECUTING_SESSIONS)
.map(|index| format!("task-{index}"))
.collect();
for session in &sessions {
let first_edit = request(
vec![tool_call("Write", json!({"file_path": "task.py"}))],
Some(session),
);
let (selected, _) = route_and_capture(Arc::clone(&algorithm), first_edit).await;
assert_eq!(selected, EFFICIENT);
}

let overflow = request(
vec![tool_call("Write", json!({"file_path": "task.py"}))],
Some("overflow"),
);
let (driver, _) = Driver::new("plan_execute", Arc::new(models()));
let result = Arc::clone(&algorithm).route(driver, overflow.clone()).await;
assert!(matches!(
result,
Err(LibsyError::AlgorithmError { message })
if message.contains("limit of 4096 executing sessions")
));

for session in &sessions {
let mut compacted = request(
vec![Message::text(Role::User, "Continue after compaction")],
Some(session),
);
if Some(session) == sessions.last() {
compacted.metadata.as_mut().unwrap().session_final = Some(true);
}
let (selected, routed) = route_and_capture(Arc::clone(&algorithm), compacted).await;
assert_eq!(selected, EFFICIENT, "session {session} lost its latch");
assert!(routed.llm_request.instructions.is_empty());
}

let (selected, _) = route_and_capture(Arc::clone(&algorithm), overflow).await;
assert_eq!(selected, EFFICIENT);
let compacted = request(
vec![Message::text(Role::User, "Continue after compaction")],
Some("overflow"),
);
let (selected, routed) = route_and_capture(algorithm, compacted).await;
assert_eq!(selected, EFFICIENT);
assert!(routed.llm_request.instructions.is_empty());
}

#[tokio::test]
async fn final_handoff_does_not_evict_an_executing_session_at_capacity() {
let algorithm = Arc::new(PlanExecute::new(PlanExecuteConfig::default()).unwrap());
algorithm.executing_sessions.lock().extend(
(0..MAX_EXECUTING_SESSIONS)
.map(|index| RoutingIdentity::Session(format!("task-{index}"))),
);
let mut final_request = request(
vec![tool_call("Write", json!({"file_path": "task.py"}))],
Some("final-handoff"),
);
final_request.metadata.as_mut().unwrap().session_final = Some(true);

let (selected, routed) = route_and_capture(algorithm.clone(), final_request).await;
assert_eq!(selected, EFFICIENT);
assert!(routed.llm_request.instructions.is_empty());
let sessions = algorithm.executing_sessions.lock();
assert_eq!(sessions.len(), MAX_EXECUTING_SESSIONS);
assert!(!sessions.contains(&RoutingIdentity::Session("final-handoff".to_string())));
}

#[tokio::test]
async fn mutation_without_a_session_uses_the_efficient_tier() {
let messages = vec![tool_call(
Expand Down
16 changes: 16 additions & 0 deletions docs/routing_algorithms/plan_execute_routing.md
Original file line number Diff line number Diff line change
Expand Up @@ -32,6 +32,22 @@ The first edit or write routes the full trajectory to the efficient target and
latches that choice by session ID. A failed edit still triggers the handoff.
Without a session ID, the first mutation must remain in the request history.

## Session capacity

Each router instance retains up to 4,096 executing identities in memory. Root
requests use the session ID. Subagent requests use the session ID and agent ID.
Once retained, an identity stays on the efficient target after history compaction
until a request marks it with `session_final: true` or the router restarts.

At capacity, a new identity's handoff returns an error before calling a model.
Existing identities keep their execution phase. Mark an existing identity's last
request with `session_final: true` to free its slot, then retry the handoff with
the mutation history intact. The server accepts this flag through the
`x-switchyard-session-final: true` header.

A handoff marked final needs no saved slot and can still run at capacity.
Requests that are still planning or have no stable identity also use no slots.

## Responses API history requirement

Plan/execute needs the conversation history to detect edits and hand the task to
Expand Down
Loading