Files
ReinLoopTest/server/src/postgres-store.js
T
2026-07-31 10:20:33 +08:00

116 lines
11 KiB
JavaScript

const path = require("node:path");
const fs = require("node:fs");
const { randomUUID } = require("node:crypto");
const { Pool } = require("pg");
const { EMPTY_DATABASE } = require("./store");
function asIso(value) {
return value instanceof Date ? value.toISOString() : value;
}
function formatShanghaiTimestamp(value) {
return new Intl.DateTimeFormat("sv-SE", {
timeZone: "Asia/Shanghai", year: "numeric", month: "2-digit", day: "2-digit",
hour: "2-digit", minute: "2-digit", hourCycle: "h23"
}).format(new Date(value));
}
class PostgresStore {
constructor(connectionString, filesDirectory) {
this.pool = new Pool({ connectionString });
this.filesDirectory = filesDirectory;
this.modelsDirectory = path.join(path.dirname(filesDirectory), "models");
}
async initialize() {
const migrationPath = path.join(__dirname, "..", "migrations", "001_normalized_schema.sql");
await this.pool.query(await fs.promises.readFile(migrationPath, "utf8"));
}
async read() {
return this.readDatabase(this.pool);
}
async readDatabase(queryable) {
const companies = await queryable.query("SELECT id, name, code, created_at FROM companies");
const productionLines = await queryable.query("SELECT id, company_id, name, code, device_id, created_at, last_seen_at FROM production_lines");
const licenses = await queryable.query("SELECT licenses.license_id, licenses.company_id, licenses.production_line_id, licenses.device_id, licenses.customer, licenses.issued_at, licenses.expiry_at, licenses.features, licenses.license, licenses.status, licenses.created_at, licenses.revoked_at, licenses.revocation_reason, companies.name AS company_name, production_lines.name AS production_line_name FROM licenses JOIN companies ON companies.id = licenses.company_id JOIN production_lines ON production_lines.id = licenses.production_line_id");
const fileRecords = await queryable.query("SELECT id, file_name, original_file_name, folder, cloud_path, file_id, upload_time, size_bytes FROM file_records");
const functionConfigs = await queryable.query("SELECT config_type, parameters, version, update_time FROM function_configs");
const panelInbox = await queryable.query("SELECT file_id, device_id, file_name, media_type, upload_time FROM panel_inbox");
const identificationFiles = await queryable.query("SELECT file_id, device_id, file_name, media_type, upload_time, size_bytes, status, processed_at, expires_at FROM identification_files");
const feedback = await queryable.query("SELECT device_id, run_id, file_name, status, result, update_time FROM identification_feedback");
const requests = await queryable.query("SELECT device_id, request_id, status, created_at_ms, expires_at_ms, config_file_id, config_file_name, uploaded_at_ms, update_time FROM volume_config_requests");
const volumeConfigs = await queryable.query("SELECT device_id, file_id, file_name, request_id, updated_at_ms FROM volume_configs");
const panelNotifications = await queryable.query("SELECT notification_id, device_id, type, title, message, request_id, file_id, created_at_ms FROM panel_notifications");
return {
...structuredClone(EMPTY_DATABASE),
companies: companies.rows.map((row) => ({ id: row.id, name: row.name, code: row.code, createdAt: asIso(row.created_at) })),
productionLines: productionLines.rows.map((row) => ({ id: row.id, companyId: row.company_id, name: row.name, code: row.code, deviceId: row.device_id, createdAt: asIso(row.created_at), lastSeenAt: row.last_seen_at && asIso(row.last_seen_at) })),
licenses: licenses.rows.map((row) => ({ licenseId: row.license_id, companyId: row.company_id, productionLineId: row.production_line_id, companyName: row.company_name, productionLineName: row.production_line_name, deviceId: row.device_id, customer: row.customer, issued: formatShanghaiTimestamp(row.issued_at), expiry: formatShanghaiTimestamp(row.expiry_at), issuedAt: asIso(row.issued_at), expiryAt: asIso(row.expiry_at), features: row.features, license: row.license, status: row.status, createdAt: asIso(row.created_at), revokedAt: row.revoked_at && asIso(row.revoked_at), revocationReason: row.revocation_reason })),
fileRecords: fileRecords.rows.map((row) => ({ _id: row.id, fileName: row.file_name, originalFileName: row.original_file_name || undefined, folder: row.folder, cloudPath: row.cloud_path, fileID: row.file_id, uploadTime: asIso(row.upload_time), size: Number(row.size_bytes) })),
functionConfigs: functionConfigs.rows.map((row) => ({ configType: row.config_type, parameters: row.parameters, version: Number(row.version), updateTime: asIso(row.update_time) })),
panelInbox: panelInbox.rows.map((row) => ({ fileID: row.file_id, deviceId: row.device_id, fileName: row.file_name, mediaType: row.media_type, uploadTime: asIso(row.upload_time) })),
identificationFiles: identificationFiles.rows.map((row) => ({ fileID: row.file_id, deviceId: row.device_id, fileName: row.file_name, mediaType: row.media_type, uploadTime: asIso(row.upload_time), size: Number(row.size_bytes), status: row.status, processedAt: row.processed_at && asIso(row.processed_at), expiresAt: row.expires_at && asIso(row.expires_at) })),
identificationFeedback: feedback.rows.map((row) => ({ deviceId: row.device_id, runId: row.run_id, fileName: row.file_name, status: row.status, result: row.result, updateTime: asIso(row.update_time) })),
volumeConfigRequests: requests.rows.map((row) => ({ deviceId: row.device_id, requestId: row.request_id, status: row.status, createdAtMs: Number(row.created_at_ms), expiresAtMs: Number(row.expires_at_ms), configFileID: row.config_file_id, configFileName: row.config_file_name, uploadedAtMs: row.uploaded_at_ms && Number(row.uploaded_at_ms), updateTime: asIso(row.update_time) })),
volumeConfigs: volumeConfigs.rows.map((row) => ({ deviceId: row.device_id, fileID: row.file_id, fileName: row.file_name, requestId: row.request_id, updatedAtMs: Number(row.updated_at_ms) })),
panelNotifications: panelNotifications.rows.map((row) => ({ notificationId: row.notification_id, deviceId: row.device_id, type: row.type, title: row.title, message: row.message, requestId: row.request_id, fileID: row.file_id, createdAtMs: Number(row.created_at_ms) }))
};
}
async update(mutator) {
const client = await this.pool.connect();
try {
await client.query("BEGIN");
await client.query("SELECT pg_advisory_xact_lock(81720260725)");
const database = await this.readDatabase(client);
const response = await mutator(database);
await this.writeDatabase(client, database);
await client.query("COMMIT");
return response;
} catch (error) {
await client.query("ROLLBACK");
throw error;
} finally {
client.release();
}
}
async writeDatabase(client, database) {
await client.query("DELETE FROM panel_notifications; DELETE FROM volume_configs; DELETE FROM volume_config_requests; DELETE FROM identification_feedback; DELETE FROM panel_inbox; DELETE FROM identification_files; DELETE FROM function_configs; DELETE FROM file_records; DELETE FROM licenses; DELETE FROM production_lines; DELETE FROM companies;");
for (const item of database.companies) await client.query("INSERT INTO companies (id, name, code, created_at) VALUES ($1, $2, $3, $4)", [item.id, item.name, item.code, item.createdAt]);
for (const item of database.productionLines) await client.query("INSERT INTO production_lines (id, company_id, name, code, device_id, created_at, last_seen_at) VALUES ($1, $2, $3, $4, $5, $6, $7)", [item.id, item.companyId, item.name, item.code, item.deviceId, item.createdAt, item.lastSeenAt]);
for (const item of database.licenses) await client.query("INSERT INTO licenses (license_id, company_id, production_line_id, device_id, customer, issued_at, expiry_at, features, license, status, created_at, revoked_at, revocation_reason) VALUES ($1,$2,$3,$4,$5,$6,$7,$8,$9,$10,$11,$12,$13)", [item.licenseId, item.companyId, item.productionLineId, item.deviceId, item.customer, item.issuedAt, item.expiryAt, item.features, item.license, item.status, item.createdAt, item.revokedAt, item.revocationReason]);
for (const item of database.fileRecords) await client.query("INSERT INTO file_records (id, file_name, original_file_name, folder, cloud_path, file_id, upload_time, size_bytes) VALUES ($1,$2,$3,$4,$5,$6,$7,$8)", [item._id, item.fileName, item.originalFileName || null, item.folder, item.cloudPath, item.fileID, item.uploadTime, item.size]);
for (const item of database.functionConfigs) await client.query("INSERT INTO function_configs (config_type, parameters, version, update_time) VALUES ($1,$2,$3,$4)", [item.configType, item.parameters, item.version, item.updateTime]);
for (const item of database.panelInbox) await client.query("INSERT INTO panel_inbox (file_id, device_id, file_name, media_type, upload_time) VALUES ($1,$2,$3,$4,$5)", [item.fileID, item.deviceId, item.fileName, item.mediaType, item.uploadTime]);
for (const item of database.identificationFiles) await client.query("INSERT INTO identification_files (file_id, device_id, file_name, media_type, upload_time, size_bytes, status, processed_at, expires_at) VALUES ($1,$2,$3,$4,$5,$6,$7,$8,$9)", [item.fileID, item.deviceId, item.fileName, item.mediaType, item.uploadTime, item.size, item.status, item.processedAt, item.expiresAt]);
for (const item of database.identificationFeedback) await client.query("INSERT INTO identification_feedback (device_id, run_id, file_name, status, result, update_time) VALUES ($1,$2,$3,$4,$5,$6)", [item.deviceId, item.runId, item.fileName, item.status, item.result, item.updateTime]);
for (const item of database.volumeConfigRequests) await client.query("INSERT INTO volume_config_requests (device_id, request_id, status, created_at_ms, expires_at_ms, config_file_id, config_file_name, uploaded_at_ms, update_time) VALUES ($1,$2,$3,$4,$5,$6,$7,$8,$9)", [item.deviceId, item.requestId, item.status, item.createdAtMs, item.expiresAtMs, item.configFileID, item.configFileName, item.uploadedAtMs, item.updateTime]);
for (const item of database.volumeConfigs) await client.query("INSERT INTO volume_configs (device_id, file_id, file_name, request_id, updated_at_ms) VALUES ($1,$2,$3,$4,$5)", [item.deviceId, item.fileID, item.fileName, item.requestId, item.updatedAtMs]);
for (const item of database.panelNotifications) await client.query("INSERT INTO panel_notifications (notification_id, device_id, type, title, message, request_id, file_id, created_at_ms) VALUES ($1,$2,$3,$4,$5,$6,$7,$8)", [item.notificationId, item.deviceId, item.type, item.title, item.message, item.requestId, item.fileID, item.createdAtMs]);
}
createId(prefix) {
return `${prefix}_${randomUUID()}`;
}
resolveStoredFile(fileID) {
if (typeof fileID !== "string") return null;
const isModel = fileID.startsWith("model://");
if (!isModel && !fileID.startsWith("local://")) return null;
const rootDirectory = isModel ? this.modelsDirectory : this.filesDirectory;
const relativePath = fileID.slice(isModel ? "model://".length : "local://".length).replace(/\\/g, "/");
const absolutePath = path.resolve(rootDirectory, relativePath);
const relativeToRoot = path.relative(rootDirectory, absolutePath);
if (relativeToRoot.startsWith("..") || path.isAbsolute(relativeToRoot)) return null;
return absolutePath;
}
async close() {
await this.pool.end();
}
}
module.exports = { PostgresStore };