Compare commits
14
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
f7cac5dabf | ||
|
|
a7f62adad4 | ||
|
|
b1f143c238 | ||
|
|
99fd3b1802 | ||
|
|
fd0003a804 | ||
|
|
bfdfe44f87 | ||
|
|
cbb7f658a3 | ||
|
|
43d946bb5c | ||
|
|
79b01ebc81 | ||
|
|
90f812fae5 | ||
|
|
1155a81c3b | ||
|
|
3dca4ad77c | ||
|
|
99fe07a761 | ||
|
|
443e8b94f0 |
+42
-4
@@ -1,6 +1,6 @@
|
|||||||
import { randomBytes, randomUUID } from "node:crypto";
|
import { randomBytes, randomUUID } from "node:crypto";
|
||||||
import { createReadStream, readFileSync } from "node:fs";
|
import { createReadStream, existsSync, readFileSync } from "node:fs";
|
||||||
import { resolve } from "node:path";
|
import { extname, resolve } from "node:path";
|
||||||
|
|
||||||
import {
|
import {
|
||||||
AccountDeletionCompleteRequestSchema,
|
AccountDeletionCompleteRequestSchema,
|
||||||
@@ -218,6 +218,24 @@ import type { ManagedStorage } from "./managed-storage.js";
|
|||||||
import { PrivateContentError, PrivateContentService } from "./private-content.js";
|
import { PrivateContentError, PrivateContentService } from "./private-content.js";
|
||||||
import { assertSafeAdminDiagnostics, assertSafeAdminServicesStorage } from "./admin-state.js";
|
import { assertSafeAdminDiagnostics, assertSafeAdminServicesStorage } from "./admin-state.js";
|
||||||
|
|
||||||
|
const productAssetContentTypes: Readonly<Record<string, string>> = {
|
||||||
|
".css": "text/css; charset=utf-8",
|
||||||
|
".jpeg": "image/jpeg",
|
||||||
|
".jpg": "image/jpeg",
|
||||||
|
".js": "text/javascript; charset=utf-8",
|
||||||
|
".json": "application/json; charset=utf-8",
|
||||||
|
".mjs": "text/javascript; charset=utf-8",
|
||||||
|
".png": "image/png",
|
||||||
|
".svg": "image/svg+xml",
|
||||||
|
".webp": "image/webp",
|
||||||
|
".woff": "font/woff",
|
||||||
|
".woff2": "font/woff2",
|
||||||
|
};
|
||||||
|
|
||||||
|
function productAssetContentType(path: string) {
|
||||||
|
return productAssetContentTypes[extname(path).toLowerCase()] ?? "application/octet-stream";
|
||||||
|
}
|
||||||
|
|
||||||
const defaultBootstrap: BootstrapResponse = {
|
const defaultBootstrap: BootstrapResponse = {
|
||||||
app_version: "0.0.0",
|
app_version: "0.0.0",
|
||||||
dependencies: [],
|
dependencies: [],
|
||||||
@@ -238,6 +256,7 @@ export interface CreateAppOptions {
|
|||||||
assetReleases?: AssetReleaseReader;
|
assetReleases?: AssetReleaseReader;
|
||||||
bootstrap?: () => BootstrapResponse | Promise<BootstrapResponse>;
|
bootstrap?: () => BootstrapResponse | Promise<BootstrapResponse>;
|
||||||
browserGate?: boolean;
|
browserGate?: boolean;
|
||||||
|
productIndexHtml?: string;
|
||||||
browserSupportRelease?: BrowserSupportRelease;
|
browserSupportRelease?: BrowserSupportRelease;
|
||||||
browserSupportSecret?: Buffer;
|
browserSupportSecret?: Buffer;
|
||||||
credits?: CreditService;
|
credits?: CreditService;
|
||||||
@@ -272,6 +291,10 @@ const supportGateDirectory = resolve(process.env.DADA_SUPPORT_GATE_ROOT ?? "apps
|
|||||||
const supportGateHtml = readFileSync(resolve(supportGateDirectory, "index.html"), "utf8");
|
const supportGateHtml = readFileSync(resolve(supportGateDirectory, "index.html"), "utf8");
|
||||||
const supportGateCss = readFileSync(resolve(supportGateDirectory, "support-gate.css"), "utf8");
|
const supportGateCss = readFileSync(resolve(supportGateDirectory, "support-gate.css"), "utf8");
|
||||||
const supportGateJavaScript = readFileSync(resolve(supportGateDirectory, "support-gate.js"), "utf8");
|
const supportGateJavaScript = readFileSync(resolve(supportGateDirectory, "support-gate.js"), "utf8");
|
||||||
|
const productWebRoot = resolve(process.env.DADA_WEB_ROOT ?? "apps/web/dist");
|
||||||
|
const packagedProductIndexHtml = existsSync(resolve(productWebRoot, "index.html"))
|
||||||
|
? readFileSync(resolve(productWebRoot, "index.html"), "utf8")
|
||||||
|
: undefined;
|
||||||
const clientHints = "Sec-CH-UA, Sec-CH-UA-Full-Version-List, Sec-CH-UA-Platform";
|
const clientHints = "Sec-CH-UA, Sec-CH-UA-Full-Version-List, Sec-CH-UA-Platform";
|
||||||
const contentSecurityPolicy = [
|
const contentSecurityPolicy = [
|
||||||
"default-src 'self'",
|
"default-src 'self'",
|
||||||
@@ -710,6 +733,7 @@ export async function createApp(options: CreateAppOptions = {}) {
|
|||||||
)
|
)
|
||||||
: undefined);
|
: undefined);
|
||||||
const browserGate = options.browserGate ?? true;
|
const browserGate = options.browserGate ?? true;
|
||||||
|
const productIndexHtml = options.productIndexHtml ?? packagedProductIndexHtml;
|
||||||
const browserSupportSecret = options.browserSupportSecret ?? randomBytes(32);
|
const browserSupportSecret = options.browserSupportSecret ?? randomBytes(32);
|
||||||
const browserSupportRelease = options.browserSupportRelease;
|
const browserSupportRelease = options.browserSupportRelease;
|
||||||
const app = Fastify({
|
const app = Fastify({
|
||||||
@@ -901,11 +925,25 @@ export async function createApp(options: CreateAppOptions = {}) {
|
|||||||
});
|
});
|
||||||
|
|
||||||
for (const route of ["/", "/app", "/app/*", "/admin", "/admin/*"]) {
|
for (const route of ["/", "/app", "/app/*", "/admin", "/admin/*"]) {
|
||||||
app.get(route, { schema: { hide: true } }, async (_request, reply) => {
|
app.get(route, { schema: { hide: true } }, async (request, reply) => {
|
||||||
reply.type("text/html; charset=utf-8");
|
reply.type("text/html; charset=utf-8");
|
||||||
return supportGateHtml;
|
if (!browserGate) return productIndexHtml ?? supportGateHtml;
|
||||||
|
const verified = verifyBrowserSupportCookie({
|
||||||
|
cookieHeader: headerValue(request.headers.cookie),
|
||||||
|
release: browserSupportRelease,
|
||||||
|
secChUa: headerValue(request.headers["sec-ch-ua"]),
|
||||||
|
secret: browserSupportSecret,
|
||||||
|
});
|
||||||
|
return verified.supported && productIndexHtml ? productIndexHtml : supportGateHtml;
|
||||||
});
|
});
|
||||||
}
|
}
|
||||||
|
app.get("/assets/*", { schema: { hide: true } }, async (request, reply) => {
|
||||||
|
const relativePath = decodeURIComponent(request.url.split("?", 1)[0]!.slice("/assets/".length));
|
||||||
|
const assetPath = resolve(productWebRoot, "assets", relativePath);
|
||||||
|
if (!assetPath.startsWith(resolve(productWebRoot, "assets")) || !existsSync(assetPath)) return reply.code(404).send();
|
||||||
|
reply.type(productAssetContentType(assetPath));
|
||||||
|
return reply.send(readFileSync(assetPath));
|
||||||
|
});
|
||||||
app.get("/support-gate.css", { schema: { hide: true } }, async (_request, reply) => {
|
app.get("/support-gate.css", { schema: { hide: true } }, async (_request, reply) => {
|
||||||
reply.type("text/css; charset=utf-8");
|
reply.type("text/css; charset=utf-8");
|
||||||
return supportGateCss;
|
return supportGateCss;
|
||||||
|
|||||||
@@ -33,6 +33,12 @@ const fixedDirectories = [
|
|||||||
"logs/supervisor",
|
"logs/supervisor",
|
||||||
] as const;
|
] as const;
|
||||||
|
|
||||||
|
export function ensureLocalDataRuntimeDirectories(dataRoot: string) {
|
||||||
|
for (const directory of fixedDirectories) {
|
||||||
|
mkdirSync(join(resolve(dataRoot), directory), { recursive: true });
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
export const DATA_TRANSFER_POLICY = {
|
export const DATA_TRANSFER_POLICY = {
|
||||||
allowed_downloads: ["original_generation", "jpg", "png"],
|
allowed_downloads: ["original_generation", "jpg", "png"],
|
||||||
application_backup: false,
|
application_backup: false,
|
||||||
@@ -219,9 +225,7 @@ export function initializeLocalDataRoot(input: {
|
|||||||
|
|
||||||
const createdRoot = !existsSync(validation.normalized_path);
|
const createdRoot = !existsSync(validation.normalized_path);
|
||||||
try {
|
try {
|
||||||
for (const directory of fixedDirectories) {
|
ensureLocalDataRuntimeDirectories(validation.normalized_path);
|
||||||
mkdirSync(join(validation.normalized_path, directory), { recursive: true });
|
|
||||||
}
|
|
||||||
openInstanceDatabase(join(validation.normalized_path, "db", "dada.sqlite3"));
|
openInstanceDatabase(join(validation.normalized_path, "db", "dada.sqlite3"));
|
||||||
const configuration: InstanceConfiguration = {
|
const configuration: InstanceConfiguration = {
|
||||||
data_root: validation.normalized_path,
|
data_root: validation.normalized_path,
|
||||||
|
|||||||
@@ -5,7 +5,7 @@ import { registrationNotice } from "@dada/shared-contracts";
|
|||||||
|
|
||||||
import { createApp } from "./app.js";
|
import { createApp } from "./app.js";
|
||||||
import { readBrowserSupportRelease } from "./browser-support.js";
|
import { readBrowserSupportRelease } from "./browser-support.js";
|
||||||
import { defaultInstanceConfigPath, readConfiguredLocalDataRoot } from "./local-data-root.js";
|
import { defaultInstanceConfigPath, ensureLocalDataRuntimeDirectories, readConfiguredLocalDataRoot } from "./local-data-root.js";
|
||||||
import { ManagedStorage } from "./managed-storage.js";
|
import { ManagedStorage } from "./managed-storage.js";
|
||||||
import { LatestExportService } from "./latest-exports.js";
|
import { LatestExportService } from "./latest-exports.js";
|
||||||
import { CreditService } from "./credits.js";
|
import { CreditService } from "./credits.js";
|
||||||
@@ -16,7 +16,7 @@ import { MockResendAdapter } from "./resend-adapter.js";
|
|||||||
import { readSecureConfigCandidate } from "./secure-config.js";
|
import { readSecureConfigCandidate } from "./secure-config.js";
|
||||||
import { StructuredJsonlLogger } from "./structured-log.js";
|
import { StructuredJsonlLogger } from "./structured-log.js";
|
||||||
import { attachApiSupervisorControl, initializeApiCredentialClients, receiveApiCredentials } from "./supervisor-channel.js";
|
import { attachApiSupervisorControl, initializeApiCredentialClients, receiveApiCredentials } from "./supervisor-channel.js";
|
||||||
import { ModelConfigurationService } from "./model-configuration.js";
|
import { ModelConfigurationService, portableRuntimeModelCandidates } from "./model-configuration.js";
|
||||||
import { MockAmapAdapter, type AmapAdapter } from "./amap-adapter.js";
|
import { MockAmapAdapter, type AmapAdapter } from "./amap-adapter.js";
|
||||||
import { StickerReleaseService } from "./sticker-releases.js";
|
import { StickerReleaseService } from "./sticker-releases.js";
|
||||||
import { createAdminDiagnosticsProvider, createAdminServicesStorageProvider } from "./admin-state.js";
|
import { createAdminDiagnosticsProvider, createAdminServicesStorageProvider } from "./admin-state.js";
|
||||||
@@ -40,6 +40,7 @@ if (credentialChannelEnabled) {
|
|||||||
.update(`Dada/P0A/${purpose}/v1`, "utf8")
|
.update(`Dada/P0A/${purpose}/v1`, "utf8")
|
||||||
.digest();
|
.digest();
|
||||||
const dataRoot = readConfiguredLocalDataRoot(instanceConfigPath);
|
const dataRoot = readConfiguredLocalDataRoot(instanceConfigPath);
|
||||||
|
ensureLocalDataRuntimeDirectories(dataRoot);
|
||||||
const databasePath = join(dataRoot, "db", "dada.sqlite3");
|
const databasePath = join(dataRoot, "db", "dada.sqlite3");
|
||||||
registration = new RegistrationService({
|
registration = new RegistrationService({
|
||||||
adminAllowlistPepper: Buffer.from(clients.adminAllowlistPepper),
|
adminAllowlistPepper: Buffer.from(clients.adminAllowlistPepper),
|
||||||
@@ -55,7 +56,7 @@ if (credentialChannelEnabled) {
|
|||||||
storage = new ManagedStorage({ dataRoot, databasePath });
|
storage = new ManagedStorage({ dataRoot, databasePath });
|
||||||
stickers = new StickerReleaseService({ databasePath, storage });
|
stickers = new StickerReleaseService({ databasePath, storage });
|
||||||
latestExports = new LatestExportService({ databasePath, storage });
|
latestExports = new LatestExportService({ databasePath, storage });
|
||||||
models = new ModelConfigurationService({ database: registration.database });
|
models = new ModelConfigurationService({ database: registration.database, seedCandidates: portableRuntimeModelCandidates });
|
||||||
recentAssets = new RecentAssetService({ database: registration.database });
|
recentAssets = new RecentAssetService({ database: registration.database });
|
||||||
registration.applySecureConfig(readSecureConfigCandidate(instanceConfigPath));
|
registration.applySecureConfig(readSecureConfigCandidate(instanceConfigPath));
|
||||||
} catch (error) {
|
} catch (error) {
|
||||||
|
|||||||
@@ -111,7 +111,7 @@ const defaultErrorMapping: Record<string, string> = {
|
|||||||
upstream_timeout: "upstream_timeout",
|
upstream_timeout: "upstream_timeout",
|
||||||
};
|
};
|
||||||
|
|
||||||
const seedCandidates: ModelConfigCandidate[] = [
|
const defaultSeedCandidates: ModelConfigCandidate[] = [
|
||||||
{
|
{
|
||||||
model_id: modelIds[0], display_name: "Gemini 3.1 Flash Image Preview", enabled: true, is_default: true,
|
model_id: modelIds[0], display_name: "Gemini 3.1 Flash Image Preview", enabled: true, is_default: true,
|
||||||
recommendation_priority: 1, route_profile: { endpoint: "https://mock.invalid/v1/images", mode: "sync" },
|
recommendation_priority: 1, route_profile: { endpoint: "https://mock.invalid/v1/images", mode: "sync" },
|
||||||
@@ -138,6 +138,42 @@ const seedCandidates: ModelConfigCandidate[] = [
|
|||||||
},
|
},
|
||||||
];
|
];
|
||||||
|
|
||||||
|
export const portableRuntimeModelCandidates: ModelConfigCandidate[] = [
|
||||||
|
{
|
||||||
|
...defaultSeedCandidates[0]!,
|
||||||
|
display_name: "Gemini 3.1 Flash Image",
|
||||||
|
route_profile: {
|
||||||
|
endpoint: "https://oneapi.intelligrow.cn/v1/chat/completions",
|
||||||
|
mode: "sync",
|
||||||
|
protocol_version: "gemini-openai-chat-v1",
|
||||||
|
provider_model_id: "gemini-3.1-flash-image",
|
||||||
|
},
|
||||||
|
gateway_account_ref: "oneapi-intelligrow-test",
|
||||||
|
contract_validation_status: "verified",
|
||||||
|
contract_evidence_ref: "contract:wp7-02:gemini-3.1-flash-image:v7",
|
||||||
|
},
|
||||||
|
{
|
||||||
|
...defaultSeedCandidates[1]!,
|
||||||
|
enabled: false,
|
||||||
|
route_profile: { endpoint: "https://oneapi.intelligrow.cn/unsupported", mode: "disabled" },
|
||||||
|
gateway_account_ref: "oneapi-intelligrow-test",
|
||||||
|
contract_validation_status: "unverified",
|
||||||
|
contract_evidence_ref: null,
|
||||||
|
},
|
||||||
|
{
|
||||||
|
...defaultSeedCandidates[2]!,
|
||||||
|
route_profile: {
|
||||||
|
endpoint: "https://oneapi.intelligrow.cn/v1/images/generations",
|
||||||
|
mode: "sync",
|
||||||
|
protocol_version: "openai-images-v1",
|
||||||
|
reference_endpoint: "https://oneapi.intelligrow.cn/v1/images/edits",
|
||||||
|
},
|
||||||
|
gateway_account_ref: "oneapi-intelligrow-test",
|
||||||
|
contract_validation_status: "verified",
|
||||||
|
contract_evidence_ref: "contract:wp7-02:gpt-image-2:v2",
|
||||||
|
},
|
||||||
|
];
|
||||||
|
|
||||||
function stableJson(value: unknown): string {
|
function stableJson(value: unknown): string {
|
||||||
if (Array.isArray(value)) return `[${value.map(stableJson).join(",")}]`;
|
if (Array.isArray(value)) return `[${value.map(stableJson).join(",")}]`;
|
||||||
if (value && typeof value === "object") {
|
if (value && typeof value === "object") {
|
||||||
@@ -207,17 +243,20 @@ export interface ModelConfigurationServiceOptions {
|
|||||||
clock?: () => number;
|
clock?: () => number;
|
||||||
database: BetterSqlite3.Database;
|
database: BetterSqlite3.Database;
|
||||||
onChanged?: (configSetVersion: number) => void;
|
onChanged?: (configSetVersion: number) => void;
|
||||||
|
seedCandidates?: ModelConfigCandidate[];
|
||||||
}
|
}
|
||||||
|
|
||||||
export class ModelConfigurationService {
|
export class ModelConfigurationService {
|
||||||
readonly database: BetterSqlite3.Database;
|
readonly database: BetterSqlite3.Database;
|
||||||
readonly #clock: () => number;
|
readonly #clock: () => number;
|
||||||
readonly #onChanged: ((configSetVersion: number) => void) | undefined;
|
readonly #onChanged: ((configSetVersion: number) => void) | undefined;
|
||||||
|
readonly #seedCandidates: ModelConfigCandidate[];
|
||||||
|
|
||||||
constructor(options: ModelConfigurationServiceOptions) {
|
constructor(options: ModelConfigurationServiceOptions) {
|
||||||
this.database = options.database;
|
this.database = options.database;
|
||||||
this.#clock = options.clock ?? Date.now;
|
this.#clock = options.clock ?? Date.now;
|
||||||
this.#onChanged = options.onChanged;
|
this.#onChanged = options.onChanged;
|
||||||
|
this.#seedCandidates = structuredClone(options.seedCandidates ?? defaultSeedCandidates);
|
||||||
this.ensureSchema();
|
this.ensureSchema();
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -503,7 +542,7 @@ export class ModelConfigurationService {
|
|||||||
const current = this.database.prepare("SELECT config_set_id FROM model_config_current WHERE singleton = 1").get() as { config_set_id: string } | undefined;
|
const current = this.database.prepare("SELECT config_set_id FROM model_config_current WHERE singleton = 1").get() as { config_set_id: string } | undefined;
|
||||||
if (current) return;
|
if (current) return;
|
||||||
const seed = this.database.transaction(() => {
|
const seed = this.database.transaction(() => {
|
||||||
validateModelConfigurationCandidateSet(seedCandidates);
|
validateModelConfigurationCandidateSet(this.#seedCandidates);
|
||||||
const now = this.#clock();
|
const now = this.#clock();
|
||||||
const setId = randomUUID();
|
const setId = randomUUID();
|
||||||
this.database.prepare("INSERT INTO model_config_sets (config_set_id, config_set_version, created_at, created_by) VALUES (?, 1, ?, 'system_seed')")
|
this.database.prepare("INSERT INTO model_config_sets (config_set_id, config_set_version, created_at, created_by) VALUES (?, 1, ?, 'system_seed')")
|
||||||
@@ -520,7 +559,9 @@ export class ModelConfigurationService {
|
|||||||
INSERT INTO model_config_set_members (config_set_id, model_id, config_version, enabled, is_default, recommendation_priority)
|
INSERT INTO model_config_set_members (config_set_id, model_id, config_version, enabled, is_default, recommendation_priority)
|
||||||
VALUES (?, ?, 1, ?, ?, ?)
|
VALUES (?, ?, 1, ?, ?, ?)
|
||||||
`);
|
`);
|
||||||
for (const candidate of seedCandidates) {
|
for (const candidate of this.#seedCandidates) {
|
||||||
|
const contractStatus = candidate.contract_validation_status ?? "unverified";
|
||||||
|
const contractEvidenceRef = contractStatus === "unverified" ? null : candidate.contract_evidence_ref ?? null;
|
||||||
const routeProfileId = profileRef("route", candidate.route_profile);
|
const routeProfileId = profileRef("route", candidate.route_profile);
|
||||||
const errorMappingProfileId = profileRef("error", candidate.error_mapping_profile);
|
const errorMappingProfileId = profileRef("error", candidate.error_mapping_profile);
|
||||||
this.database.prepare("INSERT OR IGNORE INTO gateway_route_profiles (route_profile_id, profile_json, created_at) VALUES (?, ?, ?)")
|
this.database.prepare("INSERT OR IGNORE INTO gateway_route_profiles (route_profile_id, profile_json, created_at) VALUES (?, ?, ?)")
|
||||||
@@ -530,12 +571,14 @@ export class ModelConfigurationService {
|
|||||||
insertVersion.run(candidate.model_id, candidate.display_name, candidate.enabled ? 1 : 0, candidate.is_default ? 1 : 0,
|
insertVersion.run(candidate.model_id, candidate.display_name, candidate.enabled ? 1 : 0, candidate.is_default ? 1 : 0,
|
||||||
candidate.recommendation_priority, routeProfileId, stableJson(candidate.route_profile), candidate.gateway_account_ref,
|
candidate.recommendation_priority, routeProfileId, stableJson(candidate.route_profile), candidate.gateway_account_ref,
|
||||||
errorMappingProfileId, stableJson(candidate.error_mapping_profile), candidate.credit_cost, stableJson(candidate.supported_ratios), stableJson(candidate.reference_limits),
|
errorMappingProfileId, stableJson(candidate.error_mapping_profile), candidate.credit_cost, stableJson(candidate.supported_ratios), stableJson(candidate.reference_limits),
|
||||||
candidate.prompt_max_length, candidate.safety_source, "unverified", null, fingerprint(candidate), now);
|
candidate.prompt_max_length, candidate.safety_source, contractStatus, contractEvidenceRef, fingerprint(candidate), now);
|
||||||
insertMember.run(setId, candidate.model_id, candidate.enabled ? 1 : 0, candidate.is_default ? 1 : 0, candidate.recommendation_priority);
|
insertMember.run(setId, candidate.model_id, candidate.enabled ? 1 : 0, candidate.is_default ? 1 : 0, candidate.recommendation_priority);
|
||||||
|
const available = candidate.enabled && contractStatus === "verified";
|
||||||
|
const runtimeReason = !candidate.enabled ? "configured_disabled" : available ? "available" : "contract_unverified";
|
||||||
this.database.prepare(`
|
this.database.prepare(`
|
||||||
INSERT INTO model_runtime_availability (model_id, available_for_new_jobs, reason, checked_at, runtime_availability_version)
|
INSERT INTO model_runtime_availability (model_id, available_for_new_jobs, reason, checked_at, runtime_availability_version)
|
||||||
VALUES (?, 0, 'contract_unverified', ?, 0)
|
VALUES (?, ?, ?, ?, 0)
|
||||||
`).run(candidate.model_id, now);
|
`).run(candidate.model_id, available ? 1 : 0, runtimeReason, now);
|
||||||
}
|
}
|
||||||
this.database.prepare("INSERT INTO model_config_current (singleton, config_set_id) VALUES (1, ?)").run(setId);
|
this.database.prepare("INSERT INTO model_config_current (singleton, config_set_id) VALUES (1, ?)").run(setId);
|
||||||
});
|
});
|
||||||
|
|||||||
@@ -1,6 +1,6 @@
|
|||||||
import { createConnection } from "node:net";
|
import { createConnection } from "node:net";
|
||||||
|
|
||||||
import { RealAmapAdapter } from "./amap-adapter.js";
|
import { MockAmapAdapter, RealAmapAdapter } from "./amap-adapter.js";
|
||||||
|
|
||||||
const API_CREDENTIALS = ["Dada/P0A/api/resend", "Dada/P0A/api/amap", "Dada/P0A/admin/pepper"] as const;
|
const API_CREDENTIALS = ["Dada/P0A/api/resend", "Dada/P0A/api/amap", "Dada/P0A/admin/pepper"] as const;
|
||||||
|
|
||||||
@@ -15,7 +15,7 @@ export async function receiveApiCredentials(input: NodeJS.ReadableStream = proce
|
|||||||
if (names.length !== expected.length || names.some((name, index) => name !== expected[index])) {
|
if (names.length !== expected.length || names.some((name, index) => name !== expected[index])) {
|
||||||
throw new Error("API credential channel contains an unexpected credential scope.");
|
throw new Error("API credential channel contains an unexpected credential scope.");
|
||||||
}
|
}
|
||||||
if (expected.some((name) => typeof parsed[name] !== "string" || parsed[name] === "")) {
|
if (expected.some((name) => typeof parsed[name] !== "string")) {
|
||||||
throw new Error("API credential channel contains an invalid credential value.");
|
throw new Error("API credential channel contains an invalid credential value.");
|
||||||
}
|
}
|
||||||
return parsed as Record<(typeof API_CREDENTIALS)[number], string>;
|
return parsed as Record<(typeof API_CREDENTIALS)[number], string>;
|
||||||
@@ -27,12 +27,12 @@ export async function receiveApiCredentials(input: NodeJS.ReadableStream = proce
|
|||||||
}
|
}
|
||||||
|
|
||||||
export function initializeApiCredentialClients(credentials: Record<(typeof API_CREDENTIALS)[number], string>) {
|
export function initializeApiCredentialClients(credentials: Record<(typeof API_CREDENTIALS)[number], string>) {
|
||||||
const configured = API_CREDENTIALS.every((name) => credentials[name].length > 0);
|
try {
|
||||||
try {
|
const adminPepper = credentials["Dada/P0A/admin/pepper"];
|
||||||
if (!configured) throw new Error("API credential client initialization failed.");
|
if (!adminPepper) throw new Error("admin_pepper_not_configured");
|
||||||
return {
|
return {
|
||||||
adminAllowlistPepper: Buffer.from(credentials["Dada/P0A/admin/pepper"], "utf8"),
|
adminAllowlistPepper: Buffer.from(adminPepper, "utf8"),
|
||||||
amap: new RealAmapAdapter(credentials["Dada/P0A/api/amap"]),
|
amap: credentials["Dada/P0A/api/amap"] ? new RealAmapAdapter(credentials["Dada/P0A/api/amap"]) : new MockAmapAdapter(),
|
||||||
};
|
};
|
||||||
} finally {
|
} finally {
|
||||||
for (const name of API_CREDENTIALS) credentials[name] = "";
|
for (const name of API_CREDENTIALS) credentials[name] = "";
|
||||||
|
|||||||
@@ -69,6 +69,7 @@ async function checkSupport() {
|
|||||||
browserValue.textContent = `${result.browser.brand} ${result.browser.major}`;
|
browserValue.textContent = `${result.browser.brand} ${result.browser.major}`;
|
||||||
supportedValue.textContent = supportedLabel(result.supported_browsers);
|
supportedValue.textContent = supportedLabel(result.supported_browsers);
|
||||||
window.dispatchEvent(new CustomEvent("dada:support-ready"));
|
window.dispatchEvent(new CustomEvent("dada:support-ready"));
|
||||||
|
window.location.replace(window.location.pathname.startsWith("/admin") ? "/admin" : "/app");
|
||||||
return;
|
return;
|
||||||
}
|
}
|
||||||
showBlocked(
|
showBlocked(
|
||||||
|
|||||||
@@ -7,6 +7,11 @@ export interface GenerationAdapterRequest {
|
|||||||
prompt: string;
|
prompt: string;
|
||||||
ratio: "3:4" | "1:1" | "4:3" | "9:16";
|
ratio: "3:4" | "1:1" | "4:3" | "9:16";
|
||||||
referenceAssetIds: readonly string[];
|
referenceAssetIds: readonly string[];
|
||||||
|
referenceImages?: readonly {
|
||||||
|
assetId: string;
|
||||||
|
bytes: Buffer;
|
||||||
|
mimeType: "image/jpeg" | "image/png" | "image/webp";
|
||||||
|
}[];
|
||||||
}
|
}
|
||||||
|
|
||||||
export interface NormalizedGenerationOutput {
|
export interface NormalizedGenerationOutput {
|
||||||
@@ -27,6 +32,7 @@ export type GenerationAdapterResult =
|
|||||||
};
|
};
|
||||||
|
|
||||||
export interface GenerationAdapter {
|
export interface GenerationAdapter {
|
||||||
|
dispose?(): void;
|
||||||
start(request: GenerationAdapterRequest): Promise<GenerationAdapterResult>;
|
start(request: GenerationAdapterRequest): Promise<GenerationAdapterResult>;
|
||||||
poll?(upstreamJobReference: string): Promise<GenerationAdapterResult>;
|
poll?(upstreamJobReference: string): Promise<GenerationAdapterResult>;
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -0,0 +1,47 @@
|
|||||||
|
import type { GenerationAdapter } from "./ai-adapter-contract.js";
|
||||||
|
|
||||||
|
export type AiRuntimeProbeResult =
|
||||||
|
| {
|
||||||
|
code: "ai_probe_passed";
|
||||||
|
mime_type: "image/jpeg" | "image/png" | "image/webp";
|
||||||
|
pixel_height: number;
|
||||||
|
pixel_width: number;
|
||||||
|
real_calls: 1;
|
||||||
|
success: true;
|
||||||
|
}
|
||||||
|
| {
|
||||||
|
code: "ai_probe_failed";
|
||||||
|
error_category: string;
|
||||||
|
real_calls: 1;
|
||||||
|
success: false;
|
||||||
|
};
|
||||||
|
|
||||||
|
export async function runAiRuntimeProbe(adapter: GenerationAdapter): Promise<AiRuntimeProbeResult> {
|
||||||
|
const result = await adapter.start({
|
||||||
|
configSnapshot: { probe: true },
|
||||||
|
generationId: "00000000-0000-4000-8000-000000000002",
|
||||||
|
modelId: "gemini-3.1-flash-image-preview",
|
||||||
|
prompt: "生成一张简洁的红蓝几何色块测试图,不含文字。",
|
||||||
|
ratio: "1:1",
|
||||||
|
referenceAssetIds: [],
|
||||||
|
});
|
||||||
|
if (result.status === "failed") {
|
||||||
|
return { code: "ai_probe_failed", error_category: result.category, real_calls: 1, success: false };
|
||||||
|
}
|
||||||
|
if (result.status !== "completed" || result.outputs.length !== 1) {
|
||||||
|
return { code: "ai_probe_failed", error_category: "gateway_contract_invalid", real_calls: 1, success: false };
|
||||||
|
}
|
||||||
|
const output = result.outputs[0]!;
|
||||||
|
try {
|
||||||
|
return {
|
||||||
|
code: "ai_probe_passed",
|
||||||
|
mime_type: output.mimeType,
|
||||||
|
pixel_height: output.pixelHeight,
|
||||||
|
pixel_width: output.pixelWidth,
|
||||||
|
real_calls: 1,
|
||||||
|
success: true,
|
||||||
|
};
|
||||||
|
} finally {
|
||||||
|
output.bytes.fill(0);
|
||||||
|
}
|
||||||
|
}
|
||||||
@@ -1,11 +1,11 @@
|
|||||||
import { createHash, randomUUID } from "node:crypto";
|
import { createHash, randomUUID } from "node:crypto";
|
||||||
import { existsSync, mkdirSync, renameSync, rmSync, writeFileSync } from "node:fs";
|
import { existsSync, mkdirSync, readFileSync, renameSync, rmSync, statSync, writeFileSync } from "node:fs";
|
||||||
import { dirname, join, resolve } from "node:path";
|
import { dirname, isAbsolute, join, relative, resolve, sep } from "node:path";
|
||||||
|
|
||||||
import Database from "better-sqlite3";
|
import Database from "better-sqlite3";
|
||||||
import type BetterSqlite3 from "better-sqlite3";
|
import type BetterSqlite3 from "better-sqlite3";
|
||||||
|
|
||||||
import type { GenerationAdapter, GenerationAdapterResult, NormalizedGenerationOutput } from "./ai-adapter-contract.js";
|
import type { GenerationAdapter, GenerationAdapterRequest, GenerationAdapterResult, NormalizedGenerationOutput } from "./ai-adapter-contract.js";
|
||||||
import { GatewayBalanceRuntime } from "./gateway-balance-runtime.js";
|
import { GatewayBalanceRuntime } from "./gateway-balance-runtime.js";
|
||||||
import { generationErrorRegistry, type GenerationErrorCategory } from "./generation-error-registry.js";
|
import { generationErrorRegistry, type GenerationErrorCategory } from "./generation-error-registry.js";
|
||||||
import { configureWorkerDatabase } from "./sqlite-connection.js";
|
import { configureWorkerDatabase } from "./sqlite-connection.js";
|
||||||
@@ -91,7 +91,8 @@ export class GenerationProcessor {
|
|||||||
this.clock = input.clock ?? Date.now;
|
this.clock = input.clock ?? Date.now;
|
||||||
this.dataRoot = resolve(input.dataRoot);
|
this.dataRoot = resolve(input.dataRoot);
|
||||||
this.workerId = input.workerId;
|
this.workerId = input.workerId;
|
||||||
this.database = new Database(input.databasePath);
|
const nativeBinding = process.env.DADA_SQLITE_NATIVE_BINDING;
|
||||||
|
this.database = new Database(input.databasePath, nativeBinding ? { nativeBinding } : undefined);
|
||||||
configureWorkerDatabase(this.database);
|
configureWorkerDatabase(this.database);
|
||||||
this.migrate();
|
this.migrate();
|
||||||
this.gatewayBalance = new GatewayBalanceRuntime({ clock: this.clock, database: this.database });
|
this.gatewayBalance = new GatewayBalanceRuntime({ clock: this.clock, database: this.database });
|
||||||
@@ -119,6 +120,7 @@ export class GenerationProcessor {
|
|||||||
.run("worker_stopped", now, this.workerId);
|
.run("worker_stopped", now, this.workerId);
|
||||||
});
|
});
|
||||||
this.gatewayBalance.close();
|
this.gatewayBalance.close();
|
||||||
|
this.adapter.dispose?.();
|
||||||
this.database.close();
|
this.database.close();
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -143,6 +145,12 @@ export class GenerationProcessor {
|
|||||||
SELECT managed_file_id FROM generation_reference_snapshots WHERE generation_id = ? ORDER BY position
|
SELECT managed_file_id FROM generation_reference_snapshots WHERE generation_id = ? ORDER BY position
|
||||||
`).all(generationId) as Array<{ managed_file_id: string }>;
|
`).all(generationId) as Array<{ managed_file_id: string }>;
|
||||||
let adapterResult: GenerationAdapterResult;
|
let adapterResult: GenerationAdapterResult;
|
||||||
|
let referenceImages: NonNullable<GenerationAdapterRequest["referenceImages"]>;
|
||||||
|
try {
|
||||||
|
referenceImages = this.loadReferenceImages(references.map((row) => row.managed_file_id));
|
||||||
|
} catch {
|
||||||
|
return this.completeFailure(job, "reference_invalid", "reference_load_failed");
|
||||||
|
}
|
||||||
try {
|
try {
|
||||||
if (job.upstream_job_reference) {
|
if (job.upstream_job_reference) {
|
||||||
if (!this.adapter.poll) return this.completeFailure(job, "unknown_retryable", "poll_unsupported", undefined, false, "pending_manual_review");
|
if (!this.adapter.poll) return this.completeFailure(job, "unknown_retryable", "poll_unsupported", undefined, false, "pending_manual_review");
|
||||||
@@ -155,10 +163,13 @@ export class GenerationProcessor {
|
|||||||
prompt: job.prompt,
|
prompt: job.prompt,
|
||||||
ratio: job.ratio,
|
ratio: job.ratio,
|
||||||
referenceAssetIds: references.map((row) => row.managed_file_id),
|
referenceAssetIds: references.map((row) => row.managed_file_id),
|
||||||
|
referenceImages,
|
||||||
}));
|
}));
|
||||||
}
|
}
|
||||||
} catch {
|
} catch {
|
||||||
return this.completeFailure(job, "unknown_retryable", "adapter_exception", undefined, false, "pending_manual_review");
|
return this.completeFailure(job, "unknown_retryable", "adapter_exception", undefined, false, "pending_manual_review");
|
||||||
|
} finally {
|
||||||
|
for (const reference of referenceImages) reference.bytes.fill(0);
|
||||||
}
|
}
|
||||||
|
|
||||||
if (adapterResult.status === "failed") return this.completeFailure(job, adapterResult.category, adapterResult.sourceCategory, adapterResult.balanceSignal);
|
if (adapterResult.status === "failed") return this.completeFailure(job, adapterResult.category, adapterResult.sourceCategory, adapterResult.balanceSignal);
|
||||||
@@ -174,6 +185,29 @@ export class GenerationProcessor {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
private loadReferenceImages(referenceAssetIds: string[]): NonNullable<GenerationAdapterRequest["referenceImages"]> {
|
||||||
|
return referenceAssetIds.map((assetId) => {
|
||||||
|
const row = this.database.prepare(`
|
||||||
|
SELECT relative_path, mime_type FROM managed_files
|
||||||
|
WHERE file_id = ? AND file_kind = 'reference' AND status = 'committed'
|
||||||
|
`).get(assetId) as { mime_type: string; relative_path: string } | undefined;
|
||||||
|
if (!row || !["image/jpeg", "image/png", "image/webp"].includes(row.mime_type) || isAbsolute(row.relative_path)) {
|
||||||
|
throw new Error("reference_invalid");
|
||||||
|
}
|
||||||
|
const path = resolve(this.dataRoot, row.relative_path);
|
||||||
|
const child = relative(this.dataRoot, path);
|
||||||
|
if (!child || child === ".." || child.startsWith(`..${sep}`) || isAbsolute(child)
|
||||||
|
|| !existsSync(path) || !statSync(path).isFile()) {
|
||||||
|
throw new Error("reference_invalid");
|
||||||
|
}
|
||||||
|
return {
|
||||||
|
assetId,
|
||||||
|
bytes: readFileSync(path),
|
||||||
|
mimeType: row.mime_type as "image/jpeg" | "image/png" | "image/webp",
|
||||||
|
};
|
||||||
|
});
|
||||||
|
}
|
||||||
|
|
||||||
private claim(generationId: string) {
|
private claim(generationId: string) {
|
||||||
return this.immediate(() => {
|
return this.immediate(() => {
|
||||||
const row = this.readJob(generationId);
|
const row = this.readJob(generationId);
|
||||||
|
|||||||
@@ -0,0 +1,217 @@
|
|||||||
|
import sharp from "sharp";
|
||||||
|
|
||||||
|
import type {
|
||||||
|
GenerationAdapter,
|
||||||
|
GenerationAdapterRequest,
|
||||||
|
GenerationAdapterResult,
|
||||||
|
NormalizedGenerationOutput,
|
||||||
|
} from "./ai-adapter-contract.js";
|
||||||
|
import { gptImageRequestSizeForRatio, normalizeImageOutputToRatio } from "./image-output-normalizer.mjs";
|
||||||
|
|
||||||
|
const geminiProductModelId = "gemini-3.1-flash-image-preview";
|
||||||
|
const geminiProviderModelId = "gemini-3.1-flash-image";
|
||||||
|
const gptImageModelId = "gpt-image-2";
|
||||||
|
const geminiEndpoint = "https://oneapi.intelligrow.cn/v1/chat/completions";
|
||||||
|
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;
|
||||||
|
|
||||||
|
type FetchLike = typeof fetch;
|
||||||
|
|
||||||
|
class OneApiRuntimeError extends Error {
|
||||||
|
constructor(
|
||||||
|
readonly category: "gateway_balance_insufficient" | "gateway_contract_invalid" | "model_disabled" | "reference_invalid" | "upstream_failed" | "upstream_timeout" | "unknown_non_retryable",
|
||||||
|
readonly sourceCategory: string,
|
||||||
|
) {
|
||||||
|
super(sourceCategory);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
function failure(error: unknown): GenerationAdapterResult {
|
||||||
|
if (error instanceof OneApiRuntimeError) {
|
||||||
|
return { category: error.category, sourceCategory: error.sourceCategory, status: "failed" };
|
||||||
|
}
|
||||||
|
if (error instanceof Error && error.name === "AbortError") {
|
||||||
|
return { category: "upstream_timeout", sourceCategory: "upstream_timeout", status: "failed" };
|
||||||
|
}
|
||||||
|
return { category: "upstream_failed", sourceCategory: "upstream_failed", status: "failed" };
|
||||||
|
}
|
||||||
|
|
||||||
|
function mapHttpFailure(status: number) {
|
||||||
|
if (status === 408 || status === 504) return new OneApiRuntimeError("upstream_timeout", `upstream_http_${status}`);
|
||||||
|
if (status === 429) return new OneApiRuntimeError("gateway_balance_insufficient", "upstream_http_429");
|
||||||
|
if (status >= 500) return new OneApiRuntimeError("upstream_failed", `upstream_http_${status}`);
|
||||||
|
if (status === 400 || status === 404 || status === 422) return new OneApiRuntimeError("gateway_contract_invalid", `upstream_http_${status}`);
|
||||||
|
return new OneApiRuntimeError("unknown_non_retryable", `upstream_http_${status}`);
|
||||||
|
}
|
||||||
|
|
||||||
|
async function readBoundedJson(response: Response) {
|
||||||
|
const declaredLength = Number(response.headers.get("content-length") ?? 0);
|
||||||
|
if (Number.isFinite(declaredLength) && declaredLength > maximumResponseBytes) {
|
||||||
|
throw new OneApiRuntimeError("gateway_contract_invalid", "upstream_response_too_large");
|
||||||
|
}
|
||||||
|
if (!response.body) throw new OneApiRuntimeError("gateway_contract_invalid", "upstream_response_empty");
|
||||||
|
const reader = response.body.getReader();
|
||||||
|
const chunks: Buffer[] = [];
|
||||||
|
let total = 0;
|
||||||
|
try {
|
||||||
|
while (true) {
|
||||||
|
const next = await reader.read();
|
||||||
|
if (next.done) break;
|
||||||
|
const chunk = Buffer.from(next.value);
|
||||||
|
total += chunk.length;
|
||||||
|
if (total > maximumResponseBytes) {
|
||||||
|
await reader.cancel();
|
||||||
|
throw new OneApiRuntimeError("gateway_contract_invalid", "upstream_response_too_large");
|
||||||
|
}
|
||||||
|
chunks.push(chunk);
|
||||||
|
}
|
||||||
|
try {
|
||||||
|
return JSON.parse(Buffer.concat(chunks).toString("utf8")) as unknown;
|
||||||
|
} catch {
|
||||||
|
throw new OneApiRuntimeError("gateway_contract_invalid", "upstream_response_invalid");
|
||||||
|
}
|
||||||
|
} finally {
|
||||||
|
for (const chunk of chunks) chunk.fill(0);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
function extractGeminiImage(response: unknown) {
|
||||||
|
if (!response || typeof response !== "object" || !("choices" in response) || !Array.isArray(response.choices)) {
|
||||||
|
throw new OneApiRuntimeError("gateway_contract_invalid", "response_shape_invalid");
|
||||||
|
}
|
||||||
|
const choice = response.choices[0];
|
||||||
|
const content = choice && typeof choice === "object" && "message" in choice && choice.message && typeof choice.message === "object"
|
||||||
|
&& "content" in choice.message && typeof choice.message.content === "string" ? choice.message.content : "";
|
||||||
|
const matches = [...content.matchAll(/!\[[^\]]*\]\(\s*data:(image\/(?:jpeg|png|webp));base64,([A-Za-z0-9+/=\r\n]+)\s*\)/gi)];
|
||||||
|
if (matches.length !== 1) throw new OneApiRuntimeError("gateway_contract_invalid", "response_single_image_required");
|
||||||
|
return { bytes: Buffer.from(matches[0]![2]!, "base64"), declaredMimeType: matches[0]![1]!.toLowerCase() };
|
||||||
|
}
|
||||||
|
|
||||||
|
function extractGptImage(response: unknown) {
|
||||||
|
if (!response || typeof response !== "object" || !("data" in response) || !Array.isArray(response.data)
|
||||||
|
|| response.data.length !== 1 || !response.data[0] || typeof response.data[0] !== "object"
|
||||||
|
|| !("b64_json" in response.data[0]) || typeof response.data[0].b64_json !== "string") {
|
||||||
|
throw new OneApiRuntimeError("gateway_contract_invalid", "response_single_image_required");
|
||||||
|
}
|
||||||
|
return { bytes: Buffer.from(response.data[0].b64_json, "base64"), declaredMimeType: undefined };
|
||||||
|
}
|
||||||
|
|
||||||
|
async function normalizeOutput(bytes: Buffer, declaredMimeType: string | undefined, ratio: GenerationAdapterRequest["ratio"]): Promise<NormalizedGenerationOutput> {
|
||||||
|
try {
|
||||||
|
const metadata = await sharp(bytes, { failOn: "error", limitInputPixels: 40_000_000 }).metadata();
|
||||||
|
const mimeType = metadata.format === "png" ? "image/png" : metadata.format === "jpeg" ? "image/jpeg" : metadata.format === "webp" ? "image/webp" : undefined;
|
||||||
|
if (!mimeType || !metadata.width || !metadata.height || (declaredMimeType && declaredMimeType !== mimeType)) {
|
||||||
|
throw new OneApiRuntimeError("gateway_contract_invalid", "response_media_invalid");
|
||||||
|
}
|
||||||
|
const normalized = await normalizeImageOutputToRatio({ bytes, mimeType, pixelHeight: metadata.height, pixelWidth: metadata.width, ratio });
|
||||||
|
return { bytes: normalized.bytes, mimeType: normalized.mimeType, pixelHeight: normalized.pixelHeight, pixelWidth: normalized.pixelWidth };
|
||||||
|
} catch (error) {
|
||||||
|
if (error instanceof OneApiRuntimeError) throw error;
|
||||||
|
throw new OneApiRuntimeError("gateway_contract_invalid", "response_media_invalid");
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
function validateRequest(request: GenerationAdapterRequest) {
|
||||||
|
if (!request.prompt.trim() || request.prompt.length > 1_000) throw new OneApiRuntimeError("gateway_contract_invalid", "prompt_invalid");
|
||||||
|
const references = request.referenceImages ?? [];
|
||||||
|
if (references.length !== request.referenceAssetIds.length || references.length > 2) {
|
||||||
|
throw new OneApiRuntimeError("reference_invalid", "reference_count_invalid");
|
||||||
|
}
|
||||||
|
const totalBytes = references.reduce((total, reference) => total + reference.bytes.length, 0);
|
||||||
|
if (totalBytes > 20 * 1024 * 1024 || references.some((reference) => reference.bytes.length === 0 || reference.bytes.length > 10 * 1024 * 1024)) {
|
||||||
|
throw new OneApiRuntimeError("reference_invalid", "reference_size_invalid");
|
||||||
|
}
|
||||||
|
return references;
|
||||||
|
}
|
||||||
|
|
||||||
|
function buildRequest(request: GenerationAdapterRequest) {
|
||||||
|
const references = validateRequest(request);
|
||||||
|
if (request.modelId === geminiProductModelId) {
|
||||||
|
const content = references.length === 0
|
||||||
|
? request.prompt
|
||||||
|
: [
|
||||||
|
{ text: request.prompt, type: "text" },
|
||||||
|
...references.map((reference) => ({
|
||||||
|
image_url: { url: `data:${reference.mimeType};base64,${reference.bytes.toString("base64")}` },
|
||||||
|
type: "image_url",
|
||||||
|
})),
|
||||||
|
];
|
||||||
|
return {
|
||||||
|
body: JSON.stringify({
|
||||||
|
extra_body: { google: { image_config: { aspect_ratio: request.ratio, image_size: "1K" } } },
|
||||||
|
messages: [{ content, role: "user" }],
|
||||||
|
model: geminiProviderModelId,
|
||||||
|
stream: false,
|
||||||
|
}),
|
||||||
|
contentType: "application/json",
|
||||||
|
endpoint: geminiEndpoint,
|
||||||
|
parser: extractGeminiImage,
|
||||||
|
};
|
||||||
|
}
|
||||||
|
if (request.modelId !== gptImageModelId) throw new OneApiRuntimeError("model_disabled", "model_not_supported");
|
||||||
|
if (references.length > 0) {
|
||||||
|
const form = new FormData();
|
||||||
|
form.append("model", gptImageModelId);
|
||||||
|
form.append("prompt", request.prompt);
|
||||||
|
form.append("response_format", "b64_json");
|
||||||
|
form.append("size", gptImageRequestSizeForRatio(request.ratio));
|
||||||
|
references.forEach((reference, index) => form.append("image[]", new Blob([reference.bytes], { type: reference.mimeType }), `reference-${index + 1}.png`));
|
||||||
|
return { body: form, contentType: undefined, endpoint: gptImageReferenceEndpoint, parser: extractGptImage };
|
||||||
|
}
|
||||||
|
return {
|
||||||
|
body: JSON.stringify({ model: gptImageModelId, prompt: request.prompt, response_format: "b64_json", size: gptImageRequestSizeForRatio(request.ratio) }),
|
||||||
|
contentType: "application/json",
|
||||||
|
endpoint: gptImageEndpoint,
|
||||||
|
parser: extractGptImage,
|
||||||
|
};
|
||||||
|
}
|
||||||
|
|
||||||
|
export class OneApiGenerationAdapter implements GenerationAdapter {
|
||||||
|
private readonly credential: Buffer;
|
||||||
|
private readonly fetchImpl: FetchLike;
|
||||||
|
private disposed = false;
|
||||||
|
|
||||||
|
constructor(input: { credential: Buffer; fetch?: FetchLike }) {
|
||||||
|
if (input.credential.length < 8) throw new Error("ai_gateway_credential_invalid");
|
||||||
|
this.credential = Buffer.from(input.credential);
|
||||||
|
this.fetchImpl = input.fetch ?? fetch;
|
||||||
|
}
|
||||||
|
|
||||||
|
async start(request: GenerationAdapterRequest): Promise<GenerationAdapterResult> {
|
||||||
|
if (this.disposed) return { category: "upstream_failed", sourceCategory: "adapter_disposed", status: "failed" };
|
||||||
|
const controller = new AbortController();
|
||||||
|
const timeout = setTimeout(() => controller.abort(), requestTimeoutMilliseconds);
|
||||||
|
let sourceBytes: Buffer | undefined;
|
||||||
|
try {
|
||||||
|
const providerRequest = buildRequest(request);
|
||||||
|
const headers = new Headers({ authorization: `Bearer ${this.credential.toString("utf8")}` });
|
||||||
|
if (providerRequest.contentType) headers.set("content-type", providerRequest.contentType);
|
||||||
|
const response = await this.fetchImpl(providerRequest.endpoint, {
|
||||||
|
body: providerRequest.body,
|
||||||
|
headers,
|
||||||
|
method: "POST",
|
||||||
|
redirect: "error",
|
||||||
|
signal: controller.signal,
|
||||||
|
});
|
||||||
|
if (!response.ok) throw mapHttpFailure(response.status);
|
||||||
|
const parsed = await readBoundedJson(response);
|
||||||
|
const extracted = providerRequest.parser(parsed);
|
||||||
|
sourceBytes = extracted.bytes;
|
||||||
|
const output = await normalizeOutput(sourceBytes, extracted.declaredMimeType, request.ratio);
|
||||||
|
return { outputs: [output], status: "completed" };
|
||||||
|
} catch (error) {
|
||||||
|
return failure(error);
|
||||||
|
} finally {
|
||||||
|
clearTimeout(timeout);
|
||||||
|
sourceBytes?.fill(0);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
dispose() {
|
||||||
|
if (this.disposed) return;
|
||||||
|
this.disposed = true;
|
||||||
|
this.credential.fill(0);
|
||||||
|
}
|
||||||
|
}
|
||||||
@@ -13,7 +13,7 @@ export async function receiveWorkerCredentials(input: NodeJS.ReadableStream = pr
|
|||||||
if (names.length !== expected.length || names.some((name, index) => name !== expected[index])) {
|
if (names.length !== expected.length || names.some((name, index) => name !== expected[index])) {
|
||||||
throw new Error("Worker credential channel contains an unexpected credential scope.");
|
throw new Error("Worker credential channel contains an unexpected credential scope.");
|
||||||
}
|
}
|
||||||
if (expected.some((name) => typeof parsed[name] !== "string" || parsed[name] === "")) {
|
if (expected.some((name) => typeof parsed[name] !== "string")) {
|
||||||
throw new Error("Worker credential channel contains an invalid credential value.");
|
throw new Error("Worker credential channel contains an invalid credential value.");
|
||||||
}
|
}
|
||||||
return parsed as Record<(typeof WORKER_CREDENTIALS)[number], string>;
|
return parsed as Record<(typeof WORKER_CREDENTIALS)[number], string>;
|
||||||
@@ -25,9 +25,13 @@ export async function receiveWorkerCredentials(input: NodeJS.ReadableStream = pr
|
|||||||
}
|
}
|
||||||
|
|
||||||
export function initializeWorkerCredentialClient(credentials: Record<(typeof WORKER_CREDENTIALS)[number], string>) {
|
export function initializeWorkerCredentialClient(credentials: Record<(typeof WORKER_CREDENTIALS)[number], string>) {
|
||||||
const configured = WORKER_CREDENTIALS.every((name) => credentials[name].length > 0);
|
const value = credentials["Dada/P0A/worker/ai-gateway"];
|
||||||
for (const name of WORKER_CREDENTIALS) credentials[name] = "";
|
try {
|
||||||
if (!configured) throw new Error("Worker credential client initialization failed.");
|
if (!value) throw new Error("worker_ai_gateway_not_configured");
|
||||||
|
return { aiGatewayCredential: Buffer.from(value, "utf8") };
|
||||||
|
} finally {
|
||||||
|
for (const name of WORKER_CREDENTIALS) credentials[name] = "";
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
export function attachWorkerSupervisorControl(pipeName: string, shutdown: () => Promise<void> | void) {
|
export function attachWorkerSupervisorControl(pipeName: string, shutdown: () => Promise<void> | void) {
|
||||||
|
|||||||
@@ -2,6 +2,9 @@ import { parentPort } from "node:worker_threads";
|
|||||||
import { join } from "node:path";
|
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 { GenerationProcessor } from "./generation-processor.js";
|
||||||
|
import { OneApiGenerationAdapter } from "./oneapi-generation-adapter.js";
|
||||||
import { readConfiguredLocalDataRoot } from "./runtime-config.js";
|
import { readConfiguredLocalDataRoot } from "./runtime-config.js";
|
||||||
import { RetentionCleanup } from "./retention-cleanup.js";
|
import { RetentionCleanup } from "./retention-cleanup.js";
|
||||||
import { ProjectPurgeCleanup } from "./project-purge-cleanup.js";
|
import { ProjectPurgeCleanup } from "./project-purge-cleanup.js";
|
||||||
@@ -21,8 +24,25 @@ if (workerPort) {
|
|||||||
});
|
});
|
||||||
}
|
}
|
||||||
|
|
||||||
if (!workerPort && process.argv.includes("--dada-credential-stdin")) {
|
if (!workerPort && process.argv.includes("--dada-ai-probe")) {
|
||||||
initializeWorkerCredentialClient(await receiveWorkerCredentials());
|
const credentialClient = initializeWorkerCredentialClient(await receiveWorkerCredentials());
|
||||||
|
let adapter: OneApiGenerationAdapter | undefined;
|
||||||
|
let probeResult: Awaited<ReturnType<typeof runAiRuntimeProbe>> | { code: "ai_probe_failed"; error_category: "upstream_failed"; real_calls: 0; success: false };
|
||||||
|
try {
|
||||||
|
adapter = new OneApiGenerationAdapter({ credential: credentialClient.aiGatewayCredential });
|
||||||
|
probeResult = await runAiRuntimeProbe(adapter);
|
||||||
|
} catch {
|
||||||
|
probeResult = { code: "ai_probe_failed", error_category: "upstream_failed", real_calls: 0, success: false };
|
||||||
|
} finally {
|
||||||
|
credentialClient.aiGatewayCredential.fill(0);
|
||||||
|
adapter?.dispose();
|
||||||
|
}
|
||||||
|
await new Promise<void>((resolveWrite, rejectWrite) => {
|
||||||
|
process.stdout.write(JSON.stringify(probeResult), (error) => error ? rejectWrite(error) : resolveWrite());
|
||||||
|
});
|
||||||
|
process.exit(probeResult.success ? 0 : 2);
|
||||||
|
} else if (!workerPort && process.argv.includes("--dada-credential-stdin")) {
|
||||||
|
const credentialClient = initializeWorkerCredentialClient(await receiveWorkerCredentials());
|
||||||
const controlPipeIndex = process.argv.indexOf("--dada-control-pipe");
|
const controlPipeIndex = process.argv.indexOf("--dada-control-pipe");
|
||||||
const controlPipe = process.argv[controlPipeIndex + 1];
|
const controlPipe = process.argv[controlPipeIndex + 1];
|
||||||
if (controlPipeIndex < 0 || !controlPipe) throw new Error("Supervisor control pipe name is required.");
|
if (controlPipeIndex < 0 || !controlPipe) throw new Error("Supervisor control pipe name is required.");
|
||||||
@@ -31,11 +51,15 @@ if (!workerPort && process.argv.includes("--dada-credential-stdin")) {
|
|||||||
let retention: RetentionCleanup | undefined;
|
let retention: RetentionCleanup | undefined;
|
||||||
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 generationTimer: ReturnType<typeof setInterval> | 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);
|
||||||
|
processor?.close();
|
||||||
storage?.close();
|
storage?.close();
|
||||||
});
|
});
|
||||||
let storageStatus: "active" | "unavailable" = "active";
|
let storageStatus: "active" | "unavailable" = "active";
|
||||||
@@ -45,6 +69,13 @@ if (!workerPort && process.argv.includes("--dada-credential-stdin")) {
|
|||||||
storage = new WorkerStorageStatus(databasePath);
|
storage = new WorkerStorageStatus(databasePath);
|
||||||
retention = new RetentionCleanup({ databasePath });
|
retention = new RetentionCleanup({ databasePath });
|
||||||
projectCleanup = new ProjectPurgeCleanup({ dataRoot, databasePath });
|
projectCleanup = new ProjectPurgeCleanup({ dataRoot, databasePath });
|
||||||
|
let adapter: OneApiGenerationAdapter;
|
||||||
|
try {
|
||||||
|
adapter = new OneApiGenerationAdapter({ credential: credentialClient.aiGatewayCredential });
|
||||||
|
} finally {
|
||||||
|
credentialClient.aiGatewayCredential.fill(0);
|
||||||
|
}
|
||||||
|
processor = new GenerationProcessor({ adapter, dataRoot, databasePath, workerId: `portable-oneapi-worker-${process.pid}` });
|
||||||
const runRetentionCleanup = () => {
|
const runRetentionCleanup = () => {
|
||||||
try {
|
try {
|
||||||
retention?.purgeExpired();
|
retention?.purgeExpired();
|
||||||
@@ -71,6 +102,7 @@ if (!workerPort && process.argv.includes("--dada-credential-stdin")) {
|
|||||||
});
|
});
|
||||||
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);
|
||||||
} catch {
|
} catch {
|
||||||
storageStatus = "unavailable";
|
storageStatus = "unavailable";
|
||||||
control.reportStatus("storage_unavailable");
|
control.reportStatus("storage_unavailable");
|
||||||
|
|||||||
@@ -1,10 +1,20 @@
|
|||||||
|
import { execFileSync } from "node:child_process";
|
||||||
|
import { readFileSync } from "node:fs";
|
||||||
import { resolve } from "node:path";
|
import { resolve } from "node:path";
|
||||||
|
|
||||||
import { buildAndValidatePortablePackage } from "./lib/portable-package.mjs";
|
import { buildAndValidatePortablePackage } from "./lib/portable-package.mjs";
|
||||||
|
import { validateFinalReleaseRecord } from "./lib/wp7-07-final-release.mjs";
|
||||||
|
|
||||||
const outputIndex = process.argv.indexOf("--output");
|
const outputIndex = process.argv.indexOf("--output");
|
||||||
const outputRoot = outputIndex >= 0 ? resolve(process.argv[outputIndex + 1]) : resolve(".build", "portable-release");
|
const outputRoot = outputIndex >= 0 ? resolve(process.argv[outputIndex + 1]) : resolve(".build", "portable-release");
|
||||||
const result = await buildAndValidatePortablePackage({ outputRoot });
|
const previousRelease = JSON.parse(readFileSync(resolve("RELEASE.json"), "utf8"));
|
||||||
|
const releaseRecord = validateFinalReleaseRecord({
|
||||||
|
...previousRelease,
|
||||||
|
buildCommit: execFileSync("git", ["rev-parse", "HEAD"], { encoding: "utf8" }).trim(),
|
||||||
|
maintenanceFromCommit: previousRelease.buildCommit,
|
||||||
|
recordedAt: new Date().toISOString(),
|
||||||
|
});
|
||||||
|
const result = await buildAndValidatePortablePackage({ outputRoot, releaseRecord });
|
||||||
console.log(JSON.stringify({
|
console.log(JSON.stringify({
|
||||||
package: result.packageManifest.package_name,
|
package: result.packageManifest.package_name,
|
||||||
sha256: result.packageManifest.zip_sha256,
|
sha256: result.packageManifest.zip_sha256,
|
||||||
|
|||||||
@@ -336,7 +336,7 @@ export async function buildAndValidatePortablePackage({ evidenceDirectory, outpu
|
|||||||
debug("copy API application");
|
debug("copy API application");
|
||||||
const apiDependencies = copyApplication(join(repositoryRoot, "apps", "api"), join(serverRoot, "api"), ["@fastify/multipart", "@fastify/swagger", "@sinclair/typebox", "better-sqlite3", "fastify", "sharp"]);
|
const apiDependencies = copyApplication(join(repositoryRoot, "apps", "api"), join(serverRoot, "api"), ["@fastify/multipart", "@fastify/swagger", "@sinclair/typebox", "better-sqlite3", "fastify", "sharp"]);
|
||||||
debug("copy Worker application");
|
debug("copy Worker application");
|
||||||
const workerDependencies = copyApplication(join(repositoryRoot, "apps", "worker"), join(serverRoot, "worker"), ["better-sqlite3"]);
|
const workerDependencies = copyApplication(join(repositoryRoot, "apps", "worker"), join(serverRoot, "worker"), ["better-sqlite3", "sharp"]);
|
||||||
const sharedDestination = join(serverRoot, "api", "node_modules", "@dada", "shared-contracts");
|
const sharedDestination = join(serverRoot, "api", "node_modules", "@dada", "shared-contracts");
|
||||||
mkdirSync(sharedDestination, { recursive: true });
|
mkdirSync(sharedDestination, { recursive: true });
|
||||||
copyTree(join(repositoryRoot, "packages", "shared-contracts", "dist"), join(sharedDestination, "dist"));
|
copyTree(join(repositoryRoot, "packages", "shared-contracts", "dist"), join(sharedDestination, "dist"));
|
||||||
|
|||||||
@@ -40,6 +40,7 @@ internal static class Program
|
|||||||
var supervisor = await TestSupervisorLifecycleAsync();
|
var supervisor = await TestSupervisorLifecycleAsync();
|
||||||
await TestAmapProbeSecurityAsync();
|
await TestAmapProbeSecurityAsync();
|
||||||
TestSecureConfigurationPersistence();
|
TestSecureConfigurationPersistence();
|
||||||
|
TestRuntimeDirectoryBootstrap();
|
||||||
TestStructuredLogging();
|
TestStructuredLogging();
|
||||||
WriteEvidence(Environment.GetEnvironmentVariable("DADA_EVIDENCE_DIR_SEC"), security);
|
WriteEvidence(Environment.GetEnvironmentVariable("DADA_EVIDENCE_DIR_SEC"), security);
|
||||||
WriteEvidence(Environment.GetEnvironmentVariable("DADA_EVIDENCE_DIR_SUP"), supervisor);
|
WriteEvidence(Environment.GetEnvironmentVariable("DADA_EVIDENCE_DIR_SUP"), supervisor);
|
||||||
@@ -125,6 +126,29 @@ internal static class Program
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
private static void TestRuntimeDirectoryBootstrap()
|
||||||
|
{
|
||||||
|
var root = Path.Combine(Path.GetTempPath(), $"dada-runtime-root-{Guid.NewGuid():N}");
|
||||||
|
try
|
||||||
|
{
|
||||||
|
Directory.CreateDirectory(root);
|
||||||
|
SupervisorRuntime.EnsureRuntimeDirectories(root);
|
||||||
|
foreach (var relativePath in new[]
|
||||||
|
{
|
||||||
|
"db", "content/references", "content/generated", "content/exports",
|
||||||
|
"managed-assets", "derived-assets", "staging",
|
||||||
|
"logs/api", "logs/worker", "logs/supervisor",
|
||||||
|
})
|
||||||
|
{
|
||||||
|
True(Directory.Exists(Path.Combine(root, relativePath)), $"runtime directory missing: {relativePath}");
|
||||||
|
}
|
||||||
|
}
|
||||||
|
finally
|
||||||
|
{
|
||||||
|
if (Directory.Exists(root)) Directory.Delete(root, recursive: true);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
private static async Task<object> TestCredentialBoundaryAsync()
|
private static async Task<object> TestCredentialBoundaryAsync()
|
||||||
{
|
{
|
||||||
var store = new TestCredentialStore();
|
var store = new TestCredentialStore();
|
||||||
@@ -146,6 +170,8 @@ internal static class Program
|
|||||||
True(leakProbe.SensitiveOutputDetected, "credential echo must be detected");
|
True(leakProbe.SensitiveOutputDetected, "credential echo must be detected");
|
||||||
Equal(string.Empty, leakProbe.StandardOutput, "credential echo output discarded");
|
Equal(string.Empty, leakProbe.StandardOutput, "credential echo output discarded");
|
||||||
Equal(string.Empty, leakProbe.StandardError, "credential echo error discarded");
|
Equal(string.Empty, leakProbe.StandardError, "credential echo error discarded");
|
||||||
|
True(AiGatewayProbe.TryValidateOutput("{\"code\":\"ai_probe_passed\",\"mime_type\":\"image/png\",\"pixel_height\":1080,\"pixel_width\":1080,\"real_calls\":1,\"success\":true}", out _), "AI probe success output accepted");
|
||||||
|
False(AiGatewayProbe.TryValidateOutput("{\"code\":\"ai_probe_passed\",\"raw_body\":\"private\",\"real_calls\":1,\"success\":true}", out _), "AI probe private output rejected");
|
||||||
|
|
||||||
var externalArguments = new[]
|
var externalArguments = new[]
|
||||||
{
|
{
|
||||||
|
|||||||
@@ -0,0 +1,74 @@
|
|||||||
|
using System.Diagnostics;
|
||||||
|
using System.Text.Json;
|
||||||
|
|
||||||
|
namespace Dada.Supervisor;
|
||||||
|
|
||||||
|
internal static class AiGatewayProbe
|
||||||
|
{
|
||||||
|
internal static async Task<int> RunAsync(ICredentialStore credentials, CancellationToken cancellationToken = default)
|
||||||
|
{
|
||||||
|
var node = Path.Combine(AppContext.BaseDirectory, "runtime", "node.exe");
|
||||||
|
var worker = Path.Combine(AppContext.BaseDirectory, "server", "worker.mjs");
|
||||||
|
if (!File.Exists(node) || !File.Exists(worker)) return WriteFailure("ai_probe_runtime_missing", 0);
|
||||||
|
var startInfo = new ProcessStartInfo(node) { WorkingDirectory = AppContext.BaseDirectory };
|
||||||
|
startInfo.Environment["DADA_SQLITE_NATIVE_BINDING"] = Path.Combine(AppContext.BaseDirectory, "server", "native", "better_sqlite3.node");
|
||||||
|
startInfo.ArgumentList.Add(worker);
|
||||||
|
startInfo.ArgumentList.Add("--dada-ai-probe");
|
||||||
|
startInfo.ArgumentList.Add("--dada-credential-stdin");
|
||||||
|
var result = await CredentialProcessLauncher.RunToCompletionAsync(startInfo, ChildRole.Worker, credentials, cancellationToken);
|
||||||
|
if (result.SensitiveOutputDetected || result.StandardError.Length > 0 || !TryValidateOutput(result.StandardOutput, out var sanitized))
|
||||||
|
{
|
||||||
|
return WriteFailure("ai_probe_runtime_failed", 0);
|
||||||
|
}
|
||||||
|
Console.WriteLine(sanitized);
|
||||||
|
return result.ExitCode;
|
||||||
|
}
|
||||||
|
|
||||||
|
internal static bool TryValidateOutput(string output, out string sanitized)
|
||||||
|
{
|
||||||
|
sanitized = string.Empty;
|
||||||
|
try
|
||||||
|
{
|
||||||
|
using var document = JsonDocument.Parse(output);
|
||||||
|
var root = document.RootElement;
|
||||||
|
if (root.ValueKind != JsonValueKind.Object) return false;
|
||||||
|
var allowed = new HashSet<string>(StringComparer.Ordinal)
|
||||||
|
{
|
||||||
|
"code", "error_category", "mime_type", "pixel_height", "pixel_width", "real_calls", "success",
|
||||||
|
};
|
||||||
|
if (root.EnumerateObject().Any(property => !allowed.Contains(property.Name))) return false;
|
||||||
|
if (!root.TryGetProperty("success", out var success) || success.ValueKind is not (JsonValueKind.True or JsonValueKind.False)) return false;
|
||||||
|
if (!root.TryGetProperty("real_calls", out var realCalls) || realCalls.ValueKind != JsonValueKind.Number || !realCalls.TryGetInt32(out var count) || count is < 0 or > 1) return false;
|
||||||
|
var passed = success.GetBoolean();
|
||||||
|
var code = root.GetProperty("code").GetString();
|
||||||
|
if (passed)
|
||||||
|
{
|
||||||
|
if (code != "ai_probe_passed" || count != 1) return false;
|
||||||
|
var mime = root.GetProperty("mime_type").GetString();
|
||||||
|
if (mime is not ("image/jpeg" or "image/png" or "image/webp")) return false;
|
||||||
|
if (!PositiveDimension(root, "pixel_width") || !PositiveDimension(root, "pixel_height")) return false;
|
||||||
|
}
|
||||||
|
else
|
||||||
|
{
|
||||||
|
if (code != "ai_probe_failed" || !root.TryGetProperty("error_category", out var category)
|
||||||
|
|| category.ValueKind != JsonValueKind.String || (category.GetString()?.Length ?? 0) is < 1 or > 64) return false;
|
||||||
|
}
|
||||||
|
sanitized = JsonSerializer.Serialize(root);
|
||||||
|
return true;
|
||||||
|
}
|
||||||
|
catch (Exception exception) when (exception is JsonException or InvalidOperationException or KeyNotFoundException)
|
||||||
|
{
|
||||||
|
return false;
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
private static bool PositiveDimension(JsonElement root, string name) =>
|
||||||
|
root.TryGetProperty(name, out var value) && value.ValueKind == JsonValueKind.Number
|
||||||
|
&& value.TryGetInt32(out var dimension) && dimension is > 0 and <= 4096;
|
||||||
|
|
||||||
|
private static int WriteFailure(string code, int realCalls)
|
||||||
|
{
|
||||||
|
Console.WriteLine(JsonSerializer.Serialize(new { code, real_calls = realCalls, success = false }));
|
||||||
|
return 2;
|
||||||
|
}
|
||||||
|
}
|
||||||
@@ -50,7 +50,9 @@ internal static class CredentialProcessLauncher
|
|||||||
var credentials = new Dictionary<string, string>(StringComparer.Ordinal);
|
var credentials = new Dictionary<string, string>(StringComparer.Ordinal);
|
||||||
foreach (var target in CredentialCatalog.RequiredFor(role))
|
foreach (var target in CredentialCatalog.RequiredFor(role))
|
||||||
{
|
{
|
||||||
credentials[target] = store.Read(target) ?? throw new MissingCredentialException(target);
|
var value = store.Read(target);
|
||||||
|
if (role == ChildRole.Worker && string.IsNullOrWhiteSpace(value)) throw new MissingCredentialException(target);
|
||||||
|
credentials[target] = value ?? string.Empty;
|
||||||
}
|
}
|
||||||
|
|
||||||
startInfo.UseShellExecute = false;
|
startInfo.UseShellExecute = false;
|
||||||
@@ -98,7 +100,9 @@ internal static class CredentialProcessLauncher
|
|||||||
var credentials = new Dictionary<string, string>(StringComparer.Ordinal);
|
var credentials = new Dictionary<string, string>(StringComparer.Ordinal);
|
||||||
foreach (var target in CredentialCatalog.RequiredFor(role))
|
foreach (var target in CredentialCatalog.RequiredFor(role))
|
||||||
{
|
{
|
||||||
credentials[target] = store.Read(target) ?? throw new MissingCredentialException(target);
|
var value = store.Read(target);
|
||||||
|
if (role == ChildRole.Worker && string.IsNullOrWhiteSpace(value)) throw new MissingCredentialException(target);
|
||||||
|
credentials[target] = value ?? string.Empty;
|
||||||
}
|
}
|
||||||
|
|
||||||
startInfo.UseShellExecute = false;
|
startInfo.UseShellExecute = false;
|
||||||
|
|||||||
@@ -26,7 +26,7 @@ internal static class OfflineCommandRouter
|
|||||||
return args[0] switch
|
return args[0] switch
|
||||||
{
|
{
|
||||||
"configure" => RunConfigure(args.Skip(1).ToArray()),
|
"configure" => RunConfigure(args.Skip(1).ToArray()),
|
||||||
"secrets" => RunSecrets(args.Skip(1).ToArray(), credentials),
|
"secrets" => await RunSecretsAsync(args.Skip(1).ToArray(), credentials),
|
||||||
"admin-allowlist" => RunAdminAllowlist(args.Skip(1).ToArray(), credentials),
|
"admin-allowlist" => RunAdminAllowlist(args.Skip(1).ToArray(), credentials),
|
||||||
"doctor" when args.Length == 1 => RunDoctor(credentials),
|
"doctor" when args.Length == 1 => RunDoctor(credentials),
|
||||||
"validate-external" => await ControlledExternalValidationLauncher.RunAsync(args.Skip(1).ToArray(), credentials),
|
"validate-external" => await ControlledExternalValidationLauncher.RunAsync(args.Skip(1).ToArray(), credentials),
|
||||||
@@ -66,7 +66,7 @@ internal static class OfflineCommandRouter
|
|||||||
return 0;
|
return 0;
|
||||||
}
|
}
|
||||||
|
|
||||||
private static int RunSecrets(string[] args, ICredentialStore store)
|
private static async Task<int> RunSecretsAsync(string[] args, ICredentialStore store)
|
||||||
{
|
{
|
||||||
if (args.Length != 2 || !TryResolveCredential(args[1], out var target)) return Usage();
|
if (args.Length != 2 || !TryResolveCredential(args[1], out var target)) return Usage();
|
||||||
switch (args[0])
|
switch (args[0])
|
||||||
@@ -86,6 +86,8 @@ internal static class OfflineCommandRouter
|
|||||||
return 0;
|
return 0;
|
||||||
case "probe" when target == CredentialCatalog.ApiAmap:
|
case "probe" when target == CredentialCatalog.ApiAmap:
|
||||||
return AmapProbe.Run(store.Read(target));
|
return AmapProbe.Run(store.Read(target));
|
||||||
|
case "probe" when target == CredentialCatalog.WorkerAiGateway:
|
||||||
|
return await AiGatewayProbe.RunAsync(store);
|
||||||
default:
|
default:
|
||||||
return Usage();
|
return Usage();
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -1,9 +1,23 @@
|
|||||||
using System.Diagnostics;
|
using System.Diagnostics;
|
||||||
|
using System.Security.Cryptography;
|
||||||
|
|
||||||
namespace Dada.Supervisor;
|
namespace Dada.Supervisor;
|
||||||
|
|
||||||
internal sealed class SupervisorRuntime : IAsyncDisposable
|
internal sealed class SupervisorRuntime : IAsyncDisposable
|
||||||
{
|
{
|
||||||
|
private static readonly string[] RequiredDataDirectories =
|
||||||
|
[
|
||||||
|
"db",
|
||||||
|
Path.Combine("content", "references"),
|
||||||
|
Path.Combine("content", "generated"),
|
||||||
|
Path.Combine("content", "exports"),
|
||||||
|
"managed-assets",
|
||||||
|
"derived-assets",
|
||||||
|
"staging",
|
||||||
|
Path.Combine("logs", "api"),
|
||||||
|
Path.Combine("logs", "worker"),
|
||||||
|
Path.Combine("logs", "supervisor"),
|
||||||
|
];
|
||||||
private readonly ICredentialStore credentials;
|
private readonly ICredentialStore credentials;
|
||||||
private ManagedComponentSupervisor? api;
|
private ManagedComponentSupervisor? api;
|
||||||
private ManagedComponentSupervisor? worker;
|
private ManagedComponentSupervisor? worker;
|
||||||
@@ -25,6 +39,7 @@ internal sealed class SupervisorRuntime : IAsyncDisposable
|
|||||||
{
|
{
|
||||||
return SupervisorState.StartupFailed;
|
return SupervisorState.StartupFailed;
|
||||||
}
|
}
|
||||||
|
EnsureRuntimeDirectories(configuration.LocalDataRoot);
|
||||||
try
|
try
|
||||||
{
|
{
|
||||||
logger = new StructuredJsonlLogger(Path.Combine(configuration.LocalDataRoot, "logs", "supervisor"), "supervisor");
|
logger = new StructuredJsonlLogger(Path.Combine(configuration.LocalDataRoot, "logs", "supervisor"), "supervisor");
|
||||||
@@ -34,22 +49,45 @@ internal sealed class SupervisorRuntime : IAsyncDisposable
|
|||||||
{
|
{
|
||||||
return SupervisorState.StorageUnavailable;
|
return SupervisorState.StorageUnavailable;
|
||||||
}
|
}
|
||||||
if (CredentialCatalog.RequiredFor(ChildRole.Api).Concat(CredentialCatalog.RequiredFor(ChildRole.Worker)).Any(target => !credentials.IsConfigured(target)))
|
EnsureAdminPepper();
|
||||||
{
|
|
||||||
return SupervisorState.StartupFailed;
|
|
||||||
}
|
|
||||||
|
|
||||||
var node = Path.Combine(AppContext.BaseDirectory, "runtime", "node.exe");
|
var node = Path.Combine(AppContext.BaseDirectory, "runtime", "node.exe");
|
||||||
var apiEntry = Path.Combine(AppContext.BaseDirectory, "server", "api.mjs");
|
var apiEntry = Path.Combine(AppContext.BaseDirectory, "server", "api.mjs");
|
||||||
var workerEntry = Path.Combine(AppContext.BaseDirectory, "server", "worker.mjs");
|
var workerEntry = Path.Combine(AppContext.BaseDirectory, "server", "worker.mjs");
|
||||||
if (!File.Exists(node) || !File.Exists(apiEntry) || !File.Exists(workerEntry)) return SupervisorState.StartupFailed;
|
if (!File.Exists(node) || !File.Exists(apiEntry) || !File.Exists(workerEntry)) return SupervisorState.StartupFailed;
|
||||||
|
|
||||||
api = CreateComponent(node, apiEntry, ChildRole.Api, SupervisorState.ApiDegraded);
|
try
|
||||||
await api.StartAsync(cancellationToken);
|
{
|
||||||
|
api = CreateComponent(node, apiEntry, ChildRole.Api, SupervisorState.ApiDegraded);
|
||||||
|
await api.StartAsync(cancellationToken);
|
||||||
|
|
||||||
worker = CreateComponent(node, workerEntry, ChildRole.Worker, SupervisorState.WorkerDegraded);
|
worker = CreateComponent(node, workerEntry, ChildRole.Worker, SupervisorState.WorkerDegraded);
|
||||||
await worker.StartAsync(cancellationToken);
|
await worker.StartAsync(cancellationToken);
|
||||||
return TryLog(new StructuredLogEvent("ready", ErrorCategory: "none")) ? SupervisorState.Ready : SupervisorState.StorageUnavailable;
|
return TryLog(new StructuredLogEvent("ready", ErrorCategory: "none")) ? SupervisorState.Ready : SupervisorState.StorageUnavailable;
|
||||||
|
}
|
||||||
|
catch
|
||||||
|
{
|
||||||
|
await StopComponentsAsync();
|
||||||
|
return TryLog(new StructuredLogEvent("failed", ErrorCategory: "service_unavailable"))
|
||||||
|
? SupervisorState.StartupFailed
|
||||||
|
: SupervisorState.StorageUnavailable;
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
internal static void EnsureRuntimeDirectories(string dataRoot)
|
||||||
|
{
|
||||||
|
foreach (var directory in RequiredDataDirectories)
|
||||||
|
{
|
||||||
|
Directory.CreateDirectory(Path.Combine(dataRoot, directory));
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
private void EnsureAdminPepper()
|
||||||
|
{
|
||||||
|
if (credentials.IsConfigured(CredentialCatalog.AdminPepper)) return;
|
||||||
|
var pepper = Convert.ToBase64String(RandomNumberGenerator.GetBytes(32));
|
||||||
|
credentials.Write(CredentialCatalog.AdminPepper, pepper);
|
||||||
|
Array.Clear(System.Text.Encoding.UTF8.GetBytes(pepper));
|
||||||
}
|
}
|
||||||
|
|
||||||
private ManagedComponentSupervisor CreateComponent(string node, string entry, ChildRole role, SupervisorState degradedState)
|
private ManagedComponentSupervisor CreateComponent(string node, string entry, ChildRole role, SupervisorState degradedState)
|
||||||
@@ -60,6 +98,7 @@ internal sealed class SupervisorRuntime : IAsyncDisposable
|
|||||||
startInfo.WorkingDirectory = AppContext.BaseDirectory;
|
startInfo.WorkingDirectory = AppContext.BaseDirectory;
|
||||||
startInfo.Environment["DADA_SQLITE_NATIVE_BINDING"] = Path.Combine(AppContext.BaseDirectory, "server", "native", "better_sqlite3.node");
|
startInfo.Environment["DADA_SQLITE_NATIVE_BINDING"] = Path.Combine(AppContext.BaseDirectory, "server", "native", "better_sqlite3.node");
|
||||||
startInfo.Environment["DADA_SUPPORT_GATE_ROOT"] = Path.Combine(AppContext.BaseDirectory, "web", "support-gate");
|
startInfo.Environment["DADA_SUPPORT_GATE_ROOT"] = Path.Combine(AppContext.BaseDirectory, "web", "support-gate");
|
||||||
|
startInfo.Environment["DADA_WEB_ROOT"] = Path.Combine(AppContext.BaseDirectory, "web");
|
||||||
startInfo.Environment["DADA_INSTANCE_CONFIG_PATH"] = Path.Combine(Environment.GetFolderPath(Environment.SpecialFolder.LocalApplicationData), "Dada", "P0A", "config", "instance.json");
|
startInfo.Environment["DADA_INSTANCE_CONFIG_PATH"] = Path.Combine(Environment.GetFolderPath(Environment.SpecialFolder.LocalApplicationData), "Dada", "P0A", "config", "instance.json");
|
||||||
startInfo.ArgumentList.Add(entry);
|
startInfo.ArgumentList.Add(entry);
|
||||||
var child = await ManagedChildProcess.StartAsync(startInfo, role, credentials, cancellationToken);
|
var child = await ManagedChildProcess.StartAsync(startInfo, role, credentials, cancellationToken);
|
||||||
@@ -95,10 +134,26 @@ internal sealed class SupervisorRuntime : IAsyncDisposable
|
|||||||
}
|
}
|
||||||
|
|
||||||
public async ValueTask DisposeAsync()
|
public async ValueTask DisposeAsync()
|
||||||
|
{
|
||||||
|
await StopComponentsAsync();
|
||||||
|
}
|
||||||
|
|
||||||
|
private async Task StopComponentsAsync()
|
||||||
{
|
{
|
||||||
var stops = new List<Task>();
|
var stops = new List<Task>();
|
||||||
if (worker is not null) stops.Add(worker.DisposeAsync().AsTask());
|
if (worker is not null) stops.Add(worker.DisposeAsync().AsTask());
|
||||||
if (api is not null) stops.Add(api.DisposeAsync().AsTask());
|
if (api is not null) stops.Add(api.DisposeAsync().AsTask());
|
||||||
await Task.WhenAll(stops);
|
try
|
||||||
|
{
|
||||||
|
await Task.WhenAll(stops);
|
||||||
|
}
|
||||||
|
catch
|
||||||
|
{
|
||||||
|
}
|
||||||
|
finally
|
||||||
|
{
|
||||||
|
worker = null;
|
||||||
|
api = null;
|
||||||
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -0,0 +1,44 @@
|
|||||||
|
import { createRequire } from "node:module";
|
||||||
|
import { describe, expect, it } from "vitest";
|
||||||
|
|
||||||
|
import {
|
||||||
|
ModelConfigurationService,
|
||||||
|
portableRuntimeModelCandidates,
|
||||||
|
} from "../../apps/api/src/model-configuration.js";
|
||||||
|
|
||||||
|
const requireFromApi = createRequire(new URL("../../apps/api/package.json", import.meta.url));
|
||||||
|
const Database = requireFromApi("better-sqlite3") as new (path: string) => {
|
||||||
|
close(): void;
|
||||||
|
};
|
||||||
|
|
||||||
|
describe("POSTV1-02 portable runtime model seed", () => {
|
||||||
|
it("enables only models backed by the real OneAPI contract", () => {
|
||||||
|
const database = new Database(":memory:");
|
||||||
|
try {
|
||||||
|
const models = new ModelConfigurationService({ database, seedCandidates: portableRuntimeModelCandidates }).read();
|
||||||
|
const flash = models.models.find((model) => model.model_id === "gemini-3.1-flash-image-preview");
|
||||||
|
const pro = models.models.find((model) => model.model_id === "gemini-3-pro-image-preview");
|
||||||
|
const gpt = models.models.find((model) => model.model_id === "gpt-image-2");
|
||||||
|
|
||||||
|
expect(models.configured_default_model_id).toBe("gemini-3.1-flash-image-preview");
|
||||||
|
expect(flash).toMatchObject({
|
||||||
|
contract_validation_status: "verified",
|
||||||
|
enabled: true,
|
||||||
|
runtime_availability: { available_for_new_jobs: true, reason: "available" },
|
||||||
|
});
|
||||||
|
expect(flash?.route_profile).toMatchObject({ endpoint: "https://oneapi.intelligrow.cn/v1/chat/completions" });
|
||||||
|
expect(pro).toMatchObject({
|
||||||
|
contract_validation_status: "unverified",
|
||||||
|
enabled: false,
|
||||||
|
runtime_availability: { available_for_new_jobs: false, reason: "configured_disabled" },
|
||||||
|
});
|
||||||
|
expect(gpt).toMatchObject({
|
||||||
|
contract_validation_status: "verified",
|
||||||
|
enabled: true,
|
||||||
|
runtime_availability: { available_for_new_jobs: true, reason: "available" },
|
||||||
|
});
|
||||||
|
} finally {
|
||||||
|
database.close();
|
||||||
|
}
|
||||||
|
});
|
||||||
|
});
|
||||||
@@ -125,7 +125,12 @@ afterAll(async () => {
|
|||||||
|
|
||||||
describe("TDD-WP0-BRW-001 supported browser contract", () => {
|
describe("TDD-WP0-BRW-001 supported browser contract", () => {
|
||||||
it("cross-checks UA-CH and issues only a short-lived signed support cookie", async () => {
|
it("cross-checks UA-CH and issues only a short-lived signed support cookie", async () => {
|
||||||
const app = await createApp({ browserSupportRelease: testBrowserSupportRelease } as never);
|
const app = await createApp({
|
||||||
|
browserSupportRelease: testBrowserSupportRelease,
|
||||||
|
productIndexHtml: "<!doctype html><title>Dada product test</title><div id=\"root\"></div>",
|
||||||
|
} as never);
|
||||||
|
const gate = await app.inject({ headers: { host: "127.0.0.1:43121" }, method: "GET", url: "/" });
|
||||||
|
expect(gate.body).toContain("当前浏览器无法使用 Dada");
|
||||||
const checked = await app.inject({
|
const checked = await app.inject({
|
||||||
headers: supportedEdge.headers,
|
headers: supportedEdge.headers,
|
||||||
method: "POST",
|
method: "POST",
|
||||||
@@ -150,6 +155,14 @@ describe("TDD-WP0-BRW-001 supported browser contract", () => {
|
|||||||
|
|
||||||
const cookie = supportCookie(checked);
|
const cookie = supportCookie(checked);
|
||||||
expect(cookie).toBeDefined();
|
expect(cookie).toBeDefined();
|
||||||
|
const productHtml = await app.inject({
|
||||||
|
headers: { cookie, host: "127.0.0.1:43121", "sec-ch-ua": supportedEdge.headers["sec-ch-ua"] },
|
||||||
|
method: "GET",
|
||||||
|
url: "/app",
|
||||||
|
});
|
||||||
|
expect(productHtml.statusCode).toBe(200);
|
||||||
|
expect(productHtml.body).toContain("Dada product test");
|
||||||
|
expect(productHtml.body).not.toContain("当前浏览器无法使用 Dada");
|
||||||
const product = await app.inject({
|
const product = await app.inject({
|
||||||
headers: {
|
headers: {
|
||||||
cookie,
|
cookie,
|
||||||
|
|||||||
@@ -0,0 +1,167 @@
|
|||||||
|
import test from "node:test";
|
||||||
|
import assert from "node:assert/strict";
|
||||||
|
import { mkdtemp, readFile, rm, writeFile } from "node:fs/promises";
|
||||||
|
import { existsSync } from "node:fs";
|
||||||
|
import { join, resolve } from "node:path";
|
||||||
|
import { tmpdir } from "node:os";
|
||||||
|
import { spawn } from "node:child_process";
|
||||||
|
import { createServer } from "node:net";
|
||||||
|
|
||||||
|
const packageRoot = resolve(process.env.DADA_POSTV1_PACKAGE_ROOT ?? ".build/portable-release/Dada-P0A-0.0.0-win-x64");
|
||||||
|
const port = 43121;
|
||||||
|
|
||||||
|
async function waitForHealth(child) {
|
||||||
|
const deadline = Date.now() + 15_000;
|
||||||
|
while (Date.now() < deadline) {
|
||||||
|
if (child.exitCode !== null) throw new Error(`packaged api exited: ${child.exitCode}: ${child.errorOutput ?? ""}`);
|
||||||
|
try {
|
||||||
|
const response = await fetch(`http://127.0.0.1:${port}/healthz`, { headers: { host: `127.0.0.1:${port}` } });
|
||||||
|
if (response.ok) return;
|
||||||
|
} catch {}
|
||||||
|
await new Promise((resolveDelay) => setTimeout(resolveDelay, 100));
|
||||||
|
}
|
||||||
|
throw new Error("packaged api health timeout");
|
||||||
|
}
|
||||||
|
|
||||||
|
function startApi(configPath, dataRoot) {
|
||||||
|
const child = spawn(join(packageRoot, "runtime", "node.exe"), [join(packageRoot, "server", "api.mjs"), "--dada-credential-stdin"], {
|
||||||
|
cwd: packageRoot,
|
||||||
|
env: { ...process.env, DADA_INSTANCE_CONFIG_PATH: configPath, DADA_SUPPORT_GATE_ROOT: join(packageRoot, "web", "support-gate"), DADA_WEB_ROOT: join(packageRoot, "web") },
|
||||||
|
stdio: ["pipe", "ignore", "pipe"],
|
||||||
|
windowsHide: true,
|
||||||
|
});
|
||||||
|
child.errorOutput = "";
|
||||||
|
child.stderr.setEncoding("utf8");
|
||||||
|
child.stderr.on("data", (chunk) => { child.errorOutput += chunk; });
|
||||||
|
child.stdin.end(JSON.stringify({ "Dada/P0A/api/amap": "", "Dada/P0A/api/resend": "", "Dada/P0A/admin/pepper": "portable-test-pepper-00000000000000000000000000000000" }));
|
||||||
|
return child;
|
||||||
|
}
|
||||||
|
|
||||||
|
async function stop(child) {
|
||||||
|
if (child.exitCode === null) {
|
||||||
|
child.kill();
|
||||||
|
await new Promise((resolveExit) => child.once("exit", resolveExit));
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
async function verifyWorkerStartup(configPath) {
|
||||||
|
const pipeName = `Dada.P0A.PostV1.${process.pid}.${Date.now()}`;
|
||||||
|
const pipePath = `\\\\.\\pipe\\${pipeName}`;
|
||||||
|
const server = createServer();
|
||||||
|
await new Promise((resolveListen, rejectListen) => {
|
||||||
|
server.once("error", rejectListen);
|
||||||
|
server.listen(pipePath, resolveListen);
|
||||||
|
});
|
||||||
|
const child = spawn(join(packageRoot, "runtime", "node.exe"), [
|
||||||
|
join(packageRoot, "server", "worker.mjs"),
|
||||||
|
"--dada-control-pipe", pipeName,
|
||||||
|
"--dada-credential-stdin",
|
||||||
|
], {
|
||||||
|
cwd: packageRoot,
|
||||||
|
env: {
|
||||||
|
...process.env,
|
||||||
|
DADA_INSTANCE_CONFIG_PATH: configPath,
|
||||||
|
DADA_SQLITE_NATIVE_BINDING: join(packageRoot, "server", "native", "better_sqlite3.node"),
|
||||||
|
},
|
||||||
|
stdio: ["pipe", "ignore", "ignore"],
|
||||||
|
windowsHide: true,
|
||||||
|
});
|
||||||
|
child.stdin.end(JSON.stringify({ "Dada/P0A/worker/ai-gateway": "synthetic-runtime-token" }));
|
||||||
|
try {
|
||||||
|
await new Promise((resolveReady, rejectReady) => {
|
||||||
|
let settled = false;
|
||||||
|
const finish = (error) => {
|
||||||
|
if (settled) return;
|
||||||
|
settled = true;
|
||||||
|
clearTimeout(deadline);
|
||||||
|
child.off("exit", onExit);
|
||||||
|
if (error) rejectReady(error); else resolveReady();
|
||||||
|
};
|
||||||
|
const deadline = setTimeout(() => finish(new Error("packaged worker ready timeout")), 15_000);
|
||||||
|
const onExit = (code) => finish(new Error(`packaged worker exited before ready: ${code}`));
|
||||||
|
child.once("exit", onExit);
|
||||||
|
server.once("connection", (connection) => {
|
||||||
|
connection.setEncoding("utf8");
|
||||||
|
let pending = "";
|
||||||
|
connection.on("data", (chunk) => {
|
||||||
|
pending += chunk;
|
||||||
|
while (pending.includes("\n")) {
|
||||||
|
const newline = pending.indexOf("\n");
|
||||||
|
const status = pending.slice(0, newline).trim();
|
||||||
|
pending = pending.slice(newline + 1);
|
||||||
|
if (status === "storage_unavailable") finish(new Error("packaged worker reported storage_unavailable"));
|
||||||
|
if (status === "ready") {
|
||||||
|
setTimeout(() => {
|
||||||
|
if (settled) return;
|
||||||
|
connection.write("shutdown\n");
|
||||||
|
finish();
|
||||||
|
}, 500);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
});
|
||||||
|
});
|
||||||
|
});
|
||||||
|
const exitCode = await new Promise((resolveExit) => child.once("exit", resolveExit));
|
||||||
|
assert.equal(exitCode, 0);
|
||||||
|
} finally {
|
||||||
|
if (child.exitCode === null) child.kill();
|
||||||
|
await new Promise((resolveClose) => server.close(resolveClose));
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
test("portable package serves the product and keeps SQLite data across API restart", async () => {
|
||||||
|
assert.ok(existsSync(join(packageRoot, "Dada.exe")));
|
||||||
|
assert.ok(existsSync(join(packageRoot, "web", "index.html")));
|
||||||
|
const packagedWorker = await readFile(join(packageRoot, "server", "worker", "dist", "worker.js"), "utf8");
|
||||||
|
const packagedOneApiAdapter = await readFile(join(packageRoot, "server", "worker", "dist", "oneapi-generation-adapter.js"), "utf8");
|
||||||
|
assert.match(packagedWorker, /GenerationProcessor/);
|
||||||
|
assert.match(packagedWorker, /OneApiGenerationAdapter/);
|
||||||
|
assert.match(packagedOneApiAdapter, /oneapi\.intelligrow\.cn/);
|
||||||
|
assert.doesNotMatch(packagedWorker, /portable-mock-worker/);
|
||||||
|
const root = await mkdtemp(join(tmpdir(), "dada-postv1-"));
|
||||||
|
const dataRoot = join(root, "data");
|
||||||
|
const configPath = join(root, "instance.json");
|
||||||
|
await writeFile(configPath, JSON.stringify({ data_root: dataRoot, initialized: true, instance_id: "portable-test", schema_version: 1, secure_config_revision: 1, admin_allowlist_hashes: [], admin_recovery_hashes: [] }));
|
||||||
|
let api = startApi(configPath, dataRoot);
|
||||||
|
try {
|
||||||
|
await waitForHealth(api);
|
||||||
|
const initialPage = await fetch(`http://127.0.0.1:${port}/`, { headers: { host: `127.0.0.1:${port}` } });
|
||||||
|
assert.equal(initialPage.status, 200);
|
||||||
|
assert.match(await initialPage.text(), /当前浏览器无法使用 Dada/);
|
||||||
|
const support = await fetch(`http://127.0.0.1:${port}/api/v1/support/check`, {
|
||||||
|
method: "POST",
|
||||||
|
headers: { host: `127.0.0.1:${port}`, origin: `http://127.0.0.1:${port}`, "content-type": "application/json", "sec-ch-ua": '"Google Chrome";v="150"', "sec-ch-ua-full-version-list": '"Google Chrome";v="150.0.0.0"', "sec-ch-ua-platform": '"Windows"' },
|
||||||
|
body: JSON.stringify({ brands: [{ brand: "Google Chrome", version: "150" }], full_version_list: [{ brand: "Google Chrome", version: "150.0.0.0" }], platform: "Windows" }),
|
||||||
|
});
|
||||||
|
assert.equal(support.status, 200);
|
||||||
|
const cookie = support.headers.get("set-cookie")?.split(";", 1)[0];
|
||||||
|
assert.ok(cookie);
|
||||||
|
const page = await fetch(`http://127.0.0.1:${port}/app`, { headers: { host: `127.0.0.1:${port}`, cookie, "sec-ch-ua": '"Google Chrome";v="150"' } });
|
||||||
|
assert.equal(page.status, 200);
|
||||||
|
const pageHtml = await page.text();
|
||||||
|
assert.match(pageHtml, /<div id="root"><\/div>/);
|
||||||
|
const scriptPath = pageHtml.match(/<script[^>]+src="([^"]+)"/)?.[1];
|
||||||
|
const stylesheetPath = pageHtml.match(/<link[^>]+href="([^"]+)"/)?.[1];
|
||||||
|
assert.ok(scriptPath);
|
||||||
|
assert.ok(stylesheetPath);
|
||||||
|
const assetHeaders = { host: `127.0.0.1:${port}`, cookie, "sec-ch-ua": '"Google Chrome";v="150"' };
|
||||||
|
const script = await fetch(`http://127.0.0.1:${port}${scriptPath}`, { headers: assetHeaders });
|
||||||
|
const stylesheet = await fetch(`http://127.0.0.1:${port}${stylesheetPath}`, { headers: assetHeaders });
|
||||||
|
assert.equal(script.status, 200);
|
||||||
|
assert.match(script.headers.get("content-type") ?? "", /^(?:application|text)\/javascript\b/);
|
||||||
|
assert.equal(stylesheet.status, 200);
|
||||||
|
assert.match(stylesheet.headers.get("content-type") ?? "", /^text\/css\b/);
|
||||||
|
assert.ok(existsSync(join(dataRoot, "db", "dada.sqlite3")));
|
||||||
|
await verifyWorkerStartup(configPath);
|
||||||
|
} finally {
|
||||||
|
await stop(api);
|
||||||
|
}
|
||||||
|
api = startApi(configPath, dataRoot);
|
||||||
|
try {
|
||||||
|
await waitForHealth(api);
|
||||||
|
assert.ok(existsSync(join(dataRoot, "db", "dada.sqlite3")));
|
||||||
|
} finally {
|
||||||
|
await stop(api);
|
||||||
|
await rm(root, { recursive: true, force: true });
|
||||||
|
}
|
||||||
|
});
|
||||||
@@ -0,0 +1,85 @@
|
|||||||
|
import { describe, expect, it, vi } from "vitest";
|
||||||
|
|
||||||
|
import type { GenerationAdapterRequest } from "../../apps/worker/src/ai-adapter-contract.js";
|
||||||
|
import { runAiRuntimeProbe } from "../../apps/worker/src/ai-runtime-probe.js";
|
||||||
|
import { OneApiGenerationAdapter } from "../../apps/worker/src/oneapi-generation-adapter.js";
|
||||||
|
|
||||||
|
function request(overrides: Partial<GenerationAdapterRequest> = {}): GenerationAdapterRequest {
|
||||||
|
return {
|
||||||
|
configSnapshot: {},
|
||||||
|
generationId: "00000000-0000-4000-8000-000000000001",
|
||||||
|
modelId: "gemini-3.1-flash-image-preview",
|
||||||
|
prompt: "一张用于本机验收的抽象色彩图",
|
||||||
|
ratio: "1:1",
|
||||||
|
referenceAssetIds: [],
|
||||||
|
...overrides,
|
||||||
|
};
|
||||||
|
}
|
||||||
|
|
||||||
|
describe("POSTV1-02 OneAPI runtime adapter", () => {
|
||||||
|
it("uses the fixed Gemini gateway and normalizes one real-shaped response", async () => {
|
||||||
|
const source = Buffer.from(
|
||||||
|
"iVBORw0KGgoAAAANSUhEUgAAAAEAAAABCAQAAAC1HAwCAAAAC0lEQVR42mNk+A8AAQUBAScY42YAAAAASUVORK5CYII=",
|
||||||
|
"base64",
|
||||||
|
);
|
||||||
|
const gateway = vi.fn(async (_url: string | URL | Request, init?: RequestInit) => {
|
||||||
|
const headers = new Headers(init?.headers);
|
||||||
|
expect(headers.get("authorization")).toBe("Bearer synthetic-runtime-token");
|
||||||
|
expect(init?.redirect).toBe("error");
|
||||||
|
return new Response(JSON.stringify({
|
||||||
|
choices: [{ message: { content: `})` } }],
|
||||||
|
}), { headers: { "content-type": "application/json" }, status: 200 });
|
||||||
|
});
|
||||||
|
const credential = Buffer.from("synthetic-runtime-token");
|
||||||
|
const adapter = new OneApiGenerationAdapter({ credential, fetch: gateway as typeof fetch });
|
||||||
|
const result = await adapter.start(request());
|
||||||
|
|
||||||
|
expect(gateway).toHaveBeenCalledOnce();
|
||||||
|
expect(gateway.mock.calls[0]?.[0]).toBe("https://oneapi.intelligrow.cn/v1/chat/completions");
|
||||||
|
expect(result.status === "failed" ? result.sourceCategory : "completed").toBe("completed");
|
||||||
|
if (result.status === "completed") {
|
||||||
|
expect(result.outputs).toHaveLength(1);
|
||||||
|
expect(result.outputs[0]).toMatchObject({ mimeType: "image/png", pixelHeight: 1080, pixelWidth: 1080 });
|
||||||
|
expect(result.outputs[0]?.bytes.length).toBeGreaterThan(0);
|
||||||
|
}
|
||||||
|
expect(JSON.stringify(result)).not.toContain("synthetic-runtime-token");
|
||||||
|
adapter.dispose();
|
||||||
|
credential.fill(0);
|
||||||
|
});
|
||||||
|
|
||||||
|
it("fails closed instead of returning a mock image", async () => {
|
||||||
|
const adapter = new OneApiGenerationAdapter({
|
||||||
|
credential: Buffer.from("synthetic-runtime-token"),
|
||||||
|
fetch: vi.fn(async () => new Response(null, { status: 503 })) as typeof fetch,
|
||||||
|
});
|
||||||
|
await expect(adapter.start(request())).resolves.toEqual({
|
||||||
|
category: "upstream_failed",
|
||||||
|
sourceCategory: "upstream_http_503",
|
||||||
|
status: "failed",
|
||||||
|
});
|
||||||
|
adapter.dispose();
|
||||||
|
await expect(adapter.start(request())).resolves.toEqual({
|
||||||
|
category: "upstream_failed",
|
||||||
|
sourceCategory: "adapter_disposed",
|
||||||
|
status: "failed",
|
||||||
|
});
|
||||||
|
});
|
||||||
|
|
||||||
|
it("returns only a bounded probe summary and wipes generated bytes", async () => {
|
||||||
|
const bytes = Buffer.from("probe-output");
|
||||||
|
const result = await runAiRuntimeProbe({
|
||||||
|
async start() {
|
||||||
|
return { outputs: [{ bytes, mimeType: "image/png", pixelHeight: 1080, pixelWidth: 1080 }], status: "completed" };
|
||||||
|
},
|
||||||
|
});
|
||||||
|
expect(result).toEqual({
|
||||||
|
code: "ai_probe_passed",
|
||||||
|
mime_type: "image/png",
|
||||||
|
pixel_height: 1080,
|
||||||
|
pixel_width: 1080,
|
||||||
|
real_calls: 1,
|
||||||
|
success: true,
|
||||||
|
});
|
||||||
|
expect(bytes.every((value) => value === 0)).toBe(true);
|
||||||
|
});
|
||||||
|
});
|
||||||
Reference in New Issue
Block a user