mirror of
https://github.com/Nezumi-2711/9router.git
synced 2026-09-22 13:38:31 +00:00
Async video job proxy mirroring the existing image-generation layer split:
Next routes → src/sse/handlers/videoGeneration.js (auth gate, account
fallback loop, refresh persistence) → open-sse/handlers/videoCore.js
(transparent upstream proxy, 401 refresh-once/retry-once, secret sanitization).
- POST /v1/videos/{generations,edits,extensions}: byte-exact body forward
(JSON + multipart), request_id passthrough, Idempotency-Key forwarded
- GET /v1/videos/{request_id}: status/progress/video.url passthrough
- Register grok-imagine-video (kind: "video"); add "video" to MODEL_TYPE_TO_KIND
so video models stay out of chat lists (also fixes runwayml leak)
- 9router xai video CLI: submit → poll → atomic MP4 download
- No auto-retry of creation POSTs (billable jobs); rotate accounts only on
401/403/429; sanitize Bearer tokens + credential values from errors/logs
Closes #1285
167 lines
6.9 KiB
JavaScript
167 lines
6.9 KiB
JavaScript
import { createErrorResult } from "../utils/error.js";
|
|
import { HTTP_STATUS } from "../config/runtimeConfig.js";
|
|
import { refreshTokenByProvider } from "../services/tokenRefresh.js";
|
|
import { PROVIDER_MEDIA } from "../providers/index.js";
|
|
|
|
// Upstream fetch deadline for video job submission/polling (the job itself is
|
|
// async upstream — this only bounds the HTTP round-trip, not video rendering).
|
|
const VIDEO_FETCH_TIMEOUT_MS = Number(process.env.VIDEO_FETCH_TIMEOUT_MS || 120000);
|
|
|
|
// POST /videos/* creates a billable upstream job. A network error after the
|
|
// request left the socket may still have created the job, so creation is NEVER
|
|
// auto-retried (the only re-send is the auth retry after a 401/403 refresh,
|
|
// which upstream rejects before job creation).
|
|
export const VIDEO_ACTIONS = new Set(["generations", "edits", "extensions"]);
|
|
|
|
export function getVideoConfig(provider) {
|
|
return PROVIDER_MEDIA[provider]?.videoConfig || null;
|
|
}
|
|
|
|
/** Strip bearer tokens / obvious secrets from text destined for clients or logs. */
|
|
export function sanitizeSecrets(text, credentials = null) {
|
|
if (!text) return text;
|
|
let out = String(text).replace(/Bearer\s+[A-Za-z0-9._~+/=-]{8,}/gi, "Bearer [redacted]");
|
|
for (const key of ["accessToken", "refreshToken", "apiKey"]) {
|
|
const secret = credentials?.[key];
|
|
if (typeof secret === "string" && secret.length >= 8) {
|
|
out = out.split(secret).join("[redacted]");
|
|
}
|
|
}
|
|
return out;
|
|
}
|
|
|
|
function buildUpstreamUrl(config, action, requestId) {
|
|
const base = config.baseUrl.replace(/\/$/, "");
|
|
return requestId ? `${base}/${encodeURIComponent(requestId)}` : `${base}/${action}`;
|
|
}
|
|
|
|
function buildHeaders({ token, contentType, idempotencyKey }) {
|
|
const headers = { Accept: "application/json" };
|
|
if (token) headers.Authorization = `Bearer ${token}`;
|
|
if (contentType) headers["Content-Type"] = contentType;
|
|
if (idempotencyKey) headers["Idempotency-Key"] = idempotencyKey;
|
|
return headers;
|
|
}
|
|
|
|
function combineSignals(signal, timeoutMs) {
|
|
const timeoutSignal = typeof AbortSignal?.timeout === "function" ? AbortSignal.timeout(timeoutMs) : null;
|
|
if (signal && timeoutSignal && typeof AbortSignal.any === "function") {
|
|
return AbortSignal.any([signal, timeoutSignal]);
|
|
}
|
|
return signal || timeoutSignal || undefined;
|
|
}
|
|
|
|
/**
|
|
* Transparent proxy for async video jobs (xAI Grok Imagine shape).
|
|
*
|
|
* - Forwards the raw body byte-for-byte (JSON or multipart) — no reshaping.
|
|
* - Passes upstream JSON (request_id, status, video.url, error) back verbatim.
|
|
* - 401/403 with a refresh token: refresh ONCE, retry ONCE. No other retry.
|
|
* - Upstream error text is sanitized before it reaches the client.
|
|
*
|
|
* @param {object} options
|
|
* @param {string} options.provider - Provider id (must have registry videoConfig)
|
|
* @param {"generations"|"edits"|"extensions"|null} options.action - Creation action (POST)
|
|
* @param {string|null} [options.requestId] - Poll target (GET /videos/{id})
|
|
* @param {Buffer|string|null} [options.rawBody] - Exact body to forward
|
|
* @param {string|null} [options.contentType] - Original Content-Type header
|
|
* @param {string|null} [options.idempotencyKey] - Forwarded Idempotency-Key
|
|
* @param {object} options.credentials - { accessToken?, apiKey?, refreshToken?, authType? }
|
|
* @param {AbortSignal} [options.signal] - Client cancellation signal
|
|
* @param {number} [options.timeoutMs]
|
|
* @param {object} [options.log]
|
|
* @param {function} [options.onCredentialsRefreshed]
|
|
* @returns {Promise<{ success: boolean, response: Response, status?: number, error?: string }>}
|
|
*/
|
|
export async function handleVideoProxyCore({
|
|
provider,
|
|
action = null,
|
|
requestId = null,
|
|
rawBody = null,
|
|
contentType = null,
|
|
idempotencyKey = null,
|
|
credentials,
|
|
signal,
|
|
timeoutMs = VIDEO_FETCH_TIMEOUT_MS,
|
|
log,
|
|
onCredentialsRefreshed,
|
|
}) {
|
|
const config = getVideoConfig(provider);
|
|
if (!config) {
|
|
return createErrorResult(HTTP_STATUS.BAD_REQUEST, `Provider '${provider}' does not support video generation`);
|
|
}
|
|
if (!requestId && !VIDEO_ACTIONS.has(action)) {
|
|
return createErrorResult(HTTP_STATUS.BAD_REQUEST, `Unknown video action: ${action}`);
|
|
}
|
|
|
|
const method = requestId ? "GET" : "POST";
|
|
const url = buildUpstreamUrl(config, action, requestId);
|
|
const fetchSignal = combineSignals(signal, timeoutMs);
|
|
|
|
const doFetch = (token) =>
|
|
fetch(url, {
|
|
method,
|
|
headers: buildHeaders({ token, contentType: method === "POST" ? contentType : null, idempotencyKey: method === "POST" ? idempotencyKey : null }),
|
|
body: method === "POST" ? rawBody : undefined,
|
|
signal: fetchSignal,
|
|
});
|
|
|
|
let upstream;
|
|
try {
|
|
upstream = await doFetch(credentials?.accessToken || credentials?.apiKey);
|
|
} catch (error) {
|
|
if (error?.name === "AbortError" || error?.name === "TimeoutError") {
|
|
return createErrorResult(HTTP_STATUS.REQUEST_TIMEOUT, `[${provider}] video ${method} aborted: ${error.message}`);
|
|
}
|
|
// Never re-send a creation POST on network error — the job may already exist upstream.
|
|
return createErrorResult(HTTP_STATUS.BAD_GATEWAY, sanitizeSecrets(`[${provider}] video upstream fetch failed: ${error.message}`, credentials));
|
|
}
|
|
|
|
// 401/403 → refresh once → retry once (OAuth accounts only; API keys can't refresh)
|
|
if (
|
|
(upstream.status === HTTP_STATUS.UNAUTHORIZED || upstream.status === HTTP_STATUS.FORBIDDEN) &&
|
|
credentials?.refreshToken
|
|
) {
|
|
let refreshed = null;
|
|
try {
|
|
refreshed = await refreshTokenByProvider(provider, credentials, log);
|
|
} catch (error) {
|
|
log?.warn?.("TOKEN", `${provider} | video refresh error: ${sanitizeSecrets(error.message, credentials)}`);
|
|
}
|
|
if (refreshed?.accessToken) {
|
|
log?.info?.("TOKEN", `${provider.toUpperCase()} | refreshed for video ${method}`);
|
|
Object.assign(credentials, refreshed);
|
|
if (onCredentialsRefreshed) await onCredentialsRefreshed(refreshed);
|
|
try {
|
|
await upstream.body?.cancel?.();
|
|
} catch { /* noop */ }
|
|
try {
|
|
upstream = await doFetch(credentials.accessToken || credentials.apiKey);
|
|
} catch (error) {
|
|
return createErrorResult(HTTP_STATUS.BAD_GATEWAY, sanitizeSecrets(`[${provider}] video retry after refresh failed: ${error.message}`, credentials));
|
|
}
|
|
} else {
|
|
log?.warn?.("TOKEN", `${provider.toUpperCase()} | video refresh failed — account needs re-auth`);
|
|
}
|
|
}
|
|
|
|
const bodyText = await upstream.text().catch(() => "");
|
|
|
|
if (!upstream.ok) {
|
|
const message = sanitizeSecrets(bodyText || `HTTP ${upstream.status}`, credentials);
|
|
return createErrorResult(upstream.status, `[${provider}] ${message.slice(0, 2000)}`);
|
|
}
|
|
|
|
// Success: pass the upstream JSON through untouched (request_id / status / video.url).
|
|
return {
|
|
success: true,
|
|
response: new Response(bodyText, {
|
|
status: upstream.status,
|
|
headers: {
|
|
"Content-Type": upstream.headers.get("content-type") || "application/json",
|
|
"Access-Control-Allow-Origin": "*",
|
|
},
|
|
}),
|
|
};
|
|
}
|