mirror of
https://github.com/fawney19/Aether.git
synced 2026-10-04 00:17:45 +08:00
feat(xai): add native image and video endpoints
Expose the xAI Imagine image and video surfaces on top of the `xai` provider, and make the shared OpenAI video-task layer survive the production configuration they need. Native video requests live under /v1 (generations, edits, extensions, with /v1/videos as a creation alias that only selects xAI candidates); the OpenAI-compatible adapter stays under /openai/v1/videos and maps `seconds` / `size` onto numeric duration, aspect ratio and resolution. Clients receive an opaque Aether task ID scoped to the owning user; polling uses the upstream task ID and the original credential, and completed downloads fetch the returned media URL without forwarding provider authorization to the media host. Three fixes to the shared video layer are required for this to work outside tests: - OpenAI/xAI task persistence now supplies a stable 16-character short_id, which the PostgreSQL schema requires. Existing rows keep their original value across reconstruction, so no schema change or historical rewrite is needed. - Task retrieval and content downloads are admitted by the production GET execution gate, and reconstructed tasks resolve proxy nodes, system proxy defaults, tunnel affinity and transport profiles through the same deployment resolver used for creation. A configured proxy route no longer silently becomes a direct request after restart. - When the gateway also serves the frontend, /openai/v1/videos and its subpaths bypass the static SPA handler. Otherwise a video query returns HTTP 200 with text/html instead of the task JSON. Co-Authored-By: Claude Opus 5 <[email protected]>
This commit is contained in:
@@ -12,5 +12,6 @@ aether-data-contracts.workspace = true
|
||||
async-trait.workspace = true
|
||||
serde.workspace = true
|
||||
serde_json.workspace = true
|
||||
sha2.workspace = true
|
||||
url.workspace = true
|
||||
uuid.workspace = true
|
||||
|
||||
@@ -43,6 +43,27 @@ pub fn map_openai_stored_task_to_read_response(
|
||||
}
|
||||
|
||||
fn build_openai_stored_task_body(task: StoredVideoTask, status: VideoTaskStatus) -> Value {
|
||||
if task.client_api_format.as_deref() == Some("xai:video") {
|
||||
let mut body = json!({"status":match status {
|
||||
VideoTaskStatus::Completed => "done",
|
||||
VideoTaskStatus::Expired => "expired",
|
||||
VideoTaskStatus::Failed | VideoTaskStatus::Cancelled | VideoTaskStatus::Deleted => "failed",
|
||||
_ => "pending",
|
||||
}});
|
||||
if let Some(model) = task.model {
|
||||
body["model"] = json!(model);
|
||||
}
|
||||
if let Some(url) = task.video_url {
|
||||
body["video"] = json!({"url":url});
|
||||
if let Some(duration) = task.duration_seconds {
|
||||
body["video"]["duration"] = json!(duration);
|
||||
}
|
||||
}
|
||||
if status == VideoTaskStatus::Failed {
|
||||
body["error"] = json!({"code":sanitize_video_task_error_code(task.error_code).unwrap_or_else(|| "unknown".into()),"message":"Video generation failed"});
|
||||
}
|
||||
return body;
|
||||
}
|
||||
let mut body = json!({
|
||||
"id": task.id,
|
||||
"object": "video",
|
||||
@@ -57,6 +78,9 @@ fn build_openai_stored_task_body(task: StoredVideoTask, status: VideoTaskStatus)
|
||||
if let Some(prompt) = task.prompt {
|
||||
body["prompt"] = Value::String(prompt);
|
||||
}
|
||||
if let Some(seconds) = task.duration_seconds {
|
||||
body["seconds"] = json!(seconds.to_string());
|
||||
}
|
||||
if let Some(size) = task.size {
|
||||
body["size"] = Value::String(size);
|
||||
}
|
||||
@@ -91,21 +115,97 @@ fn map_openai_stored_task_status(status: VideoTaskStatus) -> &'static str {
|
||||
}
|
||||
|
||||
impl OpenAiVideoTaskSeed {
|
||||
pub fn uses_xai_provider(&self) -> bool {
|
||||
self.xai_provider || self.is_xai_native()
|
||||
}
|
||||
|
||||
pub fn is_xai_native(&self) -> bool {
|
||||
self.persistence.client_api_format == "xai:video"
|
||||
}
|
||||
|
||||
pub fn native_create_body_json(&self) -> Value {
|
||||
let mut body = self.native_response.clone().unwrap_or_else(|| json!({}));
|
||||
body["request_id"] = json!(self.local_task_id);
|
||||
if body.get("id").is_some() {
|
||||
body["id"] = json!(self.local_task_id);
|
||||
}
|
||||
body
|
||||
}
|
||||
|
||||
fn native_read_body_json(&self) -> Value {
|
||||
if let Some(mut body) = self.native_response.clone().filter(|body| {
|
||||
body.get("status").is_some()
|
||||
|| body.get("error").is_some()
|
||||
|| body.get("code").is_some()
|
||||
}) {
|
||||
if body.get("request_id").is_some() {
|
||||
body["request_id"] = json!(self.local_task_id);
|
||||
}
|
||||
if body.get("id").is_some() {
|
||||
body["id"] = json!(self.local_task_id);
|
||||
}
|
||||
return body;
|
||||
}
|
||||
let mut body = json!({"status":match self.status {
|
||||
LocalVideoTaskStatus::Completed => "done",
|
||||
LocalVideoTaskStatus::Expired => "expired",
|
||||
LocalVideoTaskStatus::Failed | LocalVideoTaskStatus::Cancelled | LocalVideoTaskStatus::Deleted => "failed",
|
||||
_ => "pending",
|
||||
}});
|
||||
if let Some(model) = &self.model {
|
||||
body["model"] = json!(model);
|
||||
}
|
||||
if let Some(url) = &self.video_url {
|
||||
body["video"] = json!({"url":url});
|
||||
if let Some(duration) = self.seconds.as_deref().and_then(|v| v.parse::<u64>().ok()) {
|
||||
body["video"]["duration"] = json!(duration);
|
||||
}
|
||||
}
|
||||
if self.error_code.is_some() {
|
||||
body["error"] = json!({"code":self.error_code,"message":"Video generation failed"});
|
||||
}
|
||||
body
|
||||
}
|
||||
|
||||
pub fn apply_provider_body(&mut self, provider_body: &Map<String, Value>) {
|
||||
if self.uses_xai_provider() {
|
||||
self.native_response = Some(Value::Object(provider_body.clone()));
|
||||
}
|
||||
|
||||
let raw_status = provider_body
|
||||
.get("status")
|
||||
.and_then(Value::as_str)
|
||||
.map(str::trim)
|
||||
.unwrap_or_default();
|
||||
self.status = match raw_status {
|
||||
"queued" => LocalVideoTaskStatus::Queued,
|
||||
"processing" => LocalVideoTaskStatus::Processing,
|
||||
"completed" => LocalVideoTaskStatus::Completed,
|
||||
"failed" => LocalVideoTaskStatus::Failed,
|
||||
"cancelled" => LocalVideoTaskStatus::Cancelled,
|
||||
// Accept xAI's native lifecycle vocabulary alongside OpenAI's fields.
|
||||
self.status = match raw_status.to_ascii_lowercase().as_str() {
|
||||
"queued" | "pending" => LocalVideoTaskStatus::Queued,
|
||||
"processing" | "in_progress" | "running" => LocalVideoTaskStatus::Processing,
|
||||
"completed" | "done" | "succeeded" | "success" => LocalVideoTaskStatus::Completed,
|
||||
"failed" | "error" => LocalVideoTaskStatus::Failed,
|
||||
"cancelled" | "canceled" => LocalVideoTaskStatus::Cancelled,
|
||||
"expired" => LocalVideoTaskStatus::Expired,
|
||||
_ => LocalVideoTaskStatus::Submitted,
|
||||
};
|
||||
let error = provider_body.get("error").filter(|value| !value.is_null());
|
||||
let error_code = provider_body
|
||||
.get("code")
|
||||
.and_then(Value::as_str)
|
||||
.filter(|value| !value.trim().is_empty())
|
||||
.or_else(|| {
|
||||
error
|
||||
.and_then(|value| value.get("code"))
|
||||
.and_then(Value::as_str)
|
||||
});
|
||||
// xAI may report a failed job as a 200 response with code/error only.
|
||||
if (error.is_some() || error_code.is_some())
|
||||
&& !matches!(
|
||||
self.status,
|
||||
LocalVideoTaskStatus::Cancelled | LocalVideoTaskStatus::Expired
|
||||
)
|
||||
{
|
||||
self.status = LocalVideoTaskStatus::Failed;
|
||||
}
|
||||
self.progress_percent = provider_body
|
||||
.get("progress")
|
||||
.and_then(Value::as_u64)
|
||||
@@ -117,20 +217,35 @@ impl OpenAiVideoTaskSeed {
|
||||
});
|
||||
self.completed_at_unix_secs = provider_body.get("completed_at").and_then(Value::as_u64);
|
||||
self.expires_at_unix_secs = provider_body.get("expires_at").and_then(Value::as_u64);
|
||||
let error = provider_body.get("error").and_then(Value::as_object);
|
||||
self.error_code = sanitize_video_task_error_code(
|
||||
error
|
||||
.and_then(|value| value.get("code"))
|
||||
.and_then(Value::as_str)
|
||||
.map(str::to_string),
|
||||
);
|
||||
self.error_code = sanitize_video_task_error_code(error_code.map(str::to_string));
|
||||
self.error_message = None;
|
||||
self.video_url = provider_body
|
||||
.get("video_url")
|
||||
.or_else(|| provider_body.get("url"))
|
||||
.or_else(|| provider_body.get("result_url"))
|
||||
.or_else(|| {
|
||||
provider_body
|
||||
.get("video")
|
||||
.and_then(|video| video.get("url"))
|
||||
})
|
||||
.and_then(Value::as_str)
|
||||
.map(str::to_string);
|
||||
if let Some(seconds) = provider_body
|
||||
.get("seconds")
|
||||
.or_else(|| {
|
||||
provider_body
|
||||
.get("video")
|
||||
.and_then(|video| video.get("duration"))
|
||||
})
|
||||
.filter(|value| value.is_string() || value.is_number())
|
||||
{
|
||||
self.seconds = Some(
|
||||
seconds
|
||||
.as_str()
|
||||
.map(str::to_string)
|
||||
.unwrap_or_else(|| seconds.to_string()),
|
||||
);
|
||||
}
|
||||
}
|
||||
|
||||
pub fn build_content_stream_action(
|
||||
@@ -239,6 +354,9 @@ impl OpenAiVideoTaskSeed {
|
||||
}
|
||||
|
||||
pub fn client_body_json(&self) -> Value {
|
||||
if self.is_xai_native() {
|
||||
return self.native_read_body_json();
|
||||
}
|
||||
let mut body = json!({
|
||||
"id": self.local_task_id,
|
||||
"object": "video",
|
||||
@@ -259,6 +377,9 @@ impl OpenAiVideoTaskSeed {
|
||||
if let Some(seconds) = &self.seconds {
|
||||
body["seconds"] = Value::String(seconds.clone());
|
||||
}
|
||||
if let Some(video_url) = &self.video_url {
|
||||
body["video_url"] = Value::String(video_url.clone());
|
||||
}
|
||||
if let Some(remixed_from_video_id) = &self.remixed_from_video_id {
|
||||
body["remixed_from_video_id"] = Value::String(remixed_from_video_id.clone());
|
||||
}
|
||||
@@ -357,12 +478,20 @@ impl OpenAiVideoTaskSeed {
|
||||
}
|
||||
|
||||
pub fn build_get_follow_up_plan(&self, trace_id: &str) -> Option<ExecutionPlan> {
|
||||
if !matches!(
|
||||
let refreshable = matches!(
|
||||
self.status,
|
||||
LocalVideoTaskStatus::Submitted
|
||||
| LocalVideoTaskStatus::Queued
|
||||
| LocalVideoTaskStatus::Processing
|
||||
) {
|
||||
) || (self.uses_xai_provider()
|
||||
&& self.native_response.is_none()
|
||||
&& matches!(
|
||||
self.status,
|
||||
LocalVideoTaskStatus::Completed
|
||||
| LocalVideoTaskStatus::Failed
|
||||
| LocalVideoTaskStatus::Expired
|
||||
));
|
||||
if !refreshable {
|
||||
return None;
|
||||
}
|
||||
|
||||
@@ -573,7 +702,12 @@ impl OpenAiVideoTaskSeed {
|
||||
};
|
||||
let mut record = UpsertVideoTask {
|
||||
id: self.local_task_id.clone(),
|
||||
short_id: None,
|
||||
// The production schema requires a unique, non-null short_id (at most 16 chars).
|
||||
// Derive it deterministically so repeated capture and legacy snapshot reloads agree.
|
||||
short_id: Some(self.local_short_id.clone().unwrap_or_else(|| {
|
||||
use sha2::{Digest, Sha256};
|
||||
format!("{:x}", Sha256::digest(self.local_task_id.as_bytes()))[..16].to_string()
|
||||
})),
|
||||
request_id: self.persistence.request_id.clone(),
|
||||
user_id: self.user_id.clone(),
|
||||
api_key_id: self.api_key_id.clone(),
|
||||
@@ -589,7 +723,11 @@ impl OpenAiVideoTaskSeed {
|
||||
model: self.model.clone().or_else(|| Some(String::new())),
|
||||
prompt: self.prompt.clone().or_else(|| Some(String::new())),
|
||||
original_request_body: None,
|
||||
duration_seconds: request_body_u32(&self.persistence.original_request_body, "seconds"),
|
||||
duration_seconds: self
|
||||
.seconds
|
||||
.as_deref()
|
||||
.and_then(|value| value.parse().ok())
|
||||
.or_else(|| request_body_u32(&self.persistence.original_request_body, "seconds")),
|
||||
resolution: request_body_string(&self.persistence.original_request_body, "resolution"),
|
||||
aspect_ratio: request_body_string(
|
||||
&self.persistence.original_request_body,
|
||||
@@ -697,6 +835,9 @@ mod tests {
|
||||
#[test]
|
||||
fn builds_minimal_openai_persistence_record_without_sensitive_snapshot() {
|
||||
let seed = OpenAiVideoTaskSeed {
|
||||
local_short_id: None,
|
||||
native_response: None,
|
||||
xai_provider: false,
|
||||
local_task_id: "task-openai-sensitive".to_string(),
|
||||
upstream_task_id: "upstream-openai-sensitive".to_string(),
|
||||
created_at_unix_ms: 1_712_345_678,
|
||||
@@ -747,6 +888,12 @@ mod tests {
|
||||
|
||||
let record = seed.to_upsert_record();
|
||||
|
||||
let short_id = record
|
||||
.short_id
|
||||
.as_deref()
|
||||
.expect("database short_id is required");
|
||||
assert_eq!(short_id.len(), 16);
|
||||
assert_eq!(seed.to_upsert_record().short_id, record.short_id);
|
||||
assert_eq!(record.error_code.as_deref(), Some("provider_error"));
|
||||
assert!(record.original_request_body.is_none());
|
||||
assert!(record.progress_message.is_none());
|
||||
@@ -759,6 +906,8 @@ mod tests {
|
||||
|
||||
let mut stored = record.into_stored();
|
||||
stored.status = VideoTaskStatus::Completed;
|
||||
// Migrated tasks can already have a short ID unrelated to the derived ID.
|
||||
stored.short_id = Some("legacy-short-id".to_string());
|
||||
let snapshot =
|
||||
LocalVideoTaskSnapshot::from_stored_task_with_transport(&stored, seed.transport)
|
||||
.expect("stored task should reconstruct with current transport");
|
||||
@@ -766,6 +915,21 @@ mod tests {
|
||||
panic!("expected OpenAI snapshot");
|
||||
};
|
||||
assert_eq!(restored.prompt, stored.prompt);
|
||||
assert_eq!(restored.to_upsert_record().short_id, stored.short_id);
|
||||
let mut embedded = stored.clone();
|
||||
let mut legacy_snapshot =
|
||||
serde_json::to_value(LocalVideoTaskSnapshot::OpenAi(restored.clone())).unwrap();
|
||||
legacy_snapshot["OpenAi"]
|
||||
.as_object_mut()
|
||||
.unwrap()
|
||||
.remove("local_short_id");
|
||||
embedded.request_metadata = Some(json!({"rust_local_snapshot": legacy_snapshot}));
|
||||
let embedded_snapshot = LocalVideoTaskSnapshot::from_stored_task(&embedded)
|
||||
.expect("legacy embedded snapshot should hydrate");
|
||||
assert_eq!(
|
||||
embedded_snapshot.to_upsert_record().short_id,
|
||||
stored.short_id
|
||||
);
|
||||
assert_eq!(restored.to_upsert_record().video_url, stored.video_url);
|
||||
let Some(LocalVideoTaskContentAction::StreamPlan(plan)) =
|
||||
restored.build_content_stream_action(None, "trace-download")
|
||||
|
||||
@@ -8,7 +8,9 @@ use uuid::Uuid;
|
||||
use crate::{LocalVideoTaskRegistryMutation, LocalVideoTaskStatus, VideoTaskTruthSourceMode};
|
||||
|
||||
pub fn extract_openai_task_id_from_path(path: &str) -> Option<&str> {
|
||||
let suffix = path.strip_prefix("/v1/videos/")?;
|
||||
let suffix = path
|
||||
.strip_prefix("/v1/videos/")
|
||||
.or_else(|| path.strip_prefix("/openai/v1/videos/"))?;
|
||||
if suffix.is_empty()
|
||||
|| suffix.contains('/')
|
||||
|| suffix.ends_with(":cancel")
|
||||
@@ -29,21 +31,27 @@ pub fn extract_gemini_short_id_from_path(path: &str) -> Option<&str> {
|
||||
}
|
||||
|
||||
pub fn extract_openai_task_id_from_cancel_path(path: &str) -> Option<&str> {
|
||||
let suffix = path.strip_prefix("/v1/videos/")?;
|
||||
let suffix = path
|
||||
.strip_prefix("/v1/videos/")
|
||||
.or_else(|| path.strip_prefix("/openai/v1/videos/"))?;
|
||||
suffix
|
||||
.strip_suffix("/cancel")
|
||||
.filter(|value| !value.is_empty())
|
||||
}
|
||||
|
||||
pub fn extract_openai_task_id_from_remix_path(path: &str) -> Option<&str> {
|
||||
let suffix = path.strip_prefix("/v1/videos/")?;
|
||||
let suffix = path
|
||||
.strip_prefix("/v1/videos/")
|
||||
.or_else(|| path.strip_prefix("/openai/v1/videos/"))?;
|
||||
suffix
|
||||
.strip_suffix("/remix")
|
||||
.filter(|value| !value.is_empty())
|
||||
}
|
||||
|
||||
pub fn extract_openai_task_id_from_content_path(path: &str) -> Option<&str> {
|
||||
let suffix = path.strip_prefix("/v1/videos/")?;
|
||||
let suffix = path
|
||||
.strip_prefix("/v1/videos/")
|
||||
.or_else(|| path.strip_prefix("/openai/v1/videos/"))?;
|
||||
suffix
|
||||
.strip_suffix("/content")
|
||||
.filter(|value| !value.is_empty())
|
||||
|
||||
@@ -73,7 +73,7 @@ async fn read_openai_video_task_response(
|
||||
}
|
||||
None => state.find_stored_video_task(lookup).await?,
|
||||
};
|
||||
let Some(task) = task else {
|
||||
let Some(mut task) = task else {
|
||||
return Ok(None);
|
||||
};
|
||||
|
||||
@@ -81,6 +81,9 @@ async fn read_openai_video_task_response(
|
||||
return Ok(None);
|
||||
}
|
||||
|
||||
if request_path.starts_with("/openai/v1/videos/") {
|
||||
task.client_api_format = Some("openai:video".into());
|
||||
}
|
||||
Ok(Some(map_openai_stored_task_to_read_response(task)))
|
||||
}
|
||||
|
||||
|
||||
@@ -105,13 +105,8 @@ impl VideoTaskService {
|
||||
if self.truth_source_mode != VideoTaskTruthSourceMode::RustAuthoritative {
|
||||
return None;
|
||||
}
|
||||
match route_family {
|
||||
Some("openai") => extract_openai_task_id_from_path(request_path)
|
||||
.and_then(|task_id| self.store.read_openai(task_id)),
|
||||
Some("gemini") => extract_gemini_short_id_from_path(request_path)
|
||||
.and_then(|short_id| self.store.read_gemini(short_id)),
|
||||
_ => None,
|
||||
}
|
||||
self.snapshot_for_route(route_family, request_path)
|
||||
.map(|snapshot| snapshot.read_response_for_path(request_path))
|
||||
}
|
||||
|
||||
pub fn read_response_for_user(
|
||||
@@ -126,7 +121,7 @@ impl VideoTaskService {
|
||||
let snapshot = self.snapshot_for_route(route_family, request_path)?;
|
||||
snapshot
|
||||
.belongs_to_user(user_id)
|
||||
.then(|| snapshot.read_response())
|
||||
.then(|| snapshot.read_response_for_path(request_path))
|
||||
}
|
||||
|
||||
pub fn snapshot_for_route(
|
||||
|
||||
@@ -29,6 +29,7 @@ impl LocalVideoTaskSnapshot {
|
||||
// contain stale identity fields after a task import or repair.
|
||||
match &mut snapshot {
|
||||
Self::OpenAi(seed) => {
|
||||
seed.local_short_id = task.short_id.clone();
|
||||
seed.user_id = task.user_id.clone();
|
||||
seed.api_key_id = task.api_key_id.clone();
|
||||
}
|
||||
@@ -51,6 +52,9 @@ impl LocalVideoTaskSnapshot {
|
||||
"openai:video" => {
|
||||
let upstream_task_id = non_empty_owned(task.external_task_id.as_ref())?;
|
||||
Some(Self::OpenAi(OpenAiVideoTaskSeed {
|
||||
local_short_id: task.short_id.clone(),
|
||||
native_response: None,
|
||||
xai_provider: persistence.client_api_format == "xai:video",
|
||||
local_task_id: task.id.clone(),
|
||||
upstream_task_id,
|
||||
created_at_unix_ms: task.created_at_unix_ms,
|
||||
@@ -142,6 +146,19 @@ impl LocalVideoTaskSnapshot {
|
||||
}
|
||||
}
|
||||
|
||||
pub fn read_response_for_path(&self, path: &str) -> LocalVideoTaskReadResponse {
|
||||
if let Self::OpenAi(seed) = self {
|
||||
let mut seed = seed.clone();
|
||||
if path.starts_with("/openai/v1/videos/") {
|
||||
seed.persistence.client_api_format = "openai:video".to_string();
|
||||
} else if path.starts_with("/v1/videos/") && seed.uses_xai_provider() {
|
||||
seed.persistence.client_api_format = "xai:video".to_string();
|
||||
}
|
||||
return Self::OpenAi(seed).read_response();
|
||||
}
|
||||
self.read_response()
|
||||
}
|
||||
|
||||
pub fn read_response(&self) -> LocalVideoTaskReadResponse {
|
||||
match self {
|
||||
Self::OpenAi(seed) => match seed.status {
|
||||
|
||||
@@ -19,14 +19,17 @@ impl LocalVideoTaskSeed {
|
||||
) -> Option<Self> {
|
||||
let transport = LocalVideoTaskTransport::from_plan(plan)?;
|
||||
let persistence = LocalVideoTaskPersistence::from_report_context(report_context, plan);
|
||||
match report_kind {
|
||||
let mut seed = match report_kind {
|
||||
"openai_video_create_sync_finalize" => {
|
||||
let upstream_id = provider_body.get("id").and_then(Value::as_str)?.trim();
|
||||
if upstream_id.is_empty() {
|
||||
return None;
|
||||
}
|
||||
let upstream_id = openai_video_provider_task_id(provider_body)?;
|
||||
|
||||
Some(Self::OpenAiCreate(OpenAiVideoTaskSeed {
|
||||
local_short_id: None,
|
||||
native_response: None,
|
||||
xai_provider: report_context
|
||||
.get("video_provider_xai")
|
||||
.and_then(Value::as_bool)
|
||||
.unwrap_or(false),
|
||||
local_task_id: context_text(report_context, "local_task_id")
|
||||
.unwrap_or_else(|| Uuid::new_v4().to_string()),
|
||||
upstream_task_id: upstream_id.to_string(),
|
||||
@@ -37,8 +40,12 @@ impl LocalVideoTaskSeed {
|
||||
model: context_text(report_context, "model")
|
||||
.or_else(|| request_body_text(report_context, "model")),
|
||||
prompt: request_body_text(report_context, "prompt"),
|
||||
size: request_body_text(report_context, "size"),
|
||||
seconds: request_body_text(report_context, "seconds"),
|
||||
size: context_text(report_context, "video_size")
|
||||
.or_else(|| request_body_text(report_context, "size")),
|
||||
seconds: context_u64(report_context, "video_duration")
|
||||
.map(|v| v.to_string())
|
||||
.or_else(|| request_body_text(report_context, "seconds"))
|
||||
.or_else(|| request_body_text(report_context, "duration")),
|
||||
remixed_from_video_id: None,
|
||||
status: LocalVideoTaskStatus::Submitted,
|
||||
progress_percent: 0,
|
||||
@@ -52,12 +59,15 @@ impl LocalVideoTaskSeed {
|
||||
}))
|
||||
}
|
||||
"openai_video_remix_sync_finalize" => {
|
||||
let upstream_id = provider_body.get("id").and_then(Value::as_str)?.trim();
|
||||
if upstream_id.is_empty() {
|
||||
return None;
|
||||
}
|
||||
let upstream_id = openai_video_provider_task_id(provider_body)?;
|
||||
|
||||
Some(Self::OpenAiRemix(OpenAiVideoTaskSeed {
|
||||
local_short_id: None,
|
||||
native_response: None,
|
||||
xai_provider: report_context
|
||||
.get("video_provider_xai")
|
||||
.and_then(Value::as_bool)
|
||||
.unwrap_or(false),
|
||||
local_task_id: context_text(report_context, "local_task_id")
|
||||
.unwrap_or_else(|| Uuid::new_v4().to_string()),
|
||||
upstream_task_id: upstream_id.to_string(),
|
||||
@@ -68,8 +78,12 @@ impl LocalVideoTaskSeed {
|
||||
model: context_text(report_context, "model")
|
||||
.or_else(|| request_body_text(report_context, "model")),
|
||||
prompt: request_body_text(report_context, "prompt"),
|
||||
size: request_body_text(report_context, "size"),
|
||||
seconds: request_body_text(report_context, "seconds"),
|
||||
size: context_text(report_context, "video_size")
|
||||
.or_else(|| request_body_text(report_context, "size")),
|
||||
seconds: context_u64(report_context, "video_duration")
|
||||
.map(|v| v.to_string())
|
||||
.or_else(|| request_body_text(report_context, "seconds"))
|
||||
.or_else(|| request_body_text(report_context, "duration")),
|
||||
remixed_from_video_id: context_text(report_context, "task_id")
|
||||
.or_else(|| request_body_text(report_context, "remix_video_id")),
|
||||
status: LocalVideoTaskStatus::Submitted,
|
||||
@@ -110,7 +124,11 @@ impl LocalVideoTaskSeed {
|
||||
}))
|
||||
}
|
||||
_ => None,
|
||||
}?;
|
||||
if let Self::OpenAiCreate(task) | Self::OpenAiRemix(task) = &mut seed {
|
||||
task.apply_provider_body(provider_body);
|
||||
}
|
||||
Some(seed)
|
||||
}
|
||||
|
||||
pub fn success_report_kind(&self) -> &'static str {
|
||||
@@ -144,12 +162,28 @@ impl LocalVideoTaskSeed {
|
||||
|
||||
pub fn client_body_json(&self) -> Value {
|
||||
match self {
|
||||
Self::OpenAiCreate(seed) | Self::OpenAiRemix(seed) => seed.client_body_json(),
|
||||
Self::OpenAiCreate(seed) | Self::OpenAiRemix(seed) => {
|
||||
if seed.is_xai_native() {
|
||||
seed.native_create_body_json()
|
||||
} else {
|
||||
seed.client_body_json()
|
||||
}
|
||||
}
|
||||
Self::GeminiCreate(seed) => seed.client_body_json(),
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
fn openai_video_provider_task_id(body: &Map<String, Value>) -> Option<&str> {
|
||||
// xAI's OpenAI-compatible video creation returns request_id instead of id.
|
||||
["id", "request_id"].into_iter().find_map(|field| {
|
||||
body.get(field)
|
||||
.and_then(Value::as_str)
|
||||
.map(str::trim)
|
||||
.filter(|value| !value.is_empty())
|
||||
})
|
||||
}
|
||||
|
||||
impl VideoTaskTruthSourceMode {
|
||||
pub fn prepare_sync_success(
|
||||
self,
|
||||
@@ -353,6 +387,234 @@ mod tests {
|
||||
resolve_local_sync_success_background_report_kind,
|
||||
};
|
||||
|
||||
#[test]
|
||||
fn xai_native_video_protocol_survives_persistence_and_preserves_provider_fields() {
|
||||
use crate::{
|
||||
LocalVideoTaskContentAction, LocalVideoTaskSnapshot, VideoTaskService,
|
||||
VideoTaskTruthSourceMode,
|
||||
};
|
||||
let mut plan =
|
||||
build_internal_finalize_video_plan("native-create", "openai:video", None).unwrap();
|
||||
plan.url = "https://api.x.ai/v1/videos/generations".into();
|
||||
plan.headers
|
||||
.insert("authorization".into(), "Bearer test-key".into());
|
||||
let service = VideoTaskService::new(VideoTaskTruthSourceMode::RustAuthoritative);
|
||||
let context = json!({"local_task_id":"native-local", "user_id":"owner", "model":"grok-imagine-video", "video_client_protocol":"xai", "video_duration":6});
|
||||
let success = service
|
||||
.prepare_sync_success(
|
||||
"openai_video_create_sync_finalize",
|
||||
json!({"request_id":"native-upstream", "future_field":true})
|
||||
.as_object()
|
||||
.unwrap(),
|
||||
context.as_object().unwrap(),
|
||||
&plan,
|
||||
)
|
||||
.unwrap();
|
||||
assert_eq!(
|
||||
success.client_body_json(),
|
||||
json!({"request_id":"native-local","future_field":true})
|
||||
);
|
||||
let mut snapshot = success.to_snapshot();
|
||||
let body = json!({"status":"done","video":{"url":"https://vidgen.x.ai/video.mp4","duration":6,"respect_moderation":true},"future_field":[1,2]});
|
||||
snapshot.apply_provider_body(body.as_object().unwrap());
|
||||
assert_eq!(snapshot.read_response().body_json, body);
|
||||
assert_eq!(
|
||||
snapshot
|
||||
.read_response_for_path("/openai/v1/videos/native-local")
|
||||
.body_json["status"],
|
||||
"completed"
|
||||
);
|
||||
let LocalVideoTaskSnapshot::OpenAi(seed) = &snapshot else {
|
||||
panic!("openai task expected")
|
||||
};
|
||||
let Some(LocalVideoTaskContentAction::StreamPlan(download)) =
|
||||
seed.build_content_stream_action(None, "download")
|
||||
else {
|
||||
panic!("download expected")
|
||||
};
|
||||
assert_eq!(download.url, "https://vidgen.x.ai/video.mp4");
|
||||
assert!(download.headers.is_empty());
|
||||
let stored = snapshot.to_upsert_record().into_stored();
|
||||
assert!(stored.request_metadata.is_none());
|
||||
assert_eq!(stored.client_api_format.as_deref(), Some("xai:video"));
|
||||
let restored = LocalVideoTaskSnapshot::from_stored_task_with_transport(
|
||||
&stored,
|
||||
seed.transport.clone(),
|
||||
)
|
||||
.unwrap();
|
||||
assert_eq!(restored.read_response().body_json["status"], "done");
|
||||
service.record_snapshot(restored);
|
||||
let poll = service
|
||||
.prepare_read_refresh_sync_plan_for_user(
|
||||
Some("openai"),
|
||||
"/v1/videos/native-local",
|
||||
"owner",
|
||||
"poll",
|
||||
)
|
||||
.unwrap();
|
||||
assert_eq!(poll.plan.url, "https://api.x.ai/v1/videos/native-upstream");
|
||||
assert!(service
|
||||
.prepare_read_refresh_sync_plan_for_user(
|
||||
Some("openai"),
|
||||
"/v1/videos/native-local",
|
||||
"foreign",
|
||||
"poll"
|
||||
)
|
||||
.is_none());
|
||||
assert!(service.apply_read_refresh_projection(&poll, body.as_object().unwrap()));
|
||||
assert_eq!(
|
||||
service
|
||||
.read_response_for_user(Some("openai"), "/v1/videos/native-local", "owner")
|
||||
.unwrap()
|
||||
.body_json,
|
||||
body
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn xai_video_lifecycle_creates_polls_persists_and_downloads() {
|
||||
use crate::{
|
||||
LocalVideoTaskContentAction, LocalVideoTaskSnapshot, VideoTaskService,
|
||||
VideoTaskTruthSourceMode,
|
||||
};
|
||||
for api_root in ["https://cli-chat-proxy.grok.com/v1", "https://api.x.ai/v1"] {
|
||||
let mut plan =
|
||||
build_internal_finalize_video_plan("xai-create", "openai:video", None).unwrap();
|
||||
plan.url = format!("{api_root}/videos/generations");
|
||||
plan.headers
|
||||
.insert("authorization".into(), "Bearer test-token".into());
|
||||
let service = VideoTaskService::new(VideoTaskTruthSourceMode::RustAuthoritative);
|
||||
let context = json!({"local_task_id": "local-video", "model": "grok-imagine-video", "original_request_body": {"prompt": "A cat", "seconds": "6"}});
|
||||
let success = service
|
||||
.prepare_sync_success(
|
||||
"openai_video_create_sync_finalize",
|
||||
json!({"request_id": "xai-request"}).as_object().unwrap(),
|
||||
context.as_object().unwrap(),
|
||||
&plan,
|
||||
)
|
||||
.unwrap();
|
||||
assert_eq!(success.client_body_json()["id"], "local-video");
|
||||
assert_eq!(success.client_body_json()["status"], "queued");
|
||||
let snapshot = success.to_snapshot();
|
||||
assert_eq!(
|
||||
snapshot.to_upsert_record().external_task_id.as_deref(),
|
||||
Some("xai-request")
|
||||
);
|
||||
service.record_snapshot(snapshot.clone());
|
||||
let poll = service
|
||||
.prepare_poll_refresh_plan_for_snapshot(snapshot, "xai-poll")
|
||||
.unwrap();
|
||||
assert_eq!(poll.plan.method, "GET");
|
||||
assert_eq!(poll.plan.url, format!("{api_root}/videos/xai-request"));
|
||||
assert_eq!(
|
||||
poll.plan.headers.get("authorization"),
|
||||
plan.headers.get("authorization")
|
||||
);
|
||||
assert!(service.apply_read_refresh_projection(
|
||||
&poll,
|
||||
json!({"status": "pending"}).as_object().unwrap()
|
||||
));
|
||||
assert_eq!(
|
||||
service
|
||||
.read_response(Some("openai"), "/v1/videos/local-video")
|
||||
.unwrap()
|
||||
.body_json["status"],
|
||||
"queued"
|
||||
);
|
||||
assert!(service.apply_read_refresh_projection(&poll, json!({
|
||||
"status": "done", "video": {"url": "https://vidgen.x.ai/result.mp4", "duration": 6}
|
||||
}).as_object().unwrap()));
|
||||
let snapshot = service
|
||||
.snapshot_for_route(Some("openai"), "/v1/videos/local-video")
|
||||
.unwrap();
|
||||
assert!(!snapshot.is_active_for_refresh());
|
||||
let record = snapshot.to_upsert_record();
|
||||
assert_eq!(
|
||||
record.status,
|
||||
aether_data_contracts::repository::video_tasks::VideoTaskStatus::Completed
|
||||
);
|
||||
assert_eq!(
|
||||
record.video_url.as_deref(),
|
||||
Some("https://vidgen.x.ai/result.mp4")
|
||||
);
|
||||
assert_eq!(record.duration_seconds, Some(6));
|
||||
let response = snapshot.read_response();
|
||||
assert_eq!(response.body_json["status"], "completed");
|
||||
assert_eq!(response.body_json["progress"], 100);
|
||||
assert_eq!(
|
||||
response.body_json["video_url"],
|
||||
"https://vidgen.x.ai/result.mp4"
|
||||
);
|
||||
let LocalVideoTaskSnapshot::OpenAi(seed) = snapshot else {
|
||||
panic!("OpenAI video expected")
|
||||
};
|
||||
let Some(LocalVideoTaskContentAction::StreamPlan(download)) =
|
||||
seed.build_content_stream_action(None, "download")
|
||||
else {
|
||||
panic!("download expected")
|
||||
};
|
||||
assert_eq!(download.url, "https://vidgen.x.ai/result.mp4");
|
||||
assert!(
|
||||
download.headers.is_empty(),
|
||||
"provider credentials must not be sent to the media CDN"
|
||||
);
|
||||
}
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn xai_video_errors_are_terminal_even_without_a_status() {
|
||||
use crate::{LocalVideoTaskSnapshot, VideoTaskTruthSourceMode};
|
||||
let mut plan =
|
||||
build_internal_finalize_video_plan("xai-create", "openai:video", None).unwrap();
|
||||
plan.url = "https://cli-chat-proxy.grok.com/v1/videos/generations".into();
|
||||
for body in [
|
||||
json!({"code": "content_policy_violation", "error": "Rejected"}),
|
||||
json!({"error": {"code": "content_policy_violation", "message": "Rejected"}}),
|
||||
json!({"status": "failed", "error": "Rejected"}),
|
||||
] {
|
||||
let mut snapshot = VideoTaskTruthSourceMode::RustAuthoritative
|
||||
.prepare_sync_success(
|
||||
"openai_video_create_sync_finalize",
|
||||
json!({"request_id": "xai-request"}).as_object().unwrap(),
|
||||
&Default::default(),
|
||||
&plan,
|
||||
)
|
||||
.unwrap()
|
||||
.to_snapshot();
|
||||
snapshot.apply_provider_body(body.as_object().unwrap());
|
||||
assert!(!snapshot.is_active_for_refresh());
|
||||
assert_eq!(snapshot.read_response().body_json["status"], "failed");
|
||||
let LocalVideoTaskSnapshot::OpenAi(seed) = snapshot else {
|
||||
panic!("OpenAI video expected")
|
||||
};
|
||||
assert!(seed.error_message.is_none());
|
||||
}
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn openai_video_id_takes_precedence_over_xai_alias() {
|
||||
assert_eq!(
|
||||
super::openai_video_provider_task_id(
|
||||
json!({"id": "openai-id", "request_id": "trace-id"})
|
||||
.as_object()
|
||||
.unwrap()
|
||||
),
|
||||
Some("openai-id")
|
||||
);
|
||||
assert_eq!(
|
||||
super::openai_video_provider_task_id(
|
||||
json!({"id": " ", "request_id": "xai-id"})
|
||||
.as_object()
|
||||
.unwrap()
|
||||
),
|
||||
Some("xai-id")
|
||||
);
|
||||
assert_eq!(
|
||||
super::openai_video_provider_task_id(json!({"request_id": " "}).as_object().unwrap()),
|
||||
None
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn builds_local_sync_finalize_read_response_for_supported_video_finalize_kinds() {
|
||||
let delete_response = build_local_sync_finalize_read_response(
|
||||
|
||||
@@ -71,8 +71,16 @@ impl LocalVideoTaskPersistence {
|
||||
.unwrap_or_else(|| plan.request_id.clone()),
|
||||
username: context_text(report_context, "username"),
|
||||
api_key_name: context_text(report_context, "api_key_name"),
|
||||
client_api_format: context_text(report_context, "client_api_format")
|
||||
.unwrap_or_else(|| plan.client_api_format.clone()),
|
||||
client_api_format: if report_context
|
||||
.get("video_client_protocol")
|
||||
.and_then(Value::as_str)
|
||||
== Some("xai")
|
||||
{
|
||||
"xai:video".to_string()
|
||||
} else {
|
||||
context_text(report_context, "client_api_format")
|
||||
.unwrap_or_else(|| plan.client_api_format.clone())
|
||||
},
|
||||
provider_api_format: context_text(report_context, "provider_api_format")
|
||||
.unwrap_or_else(|| plan.provider_api_format.clone()),
|
||||
original_request_body: report_context
|
||||
|
||||
@@ -201,6 +201,13 @@ pub struct LocalVideoTaskPersistence {
|
||||
|
||||
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
|
||||
pub struct OpenAiVideoTaskSeed {
|
||||
/// Preserve existing database identity; older snapshots derive it from the local task ID.
|
||||
#[serde(default)]
|
||||
pub local_short_id: Option<String>,
|
||||
#[serde(default)]
|
||||
pub native_response: Option<Value>,
|
||||
#[serde(default)]
|
||||
pub xai_provider: bool,
|
||||
pub local_task_id: String,
|
||||
pub upstream_task_id: String,
|
||||
pub created_at_unix_ms: u64,
|
||||
|
||||
Reference in New Issue
Block a user