mirror of
https://github.com/tiennm99/ccs.git
synced 2026-10-03 20:13:02 +00:00
fix(cliproxy): retry contended state locks
This commit is contained in:
1 parent
fe90591511
commit
66794e0732
3 files changed
+133
-31
No files matched your search
@@ -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 {
|
||||
|
||||
@@ -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 */
|
||||
|
||||
@@ -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;
|
||||
}
|
||||
Reference in new issue
Block a user