fix: 修复在线更新暂存链路并增加全局 API 限流备底
- 新增 server/rate-limit.ts:进程内固定窗口限流器,无数据库写入 - server/app.ts 注册全局 preHandler,仅作用于 /api/*,超限返回 429 与 Retry-After - 提取 isApiPath 统一 onSend、preHandler 与 404 的路径判断 - 更新任务冲突判定改用 ACTIVE_UPDATE_CONFLICT_SQL,staged/download 产物不再阻塞新任务 - cancelUpdateJob 调用补上 await,避免结果恒为 pending Promise - server/cli/update.ts 增加特权工作区所有权校验与暂存路径重建逻辑 - 新增 tests/rate-limit.test.ts 与 tests/update-apply-staging.test.ts
This commit is contained in:
+43
-8
@@ -54,9 +54,11 @@ import {
|
||||
validateNewPassword,
|
||||
verifyPassword,
|
||||
} from "./security.js";
|
||||
import { createRateLimiter } from "./rate-limit.js";
|
||||
import { isNewerVersion } from "./update.js";
|
||||
import {
|
||||
ACTIVE_UPDATE_STATUSES,
|
||||
ACTIVE_UPDATE_CONFLICT_SQL,
|
||||
checkForUpdate,
|
||||
currentReleaseVersion,
|
||||
publicCheckFromCache,
|
||||
@@ -100,6 +102,11 @@ const unsafeMethods = new Set(["POST", "PUT", "PATCH", "DELETE"]);
|
||||
const sessionCookie = "tally_session";
|
||||
const csrfCookie = "tally_csrf";
|
||||
|
||||
/** API paths are the only requests the global rate limiter and cache rules own. */
|
||||
function isApiPath(url: string): boolean {
|
||||
return url.split("?", 1)[0]!.startsWith("/api/");
|
||||
}
|
||||
|
||||
type UpdateRateState = { checkedAt: number; downloadedAt: number; appliedAt: number };
|
||||
const updateRateStates = new WeakMap<DatabaseContext["sqlite"], Map<string, UpdateRateState>>();
|
||||
|
||||
@@ -644,7 +651,7 @@ export async function buildApp(database: DatabaseContext, config: AppConfig) {
|
||||
// retained by a browser, reverse proxy or shared cache. Keep this global so
|
||||
// future authenticated routes inherit the same privacy boundary.
|
||||
app.addHook("onSend", async (request, reply, payload) => {
|
||||
if (request.url.split("?", 1)[0]!.startsWith("/api/")) {
|
||||
if (isApiPath(request.url)) {
|
||||
reply.header("Cache-Control", "no-store");
|
||||
reply.header("Pragma", "no-cache");
|
||||
reply.header("Vary", "Cookie");
|
||||
@@ -1028,7 +1035,10 @@ export async function buildApp(database: DatabaseContext, config: AppConfig) {
|
||||
enforceUpdateCooldown(database.sqlite, config, request.auth!.admin.id, "apply", reply);
|
||||
const now = Date.now();
|
||||
const active = database.sqlite.transaction(() => {
|
||||
const conflictRow = database.sqlite.prepare(`SELECT id FROM update_jobs WHERE status IN (${ACTIVE_UPDATE_STATUSES.map(() => "?").join(",")}) AND id<>? LIMIT 1`).get(...ACTIVE_UPDATE_STATUSES, stagedJobId) as { id: string } | undefined;
|
||||
// A reusable `staged/download` artifact must never block applying a
|
||||
// *different* staged job: otherwise a leftover row keeps the queue
|
||||
// permanently busy and the operator can never apply an update.
|
||||
const conflictRow = database.sqlite.prepare(`SELECT id FROM update_jobs WHERE ${ACTIVE_UPDATE_CONFLICT_SQL} AND id<>? LIMIT 1`).get(...ACTIVE_UPDATE_STATUSES, stagedJobId) as { id: string } | undefined;
|
||||
if (conflictRow) throw new AppError(409, "UPDATE_IN_PROGRESS", "已有更新任务正在进行,请等待完成");
|
||||
const changed = database.sqlite.prepare("UPDATE update_jobs SET operation='apply', error_message=NULL, requested_at=?, request_id=?, updated_at=? WHERE id=? AND status='staged' AND operation='download'").run(now, request.id, now, stagedJobId);
|
||||
if (changed.changes !== 1) throw new AppError(409, "UPDATE_IN_PROGRESS", "更新任务正在处理中,请稍候");
|
||||
@@ -1053,9 +1063,11 @@ export async function buildApp(database: DatabaseContext, config: AppConfig) {
|
||||
}
|
||||
// Preserve the actionable in-progress response for duplicate clicks before
|
||||
// applying the per-admin cooldown.
|
||||
// Same rule as the download path: a reusable staged download is an
|
||||
// artifact, not a running task, and must not block apply.
|
||||
const activeBeforeCheck = database.sqlite.prepare(`
|
||||
SELECT id FROM update_jobs
|
||||
WHERE status IN (${ACTIVE_UPDATE_STATUSES.map(() => "?").join(",")})
|
||||
WHERE ${ACTIVE_UPDATE_CONFLICT_SQL}
|
||||
ORDER BY created_at DESC LIMIT 1
|
||||
`).get(...ACTIVE_UPDATE_STATUSES) as { id: string } | undefined;
|
||||
if (activeBeforeCheck) throw new AppError(409, "UPDATE_IN_PROGRESS", "已有更新任务正在进行,请等待完成");
|
||||
@@ -1079,7 +1091,7 @@ export async function buildApp(database: DatabaseContext, config: AppConfig) {
|
||||
const active = database.sqlite.transaction(() => {
|
||||
const existing = database.sqlite.prepare(`
|
||||
SELECT id, status FROM update_jobs
|
||||
WHERE status IN (${ACTIVE_UPDATE_STATUSES.map(() => "?").join(",")})
|
||||
WHERE ${ACTIVE_UPDATE_CONFLICT_SQL}
|
||||
ORDER BY created_at DESC LIMIT 1
|
||||
`).get(...ACTIVE_UPDATE_STATUSES) as { id: string; status: UpdateJobStatus } | undefined;
|
||||
if (existing) throw new AppError(409, "UPDATE_IN_PROGRESS", "已有更新任务正在进行,请等待完成");
|
||||
@@ -1169,7 +1181,9 @@ export async function buildApp(database: DatabaseContext, config: AppConfig) {
|
||||
const input = updateDownloadSchema.parse(request.body);
|
||||
reconcileOrphanedUpdateJobs(database.sqlite, config);
|
||||
if (config.updateStrategy !== "systemd") throw new AppError(503, "UPDATE_NOT_AVAILABLE", "当前安装方式未启用一键更新,请使用命令行更新");
|
||||
const active = database.sqlite.prepare(`SELECT id FROM update_jobs WHERE status IN (${ACTIVE_UPDATE_STATUSES.map(() => "?").join(",")}) LIMIT 1`).get(...ACTIVE_UPDATE_STATUSES) as { id: string } | undefined;
|
||||
// A finished `staged/download` row is a reusable artifact, not a running
|
||||
// task, so it does not block a new download. Apply-phase rows still do.
|
||||
const active = database.sqlite.prepare(`SELECT id FROM update_jobs WHERE ${ACTIVE_UPDATE_CONFLICT_SQL} LIMIT 1`).get(...ACTIVE_UPDATE_STATUSES) as { id: string } | undefined;
|
||||
if (active) throw new AppError(409, "UPDATE_IN_PROGRESS", "已有更新任务正在进行,请等待完成");
|
||||
enforceUpdateCooldown(database.sqlite, config, request.auth!.admin.id, "download", reply);
|
||||
const checked = await checkForUpdate(database.sqlite, config);
|
||||
@@ -1181,7 +1195,7 @@ export async function buildApp(database: DatabaseContext, config: AppConfig) {
|
||||
const now = Date.now();
|
||||
const id = randomUUID();
|
||||
database.sqlite.transaction(() => {
|
||||
const conflict = database.sqlite.prepare(`SELECT id FROM update_jobs WHERE status IN (${ACTIVE_UPDATE_STATUSES.map(() => "?").join(",")}) LIMIT 1`).get(...ACTIVE_UPDATE_STATUSES) as { id: string } | undefined;
|
||||
const conflict = database.sqlite.prepare(`SELECT id FROM update_jobs WHERE ${ACTIVE_UPDATE_CONFLICT_SQL} LIMIT 1`).get(...ACTIVE_UPDATE_STATUSES) as { id: string } | undefined;
|
||||
if (conflict) throw new AppError(409, "UPDATE_IN_PROGRESS", "已有更新任务正在进行,请等待完成");
|
||||
database.sqlite.prepare(`INSERT INTO update_jobs(id, admin_id, session_hash, request_id, requested_at, operation, status, version, platform, release_url, asset_name, asset_url, expected_sha256, created_at, updated_at) VALUES (?, ?, ?, ?, ?, 'download', 'queued', ?, ?, ?, ?, ?, ?, ?, ?)`).run(id, request.auth!.admin.id, request.auth!.tokenHash, request.id, now, version, checked.platform.target, cached.metadataUrl, cachedAsset.name, cachedAsset.url, cachedAsset.sha256, now, now);
|
||||
writeAudit(database.sqlite, { requestId: request.id, actorAdminId: request.auth!.admin.id, actorUsername: request.auth!.admin.username, action: "update.download_requested", targetType: "update", targetId: id, after: { version } });
|
||||
@@ -1201,7 +1215,10 @@ export async function buildApp(database: DatabaseContext, config: AppConfig) {
|
||||
app.post("/api/update/cancel", { preHandler: guard(database, config) }, async (request, reply) => {
|
||||
reconcileOrphanedUpdateJobs(database.sqlite, config);
|
||||
const body = (request.body && typeof request.body === "object" ? request.body : {}) as { jobId?: string };
|
||||
const result = cancelUpdateJob(database.sqlite, config, request.auth!.admin.id, request.id, body.jobId);
|
||||
// `cancelUpdateJob` awaits its staging-workspace cleanup, so the caller must
|
||||
// await it too; without the await `result` is a pending Promise and this
|
||||
// branch would always report failure even after a successful cancel.
|
||||
const result = await cancelUpdateJob(database.sqlite, config, request.auth!.admin.id, request.id, body.jobId);
|
||||
if (!result.cancelled) {
|
||||
throw new AppError(409, "CANNOT_CANCEL", result.message || "无法取消当前更新任务");
|
||||
}
|
||||
@@ -1870,6 +1887,24 @@ export async function buildApp(database: DatabaseContext, config: AppConfig) {
|
||||
return reply.send(await safeReadStream(config.exportsDir, job.filePath));
|
||||
});
|
||||
|
||||
// Global anti-flood backstop for the API surface. It is registered after all
|
||||
// /api/* routes so it covers every one of them, but the path check keeps
|
||||
// static assets, `/health` and the SPA fallback out of the limiter. The
|
||||
// per-feature limits (login lockout, dangerous-operation re-auth, update
|
||||
// cooldowns) stay authoritative; this only bounds raw request volume.
|
||||
const apiRateLimiter = createRateLimiter({ limit: config.apiRateLimitPerMinute, windowMs: 60 * 1000 });
|
||||
app.addHook("preHandler", async (request, reply) => {
|
||||
if (!isApiPath(request.url)) return;
|
||||
// `request.ip` already honours the validated trustProxy configuration, so
|
||||
// the counted address is the one the deployment declared. The limiter must
|
||||
// never parse X-Forwarded-For itself, otherwise a client could spoof its
|
||||
// way around the limit.
|
||||
const decision = apiRateLimiter.check(request.ip);
|
||||
if (decision.allowed) return;
|
||||
reply.header("Retry-After", decision.retryAfterSeconds);
|
||||
throw new AppError(429, "RATE_LIMITED", "请求过于频繁,请稍后再试");
|
||||
});
|
||||
|
||||
const hasWeb = existsSync(config.webDir);
|
||||
if (hasWeb) {
|
||||
// Serve the Vite asset graph as well as the SPA entry. API routes are
|
||||
@@ -1879,7 +1914,7 @@ export async function buildApp(database: DatabaseContext, config: AppConfig) {
|
||||
// Keep API errors structured even when the production frontend has not been
|
||||
// built yet (for example in a clean CI checkout or an API-only process).
|
||||
app.setNotFoundHandler((request, reply) => {
|
||||
if (request.url.split("?", 1)[0]!.startsWith("/api/")) {
|
||||
if (isApiPath(request.url)) {
|
||||
return reply.code(404).send(errorPayload(request, new AppError(404, "NOT_FOUND", "接口不存在")));
|
||||
}
|
||||
if (hasWeb) return reply.sendFile("index.html");
|
||||
|
||||
+312
-44
@@ -1,5 +1,5 @@
|
||||
import { randomUUID } from "node:crypto";
|
||||
import { lstat, mkdir, mkdtemp, readFile, realpath, rm } from "node:fs/promises";
|
||||
import { copyFile, lstat, mkdir, mkdtemp, readFile, readdir, realpath, rm } from "node:fs/promises";
|
||||
import path from "node:path";
|
||||
import { pathToFileURL } from "node:url";
|
||||
import type Database from "better-sqlite3";
|
||||
@@ -22,6 +22,7 @@ import {
|
||||
selectReleaseAsset,
|
||||
sanitizeAssetName,
|
||||
validateHttpsUrl,
|
||||
verifySha256,
|
||||
type ReleaseAsset,
|
||||
type ReleaseMetadata,
|
||||
type UrlPolicy,
|
||||
@@ -244,35 +245,243 @@ async function resolveRelease(options: UpdateRunOptions, platform: ReturnType<ty
|
||||
return { asset: { name: sanitizeAssetName(options.assetName ?? path.basename(assetUrl.pathname)), url: assetUrl.toString(), ...(options.expectedSha256 ? { sha256: options.expectedSha256 } : {}) }, version: options.version };
|
||||
}
|
||||
|
||||
async function ensurePrivilegedWorkspace(directory: string): Promise<string> {
|
||||
/**
|
||||
* Ownership expectation for a directory consumed by the privileged updater.
|
||||
*
|
||||
* `-1` disables the uid comparison while keeping the symlink and mode checks.
|
||||
* Production always passes a concrete uid (0 for the root-owned
|
||||
* `<installPrefix>/.update-work`), so the check never depends on the effective
|
||||
* uid of the current process and remains runnable from a non-root test.
|
||||
*/
|
||||
export type DirectoryOwnerUid = number;
|
||||
|
||||
/** The staging area owned by the unprivileged web process and the private
|
||||
* root-owned workspace are deliberately separate trust domains. */
|
||||
export class StagedWorkspaceError extends Error {
|
||||
readonly reason: string;
|
||||
constructor(message: string, reason: string) {
|
||||
super(message);
|
||||
this.name = "StagedWorkspaceError";
|
||||
this.reason = reason;
|
||||
}
|
||||
}
|
||||
|
||||
/** Owner uid of an existing path, or -1 when it cannot be inspected. */
|
||||
export async function directoryOwnerUid(targetPath: string): Promise<DirectoryOwnerUid> {
|
||||
const info = await lstat(path.resolve(targetPath)).catch(() => null);
|
||||
return info?.uid ?? -1;
|
||||
}
|
||||
|
||||
/**
|
||||
* Resolve a privileged workspace root to its canonical path.
|
||||
*
|
||||
* The root itself may be reached through a symlinked ancestor (for example
|
||||
* `/tmp` on macOS), so only the final component is required to be a real,
|
||||
* non-symlink directory with private permissions and the expected owner.
|
||||
*/
|
||||
async function canonicalizePrivilegedRoot(directory: string, expectedUid: DirectoryOwnerUid, message: string): Promise<string> {
|
||||
const resolved = path.resolve(directory);
|
||||
await mkdir(resolved, { recursive: true, mode: 0o700 });
|
||||
const info = await lstat(resolved).catch(() => null);
|
||||
const uid = typeof process.getuid === "function" ? process.getuid() : -1;
|
||||
if (!info?.isDirectory() || info.isSymbolicLink() || (info.mode & 0o077) !== 0 || info.uid !== 0 || uid !== 0) {
|
||||
throw new Error("更新工作目录必须是 root 拥有且权限为 0700");
|
||||
if (!info?.isDirectory() || info.isSymbolicLink() || (info.mode & 0o077) !== 0 || (expectedUid >= 0 && info.uid !== expectedUid)) {
|
||||
throw new Error(message);
|
||||
}
|
||||
return resolved;
|
||||
const real = await realpath(resolved).catch(() => { throw new Error(message); });
|
||||
const realInfo = await lstat(real).catch(() => null);
|
||||
if (!realInfo?.isDirectory() || realInfo.isSymbolicLink() || (realInfo.mode & 0o077) !== 0 || (expectedUid >= 0 && realInfo.uid !== expectedUid)) {
|
||||
throw new Error(message);
|
||||
}
|
||||
return real;
|
||||
}
|
||||
|
||||
/** Validate a queued staged directory before a root process consumes it. */
|
||||
async function validateStagedWorkspacePath(candidate: string, workspaceRoot: string): Promise<string> {
|
||||
const rootResolved = path.resolve(workspaceRoot);
|
||||
const rootInfo = await lstat(rootResolved).catch(() => null);
|
||||
const uid = typeof process.getuid === "function" ? process.getuid() : -1;
|
||||
if (!rootInfo?.isDirectory() || rootInfo.isSymbolicLink() || (rootInfo.mode & 0o077) !== 0 || rootInfo.uid !== 0 || uid !== 0) {
|
||||
throw new Error("更新工作目录权限无效");
|
||||
}
|
||||
const root = await realpath(rootResolved).catch(() => { throw new Error("更新工作目录无效"); });
|
||||
/** Assert that `candidate` is a real, private, expected-owner directory below `root`. */
|
||||
async function assertStagedDirectory(candidate: string, root: string, expectedUid: DirectoryOwnerUid): Promise<string> {
|
||||
const resolved = path.resolve(candidate);
|
||||
if (resolved === rootResolved || !resolved.startsWith(`${rootResolved}${path.sep}`)) throw new Error("更新暂存路径无效");
|
||||
if (resolved === root || !resolved.startsWith(`${root}${path.sep}`)) throw new StagedWorkspaceError("更新暂存路径无效", "staged_workspace_invalid");
|
||||
const info = await lstat(resolved).catch(() => null);
|
||||
if (!info?.isDirectory() || info.isSymbolicLink() || (info.mode & 0o077) !== 0 || info.uid !== 0) throw new Error("更新暂存目录权限无效");
|
||||
const real = await realpath(resolved).catch(() => { throw new Error("更新暂存目录无效"); });
|
||||
if (real !== resolved || !real.startsWith(`${root}${path.sep}`)) throw new Error("更新暂存路径无效");
|
||||
if (!info?.isDirectory() || info.isSymbolicLink() || (info.mode & 0o077) !== 0) throw new StagedWorkspaceError("更新暂存目录权限无效", "staged_workspace_insecure");
|
||||
if (expectedUid >= 0 && info.uid !== expectedUid) throw new StagedWorkspaceError("更新暂存目录属主无效", "staged_workspace_insecure");
|
||||
const real = await realpath(resolved).catch(() => { throw new StagedWorkspaceError("更新暂存目录无效", "staged_workspace_invalid"); });
|
||||
if (real !== resolved || !real.startsWith(`${root}${path.sep}`)) throw new StagedWorkspaceError("更新暂存路径无效", "staged_workspace_invalid");
|
||||
return real;
|
||||
}
|
||||
|
||||
async function ensurePrivilegedWorkspace(directory: string, expectedUid: DirectoryOwnerUid): Promise<string> {
|
||||
return canonicalizePrivilegedRoot(directory, expectedUid, "更新工作目录必须是 root 拥有且权限为 0700");
|
||||
}
|
||||
|
||||
/**
|
||||
* Rebuild the staged workspace location for `jobId` instead of trusting the
|
||||
* `download_path` column: the runner clears that column whenever it releases
|
||||
* a workspace (`clearTransientJobPath`), and a nulled column cannot be used to
|
||||
* find a payload that is still on disk waiting for the apply step.
|
||||
*
|
||||
* Candidate order:
|
||||
* 1. `<stagingRoot>/update-<jobId>` (web download workspace)
|
||||
* 2. `<stagingRoot>/update-<jobId>-*` (mkdtemp variant)
|
||||
* 3. the recorded `download_path`, but only while it stays inside the root
|
||||
*/
|
||||
export async function locateStagedWorkspace(options: {
|
||||
jobId: string;
|
||||
downloadPath?: string | null | undefined;
|
||||
stagingRoot: string;
|
||||
expectedUid: DirectoryOwnerUid;
|
||||
}): Promise<string> {
|
||||
const root = await canonicalizePrivilegedRoot(options.stagingRoot, options.expectedUid, "更新暂存根目录权限无效");
|
||||
const prefix = `update-${options.jobId}`;
|
||||
const candidates = [path.join(root, prefix)];
|
||||
const entries = await readdir(root, { withFileTypes: true }).catch(() => []);
|
||||
for (const entry of entries.filter((candidate) => candidate.name.startsWith(`${prefix}-`)).sort((a, b) => a.name.localeCompare(b.name))) {
|
||||
candidates.push(path.join(root, entry.name));
|
||||
}
|
||||
const recorded = options.downloadPath?.trim();
|
||||
if (recorded && path.isAbsolute(recorded)) {
|
||||
const resolvedRecorded = path.resolve(recorded);
|
||||
if (resolvedRecorded.startsWith(`${root}${path.sep}`)) candidates.push(resolvedRecorded);
|
||||
// A recorded path outside the staging root is never consumed. Reject it
|
||||
// loudly when it exists so the operator sees the real cause instead of a
|
||||
// generic "re-download" message.
|
||||
else if (await lstat(resolvedRecorded).catch(() => null)) throw new StagedWorkspaceError("更新暂存路径无效", "staged_workspace_invalid");
|
||||
}
|
||||
const seen = new Set<string>();
|
||||
for (const candidate of candidates) {
|
||||
const resolved = path.resolve(candidate);
|
||||
if (seen.has(resolved)) continue;
|
||||
seen.add(resolved);
|
||||
// An existing candidate must satisfy every constraint: skipping it would
|
||||
// hand the root process whatever else happens to sit in the staging area.
|
||||
if (!(await lstat(resolved).catch(() => null))) continue;
|
||||
return assertStagedDirectory(resolved, root, options.expectedUid);
|
||||
}
|
||||
throw new StagedWorkspaceError("暂存目录已不存在,请重新下载", "staged_workspace_missing");
|
||||
}
|
||||
|
||||
/** Copy a verified payload tree without following or preserving symlinks. */
|
||||
async function copyReleaseTree(source: string, target: string): Promise<void> {
|
||||
await mkdir(target, { recursive: false, mode: 0o700 });
|
||||
const entries = await readdir(source, { withFileTypes: true });
|
||||
for (const entry of entries) {
|
||||
const from = path.join(source, entry.name);
|
||||
const to = path.join(target, entry.name);
|
||||
if (entry.isSymbolicLink()) throw new Error("更新暂存内容包含符号链接");
|
||||
if (entry.isDirectory()) await copyReleaseTree(from, to);
|
||||
else if (entry.isFile()) await copyFile(from, to);
|
||||
else throw new Error("更新暂存内容包含不受支持的文件类型");
|
||||
}
|
||||
}
|
||||
|
||||
const STAGED_ARCHIVE_PATTERN = /\.(?:tar\.gz|tgz|tar|zip)$/i;
|
||||
|
||||
async function assertStagedArchiveIntegrity(source: string, expectedSha256: string | null | undefined): Promise<void> {
|
||||
const expected = expectedSha256?.trim().toLowerCase();
|
||||
if (!expected) return;
|
||||
if (!/^[a-f0-9]{64}$/.test(expected)) throw new Error("更新暂存校验值无效");
|
||||
const entries = await readdir(source, { withFileTypes: true }).catch(() => []);
|
||||
const archive = entries
|
||||
.filter((entry) => entry.isFile() && !entry.isSymbolicLink() && STAGED_ARCHIVE_PATTERN.test(entry.name))
|
||||
.sort((a, b) => a.name.localeCompare(b.name))[0];
|
||||
if (!archive) throw new Error("更新暂存归档缺失,无法校验完整性");
|
||||
const archivePath = path.join(source, archive.name);
|
||||
const info = await lstat(archivePath).catch(() => null);
|
||||
if (!info?.isFile() || info.isSymbolicLink()) throw new Error("更新暂存归档无效");
|
||||
if (!(await verifySha256(archivePath, expected))) throw new Error("更新文件 SHA-256 校验失败");
|
||||
}
|
||||
|
||||
/**
|
||||
* Snapshot the web-staged payload into a root-owned workspace before the apply
|
||||
* flow touches it. The copy is what closes the TOCTOU window: the unprivileged
|
||||
* web user keeps write access to its own staging directory, so the privileged
|
||||
* process must never execute content that lives there.
|
||||
*
|
||||
* A copy (not a rename) is required because the data directory and the install
|
||||
* prefix are frequently separate mounts, where `rename` fails with EXDEV.
|
||||
*/
|
||||
export async function preparePrivateApplyWorkspace(options: {
|
||||
jobId: string;
|
||||
source: string;
|
||||
privateRoot: string;
|
||||
expectedUid: DirectoryOwnerUid;
|
||||
expectedSha256?: string | null | undefined;
|
||||
}): Promise<string> {
|
||||
const root = await canonicalizePrivilegedRoot(options.privateRoot, options.expectedUid, "更新工作目录必须是 root 拥有且权限为 0700");
|
||||
const source = path.resolve(options.source);
|
||||
const sourcePayload = path.join(source, "payload");
|
||||
const sourcePayloadInfo = await lstat(sourcePayload).catch(() => null);
|
||||
if (!sourcePayloadInfo?.isDirectory() || sourcePayloadInfo.isSymbolicLink()) throw new Error("更新暂存内容无效");
|
||||
// Second integrity check right before the copy, so a payload swapped after
|
||||
// the download verification is rejected instead of promoted to a release.
|
||||
await assertStagedArchiveIntegrity(source, options.expectedSha256);
|
||||
const target = path.join(root, `apply-${options.jobId}`);
|
||||
const existing = await lstat(target).catch(() => null);
|
||||
if (existing) await rm(target, { recursive: true, force: true }).catch(() => undefined);
|
||||
try {
|
||||
// Create the container explicitly: `copyReleaseTree` intentionally uses a
|
||||
// non-recursive mkdir so a pre-existing/symlinked target can never be
|
||||
// silently reused, and the parent must therefore already exist.
|
||||
await mkdir(target, { recursive: false, mode: 0o700 });
|
||||
await copyReleaseTree(sourcePayload, path.join(target, "payload"));
|
||||
await normalizeReleasePermissions(path.join(target, "payload"));
|
||||
const copied = await lstat(path.join(target, "payload")).catch(() => null);
|
||||
if (!copied?.isDirectory() || copied.isSymbolicLink()) throw new Error("更新暂存内容复制失败");
|
||||
return target;
|
||||
} catch (error) {
|
||||
await rm(target, { recursive: true, force: true }).catch(() => undefined);
|
||||
throw error;
|
||||
}
|
||||
}
|
||||
|
||||
/** Reason code attached to failures that must reach the UI verbatim. */
|
||||
function failureReason(error: unknown): string {
|
||||
const reason = (error as { reason?: unknown } | null)?.reason;
|
||||
return typeof reason === "string" && /^[a-z0-9_]{1,64}$/.test(reason) ? reason : "apply_precheck_failed";
|
||||
}
|
||||
|
||||
/**
|
||||
* Persist a real failure reason from the privileged apply path.
|
||||
*
|
||||
* `writeJob`/`updateJob` cannot be used here: their final-state guard
|
||||
* (`WHERE update_jobs.status NOT IN (...)`) protects terminal rows, and it also
|
||||
* makes the runner's own progress writes a no-op once a row is terminal. This
|
||||
* helper issues an independent, guarded UPDATE so the true cause is visible in
|
||||
* the UI instead of the runner's generic health-check message.
|
||||
*/
|
||||
export function failUpdateJobWithReason(
|
||||
sqlite: Database.Database | undefined,
|
||||
jobId: string,
|
||||
message: string,
|
||||
options: { reason?: string } = {},
|
||||
): boolean {
|
||||
if (!sqlite) return false;
|
||||
const safe = safeErrorMessage(message.length ? new Error(message) : new Error("更新失败"));
|
||||
try {
|
||||
return sqlite.transaction(() => {
|
||||
const row = sqlite.prepare("SELECT status, version, request_id AS requestId, admin_id AS adminId FROM update_jobs WHERE id=?").get(jobId) as {
|
||||
status: UpdateJobStatus; version: string; requestId: string | null; adminId: string | null;
|
||||
} | undefined;
|
||||
if (!row || row.status === "completed" || row.status === "cancelled" || row.status === "failed") return false;
|
||||
const now = Date.now();
|
||||
const updated = sqlite.prepare("UPDATE update_jobs SET status='failed', error_message=?, completed_at=COALESCE(completed_at, ?), updated_at=? WHERE id=? AND status=?").run(safe, now, now, jobId, row.status);
|
||||
if (updated.changes !== 1) return false;
|
||||
writeAudit(sqlite, {
|
||||
requestId: row.requestId || randomUUID(),
|
||||
actorAdminId: row.adminId,
|
||||
action: "update.failed",
|
||||
targetType: "update",
|
||||
targetId: jobId,
|
||||
outcome: "failure",
|
||||
before: { status: row.status, version: row.version },
|
||||
after: { status: "failed", version: row.version, reason: options.reason ?? "apply_precheck_failed", error: safe },
|
||||
});
|
||||
return true;
|
||||
})();
|
||||
} catch {
|
||||
// The database may not be open (or the row may not exist) when a request is
|
||||
// rejected during preflight. Losing the diagnostic write must never turn a
|
||||
// clean rejection into a crash.
|
||||
return false;
|
||||
}
|
||||
}
|
||||
|
||||
export async function runUpdate(options: UpdateRunOptions): Promise<UpdateRunResult> {
|
||||
const platform = options.platform ?? detectPlatform();
|
||||
const jobId = options.jobId ?? randomUUID();
|
||||
@@ -437,13 +646,17 @@ export async function applyStagedUpdate(options: {
|
||||
maxBytes?: number;
|
||||
dataBackupMaxBytes?: number;
|
||||
workspaceRoot?: string;
|
||||
/** Expected owner of `workspaceRoot`. Defaults to uid 0 (the installer
|
||||
* provisions `<installPrefix>/.update-work` as root-owned 0700). Tests inject
|
||||
* the current user so the check never depends on `process.getuid()`. */
|
||||
workspaceOwnerUid?: DirectoryOwnerUid;
|
||||
}): Promise<void> {
|
||||
const row = options.sqlite.prepare(`SELECT status, operation, version, platform, release_url AS releaseUrl, asset_name AS assetName, asset_url AS assetUrl, expected_sha256 AS expectedSha256, actual_sha256 AS actualSha256, size_bytes AS sizeBytes FROM update_jobs WHERE id=?`).get(options.jobId) as Record<string, unknown> | undefined;
|
||||
if (!row || row.status !== "staged" || row.operation !== "apply") throw new Error("更新任务未处于待应用状态");
|
||||
if (typeof row.version === "string" && row.version !== options.version) throw new Error("更新版本不一致");
|
||||
const stagedPath = options.workspaceRoot
|
||||
? await validateStagedWorkspacePath(options.stagedPath, options.workspaceRoot)
|
||||
: options.stagedPath;
|
||||
? await canonicalizePrivilegedRoot(options.workspaceRoot, options.workspaceOwnerUid ?? 0, "更新工作目录权限无效")
|
||||
: path.resolve(options.stagedPath);
|
||||
const payload = path.join(stagedPath, "payload");
|
||||
const payloadInfo = await lstat(payload).catch(() => null);
|
||||
if (!payloadInfo?.isDirectory() || payloadInfo.isSymbolicLink()) throw new Error("更新暂存内容无效");
|
||||
@@ -481,7 +694,21 @@ function arg(name: string): string | undefined {
|
||||
return index >= 0 ? process.argv[index + 1] : undefined;
|
||||
}
|
||||
|
||||
export async function main(config: AppConfig = loadConfig()): Promise<void> {
|
||||
/**
|
||||
* Ownership expectations the privileged entry point uses for the two trust
|
||||
* domains it consumes. They are injectable so the apply flow can be exercised
|
||||
* end-to-end from a non-root test process: production always uses the defaults
|
||||
* (root-owned `<installPrefix>/.update-work` and the `dataDir` owner for the
|
||||
* unprivileged staging area) and never consults `process.getuid()`.
|
||||
*/
|
||||
export type UpdateMainOverrides = {
|
||||
/** Expected owner of the unprivileged staging root (`config.stagingDir`). */
|
||||
stagingOwnerUid?: DirectoryOwnerUid;
|
||||
/** Expected owner of the root-only private workspace (`config.updateWorkspaceDir`). */
|
||||
workspaceOwnerUid?: DirectoryOwnerUid;
|
||||
};
|
||||
|
||||
export async function main(config: AppConfig = loadConfig(), overrides: UpdateMainOverrides = {}): Promise<void> {
|
||||
const finalizeJobId = arg("--finalize-job");
|
||||
if (finalizeJobId) {
|
||||
const finalStatus = arg("--finalize-status");
|
||||
@@ -518,7 +745,9 @@ export async function main(config: AppConfig = loadConfig()): Promise<void> {
|
||||
const dataBackupArchive = arg("--data-backup") ?? (request ? path.join(path.dirname(config.dataDir), "tallynote-backups", `data-${request.jobId}.tar.gz`) : undefined);
|
||||
const allowedHosts = process.argv.flatMap((value, index) => value === "--allow-host" && process.argv[index + 1] ? [process.argv[index + 1]!] : []);
|
||||
prepareDataDirectories(config);
|
||||
if (request) await ensurePrivilegedWorkspace(stagingDir);
|
||||
const workspaceOwnerUid = overrides.workspaceOwnerUid ?? 0;
|
||||
const stagingOwnerUid = overrides.stagingOwnerUid ?? await directoryOwnerUid(config.dataDir);
|
||||
if (request) await ensurePrivilegedWorkspace(stagingDir, workspaceOwnerUid);
|
||||
else await mkdir(stagingDir, { recursive: true, mode: 0o700 });
|
||||
// The download phase intentionally runs beside the live app so users keep
|
||||
// access while the archive is fetched and staged. SQLite WAL plus the
|
||||
@@ -528,28 +757,67 @@ export async function main(config: AppConfig = loadConfig()): Promise<void> {
|
||||
const database = openDatabase(config);
|
||||
try {
|
||||
if (request?.operation === "apply") {
|
||||
const staged = database.sqlite.prepare("SELECT status, operation, download_path AS downloadPath, version FROM update_jobs WHERE id=?").get(request.jobId) as { status: UpdateJobStatus; operation: "download" | "apply"; downloadPath: string | null; version: string } | undefined;
|
||||
const staged = database.sqlite.prepare("SELECT status, operation, download_path AS downloadPath, version, expected_sha256 AS expectedSha256 FROM update_jobs WHERE id=?").get(request.jobId) as { status: UpdateJobStatus; operation: "download" | "apply"; downloadPath: string | null; version: string; expectedSha256: string | null } | undefined;
|
||||
if (staged?.status === "staged" && staged.operation === "apply") {
|
||||
if (!staged.downloadPath || staged.version !== request.version) throw new Error("更新暂存任务无效");
|
||||
const root = path.resolve(config.updateWorkspaceDir);
|
||||
const candidate = await validateStagedWorkspacePath(staged.downloadPath, root);
|
||||
await applyStagedUpdate({
|
||||
sqlite: database.sqlite,
|
||||
jobId: request.jobId,
|
||||
version: request.version,
|
||||
stagedPath: candidate,
|
||||
currentDir,
|
||||
currentLink: request.currentLink,
|
||||
releasesDir: request.releasesDir,
|
||||
workspaceRoot: root,
|
||||
...(backupArchive ? { backupArchivePath: backupArchive } : {}),
|
||||
...(dataBackupArchive ? { dataBackupArchivePath: dataBackupArchive } : {}),
|
||||
dataBackupSource: config.dataDir,
|
||||
maxBytes: config.updateMaxBytes,
|
||||
dataBackupMaxBytes: config.maxTotalBytes,
|
||||
});
|
||||
console.log(`更新已切换:${request.version}`);
|
||||
return;
|
||||
// `download_path` is intentionally NOT required here. The runner NULLs
|
||||
// that column as soon as it releases a workspace, so it can never be the
|
||||
// source of truth for a payload that still exists on disk. The job id is
|
||||
// the stable key; the column survives only as a last-resort candidate in
|
||||
// `locateStagedWorkspace`.
|
||||
if (staged.version !== request.version) throw new Error("更新暂存任务无效");
|
||||
let privateWorkspace: string | undefined;
|
||||
try {
|
||||
const source = await locateStagedWorkspace({
|
||||
jobId: request.jobId,
|
||||
downloadPath: staged.downloadPath,
|
||||
stagingRoot: config.stagingDir,
|
||||
expectedUid: stagingOwnerUid,
|
||||
});
|
||||
// Snapshot into the root-only workspace before applying. The web user
|
||||
// keeps write access to the staging tree, so content that is executed
|
||||
// by the privileged process must never live there (TOCTOU).
|
||||
privateWorkspace = await preparePrivateApplyWorkspace({
|
||||
jobId: request.jobId,
|
||||
source,
|
||||
privateRoot: config.updateWorkspaceDir,
|
||||
expectedUid: workspaceOwnerUid,
|
||||
expectedSha256: staged.expectedSha256,
|
||||
});
|
||||
await applyStagedUpdate({
|
||||
sqlite: database.sqlite,
|
||||
jobId: request.jobId,
|
||||
version: request.version,
|
||||
stagedPath: privateWorkspace,
|
||||
workspaceRoot: privateWorkspace,
|
||||
workspaceOwnerUid,
|
||||
currentDir,
|
||||
currentLink: request.currentLink,
|
||||
releasesDir: request.releasesDir,
|
||||
...(backupArchive ? { backupArchivePath: backupArchive } : {}),
|
||||
...(dataBackupArchive ? { dataBackupArchivePath: dataBackupArchive } : {}),
|
||||
dataBackupSource: config.dataDir,
|
||||
maxBytes: config.updateMaxBytes,
|
||||
dataBackupMaxBytes: config.maxTotalBytes,
|
||||
});
|
||||
// The private copy has been consumed by the release switch and the
|
||||
// payload is now the live release, so the web-owned source tree is
|
||||
// redundant. Best-effort cleanup must not fail an applied update.
|
||||
await rm(source, { recursive: true, force: true }).catch(() => undefined);
|
||||
console.log(`更新已切换:${request.version}`);
|
||||
return;
|
||||
} catch (error) {
|
||||
// Covers failures raised before `applyStagedUpdate` took ownership of
|
||||
// the private copy, and the "already committed" case where the failing
|
||||
// path deliberately skips its own cleanup.
|
||||
if (privateWorkspace) await rm(privateWorkspace, { recursive: true, force: true }).catch(() => undefined);
|
||||
// The runner can only report its fixed health-check message. Record the
|
||||
// real pre-flight cause so the UI and the audit trail show why the
|
||||
// update was rejected. A terminal row is only reachable through an
|
||||
// independent guarded UPDATE (`writeJob` refuses to mutate terminal
|
||||
// rows), which is exactly what this helper issues.
|
||||
failUpdateJobWithReason(database.sqlite, request.jobId, safeErrorMessage(error), { reason: failureReason(error) });
|
||||
throw error;
|
||||
}
|
||||
}
|
||||
if (staged && !(staged.status === "queued" && staged.operation === "apply")) throw new Error("更新任务状态无效");
|
||||
// A direct one-click request starts in queued/apply. Older clients do
|
||||
|
||||
@@ -162,6 +162,11 @@ export function loadConfig() {
|
||||
exportsDir: path.join(dataDir, "exports"),
|
||||
migrationsDir: path.join(projectRoot, "migrations"),
|
||||
webDir: path.join(projectRoot, "dist", "web"),
|
||||
// Coarse per-IP request ceiling applied to every /api/* request. It is a
|
||||
// backstop against request floods, not a replacement for the stricter
|
||||
// per-feature limits (login lockout, dangerous-operation re-auth, update
|
||||
// cooldowns), so the default is deliberately generous.
|
||||
apiRateLimitPerMinute: integerEnv("TALLYNOTE_RATE_LIMIT_PER_MINUTE", 600),
|
||||
maxFileBytes: integerEnv("TALLYNOTE_MAX_FILE_MB", 20) * 1024 * 1024,
|
||||
maxFilesPerRequest: integerEnv("TALLYNOTE_MAX_FILES_PER_REQUEST", 20),
|
||||
maxRecordBytes: integerEnv("TALLYNOTE_MAX_RECORD_MB", 100) * 1024 * 1024,
|
||||
|
||||
@@ -0,0 +1,80 @@
|
||||
/**
|
||||
* In-memory, per-key request limiter used as a coarse anti-flood backstop for
|
||||
* the whole HTTP API.
|
||||
*
|
||||
* The semantics are a fixed window per key: the first request of a window
|
||||
* starts the clock, every later request in the same window increments the
|
||||
* counter, and an expired window is reset on the next request. This mirrors
|
||||
* the `login_attempts` window logic already used for login lockouts
|
||||
* (`server/app.ts`), but it never touches the database: a rate limit decision
|
||||
* must stay cheap enough to run on every request.
|
||||
*
|
||||
* Precise controls (per-IP login lockout, dangerous-operation re-auth) remain
|
||||
* in place on top of this limiter; it only stops a client from issuing an
|
||||
* abusive number of requests across all endpoints.
|
||||
*/
|
||||
|
||||
export type RateLimiterOptions = {
|
||||
/** Maximum number of requests allowed per key inside one window. */
|
||||
limit: number;
|
||||
/** Window length in milliseconds. */
|
||||
windowMs: number;
|
||||
/** Injectable clock so tests can advance time without waiting. */
|
||||
now?: () => number;
|
||||
};
|
||||
|
||||
export type RateLimitDecision = {
|
||||
allowed: boolean;
|
||||
/** Seconds the caller should wait before retrying; 0 when allowed. */
|
||||
retryAfterSeconds: number;
|
||||
};
|
||||
|
||||
type Bucket = {
|
||||
count: number;
|
||||
windowStart: number;
|
||||
};
|
||||
|
||||
/** Run a full sweep every N checks instead of on every call. */
|
||||
const SWEEP_INTERVAL_CHECKS = 1000;
|
||||
|
||||
export function createRateLimiter(options: RateLimiterOptions) {
|
||||
const { limit, windowMs } = options;
|
||||
if (!Number.isInteger(limit) || limit < 1) throw new Error("rate limit 必须是大于等于 1 的整数");
|
||||
if (!Number.isInteger(windowMs) || windowMs < 1) throw new Error("rate limit 窗口必须是大于等于 1 的整数毫秒数");
|
||||
const now = options.now ?? Date.now;
|
||||
const buckets = new Map<string, Bucket>();
|
||||
let checksSinceSweep = 0;
|
||||
|
||||
return {
|
||||
check(key: string): RateLimitDecision {
|
||||
const current = now();
|
||||
let bucket = buckets.get(key);
|
||||
// A key that is unknown or whose window has already elapsed starts a
|
||||
// fresh window. This also recycles the single key being hit, so an
|
||||
// idle client never leaves a stale counter behind.
|
||||
if (!bucket || current - bucket.windowStart >= windowMs) {
|
||||
bucket = { count: 0, windowStart: current };
|
||||
buckets.set(key, bucket);
|
||||
}
|
||||
// Bounds long-running memory growth: keys that stopped sending traffic
|
||||
// are dropped by an amortized periodic sweep rather than on every call.
|
||||
if (++checksSinceSweep >= SWEEP_INTERVAL_CHECKS) {
|
||||
checksSinceSweep = 0;
|
||||
for (const [candidateKey, candidate] of buckets) {
|
||||
if (current - candidate.windowStart >= windowMs) buckets.delete(candidateKey);
|
||||
}
|
||||
}
|
||||
if (bucket.count >= limit) {
|
||||
return { allowed: false, retryAfterSeconds: Math.max(1, Math.ceil((bucket.windowStart + windowMs - current) / 1000)) };
|
||||
}
|
||||
bucket.count += 1;
|
||||
return { allowed: true, retryAfterSeconds: 0 };
|
||||
},
|
||||
/** Number of tracked keys; used to observe lazy cleanup. */
|
||||
size(): number {
|
||||
return buckets.size;
|
||||
},
|
||||
};
|
||||
}
|
||||
|
||||
export type RateLimiter = ReturnType<typeof createRateLimiter>;
|
||||
+139
-13
@@ -1,4 +1,4 @@
|
||||
import { lstatSync, realpathSync, readFileSync, unlinkSync } from "node:fs";
|
||||
import { lstatSync, readdirSync, realpathSync, readFileSync, unlinkSync } from "node:fs";
|
||||
import { chmod, lstat, mkdir, mkdtemp, rename, rm, writeFile } from "node:fs/promises";
|
||||
import path from "node:path";
|
||||
import { createPublicKey, randomUUID, verify as verifySignature } from "node:crypto";
|
||||
@@ -40,6 +40,20 @@ export const ACTIVE_UPDATE_STATUSES: readonly UpdateJobStatus[] = [
|
||||
// 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.
|
||||
/**
|
||||
* Conflict predicate for "another update is already running".
|
||||
*
|
||||
* A row that is `staged` with `operation='download'` is a finished artifact
|
||||
* waiting for an explicit apply, not a running task: the privileged runner only
|
||||
* starts working after the apply request is written. It must therefore not
|
||||
* block a new download. Real in-flight work (queued/downloading/verifying and
|
||||
* the apply phases) remains protected, which is what keeps the apply path's
|
||||
* concurrency guard intact.
|
||||
*
|
||||
* The SQL fragment expects ACTIVE_UPDATE_STATUSES bound as positional params.
|
||||
*/
|
||||
export const ACTIVE_UPDATE_CONFLICT_SQL = `status IN (${ACTIVE_UPDATE_STATUSES.map(() => "?").join(",")}) AND NOT (status='staged' AND operation='download')`;
|
||||
|
||||
export const ORPHANED_UPDATE_TIMEOUT_MS = 5 * 60 * 1000;
|
||||
export const QUEUED_UPDATE_TIMEOUT_MS = 25 * 1000;
|
||||
|
||||
@@ -447,6 +461,61 @@ export function currentReleaseVersion(config: AppConfig): string | null {
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Absolute paths that can hold a job's staging workspace. The web download flow
|
||||
* always creates `update-<jobId>`; the privileged runner may additionally use a
|
||||
* `mkdtemp` variant named `update-<jobId>-XXXXXX`.
|
||||
*
|
||||
* `update_jobs.download_path` is deliberately NOT used to rebuild these paths:
|
||||
* it held a bare basename while a download was in flight (rows written by older
|
||||
* versions still store that basename) and the privileged runner NULLs the column
|
||||
* after finalizing a row. Rebuilding from it could delete an unrelated staging
|
||||
* entry that merely shares the basename.
|
||||
*/
|
||||
function jobWorkspaceCandidates(stagingDir: string, jobId: string): string[] {
|
||||
// Job ids are UUIDs; reject anything that could escape the staging root.
|
||||
if (!jobId || jobId !== path.basename(jobId) || jobId.includes("..")) return [];
|
||||
const stagingRoot = path.resolve(stagingDir);
|
||||
const prefix = `update-${jobId}`;
|
||||
const names = [prefix];
|
||||
try {
|
||||
for (const entry of readdirSync(stagingRoot)) {
|
||||
if (entry.startsWith(`${prefix}-`)) names.push(entry);
|
||||
}
|
||||
} catch {
|
||||
// A missing or unreadable staging directory still leaves the fixed-name
|
||||
// candidate, which is what the web download path uses.
|
||||
}
|
||||
return names.map((name) => path.join(stagingRoot, name));
|
||||
}
|
||||
|
||||
function isDirectoryNotSymlink(target: string): boolean {
|
||||
try {
|
||||
const info = lstatSync(target);
|
||||
return info.isDirectory() && !info.isSymbolicLink();
|
||||
} catch {
|
||||
return false;
|
||||
}
|
||||
}
|
||||
|
||||
/** True while at least one staging workspace for the job still exists. */
|
||||
export function jobWorkspaceExists(stagingDir: string, jobId: string): boolean {
|
||||
return jobWorkspaceCandidates(stagingDir, jobId).some(isDirectoryNotSymlink);
|
||||
}
|
||||
|
||||
/**
|
||||
* Remove every staging workspace owned by a job. Deletion is awaited so callers
|
||||
* (and tests) observe a settled filesystem when they return.
|
||||
*/
|
||||
async function removeJobWorkspaces(stagingDir: string, jobId: string): Promise<void> {
|
||||
for (const candidate of jobWorkspaceCandidates(stagingDir, jobId)) {
|
||||
const info = await lstat(candidate).catch(() => null);
|
||||
// Only real directories are removed; a symlink is never followed.
|
||||
if (!info?.isDirectory() || info.isSymbolicLink()) continue;
|
||||
await rm(candidate, { recursive: true, force: true }).catch(() => undefined);
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Release an update row left behind after its privileged runner lease expired.
|
||||
* This is deliberately conservative: staged downloads remain available for an
|
||||
@@ -558,6 +627,36 @@ export function reconcileOrphanedUpdateJobs(database: Database.Database, config:
|
||||
// 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") {
|
||||
// A staged row that lost its payload (the staging janitor removes
|
||||
// `update-*` entries after 24h, and a manual cleanup has the same effect)
|
||||
// can never be applied or completed. Report it instead of leaving a
|
||||
// permanently actionable row that fails at apply time.
|
||||
if (!matchingFreshRequest && !matchingFreshState && !jobWorkspaceExists(config.stagingDir, row.id)) {
|
||||
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' 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, operation: row.operation, version: row.version },
|
||||
after: { status: "failed", version: row.version, reason: "staged_workspace_missing" },
|
||||
});
|
||||
return true;
|
||||
})();
|
||||
if (changed) {
|
||||
reconciled += 1;
|
||||
reconciledIds.add(row.id);
|
||||
}
|
||||
continue;
|
||||
}
|
||||
if (row.operation !== "apply" || matchingFreshRequest || matchingFreshState) continue;
|
||||
const changed = database.transaction(() => {
|
||||
const result = database.prepare(`
|
||||
@@ -633,23 +732,39 @@ export function reconcileOrphanedUpdateJobs(database: Database.Database, config:
|
||||
return reconciled;
|
||||
}
|
||||
|
||||
export function cancelUpdateJob(
|
||||
/**
|
||||
* Rows an administrator may cancel from the web UI.
|
||||
*
|
||||
* `staged` is cancellable only while the row still belongs to the download
|
||||
* stage. A staged row whose operation is already `apply` has been handed to the
|
||||
* privileged runner (stop/backup/switch) and must not be interrupted here.
|
||||
*/
|
||||
const CANCELLABLE_JOB_SQL = "(status IN ('queued', 'downloading') OR (status='staged' AND operation='download'))";
|
||||
|
||||
export function isCancellableUpdateJob(status: UpdateJobStatus, operation: string): boolean {
|
||||
if (status === "queued" || status === "downloading") return true;
|
||||
return status === "staged" && operation === "download";
|
||||
}
|
||||
|
||||
export async function cancelUpdateJob(
|
||||
database: Database.Database,
|
||||
config: AppConfig,
|
||||
adminId: string,
|
||||
requestId: string,
|
||||
jobId?: string,
|
||||
): { cancelled: boolean; message?: string } {
|
||||
): Promise<{ cancelled: boolean; message?: string }> {
|
||||
type CancelRow = { id: string; status: UpdateJobStatus; operation: string; version: string; adminId: string | null };
|
||||
const columns = "id, status, operation, version, admin_id AS adminId";
|
||||
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;
|
||||
? database.prepare(`SELECT ${columns} FROM update_jobs WHERE id=? AND admin_id=?`).get(jobId, adminId) as CancelRow | undefined
|
||||
: database.prepare(`SELECT ${columns} FROM update_jobs WHERE admin_id=? AND ${CANCELLABLE_JOB_SQL} ORDER BY created_at DESC LIMIT 1`).get(adminId) as CancelRow | undefined;
|
||||
|
||||
if (!job) return { cancelled: false, message: "当前没有处于等待调度或下载中的更新任务" };
|
||||
if (job.status !== "queued" && job.status !== "downloading") return { cancelled: false, message: "任务已进入就绪或切换阶段,无法取消" };
|
||||
if (!isCancellableUpdateJob(job.status, job.operation)) 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);
|
||||
const result = database.prepare(`UPDATE update_jobs SET status='cancelled', error_message='已手动取消更新', completed_at=?, updated_at=? WHERE id=? AND admin_id=? AND ${CANCELLABLE_JOB_SQL}`).run(now, now, job.id, adminId);
|
||||
if (result.changes !== 1) return false;
|
||||
writeAudit(database, {
|
||||
requestId,
|
||||
@@ -668,10 +783,11 @@ export function cancelUpdateJob(
|
||||
// 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(() => {});
|
||||
}
|
||||
// Locate the workspace by job id. `download_path` is not a reliable source
|
||||
// (older rows hold a bare archive basename and the runner NULLs the column
|
||||
// after finalizing), and a basename lookup could delete an unrelated entry.
|
||||
// The await keeps the caller from racing a still-running download writer.
|
||||
await removeJobWorkspaces(config.stagingDir, job.id);
|
||||
return { cancelled: true };
|
||||
}
|
||||
return { cancelled: false, message: "取消失败,任务状态可能已改变" };
|
||||
@@ -706,9 +822,13 @@ export async function downloadAndStageUpdate(
|
||||
|
||||
// Claim the job: transition queued -> downloading. If the job was
|
||||
// cancelled or claimed by another caller, abort immediately.
|
||||
// `download_path` always holds an absolute workspace path, both while the
|
||||
// download runs and after the job is staged. Callers must not derive paths
|
||||
// from it (the privileged runner NULLs it once it finalizes the row), but a
|
||||
// single semantic keeps the column debuggable.
|
||||
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);
|
||||
).run(Date.now(), Date.now(), workspace, Date.now(), jobId);
|
||||
if (claim.changes !== 1) return;
|
||||
|
||||
const progressStartedAt = Date.now();
|
||||
@@ -764,7 +884,13 @@ export async function downloadAndStageUpdate(
|
||||
const staged = database.prepare(
|
||||
"UPDATE update_jobs SET status='staged', operation='download', 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
|
||||
if (staged.changes !== 1) {
|
||||
// The row was cancelled or claimed elsewhere (status no longer
|
||||
// verifying/downloading). This process owns the workspace it created, so
|
||||
// remove it instead of leaking the payload into the staging directory.
|
||||
await rm(workspace, { recursive: true, force: true }).catch(() => undefined);
|
||||
return;
|
||||
}
|
||||
|
||||
writeAudit(database, {
|
||||
requestId: `download:${jobId}`,
|
||||
|
||||
Reference in New Issue
Block a user