import { lstatSync, realpathSync, readFileSync, unlinkSync } from "node:fs"; import { chmod, mkdir, mkdtemp, rename, rm, 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, downloadReleaseAsset, extractSafeArchive, fetchReleaseBytes, fetchReleaseMetadata, fetchReleaseText, isNewerVersion, normalizeReleasePermissions, parseSemver, runtimeHashFromLockfile, sanitizeAssetName, selectReleaseAsset, validateHttpsUrl, RELEASE_NOTES_MAX_BYTES, type ReleaseAsset, type ReleaseMetadata, } from "./update.js"; import type { UpdateJobStatus } from "../shared/contracts.js"; 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; 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; timeoutMs?: number | undefined; 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), timeoutMs: options.timeoutMs }); 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, timeoutMs: options.timeoutMs }); 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, timeoutMs: config.updateTimeoutMs, } 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 { 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. } // Force choosing the full standalone archive so users always get a real, visible streaming download let asset = selectReleaseAsset(metadata, platform, undefined); let signatureVerified = false; if (asset) { const integrity = await attachSidecarHash(metadata, asset, { allowedHosts: config.updateAllowedHosts, baseUrl: metadataUrl, maxBytes: config.updateMaxBytes, timeoutMs: config.updateTimeoutMs, 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 { 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 | undefined): Record | 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 recoveryStateJobId(filePath: string): string | null { try { const info = lstatSync(filePath); if (!info.isFile() || info.isSymbolicLink()) return null; const match = /^job_id=([0-9a-f-]{36})$/m.exec(readFileSync(filePath, "utf8")); return match?.[1] ?? null; } catch { return null; } } export 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; // The request marker is the hand-off contract between the web process and // the privileged runner. A queued row with a matching, unexpired marker is // still owned by that hand-off even when the runner has not written its // recovery state yet (for example while systemd is starting it). const requestMarkerJobId = requestPresent ? requestJobId(config.updateRequestPath) : null; const stateMarkerJobId = statePresent ? recoveryStateJobId(statePath) : null; // 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(); for (const row of rows) { // A fresh request/state marker means the privileged runner still owns the // hand-off. Do not expire a staged/apply row while the runner is finishing // a successful switch and finalization after a service restart. const matchingFreshRequest = requestMarkerJobId === row.id && requestFresh; const matchingFreshState = stateMarkerJobId === row.id && stateFresh; // A staged archive is actionable only while it is strictly newer than the // release currently serving requests. This can become false when an // administrator upgrades the host by another path (or another operator // completes the same release) before returning to this page. Treat the // archive as an expired terminal task so it cannot keep blocking the // queue or appear as an "apply" action for the current version. const effectiveCurrentVersion = releaseVersion ?? config.appVersion; if (row.status === "staged" && !isNewerVersion(effectiveCurrentVersion, row.version) && !matchingFreshRequest && !matchingFreshState) { const changed = database.transaction(() => { const result = database.prepare(` UPDATE update_jobs SET status='failed', error_message=?, completed_at=?, updated_at=? WHERE id=? AND status='staged' `).run("暂存更新已过期,当前版本无需再次升级", now, now, row.id); if (result.changes !== 1) return false; writeAudit(database, { requestId: row.requestId || randomUUID(), actorAdminId: row.adminId, action: "update.reconciled", targetType: "update", targetId: row.id, outcome: "failure", before: { status: row.status, operation: row.operation, version: row.version }, after: { status: "failed", version: row.version, reason: "staged_version_not_newer" }, }); return true; })(); if (changed) { reconciled += 1; reconciledIds.add(row.id); } continue; } // A request that never gets claimed by the root runner must not remain in // the UI as an endless "queued" task. Once the short hand-off window has // elapsed and no recovery marker exists, release the queue explicitly; // a fresh state marker proves that the runner has already claimed it. if (row.status === "queued" && typeof row.updatedAt === "number" && !matchingFreshState && now - row.updatedAt >= QUEUED_UPDATE_TIMEOUT_MS) { if (matchingFreshRequest) continue; const changed = database.transaction(() => { const result = database.prepare(` UPDATE update_jobs SET status='failed', error_message=?, completed_at=?, updated_at=? WHERE id=? AND status='queued' AND updated_at=? `).run("更新服务未在规定时间内接管任务", now, 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: "failure", before: { status: row.status, version: row.version }, after: { status: "failed", version: row.version, reason: "runner_claim_timeout" }, }); return true; })(); if (changed) { reconciled += 1; reconciledIds.add(row.id); } continue; } 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" || matchingFreshRequest || matchingFreshState) 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 (matchingFreshRequest || matchingFreshState) 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=? AND admin_id=?").get(jobId, adminId) 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 admin_id=? AND status IN ('queued', 'downloading') ORDER BY created_at DESC LIMIT 1").get(adminId) 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 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 admin_id=? AND status IN ('queued', 'downloading')").run(now, now, job.id, adminId); 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) { // The request marker is shared by the privileged runner. Never remove a // newer/different administrator's request while cancelling this row. if (requestJobId(config.updateRequestPath) === job.id) 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: "取消失败,任务状态可能已改变" }; } /** * Download, verify and stage a release archive in the web process (non-root). * The root runner only needs to apply (stop/backup/switch/restart) afterwards. * * This function runs asynchronously outside the request lifecycle. It updates * the job row in the database so the frontend can poll progress. On success it * writes an apply request file so the systemd path unit triggers the runner. */ export async function downloadAndStageUpdate( database: Database.Database, config: AppConfig, jobId: string, adminId: string, version: string, assetUrl: string, assetName: string, expectedSha256: string, metadataUrl: string, ): Promise { const stagingBase = path.resolve(config.stagingDir); const workspace = path.join(stagingBase, `update-${jobId}`); try { await mkdir(workspace, { recursive: true, mode: 0o700 }); const archiveName = assetName.endsWith(".tar.gz") || assetName.endsWith(".tgz") ? assetName : `${assetName}.tar.gz`; const archivePath = path.join(workspace, archiveName); // Claim the job: transition queued -> downloading. If the job was // cancelled or claimed by another caller, abort immediately. const claim = database.prepare( "UPDATE update_jobs SET status='downloading', download_started_at=?, started_at=?, download_path=?, updated_at=? WHERE id=? AND status='queued'", ).run(Date.now(), Date.now(), path.basename(archivePath), Date.now(), jobId); if (claim.changes !== 1) return; const progressStartedAt = Date.now(); let lastProgressWrite = 0; const downloaded = await downloadReleaseAsset(assetUrl, archivePath, { allowedHosts: config.updateAllowedHosts, baseUrl: config.updateMetadataUrl, maxBytes: config.updateMaxBytes, timeoutMs: config.updateTimeoutMs, onProgress: (downloadedBytes, totalBytes) => { const now = Date.now(); if (now - lastProgressWrite < 250) return; lastProgressWrite = now; const elapsed = Math.max(1, now - progressStartedAt); const speedBps = Math.round(downloadedBytes * 1000 / elapsed); 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, now, jobId); }, }); // Final progress write const finishedAt = Date.now(); const elapsed = Math.max(1, finishedAt - progressStartedAt); database.prepare( "UPDATE update_jobs SET downloaded_bytes=?, size_bytes=?, download_speed_bps=?, updated_at=? WHERE id=? AND status='downloading'", ).run(downloaded.size, downloaded.size, Math.round(downloaded.size * 1000 / elapsed), finishedAt, jobId); // SHA-256 verification database.prepare( "UPDATE update_jobs SET status='verifying', actual_sha256=?, size_bytes=?, updated_at=? WHERE id=? AND status='downloading'", ).run(downloaded.sha256, downloaded.size, Date.now(), jobId); if (expectedSha256 && downloaded.sha256 !== expectedSha256) { throw new Error("更新文件 SHA-256 校验失败"); } // Extract archive to payload directory const payloadDir = path.join(workspace, "payload"); await extractSafeArchive(archivePath, payloadDir); await normalizeReleasePermissions(payloadDir); // Verify payload contains dist directory const { lstat } = await import("node:fs/promises"); const payloadInfo = await lstat(path.join(payloadDir, "dist")).catch(() => null); if (!payloadInfo?.isDirectory() || payloadInfo.isSymbolicLink()) { throw new Error("发布包缺少 dist 目录"); } // Transition to staged const staged = database.prepare( "UPDATE update_jobs SET status='staged', operation='apply', actual_sha256=?, size_bytes=?, download_path=?, updated_at=? WHERE id=? AND status IN ('verifying', 'downloading')", ).run(downloaded.sha256, downloaded.size, workspace, Date.now(), jobId); if (staged.changes !== 1) return; // cancelled // Write apply request file for the root runner await writeUpdateRequest(config, { jobId, operation: "apply", version, metadataUrl, assetUrl, assetName, expectedSha256, requestedAt: Date.now(), currentLink: config.currentLink, releasesDir: config.releasesDir, dataDir: config.dataDir, stagedPath: workspace, }); writeAudit(database, { requestId: `download:${jobId}`, actorAdminId: adminId, action: "update.staged", targetType: "update", targetId: jobId, after: { version, sha256: downloaded.sha256, size: downloaded.size }, }); } catch (error) { const message = error instanceof Error ? error.message : "下载或校验失败"; try { database.prepare( "UPDATE update_jobs SET status='failed', error_message=?, updated_at=? WHERE id=? AND status IN ('queued', 'downloading', 'verifying')", ).run(message, Date.now(), jobId); writeAudit(database, { requestId: `download:${jobId}`, actorAdminId: adminId, action: "update.download_failed", targetType: "update", targetId: jobId, outcome: "failure", metadata: { error: message }, }); } catch { // The database may be closed (e.g. during test cleanup or process // shutdown). The workspace cleanup below still runs unconditionally. } await rm(workspace, { recursive: true, force: true }).catch(() => undefined); } }