mirror of
https://github.com/fawney19/Aether.git
synced 2026-09-03 01:40:21 +08:00
feat: add sync image heartbeat toggle
This commit is contained in:
@@ -899,6 +899,12 @@ fn build_json_whitespace_heartbeat_stream(
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
pub(crate) fn build_openai_image_sync_json_whitespace_heartbeat_stream(
|
||||||
|
rx: mpsc::Receiver<Result<Bytes, IoError>>,
|
||||||
|
) -> impl futures_util::Stream<Item = Result<Bytes, IoError>> + Send + 'static {
|
||||||
|
build_json_whitespace_heartbeat_stream(rx, OPENAI_IMAGE_SYNC_JSON_HEARTBEAT_INTERVAL, None)
|
||||||
|
}
|
||||||
|
|
||||||
async fn openai_image_sync_json_heartbeat_final_bytes(
|
async fn openai_image_sync_json_heartbeat_final_bytes(
|
||||||
result: Result<Option<Response<Body>>, GatewayError>,
|
result: Result<Option<Response<Body>>, GatewayError>,
|
||||||
) -> Vec<u8> {
|
) -> Vec<u8> {
|
||||||
|
|||||||
@@ -1,6 +1,8 @@
|
|||||||
mod execution;
|
mod execution;
|
||||||
|
|
||||||
pub(crate) use execution::execute_execution_runtime_sync;
|
pub(crate) use execution::{
|
||||||
|
build_openai_image_sync_json_whitespace_heartbeat_stream, execute_execution_runtime_sync,
|
||||||
|
};
|
||||||
|
|
||||||
#[allow(unused_imports)]
|
#[allow(unused_imports)]
|
||||||
pub(crate) use execution::{
|
pub(crate) use execution::{
|
||||||
|
|||||||
@@ -1,3 +1,13 @@
|
|||||||
|
use std::collections::{BTreeMap, VecDeque};
|
||||||
|
use std::io::Error as IoError;
|
||||||
|
use std::time::Instant;
|
||||||
|
|
||||||
|
use axum::body::{to_bytes, Body, Bytes};
|
||||||
|
use axum::http::header::{CACHE_CONTROL, CONTENT_ENCODING, CONTENT_LENGTH, CONTENT_TYPE};
|
||||||
|
use axum::http::{HeaderName, HeaderValue, Response, StatusCode};
|
||||||
|
use serde_json::{json, Value};
|
||||||
|
use tokio::sync::mpsc;
|
||||||
|
|
||||||
use crate::ai_serving::api::{
|
use crate::ai_serving::api::{
|
||||||
build_local_gemini_files_stream_attempt_source_for_kind,
|
build_local_gemini_files_stream_attempt_source_for_kind,
|
||||||
build_local_gemini_files_sync_attempt_source_for_kind,
|
build_local_gemini_files_sync_attempt_source_for_kind,
|
||||||
@@ -22,13 +32,31 @@ use crate::ai_serving::api::{
|
|||||||
LocalStandardSpec, EXECUTION_RUNTIME_STREAM_DECISION_ACTION,
|
LocalStandardSpec, EXECUTION_RUNTIME_STREAM_DECISION_ACTION,
|
||||||
EXECUTION_RUNTIME_SYNC_DECISION_ACTION,
|
EXECUTION_RUNTIME_SYNC_DECISION_ACTION,
|
||||||
};
|
};
|
||||||
|
use crate::ai_serving::LocalExecutionAttemptSource;
|
||||||
|
use crate::api::response::{
|
||||||
|
attach_control_metadata_headers, build_client_response_from_parts_with_mutator,
|
||||||
|
};
|
||||||
|
use crate::constants::EXECUTION_PATH_LOCAL_EXECUTION_RUNTIME_MISS;
|
||||||
use crate::control::GatewayControlDecision;
|
use crate::control::GatewayControlDecision;
|
||||||
|
use crate::execution_runtime::sync::{
|
||||||
|
build_openai_image_sync_json_whitespace_heartbeat_stream, execute_execution_runtime_sync,
|
||||||
|
};
|
||||||
use crate::executor::candidate_loop::{
|
use crate::executor::candidate_loop::{
|
||||||
execute_stream_attempt_source, execute_sync_attempt_source, execute_sync_plan_and_reports,
|
execute_stream_attempt_source, execute_sync_attempt_source, execute_sync_plan_and_reports,
|
||||||
|
mark_unused_local_candidates,
|
||||||
};
|
};
|
||||||
use crate::executor::LocalExecutionRequestOutcome;
|
use crate::executor::{
|
||||||
|
build_local_execution_exhaustion, record_failed_usage_for_exhausted_request,
|
||||||
|
LocalExecutionRequestOutcome,
|
||||||
|
};
|
||||||
|
use crate::handlers::shared::system_config_bool;
|
||||||
use crate::{AiExecutionDecision, AppState, GatewayError};
|
use crate::{AiExecutionDecision, AppState, GatewayError};
|
||||||
|
|
||||||
|
const ENABLE_OPENAI_IMAGE_SYNC_HEARTBEAT_CONFIG_KEY: &str = "enable_openai_image_sync_heartbeat";
|
||||||
|
const OPENAI_IMAGE_SYNC_HEARTBEAT_INTERNAL_ERROR_STATUS: u16 = 502;
|
||||||
|
const OPENAI_IMAGE_SYNC_HEARTBEAT_EXHAUSTED_STATUS: u16 = 503;
|
||||||
|
const OPENAI_IMAGE_SYNC_HEARTBEAT_ERROR_MESSAGE_LIMIT: usize = 4096;
|
||||||
|
|
||||||
pub(crate) async fn maybe_execute_sync_local_path(
|
pub(crate) async fn maybe_execute_sync_local_path(
|
||||||
state: &AppState,
|
state: &AppState,
|
||||||
parts: &http::request::Parts,
|
parts: &http::request::Parts,
|
||||||
@@ -449,6 +477,226 @@ pub(crate) async fn maybe_execute_sync_via_local_gemini_files_decision(
|
|||||||
.await
|
.await
|
||||||
}
|
}
|
||||||
|
|
||||||
|
async fn openai_image_sync_heartbeat_enabled(state: &AppState) -> bool {
|
||||||
|
match state
|
||||||
|
.read_system_config_json_value(ENABLE_OPENAI_IMAGE_SYNC_HEARTBEAT_CONFIG_KEY)
|
||||||
|
.await
|
||||||
|
{
|
||||||
|
Ok(value) => system_config_bool(value.as_ref(), false),
|
||||||
|
Err(err) => {
|
||||||
|
tracing::warn!(
|
||||||
|
event_name = "openai_image_sync_heartbeat_config_read_failed",
|
||||||
|
log_type = "ops",
|
||||||
|
error = ?err,
|
||||||
|
"gateway failed to read sync image heartbeat config; defaulting disabled"
|
||||||
|
);
|
||||||
|
false
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
fn build_openai_image_sync_heartbeat_shell_response(
|
||||||
|
state: AppState,
|
||||||
|
request_path: String,
|
||||||
|
trace_id: String,
|
||||||
|
decision: GatewayControlDecision,
|
||||||
|
plan_kind: String,
|
||||||
|
attempts: Vec<AiSyncAttempt>,
|
||||||
|
) -> Result<Response<Body>, GatewayError> {
|
||||||
|
let request_id = attempts
|
||||||
|
.first()
|
||||||
|
.map(|attempt| attempt.plan.request_id.clone())
|
||||||
|
.filter(|value| !value.trim().is_empty());
|
||||||
|
let trace_id_for_response = trace_id.clone();
|
||||||
|
let decision_for_response = decision.clone();
|
||||||
|
let started_at = Instant::now();
|
||||||
|
let (tx, rx) = mpsc::channel::<Result<Bytes, IoError>>(1);
|
||||||
|
|
||||||
|
tokio::spawn(async move {
|
||||||
|
let bytes = openai_image_sync_heartbeat_final_bytes(
|
||||||
|
execute_openai_image_sync_heartbeat_attempts(
|
||||||
|
state,
|
||||||
|
request_path,
|
||||||
|
trace_id,
|
||||||
|
decision,
|
||||||
|
plan_kind,
|
||||||
|
attempts,
|
||||||
|
started_at,
|
||||||
|
)
|
||||||
|
.await,
|
||||||
|
)
|
||||||
|
.await;
|
||||||
|
let _ = tx.send(Ok(Bytes::from(bytes))).await;
|
||||||
|
});
|
||||||
|
|
||||||
|
let headers = BTreeMap::from([(
|
||||||
|
CONTENT_TYPE.as_str().to_string(),
|
||||||
|
"application/json".to_string(),
|
||||||
|
)]);
|
||||||
|
let response = build_client_response_from_parts_with_mutator(
|
||||||
|
StatusCode::OK.as_u16(),
|
||||||
|
&headers,
|
||||||
|
Body::from_stream(build_openai_image_sync_json_whitespace_heartbeat_stream(rx)),
|
||||||
|
trace_id_for_response.as_str(),
|
||||||
|
Some(&decision_for_response),
|
||||||
|
|headers| {
|
||||||
|
headers.remove(CONTENT_LENGTH);
|
||||||
|
headers.remove(CONTENT_ENCODING);
|
||||||
|
headers.insert(
|
||||||
|
CACHE_CONTROL,
|
||||||
|
HeaderValue::from_static("no-cache, no-transform"),
|
||||||
|
);
|
||||||
|
headers.insert(
|
||||||
|
HeaderName::from_static("x-accel-buffering"),
|
||||||
|
HeaderValue::from_static("no"),
|
||||||
|
);
|
||||||
|
Ok(())
|
||||||
|
},
|
||||||
|
)?;
|
||||||
|
attach_control_metadata_headers(response, request_id.as_deref(), None)
|
||||||
|
}
|
||||||
|
|
||||||
|
async fn execute_openai_image_sync_heartbeat_attempts(
|
||||||
|
state: AppState,
|
||||||
|
request_path: String,
|
||||||
|
trace_id: String,
|
||||||
|
decision: GatewayControlDecision,
|
||||||
|
plan_kind: String,
|
||||||
|
attempts: Vec<AiSyncAttempt>,
|
||||||
|
started_at: Instant,
|
||||||
|
) -> Result<LocalExecutionRequestOutcome, GatewayError> {
|
||||||
|
let mut attempts = VecDeque::from(attempts);
|
||||||
|
let mut last_attempted = None;
|
||||||
|
|
||||||
|
while let Some(attempt) = attempts.pop_front() {
|
||||||
|
let plan = attempt.plan;
|
||||||
|
let report_kind = attempt.report_kind;
|
||||||
|
let report_context = attempt.report_context;
|
||||||
|
last_attempted = Some((plan.clone(), report_context.clone()));
|
||||||
|
match execute_execution_runtime_sync(
|
||||||
|
&state,
|
||||||
|
request_path.as_str(),
|
||||||
|
plan,
|
||||||
|
trace_id.as_str(),
|
||||||
|
&decision,
|
||||||
|
plan_kind.as_str(),
|
||||||
|
report_kind,
|
||||||
|
report_context,
|
||||||
|
)
|
||||||
|
.await?
|
||||||
|
{
|
||||||
|
Some(response) => {
|
||||||
|
mark_unused_local_candidates(&state, attempts.into_iter().collect()).await;
|
||||||
|
return Ok(LocalExecutionRequestOutcome::responded(response));
|
||||||
|
}
|
||||||
|
None => continue,
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
let Some((last_plan, last_report_context)) = last_attempted else {
|
||||||
|
return Ok(LocalExecutionRequestOutcome::NoPath);
|
||||||
|
};
|
||||||
|
let exhaustion =
|
||||||
|
build_local_execution_exhaustion(&state, &last_plan, last_report_context.as_ref()).await;
|
||||||
|
record_failed_usage_for_exhausted_request(
|
||||||
|
&state,
|
||||||
|
exhaustion,
|
||||||
|
&started_at,
|
||||||
|
"OpenAI image sync heartbeat exhausted all local candidates",
|
||||||
|
EXECUTION_PATH_LOCAL_EXECUTION_RUNTIME_MISS,
|
||||||
|
None,
|
||||||
|
)
|
||||||
|
.await;
|
||||||
|
Ok(LocalExecutionRequestOutcome::NoPath)
|
||||||
|
}
|
||||||
|
|
||||||
|
async fn openai_image_sync_heartbeat_final_bytes(
|
||||||
|
result: Result<LocalExecutionRequestOutcome, GatewayError>,
|
||||||
|
) -> Vec<u8> {
|
||||||
|
match result {
|
||||||
|
Ok(LocalExecutionRequestOutcome::Responded(response)) => {
|
||||||
|
openai_image_sync_heartbeat_response_body_bytes(response).await
|
||||||
|
}
|
||||||
|
Ok(LocalExecutionRequestOutcome::Exhausted(_))
|
||||||
|
| Ok(LocalExecutionRequestOutcome::NoPath) => openai_image_sync_heartbeat_error_body(
|
||||||
|
OPENAI_IMAGE_SYNC_HEARTBEAT_EXHAUSTED_STATUS,
|
||||||
|
"OpenAI image sync exhausted all local candidates",
|
||||||
|
),
|
||||||
|
Err(err) => openai_image_sync_heartbeat_error_body(
|
||||||
|
OPENAI_IMAGE_SYNC_HEARTBEAT_INTERNAL_ERROR_STATUS,
|
||||||
|
&format!("{err:?}"),
|
||||||
|
),
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
async fn openai_image_sync_heartbeat_response_body_bytes(response: Response<Body>) -> Vec<u8> {
|
||||||
|
let status_code = response.status().as_u16();
|
||||||
|
match to_bytes(response.into_body(), usize::MAX).await {
|
||||||
|
Ok(bytes) if status_code < 400 && !bytes.is_empty() => bytes.to_vec(),
|
||||||
|
Ok(bytes) if status_code >= 400 => {
|
||||||
|
openai_image_sync_heartbeat_error_body_from_response(status_code, bytes.as_ref())
|
||||||
|
}
|
||||||
|
Ok(_) => openai_image_sync_heartbeat_error_body(
|
||||||
|
OPENAI_IMAGE_SYNC_HEARTBEAT_INTERNAL_ERROR_STATUS,
|
||||||
|
"empty sync image response",
|
||||||
|
),
|
||||||
|
Err(err) => openai_image_sync_heartbeat_error_body(
|
||||||
|
OPENAI_IMAGE_SYNC_HEARTBEAT_INTERNAL_ERROR_STATUS,
|
||||||
|
&err.to_string(),
|
||||||
|
),
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
fn openai_image_sync_heartbeat_error_body_from_response(status_code: u16, body: &[u8]) -> Vec<u8> {
|
||||||
|
if let Ok(mut value) = serde_json::from_slice::<Value>(body) {
|
||||||
|
if let Some(error) = value.get_mut("error").and_then(Value::as_object_mut) {
|
||||||
|
error.insert("upstream_status".to_string(), Value::from(status_code));
|
||||||
|
error
|
||||||
|
.entry("type".to_string())
|
||||||
|
.or_insert_with(|| Value::String("upstream_error".to_string()));
|
||||||
|
error.entry("message".to_string()).or_insert_with(|| {
|
||||||
|
Value::String(format!("upstream returned status {status_code}"))
|
||||||
|
});
|
||||||
|
return serde_json::to_vec(&value).unwrap_or_else(|_| {
|
||||||
|
openai_image_sync_heartbeat_error_body(
|
||||||
|
status_code,
|
||||||
|
&format!("upstream returned status {status_code}"),
|
||||||
|
)
|
||||||
|
});
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
let message = openai_image_sync_heartbeat_error_message_from_body(status_code, body);
|
||||||
|
openai_image_sync_heartbeat_error_body(status_code, message.as_str())
|
||||||
|
}
|
||||||
|
|
||||||
|
fn openai_image_sync_heartbeat_error_message_from_body(status_code: u16, body: &[u8]) -> String {
|
||||||
|
let text = String::from_utf8_lossy(body).trim().to_string();
|
||||||
|
if text.is_empty() {
|
||||||
|
return format!("upstream returned status {status_code}");
|
||||||
|
}
|
||||||
|
text.chars()
|
||||||
|
.take(OPENAI_IMAGE_SYNC_HEARTBEAT_ERROR_MESSAGE_LIMIT)
|
||||||
|
.collect()
|
||||||
|
}
|
||||||
|
|
||||||
|
fn openai_image_sync_heartbeat_error_body(status_code: u16, message: &str) -> Vec<u8> {
|
||||||
|
serde_json::to_vec(&json!({
|
||||||
|
"error": {
|
||||||
|
"type": "upstream_error",
|
||||||
|
"message": message,
|
||||||
|
"code": status_code,
|
||||||
|
"upstream_status": status_code,
|
||||||
|
}
|
||||||
|
}))
|
||||||
|
.unwrap_or_else(|_| {
|
||||||
|
format!(
|
||||||
|
"{{\"error\":{{\"type\":\"upstream_error\",\"code\":{status_code},\"upstream_status\":{status_code}}}}}"
|
||||||
|
)
|
||||||
|
.into_bytes()
|
||||||
|
})
|
||||||
|
}
|
||||||
|
|
||||||
pub(crate) async fn maybe_execute_sync_via_local_image_decision(
|
pub(crate) async fn maybe_execute_sync_via_local_image_decision(
|
||||||
state: &AppState,
|
state: &AppState,
|
||||||
parts: &http::request::Parts,
|
parts: &http::request::Parts,
|
||||||
@@ -458,20 +706,38 @@ pub(crate) async fn maybe_execute_sync_via_local_image_decision(
|
|||||||
decision: &GatewayControlDecision,
|
decision: &GatewayControlDecision,
|
||||||
plan_kind: &str,
|
plan_kind: &str,
|
||||||
) -> Result<LocalExecutionRequestOutcome, GatewayError> {
|
) -> Result<LocalExecutionRequestOutcome, GatewayError> {
|
||||||
let Some((attempt_source, candidate_count)) = build_local_image_sync_attempt_source_for_kind(
|
let Some((mut attempt_source, candidate_count)) =
|
||||||
state,
|
build_local_image_sync_attempt_source_for_kind(
|
||||||
parts,
|
state,
|
||||||
body_json,
|
parts,
|
||||||
body_base64,
|
body_json,
|
||||||
trace_id,
|
body_base64,
|
||||||
decision,
|
trace_id,
|
||||||
plan_kind,
|
decision,
|
||||||
)
|
plan_kind,
|
||||||
.await?
|
)
|
||||||
|
.await?
|
||||||
else {
|
else {
|
||||||
return Ok(LocalExecutionRequestOutcome::NoPath);
|
return Ok(LocalExecutionRequestOutcome::NoPath);
|
||||||
};
|
};
|
||||||
|
|
||||||
|
if openai_image_sync_heartbeat_enabled(state).await {
|
||||||
|
let mut attempts = Vec::new();
|
||||||
|
while let Some(attempt) = attempt_source.next_execution_attempt().await? {
|
||||||
|
attempts.push(attempt);
|
||||||
|
}
|
||||||
|
return Ok(LocalExecutionRequestOutcome::responded(
|
||||||
|
build_openai_image_sync_heartbeat_shell_response(
|
||||||
|
state.clone(),
|
||||||
|
parts.uri.path().to_string(),
|
||||||
|
trace_id.to_string(),
|
||||||
|
decision.clone(),
|
||||||
|
plan_kind.to_string(),
|
||||||
|
attempts,
|
||||||
|
)?,
|
||||||
|
));
|
||||||
|
}
|
||||||
|
|
||||||
let outcome = execute_sync_attempt_source::<AiSyncAttempt, _>(
|
let outcome = execute_sync_attempt_source::<AiSyncAttempt, _>(
|
||||||
state,
|
state,
|
||||||
parts,
|
parts,
|
||||||
@@ -674,3 +940,197 @@ pub(crate) fn parse_local_request_body(
|
|||||||
pub(crate) fn decision_payload_is_direct_execution(payload: &AiExecutionDecision) -> bool {
|
pub(crate) fn decision_payload_is_direct_execution(payload: &AiExecutionDecision) -> bool {
|
||||||
planner_decision_action(payload.action.as_str())
|
planner_decision_action(payload.action.as_str())
|
||||||
}
|
}
|
||||||
|
|
||||||
|
#[cfg(test)]
|
||||||
|
mod tests {
|
||||||
|
use super::*;
|
||||||
|
use std::sync::atomic::{AtomicUsize, Ordering};
|
||||||
|
use std::sync::Arc;
|
||||||
|
|
||||||
|
const TEST_OPENAI_IMAGE_SYNC_PLAN_KIND: &str = "openai_image_sync";
|
||||||
|
|
||||||
|
fn test_openai_image_heartbeat_decision() -> GatewayControlDecision {
|
||||||
|
GatewayControlDecision::synthetic(
|
||||||
|
"/v1/images/generations",
|
||||||
|
Some("ai_public".to_string()),
|
||||||
|
Some("openai".to_string()),
|
||||||
|
Some("image".to_string()),
|
||||||
|
Some("openai:image".to_string()),
|
||||||
|
)
|
||||||
|
.with_execution_runtime_candidate(true)
|
||||||
|
}
|
||||||
|
|
||||||
|
fn test_openai_image_heartbeat_plan(
|
||||||
|
endpoint_id: &str,
|
||||||
|
candidate_id: &str,
|
||||||
|
) -> aether_contracts::ExecutionPlan {
|
||||||
|
aether_contracts::ExecutionPlan {
|
||||||
|
request_id: "trace-image-heartbeat-retry".to_string(),
|
||||||
|
candidate_id: Some(candidate_id.to_string()),
|
||||||
|
provider_name: Some("OpenAI".to_string()),
|
||||||
|
provider_id: "provider-openai".to_string(),
|
||||||
|
endpoint_id: endpoint_id.to_string(),
|
||||||
|
key_id: "key-openai".to_string(),
|
||||||
|
method: "POST".to_string(),
|
||||||
|
url: "https://api.openai.com/v1/images/generations".to_string(),
|
||||||
|
headers: BTreeMap::new(),
|
||||||
|
content_type: Some("application/json".to_string()),
|
||||||
|
content_encoding: None,
|
||||||
|
body: aether_contracts::RequestBody::from_json(json!({"prompt": "test"})),
|
||||||
|
stream: false,
|
||||||
|
client_api_format: "openai:image".to_string(),
|
||||||
|
provider_api_format: "openai:image".to_string(),
|
||||||
|
model_name: Some("gpt-image-1".to_string()),
|
||||||
|
proxy: None,
|
||||||
|
transport_profile: None,
|
||||||
|
timeouts: None,
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
fn test_openai_image_heartbeat_attempt(
|
||||||
|
candidate_index: u32,
|
||||||
|
endpoint_id: &str,
|
||||||
|
candidate_id: &str,
|
||||||
|
) -> AiSyncAttempt {
|
||||||
|
AiSyncAttempt {
|
||||||
|
plan: test_openai_image_heartbeat_plan(endpoint_id, candidate_id),
|
||||||
|
report_kind: None,
|
||||||
|
report_context: Some(json!({
|
||||||
|
"candidate_index": candidate_index,
|
||||||
|
"retry_index": 0,
|
||||||
|
})),
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
fn test_openai_image_execution_result(
|
||||||
|
plan: &aether_contracts::ExecutionPlan,
|
||||||
|
status_code: u16,
|
||||||
|
body_json: Value,
|
||||||
|
) -> aether_contracts::ExecutionResult {
|
||||||
|
aether_contracts::ExecutionResult {
|
||||||
|
request_id: plan.request_id.clone(),
|
||||||
|
candidate_id: plan.candidate_id.clone(),
|
||||||
|
status_code,
|
||||||
|
headers: BTreeMap::from([(
|
||||||
|
CONTENT_TYPE.as_str().to_string(),
|
||||||
|
"application/json".to_string(),
|
||||||
|
)]),
|
||||||
|
body: Some(aether_contracts::ResponseBody {
|
||||||
|
json_body: Some(body_json),
|
||||||
|
body_bytes_b64: None,
|
||||||
|
}),
|
||||||
|
telemetry: Some(aether_contracts::ExecutionTelemetry {
|
||||||
|
ttfb_ms: None,
|
||||||
|
elapsed_ms: Some(10),
|
||||||
|
upstream_bytes: None,
|
||||||
|
}),
|
||||||
|
error: None,
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
#[tokio::test]
|
||||||
|
async fn openai_image_sync_heartbeat_success_body_is_unchanged() {
|
||||||
|
let response = Response::builder()
|
||||||
|
.status(StatusCode::OK)
|
||||||
|
.body(Body::from(r#"{"data":[{"b64_json":"x"}]}"#))
|
||||||
|
.expect("response should build");
|
||||||
|
|
||||||
|
let bytes = openai_image_sync_heartbeat_response_body_bytes(response).await;
|
||||||
|
let body: Value = serde_json::from_slice(&bytes).expect("body should decode");
|
||||||
|
|
||||||
|
assert_eq!(body, json!({"data": [{"b64_json": "x"}]}));
|
||||||
|
}
|
||||||
|
|
||||||
|
#[tokio::test]
|
||||||
|
async fn openai_image_sync_heartbeat_missing_config_defaults_disabled() {
|
||||||
|
let state = AppState::new().expect("state should build");
|
||||||
|
|
||||||
|
assert!(!openai_image_sync_heartbeat_enabled(&state).await);
|
||||||
|
}
|
||||||
|
|
||||||
|
#[tokio::test]
|
||||||
|
async fn openai_image_sync_heartbeat_error_body_includes_upstream_status() {
|
||||||
|
let response = Response::builder()
|
||||||
|
.status(StatusCode::TOO_MANY_REQUESTS)
|
||||||
|
.body(Body::from(
|
||||||
|
r#"{"error":{"type":"rate_limit","message":"slow down"}}"#,
|
||||||
|
))
|
||||||
|
.expect("response should build");
|
||||||
|
|
||||||
|
let bytes = openai_image_sync_heartbeat_response_body_bytes(response).await;
|
||||||
|
let body: Value = serde_json::from_slice(&bytes).expect("body should decode");
|
||||||
|
|
||||||
|
assert_eq!(body["error"]["type"], json!("rate_limit"));
|
||||||
|
assert_eq!(body["error"]["message"], json!("slow down"));
|
||||||
|
assert_eq!(body["error"]["upstream_status"], json!(429));
|
||||||
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn openai_image_sync_heartbeat_non_json_error_body_is_wrapped() {
|
||||||
|
let bytes =
|
||||||
|
openai_image_sync_heartbeat_error_body_from_response(502, b"bad gateway from upstream");
|
||||||
|
let body: Value = serde_json::from_slice(&bytes).expect("body should decode");
|
||||||
|
|
||||||
|
assert_eq!(body["error"]["type"], json!("upstream_error"));
|
||||||
|
assert_eq!(body["error"]["message"], json!("bad gateway from upstream"));
|
||||||
|
assert_eq!(body["error"]["upstream_status"], json!(502));
|
||||||
|
}
|
||||||
|
|
||||||
|
#[tokio::test]
|
||||||
|
async fn openai_image_sync_heartbeat_no_path_returns_json_error_body() {
|
||||||
|
let bytes =
|
||||||
|
openai_image_sync_heartbeat_final_bytes(Ok(LocalExecutionRequestOutcome::NoPath)).await;
|
||||||
|
let body: Value = serde_json::from_slice(&bytes).expect("body should decode");
|
||||||
|
|
||||||
|
assert_eq!(body["error"]["type"], json!("upstream_error"));
|
||||||
|
assert_eq!(body["error"]["upstream_status"], json!(503));
|
||||||
|
}
|
||||||
|
|
||||||
|
#[tokio::test]
|
||||||
|
async fn openai_image_sync_heartbeat_attempts_retry_first_candidate_then_return_second() {
|
||||||
|
let call_count = Arc::new(AtomicUsize::new(0));
|
||||||
|
let call_count_for_override = Arc::clone(&call_count);
|
||||||
|
let state = AppState::new()
|
||||||
|
.expect("state should build")
|
||||||
|
.with_execution_runtime_sync_override_for_tests(move |plan| {
|
||||||
|
call_count_for_override.fetch_add(1, Ordering::SeqCst);
|
||||||
|
if plan.endpoint_id == "endpoint-retry" {
|
||||||
|
Ok(test_openai_image_execution_result(
|
||||||
|
plan,
|
||||||
|
StatusCode::TOO_MANY_REQUESTS.as_u16(),
|
||||||
|
json!({"error": {"message": "retry this candidate"}}),
|
||||||
|
))
|
||||||
|
} else {
|
||||||
|
Ok(test_openai_image_execution_result(
|
||||||
|
plan,
|
||||||
|
StatusCode::OK.as_u16(),
|
||||||
|
json!({"data": [{"b64_json": "second-candidate"}]}),
|
||||||
|
))
|
||||||
|
}
|
||||||
|
});
|
||||||
|
let attempts = vec![
|
||||||
|
test_openai_image_heartbeat_attempt(0, "endpoint-retry", "candidate-retry"),
|
||||||
|
test_openai_image_heartbeat_attempt(1, "endpoint-success", "candidate-success"),
|
||||||
|
];
|
||||||
|
|
||||||
|
let outcome = execute_openai_image_sync_heartbeat_attempts(
|
||||||
|
state,
|
||||||
|
"/v1/images/generations".to_string(),
|
||||||
|
"trace-image-heartbeat-retry".to_string(),
|
||||||
|
test_openai_image_heartbeat_decision(),
|
||||||
|
TEST_OPENAI_IMAGE_SYNC_PLAN_KIND.to_string(),
|
||||||
|
attempts,
|
||||||
|
Instant::now(),
|
||||||
|
)
|
||||||
|
.await
|
||||||
|
.expect("heartbeat attempts should execute");
|
||||||
|
let LocalExecutionRequestOutcome::Responded(response) = outcome else {
|
||||||
|
panic!("second candidate should return a response");
|
||||||
|
};
|
||||||
|
let bytes = openai_image_sync_heartbeat_response_body_bytes(response).await;
|
||||||
|
let body: Value = serde_json::from_slice(&bytes).expect("body should decode");
|
||||||
|
|
||||||
|
assert_eq!(call_count.load(Ordering::SeqCst), 2);
|
||||||
|
assert_eq!(body, json!({"data": [{"b64_json": "second-candidate"}]}));
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|||||||
@@ -63,6 +63,7 @@
|
|||||||
:password-policy-level="systemConfig.password_policy_level"
|
:password-policy-level="systemConfig.password_policy_level"
|
||||||
:auto-delete-expired-keys="systemConfig.auto_delete_expired_keys"
|
:auto-delete-expired-keys="systemConfig.auto_delete_expired_keys"
|
||||||
:enable-format-conversion="systemConfig.enable_format_conversion"
|
:enable-format-conversion="systemConfig.enable_format_conversion"
|
||||||
|
:enable-openai-image-sync-heartbeat="systemConfig.enable_openai_image_sync_heartbeat"
|
||||||
:loading="basicConfigLoading"
|
:loading="basicConfigLoading"
|
||||||
:has-changes="hasBasicConfigChanges"
|
:has-changes="hasBasicConfigChanges"
|
||||||
@save="saveBasicConfig"
|
@save="saveBasicConfig"
|
||||||
@@ -72,6 +73,7 @@
|
|||||||
@update:password-policy-level="systemConfig.password_policy_level = $event"
|
@update:password-policy-level="systemConfig.password_policy_level = $event"
|
||||||
@update:auto-delete-expired-keys="systemConfig.auto_delete_expired_keys = $event"
|
@update:auto-delete-expired-keys="systemConfig.auto_delete_expired_keys = $event"
|
||||||
@update:enable-format-conversion="systemConfig.enable_format_conversion = $event"
|
@update:enable-format-conversion="systemConfig.enable_format_conversion = $event"
|
||||||
|
@update:enable-openai-image-sync-heartbeat="systemConfig.enable_openai_image_sync_heartbeat = $event"
|
||||||
/>
|
/>
|
||||||
|
|
||||||
<!-- 请求记录配置 -->
|
<!-- 请求记录配置 -->
|
||||||
|
|||||||
@@ -150,6 +150,27 @@
|
|||||||
</div>
|
</div>
|
||||||
</div>
|
</div>
|
||||||
</div>
|
</div>
|
||||||
|
|
||||||
|
<div class="flex items-center h-full">
|
||||||
|
<div class="flex items-center space-x-2">
|
||||||
|
<Checkbox
|
||||||
|
id="enable-openai-image-sync-heartbeat"
|
||||||
|
:checked="enableOpenaiImageSyncHeartbeat"
|
||||||
|
@update:checked="$emit('update:enableOpenaiImageSyncHeartbeat', $event)"
|
||||||
|
/>
|
||||||
|
<div>
|
||||||
|
<Label
|
||||||
|
for="enable-openai-image-sync-heartbeat"
|
||||||
|
class="cursor-pointer"
|
||||||
|
>
|
||||||
|
同步生图心跳
|
||||||
|
</Label>
|
||||||
|
<p class="text-xs text-muted-foreground">
|
||||||
|
开启后同步生图外层 HTTP 状态固定为 200,上游失败需读取响应体 error.upstream_status
|
||||||
|
</p>
|
||||||
|
</div>
|
||||||
|
</div>
|
||||||
|
</div>
|
||||||
</div>
|
</div>
|
||||||
</CardSection>
|
</CardSection>
|
||||||
</template>
|
</template>
|
||||||
@@ -173,6 +194,7 @@ defineProps<{
|
|||||||
passwordPolicyLevel: string
|
passwordPolicyLevel: string
|
||||||
autoDeleteExpiredKeys: boolean
|
autoDeleteExpiredKeys: boolean
|
||||||
enableFormatConversion: boolean
|
enableFormatConversion: boolean
|
||||||
|
enableOpenaiImageSyncHeartbeat: boolean
|
||||||
loading: boolean
|
loading: boolean
|
||||||
hasChanges: boolean
|
hasChanges: boolean
|
||||||
}>()
|
}>()
|
||||||
@@ -185,5 +207,6 @@ defineEmits<{
|
|||||||
'update:passwordPolicyLevel': [value: string]
|
'update:passwordPolicyLevel': [value: string]
|
||||||
'update:autoDeleteExpiredKeys': [value: boolean]
|
'update:autoDeleteExpiredKeys': [value: boolean]
|
||||||
'update:enableFormatConversion': [value: boolean]
|
'update:enableFormatConversion': [value: boolean]
|
||||||
|
'update:enableOpenaiImageSyncHeartbeat': [value: boolean]
|
||||||
}>()
|
}>()
|
||||||
</script>
|
</script>
|
||||||
|
|||||||
@@ -19,6 +19,8 @@ export interface SystemConfig {
|
|||||||
auto_delete_expired_keys: boolean
|
auto_delete_expired_keys: boolean
|
||||||
// 格式转换
|
// 格式转换
|
||||||
enable_format_conversion: boolean
|
enable_format_conversion: boolean
|
||||||
|
// 同步生图心跳
|
||||||
|
enable_openai_image_sync_heartbeat: boolean
|
||||||
// 请求记录
|
// 请求记录
|
||||||
request_record_level: string
|
request_record_level: string
|
||||||
max_request_body_size: number
|
max_request_body_size: number
|
||||||
@@ -58,6 +60,8 @@ const CONFIG_KEYS = [
|
|||||||
'auto_delete_expired_keys',
|
'auto_delete_expired_keys',
|
||||||
// 格式转换
|
// 格式转换
|
||||||
'enable_format_conversion',
|
'enable_format_conversion',
|
||||||
|
// 同步生图心跳
|
||||||
|
'enable_openai_image_sync_heartbeat',
|
||||||
// 请求记录
|
// 请求记录
|
||||||
'request_record_level',
|
'request_record_level',
|
||||||
'max_request_body_size',
|
'max_request_body_size',
|
||||||
@@ -98,6 +102,8 @@ function createDefaultConfig(): SystemConfig {
|
|||||||
auto_delete_expired_keys: false,
|
auto_delete_expired_keys: false,
|
||||||
// 格式转换
|
// 格式转换
|
||||||
enable_format_conversion: false,
|
enable_format_conversion: false,
|
||||||
|
// 同步生图心跳
|
||||||
|
enable_openai_image_sync_heartbeat: false,
|
||||||
// 请求记录
|
// 请求记录
|
||||||
request_record_level: 'basic',
|
request_record_level: 'basic',
|
||||||
max_request_body_size: 1048576,
|
max_request_body_size: 1048576,
|
||||||
@@ -160,7 +166,8 @@ export function useSystemConfig() {
|
|||||||
systemConfig.value.enable_registration !== originalConfig.value.enable_registration ||
|
systemConfig.value.enable_registration !== originalConfig.value.enable_registration ||
|
||||||
systemConfig.value.password_policy_level !== originalConfig.value.password_policy_level ||
|
systemConfig.value.password_policy_level !== originalConfig.value.password_policy_level ||
|
||||||
systemConfig.value.auto_delete_expired_keys !== originalConfig.value.auto_delete_expired_keys ||
|
systemConfig.value.auto_delete_expired_keys !== originalConfig.value.auto_delete_expired_keys ||
|
||||||
systemConfig.value.enable_format_conversion !== originalConfig.value.enable_format_conversion
|
systemConfig.value.enable_format_conversion !== originalConfig.value.enable_format_conversion ||
|
||||||
|
systemConfig.value.enable_openai_image_sync_heartbeat !== originalConfig.value.enable_openai_image_sync_heartbeat
|
||||||
)
|
)
|
||||||
})
|
})
|
||||||
|
|
||||||
@@ -340,6 +347,11 @@ export function useSystemConfig() {
|
|||||||
value: systemConfig.value.enable_format_conversion,
|
value: systemConfig.value.enable_format_conversion,
|
||||||
description: '全局格式转换开关:开启时强制允许所有提供商的格式转换',
|
description: '全局格式转换开关:开启时强制允许所有提供商的格式转换',
|
||||||
},
|
},
|
||||||
|
{
|
||||||
|
key: 'enable_openai_image_sync_heartbeat',
|
||||||
|
value: systemConfig.value.enable_openai_image_sync_heartbeat,
|
||||||
|
description: '同步生图心跳开关:开启后外层 HTTP 状态固定为 200,上游失败写入响应体',
|
||||||
|
},
|
||||||
]
|
]
|
||||||
|
|
||||||
await Promise.all(
|
await Promise.all(
|
||||||
@@ -356,6 +368,8 @@ export function useSystemConfig() {
|
|||||||
systemConfig.value.auto_delete_expired_keys
|
systemConfig.value.auto_delete_expired_keys
|
||||||
originalConfig.value.enable_format_conversion =
|
originalConfig.value.enable_format_conversion =
|
||||||
systemConfig.value.enable_format_conversion
|
systemConfig.value.enable_format_conversion
|
||||||
|
originalConfig.value.enable_openai_image_sync_heartbeat =
|
||||||
|
systemConfig.value.enable_openai_image_sync_heartbeat
|
||||||
}
|
}
|
||||||
success('基础配置已保存')
|
success('基础配置已保存')
|
||||||
} catch (err) {
|
} catch (err) {
|
||||||
|
|||||||
Reference in New Issue
Block a user