From 7de633d4c98fb4f3bd4e0b48ed86721ab607499e Mon Sep 17 00:00:00 2001 From: suyx Date: Wed, 5 Aug 2026 16:04:11 +0800 Subject: [PATCH] =?UTF-8?q?fix(POSTV1-05):=20=E4=BF=AE=E5=A4=8D=E7=94=9F?= =?UTF-8?q?=E6=88=90=E8=AF=B7=E6=B1=82=E4=B8=8E=20Worker=20=E9=87=8D?= =?UTF-8?q?=E5=85=A5?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- apps/worker/src/generation-polling-loop.ts | 35 +++++++++++++++++++ apps/worker/src/oneapi-generation-adapter.ts | 6 +++- apps/worker/src/worker.ts | 7 ++-- .../postv1-generation-polling-loop.test.ts | 28 +++++++++++++++ tests/worker/postv1-oneapi-runtime.test.ts | 8 +++++ 5 files changed, 80 insertions(+), 4 deletions(-) create mode 100644 apps/worker/src/generation-polling-loop.ts create mode 100644 tests/worker/postv1-generation-polling-loop.test.ts diff --git a/apps/worker/src/generation-polling-loop.ts b/apps/worker/src/generation-polling-loop.ts new file mode 100644 index 0000000..5f8be9e --- /dev/null +++ b/apps/worker/src/generation-polling-loop.ts @@ -0,0 +1,35 @@ +export interface GenerationPollingProcessor { + processNext(): Promise; +} + +export class GenerationPollingLoop { + private closed = false; + private inFlight = false; + private readonly timer: ReturnType; + + 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; + }); + } +} diff --git a/apps/worker/src/oneapi-generation-adapter.ts b/apps/worker/src/oneapi-generation-adapter.ts index ce97230..5b3788c 100644 --- a/apps/worker/src/oneapi-generation-adapter.ts +++ b/apps/worker/src/oneapi-generation-adapter.ts @@ -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, }), diff --git a/apps/worker/src/worker.ts b/apps/worker/src/worker.ts index bb03306..cad8c9c 100644 --- a/apps/worker/src/worker.ts +++ b/apps/worker/src/worker.ts @@ -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 | undefined; let processor: GenerationProcessor | undefined; - let generationTimer: ReturnType | 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"); diff --git a/tests/worker/postv1-generation-polling-loop.test.ts b/tests/worker/postv1-generation-polling-loop.test.ts new file mode 100644 index 0000000..c717a18 --- /dev/null +++ b/tests/worker/postv1-generation-polling-loop.test.ts @@ -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((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(); + }); +}); diff --git a/tests/worker/postv1-oneapi-runtime.test.ts b/tests/worker/postv1-oneapi-runtime.test.ts index 8feed89..02db977 100644 --- a/tests/worker/postv1-oneapi-runtime.test.ts +++ b/tests/worker/postv1-oneapi-runtime.test.ts @@ -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: `![result](data:image/png;base64,${source.toString("base64")})` } }], }), { headers: { "content-type": "application/json" }, status: 200 });