feat: integrate platform backend and application interfaces

This commit is contained in:
Qiufeng
2026-09-02 21:09:47 +08:00
parent cd9ac1af70
commit b16390dd41
648 changed files with 102488 additions and 17737 deletions
+94
View File
@@ -0,0 +1,94 @@
import { randomUUID } from "node:crypto";
import type { MessageOutboxRecord, Store } from "../store.ts";
import { decryptSecret, encryptSecret, hashChallenge } from "../shared/auth.ts";
import { validateProviderUrlResolved } from "./provider.ts";
import { safeOutboundFetch } from "../infra/outbound-url.ts";
export type MessageChannel = "email" | "sms";
export type MessagePurpose = "register" | "reset" | "notification";
const messageWorkerId = `message-worker-${process.pid}-${randomUUID().slice(0, 8)}`;
const MAX_MESSAGE_ATTEMPTS = 5;
const MESSAGE_LEASE_MS = 30_000;
export function enqueueMessage(store: Store, input: { channel: MessageChannel; target: string; purpose: MessagePurpose; templateData?: Record<string, string>; idempotencyKey?: string; deferProcessing?: boolean }): MessageOutboxRecord {
const idempotencyKey = input.idempotencyKey?.trim() || undefined;
if (idempotencyKey) {
const existing = [...store.messageOutbox.values()].find((item) => item.idempotencyKey === idempotencyKey);
if (existing) return existing;
}
const now = new Date().toISOString();
const record: MessageOutboxRecord = { id: randomUUID(), idempotencyKey, channel: input.channel, targetHash: hashChallenge(input.target, store.channelEncryptionKey), targetEncrypted: encryptSecret(input.target, store.channelEncryptionKey), payloadEncrypted: input.templateData ? encryptSecret(JSON.stringify(input.templateData), store.channelEncryptionKey) : undefined, purpose: input.purpose, status: "queued", attempts: 0, createdAt: now, updatedAt: now };
store.messageOutbox.set(record.id, record); store.persist();
if (!input.deferProcessing) setTimeout(() => { void processMessageOutbox(store); }, 0).unref?.();
return record;
}
async function sendMessage(store: Store, record: MessageOutboxRecord): Promise<{ providerMessageId?: string }> {
const payloadText = record.payloadEncrypted ? decryptSecret(record.payloadEncrypted, store.channelEncryptionKey) : undefined;
const target = record.targetEncrypted ? decryptSecret(record.targetEncrypted, store.channelEncryptionKey) : undefined;
const providers = [...store.messageProviders.values()].filter((provider) => provider.channel === record.channel && provider.enabled !== false).sort((a, b) => Number(a.priority || 0) - Number(b.priority || 0));
const candidates = providers.length ? providers : process.env.NODE_ENV === "production" ? [] : [{ id: "mock", channel: record.channel, enabled: true }];
let sentBy: Record<string, unknown> | undefined; let providerMessageId: string | undefined; let lastError = "message provider unavailable";
for (const provider of candidates) {
const providerId = String(provider.id || "mock");
if (process.env.MIRAGENFLOW_MESSAGE_FAIL === record.channel || process.env.MIRAGENFLOW_MESSAGE_FAIL === providerId || provider.fixtureFailure === "retryable") { lastError = `${providerId}: message provider unavailable`; continue; }
try {
const endpoint = typeof provider.endpoint === "string" ? provider.endpoint.trim() : "";
if (!endpoint) { if (process.env.NODE_ENV === "production") throw new Error(`${providerId}: provider endpoint is not configured`); }
else {
const url = await validateProviderUrlResolved(endpoint); const headers = new Headers({ "content-type": "application/json", accept: "application/json" });
const secret = typeof provider.secretRef === "string" ? decryptSecret(provider.secretRef, store.channelEncryptionKey) : undefined; if (secret) headers.set("authorization", `Bearer ${secret}`);
if (record.idempotencyKey) headers.set("idempotency-key", `miragenflow-message:${record.idempotencyKey}`);
const response = await safeOutboundFetch(url, { method: "POST", headers, body: JSON.stringify({ to: target, channel: record.channel, purpose: record.purpose, data: payloadText ? JSON.parse(payloadText) : undefined }) }, { timeoutMs: 15_000, maxBytes: 1024 * 1024 });
if (!response.ok) throw new Error(`${providerId}: provider HTTP ${response.status}`);
const responseBody = await response.text(); let parsed: Record<string, unknown> = {}; try { parsed = responseBody ? JSON.parse(responseBody) as Record<string, unknown> : {}; } catch { /* empty/text body */ }
providerMessageId = typeof parsed.id === "string" ? parsed.id : response.headers.get("x-message-id") || `${providerId}-${randomUUID()}`;
}
sentBy = provider; break;
} catch (error) { lastError = error instanceof Error ? error.message : `${providerId}: provider request failed`; }
}
if (!sentBy) throw new Error(lastError);
return { providerMessageId: providerMessageId || `${String(sentBy.id || "mock")}-${randomUUID()}` };
}
async function processPostgresOutbox(store: Store) {
const claim = store.repository.claimMessageOutbox; const complete = store.repository.completeMessageOutbox; const renew = store.repository.renewMessageOutbox; if (!claim || !complete) return;
const worker = messageWorkerId; const records = await claim.call(store.repository, worker, MESSAGE_LEASE_MS);
for (const record of records) {
store.messageOutbox.set(record.id, record);
let leaseLost = false;
const renewTimer = renew ? setInterval(() => {
void renew.call(store.repository, record.id, worker, MESSAGE_LEASE_MS).then((result) => {
if (!result.renewed) leaseLost = true;
else if (result.leaseExpiresAt) record.leaseExpiresAt = result.leaseExpiresAt;
}).catch(() => { leaseLost = true; });
}, Math.max(1_000, Math.floor(MESSAGE_LEASE_MS / 3))) : undefined;
renewTimer?.unref?.();
try {
const result = await sendMessage(store, record); const completed = await complete.call(store.repository, record.id, worker, { success: true, providerMessageId: result.providerMessageId }); if (completed && !leaseLost) { record.status = "sent"; record.providerMessageId = result.providerMessageId; record.lastError = undefined; record.nextAttemptAt = undefined; record.leaseOwner = undefined; record.leaseExpiresAt = undefined; record.updatedAt = new Date().toISOString(); }
} catch (error) {
const nextAttemptAt = new Date(Date.now() + Math.min(15 * 60_000, 2 ** record.attempts * 1_000)).toISOString(); const completed = await complete.call(store.repository, record.id, worker, { success: false, error: error instanceof Error ? error.message : String(error), nextAttemptAt }); if (completed && !leaseLost) { record.status = record.attempts >= MAX_MESSAGE_ATTEMPTS ? "dead" : "failed"; record.lastError = error instanceof Error ? error.message : String(error); record.nextAttemptAt = nextAttemptAt; record.leaseOwner = undefined; record.leaseExpiresAt = undefined; record.updatedAt = new Date().toISOString(); }
} finally {
if (renewTimer) clearInterval(renewTimer);
}
}
}
export async function processMessageOutbox(store: Store) {
if (store.repository.adapter === "postgres" && store.repository.claimMessageOutbox && store.repository.completeMessageOutbox) { await processPostgresOutbox(store); return; }
const now = Date.now();
for (const record of store.messageOutbox.values()) {
if (!(record.status === "queued" || record.status === "failed" || record.status === "pending") || (record.nextAttemptAt && Date.parse(record.nextAttemptAt) > now)) continue;
if (record.attempts >= MAX_MESSAGE_ATTEMPTS) { record.status = "dead"; record.updatedAt = new Date().toISOString(); continue; }
record.status = "sending"; record.attempts += 1; record.leaseOwner = messageWorkerId; record.leaseExpiresAt = new Date(now + MESSAGE_LEASE_MS).toISOString(); record.updatedAt = new Date().toISOString(); store.persist();
try { const result = await sendMessage(store, record); record.status = "sent"; record.providerMessageId = result.providerMessageId; record.lastError = undefined; record.nextAttemptAt = undefined; }
catch (error) { record.status = record.attempts >= MAX_MESSAGE_ATTEMPTS ? "dead" : "failed"; record.lastError = error instanceof Error ? error.message : String(error); record.nextAttemptAt = new Date(Date.now() + Math.min(15 * 60_000, 2 ** record.attempts * 1_000)).toISOString(); }
finally { record.leaseExpiresAt = undefined; record.leaseOwner = undefined; record.updatedAt = new Date().toISOString(); store.persist(); }
}
}
export function recoverMessageOutbox(store: Store) {
if (store.repository.adapter === "postgres") return;
const now = Date.now();
for (const record of store.messageOutbox.values()) if (record.status === "sending" && record.leaseExpiresAt && Date.parse(record.leaseExpiresAt) <= now) { record.status = "failed"; record.nextAttemptAt = new Date(now).toISOString(); record.updatedAt = new Date(now).toISOString(); }
}
+293
View File
@@ -0,0 +1,293 @@
import { readFile } from "node:fs/promises";
import type { TaskType } from "@miragenflow/contracts";
import type { ProviderChannel, Store } from "../store.ts";
import { imageModelFor, normalizeImageParams } from "../image-options.ts";
import { decryptSecret } from "../shared/auth.ts";
import { safeStagingPath } from "../infra/staging.ts";
import { isPublicAddress, resolvePublicHttpsUrl, safeOutboundFetch } from "../infra/outbound-url.ts";
export type ProviderOutput = { mimeType: string; data: string; metadata?: { width?: number; height?: number; format?: string; size?: string; revisedPrompt?: string; usage?: Record<string, number>; source?: "base64" | "url" } };
export type ProviderResult = { status: "succeeded" | "unknown" | "failed"; providerRequestId?: string; outputs?: ProviderOutput[]; errorCode?: string; retryable?: boolean };
export type ProviderInput = { taskId: string; taskType: TaskType; publicModelId?: string; prompt?: string; count?: number; references?: string[]; referenceImages?: Array<{ objectId: string; seq: number; name: string }>; maskObjectId?: string; params?: Record<string, unknown>; platformIdempotencyKey?: string; onRequestAttempted?: () => void };
export type ProviderCancelResult = "confirmed" | "pending" | "unsupported";
function resolveChannelModel(channel: ProviderChannel, publicModelId?: string) {
const mappedModel = channel.modelMappings?.find((mapping) => mapping.displayModelId === publicModelId)?.requestModelId;
const directModel = publicModelId && (channel.enabledModelIds?.includes(publicModelId) || channel.providerModelId === publicModelId) ? publicModelId : undefined;
if (channel.providerType === "openai-images" && publicModelId && !mappedModel && !directModel) {
throw Object.assign(new Error("public model is not mapped to this channel"), { errorCode: "PROVIDER_MODEL_NOT_CONFIGURED", retryable: false });
}
return { model: mappedModel || directModel || channel.enabledModelIds?.[0] || channel.providerModelId, exact: Boolean(mappedModel || directModel) };
}
const chineseImageNumbers = ["一", "二", "三", "四", "五", "六", "七", "八", "九", "十", "十一", "十二", "十三", "十四", "十五", "十六"];
const imageMentionPattern = new RegExp(`@图片(${chineseImageNumbers.slice().sort((left, right) => right.length - left.length).join("|")})`, "g");
function mentionedImageNames(prompt: string) {
return new Set([...prompt.matchAll(imageMentionPattern)].map((match) => `图片${match[1]}`));
}
export function allowLocalProviderUrls() {
const configured = process.env.MIRAGENFLOW_ALLOW_LOCAL_PROVIDER_URLS;
return configured === "true" || (configured === undefined && process.env.NODE_ENV !== "production");
}
export function validateProviderUrl(value: string) { const url = new URL(value); const allowPrivateNetwork = allowLocalProviderUrls(); const protocolAllowed = url.protocol === "https:" || (allowPrivateNetwork && url.protocol === "http:"); if (!protocolAllowed || url.username || url.password) throw new Error(allowPrivateNetwork ? "provider URL must use HTTP or HTTPS without embedded credentials" : "provider URL must use HTTPS without embedded credentials"); const literal = url.hostname.replace(/^\[|\]$/g, ""); if (!allowPrivateNetwork && /^[\d.:a-f]+$/i.test(literal) && !isPublicAddress(literal)) throw new Error("provider URL targets a non-public address"); return url; }
export async function validateProviderUrlResolved(value: string | URL) { return (await resolvePublicHttpsUrl(validateProviderUrl(String(value)), undefined, { allowPrivateNetwork: allowLocalProviderUrls() })).url; }
const validateResolvedUrl = validateProviderUrlResolved;
function providerEndpoint(base: URL, path: string) {
const url = new URL(base);
url.search = "";
url.hash = "";
const prefix = url.pathname.replace(/\/+$/, "");
const normalizedPath = path.replace(/^\/+/, "").replace(/^v1(?:\/|$)/, "");
const apiPrefix = prefix === "/v1" || prefix.endsWith("/v1") ? prefix : `${prefix}/v1`;
url.pathname = `${apiPrefix}/${normalizedPath}`.replace(/\/{2,}/g, "/").replace(/\/$/, "") || "/";
return url;
}
/** Return the path used by the model-list probe without exposing credentials. */
export function providerModelsPath(value: string | URL) {
return providerEndpoint(value instanceof URL ? value : new URL(value), "v1/models").pathname;
}
function providerRequestId(response: Response) {
for (const header of ["x-request-id", "request-id", "openai-request-id", "x-correlation-id"]) {
const value = response.headers.get(header)?.trim();
if (value) return value;
}
return undefined;
}
async function providerFetch(url: URL, secret: string | undefined, init: RequestInit, timeoutMs = 60_000, onRequestAttempted?: () => void) {
const headers = new Headers(init.headers); if (secret) headers.set("authorization", `Bearer ${secret}`);
return safeOutboundFetch(url, { ...init, headers }, { timeoutMs, maxBytes: 50 * 1024 * 1024, allowPrivateNetwork: allowLocalProviderUrls(), onRequestAttempted });
}
async function jsonResponse(response: Response) { const text = await response.text(); if (Buffer.byteLength(text) > 50 * 1024 * 1024) throw new Error("provider response too large"); let parsed: Record<string, unknown> | undefined; try { parsed = text ? JSON.parse(text) as Record<string, unknown> : undefined; } catch { /* handled below */ } if (!response.ok) throw providerHttpError(response.status, parsed); if (!parsed) throw new Error("provider response is not valid JSON"); return parsed; }
function providerHttpError(status: number, body?: Record<string, unknown>) { const error = new Error(`provider http ${status}`) as Error & { retryable?: boolean; errorCode?: string; providerStatusCode?: number }; const nested = body?.error && typeof body.error === "object" ? body.error as Record<string, unknown> : body; const code = `${nested?.code || ""} ${nested?.type || ""}`.toLowerCase(); error.providerStatusCode = status; error.retryable = status === 408 || status === 429 || status >= 500; error.errorCode = status === 401 || status === 403 ? "PROVIDER_AUTH" : /content|safety|moderation|policy/.test(code) ? "PROVIDER_CONTENT_REJECTED" : /model_not_found|invalid_request|invalid_model/.test(code) ? "PROVIDER_INVALID_REQUEST" : "PROVIDER_HTTP"; return error; }
export function imageMimeFromBytes(bytes: Buffer) {
if (bytes.length >= 24 && bytes.subarray(0, 8).equals(Buffer.from([137, 80, 78, 71, 13, 10, 26, 10])) && bytes.subarray(12, 16).toString() === "IHDR") return "image/png";
if (bytes.length >= 4 && bytes[0] === 0xff && bytes[1] === 0xd8 && bytes[2] === 0xff && bytes.lastIndexOf(Buffer.from([0xff, 0xd9])) >= 2) return "image/jpeg";
if (bytes.length >= 30 && bytes.subarray(0, 4).toString() === "RIFF" && bytes.subarray(8, 12).toString() === "WEBP" && imageDimensions(bytes, "image/webp")) return "image/webp";
return undefined;
}
function invalidProviderOutput(message: string, cause?: unknown) {
const error = Object.assign(new Error(message), { errorCode: "PROVIDER_INVALID_OUTPUT", retryable: false });
if (cause !== undefined) (error as Error & { cause?: unknown }).cause = cause;
return error;
}
function decodeImageBase64(value: string) {
const encoded = value.trim();
if (!encoded || encoded.length % 4 !== 0 || !/^[A-Za-z0-9+/]*={0,2}$/.test(encoded)) throw invalidProviderOutput("provider image base64 is invalid");
const bytes = Buffer.from(encoded, "base64");
const mimeType = imageMimeFromBytes(bytes);
if (!bytes.length || !mimeType) throw invalidProviderOutput("provider image bytes are invalid");
if (!imageDimensions(bytes, mimeType)) throw invalidProviderOutput("provider image dimensions are invalid");
return { bytes, mimeType };
}
async function downloadOutput(url: string) {
try {
const response = await providerFetch(await validateResolvedUrl(url), undefined, { method: "GET" }, 30_000);
if (!response.ok) throw new Error(`provider output http ${response.status}`);
const bytes = Buffer.from(await response.arrayBuffer());
if (bytes.byteLength > 50 * 1024 * 1024) throw new Error("provider output too large");
const detectedMime = imageMimeFromBytes(bytes);
const declaredMime = (response.headers.get("content-type") || "").split(";")[0].trim().toLowerCase();
const normalizedDeclaredMime = declaredMime === "image/jpg" ? "image/jpeg" : declaredMime;
if (!detectedMime || (normalizedDeclaredMime && normalizedDeclaredMime !== "application/octet-stream" && normalizedDeclaredMime !== detectedMime) || !imageDimensions(bytes, detectedMime)) throw new Error("provider output is not a valid image");
return { mimeType: detectedMime, data: bytes.toString("base64"), source: "url" as const };
} catch (error) {
if ((error as Error & { errorCode?: string }).errorCode === "PROVIDER_INVALID_OUTPUT") throw error;
throw invalidProviderOutput("provider output URL is invalid or unavailable", error);
}
}
export function imageDimensions(bytes: Buffer, mimeType: string) {
if (mimeType === "image/png" && bytes.length >= 24 && bytes.subarray(0, 8).equals(Buffer.from([137, 80, 78, 71, 13, 10, 26, 10])) && bytes.subarray(12, 16).toString() === "IHDR") {
const width = bytes.readUInt32BE(16); const height = bytes.readUInt32BE(20);
return width > 0 && height > 0 ? { width, height } : undefined;
}
if (mimeType === "image/webp" && bytes.length >= 30 && bytes.subarray(0, 4).toString() === "RIFF" && bytes.subarray(8, 12).toString() === "WEBP") {
const marker = bytes.subarray(12, 16).toString();
const chunkSize = bytes.length >= 20 ? bytes.readUInt32LE(16) : 0;
if (chunkSize > bytes.length - 20) return undefined;
if (marker === "VP8X" && chunkSize >= 10 && bytes.length >= 30) return { width: 1 + bytes.readUIntLE(24, 3), height: 1 + bytes.readUIntLE(27, 3) };
if (marker === "VP8L" && chunkSize >= 5 && bytes.length >= 25 && bytes[20] === 0x2f) { const bits = bytes.readUInt32LE(21); return { width: 1 + (bits & 0x3fff), height: 1 + ((bits >>> 14) & 0x3fff) }; }
if (marker === "VP8 " && chunkSize >= 10 && bytes.length >= 30 && bytes[23] === 0x9d && bytes[24] === 0x01 && bytes[25] === 0x2a) return { width: bytes.readUInt16LE(26) & 0x3fff, height: bytes.readUInt16LE(28) & 0x3fff };
}
if (mimeType === "image/jpeg" && bytes.length > 4 && bytes[0] === 0xff && bytes[1] === 0xd8) {
let offset = 2;
while (offset + 9 < bytes.length) {
if (bytes[offset] !== 0xff) { offset += 1; continue; }
const marker = bytes[offset + 1]; const length = bytes.readUInt16BE(offset + 2);
if (length < 2 || offset + 2 + length > bytes.length) break;
if ((marker >= 0xc0 && marker <= 0xc3) || (marker >= 0xc5 && marker <= 0xc7) || (marker >= 0xc9 && marker <= 0xcb) || (marker >= 0xcd && marker <= 0xcf)) { const width = bytes.readUInt16BE(offset + 7); const height = bytes.readUInt16BE(offset + 5); return width > 0 && height > 0 ? { width, height } : undefined; }
offset += 2 + length;
}
}
return undefined;
}
function imageOutputMetadata(item: Record<string, unknown>, data: Record<string, unknown>, mimeType: string, source: "base64" | "url", bytes?: Buffer) { const size = typeof item.size === "string" ? item.size : typeof data.size === "string" ? data.size : undefined; const match = size?.match(/^(\d+)x(\d+)$/); const dimensions = bytes ? imageDimensions(bytes, mimeType) : undefined; return { width: dimensions?.width || (match ? Number(match[1]) : undefined), height: dimensions?.height || (match ? Number(match[2]) : undefined), format: mimeType.split("/")[1], size: dimensions ? `${dimensions.width}x${dimensions.height}` : size, revisedPrompt: typeof item.revised_prompt === "string" ? item.revised_prompt : undefined, usage: data.usage && typeof data.usage === "object" ? Object.fromEntries(Object.entries(data.usage as Record<string, unknown>).filter(([, value]) => typeof value === "number")) as Record<string, number> : undefined, source }; }
async function referenceBlob(store: Store, objectId: string) {
const object = store.objects.get(objectId);
if (!object || object.revoked || !object.expiresAt || Date.parse(object.expiresAt) <= Date.now()) throw Object.assign(new Error("reference object unavailable"), { errorCode: "PROVIDER_REFERENCE_UNAVAILABLE", retryable: false });
let bytes: Buffer;
try { bytes = object.stagingKey ? await readFile(safeStagingPath(store.stagingDir, object.stagingKey)) : Buffer.from(object.data || "", "base64"); }
catch (error) { throw Object.assign(new Error("reference object unavailable"), { errorCode: "PROVIDER_REFERENCE_UNAVAILABLE", retryable: false, cause: error }); }
if (!bytes.length || !/^image\/(png|jpeg|webp)$/.test(object.mimeType)) throw Object.assign(new Error("reference object is not a supported image"), { errorCode: "PROVIDER_REFERENCE_INVALID", retryable: false });
if (imageMimeFromBytes(bytes) !== object.mimeType || !imageDimensions(bytes, object.mimeType)) throw Object.assign(new Error("reference object is not a valid image"), { errorCode: "PROVIDER_REFERENCE_INVALID", retryable: false });
if (bytes.byteLength > 50 * 1024 * 1024) throw Object.assign(new Error("reference object is too large"), { errorCode: "PROVIDER_REFERENCE_TOO_LARGE", retryable: false });
const copy = new Uint8Array(bytes.byteLength);
copy.set(bytes);
return new Blob([copy.buffer as ArrayBuffer], { type: object.mimeType });
}
async function invokeImage(store: Store, base: URL, secret: string, channel: ProviderChannel, input: ProviderInput) {
const count = input.count === undefined ? 1 : input.count;
if (!Number.isSafeInteger(count) || count < 1) throw Object.assign(new Error("image count is invalid"), { errorCode: "PROVIDER_INVALID_REQUEST", retryable: false });
const rawReferenceImages = input.referenceImages || [];
const fallbackReferenceIds = [...new Set(input.references || [])];
if (rawReferenceImages.length > chineseImageNumbers.length || fallbackReferenceIds.length > chineseImageNumbers.length) throw Object.assign(new Error("too many reference images"), { errorCode: "PROVIDER_REFERENCE_LIMIT", retryable: false });
const orderedReferences = rawReferenceImages.map((item) => {
const objectId = typeof item.objectId === "string" ? item.objectId.trim() : "";
if (!objectId || !Number.isSafeInteger(item.seq) || item.seq < 1 || item.seq > chineseImageNumbers.length) throw Object.assign(new Error("reference image sequence is invalid"), { errorCode: "PROVIDER_REFERENCE_INVALID", retryable: false });
return { ...item, objectId, name: `图片${chineseImageNumbers[item.seq - 1]}` };
}).sort((a, b) => a.seq - b.seq);
if (new Set(orderedReferences.map((item) => item.seq)).size !== orderedReferences.length) throw Object.assign(new Error("reference image sequences must be unique"), { errorCode: "PROVIDER_REFERENCE_INVALID", retryable: false });
const seenReferenceIds = new Set<string>();
const uniqueReferences = orderedReferences.filter((item) => {
if (seenReferenceIds.has(item.objectId)) return false;
seenReferenceIds.add(item.objectId);
return true;
});
const referenceIds = uniqueReferences.length ? uniqueReferences.map((item) => item.objectId) : fallbackReferenceIds;
const effectiveReferences = uniqueReferences.length ? uniqueReferences : referenceIds.map((objectId, index) => ({ objectId, seq: index + 1, name: `图片${chineseImageNumbers[index]}` }));
const hasReferences = referenceIds.length > 0;
const prompt = input.prompt?.trim() || "";
const mentionedNames = mentionedImageNames(prompt);
const availableNames = new Set(effectiveReferences.map((item) => item.name));
for (const name of mentionedNames) if (!availableNames.has(name)) throw Object.assign(new Error("prompt references an unavailable image"), { errorCode: "PROVIDER_REFERENCE_UNAVAILABLE", retryable: false });
const rawParams = input.params || {};
const params = channel.providerType === "openai-images" ? normalizeImageParams(rawParams, hasReferences) : { ...rawParams };
// The channel owns the upstream model identity. Do not let a public task
// parameter override the configured resolution mapping.
delete params.model;
const channelModel = resolveChannelModel(channel, input.publicModelId);
const model = channel.providerType === "openai-images" ? imageModelFor(channelModel.model, rawParams.resolution, channelModel.exact ? undefined : channel.resolutionModelMap) : channelModel.model || "configured-model";
let rewrittenPrompt = [...effectiveReferences].sort((left, right) => right.name.length - left.name.length).reduce((value, reference) => value.split("@" + reference.name).join("第" + (chineseImageNumbers[reference.seq - 1] || reference.seq) + "张图"), prompt);
if (hasReferences && mentionedNames.size === 0) rewrittenPrompt = rewrittenPrompt ? `${rewrittenPrompt}\n基于以下参考图生成。` : "基于以下参考图生成。";
let response: Response;
if (hasReferences) {
const form = new FormData();
form.set("model", model);
form.set("prompt", rewrittenPrompt);
form.set("n", String(count));
for (const [index, objectId] of referenceIds.entries()) {
const blob = await referenceBlob(store, objectId);
const extension = blob.type.split("/")[1] || "png";
form.append("image[]", blob, "reference-" + index + "." + extension);
}
if (input.maskObjectId) form.append("mask", await referenceBlob(store, input.maskObjectId), "mask.png");
for (const [key, value] of Object.entries(params)) if (typeof value === "string" || typeof value === "number") form.set(key, String(value));
response = await providerFetch(providerEndpoint(base, "v1/images/edits"), secret, { method: "POST", headers: { "Idempotency-Key": input.platformIdempotencyKey || input.taskId }, body: form }, 60_000, input.onRequestAttempted);
} else {
response = await providerFetch(providerEndpoint(base, "v1/images/generations"), secret, { method: "POST", headers: { "content-type": "application/json", "Idempotency-Key": input.platformIdempotencyKey || input.taskId }, body: JSON.stringify({ model, prompt: rewrittenPrompt, n: count, ...params }) }, 60_000, input.onRequestAttempted);
}
const providerRequestIdHeader = providerRequestId(response);
const data = await jsonResponse(response);
const items = Array.isArray(data.data) ? data.data as Array<Record<string, unknown>> : [];
const outputs: ProviderOutput[] = [];
for (const item of items) {
if (typeof item.b64_json === "string") {
const decoded = decodeImageBase64(item.b64_json);
outputs.push({ mimeType: decoded.mimeType, data: item.b64_json, metadata: imageOutputMetadata(item, data, decoded.mimeType, "base64", decoded.bytes) });
} else if (typeof item.url === "string") {
const downloaded = await downloadOutput(item.url);
outputs.push({ ...downloaded, metadata: imageOutputMetadata(item, data, downloaded.mimeType, "url", Buffer.from(downloaded.data, "base64")) });
}
}
if (!outputs.length) throw Object.assign(new Error("provider returned no image"), { retryable: false, errorCode: "PROVIDER_EMPTY_RESULT" });
return { id: typeof data.id === "string" && data.id.trim() ? data.id : providerRequestIdHeader, outputs };
}
async function invokeText(base: URL, secret: string, channel: ProviderChannel, input: ProviderInput) { const model = resolveChannelModel(channel, input.publicModelId).model || "configured-model"; const response = await providerFetch(providerEndpoint(base, "v1/responses"), secret, { method: "POST", headers: { "content-type": "application/json", "Idempotency-Key": input.platformIdempotencyKey || input.taskId }, body: JSON.stringify({ ...input.params, model, input: input.prompt || "" }) }, 60_000, input.onRequestAttempted); const data = await jsonResponse(response); const nested = Array.isArray(data.output) ? (data.output as Array<{ content?: Array<{ text?: string }> }>).flatMap((item) => item.content || []).map((item) => item.text || "").join("") : ""; const text = typeof data.output_text === "string" ? data.output_text : nested; if (!text) throw Object.assign(new Error("provider returned no text"), { retryable: false, errorCode: "PROVIDER_EMPTY_RESULT" }); return { id: typeof data.id === "string" && data.id.trim() ? data.id : providerRequestId(response), outputs: [{ mimeType: "text/plain", data: Buffer.from(text).toString("base64") }] }; }
async function invokeAudio(base: URL, secret: string, channel: ProviderChannel, input: ProviderInput) { const model = resolveChannelModel(channel, input.publicModelId).model || "configured-model"; const response = await providerFetch(providerEndpoint(base, "v1/audio/speech"), secret, { method: "POST", headers: { "content-type": "application/json", "Idempotency-Key": input.platformIdempotencyKey || input.taskId }, body: JSON.stringify({ ...input.params, model, input: input.prompt || "" }) }, 60_000, input.onRequestAttempted); if (!response.ok) throw providerHttpError(response.status); const bytes = Buffer.from(await response.arrayBuffer()); if (!bytes.length || bytes.byteLength > 50 * 1024 * 1024) throw new Error("provider audio response invalid"); const mimeType = (response.headers.get("content-type") || "audio/mpeg").split(";")[0].toLowerCase(); if (!/^audio\/(mpeg|wav|ogg|mp4)$/.test(mimeType)) throw new Error("provider audio MIME is not allowed"); return { id: providerRequestId(response), outputs: [{ mimeType, data: bytes.toString("base64") }] }; }
export async function invokeProvider(store: Store, channel: ProviderChannel, input: ProviderInput): Promise<ProviderResult> {
if (!channel.baseUrl || !channel.secretRef) return { status: "failed", errorCode: channel.fixtureFailure === "unknown" ? "PROVIDER_TIMEOUT_UNKNOWN" : "CHANNEL_UNAVAILABLE", retryable: true };
try { const secret = decryptSecret(channel.secretRef, store.channelEncryptionKey); if (!secret) return { status: "failed", errorCode: "PROVIDER_SECRET_UNAVAILABLE", retryable: false }; const base = await validateResolvedUrl(channel.baseUrl); const result = input.taskType === "audio" ? await invokeAudio(base, secret, channel, input) : input.taskType === "text" || input.taskType === "reverse-prompt" ? await invokeText(base, secret, channel, input) : await invokeImage(store, base, secret, channel, input); return { status: "succeeded", providerRequestId: result.id, outputs: result.outputs }; }
catch (error) {
const typed = error as Error & { retryable?: boolean; errorCode?: string };
if (typed.name === "AbortError" || typed.name === "TimeoutError") return { status: "unknown", errorCode: "PROVIDER_TIMEOUT_UNKNOWN", retryable: false };
// Once a request may have reached the provider, a transport failure is
// not safe to retry because it could duplicate a billable generation.
if (!typed.errorCode) return { status: "unknown", errorCode: "PROVIDER_TRANSPORT_UNKNOWN", retryable: false };
return { status: "failed", errorCode: typed.errorCode, retryable: typed.retryable ?? true };
}
}
export async function probeProviderModels(baseUrl: string, secret: string) {
let requestAttempted = false;
try {
const base = await validateResolvedUrl(baseUrl);
const response = await providerFetch(providerEndpoint(base, "v1/models"), secret, { method: "GET" }, 15_000, () => { requestAttempted = true; });
const data = await jsonResponse(response);
const items = Array.isArray(data.data) ? data.data as Array<Record<string, unknown>> : [];
const models = items.map((item) => typeof item.id === "string" ? item.id : undefined).filter((value): value is string => Boolean(value));
if (!models.length) throw new Error("provider models response is empty");
return { healthy: true, models };
} catch (error) {
const typed = error instanceof Error ? error as Error & { requestAttempted?: boolean } : Object.assign(new Error(String(error)), {} as { requestAttempted?: boolean });
typed.requestAttempted = requestAttempted;
throw typed;
}
}
/** Query/cancel hooks used by the unknown-result reconciler. Providers that do
* not expose these endpoints simply remain in manual review. */
export async function queryProvider(store: Store, channel: ProviderChannel, providerRequestId: string): Promise<ProviderResult> {
if (channel.providerType === "openai-images") return { status: "unknown", providerRequestId, errorCode: "PROVIDER_QUERY_UNAVAILABLE" };
if (!channel.baseUrl || !channel.secretRef) return { status: "unknown", providerRequestId, errorCode: "PROVIDER_QUERY_UNAVAILABLE" };
try {
const secret = decryptSecret(channel.secretRef, store.channelEncryptionKey); if (!secret) return { status: "unknown", providerRequestId, errorCode: "PROVIDER_SECRET_UNAVAILABLE" };
const response = await providerFetch(providerEndpoint(await validateResolvedUrl(channel.baseUrl), `v1/requests/${encodeURIComponent(providerRequestId)}`), secret, { method: "GET" }, 15_000);
const data = await jsonResponse(response); const status = typeof data.status === "string" ? data.status.toLowerCase() : "";
if (["succeeded", "completed", "success"].includes(status)) {
const items = Array.isArray(data.data) ? data.data as Array<Record<string, unknown>> : [];
const outputs: ProviderOutput[] = [];
for (const item of items) {
if (typeof item.b64_json === "string") {
const decoded = decodeImageBase64(item.b64_json);
outputs.push({ mimeType: decoded.mimeType, data: item.b64_json, metadata: imageOutputMetadata(item, data, decoded.mimeType, "base64", decoded.bytes) });
} else if (typeof item.url === "string") {
const downloaded = await downloadOutput(item.url);
outputs.push({ ...downloaded, metadata: imageOutputMetadata(item, data, downloaded.mimeType, "url", Buffer.from(downloaded.data, "base64")) });
}
}
const text = typeof data.output_text === "string" ? data.output_text : typeof data.text === "string" ? data.text : undefined; if (text) outputs.push({ mimeType: "text/plain", data: Buffer.from(text).toString("base64") });
if (typeof data.audio === "string") outputs.push({ mimeType: typeof data.mime_type === "string" && /^audio\//.test(data.mime_type) ? data.mime_type : "audio/mpeg", data: data.audio });
return { status: outputs.length ? "succeeded" : "unknown", providerRequestId, outputs, errorCode: outputs.length ? undefined : "PROVIDER_QUERY_EMPTY" };
}
if (["failed", "error", "canceled", "cancelled"].includes(status)) return { status: "failed", providerRequestId, errorCode: status === "canceled" || status === "cancelled" ? "PROVIDER_CANCELED" : "PROVIDER_FAILED", retryable: false };
return { status: "unknown", providerRequestId, errorCode: "PROVIDER_STILL_RUNNING" };
} catch { return { status: "unknown", providerRequestId, errorCode: "PROVIDER_QUERY_FAILED" }; }
}
export async function cancelProvider(store: Store, channel: ProviderChannel, providerRequestId: string): Promise<ProviderCancelResult> {
if (channel.providerType === "openai-images") return "unsupported";
if (!channel.baseUrl || !channel.secretRef) return "unsupported";
try {
const secret = decryptSecret(channel.secretRef, store.channelEncryptionKey); if (!secret) return "unsupported";
const response = await providerFetch(providerEndpoint(await validateResolvedUrl(channel.baseUrl), `v1/requests/${encodeURIComponent(providerRequestId)}/cancel`), secret, { method: "POST" }, 15_000);
if (response.status === 404 || response.status === 405 || response.status === 501) return "unsupported";
if (!response.ok) return "pending";
let status = ""; try { const data = JSON.parse(await response.text()) as { status?: unknown }; status = typeof data.status === "string" ? data.status.toLowerCase() : ""; } catch { return "pending"; }
return status === "canceled" || status === "cancelled" ? "confirmed" : "pending";
} catch { return "pending"; }
}
export function providerSecretIsReference(secret: string) { return Boolean(secret.trim()); }
+158
View File
@@ -0,0 +1,158 @@
import { createHash, randomUUID } from "node:crypto";
import { decryptSecret } from "../shared/auth.ts";
import type { Store } from "../store.ts";
import type { WebDavSyncMutation } from "../infra/repository.ts";
import { resolvePublicHttpsUrl, safeOutboundFetch } from "../infra/outbound-url.ts";
type WebDavConnection = { url: string; username?: string; password?: string; directory: string };
function connectionFromStore(store: Store, userId: string, secret: string): WebDavConnection | undefined {
const record = store.webdav.get(userId); if (!record) return undefined;
return { url: decryptSecret(record.encryptedUrl, secret) || "", username: decryptSecret(record.encryptedUsername, secret), password: decryptSecret(record.encryptedPassword, secret), directory: record.directory };
}
export async function validateWebDavUrl(value: string) {
try { return (await resolvePublicHttpsUrl(value)).url; } catch (error) { throw new Error(error instanceof Error ? error.message.replace("出站地址", "WebDAV 地址") : "WebDAV 地址无效"); }
}
function cleanPath(value: string) { const segments = value.replaceAll("\\", "/").split("/").filter(Boolean); if (!segments.length || segments.some((part) => part === "." || part === ".." || part.includes("\0"))) throw new Error("WebDAV 路径无效"); return segments.map(encodeURIComponent).join("/"); }
export function validateWebDavPath(value: string) { return cleanPath(value); }
async function request(connection: WebDavConnection, path: string, init: RequestInit, maxBytes = 20 * 1024 * 1024) {
const base = await validateWebDavUrl(connection.url); const root = cleanPath(connection.directory); const relative = path ? cleanPath(path) : ""; base.pathname = `${base.pathname.replace(/\/+$/, "")}/${root}${relative ? `/${relative}` : ""}`;
const headers = new Headers(init.headers); if (connection.username || connection.password) headers.set("authorization", `Basic ${Buffer.from(`${connection.username || ""}:${connection.password || ""}`).toString("base64")}`);
return safeOutboundFetch(base, { ...init, headers }, { timeoutMs: 15_000, maxBytes });
}
export async function testWebDav(connection: WebDavConnection) { const response = await request(connection, "", { method: "PROPFIND", headers: { depth: "0" } }, 1024 * 1024); if (response.status !== 207 && !response.headers.get("dav")) throw new Error(`WebDAV 能力验证失败 (${response.status})`); const probe = `.miragenflow-probe-${randomUUID()}.json`; const written = await putWebDavFile(connection, probe, Buffer.from("{}"), "application/json", undefined, true); await deleteWebDavFile(connection, probe, written.etag); }
export async function getWebDavFile(connection: WebDavConnection, path: string) { const response = await request(connection, path, { method: "GET" }); if (!response.ok) throw new Error(`WebDAV 文件读取失败 (${response.status})`); return { mimeType: response.headers.get("content-type") || "application/octet-stream", etag: response.headers.get("etag") || undefined, data: Buffer.from(await response.arrayBuffer()) }; }
export async function putWebDavFile(connection: WebDavConnection, path: string, data: Uint8Array, mimeType: string, ifMatch?: string, ifNoneMatch = false) { const headers: Record<string, string> = { "content-type": mimeType }; if (ifMatch) headers["if-match"] = ifMatch; else if (ifNoneMatch) headers["if-none-match"] = "*"; const response = await request(connection, path, { method: "PUT", headers, body: data as BodyInit }); if (!response.ok) throw new Error(`WebDAV 文件写入失败 (${response.status})`); return { etag: response.headers.get("etag") || undefined }; }
export async function deleteWebDavFile(connection: WebDavConnection, path: string, ifMatch?: string) { const headers: Record<string, string> = {}; if (ifMatch) headers["if-match"] = ifMatch; const response = await request(connection, path, { method: "DELETE", headers }); if (!response.ok && response.status !== 404) throw new Error(`WebDAV 文件删除失败 (${response.status})`); }
/** Queue remote files for retention cleanup. Only files tracked by our own
* manifest proxy are eligible; unknown remote files are never guessed or
* deleted without a listing/manifest reference. */
export function enqueueExpiredWebDavRetention(store: Store, now = Date.now()) {
let queued = 0; let changed = false;
for (const record of store.webdav.values()) {
if (!record.configured || record.retentionState === "deleted" || !record.manifestRetentionExpiresAt || Date.parse(record.manifestRetentionExpiresAt) > now) continue;
const files = [...store.webdavFiles.values()].filter((file) => file.userId === record.userId && !file.deletedAt);
for (const file of files) {
// Terminal jobs do not suppress a later scan: a succeeded job may have
// lost its file metadata before the commit, while failed/conflict jobs
// must be retried. Only an active queued/running job is a dedupe hit.
const existing = [...store.webdavJobs.values()].find((job) => job.userId === record.userId && job.intent === "retention-delete" && job.path === file.path && ["queued", "running"].includes(job.status));
if (existing) continue;
const id = randomUUID();
store.webdavJobs.set(id, { id, userId: record.userId, operation: "delete", intent: "retention-delete", path: file.path, ifMatch: file.etag, attempts: 0, status: "queued", nextAttemptAt: new Date(now).toISOString(), createdAt: new Date(now).toISOString() });
queued += 1; changed = true;
}
const nextRetentionState = files.length ? "deleting" : "deleted";
if (record.retentionState !== nextRetentionState) { record.retentionState = nextRetentionState; changed = true; }
if (files.length && record.state !== "syncing") { record.state = "syncing"; changed = true; }
}
if (changed) store.persist();
return queued;
}
/** PostgreSQL retention scanning uses a row-locked INSERT/UPDATE transaction
* instead of the local full-snapshot writer. Memory/file fixtures retain the
* synchronous helper above for compatibility with existing callers/tests. */
export async function enqueueExpiredWebDavRetentionAsync(store: Store, now = Date.now()) {
if (store.repository.adapter === "postgres" && store.repository.enqueueExpiredWebDavRetention) {
return store.repository.enqueueExpiredWebDavRetention(new Date(now).toISOString());
}
return enqueueExpiredWebDavRetention(store, now);
}
export async function processWebDavJobs(store: Store, secret: string) {
let processed = 0;
const now = Date.now();
const pg = store.repository.adapter === "postgres" && !!store.repository.claimWebDavSyncJobs && !!store.repository.completeWebDavSyncJob;
const renew = pg ? store.repository.renewWebDavSyncJob : undefined;
const leaseMs = 30_000;
const worker = `webdav-${process.pid}-${randomUUID().slice(0, 8)}`;
const jobs = pg ? await store.repository.claimWebDavSyncJobs!(worker, leaseMs) : [...store.webdavJobs.values()];
if (pg) for (const job of jobs) store.webdavJobs.set(job.id, job);
for (const job of jobs) {
if (!pg && job.status === "running" && (!job.leaseExpiresAt || Date.parse(job.leaseExpiresAt) <= now)) {
job.status = "queued";
job.nextAttemptAt = new Date(now).toISOString();
job.lastError = "上一次同步进程未完成,已恢复重试";
}
if (!pg && (job.status !== "queued" || Date.parse(job.nextAttemptAt) > now)) continue;
// A user may extend manifest retention after expiry has queued deletion
// jobs. Re-check the current expiry immediately before dispatch so a stale
// queue item cannot delete a file whose retention was just extended.
if (job.intent === "retention-delete") {
const record = store.webdav.get(job.userId);
if (record?.manifestRetentionExpiresAt && Date.parse(record.manifestRetentionExpiresAt) > now) {
if (record.retentionState === "deleting") record.retentionState = "active";
if (record.state === "syncing") record.state = "ready";
if (pg) { job.status = "succeeded"; job.leaseOwner = undefined; job.leaseExpiresAt = undefined; const completed = await store.repository.completeWebDavSyncJob!(job.id, worker, job, { config: record }); if (!completed) { await store.loadPersisted().catch(() => undefined); } }
else { store.webdavJobs.delete(job.id); store.persist(); }
continue;
}
}
const connection = connectionFromStore(store, job.userId, secret); if (!connection) { job.status = "failed"; job.lastError = "WebDAV 未配置"; job.leaseOwner = undefined; job.leaseExpiresAt = undefined; if (pg) { const completed = await store.repository.completeWebDavSyncJob!(job.id, worker, job); if (!completed) { await store.loadPersisted().catch(() => undefined); } } else store.persist(); continue; }
if (!pg) {
job.status = "running";
job.leaseExpiresAt = new Date(now + leaseMs).toISOString();
store.persist();
}
let leaseLost = false;
const renewTimer = renew ? setInterval(() => {
void renew.call(store.repository, job.id, worker, leaseMs).then((result) => {
if (!result.renewed) leaseLost = true;
else if (result.leaseExpiresAt) job.leaseExpiresAt = result.leaseExpiresAt;
}).catch(() => { leaseLost = true; });
}, Math.max(1_000, Math.floor(leaseMs / 3))) : undefined;
renewTimer?.unref?.();
let mutation: WebDavSyncMutation | undefined;
try {
if (job.operation === "put") {
const bytes = Buffer.from(job.data || "", "base64"); const result = await putWebDavFile(connection, job.path, bytes, job.mimeType || "application/octet-stream", job.ifMatch, !job.ifMatch);
const checksum = createHash("sha256").update(bytes).digest("hex"); const key = `${job.userId}:${job.path}`; const current = store.webdavFiles.get(key); const file = { userId: job.userId, path: job.path, mimeType: job.mimeType || "application/octet-stream", data: job.data || "", checksum, etag: result.etag || `\"${checksum}\"`, version: (current?.version || 0) + 1, updatedAt: new Date().toISOString(), syncState: "synced" as const }; store.webdavFiles.set(key, file); mutation = { file };
if (job.path.endsWith("manifest.json")) { const record = store.webdav.get(job.userId); if (record) { record.manifestEtag = result.etag || `\"${checksum}\"`; record.manifestChecksum = checksum; record.manifestVersion = (record.manifestVersion || 0) + 1; record.manifestRetentionExpiresAt ||= new Date(Date.now() + record.retentionDays * 86_400_000).toISOString(); record.manifestExtensionDays ||= 0; } }
} else { await deleteWebDavFile(connection, job.path, job.ifMatch); const current = store.webdavFiles.get(`${job.userId}:${job.path}`); if (current) { current.deletedAt = new Date().toISOString(); current.version += 1; current.etag = `\"deleted-${current.version}\"`; current.data = ""; current.syncState = "synced"; mutation = { file: current }; } }
job.status = "succeeded"; job.lastError = undefined; job.leaseOwner = undefined; job.leaseExpiresAt = undefined;
if (job.intent === "archive-put") {
const taskId = job.path.split("/")[2]; const task = store.tasks.get(taskId); const related = [...store.webdavJobs.values()].filter((item) => item.userId === job.userId && item.intent === "archive-put" && item.path.startsWith(`archive/tasks/${taskId}/`));
if (task && related.length && related.every((item) => item.status === "succeeded")) { task.retentionState = "archived"; mutation = { ...(mutation || {}), taskRetention: { taskId, retentionState: "archived", updatedAt: new Date().toISOString() } }; }
}
if (job.intent === "retention-delete") {
const record = store.webdav.get(job.userId);
const pending = [...store.webdavJobs.values()].some((item) => item.userId === job.userId && item.intent === "retention-delete" && item.status !== "succeeded");
const liveFiles = [...store.webdavFiles.values()].some((item) => item.userId === job.userId && !item.deletedAt);
if (record && !pending && !liveFiles) { record.retentionState = "deleted"; record.state = "ready"; }
}
const record = store.webdav.get(job.userId); if (record) { const pending = [...store.webdavJobs.values()].some((item) => item.userId === job.userId && item.id !== job.id && item.status !== "succeeded"); record.state = pending ? "error" : "ready"; record.lastSyncedAt = new Date().toISOString(); mutation = { ...(mutation || {}), config: record }; }
processed += 1;
} catch (error) {
if (!pg) job.attempts += 1;
const message = error instanceof Error ? error.message : String(error);
const conflict = /409|etag|冲突|precondition/i.test(message);
job.status = conflict ? "conflict" : job.attempts >= 8 ? "failed" : "queued";
if (job.intent === "archive-put" && (job.status === "failed" || job.status === "conflict")) { const taskId = job.path.split("/")[2]; const task = store.tasks.get(taskId); if (task) { task.retentionState = "archive_failed"; mutation = { ...(mutation || {}), taskRetention: { taskId, retentionState: "archive_failed", updatedAt: new Date().toISOString() } }; } }
if (job.intent === "retention-delete" && (job.status === "failed" || job.status === "conflict")) { const record = store.webdav.get(job.userId); if (record) { record.retentionState = "error"; record.state = "error"; mutation = { ...(mutation || {}), config: record }; } }
if (conflict && job.operation === "put" && !job.path.startsWith(".conflicts/")) {
const fileName = job.path.split("/").pop() || "manifest.json";
const conflictCopyPath = `.conflicts/${randomUUID()}-${fileName}`;
job.conflictCopyPath = conflictCopyPath;
const copyJobId = randomUUID();
const conflictJob = { id: copyJobId, userId: job.userId, operation: "put" as const, intent: "conflict-copy" as const, path: conflictCopyPath, data: job.data, mimeType: job.mimeType, attempts: 0, status: "queued" as const, nextAttemptAt: new Date().toISOString(), lastError: `冲突副本:${message}`, createdAt: new Date().toISOString() }; store.webdavJobs.set(copyJobId, conflictJob); mutation = { ...(mutation || {}), conflictJob };
}
job.leaseOwner = undefined;
job.leaseExpiresAt = undefined;
job.lastError = message;
job.nextAttemptAt = new Date(Date.now() + Math.min(60 * 60_000, 2 ** job.attempts * 1000)).toISOString();
}
if (pg) {
const completed = await store.repository.completeWebDavSyncJob!(job.id, worker, job, mutation);
if (!completed) { await store.loadPersisted().catch(() => undefined); continue; }
} else store.persist();
if (leaseLost && pg) await store.loadPersisted().catch(() => undefined);
if (renewTimer) clearInterval(renewTimer);
}
return processed;
}