Chunk orchestrator build artifact uploads

This commit is contained in:
Chirag Aggarwal
2026-05-21 17:50:19 +05:30
parent 25b7d5d574
commit 8b84d73df8
5 changed files with 254 additions and 112 deletions
+2 -26
View File
@@ -55,11 +55,7 @@ $container->set('pools', function ($register) {
return $register->get('pools');
}, ['register']);
$legacyPayloadSize = 12 * (1024 * 1024);
$payloadSize = \max(
$legacyPayloadSize,
(int) System::getEnv('_APP_COMPUTE_SIZE_LIMIT', 0)
) + (2 * 1024 * 1024); // Add a small buffer for multipart/request overhead.
$payloadSize = 12 * (1024 * 1024); // 12MB - adding slight buffer for headers and other data that might be sent with the payload - update later with valid testing
$totalWorkers = intval(System::getEnv('_APP_CPU_NUM', swoole_cpu_num())) * intval(System::getEnv('_APP_WORKER_PER_CORE', 6));
$swoole = new Server(
@@ -503,7 +499,7 @@ $http->on(Constant::EVENT_START, function ($http) use ($payloadSize, $totalWorke
});
});
$swoole->onRequest(function ($utopiaRequest, $utopiaResponse) use ($files, $swoole, $setRequestContext, $legacyPayloadSize) {
$swoole->onRequest(function ($utopiaRequest, $utopiaResponse) use ($files, $swoole, $setRequestContext) {
Span::init('http.request');
$request = new Request($utopiaRequest->getSwooleRequest());
@@ -511,26 +507,6 @@ $swoole->onRequest(function ($utopiaRequest, $utopiaResponse) use ($files, $swoo
Span::add('http.method', $request->getMethod());
$path = \parse_url($request->getURI(), PHP_URL_PATH) ?? '';
$isOrchestratorBuildArtifactUpload = $request->getMethod() === 'PUT'
&& \preg_match('#^/v1/(functions|sites)/[^/]+/deployments/[^/]+/artifacts/build$#', $path) === 1;
if (!$isOrchestratorBuildArtifactUpload && (int)$request->getHeader('content-length', '0') > $legacyPayloadSize) {
$response
->setStatusCode(Response::STATUS_CODE_REQUEST_ENTITY_TOO_LARGE)
->json([
'message' => 'Request Entity Too Large',
'code' => Response::STATUS_CODE_REQUEST_ENTITY_TOO_LARGE,
'type' => 'general_query_limit_exceeded',
'version' => APP_VERSION_STABLE,
]);
Span::add('http.response.code', $response->getStatusCode());
Span::current()?->finish();
return;
}
if ($files->isFileLoaded($request->getURI())) {
$time = (60 * 60 * 24 * 45); // 45 days cache
+2 -2
View File
@@ -1232,7 +1232,7 @@ services:
appwrite-orchestrator:
container_name: appwrite-orchestrator
image: chiragagg5k/jobs-service:temp-appwrite-pr-4f081ed-20260521-r2
image: chiragagg5k/jobs-service:temp-appwrite-pr-4f081ed-20260521-r3
user: root
restart: unless-stopped
networks:
@@ -1246,7 +1246,7 @@ services:
- ORCHESTRATOR_BACKEND=docker
- PORT=8080
- METRICS_PORT=9090
- SIDECAR_IMAGE=chiragagg5k/job-sidecar:temp-appwrite-pr-4f081ed-20260521
- SIDECAR_IMAGE=chiragagg5k/job-sidecar:temp-appwrite-pr-4f081ed-20260521-r2
- ARTIFACT_ENDPOINT=http://host.docker.internal:9080
- EXTRA_HOSTS=appwrite.test:host-gateway
- ORCHESTRATOR_NETWORK=appwrite
@@ -0,0 +1,214 @@
<?php
namespace Appwrite\Platform\Modules\Compute\Http\Deployments\Artifacts\Build;
use Appwrite\Event\Message\Build as BuildMessage;
use Appwrite\Event\Publisher\Build as BuildPublisher;
use Appwrite\Extend\Exception;
use Appwrite\Utopia\Response;
use Utopia\Cache\Cache;
use Utopia\Database\Database;
use Utopia\Database\Document;
use Utopia\Http\Adapter\Swoole\Request;
use Utopia\Lock\Exception\Contention as LockContention;
use Utopia\Storage\Device;
trait ChunkedBuildArtifact
{
protected function uploadBuildArtifact(
string $deploymentId,
Document $project,
Document $resource,
Document $deployment,
Request $request,
Response $response,
Database $dbForProject,
Device $deviceForBuilds,
BuildPublisher $publisherForBuilds,
Cache $cache,
callable $locks
): void {
$path = $deviceForBuilds->getPath($deploymentId . '.tar.gz');
$contentRange = $request->getHeader('content-range');
$chunk = 1;
$chunks = 1;
$fileSize = null;
if (!empty($contentRange)) {
$start = $request->getContentRangeStart();
$end = $request->getContentRangeEnd();
$fileSize = $request->getContentRangeSize();
// TODO make `end >= $fileSize` in next breaking version
if (\is_null($start) || \is_null($end) || \is_null($fileSize) || $end > $fileSize) {
throw new Exception(Exception::STORAGE_INVALID_CONTENT_RANGE);
}
$chunks = (int) \ceil($fileSize / APP_LIMIT_UPLOAD_CHUNK_SIZE);
$chunk = (int) ($start / APP_LIMIT_UPLOAD_CHUNK_SIZE) + 1;
}
$tmp = \tempnam(\sys_get_temp_dir(), 'appwrite-build-');
if ($tmp === false) {
throw new Exception(Exception::GENERAL_SERVER_ERROR, 'Failed creating temporary build artifact file.');
}
try {
if (\file_put_contents($tmp, $request->getRawPayload()) === false) {
throw new Exception(Exception::GENERAL_SERVER_ERROR, 'Failed writing build artifact chunk.');
}
$fileSize ??= \filesize($tmp);
if ($fileSize === false) {
throw new Exception(Exception::GENERAL_SERVER_ERROR, 'Failed reading build artifact chunk.');
}
$metadata = ['content_type' => 'application/gzip'];
$cacheKey = 'build-artifact:' . $project->getId() . ':' . $deploymentId;
$cacheTtl = 60 * 60 * 24;
$lockKey = 'builds:artifact:' . $project->getId() . ':' . $deploymentId;
$completed = false;
try {
$locks($lockKey, 600, function () use ($cache, $cacheKey, $cacheTtl, $dbForProject, $deploymentId, &$chunks, $deviceForBuilds, &$metadata, $path, $response, &$completed): void {
$deployment = $dbForProject->getAuthorization()->skip(fn () => $dbForProject->getDocument('deployments', $deploymentId));
if (!empty($deployment->getAttribute('buildPath', ''))) {
$response
->setStatusCode(Response::STATUS_CODE_ACCEPTED)
->json([
'path' => $deployment->getAttribute('buildPath'),
'size' => $deployment->getAttribute('buildSize', 0),
]);
$completed = true;
return;
}
$stored = $cache->load($cacheKey, $cacheTtl);
if (\is_array($stored)) {
$metadata = $this->mergeArtifactUploadMetadata($stored, $metadata);
$chunks = (int) ($stored['chunksTotal'] ?? $chunks);
}
$deviceForBuilds->prepareUpload($path, 'application/gzip', $chunks, $metadata);
$metadata['chunksTotal'] = $chunks;
$cache->save($cacheKey, $metadata);
}, timeout: 120.0);
} catch (LockContention) {
$response->addHeader('Retry-After', '5');
throw new Exception(Exception::GENERAL_RATE_LIMIT_EXCEEDED, 'Build artifact upload is busy. Try again.');
}
if ($completed) {
return;
}
$chunksUploaded = $deviceForBuilds->uploadChunk($tmp, $path, $chunk, $chunks, $metadata);
if (empty($chunksUploaded)) {
throw new Exception(Exception::GENERAL_SERVER_ERROR, 'Failed storing build artifact chunk.');
}
try {
$locks($lockKey, 600, function () use ($cache, $cacheKey, $cacheTtl, &$chunks, $chunksUploaded, $dbForProject, $deploymentId, $deviceForBuilds, $fileSize, &$metadata, $path, $project, $publisherForBuilds, $resource, $response): void {
$deployment = $dbForProject->getAuthorization()->skip(fn () => $dbForProject->getDocument('deployments', $deploymentId));
if (!empty($deployment->getAttribute('buildPath', ''))) {
$cache->purge($cacheKey);
$response
->setStatusCode(Response::STATUS_CODE_ACCEPTED)
->json([
'path' => $deployment->getAttribute('buildPath'),
'size' => $deployment->getAttribute('buildSize', 0),
]);
return;
}
$stored = $cache->load($cacheKey, $cacheTtl);
if (\is_array($stored)) {
$metadata = $this->mergeArtifactUploadMetadata($stored, $metadata);
$chunks = (int) ($stored['chunksTotal'] ?? $chunks);
}
$metadata['chunksTotal'] = $chunks;
$chunksUploaded = \max((int) ($metadata['chunks'] ?? 0), $chunksUploaded);
if ($chunksUploaded === $chunks) {
$deviceForBuilds->finalizeUpload($path, $chunks, $metadata);
$size = $deviceForBuilds->getFileSize($path);
$deployment = $dbForProject->getAuthorization()->skip(fn () => $dbForProject->updateDocument('deployments', $deploymentId, new Document([
'buildPath' => $path,
'buildSize' => $size,
'totalSize' => $deployment->getAttribute('sourceSize', 0) + $size,
])));
$cache->purge($cacheKey);
$publisherForBuilds->enqueue(new BuildMessage(
project: $project,
resource: $resource,
deployment: $deployment,
type: BUILD_TYPE_ORCHESTRATOR_EVENT,
event: [
'type' => 'orchestrator.job.artifact',
'data' => [
'artifactId' => 'upload',
'status' => 'success',
],
],
));
$response
->setStatusCode(Response::STATUS_CODE_ACCEPTED)
->json([
'path' => $path,
'size' => $size,
]);
return;
}
$metadata['chunks'] = $chunksUploaded;
$cache->save($cacheKey, $metadata);
$response
->setStatusCode(Response::STATUS_CODE_ACCEPTED)
->json([
'path' => $path,
'size' => $fileSize,
'chunksTotal' => $chunks,
'chunksUploaded' => $chunksUploaded,
]);
}, timeout: 120.0);
} catch (LockContention) {
$response->addHeader('Retry-After', '5');
throw new Exception(Exception::GENERAL_RATE_LIMIT_EXCEEDED, 'Build artifact upload is busy. Try again.');
}
} finally {
@\unlink($tmp);
}
}
/**
* @param array<string, mixed> $stored
* @param array<string, mixed> $current
* @return array<string, mixed>
*/
private function mergeArtifactUploadMetadata(array $stored, array $current): array
{
$merged = \array_merge($stored, $current);
if (isset($stored['parts']) || isset($current['parts'])) {
$parts = $stored['parts'] ?? [];
foreach (($current['parts'] ?? []) as $part => $value) {
$parts[(int) $part] = $value;
}
\ksort($parts);
$merged['parts'] = $parts;
$merged['chunks'] = \count($parts);
}
return $merged;
}
}
@@ -3,10 +3,11 @@
namespace Appwrite\Platform\Modules\Functions\Http\Deployments\Artifacts\Build;
use Appwrite\Builds\OrchestratorToken;
use Appwrite\Event\Message\Build as BuildMessage;
use Appwrite\Event\Publisher\Build as BuildPublisher;
use Appwrite\Extend\Exception;
use Appwrite\Platform\Modules\Compute\Http\Deployments\Artifacts\Build\ChunkedBuildArtifact;
use Appwrite\Utopia\Response;
use Utopia\Cache\Cache;
use Utopia\Database\Database;
use Utopia\Database\Document;
use Utopia\Database\Validator\UID;
@@ -19,6 +20,7 @@ use Utopia\Validator\Text;
class Update extends Action
{
use HTTP;
use ChunkedBuildArtifact;
public static function getName()
{
@@ -43,6 +45,8 @@ class Update extends Action
->inject('project')
->inject('deviceForBuilds')
->inject('publisherForBuilds')
->inject('cache')
->inject('locks')
->callback($this->action(...));
}
@@ -55,7 +59,9 @@ class Update extends Action
Database $dbForProject,
Document $project,
Device $deviceForBuilds,
BuildPublisher $publisherForBuilds
BuildPublisher $publisherForBuilds,
Cache $cache,
callable $locks
) {
$token = $token ?: $request->getQuery('token', '');
OrchestratorToken::verify($token, $project->getId(), $functionId, $deploymentId, 'build');
@@ -70,48 +76,18 @@ class Update extends Action
throw new Exception(Exception::DEPLOYMENT_NOT_FOUND);
}
$tmp = \tempnam(\sys_get_temp_dir(), 'appwrite-build-');
if ($tmp === false) {
throw new Exception(Exception::GENERAL_SERVER_ERROR, 'Failed creating temporary build artifact file.');
}
\file_put_contents($tmp, $request->getRawPayload());
$metadata = ['content_type' => 'application/gzip'];
$path = $deviceForBuilds->getPath($deploymentId . '.tar.gz');
$uploaded = $deviceForBuilds->upload($tmp, $path, 1, 1, $metadata);
@\unlink($tmp);
if ($uploaded < 1) {
throw new Exception(Exception::GENERAL_SERVER_ERROR, 'Failed storing build artifact.');
}
$size = $deviceForBuilds->getFileSize($path);
$deployment = $dbForProject->getAuthorization()->skip(fn () => $dbForProject->updateDocument('deployments', $deploymentId, new Document([
'buildPath' => $path,
'buildSize' => $size,
'totalSize' => $deployment->getAttribute('sourceSize', 0) + $size,
])));
$publisherForBuilds->enqueue(new BuildMessage(
$this->uploadBuildArtifact(
deploymentId: $deploymentId,
project: $project,
resource: $function,
deployment: $deployment,
type: BUILD_TYPE_ORCHESTRATOR_EVENT,
event: [
'type' => 'orchestrator.job.artifact',
'data' => [
'artifactId' => 'upload',
'status' => 'success',
],
],
));
$response
->setStatusCode(Response::STATUS_CODE_ACCEPTED)
->json([
'path' => $path,
'size' => $size,
]);
request: $request,
response: $response,
dbForProject: $dbForProject,
deviceForBuilds: $deviceForBuilds,
publisherForBuilds: $publisherForBuilds,
cache: $cache,
locks: $locks
);
}
}
@@ -3,10 +3,11 @@
namespace Appwrite\Platform\Modules\Sites\Http\Deployments\Artifacts\Build;
use Appwrite\Builds\OrchestratorToken;
use Appwrite\Event\Message\Build as BuildMessage;
use Appwrite\Event\Publisher\Build as BuildPublisher;
use Appwrite\Extend\Exception;
use Appwrite\Platform\Modules\Compute\Http\Deployments\Artifacts\Build\ChunkedBuildArtifact;
use Appwrite\Utopia\Response;
use Utopia\Cache\Cache;
use Utopia\Database\Database;
use Utopia\Database\Document;
use Utopia\Database\Validator\UID;
@@ -19,6 +20,7 @@ use Utopia\Validator\Text;
class Update extends Action
{
use HTTP;
use ChunkedBuildArtifact;
public static function getName()
{
@@ -43,6 +45,8 @@ class Update extends Action
->inject('project')
->inject('deviceForBuilds')
->inject('publisherForBuilds')
->inject('cache')
->inject('locks')
->callback($this->action(...));
}
@@ -55,7 +59,9 @@ class Update extends Action
Database $dbForProject,
Document $project,
Device $deviceForBuilds,
BuildPublisher $publisherForBuilds
BuildPublisher $publisherForBuilds,
Cache $cache,
callable $locks
) {
$token = $token ?: $request->getQuery('token', '');
OrchestratorToken::verify($token, $project->getId(), $siteId, $deploymentId, 'build');
@@ -70,48 +76,18 @@ class Update extends Action
throw new Exception(Exception::DEPLOYMENT_NOT_FOUND);
}
$tmp = \tempnam(\sys_get_temp_dir(), 'appwrite-build-');
if ($tmp === false) {
throw new Exception(Exception::GENERAL_SERVER_ERROR, 'Failed creating temporary build artifact file.');
}
\file_put_contents($tmp, $request->getRawPayload());
$metadata = ['content_type' => 'application/gzip'];
$path = $deviceForBuilds->getPath($deploymentId . '.tar.gz');
$uploaded = $deviceForBuilds->upload($tmp, $path, 1, 1, $metadata);
@\unlink($tmp);
if ($uploaded < 1) {
throw new Exception(Exception::GENERAL_SERVER_ERROR, 'Failed storing build artifact.');
}
$size = $deviceForBuilds->getFileSize($path);
$deployment = $dbForProject->getAuthorization()->skip(fn () => $dbForProject->updateDocument('deployments', $deploymentId, new Document([
'buildPath' => $path,
'buildSize' => $size,
'totalSize' => $deployment->getAttribute('sourceSize', 0) + $size,
])));
$publisherForBuilds->enqueue(new BuildMessage(
$this->uploadBuildArtifact(
deploymentId: $deploymentId,
project: $project,
resource: $site,
deployment: $deployment,
type: BUILD_TYPE_ORCHESTRATOR_EVENT,
event: [
'type' => 'orchestrator.job.artifact',
'data' => [
'artifactId' => 'upload',
'status' => 'success',
],
],
));
$response
->setStatusCode(Response::STATUS_CODE_ACCEPTED)
->json([
'path' => $path,
'size' => $size,
]);
request: $request,
response: $response,
dbForProject: $dbForProject,
deviceForBuilds: $deviceForBuilds,
publisherForBuilds: $publisherForBuilds,
cache: $cache,
locks: $locks
);
}
}