fix(POSTV1-05): 修复生成请求与 Worker 重入
Dada P0-A isolated Windows CI / validate-and-package (push) Waiting to run
Dada P0-A isolated Windows CI / validate-and-package (push) Waiting to run
This commit is contained in:
@@ -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;
|
||||
});
|
||||
}
|
||||
}
|
||||
@@ -16,6 +16,7 @@ const gptImageEndpoint = "https://oneapi.intelligrow.cn/v1/images/generations";
|
||||
const gptImageReferenceEndpoint = "https://oneapi.intelligrow.cn/v1/images/edits";
|
||||
const maximumResponseBytes = 32 * 1024 * 1024;
|
||||
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;
|
||||
|
||||
@@ -141,7 +142,10 @@ function buildRequest(request: GenerationAdapterRequest) {
|
||||
return {
|
||||
body: JSON.stringify({
|
||||
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,
|
||||
stream: false,
|
||||
}),
|
||||
|
||||
@@ -3,6 +3,7 @@ import { join } from "node:path";
|
||||
|
||||
import { WorkerAiCallGate } from "./ai-call-gate.js";
|
||||
import { runAiRuntimeProbe } from "./ai-runtime-probe.js";
|
||||
import { GenerationPollingLoop } from "./generation-polling-loop.js";
|
||||
import { GenerationProcessor } from "./generation-processor.js";
|
||||
import { OneApiGenerationAdapter } from "./oneapi-generation-adapter.js";
|
||||
import { readConfiguredLocalDataRoot } from "./runtime-config.js";
|
||||
@@ -52,13 +53,13 @@ if (!workerPort && process.argv.includes("--dada-ai-probe")) {
|
||||
let projectCleanup: ProjectPurgeCleanup | undefined;
|
||||
let retentionTimer: ReturnType<typeof setInterval> | undefined;
|
||||
let processor: GenerationProcessor | undefined;
|
||||
let generationTimer: ReturnType<typeof setInterval> | undefined;
|
||||
let generationLoop: GenerationPollingLoop | undefined;
|
||||
const control = attachWorkerSupervisorControl(controlPipe, () => {
|
||||
clearInterval(keepAlive);
|
||||
if (retentionTimer) clearInterval(retentionTimer);
|
||||
retention?.close();
|
||||
projectCleanup?.close();
|
||||
if (generationTimer) clearInterval(generationTimer);
|
||||
generationLoop?.close();
|
||||
processor?.close();
|
||||
storage?.close();
|
||||
});
|
||||
@@ -102,7 +103,7 @@ if (!workerPort && process.argv.includes("--dada-ai-probe")) {
|
||||
});
|
||||
logger.write({ error_category: "none", status_category: "ready" });
|
||||
new WorkerAiCallGate({ getStorageStatus: () => storageStatus === "unavailable" ? storageStatus : (storage?.getStatus() ?? "unavailable"), logger });
|
||||
generationTimer = setInterval(() => { void processor?.processNext().catch(() => undefined); }, 250);
|
||||
generationLoop = new GenerationPollingLoop(processor);
|
||||
} catch {
|
||||
storageStatus = "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);
|
||||
expect(headers.get("authorization")).toBe("Bearer synthetic-runtime-token");
|
||||
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({
|
||||
choices: [{ message: { content: `})` } }],
|
||||
}), { headers: { "content-type": "application/json" }, status: 200 });
|
||||
|
||||
Reference in New Issue
Block a user