Skip to content

Commit 11e504d

Browse files
committed
fix(deep-review): stabilize recovery and capacity waits
1 parent 7a192ea commit 11e504d

37 files changed

Lines changed: 1976 additions & 345 deletions

src/crates/core/src/agentic/deep_review/budget.rs

Lines changed: 10 additions & 37 deletions
Original file line numberDiff line numberDiff line change
@@ -556,32 +556,14 @@ impl DeepReviewBudgetTracker {
556556
parent_dialog_turn_id: &str,
557557
max_active_reviewers: usize,
558558
launch_batch: u64,
559-
packet_id: Option<&str>,
559+
_packet_id: Option<&str>,
560560
) -> Result<Option<DeepReviewActiveReviewerGuard<'a>>, DeepReviewPolicyViolation> {
561561
let now = Instant::now();
562562
let mut budget = self
563563
.turns
564564
.entry(parent_dialog_turn_id.to_string())
565565
.or_insert_with(|| DeepReviewTurnBudget::new(now));
566566

567-
if let Some((&earliest_active_batch, _)) =
568-
budget.active_reviewer_launch_batches.iter().next()
569-
{
570-
if earliest_active_batch < launch_batch {
571-
let packet_label = packet_id
572-
.map(str::trim)
573-
.filter(|value| !value.is_empty())
574-
.unwrap_or("unknown");
575-
return Err(DeepReviewPolicyViolation::new(
576-
"deep_review_launch_batch_blocked",
577-
format!(
578-
"Reviewer packet '{}' is in launch_batch {}, but launch_batch {} still has active reviewer(s). Wait for earlier-batch reviewers to finish, timeout, or be cancelled before launching this packet. If the queue remains blocked, pause or cancel the queued reviewers from the Review Team action bar and retry with a lower max parallel reviewer setting.",
579-
packet_label, launch_batch, earliest_active_batch
580-
),
581-
));
582-
}
583-
}
584-
585567
if budget.active_reviewers >= max_active_reviewers {
586568
return Ok(None);
587569
}
@@ -822,31 +804,22 @@ mod tests {
822804
use super::*;
823805

824806
#[test]
825-
fn launch_batch_admission_blocks_later_batch_while_earlier_batch_is_active() {
807+
fn launch_batch_admission_allows_later_batch_when_reviewer_capacity_is_free() {
826808
let tracker = DeepReviewBudgetTracker::default();
827-
let turn_id = "turn-launch-batch-blocked";
809+
let turn_id = "turn-launch-batch-fill-free-slot";
828810
let _first_batch = tracker
829811
.try_begin_active_reviewer_for_launch_batch(turn_id, 2, 1, Some("packet-a"))
830812
.expect("batch admission should not fail")
831813
.expect("first reviewer should start");
832814

833-
let violation = match tracker.try_begin_active_reviewer_for_launch_batch(
834-
turn_id,
835-
2,
836-
2,
837-
Some("packet-b"),
838-
) {
839-
Err(violation) => violation,
840-
Ok(_) => panic!("later launch batch should wait while earlier batch is active"),
841-
};
815+
let second_batch = tracker
816+
.try_begin_active_reviewer_for_launch_batch(turn_id, 2, 2, Some("packet-b"))
817+
.expect("later batch admission should not fail when reviewer capacity is free");
842818

843-
assert_eq!(violation.code, "deep_review_launch_batch_blocked");
844-
assert!(violation.message.contains("packet-b"));
845-
assert!(violation.message.contains("launch_batch 2"));
846-
assert!(violation.message.contains("launch_batch 1"));
847-
assert!(violation
848-
.message
849-
.contains("Wait for earlier-batch reviewers"));
819+
assert!(
820+
second_batch.is_some(),
821+
"later batch should fill a freed reviewer slot instead of waiting for the earlier batch to drain"
822+
);
850823
}
851824

852825
#[test]

src/crates/core/src/agentic/deep_review/task_adapter.rs

Lines changed: 75 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -37,6 +37,9 @@ use tokio::time::sleep;
3737
const DEEP_REVIEW_QUEUE_POLL_INTERVAL: Duration = Duration::from_millis(10);
3838
#[cfg(not(test))]
3939
const DEEP_REVIEW_QUEUE_POLL_INTERVAL: Duration = Duration::from_secs(1);
40+
pub(crate) const DEEP_REVIEW_PROVIDER_CAPACITY_MAX_RETRY_ATTEMPTS: usize = 3;
41+
const DEEP_REVIEW_PROVIDER_CAPACITY_BACKOFF_MULTIPLIER: u64 = 3;
42+
const DEEP_REVIEW_PROVIDER_CAPACITY_MAX_BACKOFF_SECONDS: u64 = 600;
4043

4144
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
4245
pub(crate) enum DeepReviewQueueWaitSkipReason {
@@ -59,6 +62,7 @@ pub(crate) enum DeepReviewQueueWaitOutcome {
5962
pub(crate) enum DeepReviewProviderQueueWaitOutcome {
6063
ReadyToRetry {
6164
queue_elapsed_ms: u64,
65+
early_capacity_probe: bool,
6266
},
6367
Skipped {
6468
queue_elapsed_ms: u64,
@@ -514,6 +518,39 @@ pub(crate) fn provider_capacity_queue_wait_seconds(
514518
.filter(|seconds| *seconds > 0)
515519
}
516520

521+
pub(crate) fn provider_capacity_queue_wait_seconds_for_attempt(
522+
decision: &DeepReviewCapacityQueueDecision,
523+
conc_policy: &DeepReviewConcurrencyPolicy,
524+
retry_attempt_index: usize,
525+
) -> Option<u64> {
526+
let base_wait_seconds = provider_capacity_queue_wait_seconds(decision, conc_policy)?;
527+
if decision.retry_after_seconds.is_some() {
528+
return Some(base_wait_seconds);
529+
}
530+
531+
let multiplier = DEEP_REVIEW_PROVIDER_CAPACITY_BACKOFF_MULTIPLIER.saturating_pow(
532+
u32::try_from(retry_attempt_index)
533+
.unwrap_or(u32::MAX)
534+
.min(8),
535+
);
536+
Some(
537+
base_wait_seconds
538+
.saturating_mul(multiplier)
539+
.min(DEEP_REVIEW_PROVIDER_CAPACITY_MAX_BACKOFF_SECONDS),
540+
)
541+
.filter(|seconds| *seconds > 0)
542+
}
543+
544+
fn provider_capacity_wait_can_wake_on_active_reviewer_release(
545+
reason: DeepReviewCapacityQueueReason,
546+
) -> bool {
547+
matches!(
548+
reason,
549+
DeepReviewCapacityQueueReason::ProviderConcurrencyLimit
550+
| DeepReviewCapacityQueueReason::TemporaryOverload
551+
)
552+
}
553+
517554
pub(crate) fn capacity_skip_result_for_provider_reason(
518555
reason: DeepReviewCapacityQueueReason,
519556
dialog_turn_id: &str,
@@ -707,6 +744,9 @@ pub(crate) async fn wait_for_provider_capacity_retry(
707744
let mut queue_timer = QueueWaitTimer::start(Instant::now());
708745
let max_wait = Duration::from_secs(max_wait_seconds);
709746
let optional_reviewer_count = is_optional_reviewer.then_some(1);
747+
let initial_active_reviewers = deep_review_active_reviewer_count(dialog_turn_id);
748+
let can_wake_on_active_reviewer_release =
749+
provider_capacity_wait_can_wake_on_active_reviewer_release(reason);
710750

711751
record_deep_review_runtime_provider_capacity_queue(dialog_turn_id, reason);
712752

@@ -737,7 +777,7 @@ pub(crate) async fn wait_for_provider_capacity_retry(
737777
optional_reviewer_count,
738778
Some(effective_parallel_instances),
739779
queue_elapsed_ms,
740-
conc_policy.max_queue_wait_seconds,
780+
max_wait_seconds,
741781
)
742782
.await;
743783
return DeepReviewProviderQueueWaitOutcome::Skipped {
@@ -764,7 +804,7 @@ pub(crate) async fn wait_for_provider_capacity_retry(
764804
optional_reviewer_count,
765805
Some(effective_parallel_instances),
766806
queue_elapsed_ms,
767-
conc_policy.max_queue_wait_seconds,
807+
max_wait_seconds,
768808
)
769809
.await;
770810
sleep(DEEP_REVIEW_QUEUE_POLL_INTERVAL).await;
@@ -788,10 +828,40 @@ pub(crate) async fn wait_for_provider_capacity_retry(
788828
optional_reviewer_count,
789829
Some(effective_parallel_instances),
790830
queue_elapsed_ms,
791-
conc_policy.max_queue_wait_seconds,
831+
max_wait_seconds,
792832
)
793833
.await;
794-
return DeepReviewProviderQueueWaitOutcome::ReadyToRetry { queue_elapsed_ms };
834+
return DeepReviewProviderQueueWaitOutcome::ReadyToRetry {
835+
queue_elapsed_ms,
836+
early_capacity_probe: false,
837+
};
838+
}
839+
840+
if can_wake_on_active_reviewer_release
841+
&& initial_active_reviewers > 0
842+
&& active_reviewers < initial_active_reviewers
843+
{
844+
record_deep_review_runtime_queue_wait(dialog_turn_id, queue_elapsed_ms);
845+
clear_deep_review_queue_control_for_tool(dialog_turn_id, tool_id);
846+
emit_queue_state(
847+
session_id,
848+
dialog_turn_id,
849+
tool_id,
850+
subagent_type,
851+
DeepReviewQueueStatus::Running,
852+
Some(reason),
853+
0,
854+
active_reviewers,
855+
optional_reviewer_count,
856+
Some(effective_parallel_instances),
857+
queue_elapsed_ms,
858+
max_wait_seconds,
859+
)
860+
.await;
861+
return DeepReviewProviderQueueWaitOutcome::ReadyToRetry {
862+
queue_elapsed_ms,
863+
early_capacity_probe: true,
864+
};
795865
}
796866

797867
emit_queue_state(
@@ -806,7 +876,7 @@ pub(crate) async fn wait_for_provider_capacity_retry(
806876
optional_reviewer_count,
807877
Some(effective_parallel_instances),
808878
queue_elapsed_ms,
809-
conc_policy.max_queue_wait_seconds,
879+
max_wait_seconds,
810880
)
811881
.await;
812882

0 commit comments

Comments
 (0)