Compare commits

...
Author SHA1 Message Date
suyx 7de633d4c9 fix(POSTV1-05): 修复生成请求与 Worker 重入
Dada P0-A isolated Windows CI / validate-and-package (push) Waiting to run
2026-08-05 16:04:11 +08:00
5 changed files with 80 additions and 4 deletions
@@ -0,0 +1,35 @@
export interface GenerationPollingProcessor {
processNext(): Promise<unknown>;
}
export class GenerationPollingLoop {
private closed = false;
private inFlight = false;
private readonly timer: ReturnType<typeof setInterval>;
constructor(
private readonly processor: GenerationPollingProcessor,
intervalMilliseconds = 250,
) {
if (!Number.isSafeInteger(intervalMilliseconds) || intervalMilliseconds <= 0) {
throw new Error("generation_polling_interval_invalid");
}
this.timer = setInterval(() => this.run(), intervalMilliseconds);
}
close() {
if (this.closed) return;
this.closed = true;
clearInterval(this.timer);
}
private run() {
if (this.closed || this.inFlight) return;
this.inFlight = true;
void this.processor.processNext()
.catch(() => undefined)
.finally(() => {
this.inFlight = false;
});
}
}
+5 -1
View File
@@ -16,6 +16,7 @@ const gptImageEndpoint = "https://oneapi.intelligrow.cn/v1/images/generations";
const gptImageReferenceEndpoint = "https://oneapi.intelligrow.cn/v1/images/edits"; const gptImageReferenceEndpoint = "https://oneapi.intelligrow.cn/v1/images/edits";
const maximumResponseBytes = 32 * 1024 * 1024; const maximumResponseBytes = 32 * 1024 * 1024;
const requestTimeoutMilliseconds = 180_000; const requestTimeoutMilliseconds = 180_000;
const geminiImageSystemInstruction = "Generate exactly one image from the user's description. Return the generated image and do not answer with text only.";
type FetchLike = typeof fetch; type FetchLike = typeof fetch;
@@ -141,7 +142,10 @@ function buildRequest(request: GenerationAdapterRequest) {
return { return {
body: JSON.stringify({ body: JSON.stringify({
extra_body: { google: { image_config: { aspect_ratio: request.ratio, image_size: "1K" } } }, extra_body: { google: { image_config: { aspect_ratio: request.ratio, image_size: "1K" } } },
messages: [{ content, role: "user" }], messages: [
{ content: geminiImageSystemInstruction, role: "system" },
{ content, role: "user" },
],
model: geminiProviderModelId, model: geminiProviderModelId,
stream: false, stream: false,
}), }),
+4 -3
View File
@@ -3,6 +3,7 @@ import { join } from "node:path";
import { WorkerAiCallGate } from "./ai-call-gate.js"; import { WorkerAiCallGate } from "./ai-call-gate.js";
import { runAiRuntimeProbe } from "./ai-runtime-probe.js"; import { runAiRuntimeProbe } from "./ai-runtime-probe.js";
import { GenerationPollingLoop } from "./generation-polling-loop.js";
import { GenerationProcessor } from "./generation-processor.js"; import { GenerationProcessor } from "./generation-processor.js";
import { OneApiGenerationAdapter } from "./oneapi-generation-adapter.js"; import { OneApiGenerationAdapter } from "./oneapi-generation-adapter.js";
import { readConfiguredLocalDataRoot } from "./runtime-config.js"; import { readConfiguredLocalDataRoot } from "./runtime-config.js";
@@ -52,13 +53,13 @@ if (!workerPort && process.argv.includes("--dada-ai-probe")) {
let projectCleanup: ProjectPurgeCleanup | undefined; let projectCleanup: ProjectPurgeCleanup | undefined;
let retentionTimer: ReturnType<typeof setInterval> | undefined; let retentionTimer: ReturnType<typeof setInterval> | undefined;
let processor: GenerationProcessor | undefined; let processor: GenerationProcessor | undefined;
let generationTimer: ReturnType<typeof setInterval> | undefined; let generationLoop: GenerationPollingLoop | undefined;
const control = attachWorkerSupervisorControl(controlPipe, () => { const control = attachWorkerSupervisorControl(controlPipe, () => {
clearInterval(keepAlive); clearInterval(keepAlive);
if (retentionTimer) clearInterval(retentionTimer); if (retentionTimer) clearInterval(retentionTimer);
retention?.close(); retention?.close();
projectCleanup?.close(); projectCleanup?.close();
if (generationTimer) clearInterval(generationTimer); generationLoop?.close();
processor?.close(); processor?.close();
storage?.close(); storage?.close();
}); });
@@ -102,7 +103,7 @@ if (!workerPort && process.argv.includes("--dada-ai-probe")) {
}); });
logger.write({ error_category: "none", status_category: "ready" }); logger.write({ error_category: "none", status_category: "ready" });
new WorkerAiCallGate({ getStorageStatus: () => storageStatus === "unavailable" ? storageStatus : (storage?.getStatus() ?? "unavailable"), logger }); new WorkerAiCallGate({ getStorageStatus: () => storageStatus === "unavailable" ? storageStatus : (storage?.getStatus() ?? "unavailable"), logger });
generationTimer = setInterval(() => { void processor?.processNext().catch(() => undefined); }, 250); generationLoop = new GenerationPollingLoop(processor);
} catch { } catch {
storageStatus = "unavailable"; storageStatus = "unavailable";
control.reportStatus("storage_unavailable"); control.reportStatus("storage_unavailable");
@@ -0,0 +1,28 @@
import { afterEach, describe, expect, it, vi } from "vitest";
import { GenerationPollingLoop } from "../../apps/worker/src/generation-polling-loop.js";
describe("POSTV1-05 generation polling loop", () => {
afterEach(() => {
vi.useRealTimers();
});
it("does not start another processor call while the current call is unresolved", async () => {
vi.useFakeTimers();
let finishCurrentCall: (() => void) | undefined;
const processNext = vi.fn(() => new Promise<void>((resolve) => {
finishCurrentCall = resolve;
}));
const loop = new GenerationPollingLoop({ processNext }, 250);
await vi.advanceTimersByTimeAsync(1_000);
expect(processNext).toHaveBeenCalledTimes(1);
finishCurrentCall?.();
await Promise.resolve();
await vi.advanceTimersByTimeAsync(250);
expect(processNext).toHaveBeenCalledTimes(2);
loop.close();
});
});
@@ -26,6 +26,14 @@ describe("POSTV1-02 OneAPI runtime adapter", () => {
const headers = new Headers(init?.headers); const headers = new Headers(init?.headers);
expect(headers.get("authorization")).toBe("Bearer synthetic-runtime-token"); expect(headers.get("authorization")).toBe("Bearer synthetic-runtime-token");
expect(init?.redirect).toBe("error"); expect(init?.redirect).toBe("error");
const payload = JSON.parse(String(init?.body)) as { messages: Array<{ content: unknown; role: string }> };
expect(payload.messages).toEqual([
{
content: "Generate exactly one image from the user's description. Return the generated image and do not answer with text only.",
role: "system",
},
{ content: "一张用于本机验收的抽象色彩图", role: "user" },
]);
return new Response(JSON.stringify({ return new Response(JSON.stringify({
choices: [{ message: { content: `![result](data:image/png;base64,${source.toString("base64")})` } }], choices: [{ message: { content: `![result](data:image/png;base64,${source.toString("base64")})` } }],
}), { headers: { "content-type": "application/json" }, status: 200 }); }), { headers: { "content-type": "application/json" }, status: 200 });