-
Notifications
You must be signed in to change notification settings - Fork 0
fix(storage): prevent concurrent chunk finalization #8
New issue
Have a question about this project? Sign up for a free GitHub account to open an issue and contact its maintainers and the community.
By clicking “Sign up for GitHub”, you agree to our terms of service and privacy statement. We’ll occasionally send you account related emails.
Already on GitHub? Sign in to your account
base: qa/agent-appwrite-appwrite/pr-08-13505/base
Are you sure you want to change the base?
Changes from all commits
File filter
Filter by extension
Conversations
Jump to
Diff view
Diff view
There are no files selected for viewing
| Original file line number | Diff line number | Diff line change |
|---|---|---|
|
|
@@ -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 | ||
|
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. Shipwright · HIGH The LOCK_TTL constant is 600 seconds, but the lock timeout passed to $locks() is 120.0 seconds. Impact: The LOCK_TTL constant is 600 seconds, but the lock timeout passed to $locks() is 120.0 seconds. A reader must understand the difference between TTL (lease duration) and timeout (wait time to acquire) to know why these differ. The comment explains the TTL but not the timeout, and the magic number 120.0 appears twice without explanation. Suggested fix: Review the cited evidence, fix the risk if confirmed, and rerun Shipwright. |
||
| * 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()) { | ||
|
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. Shipwright · CRITICAL The lock token is interpolated into an exception message: 'throw new LockContention('Upload lease lost before transfer: ' . Impact: The lock token is interpolated into an exception message: 'throw new LockContention('Upload lease lost before transfer: ' . $lock->token())'. If the token is sensitive (e.g., a Redis lock token that could be used to release or refresh the lock), leaking it in error responses or logs enables an attacker to hijack the lock. The exception is caught and converted to a generic rate-limit error, but the original messag… Suggested fix: Review the cited evidence, fix the risk if confirmed, and rerun Shipwright. |
||
| throw new LockContention('Upload lease lost before transfer: ' . $lock->token()); | ||
| } | ||
|
|
||
| $chunksUploaded = $deviceForFiles->upload( | ||
|
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. Shipwright · CRITICAL The lock is refreshed before upload but not before finalizeUpload. Impact: The lock is refreshed before upload but not before finalizeUpload. After a long transfer, the lease can expire between the isHeld() check and finalizeUpload's document write, allowing two requests to finalize concurrently and corrupt chunk accounting. The isHeld() check is TOCTOU: it verifies ownership, then finalizeUpload performs multiple non-atomic DB operations without re-verifying or refreshing the lease. Suggested fix: Review the cited evidence, fix the risk if confirmed, and rerun Shipwright. |
||
| $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.'); | ||
|
|
||
| Original file line number | Diff line number | Diff line change |
|---|---|---|
|
|
@@ -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 | ||
|
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. Shipwright · HIGH The test provider name 'duplicate chunks' is misleading. Impact: The test provider name 'duplicate chunks' is misleading. The test does not verify that duplicate chunks are handled correctly; it sends duplicate requests and asserts all return 200/201. A new hire reading this test would assume duplicate chunk handling is validated, but the assertions only check response codes and final state, not that the duplicate was deduplicated or rejected. Suggested fix: Review the cited evidence, fix the risk if confirmed, and rerun Shipwright. |
||
| { | ||
| 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) { | ||
|
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. Shipwright · CRITICAL The duplicate-chunk test path sends the same chunk twice concurrently. Impact: The duplicate-chunk test path sends the same chunk twice concurrently. If both requests pass prepareUpload before either finalizes, both will call deviceForFiles->upload() with the same chunk. Depending on the device implementation, this can double-count chunksUploaded or corrupt the multipart upload state. The lock serializes the critical section, but the second request re-reads the file document inside prepareU… Suggested fix: Review the cited evidence, fix the risk if confirmed, and rerun Shipwright. |
||
| $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. | ||
|
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. Shipwright · HIGH The retry test reuses '$requests[0]['headers']' and '$requests[0]['chunkPath']' to construct a late retry. Impact: The retry test reuses '$requests[0]['headers']' and '$requests[0]['chunkPath']' to construct a late retry. If the original request contained a one-time upload token or session-bound header, this test masks a real-world issue where retries with stale credentials would fail. More importantly, the test asserts the retry returns 200 with the same signature, but does not verify the retry did not re-trigger finalization s… Suggested fix: Review the cited evidence, fix the risk if confirmed, and rerun Shipwright. |
||
| $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'], | ||
|
|
||
There was a problem hiding this comment.
Choose a reason for hiding this comment
The reason will be displayed to describe this comment to others. Learn more.
Shipwright · LOW
The lock callback now passes the Distributed lock object to the callback, but the callback signature in Create.php is typed as 'function (Distributed $lock)'.
Impact: The lock callback now passes the Distributed lock object to the callback, but the callback signature in Create.php is typed as 'function (Distributed $lock)'. The container closure in resources.php invokes '$callback($lock)', which is correct. However, the 'withLock' callback is 'fn () => $callback($lock)', and 'withLock' may invoke the callback with no arguments; this is fine. No defect here.
Suggested fix: Fix the review finding before release.