feat(pool): 引入 pro_first 调度预设、Pool 候选持久化跳过与诊断信息优化

- 新增 pro_first 调度预设(Pro 优先),更新 plus_first 仅针对 Plus 计划,移除 free_team_first
- Pool 内部候选(pool_key_index 不为空)跳过 DB 持久化(available/skipped/unused 均适用)
- LRU 排序新增 catalog_lru_score 回退:runtime 无记录时使用 last_used_at_unix_secs
- 执行路径 miss 诊断消息细化为中文,按 reason 分类输出可读说明
- build_local_request_candidate_status_record 补充 extra_data 和 created_at_unix_ms 字段
- OpenAI CLI 计划构建流程补充候选评估进度跟踪与 terminal reason 设置
- 前端 PoolSchedulingDialog 增加 pro_first 预设展示,修复 LRU 默认预设检测逻辑
This commit is contained in:
fawney19
2026-04-24 13:29:05 +08:00
parent f29649e3a8
commit f3c9835759
27 changed files with 950 additions and 153 deletions
@@ -9,6 +9,7 @@ use crate::ai_pipeline::planner::candidate_eligibility::{
use crate::ai_pipeline::planner::runtime_miss::record_local_runtime_candidate_skip_reason;
use crate::ai_pipeline::{GatewayAuthApiKeySnapshot, PlannerAppState};
use crate::clock::current_unix_ms;
use crate::handlers::shared::provider_pool::admin_provider_pool_config_from_config_value;
use crate::orchestration::{local_attempt_slot_count, ExecutionAttemptIdentity};
use crate::AppState;
@@ -67,6 +68,16 @@ pub(crate) fn remember_first_local_candidate_affinity(
);
}
fn should_persist_available_local_candidate(eligible: &EligibleLocalExecutionCandidate) -> bool {
eligible.orchestration.pool_key_index.is_none()
}
fn should_persist_skipped_local_candidate(candidate: &SkippedLocalExecutionCandidate) -> bool {
candidate.transport.as_ref().is_none_or(|transport| {
admin_provider_pool_config_from_config_value(transport.provider.config.as_ref()).is_none()
})
}
#[allow(clippy::too_many_arguments)]
pub(crate) async fn persist_available_local_execution_candidates<F>(
state: PlannerAppState<'_>,
@@ -102,21 +113,25 @@ where
let attempt_identity = ExecutionAttemptIdentity::new(candidate_index, retry_index)
.with_pool_key_index(pool_key_index);
let generated_candidate_id = Uuid::new_v4().to_string();
let candidate_id = state
.persist_available_local_candidate(
trace_id,
user_id,
api_key_id,
&eligible.candidate,
attempt_identity.candidate_index,
attempt_identity.retry_index,
&generated_candidate_id,
required_capabilities,
extra_data.clone(),
created_at_unix_ms,
error_context,
)
.await;
let candidate_id = if should_persist_available_local_candidate(eligible) {
state
.persist_available_local_candidate(
trace_id,
user_id,
api_key_id,
&eligible.candidate,
attempt_identity.candidate_index,
attempt_identity.retry_index,
&generated_candidate_id,
required_capabilities,
extra_data.clone(),
created_at_unix_ms,
error_context,
)
.await
} else {
generated_candidate_id
};
let eligible = if retry_index + 1 == attempt_slots {
owned_eligible
@@ -235,7 +250,11 @@ pub(crate) async fn persist_skipped_local_execution_candidates(
error_context: &'static str,
record_runtime_miss_diagnostic: bool,
) {
for (skipped_offset, skipped_candidate) in skipped_candidates.into_iter().enumerate() {
let mut next_candidate_index = starting_candidate_index;
for skipped_candidate in skipped_candidates {
if !should_persist_skipped_local_candidate(&skipped_candidate) {
continue;
}
let generated_candidate_id = Uuid::new_v4().to_string();
persist_skipped_local_execution_candidate(
state,
@@ -243,7 +262,7 @@ pub(crate) async fn persist_skipped_local_execution_candidates(
user_id,
api_key_id,
&skipped_candidate.candidate,
starting_candidate_index + skipped_offset as u32,
next_candidate_index,
&generated_candidate_id,
required_capabilities,
skipped_candidate.skip_reason,
@@ -252,6 +271,7 @@ pub(crate) async fn persist_skipped_local_execution_candidates(
record_runtime_miss_diagnostic,
)
.await;
next_candidate_index = next_candidate_index.saturating_add(1);
}
}
@@ -275,3 +295,200 @@ pub(crate) async fn persist_skipped_local_execution_candidates_with_context(
)
.await;
}
#[cfg(test)]
mod tests {
use std::sync::Arc;
use aether_data::repository::candidates::InMemoryRequestCandidateRepository;
use aether_provider_transport::snapshot::{
GatewayProviderTransportEndpoint, GatewayProviderTransportKey,
GatewayProviderTransportProvider,
};
use aether_scheduler_core::SchedulerMinimalCandidateSelectionCandidate;
use serde_json::json;
use super::*;
use crate::data::GatewayDataState;
use crate::orchestration::LocalExecutionCandidateMetadata;
fn sample_candidate(key_id: &str) -> SchedulerMinimalCandidateSelectionCandidate {
SchedulerMinimalCandidateSelectionCandidate {
provider_id: "provider-1".to_string(),
provider_name: "provider-1".to_string(),
provider_type: "codex".to_string(),
provider_priority: 10,
endpoint_id: "endpoint-1".to_string(),
endpoint_api_format: "openai:chat".to_string(),
key_id: key_id.to_string(),
key_name: key_id.to_string(),
key_auth_type: "api_key".to_string(),
key_internal_priority: 10,
key_global_priority_for_format: Some(10),
key_capabilities: None,
model_id: "model-1".to_string(),
global_model_id: "global-model-1".to_string(),
global_model_name: "gpt-5".to_string(),
selected_provider_model_name: "gpt-5".to_string(),
mapping_matched_model: None,
}
}
fn sample_transport(
key_id: &str,
provider_config: Option<serde_json::Value>,
) -> Arc<crate::ai_pipeline::GatewayProviderTransportSnapshot> {
Arc::new(crate::ai_pipeline::GatewayProviderTransportSnapshot {
provider: GatewayProviderTransportProvider {
id: "provider-1".to_string(),
name: "provider-1".to_string(),
provider_type: "codex".to_string(),
website: None,
is_active: true,
keep_priority_on_conversion: false,
enable_format_conversion: false,
concurrent_limit: None,
max_retries: None,
proxy: None,
request_timeout_secs: None,
stream_first_byte_timeout_secs: None,
config: provider_config,
},
endpoint: GatewayProviderTransportEndpoint {
id: "endpoint-1".to_string(),
provider_id: "provider-1".to_string(),
api_format: "openai:chat".to_string(),
api_family: Some("openai".to_string()),
endpoint_kind: Some("chat".to_string()),
is_active: true,
base_url: "https://example.com".to_string(),
header_rules: None,
body_rules: None,
max_retries: None,
custom_path: None,
config: None,
format_acceptance_config: None,
proxy: None,
},
key: GatewayProviderTransportKey {
id: key_id.to_string(),
provider_id: "provider-1".to_string(),
name: key_id.to_string(),
auth_type: "api_key".to_string(),
is_active: true,
api_formats: Some(vec!["openai:chat".to_string()]),
allowed_models: None,
capabilities: None,
rate_multipliers: None,
global_priority_by_format: None,
expires_at_unix_secs: None,
proxy: None,
fingerprint: None,
decrypted_api_key: "secret".to_string(),
decrypted_auth_config: None,
},
})
}
fn sample_eligible(
key_id: &str,
pool_key_index: Option<u32>,
) -> EligibleLocalExecutionCandidate {
EligibleLocalExecutionCandidate {
candidate: sample_candidate(key_id),
transport: sample_transport(
key_id,
pool_key_index.map(|_| json!({ "pool_advanced": {} })),
),
provider_api_format: "openai:chat".to_string(),
orchestration: LocalExecutionCandidateMetadata {
candidate_group_id: pool_key_index.map(|_| "pool-group".to_string()),
pool_key_index,
},
}
}
#[tokio::test]
async fn pool_candidates_are_not_persisted_as_available_before_attempt() {
let repository = Arc::new(InMemoryRequestCandidateRepository::default());
let app = AppState::new()
.expect("state should build")
.with_data_state_for_tests(
GatewayDataState::with_request_candidate_repository_for_tests(Arc::clone(
&repository,
)),
);
let attempts = persist_available_local_execution_candidates(
PlannerAppState::new(&app),
"trace-pool-lazy",
"user-1",
"api-key-1",
None,
vec![
sample_eligible("pool-key", Some(0)),
sample_eligible("normal-key", None),
],
"persist should not fail",
|_| None,
)
.await;
assert_eq!(attempts.len(), 2);
let stored = app
.read_request_candidates_by_request_id("trace-pool-lazy")
.await
.expect("request candidates should read");
assert_eq!(stored.len(), 1);
assert_eq!(stored[0].key_id.as_deref(), Some("normal-key"));
}
#[tokio::test]
async fn pool_internal_skipped_candidates_are_not_persisted() {
let repository = Arc::new(InMemoryRequestCandidateRepository::default());
let app = AppState::new()
.expect("state should build")
.with_data_state_for_tests(
GatewayDataState::with_request_candidate_repository_for_tests(Arc::clone(
&repository,
)),
);
persist_skipped_local_execution_candidates(
&app,
"trace-pool-skipped",
"user-1",
"api-key-1",
None,
0,
vec![
SkippedLocalExecutionCandidate {
candidate: sample_candidate("pool-skipped"),
skip_reason: "pool_cooldown",
transport: Some(sample_transport(
"pool-skipped",
Some(json!({ "pool_advanced": {} })),
)),
extra_data: None,
},
SkippedLocalExecutionCandidate {
candidate: sample_candidate("normal-skipped"),
skip_reason: "key_inactive",
transport: None,
extra_data: None,
},
],
"persist skipped should not fail",
false,
)
.await;
let stored = app
.read_request_candidates_by_request_id("trace-pool-skipped")
.await
.expect("request candidates should read");
assert_eq!(stored.len(), 1);
assert_eq!(stored[0].key_id.as_deref(), Some("normal-skipped"));
assert_eq!(stored[0].candidate_index, 0);
}
}
@@ -47,6 +47,7 @@ struct PoolCatalogKeyContext {
quota_exhausted: bool,
health_score: Option<f64>,
latency_avg_ms: Option<f64>,
catalog_lru_score: Option<f64>,
}
#[derive(Debug, Clone, PartialEq, Eq)]
@@ -219,6 +220,7 @@ fn build_pool_catalog_key_context(
.unwrap_or(false),
health_score,
latency_avg_ms,
catalog_lru_score: Some(key.last_used_at_unix_secs.unwrap_or(0) as f64),
}
}
@@ -482,6 +484,9 @@ fn schedule_pool_group(
continue;
}
let lru_score =
runtime_lru_score(runtime, key_id.as_str()).or(key_context.catalog_lru_score);
available.push(PoolGroupCandidateOrdering {
eligible: EligibleLocalExecutionCandidate {
candidate,
@@ -491,7 +496,7 @@ fn schedule_pool_group(
},
key_context,
original_index,
lru_score: runtime_lru_score(runtime, key_id.as_str()),
lru_score,
cost_usage: runtime_cost_usage(runtime, key_id.as_str()),
});
}
@@ -593,8 +598,8 @@ fn build_pool_sort_vectors(
"cache_affinity" => cache_affinity_ranks.clone(),
"priority_first" => priority_first_ranks(items, &lru_ranks),
"single_account" => single_account_ranks(items),
"free_team_first" => plan_ranks(items, &lru_ranks, preset.mode.as_deref()),
"plus_first" => plan_ranks(items, &lru_ranks, Some("plus_only")),
"pro_first" => plan_ranks(items, &lru_ranks, Some("pro_only")),
"free_first" => plan_ranks(items, &lru_ranks, Some("free_only")),
"team_first" => plan_ranks(items, &lru_ranks, Some("team_only")),
"health_first" => health_first_ranks(items, &lru_ranks),
@@ -906,8 +911,17 @@ fn plan_priority_score(plan_type: Option<&str>, mode: Option<&str>) -> f64 {
None => 0.8,
},
"plus_only" => match plan_type {
Some("plus" | "pro") => 0.0,
Some("enterprise" | "business") => 0.3,
Some("plus") => 0.0,
Some("pro") => 0.3,
Some("enterprise" | "business") => 0.4,
Some("free" | "team") => 0.7,
Some(_) => 0.7,
None => 0.8,
},
"pro_only" => match plan_type {
Some("pro") => 0.0,
Some("plus") => 0.3,
Some("enterprise" | "business") => 0.4,
Some("free" | "team") => 0.7,
Some(_) => 0.7,
None => 0.8,
@@ -992,7 +1006,7 @@ fn normalize_enabled_pool_presets(
fn pool_preset_supported_for_provider(preset: &str, provider_type: &str) -> bool {
match preset {
"free_first" | "free_team_first" | "plus_first" | "recent_refresh" | "team_first" => {
"free_first" | "plus_first" | "pro_first" | "recent_refresh" | "team_first" => {
matches!(provider_type, "codex" | "kiro")
}
_ => true,
@@ -1090,6 +1104,56 @@ mod tests {
);
}
#[test]
fn pool_scheduler_uses_catalog_last_used_when_runtime_lru_is_missing() {
let recent_key = sample_eligible_candidate(
"provider-pool",
"endpoint-1",
"key-recent",
10,
Some(json!({ "pool_advanced": { "lru_enabled": true } })),
);
let older_key = sample_eligible_candidate(
"provider-pool",
"endpoint-1",
"key-older",
10,
Some(json!({ "pool_advanced": { "lru_enabled": true } })),
);
let key_context_by_id = BTreeMap::from([
(
"key-recent".to_string(),
PoolCatalogKeyContext {
catalog_lru_score: Some(200.0),
..PoolCatalogKeyContext::default()
},
),
(
"key-older".to_string(),
PoolCatalogKeyContext {
catalog_lru_score: Some(100.0),
..PoolCatalogKeyContext::default()
},
),
]);
let (reordered, skipped) = apply_local_execution_pool_scheduler_with_runtime_map(
vec![recent_key, older_key],
&BTreeMap::new(),
&key_context_by_id,
);
assert!(skipped.is_empty());
assert_eq!(
reordered
.iter()
.map(|item| item.candidate.key_id.as_str())
.collect::<Vec<_>>(),
vec!["key-older", "key-recent"]
);
}
#[test]
fn pool_scheduler_attaches_group_and_pool_metadata_to_ranked_candidates() {
let pool_first = sample_eligible_candidate(
@@ -1394,7 +1458,7 @@ mod tests {
}
#[test]
fn pool_scheduler_supports_free_team_first_modes() {
fn pool_scheduler_supports_pro_first_plan_preset() {
let key_plus = sample_eligible_candidate(
"provider-pool",
"endpoint-1",
@@ -1402,18 +1466,18 @@ mod tests {
10,
Some(json!({
"pool_advanced": {
"scheduling_presets": [{"preset": "free_team_first", "enabled": true, "mode": "team_only"}]
"scheduling_presets": [{"preset": "pro_first", "enabled": true}]
}
})),
);
let key_free = sample_eligible_candidate(
let key_pro = sample_eligible_candidate(
"provider-pool",
"endpoint-1",
"key-free",
"key-pro",
10,
Some(json!({
"pool_advanced": {
"scheduling_presets": [{"preset": "free_team_first", "enabled": true, "mode": "team_only"}]
"scheduling_presets": [{"preset": "pro_first", "enabled": true}]
}
})),
);
@@ -1424,7 +1488,7 @@ mod tests {
10,
Some(json!({
"pool_advanced": {
"scheduling_presets": [{"preset": "free_team_first", "enabled": true, "mode": "team_only"}]
"scheduling_presets": [{"preset": "pro_first", "enabled": true}]
}
})),
);
@@ -1438,9 +1502,9 @@ mod tests {
},
),
(
"key-free".to_string(),
"key-pro".to_string(),
PoolCatalogKeyContext {
oauth_plan_type: Some("free".to_string()),
oauth_plan_type: Some("pro".to_string()),
..PoolCatalogKeyContext::default()
},
),
@@ -1454,7 +1518,7 @@ mod tests {
]);
let (reordered, skipped) = apply_local_execution_pool_scheduler_with_runtime_map(
vec![key_plus, key_free, key_team],
vec![key_plus, key_team, key_pro],
&BTreeMap::new(),
&key_context_by_id,
);
@@ -1465,7 +1529,7 @@ mod tests {
.iter()
.map(|item| item.candidate.key_id.as_str())
.collect::<Vec<_>>(),
vec!["key-team", "key-free", "key-plus"]
vec!["key-pro", "key-plus", "key-team"]
);
}
@@ -1584,6 +1648,7 @@ mod tests {
}));
key.success_count = Some(4);
key.total_response_time_ms = Some(200);
key.last_used_at_unix_secs = Some(1_711_000_123);
let app = AppState::new()
.expect("state should build")
@@ -1601,6 +1666,7 @@ mod tests {
assert_eq!(context.quota_usage_ratio, Some(0.25));
assert_eq!(context.quota_reset_seconds, Some(3600.0));
assert_eq!(context.latency_avg_ms, Some(50.0));
assert_eq!(context.catalog_lru_score, Some(1_711_000_123.0));
}
fn sample_eligible_candidate(
@@ -28,6 +28,7 @@ use crate::ai_pipeline::planner::decision_input::{
use crate::ai_pipeline::planner::materialization_policy::{
build_local_candidate_persistence_policy, LocalCandidatePersistencePolicyKind,
};
use crate::ai_pipeline::planner::runtime_miss::set_local_runtime_miss_diagnostic_reason;
use crate::ai_pipeline::planner::spec_metadata::local_openai_cli_spec_metadata;
use crate::ai_pipeline::PlannerAppState;
use crate::ai_pipeline::{
@@ -47,28 +48,83 @@ pub(crate) async fn resolve_local_openai_cli_decision_input(
trace_id: &str,
decision: &GatewayControlDecision,
body_json: &serde_json::Value,
plan_kind: &str,
) -> Option<LocalOpenAiCliDecisionInput> {
let auth_context: ExecutionRuntimeAuthContext =
resolve_local_decision_execution_runtime_auth_context(decision)?;
let Some(auth_context) = resolve_local_decision_execution_runtime_auth_context(decision) else {
warn!(
trace_id = %trace_id,
route_class = ?decision.route_class,
route_family = ?decision.route_family,
route_kind = ?decision.route_kind,
"gateway local openai cli decision skipped: missing_auth_context"
);
set_local_runtime_miss_diagnostic_reason(
state,
trace_id,
decision,
plan_kind,
extract_standard_requested_model(body_json).as_deref(),
"missing_auth_context",
);
return None;
};
let requested_model = extract_standard_requested_model(body_json)?;
let Some(requested_model) = extract_standard_requested_model(body_json) else {
warn!(
trace_id = %trace_id,
"gateway local openai cli decision skipped: missing_requested_model"
);
set_local_runtime_miss_diagnostic_reason(
state,
trace_id,
decision,
plan_kind,
None,
"missing_requested_model",
);
return None;
};
let resolved_input = match resolve_local_authenticated_decision_input(
state,
auth_context,
auth_context.clone(),
Some(requested_model.as_str()),
None,
)
.await
{
Ok(Some(resolved_input)) => resolved_input,
Ok(None) => return None,
Ok(None) => {
warn!(
trace_id = %trace_id,
user_id = %auth_context.user_id,
api_key_id = %auth_context.api_key_id,
"gateway local openai cli decision skipped: auth_snapshot_missing"
);
set_local_runtime_miss_diagnostic_reason(
state,
trace_id,
decision,
plan_kind,
Some(requested_model.as_str()),
"auth_snapshot_missing",
);
return None;
}
Err(err) => {
warn!(
trace_id = %trace_id,
error = ?err,
"gateway local openai cli decision auth snapshot read failed"
);
set_local_runtime_miss_diagnostic_reason(
state,
trace_id,
decision,
plan_kind,
Some(requested_model.as_str()),
"auth_snapshot_read_failed",
);
return None;
}
};
@@ -85,7 +141,7 @@ pub(crate) async fn materialize_local_openai_cli_candidate_attempts(
input: &LocalOpenAiCliDecisionInput,
body_json: &serde_json::Value,
spec: LocalOpenAiCliSpec,
) -> Result<Vec<LocalOpenAiCliCandidateAttempt>, GatewayError> {
) -> Result<(Vec<LocalOpenAiCliCandidateAttempt>, usize), GatewayError> {
let spec_metadata = local_openai_cli_spec_metadata(spec);
let client_api_format = spec_metadata.api_format.trim().to_ascii_lowercase();
let planner_state = PlannerAppState::new(state);
@@ -223,6 +279,8 @@ pub(crate) async fn materialize_local_openai_cli_candidate_attempts(
})
.collect::<Vec<_>>();
let candidate_count = candidates.len() + skipped_candidates.len();
remember_first_local_candidate_affinity(
planner_state,
Some(&input.auth_snapshot),
@@ -275,7 +333,7 @@ pub(crate) async fn materialize_local_openai_cli_candidate_attempts(
)
.await;
Ok(attempts)
Ok((attempts, candidate_count))
}
pub(crate) async fn mark_skipped_local_openai_cli_candidate(
state: &AppState,
@@ -60,12 +60,13 @@ pub(crate) async fn maybe_build_sync_local_openai_cli_decision_payload(
};
let Some(input) =
resolve_local_openai_cli_decision_input(state, trace_id, decision, body_json).await
resolve_local_openai_cli_decision_input(state, trace_id, decision, body_json, plan_kind)
.await
else {
return Ok(None);
};
let attempts =
let (attempts, _) =
materialize_local_openai_cli_candidate_attempts(state, trace_id, &input, body_json, spec)
.await?;
@@ -95,12 +96,13 @@ pub(crate) async fn maybe_build_stream_local_openai_cli_decision_payload(
};
let Some(input) =
resolve_local_openai_cli_decision_input(state, trace_id, decision, body_json).await
resolve_local_openai_cli_decision_input(state, trace_id, decision, body_json, plan_kind)
.await
else {
return Ok(None);
};
let attempts =
let (attempts, _) =
materialize_local_openai_cli_candidate_attempts(state, trace_id, &input, body_json, spec)
.await?;
@@ -9,6 +9,10 @@ use crate::ai_pipeline::planner::plan_builders::{
build_openai_cli_stream_plan_from_decision, build_openai_cli_sync_plan_from_decision,
LocalStreamPlanAndReport, LocalSyncPlanAndReport,
};
use crate::ai_pipeline::planner::runtime_miss::{
apply_local_runtime_candidate_evaluation_progress,
apply_local_runtime_candidate_terminal_reason, set_local_runtime_miss_diagnostic_reason,
};
use crate::ai_pipeline::planner::spec_metadata::local_openai_cli_spec_metadata;
use crate::ai_pipeline::GatewayControlDecision;
pub(crate) use crate::ai_pipeline::{
@@ -26,15 +30,33 @@ pub(super) async fn build_local_sync_plan_and_reports(
spec: LocalOpenAiCliSpec,
) -> Result<Vec<LocalSyncPlanAndReport>, GatewayError> {
let spec_metadata = local_openai_cli_spec_metadata(spec);
let Some(input) =
resolve_local_openai_cli_decision_input(state, trace_id, decision, body_json).await
let Some(input) = resolve_local_openai_cli_decision_input(
state,
trace_id,
decision,
body_json,
spec_metadata.decision_kind,
)
.await
else {
return Ok(Vec::new());
};
set_local_runtime_miss_diagnostic_reason(
state,
trace_id,
decision,
spec_metadata.decision_kind,
Some(input.requested_model.as_str()),
"candidate_evaluation_incomplete",
);
let attempts =
let (attempts, candidate_count) =
materialize_local_openai_cli_candidate_attempts(state, trace_id, &input, body_json, spec)
.await?;
apply_local_runtime_candidate_evaluation_progress(state, trace_id, candidate_count);
if candidate_count == 0 {
return Ok(Vec::new());
}
let mut plans = Vec::new();
for attempt in attempts {
@@ -60,6 +82,7 @@ pub(super) async fn build_local_sync_plan_and_reports(
}
}
apply_local_runtime_candidate_terminal_reason(state, trace_id, "no_local_sync_plans");
Ok(plans)
}
@@ -72,15 +95,33 @@ pub(super) async fn build_local_stream_plan_and_reports(
spec: LocalOpenAiCliSpec,
) -> Result<Vec<LocalStreamPlanAndReport>, GatewayError> {
let spec_metadata = local_openai_cli_spec_metadata(spec);
let Some(input) =
resolve_local_openai_cli_decision_input(state, trace_id, decision, body_json).await
let Some(input) = resolve_local_openai_cli_decision_input(
state,
trace_id,
decision,
body_json,
spec_metadata.decision_kind,
)
.await
else {
return Ok(Vec::new());
};
set_local_runtime_miss_diagnostic_reason(
state,
trace_id,
decision,
spec_metadata.decision_kind,
Some(input.requested_model.as_str()),
"candidate_evaluation_incomplete",
);
let attempts =
let (attempts, candidate_count) =
materialize_local_openai_cli_candidate_attempts(state, trace_id, &input, body_json, spec)
.await?;
apply_local_runtime_candidate_evaluation_progress(state, trace_id, candidate_count);
if candidate_count == 0 {
return Ok(Vec::new());
}
let mut plans = Vec::new();
for attempt in attempts {
@@ -106,5 +147,6 @@ pub(super) async fn build_local_stream_plan_and_reports(
}
}
apply_local_runtime_candidate_terminal_reason(state, trace_id, "no_local_stream_plans");
Ok(plans)
}