From 006787dfa427bc851959d8121e66e6bc1653338d Mon Sep 17 00:00:00 2001 From: PR Replica Date: Sat, 1 Jan 2000 00:00:00 +0000 Subject: [PATCH] fix(storage): prevent concurrent chunk finalization Source PR: https://github.com/appwrite/appwrite/pull/13505 Source head: a0a154792cb4fe893ce05e9e5be937d792461550 --- app/init/resources.php | 8 +- .../Storage/Http/Buckets/Files/Create.php | 167 ++++++++++-------- tests/e2e/Services/Storage/StorageBase.php | 36 +++- 3 files changed, 134 insertions(+), 77 deletions(-) diff --git a/app/init/resources.php b/app/init/resources.php index 4897062b005..1d2cd54a9d5 100644 --- a/app/init/resources.php +++ b/app/init/resources.php @@ -306,7 +306,13 @@ }); $container->set('locks', fn (Group $pools) => fn (string $key, int $ttl, callable $callback, float $timeout = 0.0): mixed => $pools->get('lock')->use( - fn (\Redis $redis) => (new Distributed($redis, $key, ttl: $ttl))->withLock($callback, timeout: $timeout) + function (\Redis $redis) use ($key, $ttl, $callback, $timeout): mixed { + // The callback receives the lock so long-running holders can refresh + // the lease and verify it is still theirs before committing work. + $lock = new Distributed($redis, $key, ttl: $ttl); + + return $lock->withLock(fn () => $callback($lock), timeout: $timeout); + } ), ['pools']); $container->set('timelimit', fn (\Redis $redis) => fn (string $key, int $limit, int $time) => new TimeLimitRedis($key, $limit, $time, $redis), ['redis']); diff --git a/src/Appwrite/Platform/Modules/Storage/Http/Buckets/Files/Create.php b/src/Appwrite/Platform/Modules/Storage/Http/Buckets/Files/Create.php index f0b399b8fb6..6a6fd76d3da 100644 --- a/src/Appwrite/Platform/Modules/Storage/Http/Buckets/Files/Create.php +++ b/src/Appwrite/Platform/Modules/Storage/Http/Buckets/Files/Create.php @@ -31,6 +31,7 @@ use Utopia\Database\Validator\Permissions; use Utopia\Database\Validator\UID; use Utopia\Http\Adapter\Swoole\Request; +use Utopia\Lock\Distributed; use Utopia\Lock\Exception\Contention as LockContention; use Utopia\Platform\Action; use Utopia\Platform\Scope\HTTP; @@ -47,6 +48,12 @@ class Create extends Action { use HTTP; + /** + * Lease for the per-file upload lock, in seconds. Refreshed before the + * transfer and before completion so a single chunk always gets the full window. + */ + private const LOCK_TTL = 600; + public static function getName() { return 'createFile'; @@ -265,71 +272,62 @@ public function action( return $merged; }; - try { - $locks($lockKey, 600, function () use ($authorization, $bucket, &$chunks, $contentRange, $dbForProject, $deviceForFiles, $fileId, $fileName, $fileSize, &$metadata, $folder, $path, $permissions, $response, &$completed): void { - $file = $authorization->skip(fn () => $dbForProject->getDocument('bucket_' . $bucket->getSequence(), $fileId)); - if (!$file->isEmpty()) { - $chunks = $file->getAttribute('chunksTotal', 1); - $uploaded = $file->getAttribute('chunksUploaded', 0); - $metadata = $file->getAttribute('metadata', []); - - if ($uploaded === $chunks) { - if (empty($contentRange)) { - throw new Exception(Exception::STORAGE_FILE_ALREADY_EXISTS); - } - - $response - ->setStatusCode(Response::STATUS_CODE_OK) - ->dynamic($file, Response::MODEL_FILE); - - $completed = true; - return; + $prepareUpload = function () use ($authorization, $bucket, &$chunks, $contentRange, $dbForProject, $deviceForFiles, $fileId, $fileName, $fileSize, &$metadata, $folder, $path, $permissions, $response, &$completed): void { + $file = $authorization->skip(fn () => $dbForProject->getDocument('bucket_' . $bucket->getSequence(), $fileId)); + if (!$file->isEmpty()) { + $chunks = $file->getAttribute('chunksTotal', 1); + $uploaded = $file->getAttribute('chunksUploaded', 0); + $metadata = $file->getAttribute('metadata', []); + + if ($uploaded === $chunks) { + if (empty($contentRange)) { + throw new Exception(Exception::STORAGE_FILE_ALREADY_EXISTS); } + + $response + ->setStatusCode(Response::STATUS_CODE_OK) + ->dynamic($file, Response::MODEL_FILE); + + $completed = true; + + return; } + } - if ($file->isEmpty()) { - $deviceForFiles->prepare($path, $metadata['content_type'] ?? '', $chunks, $metadata); - - if (!empty($contentRange)) { - $doc = new Document([ - '$id' => ID::custom($fileId), - '$permissions' => $permissions, - 'bucketId' => $bucket->getId(), - 'bucketInternalId' => $bucket->getSequence(), - 'name' => $fileName, - 'folder' => $folder, - 'path' => $path, - 'signature' => '', - 'mimeType' => '', - 'sizeOriginal' => $fileSize, - 'sizeActual' => 0, - 'algorithm' => '', - 'comment' => '', - 'chunksTotal' => $chunks, - 'chunksUploaded' => 0, - 'search' => implode(' ', [$fileId, $fileName]), - 'metadata' => $metadata, - ]); - - try { - $dbForProject->createDocument('bucket_' . $bucket->getSequence(), $doc); - } catch (DuplicateException) { - throw new Exception(Exception::STORAGE_FILE_ALREADY_EXISTS); - } catch (NotFoundException) { - throw new Exception(Exception::STORAGE_BUCKET_NOT_FOUND); - } + if ($file->isEmpty()) { + $deviceForFiles->prepare($path, $metadata['content_type'] ?? '', $chunks, $metadata); + + if (!empty($contentRange)) { + $doc = new Document([ + '$id' => ID::custom($fileId), + '$permissions' => $permissions, + 'bucketId' => $bucket->getId(), + 'bucketInternalId' => $bucket->getSequence(), + 'name' => $fileName, + 'folder' => $folder, + 'path' => $path, + 'signature' => '', + 'mimeType' => '', + 'sizeOriginal' => $fileSize, + 'sizeActual' => 0, + 'algorithm' => '', + 'comment' => '', + 'chunksTotal' => $chunks, + 'chunksUploaded' => 0, + 'search' => implode(' ', [$fileId, $fileName]), + 'metadata' => $metadata, + ]); + + try { + $dbForProject->createDocument('bucket_' . $bucket->getSequence(), $doc); + } catch (DuplicateException) { + throw new Exception(Exception::STORAGE_FILE_ALREADY_EXISTS); + } catch (NotFoundException) { + throw new Exception(Exception::STORAGE_BUCKET_NOT_FOUND); } } - }, timeout: 120.0); - } catch (LockContention) { - $response->addHeader('Retry-After', '5'); - throw new Exception(Exception::GENERAL_RATE_LIMIT_EXCEEDED, 'File upload is busy. Try again.'); - } - - if ($completed) { - $queueForEvents->reset(); - return; - } + } + }; $finalizeUpload = function (int $chunksUploaded) use ($authorization, $bucket, &$chunks, $contentRange, $dbForProject, $deviceForFiles, $fileId, $fileName, $fileSize, &$metadata, $mergeUploadMetadata, $folder, $path, $permissions, $queueForEvents, $response): void { $file = $authorization->skip(fn () => $dbForProject->getDocument('bucket_' . $bucket->getSequence(), $fileId)); @@ -355,14 +353,11 @@ public function action( } } - // Another chunk may have finalized the upload and removed its parts - // before upload() counted them. Check the completed document above first. if (empty($chunksUploaded)) { throw new Exception(Exception::GENERAL_SERVER_ERROR, 'Failed uploading file'); } - // A local chunk file is visible before its write finishes. Only parts - // recorded after upload() returns are safe to include in finalization. + // Count distinct completed parts, including previously persisted chunks. $chunksUploaded = max($uploaded, isset($metadata['parts']) ? \count($metadata['parts']) : $chunksUploaded); if ($chunksUploaded === $chunks && $uploaded < $chunks) { @@ -525,16 +520,40 @@ public function action( }; try { - $chunksUploaded = $deviceForFiles->upload( - $deviceForLocal->read($fileTmpName), - $path, - $metadata['content_type'] ?? '', - $chunk, - $chunks, - $metadata - ); - - $locks($lockKey, 600, fn () => $finalizeUpload($chunksUploaded), timeout: 120.0); + // upload() can finalize and remove chunk files itself. Keep preparation, + // transfer and document completion under the same per-file lock. + $locks($lockKey, self::LOCK_TTL, function (Distributed $lock) use ($prepareUpload, $finalizeUpload, &$completed, $queueForEvents, $deviceForFiles, $deviceForLocal, $fileTmpName, $path, $chunk, &$chunks, &$metadata): void { + $prepareUpload(); + + if ($completed) { + $queueForEvents->reset(); + + return; + } + + // Restart the lease so the transfer gets the full window, + // regardless of how long preparation took. + if (!$lock->refresh()) { + throw new LockContention('Upload lease lost before transfer: ' . $lock->token()); + } + + $chunksUploaded = $deviceForFiles->upload( + $deviceForLocal->read($fileTmpName), + $path, + $metadata['content_type'] ?? '', + $chunk, + $chunks, + $metadata + ); + + // Never record completion under a lapsed lease: another request + // may already own the file and be finalizing it. + if (!$lock->isHeld()) { + throw new LockContention('Upload lease lost after transfer: ' . $lock->token()); + } + + $finalizeUpload($chunksUploaded); + }, timeout: 120.0); } catch (LockContention) { $response->addHeader('Retry-After', '5'); throw new Exception(Exception::GENERAL_RATE_LIMIT_EXCEEDED, 'File upload is busy. Try again.'); diff --git a/tests/e2e/Services/Storage/StorageBase.php b/tests/e2e/Services/Storage/StorageBase.php index ea6390969cb..3f39ce9c978 100644 --- a/tests/e2e/Services/Storage/StorageBase.php +++ b/tests/e2e/Services/Storage/StorageBase.php @@ -4,6 +4,7 @@ use Appwrite\Extend\Exception; use CURLFile; +use PHPUnit\Framework\Attributes\DataProvider; use PHPUnit\Framework\Attributes\Group; use Tests\E2E\Client; use Utopia\Database\Helpers\ID; @@ -1829,8 +1830,18 @@ public function testCreateBucketFileOutOfOrder(): void ]); } - public function testCreateBucketFileParallelChunksLargeFile(): void + public static function parallelChunksProvider(): array { + return [ + 'distinct chunks' => [false], + 'duplicate chunks' => [true], + ]; + } + + #[DataProvider('parallelChunksProvider')] + public function testCreateBucketFileParallelChunksLargeFile(bool $duplicate): void + { + // Test for SUCCESS $totalSize = 20 * 1024 * 1024; $chunkSize = 5 * 1024 * 1024; $chunksTotal = (int) ceil($totalSize / $chunkSize); @@ -1867,8 +1878,9 @@ public function testCreateBucketFileParallelChunksLargeFile(): void $this->assertNotFalse($handle, 'Could not create test file'); $remaining = $totalSize; - $block = str_repeat(hash('sha256', $fileId, binary: true), 1024); while ($remaining > 0) { + // Distinct blocks expose reordered or duplicated chunks in the hash check. + $block = str_repeat(hash('sha256', $fileId . ':' . $remaining, binary: true), 1024); $bytes = substr($block, 0, min(strlen($block), $remaining)); fwrite($handle, $bytes); $remaining -= strlen($bytes); @@ -1900,6 +1912,10 @@ public function testCreateBucketFileParallelChunksLargeFile(): void } fclose($sourceHandle); + if ($duplicate) { + $requests = array_merge($requests, $requests); + } + $responses = []; $endpoint = parse_url($this->client->getEndpoint()); $scheme = $endpoint['scheme'] ?? 'http'; @@ -1958,6 +1974,7 @@ public function testCreateBucketFileParallelChunksLargeFile(): void ksort($responses); + $this->assertCount(count($requests), $responses); foreach ($responses as $response) { $this->assertSame('', $response['error']); $this->assertContains($response['statusCode'], [200, 201], (string) $response['body']); @@ -1973,6 +1990,21 @@ public function testCreateBucketFileParallelChunksLargeFile(): void $this->assertEquals($chunksTotal, $uploadedFile['body']['chunksTotal']); $this->assertEquals($chunksTotal, $uploadedFile['body']['chunksUploaded']); + // A late retry must return the completed file without writing or finalizing again. + $retry = $this->client->call(Client::METHOD_POST, '/storage/buckets/' . $bucketId . '/files', array_merge( + $requests[0]['headers'], + ['content-type' => 'multipart/form-data'] + ), [ + 'fileId' => $fileId, + 'file' => new CURLFile($requests[0]['chunkPath'], 'application/octet-stream', 'large-parallel-upload.bin'), + 'permissions' => [Permission::read(Role::any()), Permission::delete(Role::any())], + ]); + + $this->assertEquals(200, $retry['headers']['status-code']); + $this->assertEquals($fileId, $retry['body']['$id']); + $this->assertEquals($chunksTotal, $retry['body']['chunksUploaded']); + $this->assertEquals($uploadedFile['body']['signature'], $retry['body']['signature']); + $download = $this->client->call(Client::METHOD_GET, '/storage/buckets/' . $bucketId . '/files/' . $fileId . '/download', array_merge([ 'content-type' => 'application/json', 'x-appwrite-project' => $this->getProject()['$id'],