693 lines
30 KiB
TypeScript
693 lines
30 KiB
TypeScript
import { lstatSync, realpathSync, readFileSync, unlinkSync } from "node:fs";
|
|
import { chmod, mkdir, rename, writeFile } from "node:fs/promises";
|
|
import path from "node:path";
|
|
import { createPublicKey, randomUUID, verify as verifySignature } from "node:crypto";
|
|
import type Database from "better-sqlite3";
|
|
import { writeAudit } from "./audit.js";
|
|
import { AppError } from "./errors.js";
|
|
import type { AppConfig } from "./config.js";
|
|
import {
|
|
detectPlatform,
|
|
fetchReleaseBytes,
|
|
fetchReleaseMetadata,
|
|
fetchReleaseText,
|
|
isNewerVersion,
|
|
parseSemver,
|
|
runtimeHashFromLockfile,
|
|
sanitizeAssetName,
|
|
selectReleaseAsset,
|
|
validateHttpsUrl,
|
|
RELEASE_NOTES_MAX_BYTES,
|
|
downloadReleaseAsset,
|
|
extractSafeArchive,
|
|
normalizeReleasePermissions,
|
|
applicationUpdateRuntimeHash,
|
|
type ReleaseAsset,
|
|
type ReleaseMetadata,
|
|
} from "./update.js";
|
|
import type { UpdateJobStatus } from "../shared/contracts.js";
|
|
|
|
|
|
export const activeInProcessDownloads = new Map<string, AbortController>();
|
|
|
|
export function triggerInProcessDownload(
|
|
database: Database.Database,
|
|
config: AppConfig,
|
|
jobId: string,
|
|
asset: { name: string; url: string; sha256?: string },
|
|
expectedSha256?: string,
|
|
): void {
|
|
setImmediate(async () => {
|
|
try {
|
|
const row = database.prepare("SELECT id, status, operation FROM update_jobs WHERE id=?").get(jobId) as { id: string; status: string; operation: string } | undefined;
|
|
if (!row || row.status !== "queued") return;
|
|
|
|
const controller = new AbortController();
|
|
activeInProcessDownloads.set(jobId, controller);
|
|
|
|
const workspace = path.join(path.resolve(config.stagingDir), `update-${jobId}`);
|
|
const archivePath = path.join(workspace, asset.name.endsWith(".gz") || asset.name.endsWith(".zip") ? asset.name : `${asset.name}.tar.gz`);
|
|
|
|
await mkdir(workspace, { recursive: true, mode: 0o700 });
|
|
const now = Date.now();
|
|
database.prepare("UPDATE update_jobs SET status='downloading', download_started_at=?, started_at=?, download_path=?, updated_at=? WHERE id=? AND status='queued'").run(now, now, path.basename(archivePath), now, jobId);
|
|
|
|
const progressStartedAt = Date.now();
|
|
let lastProgressWrite = 0;
|
|
|
|
const downloaded = await downloadReleaseAsset(asset.url, archivePath, {
|
|
allowedHosts: config.updateAllowedHosts,
|
|
maxBytes: config.updateMaxBytes,
|
|
fetchImpl: (input, init) => fetch(input, { ...init, signal: controller.signal }),
|
|
onProgress: (downloadedBytes, totalBytes) => {
|
|
const cur = Date.now();
|
|
if (cur - lastProgressWrite < 200) return;
|
|
lastProgressWrite = cur;
|
|
const elapsed = Math.max(1, cur - progressStartedAt);
|
|
const speedBps = Math.round(downloadedBytes * 1000 / elapsed);
|
|
try {
|
|
database.prepare("UPDATE update_jobs SET downloaded_bytes=?, size_bytes=COALESCE(?, size_bytes), download_speed_bps=?, updated_at=? WHERE id=? AND status='downloading'").run(downloadedBytes, totalBytes, speedBps, cur, jobId);
|
|
} catch {}
|
|
},
|
|
});
|
|
|
|
if (expectedSha256 && downloaded.sha256.toLowerCase() !== expectedSha256.toLowerCase()) {
|
|
throw new Error("更新文件 SHA-256 校验失败");
|
|
}
|
|
|
|
database.prepare("UPDATE update_jobs SET status='verifying', actual_sha256=?, size_bytes=?, downloaded_bytes=?, updated_at=? WHERE id=? AND status='downloading'").run(downloaded.sha256, downloaded.size, downloaded.size, Date.now(), jobId);
|
|
|
|
const stagedDir = path.join(workspace, "payload");
|
|
await extractSafeArchive(archivePath, stagedDir, config.updateMaxBytes === undefined ? {} : { maxBytes: config.updateMaxBytes });
|
|
|
|
if (applicationUpdateRuntimeHash(asset.name)) {
|
|
try {
|
|
const currentRelease = realpathSync(config.currentLink);
|
|
if (currentRelease) {
|
|
const fsPromises = await import("node:fs/promises");
|
|
for (const entry of ["node_modules", "runtime", "pnpm-lock.yaml"] as const) {
|
|
const source = path.join(currentRelease, entry);
|
|
const sourceInfo = await fsPromises.lstat(source).catch(() => null);
|
|
if (sourceInfo && !sourceInfo.isSymbolicLink()) {
|
|
await fsPromises.cp(source, path.join(stagedDir, entry), { recursive: sourceInfo.isDirectory(), errorOnExist: true, force: false }).catch(() => {});
|
|
}
|
|
}
|
|
}
|
|
} catch {}
|
|
}
|
|
|
|
await normalizeReleasePermissions(stagedDir).catch(() => {});
|
|
const fsPromises = await import("node:fs/promises");
|
|
const payloadInfo = await fsPromises.lstat(path.join(stagedDir, "dist")).catch(() => null);
|
|
if (!payloadInfo?.isDirectory() || payloadInfo.isSymbolicLink()) {
|
|
throw new Error("发布包缺少 dist 目录");
|
|
}
|
|
|
|
database.prepare("UPDATE update_jobs SET status='staged', actual_sha256=?, size_bytes=?, download_path=?, updated_at=? WHERE id=? AND status IN ('verifying', 'downloading')").run(downloaded.sha256, downloaded.size, workspace, Date.now(), jobId);
|
|
} catch (error) {
|
|
const controller = activeInProcessDownloads.get(jobId);
|
|
if (controller?.signal.aborted) return;
|
|
const rawMsg = error instanceof Error ? error.message : "更新文件下载失败";
|
|
try {
|
|
database.prepare("UPDATE update_jobs SET status='failed', error_message=?, updated_at=? WHERE id=? AND status NOT IN ('completed', 'staged', 'cancelled')").run(rawMsg, Date.now(), jobId);
|
|
} catch {}
|
|
const workspace = path.join(path.resolve(config.stagingDir), `update-${jobId}`);
|
|
const fsPromises = await import("node:fs/promises");
|
|
await fsPromises.rm(workspace, { recursive: true, force: true }).catch(() => {});
|
|
} finally {
|
|
activeInProcessDownloads.delete(jobId);
|
|
}
|
|
});
|
|
}
|
|
|
|
export const UPDATE_CACHE_KEY = "update.release.v1";
|
|
export const ACTIVE_UPDATE_STATUSES: readonly UpdateJobStatus[] = [
|
|
"queued",
|
|
"downloading",
|
|
"verifying",
|
|
"staged",
|
|
"backing_up",
|
|
"applying",
|
|
];
|
|
|
|
// A queued job normally starts within seconds and an applying job completes
|
|
// after the service health check. The runner refreshes its recovery marker as
|
|
// a lease while doing long downloads/backups; only an expired lease permits
|
|
// the server to reclaim an active row.
|
|
export const ORPHANED_UPDATE_TIMEOUT_MS = 5 * 60 * 1000;
|
|
export const QUEUED_UPDATE_TIMEOUT_MS = 25 * 1000;
|
|
|
|
export type CachedRelease = {
|
|
checkedAt: number;
|
|
metadataUrl: string;
|
|
version: string;
|
|
tagName?: string;
|
|
releaseName?: string;
|
|
publishedAt?: string;
|
|
notes?: string;
|
|
releaseUrl?: string;
|
|
platform: string;
|
|
signatureVerified?: boolean;
|
|
asset?: {
|
|
name: string;
|
|
url: string;
|
|
size?: number;
|
|
sha256?: string;
|
|
};
|
|
};
|
|
|
|
export type UpdateCheckResult = {
|
|
configured: boolean;
|
|
currentVersion: string;
|
|
platform: ReturnType<typeof detectPlatform>;
|
|
checkedAt: number;
|
|
latest: {
|
|
version: string;
|
|
tagName?: string;
|
|
releaseName?: string;
|
|
publishedAt?: string;
|
|
notes?: string;
|
|
releaseUrl?: string;
|
|
compatible: boolean;
|
|
integrityReady: boolean;
|
|
signatureReady: boolean;
|
|
isNewer: boolean;
|
|
assetName?: string;
|
|
assetSize?: number;
|
|
} | null;
|
|
};
|
|
|
|
export type UpdateRequest = {
|
|
jobId: string;
|
|
operation?: "download" | "apply";
|
|
version: string;
|
|
metadataUrl: string;
|
|
assetUrl: string;
|
|
assetName: string;
|
|
expectedSha256: string;
|
|
requestedAt: number;
|
|
// These paths are derived from the server config and are included so the
|
|
// privileged runner does not need to infer a working directory from input.
|
|
currentLink: string;
|
|
releasesDir: string;
|
|
dataDir: string;
|
|
stagedPath?: string;
|
|
};
|
|
|
|
function setting(database: Database.Database, key: string): string | undefined {
|
|
return (database.prepare("SELECT value FROM system_settings WHERE key=?").get(key) as { value: string } | undefined)?.value;
|
|
}
|
|
|
|
function saveSetting(database: Database.Database, key: string, value: unknown): void {
|
|
database.prepare(`
|
|
INSERT INTO system_settings(key, value, updated_at) VALUES (?, ?, ?)
|
|
ON CONFLICT(key) DO UPDATE SET value=excluded.value, updated_at=excluded.updated_at
|
|
`).run(key, JSON.stringify(value), Date.now());
|
|
}
|
|
|
|
function sha256FromSums(text: string, assetName: string): string | undefined {
|
|
const wanted = sanitizeAssetName(assetName);
|
|
for (const line of text.split(/\r?\n/)) {
|
|
const match = /^\s*([a-f0-9]{64})\s+[* ]?(.+?)\s*$/.exec(line);
|
|
if (!match) continue;
|
|
const name = match[2]!.replaceAll("\\", "/").split("/").pop() ?? "";
|
|
if (name === wanted) return match[1]!.toLowerCase();
|
|
}
|
|
return undefined;
|
|
}
|
|
|
|
/** Verify an Ed25519 detached signature over the exact SHA256SUMS bytes.
|
|
* The signature sidecar is accepted as either base64 or a 64-byte hex value.
|
|
*/
|
|
export function verifyReleaseSignature(payload: string, encodedSignature: string | Uint8Array, publicKey: string): boolean {
|
|
try {
|
|
const signature = (() => {
|
|
if (encodedSignature instanceof Uint8Array) {
|
|
const bytes = Buffer.from(encodedSignature);
|
|
if (bytes.length === 64) return bytes;
|
|
encodedSignature = bytes.toString("utf8");
|
|
}
|
|
const compact = encodedSignature.trim().replace(/\s+/g, "");
|
|
return /^[a-f0-9]{128}$/i.test(compact)
|
|
? Buffer.from(compact, "hex")
|
|
: Buffer.from(compact, "base64");
|
|
})();
|
|
if (signature.length !== 64) return false;
|
|
return verifySignature(null, Buffer.from(payload, "utf8"), createPublicKey(publicKey), signature);
|
|
} catch {
|
|
return false;
|
|
}
|
|
}
|
|
|
|
function signatureAssetFor(metadata: ReleaseMetadata, sums: ReleaseAsset): ReleaseAsset | undefined {
|
|
const sumsName = sums.name.toLowerCase();
|
|
return metadata.assets.find((candidate) => {
|
|
const name = candidate.name.toLowerCase();
|
|
return name === `${sumsName}.sig` || name === `${sumsName}.asc`;
|
|
});
|
|
}
|
|
|
|
export async function attachSidecarHash(
|
|
metadata: ReleaseMetadata,
|
|
asset: ReleaseAsset,
|
|
options: { allowedHosts: readonly string[]; baseUrl: string; maxBytes: number; publicKey?: string | undefined; requireSignature?: boolean | undefined },
|
|
): Promise<{ asset: ReleaseAsset; signatureVerified: boolean }> {
|
|
let signatureVerified = false;
|
|
if (asset.sha256 && (!options.publicKey || !options.requireSignature)) return { asset, signatureVerified };
|
|
const sums = metadata.assets.find((candidate) => /^(?:sha256sums?|checksums?)(?:\.txt)?$/i.test(path.basename(candidate.name)));
|
|
if (!sums) return { asset, signatureVerified };
|
|
try {
|
|
const content = await fetchReleaseText(sums.url, { allowedHosts: options.allowedHosts, baseUrl: options.baseUrl, maxBytes: Math.min(options.maxBytes, 2 * 1024 * 1024) });
|
|
const sha256 = sha256FromSums(content, asset.name);
|
|
if (options.publicKey) {
|
|
const signatureAsset = signatureAssetFor(metadata, sums);
|
|
if (signatureAsset) {
|
|
const signature = await fetchReleaseBytes(signatureAsset.url, { allowedHosts: options.allowedHosts, baseUrl: options.baseUrl, maxBytes: 64 * 1024 });
|
|
signatureVerified = verifyReleaseSignature(content, signature, options.publicKey);
|
|
}
|
|
}
|
|
return { asset: sha256 ? { ...asset, sha256 } : asset, signatureVerified };
|
|
} catch {
|
|
// A missing/unreadable sidecar makes the update unavailable; it must not
|
|
// turn into an unverified download.
|
|
return { asset, signatureVerified };
|
|
}
|
|
}
|
|
|
|
function policy(config: AppConfig) {
|
|
return {
|
|
allowedHosts: config.updateAllowedHosts,
|
|
baseUrl: config.updateMetadataUrl,
|
|
maxRedirects: 3,
|
|
} as const;
|
|
}
|
|
|
|
function safeMetadataUrl(config: AppConfig): string {
|
|
try {
|
|
return validateHttpsUrl(config.updateMetadataUrl, policy(config)).toString();
|
|
} catch {
|
|
throw new AppError(503, "UPDATE_NOT_CONFIGURED", "更新源地址配置无效");
|
|
}
|
|
}
|
|
|
|
export async function checkForUpdate(database: Database.Database, config: AppConfig): Promise<UpdateCheckResult> {
|
|
const platform = detectPlatform();
|
|
const checkedAt = Date.now();
|
|
if (config.updateStrategy === "disabled" || !config.updateMetadataUrl) {
|
|
return { configured: false, currentVersion: config.appVersion, platform, checkedAt, latest: null };
|
|
}
|
|
const metadataUrl = safeMetadataUrl(config);
|
|
let metadata: ReleaseMetadata;
|
|
try {
|
|
metadata = await fetchReleaseMetadata(metadataUrl, policy(config));
|
|
} catch {
|
|
throw new AppError(502, "UPDATE_CHECK_FAILED", "暂时无法获取最新版本,请稍后重试");
|
|
}
|
|
let runtimeHash: string | undefined;
|
|
try {
|
|
runtimeHash = runtimeHashFromLockfile(readFileSync(path.join(config.projectRoot, "pnpm-lock.yaml")));
|
|
} catch {
|
|
// Legacy or source installations may not contain the lockfile. They stay
|
|
// on the full release asset instead of risking an incompatible runtime.
|
|
}
|
|
let asset = selectReleaseAsset(metadata, platform, runtimeHash);
|
|
let signatureVerified = false;
|
|
if (asset) {
|
|
const integrity = await attachSidecarHash(metadata, asset, {
|
|
allowedHosts: config.updateAllowedHosts,
|
|
baseUrl: metadataUrl,
|
|
maxBytes: config.updateMaxBytes,
|
|
publicKey: config.updatePublicKey,
|
|
requireSignature: config.updateRequireSignature,
|
|
});
|
|
asset = integrity.asset;
|
|
signatureVerified = integrity.signatureVerified;
|
|
}
|
|
const safeVersion = metadata.version;
|
|
const cached: CachedRelease = {
|
|
checkedAt,
|
|
metadataUrl,
|
|
version: safeVersion,
|
|
...(metadata.tagName ? { tagName: metadata.tagName } : {}),
|
|
...(metadata.releaseName ? { releaseName: metadata.releaseName } : {}),
|
|
...(metadata.publishedAt ? { publishedAt: metadata.publishedAt } : {}),
|
|
...(metadata.notes ? { notes: metadata.notes } : {}),
|
|
...(metadata.releaseUrl ? { releaseUrl: metadata.releaseUrl } : {}),
|
|
platform: platform.target,
|
|
signatureVerified,
|
|
...(asset ? {
|
|
asset: {
|
|
name: sanitizeAssetName(asset.name),
|
|
url: validateHttpsUrl(asset.url, policy(config)).toString(),
|
|
...(asset.size === undefined ? {} : { size: asset.size }),
|
|
...(asset.sha256 ? { sha256: asset.sha256 } : {}),
|
|
},
|
|
} : {}),
|
|
};
|
|
saveSetting(database, UPDATE_CACHE_KEY, cached);
|
|
return {
|
|
configured: true,
|
|
currentVersion: config.appVersion,
|
|
platform,
|
|
checkedAt,
|
|
latest: {
|
|
version: safeVersion,
|
|
...(metadata.tagName ? { tagName: metadata.tagName } : {}),
|
|
...(metadata.releaseName ? { releaseName: metadata.releaseName } : {}),
|
|
...(metadata.publishedAt ? { publishedAt: metadata.publishedAt } : {}),
|
|
...(metadata.notes ? { notes: metadata.notes } : {}),
|
|
...(metadata.releaseUrl ? { releaseUrl: metadata.releaseUrl } : {}),
|
|
compatible: Boolean(asset),
|
|
integrityReady: Boolean(asset?.sha256 && (!config.updateRequireSignature || signatureVerified)),
|
|
signatureReady: !config.updateRequireSignature || signatureVerified,
|
|
isNewer: isNewerVersion(config.appVersion, safeVersion),
|
|
...(asset ? { assetName: asset.name, ...(asset.size === undefined ? {} : { assetSize: asset.size }) } : {}),
|
|
},
|
|
};
|
|
}
|
|
|
|
export function readCachedRelease(database: Database.Database, config: AppConfig): CachedRelease | null {
|
|
const raw = setting(database, UPDATE_CACHE_KEY);
|
|
if (!raw) return null;
|
|
try {
|
|
const value = JSON.parse(raw) as CachedRelease;
|
|
if (!value || typeof value !== "object" || typeof value.version !== "string" || typeof value.metadataUrl !== "string" || typeof value.platform !== "string") return null;
|
|
parseSemver(value.version);
|
|
const metadataUrl = validateHttpsUrl(value.metadataUrl, policy(config)).toString();
|
|
if (value.releaseName !== undefined && (typeof value.releaseName !== "string" || value.releaseName.length > 200 || /[\u0000-\u001f\u007f]/.test(value.releaseName))) return null;
|
|
if (value.notes !== undefined && (typeof value.notes !== "string" || Buffer.byteLength(value.notes, "utf8") > RELEASE_NOTES_MAX_BYTES)) return null;
|
|
if (value.releaseUrl !== undefined) validateHttpsUrl(value.releaseUrl, policy(config));
|
|
if (value.signatureVerified !== undefined && typeof value.signatureVerified !== "boolean") return null;
|
|
if (value.asset) {
|
|
if (typeof value.asset.name !== "string" || typeof value.asset.url !== "string") return null;
|
|
sanitizeAssetName(value.asset.name);
|
|
validateHttpsUrl(value.asset.url, policy(config));
|
|
if (value.asset.sha256 !== undefined && !/^[a-f0-9]{64}$/i.test(value.asset.sha256)) return null;
|
|
}
|
|
return { ...value, metadataUrl };
|
|
} catch {
|
|
return null;
|
|
}
|
|
}
|
|
|
|
export function publicCheckFromCache(database: Database.Database, config: AppConfig): UpdateCheckResult {
|
|
const platform = detectPlatform();
|
|
const cached = readCachedRelease(database, config);
|
|
if (!cached || cached.platform !== platform.target) {
|
|
const compatible = Boolean(cached && cached.platform === platform.target && cached.asset);
|
|
return { configured: config.updateStrategy !== "disabled", currentVersion: config.appVersion, platform, checkedAt: cached?.checkedAt ?? 0, latest: cached ? {
|
|
version: cached.version,
|
|
...(cached.tagName ? { tagName: cached.tagName } : {}),
|
|
...(cached.releaseName ? { releaseName: cached.releaseName } : {}),
|
|
...(cached.publishedAt ? { publishedAt: cached.publishedAt } : {}),
|
|
...(cached.notes ? { notes: cached.notes } : {}),
|
|
...(cached.releaseUrl ? { releaseUrl: cached.releaseUrl } : {}),
|
|
compatible,
|
|
integrityReady: compatible && Boolean(cached.asset?.sha256) && (!config.updateRequireSignature || cached.signatureVerified === true),
|
|
signatureReady: !config.updateRequireSignature || cached.signatureVerified === true,
|
|
isNewer: isNewerVersion(config.appVersion, cached.version),
|
|
...(cached.asset ? { assetName: cached.asset.name, ...(cached.asset.size === undefined ? {} : { assetSize: cached.asset.size }) } : {}),
|
|
} : null };
|
|
}
|
|
return {
|
|
configured: config.updateStrategy !== "disabled",
|
|
currentVersion: config.appVersion,
|
|
platform,
|
|
checkedAt: cached.checkedAt,
|
|
latest: {
|
|
version: cached.version,
|
|
...(cached.tagName ? { tagName: cached.tagName } : {}),
|
|
...(cached.releaseName ? { releaseName: cached.releaseName } : {}),
|
|
...(cached.publishedAt ? { publishedAt: cached.publishedAt } : {}),
|
|
...(cached.notes ? { notes: cached.notes } : {}),
|
|
...(cached.releaseUrl ? { releaseUrl: cached.releaseUrl } : {}),
|
|
compatible: Boolean(cached.asset),
|
|
integrityReady: Boolean(cached.asset?.sha256) && (!config.updateRequireSignature || cached.signatureVerified === true),
|
|
signatureReady: !config.updateRequireSignature || cached.signatureVerified === true,
|
|
isNewer: isNewerVersion(config.appVersion, cached.version),
|
|
...(cached.asset ? { assetName: cached.asset.name, ...(cached.asset.size === undefined ? {} : { assetSize: cached.asset.size }) } : {}),
|
|
},
|
|
};
|
|
}
|
|
|
|
export async function writeUpdateRequest(config: AppConfig, request: UpdateRequest): Promise<void> {
|
|
const parent = path.dirname(config.updateRequestPath);
|
|
await mkdir(parent, { recursive: true, mode: 0o700 });
|
|
const temporary = `${config.updateRequestPath}.tmp-${randomUUID()}`;
|
|
await writeFile(temporary, JSON.stringify(request), { encoding: "utf8", mode: 0o600, flag: "wx" });
|
|
try {
|
|
await chmod(temporary, 0o600);
|
|
await rename(temporary, config.updateRequestPath);
|
|
} catch (error) {
|
|
await import("node:fs/promises").then(({ rm }) => rm(temporary, { force: true })).catch(() => undefined);
|
|
throw error;
|
|
}
|
|
}
|
|
|
|
|
|
function safePublicErrorMessage(msg: unknown): string {
|
|
if (typeof msg !== "string" || !msg.trim()) return "更新失败,请查看服务器日志或重试";
|
|
if (msg.includes("/var/lib") || msg.includes("/opt/") || msg.includes("/etc/") || msg.includes("secret") || msg.includes("command-output")) {
|
|
return "更新失败,请查看服务器日志或重试";
|
|
}
|
|
return msg.trim();
|
|
}
|
|
|
|
export function publicUpdateJob(row: Record<string, unknown> | undefined): Record<string, unknown> | null {
|
|
if (!row) return null;
|
|
const hasError = typeof row.errorMessage === "string" && row.errorMessage.length > 0;
|
|
const updatedAt = typeof row.updatedAt === "number" ? row.updatedAt : null;
|
|
const expectedRecoveryAt = row.status === "applying" && updatedAt !== null ? updatedAt + 30_000 : null;
|
|
return {
|
|
id: row.id,
|
|
operation: row.operation ?? "apply",
|
|
status: row.status,
|
|
version: row.version,
|
|
platform: row.platform,
|
|
assetName: row.assetName ?? null,
|
|
assetUrl: row.assetUrl ?? null,
|
|
releaseUrl: row.releaseUrl ?? null,
|
|
sizeBytes: row.sizeBytes ?? null,
|
|
downloadedBytes: row.downloadedBytes ?? null,
|
|
downloadStartedAt: row.downloadStartedAt ?? null,
|
|
downloadSpeedBps: row.downloadSpeedBps ?? null,
|
|
errorMessage: hasError ? safePublicErrorMessage(row.errorMessage) : null,
|
|
createdAt: row.createdAt,
|
|
updatedAt: row.updatedAt,
|
|
completedAt: row.completedAt ?? null,
|
|
...(row.applyQueuedAt ? { applyQueuedAt: row.applyQueuedAt } : {}),
|
|
...(expectedRecoveryAt ? { expectedRecoveryAt } : {}),
|
|
...(row.status === "applying" ? { restartWindowSeconds: 30 } : {}),
|
|
};
|
|
}
|
|
|
|
function markerMtime(filePath: string): number | null {
|
|
try {
|
|
const info = lstatSync(filePath);
|
|
return info.isFile() ? info.mtimeMs : null;
|
|
} catch {
|
|
return null;
|
|
}
|
|
}
|
|
|
|
function forceRemoveRequest(filePath: string): void {
|
|
try {
|
|
const info = lstatSync(filePath);
|
|
if (!info.isFile() && !info.isSymbolicLink()) return;
|
|
unlinkSync(filePath);
|
|
} catch {}
|
|
}
|
|
|
|
function removeExpiredRequest(filePath: string, now: number): void {
|
|
try {
|
|
const info = lstatSync(filePath);
|
|
if (!info.isFile() && !info.isSymbolicLink()) return;
|
|
if (now - info.mtimeMs < ORPHANED_UPDATE_TIMEOUT_MS) return;
|
|
unlinkSync(filePath);
|
|
} catch {
|
|
// The root runner may own the marker during a recovery race. The DB
|
|
// transition below is still enough to release the browser queue.
|
|
}
|
|
}
|
|
|
|
function requestJobId(filePath: string): string | null {
|
|
try {
|
|
const info = lstatSync(filePath);
|
|
if (!info.isFile() || info.isSymbolicLink()) return null;
|
|
const value = JSON.parse(readFileSync(filePath, "utf8")) as { jobId?: unknown };
|
|
return typeof value.jobId === "string" && /^[0-9a-f-]{36}$/.test(value.jobId) ? value.jobId : null;
|
|
} catch {
|
|
return null;
|
|
}
|
|
}
|
|
|
|
function currentReleaseVersion(config: AppConfig): string | null {
|
|
try {
|
|
const target = realpathSync(config.currentLink);
|
|
const releases = realpathSync(config.releasesDir);
|
|
if (!target.startsWith(`${releases}${path.sep}`)) return null;
|
|
return path.basename(target);
|
|
} catch {
|
|
return null;
|
|
}
|
|
}
|
|
|
|
/**
|
|
* Release an update row left behind after its privileged runner lease expired.
|
|
* This is deliberately conservative: staged downloads remain available for an
|
|
* explicit apply, and a fresh request/state marker means the runner still owns
|
|
* recovery.
|
|
*/
|
|
export function reconcileOrphanedUpdateJobs(database: Database.Database, config: AppConfig, now = Date.now()): number {
|
|
const placeholders = ACTIVE_UPDATE_STATUSES.map(() => "?").join(",");
|
|
const rows = database.prepare(`
|
|
SELECT id, status, operation, version, admin_id AS adminId, request_id AS requestId,
|
|
updated_at AS updatedAt
|
|
FROM update_jobs
|
|
WHERE status IN (${placeholders})
|
|
ORDER BY updated_at ASC
|
|
`).all(...ACTIVE_UPDATE_STATUSES) as Array<{ id: string; status: UpdateJobStatus; operation: "download" | "apply"; version: string; adminId: string | null; requestId: string | null; updatedAt: number | null }>;
|
|
if (rows.length === 0) return 0;
|
|
const statePath = path.join(config.installPrefix, ".update-state");
|
|
const requestMtime = markerMtime(config.updateRequestPath);
|
|
const stateMtime = markerMtime(statePath);
|
|
const requestPresent = requestMtime !== null;
|
|
const statePresent = stateMtime !== null;
|
|
const requestFresh = requestPresent && now - (requestMtime ?? 0) < ORPHANED_UPDATE_TIMEOUT_MS;
|
|
const stateFresh = statePresent && now - (stateMtime ?? 0) < ORPHANED_UPDATE_TIMEOUT_MS;
|
|
// A staged download is normally kept for an explicit apply. The one
|
|
// exception is the hand-off window where the API has already changed the
|
|
// operation to `apply` but crashed before writing the request file. That
|
|
// row is still safe to retry and must not block the queue forever.
|
|
const releaseVersion = currentReleaseVersion(config);
|
|
let reconciled = 0;
|
|
const reconciledIds = new Set<string>();
|
|
for (const row of rows) {
|
|
if (typeof row.updatedAt !== "number" || now - row.updatedAt < ORPHANED_UPDATE_TIMEOUT_MS) continue;
|
|
// The runner refreshes the state marker while a download is in flight.
|
|
// A stale request/state marker therefore no longer protects an orphaned
|
|
// row forever, while a fresh marker remains owned by the runner.
|
|
if (row.status === "staged") {
|
|
if (row.operation !== "apply" || requestFresh || stateFresh) continue;
|
|
const changed = database.transaction(() => {
|
|
const result = database.prepare(`
|
|
UPDATE update_jobs
|
|
SET operation='download', error_message=NULL, updated_at=?
|
|
WHERE id=? AND status='staged' AND operation='apply' AND updated_at=?
|
|
`).run(now, row.id, row.updatedAt);
|
|
if (result.changes !== 1) return false;
|
|
writeAudit(database, {
|
|
requestId: row.requestId || randomUUID(),
|
|
actorAdminId: row.adminId,
|
|
action: "update.reconciled",
|
|
targetType: "update",
|
|
targetId: row.id,
|
|
outcome: "success",
|
|
before: { status: row.status, operation: row.operation, version: row.version },
|
|
after: { status: "staged", operation: "download", version: row.version, reason: "apply_request_missing" },
|
|
});
|
|
return true;
|
|
})();
|
|
if (changed) {
|
|
reconciled += 1;
|
|
reconciledIds.add(row.id);
|
|
}
|
|
continue;
|
|
}
|
|
if (requestFresh || stateFresh) continue;
|
|
const status: "completed" | "failed" = row.status === "applying" && releaseVersion === row.version ? "completed" : "failed";
|
|
const errorMessage = status === "failed" ? "更新任务超时,已释放更新队列" : null;
|
|
const changed = database.transaction(() => {
|
|
const result = database.prepare(`
|
|
UPDATE update_jobs
|
|
SET status=?, error_message=?, completed_at=?, updated_at=?
|
|
WHERE id=? AND status=? AND updated_at=?
|
|
`).run(status, errorMessage, now, now, row.id, row.status, row.updatedAt);
|
|
if (result.changes !== 1) return false;
|
|
writeAudit(database, {
|
|
requestId: row.requestId || randomUUID(),
|
|
actorAdminId: row.adminId,
|
|
action: "update.reconciled",
|
|
targetType: "update",
|
|
targetId: row.id,
|
|
outcome: status === "completed" ? "success" : "failure",
|
|
before: { status: row.status, version: row.version },
|
|
after: { status, version: row.version, reason: "orphaned_timeout" },
|
|
});
|
|
return true;
|
|
})();
|
|
if (changed) {
|
|
reconciled += 1;
|
|
reconciledIds.add(row.id);
|
|
}
|
|
}
|
|
// Prevent a stale request from being replayed after its DB row has been
|
|
// marked failed. The path is fixed by the server configuration and the
|
|
// operation is safe even when a root runner is racing with this call.
|
|
// A download runner refreshes the state marker while it is still using the
|
|
// request. Keep the request until that lease also expires; otherwise a
|
|
// long download can lose its job id and fail to finalize its row.
|
|
const queuedRequestId = requestPresent ? requestJobId(config.updateRequestPath) : null;
|
|
const queuedRequest = queuedRequestId ? rows.find((row) => row.id === queuedRequestId) : undefined;
|
|
const requestStillNeeded = Boolean(
|
|
queuedRequest
|
|
&& ACTIVE_UPDATE_STATUSES.includes(queuedRequest.status)
|
|
&& !reconciledIds.has(queuedRequest.id)
|
|
&& !(queuedRequest.status === "staged" && queuedRequest.operation === "download"),
|
|
);
|
|
if (!stateFresh && !requestStillNeeded) {
|
|
forceRemoveRequest(config.updateRequestPath);
|
|
} else if (!stateFresh && (!requestPresent || (requestMtime !== null && now - requestMtime >= ORPHANED_UPDATE_TIMEOUT_MS))) {
|
|
removeExpiredRequest(config.updateRequestPath, now);
|
|
}
|
|
return reconciled;
|
|
}
|
|
|
|
export function cancelUpdateJob(
|
|
database: Database.Database,
|
|
config: AppConfig,
|
|
adminId: string,
|
|
requestId: string,
|
|
jobId?: string,
|
|
): { cancelled: boolean; message?: string } {
|
|
const job = jobId
|
|
? database.prepare("SELECT id, status, operation, version, admin_id AS adminId, download_path AS downloadPath FROM update_jobs WHERE id=?").get(jobId) as { id: string; status: UpdateJobStatus; operation: string; version: string; adminId: string | null; downloadPath: string | null } | undefined
|
|
: database.prepare("SELECT id, status, operation, version, admin_id AS adminId, download_path AS downloadPath FROM update_jobs WHERE status IN ('queued', 'downloading') ORDER BY created_at DESC LIMIT 1").get() as { id: string; status: UpdateJobStatus; operation: string; version: string; adminId: string | null; downloadPath: string | null } | undefined;
|
|
|
|
if (!job) return { cancelled: false, message: "当前没有处于等待调度或下载中的更新任务" };
|
|
if (job.status !== "queued" && job.status !== "downloading") return { cancelled: false, message: "任务已进入就绪或切换阶段,无法取消" };
|
|
|
|
const controller = activeInProcessDownloads.get(job.id);
|
|
if (controller) {
|
|
controller.abort();
|
|
activeInProcessDownloads.delete(job.id);
|
|
}
|
|
|
|
const now = Date.now();
|
|
const changed = database.transaction(() => {
|
|
const result = database.prepare("UPDATE update_jobs SET status='cancelled', error_message='已手动取消更新', completed_at=?, updated_at=? WHERE id=? AND status IN ('queued', 'downloading')").run(now, now, job.id);
|
|
if (result.changes !== 1) return false;
|
|
writeAudit(database, {
|
|
requestId,
|
|
actorAdminId: adminId,
|
|
action: "update.cancelled",
|
|
targetType: "update",
|
|
targetId: job.id,
|
|
outcome: "success",
|
|
before: { status: job.status, operation: job.operation, version: job.version },
|
|
after: { status: "cancelled", version: job.version },
|
|
});
|
|
return true;
|
|
})();
|
|
|
|
if (changed) {
|
|
forceRemoveRequest(config.updateRequestPath);
|
|
if (job.downloadPath) {
|
|
const target = path.isAbsolute(job.downloadPath) ? job.downloadPath : path.join(config.stagingDir, job.downloadPath);
|
|
import("node:fs/promises").then(({ rm }) => rm(target, { recursive: true, force: true })).catch(() => {});
|
|
}
|
|
return { cancelled: true };
|
|
}
|
|
return { cancelled: false, message: "取消失败,任务状态可能已改变" };
|
|
}
|