Skip to content

Commit 40e0348

Browse files
committed
fix(queue): protect new work from retiring turns and preserve draft recovery
1 parent eb4fb8e commit 40e0348

8 files changed

Lines changed: 538 additions & 28 deletions

File tree

‎src/crates/assembly/core/src/agentic/coordination/host_message_queue.rs‎

Lines changed: 71 additions & 10 deletions
Original file line numberDiff line numberDiff line change
@@ -28,6 +28,9 @@ struct Entry {
2828
fingerprint: String,
2929
view: DialogQueueItem,
3030
held: Option<QueuedTurn>,
31+
// User work accepted after this turn settled must outlive its delayed
32+
// outcome cleanup. This is host-local admission bookkeeping, not wire data.
33+
after_terminal_turn: Option<String>,
3134
}
3235
impl Default for QueueSession {
3336
fn default() -> Self {
@@ -54,6 +57,31 @@ impl HostQueueState {
5457
.get(session)
5558
.is_some_and(|s| s.entries.contains_key(turn))
5659
}
60+
pub(super) fn mark_after_terminal_turn(&mut self, session: &str, turn: &str, previous: String) {
61+
if let Some(e) = self
62+
.sessions
63+
.get_mut(session)
64+
.and_then(|s| s.entries.get_mut(turn))
65+
{
66+
e.after_terminal_turn = Some(previous);
67+
}
68+
}
69+
pub(super) fn admitted_after(
70+
&self,
71+
session: &str,
72+
previous: &str,
73+
) -> std::collections::HashSet<String> {
74+
self.sessions
75+
.get(session)
76+
.map(|s| {
77+
s.entries
78+
.iter()
79+
.filter(|(_, e)| e.after_terminal_turn.as_deref() == Some(previous))
80+
.map(|(id, _)| id.clone())
81+
.collect()
82+
})
83+
.unwrap_or_default()
84+
}
5785
pub(super) fn pending_held(&self, session: &str) -> usize {
5886
self.sessions.get(session).map_or(0, |s| {
5987
s.entries
@@ -178,14 +206,26 @@ impl HostQueueState {
178206
}
179207
impl DialogScheduler {
180208
pub(super) fn hold_managed_queue(&self, session: &str, reason: &str) {
209+
self.hold_managed_queue_for_outcome(session, reason, None);
210+
}
211+
pub(super) fn hold_managed_queue_for_outcome(
212+
&self,
213+
session: &str,
214+
reason: &str,
215+
previous: Option<&str>,
216+
) {
181217
let ids: Vec<String> = self
182218
.queue_state()
183219
.sessions
184220
.get(session)
185221
.map(|s| {
186222
s.entries
187223
.iter()
188-
.filter(|(_, e)| e.view.status == DialogQueueStatus::Queued)
224+
.filter(|(_, e)| {
225+
e.view.status == DialogQueueStatus::Queued
226+
&& previous
227+
.is_none_or(|id| e.after_terminal_turn.as_deref() != Some(id))
228+
})
189229
.map(|(id, _)| id.clone())
190230
.collect()
191231
})
@@ -346,6 +386,7 @@ impl DialogScheduler {
346386
fingerprint: digest,
347387
view,
348388
held: None,
389+
after_terminal_turn: None,
349390
},
350391
);
351392
s.order.push(message.turn_id.clone());
@@ -445,6 +486,7 @@ impl DialogScheduler {
445486
));
446487
}
447488
}
489+
let mut interrupted_turn_to_abandon = None;
448490
if let DialogQueueAction::Promote {
449491
expected_active_turn_id,
450492
..
@@ -460,13 +502,21 @@ impl DialogScheduler {
460502
{
461503
return Err(error("queue_conflict: target is no longer active"));
462504
}
463-
if self
464-
.session_manager
465-
.latest_dialog_turn_holds_dispatch(session)
466-
.await
467-
.map_err(|e| error(e.to_string()))?
505+
if expected_active_turn_id.is_none() && self.active_turns.contains(session) {
506+
return Err(error("queue_conflict: previous turn is still retiring"));
507+
}
508+
if expected_active_turn_id.is_none()
509+
&& self
510+
.session_manager
511+
.latest_dialog_turn_holds_dispatch(session)
512+
.await
513+
.map_err(|e| error(e.to_string()))?
468514
{
469-
return Err(error("Queue is blocked by interrupted turn recovery"));
515+
interrupted_turn_to_abandon = self
516+
.session_manager
517+
.get_session(session)
518+
.and_then(|s| s.dialog_turn_ids.last().cloned());
519+
self.hold_managed_queue(session, "Turn interrupted; retry this message explicitly");
470520
}
471521
}
472522
let held = {
@@ -556,11 +606,22 @@ impl DialogScheduler {
556606
DialogQueueAction::Promote {
557607
expected_active_turn_id: None,
558608
..
559-
} => {
560-
if let Err(e) = self.start_turn(session, &turn).await {
609+
} => match self.start_turn(session, &turn).await {
610+
Err(e) => {
561611
self.queue_state().hold(session, &turn, &e.to_string());
562612
}
563-
}
613+
Ok(_) => {
614+
if let Err(e) = self
615+
.abandon_superseded_interrupted_turn(
616+
session,
617+
interrupted_turn_to_abandon.as_deref(),
618+
)
619+
.await
620+
{
621+
warn!("Failed to retire interrupted recovery after explicit queue promotion: session_id={}, error={}", session, e);
622+
}
623+
}
624+
},
564625
_ => unreachable!(),
565626
}
566627
{

0 commit comments

Comments
 (0)