feat(scheduler): candidate report 携带 header_rules 与 body_rules

各 planner 在构造 report context 时附带端点的 header/body 改写规则,scheduler 解析后写入 candidate extra_data,便于追踪请求实际应用的透传规则
This commit is contained in:
fawney19
2026-04-25 11:43:53 +08:00
parent 65e915fd1d
commit 61e12e2c17
9 changed files with 68 additions and 0 deletions
@@ -99,6 +99,8 @@ pub(crate) async fn maybe_build_local_same_format_provider_decision_payload_for_
mapped_model: Some(&resolved.mapped_model),
candidate_group_id: eligible.orchestration.candidate_group_id.as_deref(),
upstream_url: Some(&resolved.upstream_url),
header_rules: resolved.transport.endpoint.header_rules.as_ref(),
body_rules: resolved.transport.endpoint.body_rules.as_ref(),
provider_request_method: Some(serde_json::Value::Null),
provider_request_headers: Some(&resolved.provider_request_headers),
original_headers: &parts.headers,
@@ -24,6 +24,8 @@ pub(crate) struct LocalExecutionReportContextParts<'a> {
pub(crate) mapped_model: Option<&'a str>,
pub(crate) candidate_group_id: Option<&'a str>,
pub(crate) upstream_url: Option<&'a str>,
pub(crate) header_rules: Option<&'a Value>,
pub(crate) body_rules: Option<&'a Value>,
pub(crate) provider_request_method: Option<Value>,
pub(crate) provider_request_headers: Option<&'a BTreeMap<String, String>>,
pub(crate) original_headers: &'a http::HeaderMap,
@@ -172,6 +174,12 @@ pub(crate) fn build_local_execution_report_context(
Value::String(upstream_url.to_string()),
);
}
if let Some(header_rules) = parts.header_rules {
object.insert("header_rules".to_string(), header_rules.clone());
}
if let Some(body_rules) = parts.body_rules {
object.insert("body_rules".to_string(), body_rules.clone());
}
if let Some(provider_request_method) = parts.provider_request_method {
object.insert(
"provider_request_method".to_string(),
@@ -84,6 +84,8 @@ pub(super) async fn maybe_build_local_gemini_files_decision_payload_for_candidat
mapped_model: None,
candidate_group_id: eligible.orchestration.candidate_group_id.as_deref(),
upstream_url: None,
header_rules: transport.endpoint.header_rules.as_ref(),
body_rules: transport.endpoint.body_rules.as_ref(),
provider_request_method: None,
provider_request_headers: None,
original_headers: &parts.headers,
@@ -81,6 +81,8 @@ pub(super) async fn maybe_build_local_openai_image_decision_payload_for_candidat
mapped_model: Some(&resolved.mapped_model),
candidate_group_id: eligible.orchestration.candidate_group_id.as_deref(),
upstream_url: Some(&resolved.upstream_url),
header_rules: transport.endpoint.header_rules.as_ref(),
body_rules: transport.endpoint.body_rules.as_ref(),
provider_request_method: Some(serde_json::Value::String(parts.method.to_string())),
provider_request_headers: Some(&resolved.provider_request_headers),
original_headers: &parts.headers,
@@ -67,6 +67,8 @@ pub(super) async fn maybe_build_local_video_create_decision_payload_for_candidat
mapped_model: Some(&resolved.mapped_model),
candidate_group_id: eligible.orchestration.candidate_group_id.as_deref(),
upstream_url: None,
header_rules: transport.endpoint.header_rules.as_ref(),
body_rules: transport.endpoint.body_rules.as_ref(),
provider_request_method: None,
provider_request_headers: None,
original_headers: &parts.headers,
@@ -79,6 +79,8 @@ pub(super) async fn maybe_build_local_standard_decision_payload_for_candidate(
mapped_model: Some(&resolved.mapped_model),
candidate_group_id: eligible.orchestration.candidate_group_id.as_deref(),
upstream_url: Some(&resolved.upstream_url),
header_rules: resolved.transport.endpoint.header_rules.as_ref(),
body_rules: resolved.transport.endpoint.body_rules.as_ref(),
provider_request_method: Some(serde_json::Value::Null),
provider_request_headers: Some(&resolved.provider_request_headers),
original_headers: &parts.headers,
@@ -101,6 +101,8 @@ pub(crate) async fn maybe_build_local_openai_chat_decision_payload_for_candidate
mapped_model: Some(&resolved.mapped_model),
candidate_group_id: eligible.orchestration.candidate_group_id.as_deref(),
upstream_url: Some(&resolved.upstream_url),
header_rules: resolved.transport.endpoint.header_rules.as_ref(),
body_rules: resolved.transport.endpoint.body_rules.as_ref(),
provider_request_method: Some(serde_json::Value::Null),
provider_request_headers: Some(&resolved.provider_request_headers),
original_headers: &parts.headers,
@@ -99,6 +99,8 @@ pub(crate) async fn maybe_build_local_openai_cli_decision_payload_for_candidate(
mapped_model: Some(&resolved.mapped_model),
candidate_group_id: eligible.orchestration.candidate_group_id.as_deref(),
upstream_url: Some(&resolved.upstream_url),
header_rules: resolved.transport.endpoint.header_rules.as_ref(),
body_rules: resolved.transport.endpoint.body_rules.as_ref(),
provider_request_method: Some(serde_json::Value::Null),
provider_request_headers: Some(&resolved.provider_request_headers),
original_headers: &parts.headers,
@@ -20,6 +20,8 @@ pub struct SchedulerRequestCandidateReportContext {
pub upstream_url: Option<String>,
pub mapped_model: Option<String>,
pub key_name: Option<String>,
pub header_rules: Option<Value>,
pub body_rules: Option<Value>,
pub proxy: Option<Value>,
}
@@ -109,6 +111,14 @@ pub fn parse_request_candidate_report_context(
upstream_url: string_field(report_context, "upstream_url"),
mapped_model: string_field(report_context, "mapped_model"),
key_name: string_field(report_context, "key_name"),
header_rules: report_context
.get("header_rules")
.cloned()
.filter(|value| !value.is_null()),
body_rules: report_context
.get("body_rules")
.cloned()
.filter(|value| !value.is_null()),
proxy: report_context
.get("proxy")
.cloned()
@@ -138,6 +148,8 @@ pub fn resolve_report_request_candidate_slot(
upstream_url,
mapped_model,
key_name,
header_rules,
body_rules,
proxy,
} = metadata;
let request_id = request_id?;
@@ -147,6 +159,8 @@ pub fn resolve_report_request_candidate_slot(
upstream_url,
mapped_model,
key_name,
header_rules,
body_rules,
proxy,
);
let created_at_unix_ms = matched_candidate
@@ -308,6 +322,8 @@ pub fn build_local_request_candidate_status_record(
metadata.upstream_url.clone(),
metadata.mapped_model.clone(),
metadata.key_name.clone(),
metadata.header_rules.clone(),
metadata.body_rules.clone(),
metadata.proxy.clone(),
);
let created_at_unix_ms = started_at_unix_ms.or(finished_at_unix_ms);
@@ -516,6 +532,8 @@ fn build_report_candidate_extra_data(
upstream_url: Option<String>,
mapped_model: Option<String>,
key_name: Option<String>,
header_rules: Option<Value>,
body_rules: Option<Value>,
proxy: Option<Value>,
) -> Option<Value> {
let mut extra_data = Map::with_capacity(8);
@@ -542,6 +560,12 @@ fn build_report_candidate_extra_data(
if let Some(key_name) = key_name {
extra_data.insert("key_name".to_string(), Value::String(key_name));
}
if let Some(header_rules) = header_rules {
extra_data.insert("header_rules".to_string(), header_rules);
}
if let Some(body_rules) = body_rules {
extra_data.insert("body_rules".to_string(), body_rules);
}
if let Some(proxy) = proxy {
extra_data.insert("proxy".to_string(), proxy);
}
@@ -683,6 +707,12 @@ mod tests {
"key_id": "catalog-key-1",
"client_api_format": "openai:chat",
"provider_api_format": "openai:cli",
"header_rules": [
{"op": "set", "name": "x-test", "value": "1"}
],
"body_rules": [
{"op": "remove", "path": "/store"}
],
"proxy": {
"node_id": "proxy-node-1",
"node_name": "edge-1",
@@ -719,6 +749,22 @@ mod tests {
.and_then(|value| value.get("source")),
Some(&json!("provider"))
);
assert_eq!(
slot.extra_data
.as_ref()
.and_then(|value| value.get("header_rules"))
.and_then(Value::as_array)
.map(Vec::len),
Some(1)
);
assert_eq!(
slot.extra_data
.as_ref()
.and_then(|value| value.get("body_rules"))
.and_then(Value::as_array)
.map(Vec::len),
Some(1)
);
}
#[test]