fix(orchestrator): сериализовать операции жизненного цикла
У оркестратора не было никакой блокировки операций: ни flock, ни mutex, ни lockfile. Вся архитектура отката при этом опиралась на невысказанное допущение, что в каждый момент выполняется ровно одна операция HY2XS. install-state.json замком не является — это запись о состоянии, а не право на изменение. Два одновременных reconfigure спокойно доходили до конца каждый по-своему, и уникальные op-id не спасали: они разделяют резервные копии, но production paths общие — /etc/hysteria/config.yaml, unit-файлы, /etc/nftables.conf, install-state.json. Дальше любая из операций могла упасть и "восстановить" состояние поверх изменений другой, отчитавшись при этом полным успехом: со своим манифестом она действительно сверилась. Отдельно опасен firewall: обе операции независимо взводят транзиентные rollback-юниты, и guard одной способен снять правила другой. Введён эксклюзивный замок /run/lock/hy2xs-orchestrator.lock через атомарное создание с O_EXCL. Не flock(2): прямого биндинга в рантайме нет, а держать замок подпроцессом означало бы сторожевой процесс на каждую операцию. - install/reconfigure/repair берут замок как мутирующие; - doctor тоже: диагностика в середине транзакции описывает промежуточное состояние и выдаёт бессмысленные ошибки; - status и diagnostics collect замок НЕ берут — они нужны в том числе во время долгой операции, — но сообщают, что операция идёт; - preflight-install отказывает сразу, до exec в install.sh. Замок снимается в finally, а также на SIGINT/SIGTERM/SIGHUP и при выходе процесса: обрыв SSH не имеет права заблокировать сервер до перезагрузки. Замок мёртвого держателя переиспользуется, но только через увод файла переименованием со сверкой nonce — снимать его на месте означало бы риск снять живой. Непонятое содержимое не снимается автоматически: оно не доказывает отсутствие операции, и сомнение трактуется в пользу отказа.
This commit is contained in:
@@ -0,0 +1,468 @@
|
||||
import { existsSync, mkdirSync, readFileSync, unlinkSync } from "node:fs";
|
||||
import { open, readFile, rename, unlink } from "node:fs/promises";
|
||||
import { dirname } from "node:path";
|
||||
import { assertMutationAllowed } from "./guard";
|
||||
import { info } from "./log";
|
||||
|
||||
/**
|
||||
* Взаимное исключение операций жизненного цикла HY2XS.
|
||||
*
|
||||
* Вся архитектура отката держится на невысказанном допущении:
|
||||
*
|
||||
* в каждый момент времени на сервере выполняется ровно одна операция HY2XS.
|
||||
*
|
||||
* Код его никак не обеспечивал. `install-state.json` — не замок: это запись о
|
||||
* состоянии, а не право на изменение. Два одновременных `reconfigure` спокойно
|
||||
* доходили до конца каждый по-своему:
|
||||
*
|
||||
* A: прочитал installed=true B: прочитал installed=true
|
||||
* A: снял копию в op-A B: снял копию в op-B
|
||||
* A: записал конфиг B: записал конфиг
|
||||
* A: развернул юниты B: развернул юниты
|
||||
* A: применил firewall B: применил firewall
|
||||
*
|
||||
* Уникальные op-id разделяют РЕЗЕРВНЫЕ КОПИИ, но не production paths: и конфиг,
|
||||
* и юниты, и /etc/nftables.conf, и /var/lib/hy2xs/install-state.json — общие.
|
||||
* Дальше любая из операций могла упасть и «восстановить» состояние поверх
|
||||
* изменений другой, причём её откат отчитался бы полным успехом: с точки зрения
|
||||
* своего манифеста он действительно восстановил всё, что копировал.
|
||||
*
|
||||
* Отдельно опасен firewall: обе операции независимо взводят транзиентные
|
||||
* rollback-юниты, и guard одной операции способен снять правила другой.
|
||||
*
|
||||
* Почему не flock(2). Прямого биндинга в рантайме нет, а держать замок
|
||||
* подпроцессом `flock` означало бы завести сторожевой процесс на каждую
|
||||
* операцию. Здесь используется классический lockfile через O_EXCL: создание
|
||||
* файла с флагом «отказать, если существует» атомарно на всех файловых
|
||||
* системах, включая tmpfs, на которой живёт /run/lock.
|
||||
*
|
||||
* Асимметрия решений сознательная. Ложное «держатель жив» приводит к отказу,
|
||||
* который оператор снимает одной командой. Ложное «держатель мёртв» пускает
|
||||
* вторую операцию в те же production paths — то есть ровно то, ради чего
|
||||
* замок существует. Поэтому сомнение всегда трактуется в пользу отказа.
|
||||
*/
|
||||
|
||||
export const OPERATION_LOCK_PATH = "/run/lock/hy2xs-orchestrator.lock";
|
||||
|
||||
export type LockRecord = {
|
||||
/** PID процесса-держателя. */
|
||||
pid: number;
|
||||
/** Команда оркестратора: install | reconfigure | repair | doctor. */
|
||||
command: string;
|
||||
startedAt: string;
|
||||
/**
|
||||
* Одноразовое значение операции.
|
||||
*
|
||||
* Нужно там, где PID недостаточен: снятие замка обязано убрать ИМЕННО свой
|
||||
* файл, а не тот, который успел занять кто-то другой после переиспользования
|
||||
* освободившегося замка.
|
||||
*/
|
||||
nonce: string;
|
||||
};
|
||||
|
||||
export type LockHolder = LockRecord & { alive: boolean };
|
||||
|
||||
export class OperationInProgressError extends Error {
|
||||
readonly holder: LockHolder;
|
||||
|
||||
constructor(message: string, holder: LockHolder) {
|
||||
super(message);
|
||||
this.name = "OperationInProgressError";
|
||||
this.holder = holder;
|
||||
}
|
||||
}
|
||||
|
||||
export type OperationLock = {
|
||||
readonly command: string;
|
||||
readonly nonce: string;
|
||||
release(): Promise<void>;
|
||||
};
|
||||
|
||||
export type LockOptions = {
|
||||
/** Путь замка. Переопределяется тестами; на сервере всегда OPERATION_LOCK_PATH. */
|
||||
path?: string;
|
||||
/** Проверка живости держателя. Переопределяется тестами. */
|
||||
isProcessAlive?: (pid: number) => boolean;
|
||||
/** PID текущего процесса. Переопределяется тестами. */
|
||||
pid?: number;
|
||||
};
|
||||
|
||||
export function renderLockRecord(record: LockRecord): string {
|
||||
return `${JSON.stringify(record, null, 2)}\n`;
|
||||
}
|
||||
|
||||
/**
|
||||
* Разбор строгий: непонятый замок НЕ превращается в «замка нет».
|
||||
*
|
||||
* Возврат null означает «файл существует, но не является нашим замком», и
|
||||
* вызывающий обязан решить это отдельно, а не считать путь свободным.
|
||||
*/
|
||||
export function parseLockRecord(raw: string): LockRecord | null {
|
||||
let parsed: unknown;
|
||||
try {
|
||||
parsed = JSON.parse(raw);
|
||||
} catch {
|
||||
return null;
|
||||
}
|
||||
|
||||
if (!parsed || typeof parsed !== "object" || Array.isArray(parsed)) {
|
||||
return null;
|
||||
}
|
||||
|
||||
const record = parsed as Record<string, unknown>;
|
||||
if (typeof record.pid !== "number" || !Number.isInteger(record.pid) || record.pid <= 0) {
|
||||
return null;
|
||||
}
|
||||
if (typeof record.command !== "string" || record.command === "") {
|
||||
return null;
|
||||
}
|
||||
if (typeof record.nonce !== "string" || record.nonce === "") {
|
||||
return null;
|
||||
}
|
||||
|
||||
return {
|
||||
pid: record.pid,
|
||||
command: record.command,
|
||||
startedAt: typeof record.startedAt === "string" ? record.startedAt : "",
|
||||
nonce: record.nonce
|
||||
};
|
||||
}
|
||||
|
||||
/**
|
||||
* Жив ли держатель замка.
|
||||
*
|
||||
* На Linux решает наличие /proc/<pid>. Дополнительной сверки cmdline здесь
|
||||
* намеренно нет: она превратила бы «оркестратор запущен не тем способом»
|
||||
* (например, из исходников при отладке) в «держатель мёртв, замок можно
|
||||
* забрать», то есть в разрешение на параллельную мутацию. Переиспользование
|
||||
* PID теоретически возможно, но его цена — лишний отказ, а не потерянное
|
||||
* взаимное исключение.
|
||||
*/
|
||||
export function isProcessAliveByDefault(pid: number): boolean {
|
||||
if (!Number.isInteger(pid) || pid <= 0) {
|
||||
return false;
|
||||
}
|
||||
|
||||
if (process.platform === "linux" && existsSync("/proc")) {
|
||||
return existsSync(`/proc/${pid}`);
|
||||
}
|
||||
|
||||
try {
|
||||
process.kill(pid, 0);
|
||||
return true;
|
||||
} catch (error) {
|
||||
// EPERM означает «процесс есть, но он не наш» — это живой держатель.
|
||||
return (error as { code?: string }).code === "EPERM";
|
||||
}
|
||||
}
|
||||
|
||||
function lockPath(options: LockOptions): string {
|
||||
return options.path ?? OPERATION_LOCK_PATH;
|
||||
}
|
||||
|
||||
function livenessProbe(options: LockOptions): (pid: number) => boolean {
|
||||
return options.isProcessAlive ?? isProcessAliveByDefault;
|
||||
}
|
||||
|
||||
/**
|
||||
* Кто держит замок сейчас. Чтение, поэтому доступно и под read-only guard.
|
||||
*
|
||||
* Возвращает null, только если замка НЕТ. Нечитаемый или неразобранный файл
|
||||
* описывается как держатель с pid 0 и command "unknown": «здесь что-то лежит»
|
||||
* и «здесь ничего нет» — разные ответы.
|
||||
*/
|
||||
export async function readLockHolder(options: LockOptions = {}): Promise<LockHolder | null> {
|
||||
const path = lockPath(options);
|
||||
|
||||
let raw: string;
|
||||
try {
|
||||
raw = await readFile(path, "utf8");
|
||||
} catch (error) {
|
||||
if ((error as { code?: string }).code === "ENOENT") {
|
||||
return null;
|
||||
}
|
||||
throw error;
|
||||
}
|
||||
|
||||
const record = parseLockRecord(raw);
|
||||
if (!record) {
|
||||
return {
|
||||
pid: 0,
|
||||
command: "unknown",
|
||||
startedAt: "",
|
||||
nonce: "",
|
||||
alive: false
|
||||
};
|
||||
}
|
||||
|
||||
return { ...record, alive: livenessProbe(options)(record.pid) };
|
||||
}
|
||||
|
||||
export function describeLockHolder(holder: LockHolder): string {
|
||||
const started = holder.startedAt ? `, started at ${holder.startedAt}` : "";
|
||||
return `${holder.command} (pid ${holder.pid}${started})`;
|
||||
}
|
||||
|
||||
/**
|
||||
* Освобождает замок мёртвого держателя.
|
||||
*
|
||||
* Забирать замок «на месте» нельзя: между чтением записи и её удалением
|
||||
* держатель мог смениться, и тогда `unlink` снял бы ЖИВОЙ замок. Поэтому файл
|
||||
* сначала уводится в сторону переименованием, а затем проверяется, что уведён
|
||||
* именно тот замок, который был признан мёртвым. Если нет — он возвращается на
|
||||
* место ровно тем же содержимым, и операция отказывает.
|
||||
*/
|
||||
async function reclaimStaleLock(path: string, observed: LockRecord): Promise<boolean> {
|
||||
const stealPath = `${path}.stale-${observed.nonce}`;
|
||||
|
||||
try {
|
||||
await rename(path, stealPath);
|
||||
} catch (error) {
|
||||
if ((error as { code?: string }).code === "ENOENT") {
|
||||
// Замок уже снят кем-то другим: путь свободен, повторная попытка уместна.
|
||||
return true;
|
||||
}
|
||||
throw error;
|
||||
}
|
||||
|
||||
const stolen = parseLockRecord(await readFile(stealPath, "utf8").catch(() => ""));
|
||||
if (stolen && stolen.nonce !== observed.nonce) {
|
||||
await rename(stealPath, path);
|
||||
return false;
|
||||
}
|
||||
|
||||
await unlink(stealPath).catch(() => undefined);
|
||||
return true;
|
||||
}
|
||||
|
||||
/**
|
||||
* Замки, которые держит ЭТОТ процесс.
|
||||
*
|
||||
* Нужны обработчикам завершения: снятие замка обязано пережить и штатный
|
||||
* выход, и Ctrl+C, и SIGTERM от systemd, и обрыв SSH (SIGHUP). Без этого
|
||||
* прерванная операция оставляла бы замок до перезагрузки.
|
||||
*/
|
||||
const heldLocks = new Map<string, string>();
|
||||
let exitHandlersInstalled = false;
|
||||
|
||||
function releaseSync(path: string, nonce: string): void {
|
||||
try {
|
||||
const record = parseLockRecord(readFileSync(path, "utf8"));
|
||||
if (record && record.nonce !== nonce) {
|
||||
// Замок уже не наш: снимать его — значит открыть дорогу третьей операции.
|
||||
return;
|
||||
}
|
||||
unlinkSync(path);
|
||||
} catch {
|
||||
// Замка уже нет — снимать нечего.
|
||||
}
|
||||
heldLocks.delete(path);
|
||||
}
|
||||
|
||||
function installExitHandlers(): void {
|
||||
if (exitHandlersInstalled) {
|
||||
return;
|
||||
}
|
||||
exitHandlersInstalled = true;
|
||||
|
||||
process.on("exit", () => {
|
||||
for (const [path, nonce] of heldLocks) {
|
||||
releaseSync(path, nonce);
|
||||
}
|
||||
});
|
||||
|
||||
for (const signal of ["SIGINT", "SIGTERM", "SIGHUP"] as const) {
|
||||
process.on(signal, () => {
|
||||
for (const [path, nonce] of heldLocks) {
|
||||
releaseSync(path, nonce);
|
||||
}
|
||||
// Код возврата по соглашению: 128 + номер сигнала. Обработчик снимает
|
||||
// замок и НЕ подменяет собой штатное завершение.
|
||||
process.exit(signal === "SIGINT" ? 130 : signal === "SIGTERM" ? 143 : 129);
|
||||
});
|
||||
}
|
||||
}
|
||||
|
||||
function newNonce(): string {
|
||||
return `${Date.now().toString(36)}-${Math.random().toString(16).slice(2, 10)}`;
|
||||
}
|
||||
|
||||
async function writeLockFile(path: string, record: LockRecord): Promise<boolean> {
|
||||
let handle;
|
||||
try {
|
||||
// "wx" — атомарное создание с отказом, если файл уже существует.
|
||||
handle = await open(path, "wx", 0o644);
|
||||
} catch (error) {
|
||||
if ((error as { code?: string }).code === "EEXIST") {
|
||||
return false;
|
||||
}
|
||||
throw error;
|
||||
}
|
||||
|
||||
try {
|
||||
await handle.writeFile(renderLockRecord(record), "utf8");
|
||||
// Замок читают другие процессы, в том числе после внезапной перезагрузки
|
||||
// держателя: содержимое обязано быть на носителе, а не в page cache.
|
||||
await handle.sync();
|
||||
} finally {
|
||||
await handle.close();
|
||||
}
|
||||
|
||||
return true;
|
||||
}
|
||||
|
||||
/**
|
||||
* Берёт эксклюзивный замок операции или отказывает.
|
||||
*
|
||||
* Отказ неблокирующий и намеренно быстрый: очередь из операций жизненного
|
||||
* цикла — это не то, чего оператор ждёт от установщика. Он должен увидеть, что
|
||||
* на сервере уже что-то выполняется, и решить сам.
|
||||
*/
|
||||
export async function acquireOperationLock(
|
||||
command: string,
|
||||
options: LockOptions = {}
|
||||
): Promise<OperationLock> {
|
||||
// Замок пишется в /run — это ephemeral runtime state, а не persistent path,
|
||||
// и берётся ДО того, как команда включает read-only guard. Проверка стоит
|
||||
// здесь, чтобы перенос захвата внутрь guard'а отказал громко, а не записал
|
||||
// файл молча в фазе, которая обязана быть читающей.
|
||||
assertMutationAllowed(`acquireOperationLock(${command})`);
|
||||
|
||||
const path = lockPath(options);
|
||||
const pid = options.pid ?? process.pid;
|
||||
const record: LockRecord = {
|
||||
pid,
|
||||
command,
|
||||
startedAt: new Date().toISOString(),
|
||||
nonce: newNonce()
|
||||
};
|
||||
|
||||
mkdirSync(dirname(path), { recursive: true });
|
||||
|
||||
for (let attempt = 0; attempt < 2; attempt += 1) {
|
||||
if (await writeLockFile(path, record)) {
|
||||
installExitHandlers();
|
||||
heldLocks.set(path, record.nonce);
|
||||
info(`operation lock acquired: ${path} (${command}, pid ${pid})`);
|
||||
return {
|
||||
command,
|
||||
nonce: record.nonce,
|
||||
release: async () => {
|
||||
releaseSync(path, record.nonce);
|
||||
}
|
||||
};
|
||||
}
|
||||
|
||||
const holder = await readLockHolder(options);
|
||||
if (!holder) {
|
||||
// Замок исчез между отказом создания и чтением: пробуем ещё раз.
|
||||
continue;
|
||||
}
|
||||
|
||||
if (holder.alive) {
|
||||
throw new OperationInProgressError(
|
||||
`another HY2XS operation is already in progress: ${describeLockHolder(holder)}. ` +
|
||||
`Дождитесь её завершения; замок: ${path}`,
|
||||
holder
|
||||
);
|
||||
}
|
||||
|
||||
if (attempt > 0) {
|
||||
throw new OperationInProgressError(
|
||||
`operation lock ${path} is held by ${describeLockHolder(holder)}, which is no longer running, ` +
|
||||
"но освободить его не удалось. Проверьте состояние сервера и при необходимости удалите замок вручную.",
|
||||
holder
|
||||
);
|
||||
}
|
||||
|
||||
if (holder.pid === 0) {
|
||||
throw new OperationInProgressError(
|
||||
`operation lock ${path} exists but is not a valid HY2XS lock record. ` +
|
||||
"Файл не снимается автоматически: непонятое содержимое не является доказательством того, " +
|
||||
"что операции нет. Проверьте сервер и удалите замок вручную.",
|
||||
holder
|
||||
);
|
||||
}
|
||||
|
||||
info(
|
||||
`operation lock ${path} is held by ${describeLockHolder(holder)}, which is no longer running; reclaiming it`
|
||||
);
|
||||
if (!(await reclaimStaleLock(path, holder))) {
|
||||
throw new OperationInProgressError(
|
||||
`another HY2XS operation took the lock ${path} while a stale one was being reclaimed`,
|
||||
holder
|
||||
);
|
||||
}
|
||||
}
|
||||
|
||||
const holder = await readLockHolder(options);
|
||||
throw new OperationInProgressError(
|
||||
`failed to acquire the operation lock ${path}` +
|
||||
(holder ? `: held by ${describeLockHolder(holder)}` : ""),
|
||||
holder ?? { pid: 0, command: "unknown", startedAt: "", nonce: "", alive: false }
|
||||
);
|
||||
}
|
||||
|
||||
/** Замок снимается в finally: операция не имеет права оставить его за собой. */
|
||||
export async function withOperationLock<T>(
|
||||
command: string,
|
||||
run: () => Promise<T>,
|
||||
options: LockOptions = {}
|
||||
): Promise<T> {
|
||||
const lock = await acquireOperationLock(command, options);
|
||||
try {
|
||||
return await run();
|
||||
} finally {
|
||||
await lock.release();
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Отказ до первой мутации для команд, которые сами замок не берут.
|
||||
*
|
||||
* PHASE 0 установки обязана отказать здесь, а не после `exec` в install.sh:
|
||||
* иначе оператор получил бы отказ уже от `install`, потратив на проверки время
|
||||
* и увидев его вторым сообщением вместо первого.
|
||||
*/
|
||||
export async function assertNoOperationInProgress(
|
||||
what: string,
|
||||
options: LockOptions = {}
|
||||
): Promise<void> {
|
||||
const holder = await readLockHolder(options);
|
||||
if (!holder || !holder.alive) {
|
||||
return;
|
||||
}
|
||||
throw new OperationInProgressError(
|
||||
`${what} refused: another HY2XS operation is already in progress: ${describeLockHolder(holder)}`,
|
||||
holder
|
||||
);
|
||||
}
|
||||
|
||||
/**
|
||||
* Описание текущей операции для команд, которые ничего не меняют.
|
||||
*
|
||||
* `status` и `diagnostics` не берут замок сознательно: они существуют в том
|
||||
* числе ради того, чтобы посмотреть на сервер во время долгой операции.
|
||||
* Но показать, что операция идёт, они обязаны — иначе оператор будет разбирать
|
||||
* промежуточное состояние транзакции как окончательное.
|
||||
*/
|
||||
export async function describeOperationInProgress(
|
||||
options: LockOptions = {}
|
||||
): Promise<string | null> {
|
||||
const holder = await readLockHolder(options);
|
||||
if (!holder || !holder.alive) {
|
||||
return null;
|
||||
}
|
||||
return describeLockHolder(holder);
|
||||
}
|
||||
|
||||
/**
|
||||
* Только для тестов: сбрасывает учёт замков этого процесса.
|
||||
*
|
||||
* Учёт нужен обработчикам завершения и живёт на уровне модуля, поэтому между
|
||||
* тестами он обязан обнуляться — иначе обработчик выхода попытается снять
|
||||
* замок из временного каталога, которого уже нет.
|
||||
*/
|
||||
export function resetHeldLocksForTests(): void {
|
||||
heldLocks.clear();
|
||||
}
|
||||
Reference in New Issue
Block a user