diff --git a/.env.example b/.env.example index f838443..e54e7c7 100644 --- a/.env.example +++ b/.env.example @@ -31,6 +31,8 @@ TALLYNOTE_INSTALL_PREFIX=./ TALLYNOTE_UPDATE_METADATA_URL=https://git.awaioi.com/api/v1/repos/awaioi/TallyNote/releases/latest TALLYNOTE_UPDATE_ALLOWED_HOSTS=git.awaioi.com TALLYNOTE_UPDATE_MAX_MB=512 +# Per-request timeout for update metadata, checksums, signatures, and archives. +TALLYNOTE_UPDATE_TIMEOUT_SECONDS=30 # SHA-256 is always required. Detached Ed25519 signatures are optional; set # this to true only when a root-managed public key is configured below. TALLYNOTE_UPDATE_REQUIRE_SIGNATURE=false diff --git a/package.json b/package.json index 9bc31df..6e4e08b 100644 --- a/package.json +++ b/package.json @@ -1,6 +1,6 @@ { "name": "tallynote", - "version": "1.2.3", + "version": "1.2.4", "private": true, "type": "module", "packageManager": "pnpm@9.0.6", diff --git a/scripts/tallynote-update-runner.sh b/scripts/tallynote-update-runner.sh index 43217a0..afb6885 100755 --- a/scripts/tallynote-update-runner.sh +++ b/scripts/tallynote-update-runner.sh @@ -10,6 +10,8 @@ DATA_DIR=${TALLYNOTE_DATA_DIR:-/var/lib/tallynote} REQUEST_FILE="$DATA_DIR/update-request.json" CURRENT_LINK="$PREFIX/current" STATE_FILE="$PREFIX/.update-state" +LOCK_FILE="$PREFIX/.update-runner.lock" +RUNNER_LOG="$PREFIX/.update-runner.log" SERVICE_NAME=${TALLYNOTE_SERVICE_NAME:-tallynote.service} HOST=${TALLYNOTE_HOST:-127.0.0.1} PORT=${TALLYNOTE_PORT:-3000} @@ -19,13 +21,72 @@ if [[ "$HEALTH_HOST" == :: ]]; then HEALTH_HOST=::1; fi if [[ "$HEALTH_HOST" == *:* && "$HEALTH_HOST" != \[* ]]; then HEALTH_HOST="[$HEALTH_HOST]"; fi die() { printf 'tallynote update runner: %s\n' "$*" >&2; exit 1; } + +# The runner may exit during any of the checks below. Install its EXIT cleanup +# before doing privileged preflight so a partial invocation never leaves a +# heartbeat or lock behind. +STATE_CREATED=0 +heartbeat_pid='' +heartbeat_owner=$$ +RUNNER_LOCK_FD=9 +RUNNER_LOCK_MODE='' +stop_heartbeat() { + if [[ -n "$heartbeat_pid" ]]; then + kill "$heartbeat_pid" 2>/dev/null || true + wait "$heartbeat_pid" 2>/dev/null || true + heartbeat_pid='' + fi +} +# shellcheck disable=SC2329 # invoked indirectly by the EXIT trap +release_runner_lock() { + if [[ "$RUNNER_LOCK_MODE" == flock ]]; then + flock -u "$RUNNER_LOCK_FD" 2>/dev/null || true + eval "exec ${RUNNER_LOCK_FD}>&-" 2>/dev/null || true + elif [[ "$RUNNER_LOCK_MODE" == mkdir ]]; then + rmdir -- "$LOCK_FILE.d" 2>/dev/null || true + fi +} +# shellcheck disable=SC2329 # invoked indirectly by the EXIT trap +early_cleanup() { + local result=$? + stop_heartbeat + if (( result != 0 )); then + # A preflight failure happens before the normal phase-specific trap is + # installed. Remove only the one-shot request marker; never remove an + # existing recovery marker unless this invocation created it. + rm -f -- "$REQUEST_FILE" 2>/dev/null || true + if (( STATE_CREATED == 1 )); then rm -f -- "$STATE_FILE" 2>/dev/null || true; fi + fi + release_runner_lock + return "$result" +} +trap early_cleanup EXIT + [[ ${EUID:-$(id -u)} -eq 0 ]] || die 'must run as root' +[[ -d "$PREFIX" ]] || die 'install prefix is missing' +if command -v flock >/dev/null 2>&1; then + exec 9>"$LOCK_FILE" || die '无法打开更新运行锁' + flock -n "$RUNNER_LOCK_FD" || exit 0 + RUNNER_LOCK_MODE=flock +else + # macOS development fixtures do not ship util-linux; retain an atomic lock + # fallback there while Linux production uses flock above. + mkdir "$LOCK_FILE.d" 2>/dev/null || exit 0 + RUNNER_LOCK_MODE='mkdir' +fi [[ -f "$REQUEST_FILE" || -f "$STATE_FILE" ]] || exit 0 [[ -L "$CURRENT_LINK" ]] || die 'current release link is missing' old_target=$(readlink -f -- "$CURRENT_LINK") [[ "$old_target" == "$PREFIX/releases/"* && -d "$old_target" ]] || die 'current release target is invalid' +# Capture the service state before any download/apply work. The value is +# persisted in the recovery marker so a later runner process can restore the +# operator's original state after a crash (the service is normally inactive by +# the time recovery starts). +was_active=0 +if systemctl is-active --quiet "$SERVICE_NAME"; then was_active=1; fi + request_operation='apply' if [[ -f "$REQUEST_FILE" && ! -L "$REQUEST_FILE" ]]; then request_operation=$(sed -n 's/.*"operation"[[:space:]]*:[[:space:]]*"\(download\|apply\)".*/\1/p' "$REQUEST_FILE" | head -n 1) @@ -39,15 +100,11 @@ job_id='' if [[ -f "$REQUEST_FILE" && ! -L "$REQUEST_FILE" ]]; then job_id=$(sed -n 's/.*"jobId"[[:space:]]*:[[:space:]]*"\([0-9a-f-]*\)".*/\1/p' "$REQUEST_FILE" | head -n 1) fi -STATE_CREATED=0 -heartbeat_pid='' -heartbeat_owner=$$ - write_recovery_state() { local phase=$1 temporary temporary="$PREFIX/.update-state-$$-${RANDOM}.tmp" [[ ! -e "$temporary" && ! -L "$temporary" ]] || return 1 - printf 'job_id=%s\nold_target=%s\nphase=%s\n' "$job_id" "$old_target" "$phase" > "$temporary" + printf 'job_id=%s\nold_target=%s\nphase=%s\ninitial_active=%s\n' "$job_id" "$old_target" "$phase" "$was_active" > "$temporary" chmod 600 "$temporary" mv -Tf -- "$temporary" "$STATE_FILE" STATE_CREATED=1 @@ -59,14 +116,6 @@ clear_recovery_state() { STATE_CREATED=0 } -stop_heartbeat() { - if [[ -n "$heartbeat_pid" ]]; then - kill "$heartbeat_pid" 2>/dev/null || true - wait "$heartbeat_pid" 2>/dev/null || true - heartbeat_pid='' - fi -} - heartbeat() { # Keep the lease fresh during long downloads/backups, but stop on a hard # runner kill so an orphaned child cannot keep the recovery marker alive. @@ -86,9 +135,11 @@ start_heartbeat() { # This trap covers failures before the normal apply cleanup trap is installed, # including a missing runtime, an invalid current link, and a failed service # stop. It deliberately does not remove a pre-existing recovery marker. +# shellcheck disable=SC2329 # invoked indirectly by the EXIT trap preflight_cleanup() { local result=$? stop_heartbeat + release_runner_lock if (( result != 0 )); then rm -f -- "$REQUEST_FILE" 2>/dev/null || true if (( STATE_CREATED == 1 )); then clear_recovery_state || true; fi @@ -97,6 +148,33 @@ preflight_cleanup() { } trap preflight_cleanup EXIT +DOWNLOAD_TIMEOUT_SECONDS=${TALLYNOTE_UPDATE_DOWNLOAD_TIMEOUT_SECONDS:-${TALLYNOTE_UPDATE_RUNNER_DOWNLOAD_TIMEOUT_SECONDS:-1800}} +APPLY_TIMEOUT_SECONDS=${TALLYNOTE_UPDATE_APPLY_TIMEOUT_SECONDS:-${TALLYNOTE_UPDATE_RUNNER_APPLY_TIMEOUT_SECONDS:-1800}} +FINALIZE_TIMEOUT_SECONDS=${TALLYNOTE_UPDATE_FINALIZE_TIMEOUT_SECONDS:-${TALLYNOTE_UPDATE_RUNNER_FINALIZE_TIMEOUT_SECONDS:-30}} +TIMEOUT_BIN=$(command -v timeout || true) + +run_update_cli() { + local node=$1 timeout_seconds=$2 label=$3 result + shift 3 + [[ "$timeout_seconds" =~ ^[1-9][0-9]*$ ]] || die "${label} timeout must be a positive integer" + { + printf '\n[%s] %s (timeout=%ss)\ncommand:' "$(date -u '+%Y-%m-%dT%H:%M:%SZ')" "$label" "$timeout_seconds" + printf ' %q' "$node" "$CURRENT_LINK/dist/server/cli/update.js" "$@" + printf '\n' + } >>"$RUNNER_LOG" + if [[ -n "$TIMEOUT_BIN" ]]; then + "$TIMEOUT_BIN" --foreground --signal=TERM --kill-after=10s "${timeout_seconds}s" \ + "$node" "$CURRENT_LINK/dist/server/cli/update.js" "$@" >>"$RUNNER_LOG" 2>&1 + result=$? + elif "$node" "$CURRENT_LINK/dist/server/cli/update.js" "$@" >>"$RUNNER_LOG" 2>&1; then + result=0 + else + result=$? + fi + printf '[%s] %s exited with status %s\n' "$(date -u '+%Y-%m-%dT%H:%M:%SZ')" "$label" "$result" >>"$RUNNER_LOG" + return "$result" +} + # Downloading is intentionally handled while the main service remains up. # The CLI persists the validated payload under the root-owned workspace and # leaves the job staged for a later apply request. @@ -107,6 +185,7 @@ if [[ "$request_operation" == download ]]; then clear_recovery_state || die '无法清理上一次下载状态' fi write_recovery_state download || die '无法写入更新恢复状态' +# shellcheck disable=SC2329 # invoked indirectly by the EXIT trap cleanup_download() { local result=$? stop_heartbeat @@ -116,6 +195,7 @@ if [[ "$request_operation" == download ]]; then rm -f -- "$REQUEST_FILE" 2>/dev/null || true fi clear_recovery_state || true + release_runner_lock return "$result" } trap cleanup_download EXIT @@ -128,7 +208,7 @@ if [[ "$request_operation" == download ]]; then cli="$CURRENT_LINK/dist/server/cli/update.js" [[ -f "$cli" ]] || die 'update CLI not found in current release' set +e - "$node_bin" "$cli" --request-file "$REQUEST_FILE" + run_update_cli "$node_bin" "$DOWNLOAD_TIMEOUT_SECONDS" download --request-file "$REQUEST_FILE" download_result=$? set -e if (( download_result != 0 )); then @@ -139,7 +219,7 @@ if [[ "$request_operation" == download ]]; then download_job_id=$(sed -n 's/.*"jobId"[[:space:]]*:[[:space:]]*"\([0-9a-f-]*\)".*/\1/p' "$REQUEST_FILE" | head -n 1) if [[ "$download_job_id" =~ ^[0-9a-f-]{36}$ ]]; then for _ in 1 2 3; do - if "$node_bin" "$cli" --finalize-job "$download_job_id" --finalize-status failed --message '更新下载失败' >/dev/null 2>&1; then break; fi + if run_update_cli "$node_bin" "$FINALIZE_TIMEOUT_SECONDS" finalize-download --finalize-job "$download_job_id" --finalize-status failed --message '更新下载失败'; then break; fi sleep 1 done fi @@ -150,12 +230,11 @@ if [[ "$request_operation" == download ]]; then exit 0 fi -was_active=0 -if systemctl is-active --quiet "$SERVICE_NAME"; then was_active=1; fi # shellcheck disable=SC2329 # invoked indirectly by the EXIT trap below restore_initial_service() { local result=$? stop_heartbeat + release_runner_lock if (( result != 0 )); then rm -f -- "$REQUEST_FILE" 2>/dev/null || true if (( STATE_CREATED == 1 )); then clear_recovery_state || true; fi @@ -175,11 +254,11 @@ finalize_state_job() { local node=$1 status=$2 state_job=$3 [[ "$state_job" =~ ^[0-9a-f-]{36}$ && -n "$node" ]] || return 1 [[ -f "$CURRENT_LINK/dist/server/cli/update.js" ]] || return 1 - "$node" "$CURRENT_LINK/dist/server/cli/update.js" --finalize-job "$state_job" --finalize-status "$status" --message '新版本健康检查失败,已恢复上一版本' >/dev/null 2>&1 + run_update_cli "$node" "$FINALIZE_TIMEOUT_SECONDS" finalize-recovery --finalize-job "$state_job" --finalize-status "$status" --message '新版本健康检查失败,已恢复上一版本' } recover_stale_state() { - local state_job state_old state_phase current_target recovery_node rollback_link state_mode state_uid + local state_job state_old state_phase state_initial_active current_target recovery_node rollback_link state_mode state_uid [[ -f "$STATE_FILE" && ! -L "$STATE_FILE" ]] || die 'update state file is invalid' state_uid=$(stat -c '%u' "$STATE_FILE" 2>/dev/null || stat -f '%u' "$STATE_FILE") state_mode=$(stat -c '%a' "$STATE_FILE" 2>/dev/null || stat -f '%Lp' "$STATE_FILE") @@ -187,8 +266,16 @@ recover_stale_state() { state_job=$(sed -n 's/^job_id=//p' "$STATE_FILE" | head -n 1) state_old=$(sed -n 's/^old_target=//p' "$STATE_FILE" | head -n 1) state_phase=$(sed -n 's/^phase=//p' "$STATE_FILE" | head -n 1) + state_initial_active=$(sed -n 's/^initial_active=//p' "$STATE_FILE" | head -n 1) [[ "$state_job" =~ ^[0-9a-f-]{36}$ ]] || die 'update state job id is invalid' [[ "$state_old" == "$PREFIX/releases/"* && -d "$state_old" && ! -L "$state_old" ]] || die 'update state target is invalid' + if [[ -z "$state_initial_active" ]]; then + # Markers from older releases did not persist this field. Preserve their + # historical conservative behavior instead of rejecting recovery. + state_initial_active=0 + fi + [[ "$state_initial_active" == 0 || "$state_initial_active" == 1 ]] || die 'update state initial service state is invalid' + was_active=$state_initial_active current_target=$(readlink -f -- "$CURRENT_LINK" 2>/dev/null || true) if [[ "$state_phase" == download && "$current_target" == "$state_old" ]]; then # Downloading never changes the active release. If the runner was killed @@ -312,7 +399,7 @@ finalize_failed_job() { # Give SQLite a moment to release a transient lock before declaring the # recovery itself failed. for _ in 1 2 3; do - if "$old_node" "$CURRENT_LINK/dist/server/cli/update.js" --finalize-job "$job_id" --finalize-status failed --message '新版本健康检查失败,已恢复上一版本' >/dev/null 2>&1; then + if run_update_cli "$old_node" "$FINALIZE_TIMEOUT_SECONDS" finalize-failed --finalize-job "$job_id" --finalize-status failed --message '新版本健康检查失败,已恢复上一版本'; then return 0 fi sleep 1 @@ -323,7 +410,7 @@ finalize_failed_job() { finalize_completed_job() { [[ "$job_id" =~ ^[0-9a-f-]{36}$ ]] || return 0 [[ -n "$final_node" ]] || return 1 - "$final_node" "$CURRENT_LINK/dist/server/cli/update.js" --finalize-job "$job_id" --finalize-status completed >/dev/null 2>&1 + run_update_cli "$final_node" "$FINALIZE_TIMEOUT_SECONDS" finalize-completed --finalize-job "$job_id" --finalize-status completed } # shellcheck disable=SC2329 # invoked indirectly by the EXIT trap below @@ -347,6 +434,7 @@ cleanup_after_update() { else systemctl stop "$SERVICE_NAME" || true fi + release_runner_lock return "$result" } trap cleanup_after_update EXIT @@ -358,7 +446,7 @@ cli="$CURRENT_LINK/dist/server/cli/update.js" [[ -f "$cli" ]] || die 'update CLI not found in current release' set +e -"$node_bin" "$cli" --request-file "$REQUEST_FILE" --defer-completion +run_update_cli "$node_bin" "$APPLY_TIMEOUT_SECONDS" apply --request-file "$REQUEST_FILE" --defer-completion update_result=$? set -e if (( update_result != 0 )); then diff --git a/server/cli/update.ts b/server/cli/update.ts index 8f221e3..ab65384 100644 --- a/server/cli/update.ts +++ b/server/cli/update.ts @@ -163,7 +163,10 @@ function writeJob(sqlite: Database.Database | undefined, jobId: string, values: requested_at=COALESCE(excluded.requested_at, update_jobs.requested_at), started_at=COALESCE(excluded.started_at, update_jobs.started_at), operation=excluded.operation, - status=CASE WHEN update_jobs.status='cancelled' THEN update_jobs.status ELSE excluded.status END, + -- Terminal rows are immutable from the runner's ordinary progress + -- writes. In particular, a stale/replayed request must not resurrect a + -- failed job as queued/downloading/etc. + status=CASE WHEN update_jobs.status IN ('cancelled', 'failed', 'completed') THEN update_jobs.status ELSE excluded.status END, version=excluded.version, platform=excluded.platform, release_url=COALESCE(excluded.release_url, update_jobs.release_url), asset_name=COALESCE(excluded.asset_name, update_jobs.asset_name), @@ -176,6 +179,11 @@ function writeJob(sqlite: Database.Database | undefined, jobId: string, values: error_message=COALESCE(excluded.error_message, update_jobs.error_message), updated_at=excluded.updated_at, completed_at=COALESCE(excluded.completed_at, update_jobs.completed_at) + -- Do not let a delayed runner replay overwrite any field on a terminal + -- row. The predicate is part of the same SQLite upsert, so a finalizer + -- racing this write still wins atomically instead of leaving a partially + -- mutated completed/failed/cancelled record. + WHERE update_jobs.status NOT IN ('cancelled', 'failed', 'completed') `).run( jobId, values.adminId ?? null, @@ -230,6 +238,7 @@ async function resolveRelease(options: UpdateRunOptions, platform: ReturnType { - sqlite.prepare("UPDATE update_jobs SET status=?, error_message=?, completed_at=?, updated_at=? WHERE id=?").run(status, safeFailureMessage, now, now, jobId); + const row = sqlite.prepare(` + SELECT id, status, version, platform, admin_id AS adminId, + request_id AS requestId, session_hash AS sessionHash + FROM update_jobs WHERE id=? + `).get(jobId) as { id: string; status: UpdateJobStatus; version: string; platform: string; adminId: string | null; requestId: string | null; sessionHash: string | null } | undefined; + if (!row) throw new Error("更新任务不存在"); + // A failed finalization can be retried by the runner. Once it has been + // committed, make retries a no-op so the error and audit trail stay stable. + if (row.status === status) return; + // A completed release is terminal. A delayed recovery process must never + // be able to downgrade it to failed after the service was healthy. + if (row.status === "completed" && status === "failed") throw new Error("更新任务状态不允许完成"); + const canComplete = row.status === "applying" || row.status === "completed"; + const canFail = ACTIVE_UPDATE_STATUSES.includes(row.status) || row.status === "completed" || row.status === "failed"; + if (status === "completed" ? !canComplete : !canFail) throw new Error("更新任务状态不允许完成"); + const now = Date.now(); + const safeFailureMessage = status === "failed" + ? (message?.trim() ? safeErrorMessage(new Error(message)) : "新版本健康检查失败,已恢复上一版本") + : null; + const result = sqlite.prepare("UPDATE update_jobs SET status=?, error_message=?, completed_at=?, updated_at=? WHERE id=? AND status=?").run(status, safeFailureMessage, now, now, jobId, row.status); + if (result.changes !== 1) return; writeAudit(sqlite, { requestId: row.requestId || randomUUID(), actorAdminId: row.adminId, @@ -569,6 +587,7 @@ export async function main(config: AppConfig = loadConfig()): Promise { ...(dataBackupArchive ? { dataBackupArchivePath: dataBackupArchive, dataBackupSource: config.dataDir } : {}), ...((arg("--backup-dir")) ? { backupDir: arg("--backup-dir") } : {}), allowedHosts: allowedHosts.length ? allowedHosts : config.updateAllowedHosts, + timeoutMs: config.updateTimeoutMs, maxBytes: config.updateMaxBytes, dataBackupMaxBytes: config.maxTotalBytes, currentVersion: config.appVersion, diff --git a/server/config.ts b/server/config.ts index 4440cb9..0addafc 100644 --- a/server/config.ts +++ b/server/config.ts @@ -147,6 +147,7 @@ export function loadConfig() { // as 0700 root:root; development/test callers may override --staging-dir. updateWorkspaceDir: path.join(installPrefix, ".update-work"), updateMaxBytes: integerEnv("TALLYNOTE_UPDATE_MAX_MB", 512) * 1024 * 1024, + updateTimeoutMs: integerEnv("TALLYNOTE_UPDATE_TIMEOUT_SECONDS", 30) * 1000, // Update checks hit an external release endpoint. Keep a short local // cooldown so an authenticated account cannot turn the endpoint into an // outbound request flood; set to 0 only for controlled test environments. diff --git a/server/update-service.ts b/server/update-service.ts index 39fe253..1f692e8 100644 --- a/server/update-service.ts +++ b/server/update-service.ts @@ -154,19 +154,19 @@ function signatureAssetFor(metadata: ReleaseMetadata, sums: ReleaseAsset): Relea export async function attachSidecarHash( metadata: ReleaseMetadata, asset: ReleaseAsset, - options: { allowedHosts: readonly string[]; baseUrl: string; maxBytes: number; publicKey?: string | undefined; requireSignature?: boolean | undefined }, + 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) }); + 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 }); + const signature = await fetchReleaseBytes(signatureAsset.url, { allowedHosts: options.allowedHosts, baseUrl: options.baseUrl, maxBytes: 64 * 1024, timeoutMs: options.timeoutMs }); signatureVerified = verifyReleaseSignature(content, signature, options.publicKey); } } @@ -183,6 +183,7 @@ function policy(config: AppConfig) { allowedHosts: config.updateAllowedHosts, baseUrl: config.updateMetadataUrl, maxRedirects: 3, + timeoutMs: config.updateTimeoutMs, } as const; } @@ -222,6 +223,7 @@ export async function checkForUpdate(database: Database.Database, config: AppCon allowedHosts: config.updateAllowedHosts, baseUrl: metadataUrl, maxBytes: config.updateMaxBytes, + timeoutMs: config.updateTimeoutMs, publicKey: config.updatePublicKey, requireSignature: config.updateRequireSignature, }); @@ -426,6 +428,17 @@ function requestJobId(filePath: string): string | 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; + } +} + function currentReleaseVersion(config: AppConfig): string | null { try { const target = realpathSync(config.currentLink); @@ -460,6 +473,12 @@ export function reconcileOrphanedUpdateJobs(database: Database.Database, config: 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 @@ -472,7 +491,10 @@ export function reconcileOrphanedUpdateJobs(database: Database.Database, config: // 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" && !stateFresh && now - row.updatedAt >= QUEUED_UPDATE_TIMEOUT_MS) { + const matchingFreshRequest = requestMarkerJobId === row.id && requestFresh; + const matchingFreshState = stateMarkerJobId === row.id && stateFresh; + 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 @@ -503,7 +525,7 @@ 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") { - if (row.operation !== "apply" || requestFresh || stateFresh) continue; + if (row.operation !== "apply" || matchingFreshRequest || matchingFreshState) continue; const changed = database.transaction(() => { const result = database.prepare(` UPDATE update_jobs @@ -529,7 +551,7 @@ export function reconcileOrphanedUpdateJobs(database: Database.Database, config: } continue; } - if (requestFresh || stateFresh) 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(() => { @@ -586,15 +608,15 @@ export function cancelUpdateJob( 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; + ? 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 status IN ('queued', 'downloading')").run(now, now, job.id); + 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, @@ -610,7 +632,9 @@ export function cancelUpdateJob( })(); if (changed) { - forceRemoveRequest(config.updateRequestPath); + // 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(() => {}); diff --git a/server/update.ts b/server/update.ts index 16b7d70..a6e07d6 100644 --- a/server/update.ts +++ b/server/update.ts @@ -59,8 +59,22 @@ export type UrlPolicy = { /** When allowedHosts is omitted, requests are constrained to this URL's host. */ baseUrl?: string | URL | undefined; maxRedirects?: number | undefined; + /** Maximum time allowed for one metadata/sidecar/archive request. */ + timeoutMs?: number | undefined; }; +/** Release an unread response body before following a redirect or returning + * an error. Undici keeps the underlying connection associated with a body + * until it is consumed or cancelled; leaving it open can exhaust sockets when + * an update feed repeatedly returns errors or oversized responses. */ +async function cancelResponseBody(response: Response): Promise { + try { + await response.body?.cancel(); + } catch { + // The body may already be consumed/closed. Cancellation is best effort. + } +} + function invalidVersion(): never { throw new Error("更新版本号无效"); } @@ -157,8 +171,37 @@ function metadataError(): Error { } const DEFAULT_METADATA_MAX_BYTES = 2 * 1024 * 1024; +/** Maximum time allowed for one update HTTP request, including its body. */ +export const DEFAULT_UPDATE_TIMEOUT_MS = 30_000; export const RELEASE_NOTES_MAX_BYTES = 64 * 1024; +type UpdateFetchOptions = { + fetchImpl?: typeof fetch | undefined; + maxBytes?: number | undefined; + timeoutMs?: number | undefined; +}; + +function updateTimeoutMs(options: UpdateFetchOptions): number { + if (options.timeoutMs !== undefined) { + if (!Number.isSafeInteger(options.timeoutMs) || options.timeoutMs <= 0) throw new Error("更新请求超时配置无效"); + return options.timeoutMs; + } + const configuredSeconds = process.env.TALLYNOTE_UPDATE_TIMEOUT_SECONDS; + if (configuredSeconds !== undefined && configuredSeconds.trim() !== "") { + const seconds = Number(configuredSeconds); + if (!Number.isSafeInteger(seconds) || seconds <= 0) throw new Error("TALLYNOTE_UPDATE_TIMEOUT_SECONDS 必须是大于 0 的整数"); + return seconds * 1000; + } + return DEFAULT_UPDATE_TIMEOUT_MS; +} + +function beginUpdateRequest(options: UpdateFetchOptions): { signal: AbortSignal; clear: () => void } { + const controller = new AbortController(); + const timer = setTimeout(() => controller.abort(), updateTimeoutMs(options)); + timer.unref?.(); + return { signal: controller.signal, clear: () => clearTimeout(timer) }; +} + function releaseNotesText(value: unknown): string | undefined { if (typeof value !== "string" || value.length === 0) return undefined; // Gitea exposes both Markdown (body/body_html) and releaseNotes depending on @@ -210,20 +253,29 @@ function releaseResourceUrl(value: string, current: URL, options: UrlPolicy): st } /** Read a fetch body without ever buffering more than the caller's bound. */ -async function readBoundedResponse(response: Response, maxBytes: number, tooLargeMessage: string): Promise { +async function readBoundedResponse(response: Response, maxBytes: number, tooLargeMessage: string, signal?: AbortSignal): Promise { if (!Number.isSafeInteger(maxBytes) || maxBytes <= 0) throw new Error("响应大小限制无效"); const contentLength = response.headers.get("content-length"); if (contentLength !== null) { const declared = Number(contentLength); - if (Number.isFinite(declared) && declared > maxBytes) throw new Error(tooLargeMessage); + if (Number.isFinite(declared) && declared > maxBytes) { + await cancelResponseBody(response); + throw new Error(tooLargeMessage); + } } if (!response.body) return Buffer.alloc(0); const reader = response.body.getReader(); const chunks: Buffer[] = []; let total = 0; + let onAbort: (() => void) | undefined; + const abort = signal ? new Promise((_, reject) => { + onAbort = () => reject(new Error("更新请求超时")); + if (signal.aborted) onAbort(); + else signal.addEventListener("abort", onAbort, { once: true }); + }) : undefined; try { for (;;) { - const result = await reader.read(); + const result = await (abort ? Promise.race([reader.read(), abort]) : reader.read()); if (result.done) break; const chunk = Buffer.from(result.value); if (chunk.length > maxBytes - total) { @@ -233,7 +285,11 @@ async function readBoundedResponse(response: Response, maxBytes: number, tooLarg total += chunk.length; chunks.push(chunk); } + } catch (error) { + await reader.cancel().catch(() => undefined); + throw error; } finally { + if (signal && onAbort) signal.removeEventListener("abort", onAbort); reader.releaseLock(); } return Buffer.concat(chunks, total); @@ -241,84 +297,78 @@ async function readBoundedResponse(response: Response, maxBytes: number, tooLarg export async function fetchReleaseMetadata( metadataUrl: string | URL, - options: UrlPolicy & { fetchImpl?: typeof fetch | undefined; maxBytes?: number | undefined } = {}, + options: UrlPolicy & UpdateFetchOptions = {}, ): Promise { const fetchImpl = options.fetchImpl ?? fetch; let current = validateHttpsUrl(metadataUrl, options); const maxRedirects = options.maxRedirects ?? 3; let response: Response; for (let redirects = 0; ; redirects += 1) { + const request = beginUpdateRequest(options); try { - response = await fetchImpl(current, { method: "GET", redirect: "manual", headers: { accept: "application/json" } }); + response = await fetchImpl(current, { method: "GET", redirect: "manual", headers: { accept: "application/json" }, signal: request.signal }); } catch { + request.clear(); throw metadataError(); } - if (response.status < 300 || response.status >= 400) break; + if (response.status < 300 || response.status >= 400) { + try { + if (response.status < 200 || response.status >= 300) { + await cancelResponseBody(response); + throw metadataError(); + } + const maxBytes = Math.min(options.maxBytes ?? DEFAULT_METADATA_MAX_BYTES, DEFAULT_METADATA_MAX_BYTES); + const body = await readBoundedResponse(response, maxBytes, "更新发布信息过大", request.signal); + const payload: unknown = JSON.parse(body.toString("utf8")); + if (!payload || typeof payload !== "object") throw metadataError(); + const item = payload as Record; + const rawVersion = typeof item.version === "string" ? item.version : typeof item.tag_name === "string" ? item.tag_name : typeof item.tagName === "string" ? item.tagName : undefined; + if (!rawVersion) throw metadataError(); + const version = parseSemver(rawVersion); + if (typeof item.tag_name === "string" && compareSemver(version, item.tag_name) !== 0) throw metadataError(); + const assetsRaw = Array.isArray(item.assets) ? item.assets : []; + const assets: ReleaseAsset[] = []; + for (const raw of assetsRaw) { + if (!raw || typeof raw !== "object") continue; + const asset = raw as Record; + const name = typeof asset.name === "string" ? asset.name : undefined; + const url = typeof asset.url === "string" ? asset.url : typeof asset.browser_download_url === "string" ? asset.browser_download_url : undefined; + if (!name || !url) continue; + let sha256: string | undefined; + const digest = typeof asset.sha256 === "string" ? asset.sha256 : typeof asset.digest === "string" ? asset.digest : undefined; + if (digest) { + const candidate = digest.replace(/^sha256:/i, "").toLowerCase(); + if (/^[a-f0-9]{64}$/.test(candidate)) sha256 = candidate; + } + assets.push({ name, url: releaseResourceUrl(url, current, options), ...(sha256 ? { sha256 } : {}), ...(typeof asset.size === "number" && Number.isSafeInteger(asset.size) && asset.size >= 0 ? { size: asset.size } : {}) }); + } + const notes = releaseNotesText(item.body ?? item.releaseNotes ?? item.release_notes ?? item.body_html); + const releaseName = releaseNameText(item.name ?? item.releaseName); + let releaseUrl: string | undefined; + if (typeof item.html_url === "string" || typeof item.url === "string") { + try { releaseUrl = releaseResourceUrl(typeof item.html_url === "string" ? item.html_url : item.url as string, current, options); } catch { /* optional */ } + } + return { + version: `${version.major}.${version.minor}.${version.patch}${version.prerelease.length ? `-${version.prerelease.join(".")}` : ""}${version.build.length ? `+${version.build.join(".")}` : ""}`, + ...(typeof item.tag_name === "string" ? { tagName: item.tag_name } : {}), ...(releaseName ? { releaseName } : {}), ...(typeof item.published_at === "string" ? { publishedAt: item.published_at } : {}), ...(notes ? { notes } : {}), ...(releaseUrl ? { releaseUrl } : {}), assets, + }; + } catch { throw metadataError(); } + finally { request.clear(); } + } + await cancelResponseBody(response); + request.clear(); if (redirects >= maxRedirects) throw metadataError(); const location = response.headers.get("location"); if (!location) throw metadataError(); current = validateHttpsUrl(new URL(location, current), options.baseUrl ? options : { ...options, baseUrl: current }); } - if (response.status < 200 || response.status >= 300) throw metadataError(); - let payload: unknown; - try { - const maxBytes = Math.min(options.maxBytes ?? DEFAULT_METADATA_MAX_BYTES, DEFAULT_METADATA_MAX_BYTES); - const body = await readBoundedResponse(response, maxBytes, "更新发布信息过大"); - payload = JSON.parse(body.toString("utf8")); - } catch { throw metadataError(); } - if (!payload || typeof payload !== "object") throw metadataError(); - const item = payload as Record; - const rawVersion = typeof item.version === "string" ? item.version : typeof item.tag_name === "string" ? item.tag_name : typeof item.tagName === "string" ? item.tagName : undefined; - if (!rawVersion) throw metadataError(); - const version = parseSemver(rawVersion); - if (typeof item.tag_name === "string") { - try { - if (compareSemver(version, item.tag_name) !== 0) throw metadataError(); - } catch { - throw metadataError(); - } - } - const assetsRaw = Array.isArray(item.assets) ? item.assets : []; - const assets: ReleaseAsset[] = []; - for (const raw of assetsRaw) { - if (!raw || typeof raw !== "object") continue; - const asset = raw as Record; - const name = typeof asset.name === "string" ? asset.name : undefined; - const url = typeof asset.url === "string" ? asset.url : typeof asset.browser_download_url === "string" ? asset.browser_download_url : undefined; - if (!name || !url) continue; - let sha256: string | undefined; - const digest = typeof asset.sha256 === "string" ? asset.sha256 : typeof asset.digest === "string" ? asset.digest : undefined; - if (digest) { - const candidate = digest.replace(/^sha256:/i, "").toLowerCase(); - if (/^[a-f0-9]{64}$/.test(candidate)) sha256 = candidate; - } - assets.push({ name, url: releaseResourceUrl(url, current, options), ...(sha256 ? { sha256 } : {}), ...(typeof asset.size === "number" && Number.isSafeInteger(asset.size) && asset.size >= 0 ? { size: asset.size } : {}) }); - } - const notes = releaseNotesText(item.body ?? item.releaseNotes ?? item.release_notes ?? item.body_html); - const releaseName = releaseNameText(item.name ?? item.releaseName); - let releaseUrl: string | undefined; - if (typeof item.html_url === "string" || typeof item.url === "string") { - try { - const candidate = typeof item.html_url === "string" ? item.html_url : item.url as string; - releaseUrl = releaseResourceUrl(candidate, current, options); - } catch { /* omit invalid optional release page URL */ } - } - return { - version: `${version.major}.${version.minor}.${version.patch}${version.prerelease.length ? `-${version.prerelease.join(".")}` : ""}${version.build.length ? `+${version.build.join(".")}` : ""}`, - ...(typeof item.tag_name === "string" ? { tagName: item.tag_name } : {}), - ...(releaseName ? { releaseName } : {}), - ...(typeof item.published_at === "string" ? { publishedAt: item.published_at } : {}), - ...(notes ? { notes } : {}), - ...(releaseUrl ? { releaseUrl } : {}), - assets, - }; } /** Fetch a small text sidecar (for example SHA256SUMS) with the same * redirect, HTTPS and host policy used for release metadata. */ export async function fetchReleaseText( textUrl: string | URL, - options: UrlPolicy & { fetchImpl?: typeof fetch | undefined; maxBytes?: number | undefined } = {}, + options: UrlPolicy & UpdateFetchOptions = {}, ): Promise { const fetchImpl = options.fetchImpl ?? fetch; let current = validateHttpsUrl(textUrl, options); @@ -328,27 +378,32 @@ export async function fetchReleaseText( const maxRedirects = options.maxRedirects ?? 3; let response: Response; for (let redirects = 0; ; redirects += 1) { + const request = beginUpdateRequest(options); try { - response = await fetchImpl(current, { method: "GET", redirect: "manual" }); + response = await fetchImpl(current, { method: "GET", redirect: "manual", signal: request.signal }); } catch { + request.clear(); throw new Error("更新校验文件下载失败"); } - if (response.status < 300 || response.status >= 400) break; + if (response.status < 300 || response.status >= 400) { + if (response.status < 200 || response.status >= 300) { await cancelResponseBody(response); request.clear(); throw new Error("更新校验文件下载失败"); } + const declared = Number(response.headers.get("content-length") ?? 0); + const maxBytes = options.maxBytes ?? 1024 * 1024; + if (declared > maxBytes) { await cancelResponseBody(response); request.clear(); throw new Error("更新校验文件过大"); } + try { + return (await readBoundedResponse(response, maxBytes, "更新校验文件过大", request.signal)).toString("utf8"); + } catch (error) { + if (error instanceof Error && error.message === "更新校验文件过大") throw error; + throw new Error("更新校验文件下载失败"); + } finally { request.clear(); } + } + await cancelResponseBody(response); + request.clear(); if (redirects >= maxRedirects) throw new Error("更新校验文件下载失败"); const location = response.headers.get("location"); if (!location) throw new Error("更新校验文件下载失败"); current = validateHttpsUrl(new URL(location, current), redirectPolicy); } - if (response.status < 200 || response.status >= 300) throw new Error("更新校验文件下载失败"); - const declared = Number(response.headers.get("content-length") ?? 0); - const maxBytes = options.maxBytes ?? 1024 * 1024; - if (declared > maxBytes) throw new Error("更新校验文件过大"); - try { - return (await readBoundedResponse(response, maxBytes, "更新校验文件过大")).toString("utf8"); - } catch (error) { - if (error instanceof Error && error.message === "更新校验文件过大") throw error; - throw new Error("更新校验文件下载失败"); - } } /** Fetch a bounded binary sidecar (for example an Ed25519 detached @@ -356,7 +411,7 @@ export async function fetchReleaseText( * this separate from fetchReleaseText. */ export async function fetchReleaseBytes( bytesUrl: string | URL, - options: UrlPolicy & { fetchImpl?: typeof fetch | undefined; maxBytes?: number | undefined } = {}, + options: UrlPolicy & UpdateFetchOptions = {}, ): Promise { const fetchImpl = options.fetchImpl ?? fetch; let current = validateHttpsUrl(bytesUrl, options); @@ -366,27 +421,32 @@ export async function fetchReleaseBytes( const maxRedirects = options.maxRedirects ?? 3; let response: Response; for (let redirects = 0; ; redirects += 1) { + const request = beginUpdateRequest(options); try { - response = await fetchImpl(current, { method: "GET", redirect: "manual" }); + response = await fetchImpl(current, { method: "GET", redirect: "manual", signal: request.signal }); } catch { + request.clear(); throw new Error("更新签名下载失败"); } - if (response.status < 300 || response.status >= 400) break; + if (response.status < 300 || response.status >= 400) { + if (response.status < 200 || response.status >= 300) { await cancelResponseBody(response); request.clear(); throw new Error("更新签名下载失败"); } + const declared = Number(response.headers.get("content-length") ?? 0); + const maxBytes = options.maxBytes ?? 64 * 1024; + if (declared > maxBytes) { await cancelResponseBody(response); request.clear(); throw new Error("更新签名文件过大"); } + try { + return await readBoundedResponse(response, maxBytes, "更新签名文件过大", request.signal); + } catch (error) { + if (error instanceof Error && error.message === "更新签名文件过大") throw error; + throw new Error("更新签名下载失败"); + } finally { request.clear(); } + } + await cancelResponseBody(response); + request.clear(); if (redirects >= maxRedirects) throw new Error("更新签名下载失败"); const location = response.headers.get("location"); if (!location) throw new Error("更新签名下载失败"); current = validateHttpsUrl(new URL(location, current), redirectPolicy); } - if (response.status < 200 || response.status >= 300) throw new Error("更新签名下载失败"); - const declared = Number(response.headers.get("content-length") ?? 0); - const maxBytes = options.maxBytes ?? 64 * 1024; - if (declared > maxBytes) throw new Error("更新签名文件过大"); - try { - return await readBoundedResponse(response, maxBytes, "更新签名文件过大"); - } catch (error) { - if (error instanceof Error && error.message === "更新签名文件过大") throw error; - throw new Error("更新签名下载失败"); - } } export function selectReleaseAsset(release: ReleaseMetadata, platform = detectPlatform(), runtimeHash?: string): ReleaseAsset | undefined { @@ -438,7 +498,7 @@ export async function verifySha256(filePath: string, expected: string): Promise< export async function downloadReleaseAsset( url: string | URL, destination: string, - options: UrlPolicy & { fetchImpl?: typeof fetch | undefined; maxBytes?: number | undefined; onProgress?: ((downloadedBytes: number, totalBytes: number | null) => void) | undefined } = {}, + options: UrlPolicy & UpdateFetchOptions & { onProgress?: ((downloadedBytes: number, totalBytes: number | null) => void) | undefined } = {}, ): Promise<{ size: number; sha256: string }> { const fetchImpl = options.fetchImpl ?? fetch; let current = validateHttpsUrl(url, options); @@ -448,22 +508,32 @@ export async function downloadReleaseAsset( const maxRedirects = options.maxRedirects ?? 3; let response: Response; for (let redirects = 0; ; redirects += 1) { + const request = beginUpdateRequest(options); try { - response = await fetchImpl(current, { method: "GET", redirect: "manual" }); + response = await fetchImpl(current, { method: "GET", redirect: "manual", signal: request.signal }); } catch { + request.clear(); throw new Error("更新文件下载失败"); } - if (response.status < 300 || response.status >= 400) break; + if (response.status < 300 || response.status >= 400) { request.clear(); break; } + await cancelResponseBody(response); + request.clear(); if (redirects >= maxRedirects) throw new Error("更新文件下载失败"); const location = response.headers.get("location"); if (!location) throw new Error("更新文件下载失败"); current = validateHttpsUrl(new URL(location, current), redirectPolicy); } - if (response.status < 200 || response.status >= 300 || !response.body) throw new Error("更新文件下载失败"); + if (response.status < 200 || response.status >= 300 || !response.body) { + await cancelResponseBody(response); + throw new Error("更新文件下载失败"); + } const declared = Number(response.headers.get("content-length") ?? 0); const totalBytes = Number.isSafeInteger(declared) && declared > 0 ? declared : null; const maxBytes = options.maxBytes ?? 512 * 1024 * 1024; - if (declared > maxBytes) throw new Error("更新文件超过大小限制"); + if (declared > maxBytes) { + await cancelResponseBody(response); + throw new Error("更新文件超过大小限制"); + } await mkdir(path.dirname(destination), { recursive: true, mode: 0o700 }); const temporary = `${destination}.part-${randomUUID()}`; let size = 0; @@ -475,8 +545,10 @@ export async function downloadReleaseAsset( hash.update(chunk); callback(null, chunk); } }); + const request = beginUpdateRequest(options); try { - await pipeline(Readable.fromWeb(response.body as import("node:stream/web").ReadableStream), meter, createWriteStream(temporary, { flags: "wx", mode: 0o600 })); + const source = Readable.fromWeb(response.body as import("node:stream/web").ReadableStream, { signal: request.signal }); + await pipeline(source, meter, createWriteStream(temporary, { flags: "wx", mode: 0o600 })); const fd = await open(temporary, "r"); await fd.sync(); await fd.close(); @@ -484,6 +556,8 @@ export async function downloadReleaseAsset( } catch (error) { await import("node:fs/promises").then(({ rm }) => rm(temporary, { force: true })).catch(() => undefined); throw error instanceof Error && error.message.startsWith("更新文件") ? error : new Error("更新文件下载失败"); + } finally { + request.clear(); } return { size, sha256: hash.digest("hex") }; } diff --git a/systemd/tallynote-update.service b/systemd/tallynote-update.service index aa8ee4e..8333261 100644 --- a/systemd/tallynote-update.service +++ b/systemd/tallynote-update.service @@ -1,8 +1,5 @@ [Unit] Description=TallyNote privileged release updater -After=network-online.target -Wants=network-online.target - [Service] Type=oneshot User=root @@ -11,10 +8,12 @@ WorkingDirectory=/opt/tallynote/current EnvironmentFile=-/etc/tallynote/tallynote.env ExecStart=/usr/local/libexec/tallynote-update-runner Environment=PATH=/usr/sbin:/usr/bin:/sbin:/bin +# The runner consumes queued requests immediately and applies its own bounded +# phase timeouts while keeping full CLI diagnostics in the runner log. # Downloads, archive validation and data backups can exceed systemd's 90s # default start timeout on a slower server. Keep one update job alive long # enough to finish or reach its own health-check/recovery path. -TimeoutStartSec=30min +TimeoutStartSec=32min NoNewPrivileges=true # Keep the updater compatible with the same Node/libuv interface discovery # path while retaining an explicit socket-family allowlist. diff --git a/systemd/tallynote.env.example b/systemd/tallynote.env.example index 1cc6d8c..20ca1d0 100644 --- a/systemd/tallynote.env.example +++ b/systemd/tallynote.env.example @@ -9,6 +9,7 @@ TALLYNOTE_TIMEZONE=Asia/Shanghai TALLYNOTE_UPDATE_STRATEGY=systemd TALLYNOTE_UPDATE_METADATA_URL=https://git.awaioi.com/api/v1/repos/awaioi/TallyNote/releases/latest TALLYNOTE_UPDATE_ALLOWED_HOSTS=git.awaioi.com +TALLYNOTE_UPDATE_TIMEOUT_SECONDS=30 TALLYNOTE_UPDATE_REQUIRE_SIGNATURE=false TALLYNOTE_UPDATE_CHECK_COOLDOWN_SECONDS=60 TALLYNOTE_UPDATE_DOWNLOAD_COOLDOWN_SECONDS=15 diff --git a/tests/update-api.test.ts b/tests/update-api.test.ts index 384f05a..8e8d6be 100644 --- a/tests/update-api.test.ts +++ b/tests/update-api.test.ts @@ -176,6 +176,42 @@ describe("更新 API", () => { expect(ownDetail.json().job.errorMessage).toBe("更新失败,请查看服务器日志或重试"); }); + it("取消任务按管理员隔离,并只删除匹配任务的请求文件", async () => { + const owner = await login("cancel-owner"); + const other = await login("cancel-other"); + const ownerId = (database.sqlite.prepare("SELECT id FROM admins WHERE username=?").get("cancel-owner") as { id: string }).id; + const otherId = (database.sqlite.prepare("SELECT id FROM admins WHERE username=?").get("cancel-other") as { id: string }).id; + const now = Date.now(); + const ownerJobId = randomUUID(); + const otherJobId = randomUUID(); + const insert = database.sqlite.prepare(` + INSERT INTO update_jobs(id, admin_id, operation, status, version, platform, asset_url, created_at, updated_at) + VALUES (?, ?, 'download', 'queued', '1.3.0', ?, 'https://updates.example/update.tar.gz', ?, ?) + `); + insert.run(ownerJobId, ownerId, detectPlatform().target, now, now); + insert.run(otherJobId, otherId, detectPlatform().target, now + 1, now + 1); + await import("node:fs/promises").then(({ writeFile }) => writeFile(config.updateRequestPath, JSON.stringify({ jobId: otherJobId }), { encoding: "utf8", mode: 0o600 })); + + const ownerCancel = await app.inject({ + method: "POST", url: "/api/update/cancel", + headers: { origin: config.publicOrigin, cookie: owner.cookies, "x-csrf-token": owner.csrf }, + payload: { jobId: ownerJobId }, + }); + expect(ownerCancel.statusCode).toBe(200); + expect((database.sqlite.prepare("SELECT status FROM update_jobs WHERE id=?").get(ownerJobId) as { status: string }).status).toBe("cancelled"); + expect((database.sqlite.prepare("SELECT status FROM update_jobs WHERE id=?").get(otherJobId) as { status: string }).status).toBe("queued"); + expect(existsSync(config.updateRequestPath)).toBe(true); + + const otherCancel = await app.inject({ + method: "POST", url: "/api/update/cancel", + headers: { origin: config.publicOrigin, cookie: other.cookies, "x-csrf-token": other.csrf }, + payload: {}, + }); + expect(otherCancel.statusCode).toBe(200); + expect((database.sqlite.prepare("SELECT status FROM update_jobs WHERE id=?").get(otherJobId) as { status: string }).status).toBe("cancelled"); + expect(existsSync(config.updateRequestPath)).toBe(false); + }); + it("应用前重新校验失败时写入失败审计", async () => { const session = await login("update-audit"); globalThis.fetch = (async () => new Response("upstream unavailable", { status: 503 })) as typeof fetch; diff --git a/tests/update.test.ts b/tests/update.test.ts index 8023991..3a4cbcb 100644 --- a/tests/update.test.ts +++ b/tests/update.test.ts @@ -286,8 +286,27 @@ describe("更新安全工具", () => { expect(await readFile(path.join(config.currentLink, "dist", "marker"), "utf8")).toBe("new"); const row = database.sqlite.prepare("SELECT operation, status FROM update_jobs WHERE id=?").get(jobId); expect(row).toEqual({ operation: "apply", status: "applying" }); - finalizeUpdateJob(database.sqlite, jobId, "failed"); - expect(database.sqlite.prepare("SELECT status FROM update_jobs WHERE id=?").get(jobId)).toEqual({ status: "failed" }); + finalizeUpdateJob(database.sqlite, jobId, "failed", "健康检查失败(自定义)"); + expect(database.sqlite.prepare("SELECT status, error_message AS errorMessage FROM update_jobs WHERE id=?").get(jobId)).toEqual({ status: "failed", errorMessage: "健康检查失败(自定义)" }); + finalizeUpdateJob(database.sqlite, jobId, "failed", "第二次 finalize 不应覆盖原消息"); + expect(database.sqlite.prepare("SELECT status, error_message AS errorMessage FROM update_jobs WHERE id=?").get(jobId)).toEqual({ status: "failed", errorMessage: "健康检查失败(自定义)" }); + expect(database.sqlite.prepare("SELECT COUNT(*) AS count FROM audit_events WHERE action='update.failed' AND target_id=?").get(jobId)).toEqual({ count: 1 }); + const defaultJobId = randomUUID(); + database.sqlite.prepare(` + INSERT INTO update_jobs(id, operation, status, version, platform, asset_url, created_at, updated_at) + VALUES (?, 'apply', 'applying', '1.3.0', ?, ?, ?, ?) + `).run(defaultJobId, detectPlatform().target, "https://updates.example/" + assetName, now, now); + finalizeUpdateJob(database.sqlite, defaultJobId, "failed", ""); + expect(database.sqlite.prepare("SELECT error_message AS errorMessage FROM update_jobs WHERE id=?").get(defaultJobId)).toEqual({ errorMessage: "新版本健康检查失败,已恢复上一版本" }); + const completedJobId = randomUUID(); + database.sqlite.prepare(` + INSERT INTO update_jobs(id, operation, status, version, platform, asset_url, created_at, updated_at) + VALUES (?, 'apply', 'completed', '1.3.0', ?, ?, ?, ?) + `).run(completedJobId, detectPlatform().target, "https://updates.example/" + assetName, now, now); + finalizeUpdateJob(database.sqlite, completedJobId, "completed"); + expect(database.sqlite.prepare("SELECT COUNT(*) AS count FROM audit_events WHERE action='update.completed' AND target_id=?").get(completedJobId)).toEqual({ count: 0 }); + expect(() => finalizeUpdateJob(database.sqlite, completedJobId, "failed", "不能降级已完成任务")).toThrow("更新任务状态不允许完成"); + expect(database.sqlite.prepare("SELECT status FROM update_jobs WHERE id=?").get(completedJobId)).toEqual({ status: "completed" }); } finally { globalThis.fetch = previousFetch; if (database) database.sqlite.close(); @@ -385,7 +404,7 @@ describe("更新安全工具", () => { } }); - it("队列任务有新请求标记时可被重新检查,标记过期后才回收", async () => { + it("队列任务有匹配请求标记时保留到租约过期,过期后才回收", async () => { const root = await mkdtemp(path.join(tmpdir(), "tallynote-update-queued-marker-")); let database: ReturnType | undefined; try { @@ -412,12 +431,14 @@ describe("更新安全工具", () => { const now = Date.now(); await utimes(config.updateRequestPath, new Date(now), new Date(now)); - expect(reconcileOrphanedUpdateJobs(database.sqlite, config, now)).toBe(1); - expect(database.sqlite.prepare("SELECT status FROM update_jobs WHERE id=?").get(jobId)).toEqual({ status: "failed" }); - await expect(stat(config.updateRequestPath)).rejects.toThrow(); + // The DB row is old, but the request marker is fresh and names this + // exact job. Keep it queued while systemd has a chance to consume it. + expect(reconcileOrphanedUpdateJobs(database.sqlite, config, now)).toBe(0); + expect(database.sqlite.prepare("SELECT status FROM update_jobs WHERE id=?").get(jobId)).toEqual({ status: "queued" }); + await expect(stat(config.updateRequestPath)).resolves.toBeTruthy(); const expiredNow = now + ORPHANED_UPDATE_TIMEOUT_MS + 1; - expect(reconcileOrphanedUpdateJobs(database.sqlite, config, expiredNow)).toBe(0); + expect(reconcileOrphanedUpdateJobs(database.sqlite, config, expiredNow)).toBe(1); expect(database.sqlite.prepare("SELECT status FROM update_jobs WHERE id=?").get(jobId)).toEqual({ status: "failed" }); await expect(stat(config.updateRequestPath)).rejects.toThrow(); } finally { diff --git a/web-next/src/pages/update/UpdatePage.tsx b/web-next/src/pages/update/UpdatePage.tsx index f16f8b1..f7de0e4 100644 --- a/web-next/src/pages/update/UpdatePage.tsx +++ b/web-next/src/pages/update/UpdatePage.tsx @@ -248,22 +248,48 @@ export default function UpdatePage({ const announced = useRef(null); const checkInFlight = useRef(false); const actionInFlight = useRef(false); + const cancelInFlight = useRef(false); + const loadInFlight = useRef(false); const disconnected = useRef(false); const recoveredNotice = useRef(false); const info = liveInfo; - const load = async () => { - setLoading(true); + const mergeInfo = (incoming: UpdateInfo, preserveJob = false) => { + setLiveInfo((current) => { + if (!current) return incoming; + const next = { + ...current, + ...incoming, + // A check response describes the release, but must not erase the job + // that is already being displayed while that check is in flight. + job: preserveJob + ? (current.job ?? incoming.job) + : (!Object.prototype.hasOwnProperty.call(incoming, "job") ? current.job : incoming.job), + }; + // Some status/check responses are intentionally partial. Keep fields + // from the last successful response when a field is omitted. + if (!Object.prototype.hasOwnProperty.call(incoming, "latest")) next.latest = current.latest; + if (!Object.prototype.hasOwnProperty.call(incoming, "checkedAt")) next.checkedAt = current.checkedAt; + if (!Object.prototype.hasOwnProperty.call(incoming, "strategy")) next.strategy = current.strategy; + return next; + }); + }; + + const load = async (showLoading = true) => { + if (loadInFlight.current) return; + loadInFlight.current = true; + if (showLoading) setLoading(true); setError(""); try { const res = await api("/api/update/status"); - setLiveInfo(res); + mergeInfo(res); updateInfoCache = res; } catch (e) { setError((e as Error).message); } finally { - setLoading(false); + if (showLoading) setLoading(false); + loadInFlight.current = false; } }; @@ -293,7 +319,8 @@ export default function UpdatePage({ return () => window.clearInterval(timer); }, [restartAt]); - // Live polling effect: only poll when active job exists + // Live polling effect for an active job. Completed responses are refreshed + // once so latest/checkedAt/strategy stay current; other terminal states stay visible. useEffect(() => { const shouldPoll = Boolean( job && (pollableStatuses.has(job.status) || (job.status === "staged" && job.operation === "apply")) @@ -326,23 +353,53 @@ export default function UpdatePage({ setReloadReady(true); notify?.("系统升级完成,请刷新页面", "success"); } - if ( + if (result.job.status === "completed") { + void load(false); + } else if ( pollableStatuses.has(result.job.status) || (result.job.status === "staged" && result.job.operation === "apply") ) { schedule(1500); } - } catch { + } catch (caught) { if (disposed) return; + const status = caught instanceof ApiError ? caught.status : 0; + const message = caught instanceof Error ? caught.message : "读取更新任务状态失败"; + if (status === 404) { + setPollError("更新任务已结束,正在刷新状态。"); + void load(false); + return; + } + if (status === 401) { + setPollError("登录已失效,请重新登录。"); + return; + } + if (status === 403) { + setPollError("没有权限读取更新任务状态。"); + return; + } + failures += 1; - disconnected.current = true; + disconnected.current = status === 0 || status === 408; recoveredNotice.current = false; - setPollError( - `服务正在平滑重启${ - restartSeconds !== null ? `,预计 ${restartSeconds} 秒后恢复` : ",页面正在自动探活重试" - }。升级任务仍在后台安全执行。` - ); - schedule(Math.min(1500 * 2 ** Math.min(failures, 3), 10_000)); + if (status === 409) { + setPollError(`${message},正在刷新状态。`); + void load(false); + return; + } + if (status === 429) { + setPollError(`${message},稍后继续同步。`); + } else if (status === 408) { + setPollError("读取更新任务状态超时,正在重试。"); + } else if (status === 0) { + setPollError("更新服务连接异常,正在重试。"); + } else { + setPollError(`${message},正在重试。`); + } + const retryAfter = caught instanceof ApiError && caught.retryAfter + ? Math.max(1000, caught.retryAfter * 1000) + : Math.min(1500 * 2 ** Math.min(failures, 3), 10_000); + schedule(retryAfter); } }; void poll(); @@ -352,6 +409,42 @@ export default function UpdatePage({ }; }, [job?.id, job?.status, job?.operation, notify]); + // With no visible job, make a low-frequency status request so a task created + // elsewhere can still appear without creating a request storm. + useEffect(() => { + if (job) return; + let timer: number | undefined; + let disposed = false; + const discover = () => { + if (disposed || document.visibilityState !== "visible") return; + void load(false); + }; + const schedule = () => { + if (disposed) return; + if (timer !== undefined) window.clearTimeout(timer); + timer = window.setTimeout(() => { + discover(); + schedule(); + }, 5000); + }; + const onVisibilityChange = () => { + if (document.visibilityState === "visible") { + discover(); + schedule(); + } else if (timer !== undefined) { + window.clearTimeout(timer); + timer = undefined; + } + }; + if (document.visibilityState === "visible") schedule(); + document.addEventListener("visibilitychange", onVisibilityChange); + return () => { + disposed = true; + if (timer !== undefined) window.clearTimeout(timer); + document.removeEventListener("visibilitychange", onVisibilityChange); + }; + }, [job?.id, job?.status, job?.operation]); + // Check update handler const check = async () => { if (checkInFlight.current) return; @@ -359,7 +452,8 @@ export default function UpdatePage({ setChecking(true); setError(""); try { - setLiveInfo(await api("/api/update/check", { method: "POST" })); + const result = await api("/api/update/check", { method: "POST" }); + mergeInfo(result, true); notify?.("版本检查完成", "info"); } catch (e) { if (e instanceof ApiError && e.status === 429) { @@ -375,7 +469,8 @@ export default function UpdatePage({ // Cancel queued job handler const cancelJob = async () => { - if (cancelling) return; + if (cancelInFlight.current || !job?.id) return; + cancelInFlight.current = true; setCancelling(true); try { await api("/api/update/cancel", { method: "POST", body: JSON.stringify({ jobId: job?.id }) }); @@ -387,6 +482,7 @@ export default function UpdatePage({ notify?.((e as Error).message, "error"); } finally { setCancelling(false); + cancelInFlight.current = false; } }; @@ -487,7 +583,7 @@ export default function UpdatePage({ if (!job) return 0; if (job.status === "completed" || job.status === "staged") return 100; if (job.status === "downloading") { - if (!job.sizeBytes || !job.downloadedBytes) return 10; + if (!job.sizeBytes || !job.downloadedBytes) return 0; return Math.min(99, Math.max(1, Math.round((job.downloadedBytes / job.sizeBytes) * 100))); } if (job.status === "verifying") return 99;