fix: make online updates recoverable
TallyNote release / linux-x64 (push) Successful in 7m45s

This commit is contained in:
Qiufeng
2026-09-05 14:56:15 +08:00
parent a080f531cd
commit ed0b492461
12 changed files with 526 additions and 165 deletions
+2
View File
@@ -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
+1 -1
View File
@@ -1,6 +1,6 @@
{
"name": "tallynote",
"version": "1.2.3",
"version": "1.2.4",
"private": true,
"type": "module",
"packageManager": "pnpm@9.0.6",
+110 -22
View File
@@ -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
+32 -13
View File
@@ -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<ty
allowedHosts: options.allowedHosts ?? [],
baseUrl: metadataUrl.toString(),
maxBytes: options.maxBytes ?? 512 * 1024 * 1024,
timeoutMs: options.timeoutMs,
publicKey: options.publicKey,
requireSignature: options.requireSignature,
});
@@ -398,19 +407,28 @@ export function finalizeUpdateJob(
status: "completed" | "failed",
message?: string,
): void {
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("更新任务不存在");
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" ? "新版本健康检查失败,已恢复上一版本" : null;
sqlite.transaction(() => {
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<void> {
...(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,
+1
View File
@@ -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.
+34 -10
View File
@@ -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(() => {});
+165 -91
View File
@@ -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<void> {
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<Buffer> {
async function readBoundedResponse(response: Response, maxBytes: number, tooLargeMessage: string, signal?: AbortSignal): Promise<Buffer> {
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<never>((_, 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<ReleaseMetadata> {
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<string, unknown>;
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<string, unknown>;
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<string, unknown>;
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<string, unknown>;
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<string> {
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<Buffer> {
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") };
}
+3 -4
View File
@@ -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.
+1
View File
@@ -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
+36
View File
@@ -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;
+28 -7
View File
@@ -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<typeof openDatabase> | 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 {
+113 -17
View File
@@ -248,22 +248,48 @@ export default function UpdatePage({
const announced = useRef<string | null>(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<UpdateInfo>("/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<UpdateInfo>("/api/update/check", { method: "POST" }));
const result = await api<UpdateInfo>("/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;