fix(cliproxy): retry contended state locks

This commit is contained in:
Tam Nhu Tran committed 2026-07-29 12:52:03 -04:00
1 parent fe90591511
commit 66794e0732
3 files changed
+133 -31

No files matched your search

+14 -14
View File
@@ -5,11 +5,11 @@
import * as fs from 'fs';
import * as path from 'path';
import * as lockfile from 'proper-lockfile';
import { CLIProxyProvider } from '../types';
import { PROVIDER_CAPABILITIES } from '../provider-capabilities';
import { PROVIDER_TYPE_VALUES } from '../auth/auth-types';
import { getAuthDir, getCliproxyDir } from '../config/config-generator';
import { withSyncLockRetry } from '../../utils/sync-lock-retry';
import { AccountsRegistry, AccountInfo, PROVIDERS_WITHOUT_EMAIL, DrainOrderConfig } from './types';
import {
getAccountsRegistryPath,
@@ -338,24 +338,24 @@ function recoverAccountsRegistryFromCorruption(registryPath: string): AccountsRe
function withAccountsRegistryLock<T>(callback: () => T): T {
const lockTarget = getCliproxyDir();
let release: (() => void) | undefined;
if (!fs.existsSync(lockTarget)) {
fs.mkdirSync(lockTarget, { recursive: true, mode: 0o700 });
}
try {
release = lockfile.lockSync(lockTarget, { stale: 10000 }) as () => void;
return callback();
} finally {
if (release) {
try {
release();
} catch {
// Best-effort release
}
}
}
// proper-lockfile rejects built-in retries in synchronous mode.
// Keep this API sync while waiting for short-lived concurrent launches.
// Preserve the existing stale-lock window.
// Retry only contention failures.
// Wait in short intervals so released locks are picked up promptly.
// Bound the wait so a genuinely stuck process still produces an error.
// The helper adds recovery guidance when that timeout is reached.
return withSyncLockRetry(lockTarget, callback, {
description: 'CLIProxy account registry lock',
staleMs: 10000,
retryDelayMs: 200,
retryTimeoutMs: 10000,
});
}
function readAccountsRegistryFromDisk(): AccountsRegistry {
+4 -17
View File
@@ -17,10 +17,10 @@
import * as fs from 'fs';
import * as path from 'path';
import * as crypto from 'crypto';
import * as lockfile from 'proper-lockfile';
import { getCliproxyDir } from './config/config-generator';
import { getPortProcess, isCLIProxyProcess } from '../utils/port-utils';
import { CLIPROXY_DEFAULT_PORT } from './config/config-generator';
import { withSyncLockRetry } from '../utils/sync-lock-retry';
/** Session lock file structure */
interface SessionLock {
@@ -61,22 +61,9 @@ function ensureCliproxyDir(): string {
function withSessionTrackerLock<T>(fn: () => T): T {
const dir = ensureCliproxyDir();
let release: (() => void) | undefined;
try {
release = lockfile.lockSync(dir, {
stale: 10000,
}) as () => void;
return fn();
} finally {
if (release) {
try {
release();
} catch {
// Best-effort release
}
}
}
return withSyncLockRetry(dir, fn, {
description: 'CLIProxy session tracker lock',
});
}
/** Get path to session lock file (default port) - kept for future use */
+115
View File
@@ -0,0 +1,115 @@
import * as fs from 'fs';
import { performance } from 'perf_hooks';
import * as lockfile from 'proper-lockfile';
const DEFAULT_STALE_MS = 10000;
const DEFAULT_RETRY_DELAY_MS = 200;
const DEFAULT_RETRY_TIMEOUT_MS = 10000;
interface SyncLockRetryOptions {
description: string;
staleMs?: number;
retryDelayMs?: number;
retryTimeoutMs?: number;
}
export function withSyncLockRetry<T>(
lockTarget: string,
callback: () => T,
options: SyncLockRetryOptions
): T {
const staleMs = options.staleMs ?? DEFAULT_STALE_MS;
const retryDelayMs = options.retryDelayMs ?? DEFAULT_RETRY_DELAY_MS;
const retryTimeoutMs = options.retryTimeoutMs ?? DEFAULT_RETRY_TIMEOUT_MS;
const startedAt = performance.now();
let release: (() => void) | undefined;
let contendedLockPath: string | undefined;
let lastContentionError: unknown;
for (;;) {
if (contendedLockPath && fs.existsSync(contendedLockPath)) {
sleepBeforeRetry(
startedAt,
retryDelayMs,
retryTimeoutMs,
options.description,
lastContentionError
);
continue;
}
if (contendedLockPath && performance.now() - startedAt >= retryTimeoutMs) {
throw buildLockTimeoutError(options.description, retryTimeoutMs, lastContentionError);
}
try {
release = lockfile.lockSync(lockTarget, { stale: staleMs }) as () => void;
break;
} catch (error) {
if (!isLockContentionError(error)) {
throw error;
}
contendedLockPath ??= `${fs.realpathSync(lockTarget)}.lock`;
lastContentionError = error;
sleepBeforeRetry(startedAt, retryDelayMs, retryTimeoutMs, options.description, error);
}
}
try {
return callback();
} finally {
try {
release();
} catch {
// Best-effort release.
}
}
}
function isLockContentionError(error: unknown): boolean {
const code = (error as NodeJS.ErrnoException | undefined)?.code;
return code === 'ELOCKED' || code === 'ENOTACQUIRED';
}
function sleepSync(ms: number): void {
if (ms <= 0) {
return;
}
try {
Atomics.wait(new Int32Array(new SharedArrayBuffer(4)), 0, 0, ms);
} catch {
const deadline = performance.now() + ms;
while (performance.now() < deadline) {
// Fall back for runtimes without Atomics.wait.
}
}
}
function sleepBeforeRetry(
startedAt: number,
retryDelayMs: number,
retryTimeoutMs: number,
description: string,
cause: unknown
): void {
const elapsedMs = performance.now() - startedAt;
if (elapsedMs >= retryTimeoutMs) {
throw buildLockTimeoutError(description, retryTimeoutMs, cause);
}
sleepSync(Math.min(retryDelayMs, retryTimeoutMs - elapsedMs));
}
function buildLockTimeoutError(
description: string,
timeoutMs: number,
cause: unknown
): NodeJS.ErrnoException {
const causeCode = (cause as NodeJS.ErrnoException | undefined)?.code;
const error = new Error(
`Failed to acquire ${description} after ${timeoutMs}ms; another CCS process may still be updating CLIProxy state. Wait for it to finish, then retry.`
) as NodeJS.ErrnoException & { cause?: unknown };
error.code = causeCode === 'ENOTACQUIRED' ? 'ENOTACQUIRED' : 'ELOCKED';
error.cause = cause;
return error;
}