mirror of
https://github.com/Nezumi-2711/9router.git
synced 2026-09-22 20:00:47 +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
224 lines
8.7 KiB
JavaScript
224 lines
8.7 KiB
JavaScript
import {
|
|
getProviderCredentials,
|
|
markAccountUnavailable,
|
|
clearAccountError,
|
|
extractApiKey,
|
|
isValidApiKey,
|
|
} from "../services/auth.js";
|
|
import { getSettings } from "@/lib/localDb";
|
|
import { getModelInfo } from "../services/model.js";
|
|
import { handleVideoProxyCore, getVideoConfig, sanitizeSecrets } from "open-sse/handlers/videoCore.js";
|
|
import { errorResponse, unavailableResponse } from "open-sse/utils/error.js";
|
|
import { HTTP_STATUS } from "open-sse/config/runtimeConfig.js";
|
|
import { updateProviderCredentials, checkAndRefreshToken } from "../services/tokenRefresh.js";
|
|
import * as log from "../utils/logger.js";
|
|
|
|
// Video generation is xAI-only today; requests without a provider prefix
|
|
// (bare model id, or multipart bodies we deliberately don't parse) land here.
|
|
const DEFAULT_VIDEO_PROVIDER = "xai";
|
|
|
|
// Creation POSTs are billable jobs — only rotate to another account for
|
|
// errors that upstream rejects BEFORE creating a job (auth/quota). A 5xx may
|
|
// have created the job, so it is returned to the caller instead of re-sent.
|
|
const CREATE_ROTATION_STATUSES = new Set([
|
|
HTTP_STATUS.UNAUTHORIZED,
|
|
HTTP_STATUS.FORBIDDEN,
|
|
HTTP_STATUS.RATE_LIMITED,
|
|
]);
|
|
|
|
async function requireValidApiKey(request) {
|
|
const apiKey = extractApiKey(request);
|
|
const settings = await getSettings();
|
|
if (settings.requireApiKey) {
|
|
if (!apiKey) return errorResponse(HTTP_STATUS.UNAUTHORIZED, "Missing API key");
|
|
const valid = await isValidApiKey(apiKey);
|
|
if (!valid) return errorResponse(HTTP_STATUS.UNAUTHORIZED, "Invalid API key");
|
|
}
|
|
return null;
|
|
}
|
|
|
|
/**
|
|
* Read the request body once, byte-preserving.
|
|
* JSON bodies are additionally parsed so the `model` provider prefix can be
|
|
* resolved (and stripped) — everything else is forwarded exactly as received.
|
|
*/
|
|
async function readForwardableBody(request) {
|
|
const contentType = request.headers.get("content-type") || "";
|
|
if (contentType.includes("application/json")) {
|
|
const raw = await request.text();
|
|
let parsed;
|
|
try {
|
|
parsed = JSON.parse(raw);
|
|
} catch {
|
|
return { error: errorResponse(HTTP_STATUS.BAD_REQUEST, "Invalid JSON body") };
|
|
}
|
|
return { raw, parsed, contentType };
|
|
}
|
|
// Multipart (or any other content type): forward the exact bytes — parsing
|
|
// and re-encoding FormData would change the multipart boundary.
|
|
const buf = Buffer.from(await request.arrayBuffer());
|
|
return { raw: buf, parsed: null, contentType };
|
|
}
|
|
|
|
async function resolveVideoProvider(parsedBody) {
|
|
if (!parsedBody?.model) return { provider: DEFAULT_VIDEO_PROVIDER, model: null };
|
|
|
|
const modelStr = String(parsedBody.model);
|
|
const modelInfo = await getModelInfo(modelStr);
|
|
if (!modelInfo.provider) {
|
|
return { error: errorResponse(HTTP_STATUS.BAD_REQUEST, "Combos are not supported for video generation") };
|
|
}
|
|
if (!getVideoConfig(modelInfo.provider)) {
|
|
// Bare model ids (no explicit "provider/" prefix) fall back to the default
|
|
// video provider — the prefix-less inference targets chat providers only.
|
|
if (!modelStr.includes("/")) {
|
|
return { provider: DEFAULT_VIDEO_PROVIDER, model: modelStr };
|
|
}
|
|
return { error: errorResponse(HTTP_STATUS.BAD_REQUEST, `Provider '${modelInfo.provider}' does not support video generation`) };
|
|
}
|
|
return { provider: modelInfo.provider, model: modelInfo.model };
|
|
}
|
|
|
|
function withConnectionHeader(response, connectionId) {
|
|
if (!connectionId) return response;
|
|
const headers = new Headers(response.headers);
|
|
// Video jobs are account-bound upstream — clients echo this back as
|
|
// `x-connection-id` on GET polls so the same account is used.
|
|
headers.set("x-9router-connection-id", String(connectionId));
|
|
return new Response(response.body, { status: response.status, headers });
|
|
}
|
|
|
|
/**
|
|
* POST /v1/videos/{generations|edits|extensions} — async job creation proxy.
|
|
*/
|
|
export async function handleVideoCreate(request, action) {
|
|
const authError = await requireValidApiKey(request);
|
|
if (authError) return authError;
|
|
|
|
const bodyInfo = await readForwardableBody(request);
|
|
if (bodyInfo.error) return bodyInfo.error;
|
|
|
|
const resolved = await resolveVideoProvider(bodyInfo.parsed);
|
|
if (resolved.error) return resolved.error;
|
|
const { provider, model } = resolved;
|
|
|
|
// Strip the provider prefix (e.g. "xai/grok-imagine-video") before forwarding;
|
|
// otherwise forward the original bytes untouched.
|
|
let forwardBody = bodyInfo.raw;
|
|
if (bodyInfo.parsed && model && bodyInfo.parsed.model !== model) {
|
|
forwardBody = JSON.stringify({ ...bodyInfo.parsed, model });
|
|
}
|
|
|
|
const preferredConnectionId = request.headers.get("x-connection-id") || null;
|
|
const idempotencyKey = request.headers.get("idempotency-key") || null;
|
|
|
|
const excludeConnectionIds = new Set();
|
|
let lastError = null;
|
|
let lastStatus = null;
|
|
|
|
while (true) {
|
|
const credentials = await getProviderCredentials(provider, excludeConnectionIds, model, { preferredConnectionId });
|
|
|
|
if (!credentials || credentials.allRateLimited) {
|
|
if (credentials?.allRateLimited) {
|
|
const errorMsg = lastError || credentials.lastError || "Unavailable";
|
|
const status = lastStatus || Number(credentials.lastErrorCode) || HTTP_STATUS.SERVICE_UNAVAILABLE;
|
|
return unavailableResponse(status, `[${provider}/${model || "video"}] ${errorMsg}`, credentials.retryAfter, credentials.retryAfterHuman);
|
|
}
|
|
if (excludeConnectionIds.size === 0) {
|
|
return errorResponse(HTTP_STATUS.BAD_REQUEST, `No credentials for provider: ${provider}`);
|
|
}
|
|
return errorResponse(lastStatus || HTTP_STATUS.SERVICE_UNAVAILABLE, lastError || "All accounts unavailable");
|
|
}
|
|
|
|
const refreshedCredentials = await checkAndRefreshToken(provider, credentials);
|
|
|
|
const result = await handleVideoProxyCore({
|
|
provider,
|
|
action,
|
|
rawBody: forwardBody,
|
|
contentType: bodyInfo.contentType || null,
|
|
idempotencyKey,
|
|
credentials: refreshedCredentials,
|
|
signal: request.signal,
|
|
log,
|
|
onCredentialsRefreshed: async (newCreds) => {
|
|
await updateProviderCredentials(credentials.connectionId, {
|
|
accessToken: newCreds.accessToken,
|
|
refreshToken: newCreds.refreshToken,
|
|
providerSpecificData: newCreds.providerSpecificData,
|
|
testStatus: "active",
|
|
});
|
|
},
|
|
});
|
|
|
|
if (result.success) {
|
|
await clearAccountError(credentials.connectionId, credentials, model);
|
|
log.info("VIDEO", `${provider.toUpperCase()} | ${action} accepted (connection ${credentials.connectionId})`);
|
|
return withConnectionHeader(result.response, credentials.connectionId);
|
|
}
|
|
|
|
// Record the failure (dashboard shows lastError/errorCode → user sees re-auth is needed)
|
|
const { shouldFallback } = await markAccountUnavailable(
|
|
credentials.connectionId, result.status, sanitizeSecrets(result.error, refreshedCredentials), provider, model
|
|
);
|
|
|
|
if (shouldFallback && CREATE_ROTATION_STATUSES.has(result.status)) {
|
|
excludeConnectionIds.add(credentials.connectionId);
|
|
lastError = result.error;
|
|
lastStatus = result.status;
|
|
continue;
|
|
}
|
|
|
|
return result.response;
|
|
}
|
|
}
|
|
|
|
/**
|
|
* GET /v1/videos/{request_id} — poll job status.
|
|
* Jobs are account-bound upstream, so no cross-account rotation here: the
|
|
* caller pins the creating account via `x-connection-id` (returned on create).
|
|
*/
|
|
export async function handleVideoGet(request, requestId) {
|
|
const authError = await requireValidApiKey(request);
|
|
if (authError) return authError;
|
|
|
|
if (!requestId) return errorResponse(HTTP_STATUS.BAD_REQUEST, "Missing video request id");
|
|
|
|
const provider = DEFAULT_VIDEO_PROVIDER;
|
|
const preferredConnectionId = request.headers.get("x-connection-id") || null;
|
|
|
|
const credentials = await getProviderCredentials(provider, null, null, { preferredConnectionId });
|
|
if (!credentials || credentials.allRateLimited) {
|
|
return errorResponse(HTTP_STATUS.BAD_REQUEST, `No credentials for provider: ${provider}`);
|
|
}
|
|
|
|
const refreshedCredentials = await checkAndRefreshToken(provider, credentials);
|
|
|
|
const result = await handleVideoProxyCore({
|
|
provider,
|
|
requestId,
|
|
credentials: refreshedCredentials,
|
|
signal: request.signal,
|
|
log,
|
|
onCredentialsRefreshed: async (newCreds) => {
|
|
await updateProviderCredentials(credentials.connectionId, {
|
|
accessToken: newCreds.accessToken,
|
|
refreshToken: newCreds.refreshToken,
|
|
providerSpecificData: newCreds.providerSpecificData,
|
|
testStatus: "active",
|
|
});
|
|
},
|
|
});
|
|
|
|
if (result.success) {
|
|
await clearAccountError(credentials.connectionId, credentials, null);
|
|
return withConnectionHeader(result.response, credentials.connectionId);
|
|
}
|
|
|
|
await markAccountUnavailable(
|
|
credentials.connectionId, result.status, sanitizeSecrets(result.error, refreshedCredentials), provider, null
|
|
);
|
|
return result.response;
|
|
}
|