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
22 changes: 13 additions & 9 deletions src-tauri/src/router/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -139,7 +139,7 @@ impl StdinInjector for SessionManager {
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum DeliveryReservation {
Ready(u64),
PendingInput,
PendingInput(Duration),
RecentlyTyping(Duration),
InFlight,
}
Expand Down Expand Up @@ -698,7 +698,7 @@ impl Router {
}
(session_id, None)
}
DeliveryReservation::PendingInput => {
DeliveryReservation::PendingInput(delay) => {
if reconciliation.is_some() {
return Ok(false);
}
Expand All @@ -710,6 +710,7 @@ impl Router {
outbox.handle = handle.to_string();
outbox.pending_input_blocked = true;
outbox.enqueue(delivery);
retry = Some((session_id.clone(), delay));
self.blocked_transition(&mut state, &session_id);
}
(session_id, None)
Expand Down Expand Up @@ -1035,13 +1036,16 @@ impl Router {
}
self.schedule_outbox_retry(session_id.to_string(), delay);
}
DeliveryReservation::PendingInput => {
let mut state = self.state.lock().unwrap();
let Some(outbox) = state.outbox_by_session.get_mut(session_id) else {
return;
};
outbox.pending_input_blocked = true;
self.blocked_transition(&mut state, session_id);
DeliveryReservation::PendingInput(delay) => {
{
let mut state = self.state.lock().unwrap();
let Some(outbox) = state.outbox_by_session.get_mut(session_id) else {
return;
};
outbox.pending_input_blocked = true;
self.blocked_transition(&mut state, session_id);
}
self.schedule_outbox_retry(session_id.to_string(), delay);
}
DeliveryReservation::InFlight => {
let mut state = self.state.lock().unwrap();
Expand Down
51 changes: 50 additions & 1 deletion src-tauri/src/router/tests.rs
Original file line number Diff line number Diff line change
Expand Up @@ -48,6 +48,7 @@ struct RecordingInjector {
#[derive(Default)]
struct RecordingInputState {
pending: bool,
pending_until: Option<Instant>,
last_input_at: Option<Instant>,
in_flight: bool,
generation: u64,
Expand Down Expand Up @@ -122,16 +123,22 @@ impl RecordingInjector {
}

fn set_pending(&self, session_id: &str) {
self.set_pending_for(session_id, Duration::from_secs(10 * 60));
}

fn set_pending_for(&self, session_id: &str, duration: Duration) {
let mut input = self.input.lock().unwrap();
let input = input.entry(session_id.to_string()).or_default();
input.pending = true;
input.pending_until = Some(Instant::now() + duration);
input.last_input_at = Some(Instant::now());
}

fn set_recent_typing(&self, session_id: &str) {
let mut input = self.input.lock().unwrap();
let input = input.entry(session_id.to_string()).or_default();
input.pending = false;
input.pending_until = None;
input.last_input_at = Some(Instant::now() - Duration::from_millis(1950));
}

Expand All @@ -147,6 +154,7 @@ impl RecordingInjector {
fn clear_pending(&self, session_id: &str) {
if let Some(input) = self.input.lock().unwrap().get_mut(session_id) {
input.pending = false;
input.pending_until = None;
input.last_input_at = None;
}
self.notify(session_id, SessionDeliveryEvent::InputCleared);
Expand All @@ -160,6 +168,7 @@ impl RecordingInjector {
input.generation = input.generation.wrapping_add(1);
input.in_flight = false;
input.pending = false;
input.pending_until = None;
input.last_input_at = None;
}
self.notify(session_id, SessionDeliveryEvent::Respawned);
Expand All @@ -171,6 +180,7 @@ impl RecordingInjector {
input.generation = input.generation.wrapping_add(1);
input.in_flight = false;
input.pending = false;
input.pending_until = None;
input.last_input_at = None;
}
self.notify(session_id, SessionDeliveryEvent::Exited);
Expand Down Expand Up @@ -276,7 +286,16 @@ impl StdinInjector for RecordingInjector {
return Ok(DeliveryReservation::InFlight);
}
if input.pending {
return Ok(DeliveryReservation::PendingInput);
let remaining = input
.pending_until
.expect("pending test input must have an abandonment deadline")
.saturating_duration_since(Instant::now());
if !remaining.is_zero() {
return Ok(DeliveryReservation::PendingInput(remaining));
}
input.pending = false;
input.pending_until = None;
input.last_input_at = None;
}
if let Some(last) = input.last_input_at {
let elapsed = last.elapsed();
Expand Down Expand Up @@ -1124,6 +1143,36 @@ fn recent_typing_retries_after_quiet_window() {
});
}

#[test]
fn pending_input_retries_after_abandonment_deadline() {
let (router, injector, _log, _dir) = fixture(
vec![
slot_with_runner("lead", true),
slot_with_runner("impl", false),
],
&[("lead", "S-LEAD"), ("impl", "S-IMPL")],
);
injector.set_pending_for("S-IMPL", Duration::from_millis(25));
router
.inject_and_submit("impl", b"after abandonment")
.unwrap();
assert!(injector.pushes_for("S-IMPL").is_empty());
assert!(
router
.state
.lock()
.unwrap()
.outbox_by_session
.get("S-IMPL")
.is_some_and(|outbox| outbox.retry_scheduled),
"a draft-gated outbox must always have a scheduled retry"
);

wait_until(Duration::from_millis(300), || {
injector.submitted_bodies_for("S-IMPL") == ["after abandonment"]
});
}

#[test]
fn deferred_nudges_coalesce_while_relays_preserve_order() {
let (router, injector, log, _dir) = fixture(
Expand Down
4 changes: 2 additions & 2 deletions src-tauri/src/session/manager/lifecycle.rs
Original file line number Diff line number Diff line change
Expand Up @@ -73,7 +73,7 @@ impl SessionManager {
gate.ready.notify_all();
state.activity = None;
state.suppress_local_input_busy = false;
state.local_input_pending = false;
state.draft.clear();
state.last_local_input_at = None;
state.mission_status_sink = None;
state.completion_armed = false;
Expand Down Expand Up @@ -312,7 +312,7 @@ impl SessionManager {
state.handle = None;
state.activity = None;
state.suppress_local_input_busy = false;
state.local_input_pending = false;
state.draft.clear();
state.last_local_input_at = None;
state.mission_status_sink = None;
state.completion_armed = false;
Expand Down
Loading
Loading