From e63414561206ea7b2420737ef4fe79455f24701f Mon Sep 17 00:00:00 2001 From: Prem Palanisamy Date: Wed, 29 Apr 2026 07:41:54 +0100 Subject: [PATCH] refactor: consolidate lock implementation into Lock class MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 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. --- app/init/resources/request.php | 153 +---------------------------- src/Appwrite/Locking/Lock.php | 169 ++++++++++++++++++++++++++++----- 2 files changed, 150 insertions(+), 172 deletions(-) diff --git a/app/init/resources/request.php b/app/init/resources/request.php index c5fafbef47..d927c1dff8 100644 --- a/app/init/resources/request.php +++ b/app/init/resources/request.php @@ -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(); diff --git a/src/Appwrite/Locking/Lock.php b/src/Appwrite/Locking/Lock.php index ec85fe0fa6..288b50c6ea 100644 --- a/src/Appwrite/Locking/Lock.php +++ b/src/Appwrite/Locking/Lock.php @@ -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 */ + 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) { + } } }