refactor: consolidate lock implementation into Lock class

Lock now uses Utopia\Lock\Distributed directly and owns the full
acquire/release/telemetry/error-reporting/fail-open/kill-switch logic
that previously lived in two inline DI factory closures.

Adds withKey($key, $fn, $ttl, $orFail, $waitTimeout) as a generic
escape hatch for non-platform key shapes (cache, queue, edge) and
unusual TTL/timeout requirements.

Per-attribute lock keys for set() so that an accessedAt bump and a
mcpAccessedAt bump on the same projects:{id} document don't compete.
Whole-document operations (run, runOrFail) keep document-level keys.

Removes the standalone distributedLock and distributedLockOrFail DI
factories — Lock is the single API.

request.php shrinks ~150 LOC; Lock.php grows to ~190 LOC.
This commit is contained in:
Prem Palanisamy
2026-04-29 07:41:54 +01:00
parent ce15eeb722
commit e634145612
2 changed files with 150 additions and 172 deletions
+3 -150
View File
@@ -38,7 +38,6 @@ use Utopia\Auth\Proofs\Token;
use Utopia\Auth\Store;
use Utopia\Cache\Cache;
use Utopia\Config\Config;
use Utopia\Console;
use Utopia\Database\Adapter\Pool as DatabasePool;
use Utopia\Database\Database;
use Utopia\Database\DateTime as DatabaseDateTime;
@@ -50,7 +49,6 @@ use Utopia\Domains\Domain;
use Utopia\DSN\DSN;
use Utopia\Http\Http;
use Utopia\Locale\Locale;
use Utopia\Lock\Distributed as DistributedLock;
use Utopia\Logger\Log;
use Utopia\Logger\Logger;
use Utopia\Pools\Group;
@@ -77,154 +75,9 @@ return function (Container $container): void {
return $register->get('logger');
}, ['register']);
// Rate-limited to one push per 60s per (action, target) so a sustained
// backend outage doesn't flood Sentry across the pod fleet.
$lockErrorReporter = function (Log $log, ?Logger $logger, string $action, string $key, string $target, Throwable $e): void {
static $lastReportAt = [];
Console::warning("Lock {$action} for {$key}: {$e->getMessage()}");
if ($logger === null) {
return;
}
$bucket = $action . ':' . $target;
$now = time();
if (($lastReportAt[$bucket] ?? 0) + 60 > $now) {
return;
}
$lastReportAt[$bucket] = $now;
$log->setNamespace('http');
$log->setServer(System::getEnv('_APP_LOGGING_SERVICE_IDENTIFIER', \gethostname()));
$log->setVersion(APP_VERSION_STABLE);
$log->setType(Log::TYPE_WARNING);
$log->setMessage('Distributed lock ' . $action . ': ' . $e->getMessage());
$log->setAction("lock.{$action}");
$log->setEnvironment(System::getEnv('_APP_ENV', 'development') === 'production'
? Log::ENVIRONMENT_PRODUCTION
: Log::ENVIRONMENT_STAGING);
$log->addTag('lock.target', $target);
// Strip trailing document ID to keep aggregator cardinality bounded.
$log->addTag('lock.key_pattern', preg_replace('/:[^:]+$/', ':*', $key));
$log->addTag('code', $e->getCode());
$log->addExtra('file', $e->getFile());
$log->addExtra('line', $e->getLine());
$log->addExtra('trace', $e->getTraceAsString());
try {
$logger->addLog($log);
} catch (Throwable) {
}
};
$lockTargetOf = function (string $key): string {
$parts = explode(':', $key, 4);
return $parts[2] ?? 'unknown';
};
/**
* Skip-on-contention. For idempotent writes where losing the race is correct
* (e.g., the winning pod's update covers ours). Fail-open on backend error.
*/
$container->set('distributedLock', function (\Redis $redis, Telemetry $telemetry, Log $log, ?Logger $logger) use ($lockTargetOf, $lockErrorReporter) {
$enabled = System::getEnv('_APP_LOCKING_ENABLED', 'enabled') !== 'disabled';
$attempts = $telemetry->createCounter('lock.attempts', null, 'Distributed lock acquire outcomes');
if (! $enabled) {
return function (string $key, \Closure $fn, float $ttl = 5.0): void {
$fn();
};
}
return function (string $key, \Closure $fn, float $ttl = 5.0) use ($redis, $attempts, $log, $logger, $lockTargetOf, $lockErrorReporter): void {
$target = $lockTargetOf($key);
$lock = new DistributedLock($redis, $key, (int) $ttl);
try {
$acquired = $lock->tryAcquire();
} catch (\RedisException $e) {
$attempts->add(1, ['outcome' => 'backend_error', 'target' => $target]);
$lockErrorReporter($log, $logger, 'backend_error', $key, $target, $e);
$fn();
return;
}
if (! $acquired) {
$attempts->add(1, ['outcome' => 'skipped', 'target' => $target]);
return;
}
$attempts->add(1, ['outcome' => 'acquired', 'target' => $target]);
try {
$fn();
} finally {
try {
$lock->release();
} catch (\Throwable $e) {
$attempts->add(1, ['outcome' => 'release_error', 'target' => $target]);
$lockErrorReporter($log, $logger, 'release_error', $key, $target, $e);
}
}
};
}, ['redis', 'telemetry', 'log', 'logger']);
/**
* Block-then-409 on contention. For read-modify-write on shared mutable state
* where silently dropping a request is wrong. Fail-open on backend error.
*/
$container->set('distributedLockOrFail', function (\Redis $redis, Telemetry $telemetry, Log $log, ?Logger $logger) use ($lockTargetOf, $lockErrorReporter) {
$enabled = System::getEnv('_APP_LOCKING_ENABLED', 'enabled') !== 'disabled';
$attempts = $telemetry->createCounter('lock.attempts', null, 'Distributed lock acquire outcomes');
if (! $enabled) {
return function (string $key, \Closure $fn, float $ttl = 10.0, float $waitTimeout = 3.0): mixed {
return $fn();
};
}
return function (string $key, \Closure $fn, float $ttl = 10.0, float $waitTimeout = 3.0) use ($redis, $attempts, $log, $logger, $lockTargetOf, $lockErrorReporter): mixed {
$target = $lockTargetOf($key);
$lock = new DistributedLock($redis, $key, (int) $ttl);
try {
$acquired = $lock->acquire($waitTimeout);
} catch (\RedisException $e) {
$attempts->add(1, ['outcome' => 'backend_error', 'target' => $target]);
$lockErrorReporter($log, $logger, 'backend_error', $key, $target, $e);
return $fn();
}
if (! $acquired) {
$attempts->add(1, ['outcome' => 'contended', 'target' => $target]);
// No custom message — the lock key embeds collection + document id.
throw new Exception(Exception::GENERAL_RESOURCE_LOCKED);
}
$attempts->add(1, ['outcome' => 'acquired', 'target' => $target]);
try {
return $fn();
} finally {
try {
$lock->release();
} catch (\Throwable $e) {
$attempts->add(1, ['outcome' => 'release_error', 'target' => $target]);
$lockErrorReporter($log, $logger, 'release_error', $key, $target, $e);
}
}
};
}, ['redis', 'telemetry', 'log', 'logger']);
$container->set('lock', function (callable $distributedLock, callable $distributedLockOrFail, Database $dbForPlatform, Authorization $authorization): Lock {
return new Lock(
\Closure::fromCallable($distributedLock),
\Closure::fromCallable($distributedLockOrFail),
$dbForPlatform,
$authorization,
);
}, ['distributedLock', 'distributedLockOrFail', 'dbForPlatform', 'authorization']);
$container->set('lock', function (\Redis $redis, Telemetry $telemetry, Database $dbForPlatform, Authorization $authorization, Log $log, ?Logger $logger): Lock {
return new Lock($redis, $telemetry, $dbForPlatform, $authorization, $log, $logger);
}, ['redis', 'telemetry', 'dbForPlatform', 'authorization', 'log', 'logger']);
$container->set('authorization', function () {
return new Authorization();
+147 -22
View File
@@ -2,26 +2,46 @@
namespace Appwrite\Locking;
use Appwrite\Extend\Exception;
use Closure;
use Throwable;
use Utopia\Console;
use Utopia\Database\Database;
use Utopia\Database\DateTime;
use Utopia\Database\Document;
use Utopia\Database\Validator\Authorization;
use Utopia\Lock\Distributed as DistributedLock;
use Utopia\Logger\Log;
use Utopia\Logger\Logger;
use Utopia\System\System;
use Utopia\Telemetry\Adapter as Telemetry;
final class Lock
{
private readonly bool $enabled;
private readonly mixed $attempts;
/** @var array<string,int> */
private static array $lastReportAt = [];
public function __construct(
private readonly Closure $skipLock,
private readonly Closure $failLock,
private readonly \Redis $redis,
Telemetry $telemetry,
private readonly Database $dbForPlatform,
private readonly Authorization $authorization,
) {}
private readonly Log $log,
private readonly ?Logger $logger,
) {
$this->enabled = System::getEnv('_APP_LOCKING_ENABLED', 'enabled') !== 'disabled';
$this->attempts = $telemetry->createCounter('lock.attempts', null, 'Distributed lock acquire outcomes');
}
/**
* Throttled single-attribute write under a skip-on-contention lock with
* authorization bypass. For idempotent timestamp-style updates (accessedAt,
* mcpAccessedAt) where regional pods writing the same value would thrash
* the platform DB.
* Throttled single-attribute write under a per-attribute skip-on-contention
* lock with authorization bypass. For idempotent timestamp-style updates
* (accessedAt, mcpAccessedAt) where regional pods writing the same value
* would thrash the platform DB.
*/
public function set(
string $collection,
@@ -29,35 +49,140 @@ final class Lock
string $attribute = 'accessedAt',
?string $value = null,
): void {
($this->skipLock)(self::key($collection, $id), function () use ($collection, $id, $attribute, $value) {
$this->authorization->skip(fn () => $this->dbForPlatform->updateDocument(
$collection,
$id,
new Document([$attribute => $value ?? DateTime::now()])
));
});
$this->withKey(
"lock:platform:{$collection}:{$id}:{$attribute}",
function () use ($collection, $id, $attribute, $value) {
$this->authorization->skip(fn () => $this->dbForPlatform->updateDocument(
$collection,
$id,
new Document([$attribute => $value ?? DateTime::now()])
));
}
);
}
/**
* Skip-on-contention lock around an arbitrary callback. For idempotent
* writes that don't fit the set shape (e.g., updates with cache purge).
* Skip-on-contention lock around an arbitrary callback for a platform
* document. For idempotent multi-statement writes that don't fit `set`.
*/
public function run(string $collection, string $id, Closure $fn): void
{
($this->skipLock)(self::key($collection, $id), $fn);
$this->withKey("lock:platform:{$collection}:{$id}", $fn);
}
/**
* Block-then-409 lock around an arbitrary callback. For read-modify-write
* endpoints where silently dropping a concurrent request would lose data.
* Block-then-409 lock around an arbitrary callback for a platform document.
* For read-modify-write endpoints where silently dropping a concurrent
* request would lose user data.
*/
public function runOrFail(string $collection, string $id, Closure $fn): mixed
{
return ($this->failLock)(self::key($collection, $id), $fn);
return $this->withKey(
"lock:platform:{$collection}:{$id}",
$fn,
ttl: 10,
orFail: true,
);
}
private static function key(string $collection, string $id): string
/**
* Generic lock primitive with full control over key, TTL, contention
* behavior, and wait timeout. Escape hatch for non-platform keys
* (cache, queue, edge) and for unusual TTL/timeout requirements.
*/
public function withKey(
string $key,
Closure $fn,
int $ttl = 5,
bool $orFail = false,
float $waitTimeout = 3.0,
): mixed {
if (! $this->enabled) {
return $fn();
}
$target = self::targetOf($key);
$lock = new DistributedLock($this->redis, $key, $ttl);
try {
$acquired = $orFail ? $lock->acquire($waitTimeout) : $lock->tryAcquire();
} catch (\RedisException $e) {
$this->attempts->add(1, ['outcome' => 'backend_error', 'target' => $target]);
$this->reportError('backend_error', $key, $target, $e);
return $fn();
}
if (! $acquired) {
if ($orFail) {
$this->attempts->add(1, ['outcome' => 'contended', 'target' => $target]);
// No custom message — the lock key embeds collection + document id.
throw new Exception(Exception::GENERAL_RESOURCE_LOCKED);
}
$this->attempts->add(1, ['outcome' => 'skipped', 'target' => $target]);
return null;
}
$this->attempts->add(1, ['outcome' => 'acquired', 'target' => $target]);
try {
return $fn();
} finally {
try {
$lock->release();
} catch (Throwable $e) {
$this->attempts->add(1, ['outcome' => 'release_error', 'target' => $target]);
$this->reportError('release_error', $key, $target, $e);
}
}
}
private static function targetOf(string $key): string
{
return "lock:platform:{$collection}:{$id}";
$parts = explode(':', $key, 4);
return $parts[2] ?? 'unknown';
}
/**
* Rate-limited to one push per 60s per (action, target) so a sustained
* backend outage doesn't flood Sentry across the pod fleet.
*/
private function reportError(string $action, string $key, string $target, Throwable $e): void
{
Console::warning("Lock {$action} for {$key}: {$e->getMessage()}");
if ($this->logger === null) {
return;
}
$bucket = $action.':'.$target;
$now = time();
if ((self::$lastReportAt[$bucket] ?? 0) + 60 > $now) {
return;
}
self::$lastReportAt[$bucket] = $now;
$this->log->setNamespace('http');
$this->log->setServer(System::getEnv('_APP_LOGGING_SERVICE_IDENTIFIER', \gethostname()));
$this->log->setVersion(APP_VERSION_STABLE);
$this->log->setType(Log::TYPE_WARNING);
$this->log->setMessage('Distributed lock '.$action.': '.$e->getMessage());
$this->log->setAction("lock.{$action}");
$this->log->setEnvironment(System::getEnv('_APP_ENV', 'development') === 'production'
? Log::ENVIRONMENT_PRODUCTION
: Log::ENVIRONMENT_STAGING);
$this->log->addTag('lock.target', $target);
// Strip trailing document ID to keep aggregator cardinality bounded.
$this->log->addTag('lock.key_pattern', preg_replace('/:[^:]+$/', ':*', $key));
$this->log->addTag('code', $e->getCode());
$this->log->addExtra('file', $e->getFile());
$this->log->addExtra('line', $e->getLine());
$this->log->addExtra('trace', $e->getTraceAsString());
try {
$this->logger->addLog($this->log);
} catch (Throwable) {
}
}
}