From 1e5fe097fe002b20c79b294dd492fef4d7020575 Mon Sep 17 00:00:00 2001 From: XiaoSeS <87064762+XiaoSeS@users.noreply.github.com> Date: Wed, 9 Sep 2026 13:43:33 +0800 Subject: [PATCH] fix(cli): serialize concurrent stale lock recovery (#839) * test(cli): reproduce concurrent stale lock ownership Signed-off-by: XiaoSeS <87064762+XiaoSeS@users.noreply.github.com> * test(cli): wait for every stale lock contender Signed-off-by: XiaoSeS <87064762+XiaoSeS@users.noreply.github.com> * fix(cli): serialize stale target lock recovery Signed-off-by: XiaoSeS <87064762+XiaoSeS@users.noreply.github.com> * fix(cli): recover abandoned acquisition gates Signed-off-by: XiaoSeS <87064762+XiaoSeS@users.noreply.github.com> * test(cli): cover acquisition gate crash recovery Signed-off-by: XiaoSeS <87064762+XiaoSeS@users.noreply.github.com> * test(cli): keep live acquisition tickets Signed-off-by: XiaoSeS <87064762+XiaoSeS@users.noreply.github.com> * fix(cli): use crash-recoverable acquisition queue Signed-off-by: XiaoSeS <87064762+XiaoSeS@users.noreply.github.com> * perf(cli): minimize acquisition queue scans Signed-off-by: XiaoSeS <87064762+XiaoSeS@users.noreply.github.com> * fix(cli): stabilize acquisition queue election Signed-off-by: XiaoSeS <87064762+XiaoSeS@users.noreply.github.com> * test(cli): verify stalled contenders retain ownership Signed-off-by: XiaoSeS <87064762+XiaoSeS@users.noreply.github.com> * fix(cli): wait for acquisition queue turn Signed-off-by: XiaoSeS <87064762+XiaoSeS@users.noreply.github.com> * test(cli): recover stale target after gate owner exit Signed-off-by: XiaoSeS <87064762+XiaoSeS@users.noreply.github.com> * test(cli): assert stale target cleanup after recovery Signed-off-by: XiaoSeS <87064762+XiaoSeS@users.noreply.github.com> --------- Signed-off-by: XiaoSeS <87064762+XiaoSeS@users.noreply.github.com> --- cli/src/services/skill-target-lock.ts | 204 ++++++++++++++++-- .../unit/services/skill-target-lock.test.ts | 99 ++++++++- 2 files changed, 283 insertions(+), 20 deletions(-) diff --git a/cli/src/services/skill-target-lock.ts b/cli/src/services/skill-target-lock.ts index cd32cfd1..b2b3ad85 100644 --- a/cli/src/services/skill-target-lock.ts +++ b/cli/src/services/skill-target-lock.ts @@ -1,5 +1,5 @@ -import { createHash } from 'node:crypto' -import { chmod, lstat, mkdir } from 'node:fs/promises' +import { createHash, randomUUID } from 'node:crypto' +import { chmod, lstat, mkdir, readdir, rename, unlink, writeFile } from 'node:fs/promises' import { tmpdir } from 'node:os' import { join, resolve } from 'node:path' import { lock } from 'proper-lockfile' @@ -7,23 +7,203 @@ import { canonicalizeExistingPath } from '../platform/paths' import { CliError } from '../shared/errors' import { EXIT } from '../shared/constants' +const ACQUISITION_GATE_MAX_POLL_MS = 100 + /** Serializes every local lifecycle mutation for one Skill target directory. */ export async function acquireSkillTargetLock(rootDir: string, slug: string): Promise<() => Promise> { const lockPath = await skillTargetLockPath(rootDir, slug) + // proper-lockfile's stale deletion is not serialized. Gate only the acquisition attempt so two + // recoverers cannot remove and replace the same target lock concurrently. + const acquisitionGatePath = `${lockPath}.acquire` + let releaseAcquisitionGate: () => Promise try { - return await lock(lockPath, { - lockfilePath: lockPath, - realpath: false, - stale: 10_000, - update: 3_000, - retries: 0 - }) + releaseAcquisitionGate = await acquireAcquisitionGate(acquisitionGatePath) } catch (error) { - if (error instanceof Error && 'code' in error && error.code === 'ELOCKED') { - throw targetBusyError(rootDir, slug) - } + if (hasErrorCode(error, 'EEXIST')) throw targetBusyError(rootDir, slug) throw error } + + let releaseTarget: () => Promise + try { + releaseTarget = await acquireTargetLock(lockPath) + } catch (operationError) { + try { + await releaseAcquisitionGate() + } catch (cleanupError) { + throw new AggregateError([operationError, cleanupError], 'target lock acquisition and gate cleanup both failed') + } + if (hasErrorCode(operationError, 'ELOCKED')) throw targetBusyError(rootDir, slug) + throw operationError + } + + try { + await releaseAcquisitionGate() + } catch (gateCleanupError) { + try { + await releaseTarget() + } catch (targetCleanupError) { + throw new AggregateError([gateCleanupError, targetCleanupError], 'target and acquisition gate cleanup both failed') + } + throw gateCleanupError + } + + return releaseTarget +} + +function acquireTargetLock(lockPath: string): Promise<() => Promise> { + return lock(lockPath, { + lockfilePath: lockPath, + realpath: false, + stale: 10_000, + update: 3_000, + retries: 0 + }) +} + +async function acquireAcquisitionGate(gatePath: string): Promise<() => Promise> { + await ensureAcquisitionGateDirectory(gatePath) + // A per-target Lamport bakery queue avoids deleting a shared stale gate. Every removable path + // contains a nonce and is owned by one PID, so crash recovery cannot unlink a replacement owner. + const contenderId = `${process.pid}-${randomUUID()}` + const choosingPath = join(gatePath, `choosing.${contenderId}`) + let ticketPath: string | null = null + let ticket: number | null = null + + await writeFile(choosingPath, '', { flag: 'wx', mode: 0o600 }) + try { + const { tickets } = await readGateState(gatePath, contenderId) + ticket = Math.max(0, ...tickets.map(contender => contender.ticket)) + 1 + ticketPath = join(gatePath, `ticket.${ticket}.${contenderId}`) + await rename(choosingPath, ticketPath) + + await waitForAcquisitionTurn(gatePath, { id: contenderId, ticket }) + } catch (operationError) { + const cleanupErrors = await removeContenderFiles(choosingPath, ...(ticketPath === null ? [] : [ticketPath])) + if (cleanupErrors.length > 0) { + throw new AggregateError([operationError, ...cleanupErrors], 'acquisition gate attempt and cleanup both failed') + } + throw operationError + } + + return async () => { + try { + await unlink(ticketPath!) + } catch (error) { + if (hasErrorCode(error, 'ENOENT')) throw compromisedGateError(ticketPath!) + throw error + } + } +} + +async function ensureAcquisitionGateDirectory(gatePath: string): Promise { + try { + await mkdir(gatePath, { mode: 0o700 }) + } catch (error) { + if (!hasErrorCode(error, 'EEXIST')) throw error + } + const details = await lstat(gatePath) + if (!details.isDirectory() || details.isSymbolicLink()) { + throw new Error(`unsafe SkillHub CLI acquisition gate directory: ${gatePath}`) + } +} + +interface AcquisitionContender { + id: string + ticket: number +} + +interface AcquisitionGateState { + tickets: AcquisitionContender[] + hasLiveChoosing: boolean +} + +async function readGateState(gatePath: string, contenderId: string): Promise { + const entries = await readdir(gatePath) + const tickets: AcquisitionContender[] = [] + let hasLiveChoosing = false + for (const entry of entries) { + const path = join(gatePath, entry) + const ticketMatch = /^ticket\.(\d+)\.(\d+)-([^.]+)$/.exec(entry) + if (ticketMatch) { + const ticket = Number(ticketMatch[1]) + const pid = Number(ticketMatch[2]) + if (!isProcessAlive(pid)) { + await unlinkIfPresent(path) + continue + } + if (!Number.isSafeInteger(ticket) || ticket < 1) throw compromisedGateError(path) + tickets.push({ id: `${ticketMatch[2]}-${ticketMatch[3]}`, ticket }) + continue + } + const choosingMatch = /^choosing\.(\d+)-([^.]+)$/.exec(entry) + if (choosingMatch && `${choosingMatch[1]}-${choosingMatch[2]}` !== contenderId) { + if (!isProcessAlive(Number(choosingMatch[1]))) { + await unlinkIfPresent(path) + } else { + hasLiveChoosing = true + } + } + } + return { tickets, hasLiveChoosing } +} + +async function waitForAcquisitionTurn( + gatePath: string, + contender: AcquisitionContender +): Promise { + let delayMs = 5 + for (;;) { + const state = await readGateState(gatePath, contender.id) + const hasEarlierTicket = state.tickets.some(candidate => compareContenders(candidate, contender) < 0) + if (!state.hasLiveChoosing && !hasEarlierTicket) return + await new Promise(resolve => setTimeout(resolve, delayMs)) + delayMs = Math.min(delayMs * 2, ACQUISITION_GATE_MAX_POLL_MS) + } +} + +function compareContenders(left: AcquisitionContender, right: AcquisitionContender): number { + if (left.ticket !== right.ticket) return left.ticket < right.ticket ? -1 : 1 + if (left.id === right.id) return 0 + return left.id < right.id ? -1 : 1 +} + +function isProcessAlive(pid: number): boolean { + try { + process.kill(pid, 0) + return true + } catch (error) { + return !hasErrorCode(error, 'ESRCH') + } +} + +async function removeContenderFiles(...paths: string[]): Promise { + const errors: Error[] = [] + for (const path of paths) { + try { + await unlink(path) + } catch (error) { + if (!hasErrorCode(error, 'ENOENT')) errors.push(error instanceof Error ? error : new Error(String(error))) + } + } + return errors +} + +async function unlinkIfPresent(path: string): Promise { + try { + await unlink(path) + } catch (error) { + if (!hasErrorCode(error, 'ENOENT')) throw error + } +} + +function compromisedGateError(gatePath: string): Error { + return Object.assign(new Error(`SkillHub CLI acquisition gate was replaced: ${gatePath}`), { + code: 'ECOMPROMISED' + }) +} + +function hasErrorCode(error: unknown, code: string): boolean { + return error instanceof Error && 'code' in error && error.code === code } export async function skillTargetLockPath(rootDir: string, slug: string): Promise { diff --git a/cli/test/unit/services/skill-target-lock.test.ts b/cli/test/unit/services/skill-target-lock.test.ts index ddee6418..7ceae909 100644 --- a/cli/test/unit/services/skill-target-lock.test.ts +++ b/cli/test/unit/services/skill-target-lock.test.ts @@ -1,4 +1,4 @@ -import { access, chmod, lstat, mkdir, mkdtemp, symlink, unlink, utimes, writeFile } from 'node:fs/promises' +import { access, chmod, lstat, mkdir, mkdtemp, readdir, symlink, unlink, utimes, writeFile } from 'node:fs/promises' import { tmpdir } from 'node:os' import { dirname, join } from 'node:path' import { fileURLToPath } from 'node:url' @@ -27,6 +27,21 @@ async function waitForFile(path: string): Promise { throw new Error(`timed out waiting for ${path}`) } +async function waitForExitCount( + processes: Array<{ exited: Promise }>, + count: number +): Promise { + const exitCodes: number[] = [] + for (const process of processes) { + void process.exited.then(exitCode => exitCodes.push(exitCode)) + } + for (let attempt = 0; attempt < 500; attempt++) { + if (exitCodes.length >= count) return exitCodes + await Bun.sleep(10) + } + throw new Error(`timed out waiting for ${count} lock contenders to exit`) +} + describe('skill target lifecycle lock', () => { test('creates or repairs a private lock root and rejects unsafe roots', async () => { const parent = await mkdtemp(join(tmpdir(), 'skillhub-lock-root-')) @@ -65,12 +80,14 @@ describe('skill target lifecycle lock', () => { await mkdir(lockPath) const staleTime = new Date(Date.now() - 60_000) await utimes(lockPath, staleTime, staleTime) + const acquisitionGatePath = `${lockPath}.acquire` const worker = fileURLToPath(new URL('../../helpers/target-lock-worker.ts', import.meta.url)) const bunPath = (await Bun.which('bun')) ?? process.execPath const acquiredPath = join(rootDir, 'acquired') const releasePath = join(rootDir, 'release') const startPath = join(rootDir, 'start') - const readyPaths = [join(rootDir, 'ready-0'), join(rootDir, 'ready-1')] + const workerCount = 8 + const readyPaths = Array.from({ length: workerCount }, (_, index) => join(rootDir, `ready-${index}`)) const processes = readyPaths.map(readyPath => Bun.spawn({ cmd: [bunPath, worker, rootDir, 'demo', readyPath, startPath, acquiredPath, releasePath], @@ -81,11 +98,8 @@ describe('skill target lifecycle lock', () => { await Promise.all(readyPaths.map(waitForFile)) await writeFile(startPath, 'start') await waitForFile(acquiredPath) - const loserExitCode = await Promise.race([ - ...processes.map(process => process.exited), - Bun.sleep(5_000).then(() => { throw new Error('timed out waiting for the lock loser') }) - ]) - expect(loserExitCode).toBe(4) + const loserExitCodes = await waitForExitCount(processes, workerCount - 1) + expect(loserExitCodes).toEqual(Array(workerCount - 1).fill(4)) } finally { try { await writeFile(releasePath, 'release') @@ -106,14 +120,83 @@ describe('skill target lifecycle lock', () => { stderr: (await new Response(process.stderr).text()).trim() }))) - expect(results.map(result => result.exitCode).sort()).toEqual([0, 4]) + expect(results.map(result => result.exitCode).sort()).toEqual([0, ...Array(workerCount - 1).fill(4)]) expect(results.filter(result => result.stdout === 'acquired')).toHaveLength(1) + for (const loser of results.filter(result => result.exitCode === 4)) { + expect(loser.stderr).toContain('install target is busy') + } expect(await exists(lockPath)).toBe(false) + expect(await readdir(acquisitionGatePath)).toEqual([]) + const releaseAfterContention = await acquireSkillTargetLock(rootDir, 'demo') + await releaseAfterContention() if (process.platform !== 'win32') { expect((await lstat(dirname(lockPath))).mode & 0o077).toBe(0) } }, 15_000) + test('does not recover a live acquisition contender solely because its file is old', async () => { + const rootDir = await mkdtemp(join(tmpdir(), 'skillhub-target-live-gate-')) + const lockPath = await skillTargetLockPath(rootDir, 'demo') + const acquisitionGatePath = `${lockPath}.acquire` + await mkdir(acquisitionGatePath) + const liveContenderPath = join(acquisitionGatePath, `choosing.${process.pid}-suspended`) + await writeFile(liveContenderPath, '') + const staleTime = new Date(Date.now() - 60_000) + await utimes(liveContenderPath, staleTime, staleTime) + + const acquisition = acquireSkillTargetLock(rootDir, 'demo') + const state = await Promise.race([ + acquisition.then(() => 'acquired', () => 'rejected'), + Bun.sleep(50).then(() => 'waiting') + ]) + expect(state).toBe('waiting') + expect(await exists(liveContenderPath)).toBe(true) + await unlink(liveContenderPath) + const release = await acquisition + await release() + }) + + test('does not pass a live acquisition ticket', async () => { + const rootDir = await mkdtemp(join(tmpdir(), 'skillhub-target-live-ticket-')) + const lockPath = await skillTargetLockPath(rootDir, 'demo') + const acquisitionGatePath = `${lockPath}.acquire` + await mkdir(acquisitionGatePath) + const liveTicketPath = join(acquisitionGatePath, `ticket.1.${process.pid}-suspended`) + await writeFile(liveTicketPath, '') + + const acquisition = acquireSkillTargetLock(rootDir, 'demo') + const state = await Promise.race([ + acquisition.then(() => 'acquired', () => 'rejected'), + Bun.sleep(50).then(() => 'waiting') + ]) + expect(state).toBe('waiting') + expect(await exists(liveTicketPath)).toBe(true) + await unlink(liveTicketPath) + const release = await acquisition + await release() + }) + + test('recovers acquisition contenders whose owner process exited', async () => { + const rootDir = await mkdtemp(join(tmpdir(), 'skillhub-target-dead-gate-')) + const lockPath = await skillTargetLockPath(rootDir, 'demo') + await mkdir(lockPath) + const staleTime = new Date(Date.now() - 60_000) + await utimes(lockPath, staleTime, staleTime) + const acquisitionGatePath = `${lockPath}.acquire` + await mkdir(acquisitionGatePath) + const bunPath = (await Bun.which('bun')) ?? process.execPath + const exitedOwner = Bun.spawn({ cmd: [bunPath, '-e', ''], stdout: 'ignore', stderr: 'ignore' }) + const deadPid = exitedOwner.pid + expect(await exitedOwner.exited).toBe(0) + await writeFile(join(acquisitionGatePath, `choosing.${deadPid}-abandoned`), '') + await writeFile(join(acquisitionGatePath, `ticket.1.${deadPid}-abandoned`), '') + + const release = await acquireSkillTargetLock(rootDir, 'demo') + await release() + expect(await exists(lockPath)).toBe(false) + expect(await readdir(acquisitionGatePath)).toEqual([]) + }) + test('keeps one lock identity when a symlink target is removed', async () => { const rootDir = await mkdtemp(join(tmpdir(), 'skillhub-target-symlink-root-')) const linkedDir = await mkdtemp(join(tmpdir(), 'skillhub-target-symlink-value-'))