Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
8 changes: 7 additions & 1 deletion app/init/resources.php
Original file line number Diff line number Diff line change
Expand Up @@ -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 {

Copy link
Copy Markdown

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.

// 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']);
Expand Down
167 changes: 93 additions & 74 deletions src/Appwrite/Platform/Modules/Storage/Http/Buckets/Files/Create.php
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand All @@ -47,6 +48,12 @@ class Create extends Action
{
use HTTP;

/**
* Lease for the per-file upload lock, in seconds. Refreshed before the

Copy link
Copy Markdown

Choose a reason for hiding this comment

The 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';
Expand Down Expand Up @@ -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));
Expand All @@ -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) {
Expand Down Expand Up @@ -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()) {

Copy link
Copy Markdown

Choose a reason for hiding this comment

The 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(

Copy link
Copy Markdown

Choose a reason for hiding this comment

The 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.');
Expand Down
36 changes: 34 additions & 2 deletions tests/e2e/Services/Storage/StorageBase.php
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -1829,8 +1830,18 @@ public function testCreateBucketFileOutOfOrder(): void
]);
}

public function testCreateBucketFileParallelChunksLargeFile(): void
public static function parallelChunksProvider(): array

Copy link
Copy Markdown

Choose a reason for hiding this comment

The 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);
Expand Down Expand Up @@ -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);
Expand Down Expand Up @@ -1900,6 +1912,10 @@ public function testCreateBucketFileParallelChunksLargeFile(): void
}
fclose($sourceHandle);

if ($duplicate) {

Copy link
Copy Markdown

Choose a reason for hiding this comment

The 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';
Expand Down Expand Up @@ -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']);
Expand All @@ -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.

Copy link
Copy Markdown

Choose a reason for hiding this comment

The 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'],
Expand Down