Skip to content

Commit c2030f2

Browse files
committed
fix(cp): free the delegation slot via RAII so task abort cannot leak it
Self-review round-2 finding. serve() released its inflight slot with a plain call after the execute().await — but the client ABORTS serving tasks that outlive the 5s drain window on disconnect, and an aborted task never runs code after its await point. The leaked entry then survives reconnection (executor state is process-lifetime), permanently consuming capacity: with the default max_delegated_sessions = 1, one abort turns the worker into a zombie that refuses every delegation while heartbeating happily. The window is realistic, not theoretical: the bounded (5s) session/cancel added in the previous fix means a wedged agent stdin makes the task overrun the equally-sized drain window by construction. release() now lives in a Drop guard taken at admission, so the abort's unwind frees the slot. Regression test aborts a serving task mid-run and asserts active() == 0; mutation-verified (plain-release pattern fails it 5/5, the guard passes).
1 parent a9906f8 commit c2030f2

1 file changed

Lines changed: 53 additions & 3 deletions

File tree

‎crates/openab-core/src/control_plane/executor.rs‎

Lines changed: 53 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -219,9 +219,16 @@ impl DelegationExecutor {
219219
return failed(&id, error);
220220
}
221221
};
222-
let outcome = self.execute(&forward, cancel).await;
223-
self.release(&id);
224-
outcome
222+
// RAII: the slot must free even if this task is ABORTED mid-await —
223+
// the client aborts serving tasks that outlive the drain window on
224+
// disconnect, and a plain post-await release would be skipped there,
225+
// leaking the inflight entry forever (with the default cap of 1, the
226+
// worker would refuse every delegation after reconnecting).
227+
let _slot = SlotGuard {
228+
executor: self.as_ref(),
229+
id: id.clone(),
230+
};
231+
self.execute(&forward, cancel).await
225232
}
226233

227234
async fn execute(
@@ -353,6 +360,19 @@ fn cap_result(text: String) -> String {
353360
out
354361
}
355362

363+
/// Frees a delegation's inflight slot on drop — including the drop that
364+
/// happens when the serving task is aborted at an await point.
365+
struct SlotGuard<'a> {
366+
executor: &'a DelegationExecutor,
367+
id: String,
368+
}
369+
370+
impl Drop for SlotGuard<'_> {
371+
fn drop(&mut self) {
372+
self.executor.release(&self.id);
373+
}
374+
}
375+
356376
fn failed(delegation_id: &str, error: impl Into<String>) -> DelegateResultParams {
357377
DelegateResultParams {
358378
delegation_id: delegation_id.to_string(),
@@ -863,4 +883,34 @@ mod tests {
863883
// Under the cap: untouched.
864884
assert_eq!(cap_result("small".into()), "small");
865885
}
886+
887+
#[tokio::test(start_paused = true)]
888+
async fn an_aborted_serving_task_still_frees_its_slot() {
889+
// The client aborts serving tasks that outlive the drain window on
890+
// disconnect. A plain post-await release would be skipped by the
891+
// abort, leaking the inflight entry: with the default cap of 1 the
892+
// worker would then refuse every delegation after reconnecting.
893+
let runner = Arc::new(FakeRunner {
894+
delay: Some(Duration::from_secs(300)),
895+
..Default::default()
896+
});
897+
let ex = executor(Arc::clone(&runner), 1);
898+
let task = {
899+
let ex = Arc::clone(&ex);
900+
tokio::spawn(async move { ex.serve(forward("d-abort", 600)).await })
901+
};
902+
while ex.active() < 1 {
903+
tokio::task::yield_now().await;
904+
}
905+
906+
task.abort();
907+
let _ = task.await; // JoinError::Cancelled — the abort landed
908+
909+
assert_eq!(ex.active(), 0, "abort must free the slot via the guard");
910+
// And the freed slot is genuinely reusable.
911+
let done = FakeRunner::completing("ok");
912+
let ex2 = executor(done, 1);
913+
let r = ex2.serve(forward("d-after", 600)).await;
914+
assert_eq!(r.status, DelegationStatus::Completed);
915+
}
866916
}

0 commit comments

Comments
 (0)