mirror of
https://github.com/fawney19/Aether.git
synced 2026-10-11 19:59:50 +08:00
feat: provider api_formats 可空继承、OpenAI 图片 edit/variation 与用量配额多项补强
- 鉴权: provider_api_keys.api_formats 改为可空,OAuth 托管 key 自动继承 provider endpoints 激活格式,相关 handler/测试同步更新 - 图片 planner: OpenAI 图片路由新增 edit/variation 操作并完善参数校验、响应合并与流式处理 - 用量: user me usage 返回区分 client_requested_stream/upstream_is_stream,前端 usage 列表筛选与展示增强 - 统计: stats_daily_model 新增 cache_creation_ephemeral_5m/1h tokens 字段与回填链路 - 配额/observability: quota repository 新增内存与 SQL 扩展,admin observability usage 字段扩充 - 其它: OAuth 导入/轮询收敛、provider 汇总与 pool admin 读写链路小修、新增 system_config 缓存与 provider template handler Closes #318 Co-authored-by: Entropy.Xu <[email protected]>
This commit is contained in:
@@ -29,7 +29,8 @@ pub(crate) use crate::ai_pipeline::{
|
||||
OPENAI_CLI_SYNC_FINALIZE_REPORT_KIND, OPENAI_CLI_SYNC_PLAN_KIND,
|
||||
OPENAI_CLI_SYNC_SUCCESS_REPORT_KIND, OPENAI_COMPACT_STREAM_PLAN_KIND,
|
||||
OPENAI_COMPACT_SYNC_ERROR_REPORT_KIND, OPENAI_COMPACT_SYNC_FINALIZE_REPORT_KIND,
|
||||
OPENAI_COMPACT_SYNC_PLAN_KIND, OPENAI_IMAGE_SYNC_FINALIZE_REPORT_KIND,
|
||||
OPENAI_COMPACT_SYNC_PLAN_KIND, OPENAI_IMAGE_STREAM_PLAN_KIND,
|
||||
OPENAI_IMAGE_STREAM_SUCCESS_REPORT_KIND, OPENAI_IMAGE_SYNC_FINALIZE_REPORT_KIND,
|
||||
OPENAI_IMAGE_SYNC_PLAN_KIND, OPENAI_VIDEO_CANCEL_SYNC_PLAN_KIND,
|
||||
OPENAI_VIDEO_CONTENT_PLAN_KIND, OPENAI_VIDEO_CREATE_SYNC_FINALIZE_REPORT_KIND,
|
||||
OPENAI_VIDEO_CREATE_SYNC_PLAN_KIND, OPENAI_VIDEO_DELETE_SYNC_PLAN_KIND,
|
||||
|
||||
@@ -2,12 +2,14 @@ use serde_json::Value;
|
||||
|
||||
use crate::ai_pipeline::adaptation::private_envelope::transform_provider_private_stream_line as transform_envelope_line;
|
||||
use crate::ai_pipeline::adaptation::KiroToClaudeCliStreamState;
|
||||
use crate::ai_pipeline::finalize::sse::encode_json_sse;
|
||||
use crate::ai_pipeline::finalize::standard::StreamingStandardConversionState;
|
||||
use crate::ai_pipeline::{resolve_finalize_stream_rewrite_mode, FinalizeStreamRewriteMode};
|
||||
use crate::GatewayError;
|
||||
|
||||
enum RewriteMode {
|
||||
EnvelopeUnwrap,
|
||||
OpenAiImage(OpenAiImageStreamState),
|
||||
Standard(StreamingStandardConversionState),
|
||||
KiroToClaudeCli(KiroToClaudeCliStreamState),
|
||||
}
|
||||
@@ -24,6 +26,9 @@ pub(crate) fn maybe_build_local_stream_rewriter<'a>(
|
||||
let report_context = report_context?;
|
||||
let mode = match resolve_finalize_stream_rewrite_mode(report_context)? {
|
||||
FinalizeStreamRewriteMode::EnvelopeUnwrap => RewriteMode::EnvelopeUnwrap,
|
||||
FinalizeStreamRewriteMode::OpenAiImage => {
|
||||
RewriteMode::OpenAiImage(OpenAiImageStreamState::default())
|
||||
}
|
||||
FinalizeStreamRewriteMode::Standard => {
|
||||
RewriteMode::Standard(StreamingStandardConversionState::default())
|
||||
}
|
||||
@@ -41,6 +46,9 @@ pub(crate) fn maybe_build_local_stream_rewriter<'a>(
|
||||
|
||||
impl LocalStreamRewriter<'_> {
|
||||
pub(crate) fn push_chunk(&mut self, chunk: &[u8]) -> Result<Vec<u8>, GatewayError> {
|
||||
if let RewriteMode::OpenAiImage(state) = &mut self.mode {
|
||||
return state.push_chunk(self.report_context, chunk);
|
||||
}
|
||||
if let RewriteMode::KiroToClaudeCli(state) = &mut self.mode {
|
||||
return state.push_chunk(self.report_context, chunk);
|
||||
}
|
||||
@@ -54,12 +62,16 @@ impl LocalStreamRewriter<'_> {
|
||||
}
|
||||
|
||||
pub(crate) fn finish(&mut self) -> Result<Vec<u8>, GatewayError> {
|
||||
if let RewriteMode::OpenAiImage(state) = &mut self.mode {
|
||||
return state.finish(self.report_context);
|
||||
}
|
||||
if let RewriteMode::KiroToClaudeCli(state) = &mut self.mode {
|
||||
return state.finish(self.report_context);
|
||||
}
|
||||
if self.buffered.is_empty() {
|
||||
match &mut self.mode {
|
||||
RewriteMode::Standard(state) => return state.finish(self.report_context),
|
||||
RewriteMode::OpenAiImage(_) => {}
|
||||
RewriteMode::KiroToClaudeCli(_) => {}
|
||||
RewriteMode::EnvelopeUnwrap => {}
|
||||
}
|
||||
@@ -71,6 +83,7 @@ impl LocalStreamRewriter<'_> {
|
||||
RewriteMode::Standard(state) => {
|
||||
output.extend(state.finish(self.report_context)?);
|
||||
}
|
||||
RewriteMode::OpenAiImage(_) => {}
|
||||
RewriteMode::KiroToClaudeCli(_) => {}
|
||||
RewriteMode::EnvelopeUnwrap => {}
|
||||
}
|
||||
@@ -81,12 +94,208 @@ impl LocalStreamRewriter<'_> {
|
||||
match &mut self.mode {
|
||||
RewriteMode::EnvelopeUnwrap => transform_envelope_line(self.report_context, line)
|
||||
.map_err(|err| GatewayError::Internal(err.to_string())),
|
||||
RewriteMode::OpenAiImage(_) => Ok(Vec::new()),
|
||||
RewriteMode::Standard(state) => state.transform_line(self.report_context, line),
|
||||
RewriteMode::KiroToClaudeCli(_) => Ok(Vec::new()),
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
#[derive(Default)]
|
||||
struct OpenAiImageStreamState {
|
||||
buffered: Vec<u8>,
|
||||
latest_image: Option<OpenAiImageFrame>,
|
||||
emitted_partial_count: u64,
|
||||
}
|
||||
|
||||
#[derive(Clone)]
|
||||
struct OpenAiImageFrame {
|
||||
b64_json: String,
|
||||
}
|
||||
|
||||
impl OpenAiImageStreamState {
|
||||
fn push_chunk(
|
||||
&mut self,
|
||||
report_context: &Value,
|
||||
chunk: &[u8],
|
||||
) -> Result<Vec<u8>, GatewayError> {
|
||||
self.buffered.extend_from_slice(chunk);
|
||||
let mut output = Vec::new();
|
||||
while let Some(block_end) = find_sse_block_end(&self.buffered) {
|
||||
let block = self.buffered.drain(..block_end).collect::<Vec<_>>();
|
||||
output.extend(self.transform_block(report_context, &block)?);
|
||||
drain_sse_separator(&mut self.buffered);
|
||||
}
|
||||
Ok(output)
|
||||
}
|
||||
|
||||
fn finish(&mut self, report_context: &Value) -> Result<Vec<u8>, GatewayError> {
|
||||
if self.buffered.is_empty() {
|
||||
return Ok(Vec::new());
|
||||
}
|
||||
let block = std::mem::take(&mut self.buffered);
|
||||
self.transform_block(report_context, &block)
|
||||
}
|
||||
|
||||
fn transform_block(
|
||||
&mut self,
|
||||
report_context: &Value,
|
||||
block: &[u8],
|
||||
) -> Result<Vec<u8>, GatewayError> {
|
||||
let text =
|
||||
std::str::from_utf8(block).map_err(|err| GatewayError::Internal(err.to_string()))?;
|
||||
let mut event_name = None::<String>;
|
||||
let mut data_lines = Vec::new();
|
||||
for raw_line in text.lines() {
|
||||
let line = raw_line.trim_end_matches('\r');
|
||||
if let Some(value) = line.strip_prefix("event:") {
|
||||
event_name = Some(value.trim().to_string());
|
||||
} else if let Some(value) = line.strip_prefix("data:") {
|
||||
data_lines.push(value.trim().to_string());
|
||||
}
|
||||
}
|
||||
let data = data_lines.join("\n");
|
||||
if data.is_empty() || data == "[DONE]" {
|
||||
return Ok(Vec::new());
|
||||
}
|
||||
let event: Value =
|
||||
serde_json::from_str(&data).map_err(|err| GatewayError::Internal(err.to_string()))?;
|
||||
let event_type = event
|
||||
.get("type")
|
||||
.and_then(Value::as_str)
|
||||
.or(event_name.as_deref())
|
||||
.unwrap_or_default();
|
||||
match event_type {
|
||||
"response.output_item.done" => self.handle_output_item_done(report_context, &event),
|
||||
"response.completed" => self.handle_completed(report_context, &event),
|
||||
_ => Ok(Vec::new()),
|
||||
}
|
||||
}
|
||||
|
||||
fn handle_output_item_done(
|
||||
&mut self,
|
||||
report_context: &Value,
|
||||
event: &Value,
|
||||
) -> Result<Vec<u8>, GatewayError> {
|
||||
let Some(item) = event.get("item").and_then(Value::as_object) else {
|
||||
return Ok(Vec::new());
|
||||
};
|
||||
if item.get("type").and_then(Value::as_str) != Some("image_generation_call") {
|
||||
return Ok(Vec::new());
|
||||
}
|
||||
let Some(result) = item.get("result").and_then(Value::as_str).map(str::trim) else {
|
||||
return Ok(Vec::new());
|
||||
};
|
||||
if result.is_empty() {
|
||||
return Ok(Vec::new());
|
||||
}
|
||||
self.latest_image = Some(OpenAiImageFrame {
|
||||
b64_json: result.to_string(),
|
||||
});
|
||||
|
||||
if requested_partial_images(report_context) == 0 {
|
||||
return Ok(Vec::new());
|
||||
}
|
||||
|
||||
let partial_image_index = event
|
||||
.get("output_index")
|
||||
.and_then(Value::as_u64)
|
||||
.unwrap_or(self.emitted_partial_count);
|
||||
self.emitted_partial_count = partial_image_index.saturating_add(1);
|
||||
|
||||
encode_json_sse(
|
||||
Some(image_partial_event_name(report_context)),
|
||||
&serde_json::json!({
|
||||
"type": image_partial_event_name(report_context),
|
||||
"b64_json": result,
|
||||
"partial_image_index": partial_image_index,
|
||||
}),
|
||||
)
|
||||
}
|
||||
|
||||
fn handle_completed(
|
||||
&mut self,
|
||||
report_context: &Value,
|
||||
event: &Value,
|
||||
) -> Result<Vec<u8>, GatewayError> {
|
||||
let Some(latest_image) = self.latest_image.clone() else {
|
||||
return Ok(Vec::new());
|
||||
};
|
||||
let usage = event
|
||||
.get("response")
|
||||
.and_then(Value::as_object)
|
||||
.and_then(|response| {
|
||||
response
|
||||
.get("tool_usage")
|
||||
.and_then(|value| value.get("image_gen"))
|
||||
.cloned()
|
||||
.or_else(|| response.get("usage").cloned())
|
||||
})
|
||||
.unwrap_or(Value::Null);
|
||||
|
||||
encode_json_sse(
|
||||
Some(image_completed_event_name(report_context)),
|
||||
&serde_json::json!({
|
||||
"type": image_completed_event_name(report_context),
|
||||
"b64_json": latest_image.b64_json,
|
||||
"usage": usage,
|
||||
}),
|
||||
)
|
||||
}
|
||||
}
|
||||
|
||||
fn requested_partial_images(report_context: &Value) -> u64 {
|
||||
report_context
|
||||
.get("image_request")
|
||||
.and_then(|value| value.get("partial_images"))
|
||||
.and_then(Value::as_u64)
|
||||
.unwrap_or(0)
|
||||
}
|
||||
|
||||
fn image_partial_event_name(report_context: &Value) -> &'static str {
|
||||
if image_request_operation(report_context) == Some("edit") {
|
||||
"image_edit.partial_image"
|
||||
} else {
|
||||
"image_generation.partial_image"
|
||||
}
|
||||
}
|
||||
|
||||
fn image_completed_event_name(report_context: &Value) -> &'static str {
|
||||
if image_request_operation(report_context) == Some("edit") {
|
||||
"image_edit.completed"
|
||||
} else {
|
||||
"image_generation.completed"
|
||||
}
|
||||
}
|
||||
|
||||
fn image_request_operation(report_context: &Value) -> Option<&str> {
|
||||
report_context
|
||||
.get("image_request")
|
||||
.and_then(|value| value.get("operation"))
|
||||
.and_then(Value::as_str)
|
||||
.map(str::trim)
|
||||
.filter(|value| !value.is_empty())
|
||||
}
|
||||
|
||||
fn find_sse_block_end(buffer: &[u8]) -> Option<usize> {
|
||||
buffer
|
||||
.windows(2)
|
||||
.position(|window| window == b"\n\n")
|
||||
.map(|index| index + 2)
|
||||
.or_else(|| {
|
||||
buffer
|
||||
.windows(4)
|
||||
.position(|window| window == b"\r\n\r\n")
|
||||
.map(|index| index + 4)
|
||||
})
|
||||
}
|
||||
|
||||
fn drain_sse_separator(buffer: &mut Vec<u8>) {
|
||||
while matches!(buffer.first(), Some(b'\n' | b'\r')) {
|
||||
buffer.remove(0);
|
||||
}
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
#[path = "../tests_stream.rs"]
|
||||
mod tests;
|
||||
|
||||
@@ -1,4 +1,5 @@
|
||||
use crate::ai_pipeline::GatewayControlDecision;
|
||||
use crate::ai_pipeline::CODEX_OPENAI_IMAGE_DEFAULT_OUTPUT_FORMAT;
|
||||
use crate::ai_pipeline::{build_generated_tool_call_id, canonicalize_tool_arguments};
|
||||
use crate::{usage::GatewaySyncReportRequest, GatewayError};
|
||||
use base64::Engine as _;
|
||||
@@ -105,6 +106,20 @@ fn maybe_build_local_openai_image_sync_finalize_response(
|
||||
let Some(body_base64) = payload.body_base64.as_deref() else {
|
||||
return Ok(None);
|
||||
};
|
||||
let response_format = report_context
|
||||
.get("image_request")
|
||||
.and_then(|value| value.get("response_format"))
|
||||
.and_then(serde_json::Value::as_str)
|
||||
.map(str::trim)
|
||||
.filter(|value| !value.is_empty())
|
||||
.unwrap_or("url");
|
||||
let default_output_format = report_context
|
||||
.get("image_request")
|
||||
.and_then(|value| value.get("output_format"))
|
||||
.and_then(serde_json::Value::as_str)
|
||||
.map(str::trim)
|
||||
.filter(|value| !value.is_empty())
|
||||
.unwrap_or(CODEX_OPENAI_IMAGE_DEFAULT_OUTPUT_FORMAT);
|
||||
let body_bytes = base64::engine::general_purpose::STANDARD
|
||||
.decode(body_base64)
|
||||
.map_err(|err| GatewayError::Internal(err.to_string()))?;
|
||||
@@ -157,6 +172,7 @@ fn maybe_build_local_openai_image_sync_finalize_response(
|
||||
};
|
||||
images.push(serde_json::json!({
|
||||
"b64_json": result,
|
||||
"output_format": item.get("output_format").cloned().unwrap_or(serde_json::Value::String(default_output_format.to_string())),
|
||||
"revised_prompt": item.get("revised_prompt").cloned().unwrap_or(serde_json::Value::Null),
|
||||
}));
|
||||
}
|
||||
@@ -191,13 +207,46 @@ fn maybe_build_local_openai_image_sync_finalize_response(
|
||||
.iter()
|
||||
.map(|image| serde_json::json!({
|
||||
"type": "image_generation_call",
|
||||
"output_format": image.get("output_format").cloned().unwrap_or(serde_json::Value::Null),
|
||||
"revised_prompt": image.get("revised_prompt").cloned().unwrap_or(serde_json::Value::Null),
|
||||
}))
|
||||
.collect::<Vec<_>>(),
|
||||
});
|
||||
let client_images = images
|
||||
.iter()
|
||||
.map(|image| {
|
||||
let revised_prompt = image
|
||||
.get("revised_prompt")
|
||||
.cloned()
|
||||
.unwrap_or(serde_json::Value::Null);
|
||||
let b64_json = image
|
||||
.get("b64_json")
|
||||
.and_then(serde_json::Value::as_str)
|
||||
.unwrap_or_default();
|
||||
let output_format = image
|
||||
.get("output_format")
|
||||
.and_then(serde_json::Value::as_str)
|
||||
.unwrap_or(default_output_format);
|
||||
if response_format.eq_ignore_ascii_case("b64_json") {
|
||||
serde_json::json!({
|
||||
"b64_json": b64_json,
|
||||
"revised_prompt": revised_prompt,
|
||||
})
|
||||
} else {
|
||||
serde_json::json!({
|
||||
"url": format!(
|
||||
"data:{};base64,{}",
|
||||
image_output_mime_type(output_format),
|
||||
b64_json
|
||||
),
|
||||
"revised_prompt": revised_prompt,
|
||||
})
|
||||
}
|
||||
})
|
||||
.collect::<Vec<_>>();
|
||||
let client_body_json = serde_json::json!({
|
||||
"created": created.unwrap_or_default(),
|
||||
"data": images,
|
||||
"data": client_images,
|
||||
"usage": provider_body_json.get("usage").cloned().unwrap_or(serde_json::Value::Null),
|
||||
});
|
||||
|
||||
@@ -210,6 +259,14 @@ fn maybe_build_local_openai_image_sync_finalize_response(
|
||||
)?))
|
||||
}
|
||||
|
||||
fn image_output_mime_type(output_format: &str) -> &'static str {
|
||||
match output_format.trim().to_ascii_lowercase().as_str() {
|
||||
"jpeg" | "jpg" => "image/jpeg",
|
||||
"webp" => "image/webp",
|
||||
_ => "image/png",
|
||||
}
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
#[path = "../tests_sync.rs"]
|
||||
mod tests;
|
||||
|
||||
@@ -29,6 +29,94 @@ fn antigravity_stream_rewriter_unwraps_and_injects_tool_ids() {
|
||||
assert!(output_text.contains("\"modelVersion\":\"claude-sonnet-4-5\""));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn openai_image_stream_rewriter_emits_completed_event_for_generate() {
|
||||
let report_context = json!({
|
||||
"provider_api_format": "openai:image",
|
||||
"client_api_format": "openai:image",
|
||||
"needs_conversion": false,
|
||||
"image_request": {
|
||||
"operation": "generate"
|
||||
}
|
||||
});
|
||||
let mut rewriter =
|
||||
maybe_build_local_stream_rewriter(Some(&report_context)).expect("rewriter should exist");
|
||||
|
||||
let first = rewriter
|
||||
.push_chunk(
|
||||
concat!(
|
||||
"event: response.output_item.done\n",
|
||||
"data: {\"type\":\"response.output_item.done\",\"output_index\":0,\"item\":{\"id\":\"ig_123\",\"type\":\"image_generation_call\",\"result\":\"aGVsbG8=\"}}\n\n"
|
||||
)
|
||||
.as_bytes(),
|
||||
)
|
||||
.expect("rewrite should succeed");
|
||||
assert!(first.is_empty());
|
||||
|
||||
let second = rewriter
|
||||
.push_chunk(
|
||||
concat!(
|
||||
"event: response.completed\n",
|
||||
"data: {\"type\":\"response.completed\",\"response\":{\"tool_usage\":{\"image_gen\":{\"input_tokens\":1,\"output_tokens\":2,\"total_tokens\":3}}}}\n\n"
|
||||
)
|
||||
.as_bytes(),
|
||||
)
|
||||
.expect("rewrite should succeed");
|
||||
let output_text = utf8(second);
|
||||
assert!(output_text.contains("event: image_generation.completed"));
|
||||
assert!(output_text.contains("\"type\":\"image_generation.completed\""));
|
||||
assert!(output_text.contains("\"b64_json\":\"aGVsbG8=\""));
|
||||
assert!(output_text.contains("\"input_tokens\":1"));
|
||||
assert!(!output_text.contains("data: [DONE]"));
|
||||
assert!(rewriter.finish().expect("finish should succeed").is_empty());
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn openai_image_stream_rewriter_emits_partial_and_completed_events_for_edit() {
|
||||
let report_context = json!({
|
||||
"provider_api_format": "openai:image",
|
||||
"client_api_format": "openai:image",
|
||||
"needs_conversion": false,
|
||||
"image_request": {
|
||||
"operation": "edit",
|
||||
"partial_images": 2
|
||||
}
|
||||
});
|
||||
let mut rewriter =
|
||||
maybe_build_local_stream_rewriter(Some(&report_context)).expect("rewriter should exist");
|
||||
|
||||
let partial = rewriter
|
||||
.push_chunk(
|
||||
concat!(
|
||||
"event: response.output_item.done\n",
|
||||
"data: {\"type\":\"response.output_item.done\",\"output_index\":1,\"item\":{\"id\":\"ig_edit_123\",\"type\":\"image_generation_call\",\"result\":\"d29ybGQ=\"}}\n\n"
|
||||
)
|
||||
.as_bytes(),
|
||||
)
|
||||
.expect("rewrite should succeed");
|
||||
let partial_text = utf8(partial);
|
||||
assert!(partial_text.contains("event: image_edit.partial_image"));
|
||||
assert!(partial_text.contains("\"type\":\"image_edit.partial_image\""));
|
||||
assert!(partial_text.contains("\"b64_json\":\"d29ybGQ=\""));
|
||||
assert!(partial_text.contains("\"partial_image_index\":1"));
|
||||
|
||||
let completed = rewriter
|
||||
.push_chunk(
|
||||
concat!(
|
||||
"event: response.completed\n",
|
||||
"data: {\"type\":\"response.completed\",\"response\":{\"usage\":{\"input_tokens\":4,\"output_tokens\":5,\"total_tokens\":9}}}\n\n"
|
||||
)
|
||||
.as_bytes(),
|
||||
)
|
||||
.expect("rewrite should succeed");
|
||||
let completed_text = utf8(completed);
|
||||
assert!(completed_text.contains("event: image_edit.completed"));
|
||||
assert!(completed_text.contains("\"type\":\"image_edit.completed\""));
|
||||
assert!(completed_text.contains("\"b64_json\":\"d29ybGQ=\""));
|
||||
assert!(completed_text.contains("\"total_tokens\":9"));
|
||||
assert!(rewriter.finish().expect("finish should succeed").is_empty());
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn gemini_cli_v1internal_stream_rewriter_unwraps_response_object() {
|
||||
let report_context = json!({
|
||||
|
||||
@@ -1773,7 +1773,12 @@ async fn local_finalize_handles_openai_image_stream_response_from_output_item_do
|
||||
"client_api_format": "openai:image",
|
||||
"provider_api_format": "openai:image",
|
||||
"model": "gpt-image-2",
|
||||
"mapped_model": "gpt-5.4"
|
||||
"mapped_model": "gpt-5.4",
|
||||
"image_request": {
|
||||
"operation": "generate",
|
||||
"response_format": "b64_json",
|
||||
"output_format": "png"
|
||||
}
|
||||
})),
|
||||
status_code: 200,
|
||||
headers: BTreeMap::from([(
|
||||
@@ -1830,3 +1835,63 @@ async fn local_finalize_handles_openai_image_stream_response_from_output_item_do
|
||||
"aGVsbG8="
|
||||
);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn local_finalize_handles_openai_image_stream_response_with_url_response_format() {
|
||||
let payload = GatewaySyncReportRequest {
|
||||
trace_id: "trace-openai-image-finalize-url-123".to_string(),
|
||||
report_kind: "openai_image_sync_finalize".to_string(),
|
||||
report_context: Some(json!({
|
||||
"client_api_format": "openai:image",
|
||||
"provider_api_format": "openai:image",
|
||||
"model": "gpt-image-1",
|
||||
"mapped_model": "gpt-5.4",
|
||||
"image_request": {
|
||||
"operation": "generate",
|
||||
"response_format": "url",
|
||||
"output_format": "webp"
|
||||
}
|
||||
})),
|
||||
status_code: 200,
|
||||
headers: BTreeMap::from([(
|
||||
"content-type".to_string(),
|
||||
"text/event-stream".to_string(),
|
||||
)]),
|
||||
body_json: None,
|
||||
client_body_json: None,
|
||||
body_base64: Some(base64::engine::general_purpose::STANDARD.encode(
|
||||
concat!(
|
||||
"event: response.created\n",
|
||||
"data: {\"type\":\"response.created\",\"response\":{\"id\":\"resp_img_url_123\",\"object\":\"response\",\"created_at\":1776839946,\"status\":\"in_progress\",\"model\":\"gpt-5.4\"}}\n\n",
|
||||
"event: response.output_item.done\n",
|
||||
"data: {\"type\":\"response.output_item.done\",\"output_index\":0,\"item\":{\"id\":\"ig_url_123\",\"type\":\"image_generation_call\",\"status\":\"completed\",\"output_format\":\"webp\",\"revised_prompt\":\"revised webp prompt\",\"result\":\"aGVsbG8=\"}}\n\n",
|
||||
"event: response.completed\n",
|
||||
"data: {\"type\":\"response.completed\",\"response\":{\"id\":\"resp_img_url_123\",\"object\":\"response\",\"model\":\"gpt-5.4\",\"status\":\"completed\",\"output\":[],\"tool_usage\":{\"image_gen\":{\"input_tokens\":11,\"output_tokens\":22,\"total_tokens\":33}}}}\n\n"
|
||||
)
|
||||
.as_bytes(),
|
||||
)),
|
||||
telemetry: None,
|
||||
};
|
||||
|
||||
let outcome = maybe_build_local_core_sync_finalize_response(
|
||||
"trace-openai-image-finalize-url-123",
|
||||
&test_decision(),
|
||||
&payload,
|
||||
)
|
||||
.expect("image finalize should succeed")
|
||||
.expect("image finalize should match");
|
||||
|
||||
let response_body = to_bytes(outcome.response.into_body(), usize::MAX)
|
||||
.await
|
||||
.expect("response body should read");
|
||||
let response_json: serde_json::Value =
|
||||
serde_json::from_slice(&response_body).expect("response should be json");
|
||||
assert_eq!(
|
||||
response_json["data"][0]["url"],
|
||||
"data:image/webp;base64,aGVsbG8="
|
||||
);
|
||||
assert_eq!(
|
||||
response_json["data"][0]["revised_prompt"],
|
||||
"revised webp prompt"
|
||||
);
|
||||
}
|
||||
|
||||
@@ -26,6 +26,7 @@ pub(crate) use self::planner::{
|
||||
build_gemini_stream_plan_from_decision, build_gemini_sync_plan_from_decision,
|
||||
build_local_gemini_files_stream_plan_and_reports_for_kind,
|
||||
build_local_gemini_files_sync_plan_and_reports_for_kind,
|
||||
build_local_image_stream_plan_and_reports_for_kind,
|
||||
build_local_image_sync_plan_and_reports_for_kind,
|
||||
build_local_openai_chat_stream_plan_and_reports_for_kind,
|
||||
build_local_openai_chat_sync_plan_and_reports_for_kind,
|
||||
@@ -38,9 +39,9 @@ pub(crate) use self::planner::{
|
||||
build_standard_stream_plan_from_decision, build_standard_sync_plan_from_decision,
|
||||
extract_pool_sticky_session_token, maybe_build_stream_decision_payload,
|
||||
maybe_build_stream_plan_payload, maybe_build_sync_decision_payload,
|
||||
maybe_build_sync_plan_payload, set_local_openai_chat_execution_exhausted_diagnostic,
|
||||
GatewayAuthApiKeySnapshot, GatewayProviderTransportSnapshot, LocalResolvedOAuthRequestAuth,
|
||||
PlannerAppState,
|
||||
maybe_build_sync_plan_payload, planner_is_matching_stream_request,
|
||||
set_local_openai_chat_execution_exhausted_diagnostic, GatewayAuthApiKeySnapshot,
|
||||
GatewayProviderTransportSnapshot, LocalResolvedOAuthRequestAuth, PlannerAppState,
|
||||
};
|
||||
pub(crate) use self::pure::*;
|
||||
pub(crate) use crate::control::GatewayControlDecision;
|
||||
|
||||
@@ -1,5 +1,6 @@
|
||||
use std::sync::Arc;
|
||||
|
||||
use aether_provider_transport::provider_types::provider_type_is_fixed;
|
||||
use tracing::warn;
|
||||
|
||||
use aether_scheduler_core::SchedulerMinimalCandidateSelectionCandidate;
|
||||
@@ -295,6 +296,22 @@ fn transport_key_supports_api_format(
|
||||
transport: &GatewayProviderTransportSnapshot,
|
||||
endpoint_api_format: &str,
|
||||
) -> bool {
|
||||
let provider_type = transport.provider.provider_type.trim();
|
||||
let auth_type = transport.key.auth_type.trim();
|
||||
let inherits_provider_api_formats = provider_type_is_fixed(provider_type)
|
||||
&& (auth_type.eq_ignore_ascii_case("oauth")
|
||||
|| (provider_type.eq_ignore_ascii_case("kiro")
|
||||
&& auth_type.eq_ignore_ascii_case("bearer")
|
||||
&& transport
|
||||
.key
|
||||
.decrypted_auth_config
|
||||
.as_deref()
|
||||
.map(str::trim)
|
||||
.is_some_and(|value| !value.is_empty())));
|
||||
if inherits_provider_api_formats {
|
||||
return true;
|
||||
}
|
||||
|
||||
match transport.key.api_formats.as_deref() {
|
||||
None => true,
|
||||
Some(formats) => formats
|
||||
|
||||
@@ -10,9 +10,10 @@ pub(crate) use crate::ai_pipeline::contracts::{
|
||||
GEMINI_FILES_LIST_PLAN_KIND, GEMINI_FILES_UPLOAD_PLAN_KIND, GEMINI_VIDEO_CANCEL_SYNC_PLAN_KIND,
|
||||
GEMINI_VIDEO_CREATE_SYNC_PLAN_KIND, OPENAI_CHAT_STREAM_PLAN_KIND, OPENAI_CHAT_SYNC_PLAN_KIND,
|
||||
OPENAI_CLI_STREAM_PLAN_KIND, OPENAI_CLI_SYNC_PLAN_KIND, OPENAI_COMPACT_STREAM_PLAN_KIND,
|
||||
OPENAI_COMPACT_SYNC_PLAN_KIND, OPENAI_IMAGE_SYNC_PLAN_KIND, OPENAI_VIDEO_CANCEL_SYNC_PLAN_KIND,
|
||||
OPENAI_VIDEO_CONTENT_PLAN_KIND, OPENAI_VIDEO_CREATE_SYNC_PLAN_KIND,
|
||||
OPENAI_VIDEO_DELETE_SYNC_PLAN_KIND, OPENAI_VIDEO_REMIX_SYNC_PLAN_KIND,
|
||||
OPENAI_COMPACT_SYNC_PLAN_KIND, OPENAI_IMAGE_STREAM_PLAN_KIND, OPENAI_IMAGE_SYNC_PLAN_KIND,
|
||||
OPENAI_VIDEO_CANCEL_SYNC_PLAN_KIND, OPENAI_VIDEO_CONTENT_PLAN_KIND,
|
||||
OPENAI_VIDEO_CREATE_SYNC_PLAN_KIND, OPENAI_VIDEO_DELETE_SYNC_PLAN_KIND,
|
||||
OPENAI_VIDEO_REMIX_SYNC_PLAN_KIND,
|
||||
};
|
||||
use crate::ai_pipeline::GatewayControlDecision;
|
||||
use crate::ai_pipeline::{
|
||||
|
||||
@@ -7,9 +7,10 @@ use crate::ai_pipeline::planner::common::{
|
||||
GEMINI_FILES_GET_PLAN_KIND, GEMINI_FILES_LIST_PLAN_KIND, GEMINI_VIDEO_CANCEL_SYNC_PLAN_KIND,
|
||||
GEMINI_VIDEO_CREATE_SYNC_PLAN_KIND, OPENAI_CHAT_STREAM_PLAN_KIND, OPENAI_CHAT_SYNC_PLAN_KIND,
|
||||
OPENAI_CLI_STREAM_PLAN_KIND, OPENAI_CLI_SYNC_PLAN_KIND, OPENAI_COMPACT_STREAM_PLAN_KIND,
|
||||
OPENAI_COMPACT_SYNC_PLAN_KIND, OPENAI_IMAGE_SYNC_PLAN_KIND, OPENAI_VIDEO_CANCEL_SYNC_PLAN_KIND,
|
||||
OPENAI_VIDEO_CONTENT_PLAN_KIND, OPENAI_VIDEO_CREATE_SYNC_PLAN_KIND,
|
||||
OPENAI_VIDEO_DELETE_SYNC_PLAN_KIND, OPENAI_VIDEO_REMIX_SYNC_PLAN_KIND,
|
||||
OPENAI_COMPACT_SYNC_PLAN_KIND, OPENAI_IMAGE_STREAM_PLAN_KIND, OPENAI_IMAGE_SYNC_PLAN_KIND,
|
||||
OPENAI_VIDEO_CANCEL_SYNC_PLAN_KIND, OPENAI_VIDEO_CONTENT_PLAN_KIND,
|
||||
OPENAI_VIDEO_CREATE_SYNC_PLAN_KIND, OPENAI_VIDEO_DELETE_SYNC_PLAN_KIND,
|
||||
OPENAI_VIDEO_REMIX_SYNC_PLAN_KIND,
|
||||
};
|
||||
use crate::ai_pipeline::planner::plan_builders::{
|
||||
build_gemini_stream_plan_from_decision, build_gemini_sync_plan_from_decision,
|
||||
@@ -63,13 +64,20 @@ pub(crate) async fn maybe_build_stream_plan_payload_impl(
|
||||
trace_id: &str,
|
||||
decision: &GatewayControlDecision,
|
||||
body_json: &serde_json::Value,
|
||||
body_base64: Option<&str>,
|
||||
) -> Result<Option<GatewayControlPlanResponse>, GatewayError> {
|
||||
let Some(plan_kind) = resolve_stream_plan_kind(parts, decision) else {
|
||||
return Ok(None);
|
||||
};
|
||||
let Some(payload) =
|
||||
super::maybe_build_stream_decision_payload(state, parts, trace_id, decision, body_json)
|
||||
.await?
|
||||
let Some(payload) = super::maybe_build_stream_decision_payload(
|
||||
state,
|
||||
parts,
|
||||
trace_id,
|
||||
decision,
|
||||
body_json,
|
||||
body_base64,
|
||||
)
|
||||
.await?
|
||||
else {
|
||||
return Ok(None);
|
||||
};
|
||||
@@ -132,6 +140,9 @@ fn build_stream_plan_payload_from_decision(
|
||||
OPENAI_CLI_STREAM_PLAN_KIND => {
|
||||
build_openai_cli_stream_plan_from_decision(parts, body_json, payload, false)?
|
||||
}
|
||||
OPENAI_IMAGE_STREAM_PLAN_KIND => {
|
||||
build_standard_stream_plan_from_decision(parts, body_json, payload, false)?
|
||||
}
|
||||
OPENAI_COMPACT_STREAM_PLAN_KIND => {
|
||||
build_openai_cli_stream_plan_from_decision(parts, body_json, payload, true)?
|
||||
}
|
||||
|
||||
@@ -13,6 +13,7 @@ pub(crate) use super::passthrough::{
|
||||
};
|
||||
pub(crate) use super::specialized::{
|
||||
maybe_build_stream_local_gemini_files_decision_payload,
|
||||
maybe_build_stream_local_image_decision_payload,
|
||||
maybe_build_sync_local_gemini_files_decision_payload,
|
||||
maybe_build_sync_local_image_decision_payload, maybe_build_sync_local_video_decision_payload,
|
||||
};
|
||||
|
||||
@@ -18,12 +18,13 @@ pub(crate) async fn maybe_build_stream_decision_payload(
|
||||
trace_id: &str,
|
||||
decision: &GatewayControlDecision,
|
||||
body_json: &serde_json::Value,
|
||||
body_base64: Option<&str>,
|
||||
) -> Result<Option<GatewayControlSyncDecisionResponse>, GatewayError> {
|
||||
let Some(plan_kind) = resolve_execution_runtime_stream_plan_kind(parts, decision) else {
|
||||
return Ok(None);
|
||||
};
|
||||
|
||||
if !is_matching_stream_request(plan_kind, parts, body_json) {
|
||||
if !is_matching_stream_request(plan_kind, parts, body_json, body_base64) {
|
||||
return Ok(None);
|
||||
}
|
||||
|
||||
@@ -35,6 +36,20 @@ pub(crate) async fn maybe_build_stream_decision_payload(
|
||||
return Ok(Some(payload));
|
||||
}
|
||||
|
||||
if let Some(payload) = super::maybe_build_stream_local_image_decision_payload(
|
||||
state,
|
||||
parts,
|
||||
body_json,
|
||||
body_base64,
|
||||
trace_id,
|
||||
decision,
|
||||
plan_kind,
|
||||
)
|
||||
.await?
|
||||
{
|
||||
return Ok(Some(payload));
|
||||
}
|
||||
|
||||
if let Some(payload) = super::maybe_build_stream_local_decision_payload(
|
||||
state, parts, trace_id, decision, body_json, plan_kind,
|
||||
)
|
||||
|
||||
@@ -36,9 +36,11 @@ pub(crate) use self::plan_builders::{
|
||||
build_passthrough_sync_plan_from_decision, build_standard_stream_plan_from_decision,
|
||||
build_standard_sync_plan_from_decision, LocalStreamPlanAndReport, LocalSyncPlanAndReport,
|
||||
};
|
||||
pub(crate) use self::route::is_matching_stream_request as planner_is_matching_stream_request;
|
||||
pub(crate) use self::specialized::{
|
||||
build_local_gemini_files_stream_plan_and_reports_for_kind,
|
||||
build_local_gemini_files_sync_plan_and_reports_for_kind,
|
||||
build_local_image_stream_plan_and_reports_for_kind,
|
||||
build_local_image_sync_plan_and_reports_for_kind,
|
||||
build_local_video_sync_plan_and_reports_for_kind,
|
||||
};
|
||||
@@ -83,8 +85,17 @@ pub(crate) async fn maybe_build_stream_decision_payload(
|
||||
trace_id: &str,
|
||||
decision: &GatewayControlDecision,
|
||||
body_json: &serde_json::Value,
|
||||
body_base64: Option<&str>,
|
||||
) -> Result<Option<GatewayControlSyncDecisionResponse>, GatewayError> {
|
||||
decision::maybe_build_stream_decision_payload(state, parts, trace_id, decision, body_json).await
|
||||
decision::maybe_build_stream_decision_payload(
|
||||
state,
|
||||
parts,
|
||||
trace_id,
|
||||
decision,
|
||||
body_json,
|
||||
body_base64,
|
||||
)
|
||||
.await
|
||||
}
|
||||
|
||||
pub(crate) async fn maybe_build_sync_plan_payload(
|
||||
@@ -114,7 +125,15 @@ pub(crate) async fn maybe_build_stream_plan_payload(
|
||||
trace_id: &str,
|
||||
decision: &GatewayControlDecision,
|
||||
body_json: &serde_json::Value,
|
||||
body_base64: Option<&str>,
|
||||
) -> Result<Option<GatewayControlPlanResponse>, GatewayError> {
|
||||
decision::maybe_build_stream_plan_payload_impl(state, parts, trace_id, decision, body_json)
|
||||
.await
|
||||
decision::maybe_build_stream_plan_payload_impl(
|
||||
state,
|
||||
parts,
|
||||
trace_id,
|
||||
decision,
|
||||
body_json,
|
||||
body_base64,
|
||||
)
|
||||
.await
|
||||
}
|
||||
|
||||
@@ -101,6 +101,7 @@ pub(crate) async fn maybe_build_local_same_format_provider_decision_payload_for_
|
||||
original_headers: &parts.headers,
|
||||
original_request_body_json: Some(body_json),
|
||||
original_request_body_base64: None,
|
||||
client_requested_stream: spec_metadata.require_streaming,
|
||||
has_envelope: resolved.is_kiro || resolved.is_antigravity,
|
||||
needs_conversion: false,
|
||||
extra_fields,
|
||||
|
||||
@@ -26,6 +26,7 @@ pub(crate) struct LocalExecutionReportContextParts<'a> {
|
||||
pub(crate) original_headers: &'a http::HeaderMap,
|
||||
pub(crate) original_request_body_json: Option<&'a Value>,
|
||||
pub(crate) original_request_body_base64: Option<&'a str>,
|
||||
pub(crate) client_requested_stream: bool,
|
||||
pub(crate) has_envelope: bool,
|
||||
pub(crate) needs_conversion: bool,
|
||||
pub(crate) extra_fields: Map<String, Value>,
|
||||
@@ -117,6 +118,10 @@ pub(crate) fn build_local_execution_report_context(
|
||||
)
|
||||
.unwrap_or(Value::Null),
|
||||
);
|
||||
object.insert(
|
||||
"client_requested_stream".to_string(),
|
||||
Value::Bool(parts.client_requested_stream),
|
||||
);
|
||||
object.insert("has_envelope".to_string(), Value::Bool(parts.has_envelope));
|
||||
object.insert(
|
||||
"needs_conversion".to_string(),
|
||||
|
||||
@@ -1,3 +1,4 @@
|
||||
use super::specialized::is_openai_image_stream_request;
|
||||
use crate::ai_pipeline::GatewayControlDecision;
|
||||
use crate::ai_pipeline::{
|
||||
is_matching_stream_request as is_matching_stream_request_impl,
|
||||
@@ -5,6 +6,7 @@ use crate::ai_pipeline::{
|
||||
resolve_execution_runtime_sync_plan_kind as resolve_execution_runtime_sync_plan_kind_impl,
|
||||
supports_stream_scheduler_decision_kind as supports_stream_scheduler_decision_kind_impl,
|
||||
supports_sync_scheduler_decision_kind as supports_sync_scheduler_decision_kind_impl,
|
||||
OPENAI_IMAGE_STREAM_PLAN_KIND,
|
||||
};
|
||||
|
||||
pub(crate) fn resolve_execution_runtime_stream_plan_kind(
|
||||
@@ -37,7 +39,11 @@ pub(crate) fn is_matching_stream_request(
|
||||
plan_kind: &str,
|
||||
parts: &http::request::Parts,
|
||||
body_json: &serde_json::Value,
|
||||
body_base64: Option<&str>,
|
||||
) -> bool {
|
||||
if plan_kind == OPENAI_IMAGE_STREAM_PLAN_KIND {
|
||||
return is_openai_image_stream_request(parts, body_json, body_base64);
|
||||
}
|
||||
is_matching_stream_request_impl(plan_kind, parts.uri.path(), body_json)
|
||||
}
|
||||
|
||||
@@ -52,6 +58,7 @@ pub(crate) fn supports_stream_scheduler_decision_kind(plan_kind: &str) -> bool {
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
use axum::http::{Method, Request};
|
||||
use base64::Engine as _;
|
||||
|
||||
use super::{
|
||||
is_matching_stream_request, resolve_execution_runtime_stream_plan_kind,
|
||||
@@ -108,15 +115,45 @@ mod tests {
|
||||
"openai_chat_stream",
|
||||
&parts,
|
||||
&serde_json::json!({"stream": false}),
|
||||
None,
|
||||
));
|
||||
assert!(is_matching_stream_request(
|
||||
"openai_chat_stream",
|
||||
&parts,
|
||||
&serde_json::json!({"stream": true}),
|
||||
None,
|
||||
));
|
||||
assert!(supports_sync_scheduler_decision_kind("openai_chat_sync"));
|
||||
assert!(supports_stream_scheduler_decision_kind(
|
||||
"openai_chat_stream"
|
||||
));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn image_stream_matching_parses_multipart_stream_flag() {
|
||||
let request = Request::builder()
|
||||
.method(Method::POST)
|
||||
.uri("/v1/images/edits")
|
||||
.header(
|
||||
http::header::CONTENT_TYPE,
|
||||
"multipart/form-data; boundary=image-stream-boundary",
|
||||
)
|
||||
.body(())
|
||||
.expect("request should build");
|
||||
let (parts, _) = request.into_parts();
|
||||
let body = concat!(
|
||||
"--image-stream-boundary\r\n",
|
||||
"Content-Disposition: form-data; name=\"stream\"\r\n\r\n",
|
||||
"true\r\n",
|
||||
"--image-stream-boundary--\r\n"
|
||||
);
|
||||
let body_base64 = base64::engine::general_purpose::STANDARD.encode(body.as_bytes());
|
||||
|
||||
assert!(is_matching_stream_request(
|
||||
"openai_image_stream",
|
||||
&parts,
|
||||
&serde_json::json!({}),
|
||||
Some(body_base64.as_str()),
|
||||
));
|
||||
}
|
||||
}
|
||||
|
||||
@@ -84,7 +84,7 @@ pub(crate) fn local_openai_image_spec_metadata(
|
||||
api_format: spec.api_format,
|
||||
decision_kind: spec.decision_kind,
|
||||
report_kind: Some(spec.report_kind),
|
||||
require_streaming: false,
|
||||
require_streaming: spec.require_streaming,
|
||||
requested_model_family: Some(RequestedModelFamily::Standard),
|
||||
}
|
||||
}
|
||||
|
||||
@@ -86,6 +86,7 @@ pub(super) async fn maybe_build_local_gemini_files_decision_payload_for_candidat
|
||||
original_headers: &parts.headers,
|
||||
original_request_body_json: Some(body_json),
|
||||
original_request_body_base64: resolved.provider_request_body_base64.as_deref(),
|
||||
client_requested_stream: spec_metadata.require_streaming,
|
||||
has_envelope: false,
|
||||
needs_conversion: false,
|
||||
extra_fields,
|
||||
|
||||
@@ -5,11 +5,15 @@ mod support;
|
||||
use tracing::warn;
|
||||
|
||||
use crate::ai_pipeline::planner::plan_builders::{
|
||||
build_passthrough_sync_plan_from_decision, LocalSyncPlanAndReport,
|
||||
build_passthrough_sync_plan_from_decision, build_standard_stream_plan_from_decision,
|
||||
LocalStreamPlanAndReport, LocalSyncPlanAndReport,
|
||||
};
|
||||
use crate::ai_pipeline::planner::spec_metadata::local_openai_image_spec_metadata;
|
||||
use crate::ai_pipeline::resolve_local_image_sync_spec as resolve_sync_spec;
|
||||
use crate::ai_pipeline::GatewayControlDecision;
|
||||
use crate::ai_pipeline::{
|
||||
resolve_local_image_stream_spec as resolve_stream_spec,
|
||||
resolve_local_image_sync_spec as resolve_sync_spec,
|
||||
};
|
||||
use crate::{AppState, GatewayControlSyncDecisionResponse, GatewayError};
|
||||
|
||||
use self::decision::maybe_build_local_openai_image_decision_payload_for_candidate;
|
||||
@@ -17,6 +21,7 @@ use self::support::{
|
||||
list_local_openai_image_candidate_attempts, resolve_local_openai_image_decision_input,
|
||||
};
|
||||
|
||||
pub(crate) use self::request::is_openai_image_stream_request;
|
||||
pub(super) use crate::ai_pipeline::LocalOpenAiImageSpec;
|
||||
|
||||
pub(crate) async fn build_local_image_sync_plan_and_reports_for_kind(
|
||||
@@ -44,6 +49,31 @@ pub(crate) async fn build_local_image_sync_plan_and_reports_for_kind(
|
||||
.await
|
||||
}
|
||||
|
||||
pub(crate) async fn build_local_image_stream_plan_and_reports_for_kind(
|
||||
state: &AppState,
|
||||
parts: &http::request::Parts,
|
||||
body_json: &serde_json::Value,
|
||||
body_base64: Option<&str>,
|
||||
trace_id: &str,
|
||||
decision: &GatewayControlDecision,
|
||||
plan_kind: &str,
|
||||
) -> Result<Vec<LocalStreamPlanAndReport>, GatewayError> {
|
||||
let Some(spec) = resolve_stream_spec(plan_kind) else {
|
||||
return Ok(Vec::new());
|
||||
};
|
||||
|
||||
build_local_stream_plan_and_reports(
|
||||
state,
|
||||
parts,
|
||||
body_json,
|
||||
body_base64,
|
||||
trace_id,
|
||||
decision,
|
||||
spec,
|
||||
)
|
||||
.await
|
||||
}
|
||||
|
||||
pub(crate) async fn maybe_build_sync_local_image_decision_payload(
|
||||
state: &AppState,
|
||||
parts: &http::request::Parts,
|
||||
@@ -58,8 +88,75 @@ pub(crate) async fn maybe_build_sync_local_image_decision_payload(
|
||||
};
|
||||
let spec_metadata = local_openai_image_spec_metadata(spec);
|
||||
|
||||
let Some(input) =
|
||||
resolve_local_openai_image_decision_input(state, trace_id, decision, body_json).await
|
||||
let Some(input) = resolve_local_openai_image_decision_input(
|
||||
state,
|
||||
parts,
|
||||
body_json,
|
||||
body_base64,
|
||||
trace_id,
|
||||
decision,
|
||||
)
|
||||
.await
|
||||
else {
|
||||
return Ok(None);
|
||||
};
|
||||
|
||||
let Some(attempts) = list_local_openai_image_candidate_attempts(
|
||||
state,
|
||||
trace_id,
|
||||
&input,
|
||||
body_json,
|
||||
spec_metadata.api_format,
|
||||
spec_metadata.decision_kind,
|
||||
)
|
||||
.await
|
||||
else {
|
||||
return Ok(None);
|
||||
};
|
||||
|
||||
for attempt in attempts {
|
||||
if let Some(payload) = maybe_build_local_openai_image_decision_payload_for_candidate(
|
||||
state,
|
||||
parts,
|
||||
body_json,
|
||||
body_base64,
|
||||
trace_id,
|
||||
&input,
|
||||
attempt,
|
||||
spec,
|
||||
)
|
||||
.await
|
||||
{
|
||||
return Ok(Some(payload));
|
||||
}
|
||||
}
|
||||
|
||||
Ok(None)
|
||||
}
|
||||
|
||||
pub(crate) async fn maybe_build_stream_local_image_decision_payload(
|
||||
state: &AppState,
|
||||
parts: &http::request::Parts,
|
||||
body_json: &serde_json::Value,
|
||||
body_base64: Option<&str>,
|
||||
trace_id: &str,
|
||||
decision: &GatewayControlDecision,
|
||||
plan_kind: &str,
|
||||
) -> Result<Option<GatewayControlSyncDecisionResponse>, GatewayError> {
|
||||
let Some(spec) = resolve_stream_spec(plan_kind) else {
|
||||
return Ok(None);
|
||||
};
|
||||
let spec_metadata = local_openai_image_spec_metadata(spec);
|
||||
|
||||
let Some(input) = resolve_local_openai_image_decision_input(
|
||||
state,
|
||||
parts,
|
||||
body_json,
|
||||
body_base64,
|
||||
trace_id,
|
||||
decision,
|
||||
)
|
||||
.await
|
||||
else {
|
||||
return Ok(None);
|
||||
};
|
||||
@@ -107,8 +204,15 @@ async fn build_local_sync_plan_and_reports(
|
||||
spec: LocalOpenAiImageSpec,
|
||||
) -> Result<Vec<LocalSyncPlanAndReport>, GatewayError> {
|
||||
let spec_metadata = local_openai_image_spec_metadata(spec);
|
||||
let Some(input) =
|
||||
resolve_local_openai_image_decision_input(state, trace_id, decision, body_json).await
|
||||
let Some(input) = resolve_local_openai_image_decision_input(
|
||||
state,
|
||||
parts,
|
||||
body_json,
|
||||
body_base64,
|
||||
trace_id,
|
||||
decision,
|
||||
)
|
||||
.await
|
||||
else {
|
||||
return Ok(Vec::new());
|
||||
};
|
||||
@@ -159,3 +263,73 @@ async fn build_local_sync_plan_and_reports(
|
||||
|
||||
Ok(plans)
|
||||
}
|
||||
|
||||
async fn build_local_stream_plan_and_reports(
|
||||
state: &AppState,
|
||||
parts: &http::request::Parts,
|
||||
body_json: &serde_json::Value,
|
||||
body_base64: Option<&str>,
|
||||
trace_id: &str,
|
||||
decision: &GatewayControlDecision,
|
||||
spec: LocalOpenAiImageSpec,
|
||||
) -> Result<Vec<LocalStreamPlanAndReport>, GatewayError> {
|
||||
let spec_metadata = local_openai_image_spec_metadata(spec);
|
||||
let Some(input) = resolve_local_openai_image_decision_input(
|
||||
state,
|
||||
parts,
|
||||
body_json,
|
||||
body_base64,
|
||||
trace_id,
|
||||
decision,
|
||||
)
|
||||
.await
|
||||
else {
|
||||
return Ok(Vec::new());
|
||||
};
|
||||
|
||||
let Some(attempts) = list_local_openai_image_candidate_attempts(
|
||||
state,
|
||||
trace_id,
|
||||
&input,
|
||||
body_json,
|
||||
spec_metadata.api_format,
|
||||
spec_metadata.decision_kind,
|
||||
)
|
||||
.await
|
||||
else {
|
||||
return Ok(Vec::new());
|
||||
};
|
||||
|
||||
let mut plans = Vec::new();
|
||||
for attempt in attempts {
|
||||
let Some(payload) = maybe_build_local_openai_image_decision_payload_for_candidate(
|
||||
state,
|
||||
parts,
|
||||
body_json,
|
||||
body_base64,
|
||||
trace_id,
|
||||
&input,
|
||||
attempt,
|
||||
spec,
|
||||
)
|
||||
.await
|
||||
else {
|
||||
continue;
|
||||
};
|
||||
|
||||
match build_standard_stream_plan_from_decision(parts, body_json, payload, false) {
|
||||
Ok(Some(value)) => plans.push(value),
|
||||
Ok(None) => {}
|
||||
Err(err) => {
|
||||
warn!(
|
||||
trace_id = %trace_id,
|
||||
decision_kind = spec_metadata.decision_kind,
|
||||
error = ?err,
|
||||
"gateway local openai image stream decision plan build failed"
|
||||
);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
Ok(plans)
|
||||
}
|
||||
|
||||
@@ -78,6 +78,7 @@ pub(super) async fn maybe_build_local_openai_image_decision_payload_for_candidat
|
||||
original_headers: &parts.headers,
|
||||
original_request_body_json: Some(body_json),
|
||||
original_request_body_base64: body_base64,
|
||||
client_requested_stream: spec_metadata.require_streaming,
|
||||
has_envelope: false,
|
||||
needs_conversion: false,
|
||||
extra_fields,
|
||||
@@ -85,7 +86,7 @@ pub(super) async fn maybe_build_local_openai_image_decision_payload_for_candidat
|
||||
|
||||
Some(build_local_execution_decision_response(
|
||||
LocalExecutionDecisionResponseParts {
|
||||
decision_is_stream: false,
|
||||
decision_is_stream: spec_metadata.require_streaming,
|
||||
decision_kind: spec_metadata.decision_kind.to_string(),
|
||||
execution_strategy: ExecutionStrategy::LocalSameFormat,
|
||||
conversion_mode: ConversionMode::None,
|
||||
@@ -112,7 +113,7 @@ pub(super) async fn maybe_build_local_openai_image_decision_payload_for_candidat
|
||||
proxy,
|
||||
tls_profile,
|
||||
timeouts: resolve_transport_execution_timeouts(&transport),
|
||||
upstream_is_stream: false,
|
||||
upstream_is_stream: spec_metadata.require_streaming,
|
||||
report_kind: spec_metadata.report_kind.map(ToOwned::to_owned),
|
||||
report_context: Some(report_context),
|
||||
auth_context: input.auth_context.clone(),
|
||||
|
||||
File diff suppressed because it is too large
Load Diff
@@ -30,28 +30,24 @@ use crate::clock::current_unix_secs;
|
||||
use crate::AppState;
|
||||
use aether_scheduler_core::SchedulerMinimalCandidateSelectionCandidate;
|
||||
|
||||
pub(super) const OPENAI_IMAGE_DEFAULT_MODEL: &str = "gpt-image-2";
|
||||
|
||||
pub(super) use crate::ai_pipeline::planner::candidate_materialization::LocalExecutionCandidateAttempt as LocalOpenAiImageCandidateAttempt;
|
||||
pub(super) use crate::ai_pipeline::planner::decision_input::LocalRequestedModelDecisionInput as LocalOpenAiImageDecisionInput;
|
||||
|
||||
use super::request::resolve_requested_image_model_for_request;
|
||||
|
||||
pub(super) async fn resolve_local_openai_image_decision_input(
|
||||
state: &AppState,
|
||||
parts: &http::request::Parts,
|
||||
body_json: &serde_json::Value,
|
||||
body_base64: Option<&str>,
|
||||
trace_id: &str,
|
||||
decision: &GatewayControlDecision,
|
||||
body_json: &serde_json::Value,
|
||||
) -> Option<LocalOpenAiImageDecisionInput> {
|
||||
let Some(auth_context) = resolve_local_openai_image_auth_context(decision) else {
|
||||
return None;
|
||||
};
|
||||
|
||||
let requested_model = body_json
|
||||
.get("model")
|
||||
.and_then(serde_json::Value::as_str)
|
||||
.map(str::trim)
|
||||
.filter(|value| !value.is_empty())
|
||||
.unwrap_or(OPENAI_IMAGE_DEFAULT_MODEL)
|
||||
.to_string();
|
||||
let requested_model = resolve_requested_image_model_for_request(parts, body_json, body_base64)?;
|
||||
|
||||
let resolved_input = match resolve_local_authenticated_decision_input(
|
||||
state,
|
||||
|
||||
@@ -11,7 +11,9 @@ pub(crate) use self::files::{
|
||||
maybe_build_sync_local_gemini_files_decision_payload,
|
||||
};
|
||||
pub(crate) use self::image::{
|
||||
build_local_image_sync_plan_and_reports_for_kind, maybe_build_sync_local_image_decision_payload,
|
||||
build_local_image_stream_plan_and_reports_for_kind,
|
||||
build_local_image_sync_plan_and_reports_for_kind, is_openai_image_stream_request,
|
||||
maybe_build_stream_local_image_decision_payload, maybe_build_sync_local_image_decision_payload,
|
||||
};
|
||||
pub(crate) use self::video::{
|
||||
build_local_video_sync_plan_and_reports_for_kind, maybe_build_sync_local_video_decision_payload,
|
||||
|
||||
@@ -69,6 +69,7 @@ pub(super) async fn maybe_build_local_video_create_decision_payload_for_candidat
|
||||
original_headers: &parts.headers,
|
||||
original_request_body_json: Some(body_json),
|
||||
original_request_body_base64: None,
|
||||
client_requested_stream: false,
|
||||
has_envelope: false,
|
||||
needs_conversion: false,
|
||||
extra_fields,
|
||||
|
||||
@@ -22,7 +22,7 @@ fn applies_codex_defaults_when_body_rules_do_not_handle_fields() {
|
||||
assert!(body.get("top_p").is_none());
|
||||
assert!(body.get("metadata").is_none());
|
||||
assert_eq!(body["store"], false);
|
||||
assert_eq!(body["instructions"], "You are GPT-5.");
|
||||
assert_eq!(body["instructions"], "You are ChatGPT.");
|
||||
}
|
||||
|
||||
#[test]
|
||||
@@ -125,10 +125,12 @@ fn injects_chatgpt_account_id_and_session_headers_for_codex_requests() {
|
||||
);
|
||||
assert_eq!(
|
||||
headers.get("user-agent"),
|
||||
Some(&"codex-tui/0.122.0 (Aether; x86_64) vscode/3.0.12 (codex-tui; 0.122.0)".to_string())
|
||||
Some(
|
||||
&"codex-tui/0.122.0 (Mac OS 15.2.0; arm64) vscode/2.6.11 (codex-tui; 0.122.0)"
|
||||
.to_string()
|
||||
)
|
||||
);
|
||||
assert_eq!(headers.get("version"), Some(&"0.122.0".to_string()));
|
||||
assert_eq!(headers.get("originator"), Some(&"codex_cli_rs".to_string()));
|
||||
assert_eq!(headers.get("originator"), Some(&"codex-tui".to_string()));
|
||||
assert_eq!(
|
||||
headers.get("session_id"),
|
||||
Some(&"ab5ecce4f0d110fe".to_string())
|
||||
@@ -168,10 +170,6 @@ fn respects_existing_codex_request_and_session_headers() {
|
||||
"user-agent",
|
||||
HeaderValue::from_static("user-specified-agent"),
|
||||
);
|
||||
original_headers.insert(
|
||||
"version",
|
||||
HeaderValue::from_static("user-specified-version"),
|
||||
);
|
||||
original_headers.insert(
|
||||
"originator",
|
||||
HeaderValue::from_static("user-specified-originator"),
|
||||
@@ -192,7 +190,6 @@ fn respects_existing_codex_request_and_session_headers() {
|
||||
Some(&"kept-by-rule-request".to_string())
|
||||
);
|
||||
assert!(!headers.contains_key("user-agent"));
|
||||
assert!(!headers.contains_key("version"));
|
||||
assert!(!headers.contains_key("originator"));
|
||||
assert_eq!(headers.get("session_id"), Some(&"kept-by-rule".to_string()));
|
||||
assert!(!headers.contains_key("conversation_id"));
|
||||
@@ -226,10 +223,12 @@ fn skips_conversation_id_for_compact_codex_requests() {
|
||||
);
|
||||
assert_eq!(
|
||||
headers.get("user-agent"),
|
||||
Some(&"codex-tui/0.122.0 (Aether; x86_64) vscode/3.0.12 (codex-tui; 0.122.0)".to_string())
|
||||
Some(
|
||||
&"codex-tui/0.122.0 (Mac OS 15.2.0; arm64) vscode/2.6.11 (codex-tui; 0.122.0)"
|
||||
.to_string()
|
||||
)
|
||||
);
|
||||
assert_eq!(headers.get("version"), Some(&"0.122.0".to_string()));
|
||||
assert_eq!(headers.get("originator"), Some(&"codex_cli_rs".to_string()));
|
||||
assert_eq!(headers.get("originator"), Some(&"codex-tui".to_string()));
|
||||
assert_eq!(
|
||||
headers.get("session_id"),
|
||||
Some(&"ab5ecce4f0d110fe".to_string())
|
||||
|
||||
@@ -75,6 +75,7 @@ pub(super) async fn maybe_build_local_standard_decision_payload_for_candidate(
|
||||
original_headers: &parts.headers,
|
||||
original_request_body_json: Some(body_json),
|
||||
original_request_body_base64: None,
|
||||
client_requested_stream: spec_metadata.require_streaming,
|
||||
has_envelope: false,
|
||||
needs_conversion: true,
|
||||
extra_fields,
|
||||
|
||||
@@ -265,7 +265,7 @@ mod tests {
|
||||
|
||||
assert!(converted.get("metadata").is_none());
|
||||
assert_eq!(converted["store"], false);
|
||||
assert_eq!(converted["instructions"], "You are GPT-5.");
|
||||
assert_eq!(converted["instructions"], "You are ChatGPT.");
|
||||
}
|
||||
|
||||
#[test]
|
||||
|
||||
@@ -92,6 +92,7 @@ pub(crate) async fn maybe_build_local_openai_chat_decision_payload_for_candidate
|
||||
original_headers: &parts.headers,
|
||||
original_request_body_json: Some(body_json),
|
||||
original_request_body_base64: None,
|
||||
client_requested_stream: upstream_is_stream,
|
||||
has_envelope: false,
|
||||
needs_conversion: matches!(
|
||||
resolved.conversion_mode,
|
||||
|
||||
@@ -96,6 +96,7 @@ pub(crate) async fn maybe_build_local_openai_cli_decision_payload_for_candidate(
|
||||
original_headers: &parts.headers,
|
||||
original_request_body_json: Some(body_json),
|
||||
original_request_body_base64: None,
|
||||
client_requested_stream: spec_metadata.require_streaming,
|
||||
has_envelope: resolved.is_antigravity,
|
||||
needs_conversion: matches!(
|
||||
resolved.conversion_mode,
|
||||
|
||||
@@ -53,13 +53,14 @@ pub(crate) use aether_ai_pipeline::api::{
|
||||
resolve_execution_runtime_stream_plan_kind, resolve_execution_runtime_sync_plan_kind,
|
||||
resolve_finalize_stream_rewrite_mode, resolve_gemini_files_stream_spec,
|
||||
resolve_gemini_files_sync_spec, resolve_gemini_stream_spec, resolve_gemini_sync_spec,
|
||||
resolve_local_image_sync_spec, resolve_local_same_format_stream_spec,
|
||||
resolve_local_same_format_sync_spec, resolve_local_video_sync_spec,
|
||||
resolve_openai_chat_max_tokens, resolve_openai_cli_stream_spec, resolve_openai_cli_sync_spec,
|
||||
stream_body_contains_error_event, supports_stream_scheduler_decision_kind,
|
||||
supports_sync_scheduler_decision_kind, sync_chat_response_conversion_kind,
|
||||
sync_cli_response_conversion_kind, transform_provider_private_stream_line, value_as_u64,
|
||||
CanonicalStreamFrame, ClaudeClientEmitter, ClaudeProviderState, ExecutionRuntimeAuthContext,
|
||||
resolve_local_image_stream_spec, resolve_local_image_sync_spec,
|
||||
resolve_local_same_format_stream_spec, resolve_local_same_format_sync_spec,
|
||||
resolve_local_video_sync_spec, resolve_openai_chat_max_tokens, resolve_openai_cli_stream_spec,
|
||||
resolve_openai_cli_sync_spec, stream_body_contains_error_event,
|
||||
supports_stream_scheduler_decision_kind, supports_sync_scheduler_decision_kind,
|
||||
sync_chat_response_conversion_kind, sync_cli_response_conversion_kind,
|
||||
transform_provider_private_stream_line, value_as_u64, CanonicalStreamFrame,
|
||||
ClaudeClientEmitter, ClaudeProviderState, ExecutionRuntimeAuthContext,
|
||||
FinalizeStreamRewriteMode, GatewayControlPlanRequest, GatewayControlPlanResponse,
|
||||
GatewayControlSyncDecisionResponse, GeminiClientEmitter, GeminiProviderState,
|
||||
LocalCoreSyncErrorKind, LocalGeminiFilesSpec, LocalOpenAiCliSpec, LocalOpenAiImageSpec,
|
||||
@@ -76,12 +77,14 @@ pub(crate) use aether_ai_pipeline::api::{
|
||||
CLAUDE_CHAT_SYNC_SUCCESS_REPORT_KIND, CLAUDE_CLI_STREAM_PLAN_KIND,
|
||||
CLAUDE_CLI_STREAM_SUCCESS_REPORT_KIND, CLAUDE_CLI_SYNC_ERROR_REPORT_KIND,
|
||||
CLAUDE_CLI_SYNC_FINALIZE_REPORT_KIND, CLAUDE_CLI_SYNC_PLAN_KIND,
|
||||
CLAUDE_CLI_SYNC_SUCCESS_REPORT_KIND, EXECUTION_RUNTIME_STREAM_ACTION,
|
||||
EXECUTION_RUNTIME_STREAM_DECISION_ACTION, EXECUTION_RUNTIME_SYNC_ACTION,
|
||||
EXECUTION_RUNTIME_SYNC_DECISION_ACTION, GEMINI_CHAT_STREAM_PLAN_KIND,
|
||||
GEMINI_CHAT_STREAM_SUCCESS_REPORT_KIND, GEMINI_CHAT_SYNC_ERROR_REPORT_KIND,
|
||||
GEMINI_CHAT_SYNC_FINALIZE_REPORT_KIND, GEMINI_CHAT_SYNC_PLAN_KIND,
|
||||
GEMINI_CHAT_SYNC_SUCCESS_REPORT_KIND, GEMINI_CLI_STREAM_PLAN_KIND,
|
||||
CLAUDE_CLI_SYNC_SUCCESS_REPORT_KIND, CODEX_OPENAI_IMAGE_DEFAULT_MODEL,
|
||||
CODEX_OPENAI_IMAGE_DEFAULT_OUTPUT_FORMAT, CODEX_OPENAI_IMAGE_DEFAULT_VARIATION_MODEL,
|
||||
CODEX_OPENAI_IMAGE_DEFAULT_VARIATION_PROMPT, CODEX_OPENAI_IMAGE_INTERNAL_MODEL,
|
||||
EXECUTION_RUNTIME_STREAM_ACTION, EXECUTION_RUNTIME_STREAM_DECISION_ACTION,
|
||||
EXECUTION_RUNTIME_SYNC_ACTION, EXECUTION_RUNTIME_SYNC_DECISION_ACTION,
|
||||
GEMINI_CHAT_STREAM_PLAN_KIND, GEMINI_CHAT_STREAM_SUCCESS_REPORT_KIND,
|
||||
GEMINI_CHAT_SYNC_ERROR_REPORT_KIND, GEMINI_CHAT_SYNC_FINALIZE_REPORT_KIND,
|
||||
GEMINI_CHAT_SYNC_PLAN_KIND, GEMINI_CHAT_SYNC_SUCCESS_REPORT_KIND, GEMINI_CLI_STREAM_PLAN_KIND,
|
||||
GEMINI_CLI_STREAM_SUCCESS_REPORT_KIND, GEMINI_CLI_SYNC_ERROR_REPORT_KIND,
|
||||
GEMINI_CLI_SYNC_FINALIZE_REPORT_KIND, GEMINI_CLI_SYNC_PLAN_KIND,
|
||||
GEMINI_CLI_SYNC_SUCCESS_REPORT_KIND, GEMINI_CLI_V1INTERNAL_ENVELOPE_NAME,
|
||||
@@ -96,7 +99,8 @@ pub(crate) use aether_ai_pipeline::api::{
|
||||
OPENAI_CLI_SYNC_FINALIZE_REPORT_KIND, OPENAI_CLI_SYNC_PLAN_KIND,
|
||||
OPENAI_CLI_SYNC_SUCCESS_REPORT_KIND, OPENAI_COMPACT_STREAM_PLAN_KIND,
|
||||
OPENAI_COMPACT_SYNC_ERROR_REPORT_KIND, OPENAI_COMPACT_SYNC_FINALIZE_REPORT_KIND,
|
||||
OPENAI_COMPACT_SYNC_PLAN_KIND, OPENAI_IMAGE_SYNC_FINALIZE_REPORT_KIND,
|
||||
OPENAI_COMPACT_SYNC_PLAN_KIND, OPENAI_IMAGE_STREAM_PLAN_KIND,
|
||||
OPENAI_IMAGE_STREAM_SUCCESS_REPORT_KIND, OPENAI_IMAGE_SYNC_FINALIZE_REPORT_KIND,
|
||||
OPENAI_IMAGE_SYNC_PLAN_KIND, OPENAI_IMAGE_SYNC_SUCCESS_REPORT_KIND,
|
||||
OPENAI_VIDEO_CANCEL_SYNC_PLAN_KIND, OPENAI_VIDEO_CONTENT_PLAN_KIND,
|
||||
OPENAI_VIDEO_CREATE_SYNC_FINALIZE_REPORT_KIND, OPENAI_VIDEO_CREATE_SYNC_PLAN_KIND,
|
||||
|
||||
Reference in New Issue
Block a user